This commit is contained in:
Your Name
2026-08-12 17:18:56 +08:00
parent 28cd110dae
commit f48a66b611
24 changed files with 3095 additions and 233 deletions
@@ -130,6 +130,7 @@ def _ticket_mapping(ticket: Any) -> Mapping[str, Any]:
"patientUserId": "patient_user_id",
"diagnosisId": "diagnosis_id",
"patientId": "patient_id",
"callRecordId": "call_record_id",
}
adapted = {
json_name: getattr(ticket, attribute_name)
@@ -196,6 +197,7 @@ class VideoCallRequest:
target_user_id: str
diagnosis_id: Identifier
patient_id: Identifier | None = None
call_record_id: Identifier | None = None
backend_mode: BackendMode = BackendMode.EMBEDDED
def __post_init__(self) -> None:
@@ -218,6 +220,12 @@ class VideoCallRequest:
"patient_id",
_identifier(self.patient_id, "patientId"),
)
if self.call_record_id is not None:
object.__setattr__(
self,
"call_record_id",
_identifier(self.call_record_id, "callRecordId"),
)
object.__setattr__(self, "backend_mode", BackendMode.parse(self.backend_mode))
@classmethod
@@ -257,6 +265,7 @@ class VideoCallRequest:
return {
"diagnosis_id": self.diagnosis_id,
"patient_id": self.patient_id,
"call_record_id": self.call_record_id,
"backend_mode": self.backend_mode.value,
}
@@ -285,6 +294,13 @@ def normalize_backend_ticket(
_identifier,
required=False,
)
payload_call_record = _read_aliases(
payload,
("callRecordId", "call_record_id"),
"callRecordId",
_identifier,
required=False,
)
normalized_diagnosis = _merge_identifier(
payload_diagnosis,
@@ -323,6 +339,7 @@ def normalize_backend_ticket(
),
diagnosis_id=normalized_diagnosis,
patient_id=normalized_patient,
call_record_id=payload_call_record,
backend_mode=BackendMode.parse(backend_mode),
)
+237 -1
View File
@@ -74,6 +74,71 @@ def _call_repository_method(method: Callable[..., Any], payload: Mapping[str, An
return _resolve_result(result)
def _mapping_candidate(value: Any) -> Mapping[str, Any] | None:
if isinstance(value, Mapping):
return value
for attribute in ("raw", "data"):
candidate = getattr(value, attribute, None)
if isinstance(candidate, Mapping):
return candidate
return None
def _extract_call_record_id(result: Any) -> int | str | None:
"""Read a positive call-record identity from common backend envelopes."""
pending = [result]
visited: set[int] = set()
while pending:
candidate = pending.pop(0)
mapping = _mapping_candidate(candidate)
if mapping is None or id(mapping) in visited:
continue
visited.add(id(mapping))
for key in ("call_record_id", "callRecordId", "id"):
value = mapping.get(key)
if isinstance(value, bool) or value is None:
continue
if isinstance(value, int) and value > 0:
return value
if isinstance(value, str) and value.strip():
return value.strip()
for key in ("data", "result", "record", "call_record", "callRecord"):
nested = mapping.get(key)
if isinstance(nested, Mapping):
pending.append(nested)
return None
def _clean_transcript_segment(segment: Mapping[str, Any], session_id: str) -> dict[str, Any]:
segment_id = str(segment.get("segment_id", segment.get("segmentId", ""))).strip()
text = str(segment.get("text", segment.get("sourceText", ""))).strip()
speaker_user_id = str(
segment.get("speaker_user_id", segment.get("speakerUserId", ""))
).strip()
speaker_role = str(segment.get("speaker_role", segment.get("speakerRole", "unknown"))).strip()
if not segment_id or len(segment_id) > 160:
raise ValueError("transcript segment_id must contain 1 to 160 characters")
if not text or len(text) > 4_000:
raise ValueError("transcript text must contain 1 to 4000 characters")
if len(speaker_user_id) > 160:
raise ValueError("transcript speaker_user_id is too long")
if speaker_role not in {"doctor", "patient", "unknown"}:
speaker_role = "unknown"
try:
timestamp = int(segment.get("timestamp") or 0)
except (TypeError, ValueError):
timestamp = 0
return {
"segment_id": segment_id,
"transcription_session_id": session_id,
"speaker_user_id": speaker_user_id,
"speaker_role": speaker_role,
"timestamp": max(timestamp, 0),
"text": text,
}
@dataclass(slots=True)
class _WorkItem:
operation: str
@@ -169,11 +234,17 @@ class OrderedCallLifecycle:
self.logger = logger
self.started = False
self.ended = False
self.call_record_id: int | str | None = request.call_record_id
self.bound_room_id: str | None = None
self._claimed_room_id: str | None = None
self._start_future: Future[bool] | None = None
self._bind_future: Future[bool] | None = None
self._end_future: Future[bool] | None = None
self._transcription_start_future: Future[bool] | None = None
self._transcription_finish_future: Future[bool] | None = None
self._transcription_session_id: str | None = None
self._transcription_active = False
self._segment_futures: dict[str, Future[bool]] = {}
self._lock = threading.RLock()
self._worker = _OrderedDaemonWorker(logger, request.safe_log_context())
@@ -181,6 +252,10 @@ class OrderedCallLifecycle:
def worker_is_daemon(self) -> bool:
return self._worker.is_daemon
@property
def transcription_session_id(self) -> str | None:
return self._transcription_session_id
def start(self) -> Future[bool]:
with self._lock:
if self._start_future is not None:
@@ -196,9 +271,13 @@ class OrderedCallLifecycle:
payload["patient_id"] = self.request.patient_id
def operation() -> bool:
_call_repository_method(method, payload)
result = _call_repository_method(method, payload)
record_id = _extract_call_record_id(result)
if record_id is None:
raise ValueError("startCall response did not include the current call_record_id")
with self._lock:
self.started = True
self.call_record_id = record_id
self.logger.info(
"video call record started",
extra={"video_call": self.request.safe_log_context()},
@@ -309,10 +388,167 @@ class OrderedCallLifecycle:
return self._worker.submit("screenshot", operation)
def start_transcription(self, session_id: str, *, language: str = "zh-CN") -> Future[bool]:
"""Start persisted realtime transcription for this exact call record."""
clean_session = str(session_id or "").strip()
clean_language = str(language or "zh-CN").strip()[:32] or "zh-CN"
if not clean_session or len(clean_session) > 128:
raise ValueError("transcription session id must contain 1 to 128 characters")
with self._lock:
if self._end_future is not None:
raise RuntimeError("video call has already ended")
if self._transcription_start_future is not None:
if self._transcription_session_id != clean_session:
raise RuntimeError("another transcription session already exists")
return self._transcription_start_future
if self._start_future is None:
self.start()
method = getattr(self.repository, "start_call_transcription", None)
if not callable(method):
raise ValueError("video repository does not implement call transcription storage")
self._transcription_session_id = clean_session
def operation() -> bool:
with self._lock:
started = self.started
record_id = self.call_record_id
if not started:
return False
if record_id is None:
raise ValueError("server did not return the current call_record_id")
_call_repository_method(
method,
{
"diagnosis_id": self.request.diagnosis_id,
"call_record_id": record_id,
"transcription_session_id": clean_session,
"language": clean_language,
},
)
with self._lock:
self._transcription_active = True
self.logger.info(
"video call transcription started",
extra={"video_call": self.request.safe_log_context()},
)
return True
self._transcription_start_future = self._worker.submit(
"transcription-start", operation
)
return self._transcription_start_future
def save_transcript_segment(self, segment: Mapping[str, Any]) -> Future[bool]:
"""Upsert one completed, bounded transcript segment without logging its text."""
with self._lock:
session_id = self._transcription_session_id
if not session_id or self._transcription_start_future is None:
raise RuntimeError("call transcription has not started")
if self._transcription_finish_future is not None:
raise RuntimeError("call transcription has already stopped")
cleaned = _clean_transcript_segment(segment, session_id)
segment_id = cleaned["segment_id"]
existing = self._segment_futures.get(segment_id)
if existing is not None:
return existing
method = getattr(self.repository, "upsert_call_transcript_segments", None)
if not callable(method):
raise ValueError("video repository does not implement transcript segment storage")
def operation() -> bool:
with self._lock:
active = self._transcription_active
record_id = self.call_record_id
if not active:
return False
if record_id is None:
raise ValueError("server did not return the current call_record_id")
_call_repository_method(
method,
{
"diagnosis_id": self.request.diagnosis_id,
"call_record_id": record_id,
"transcription_session_id": session_id,
"segments": [cleaned],
},
)
return True
future = self._worker.submit("transcription-segment", operation)
self._segment_futures[segment_id] = future
def release_failed(completed: Future[bool]) -> None:
if completed.cancelled() or completed.exception() is not None:
with self._lock:
if self._segment_futures.get(segment_id) is completed:
self._segment_futures.pop(segment_id, None)
future.add_done_callback(release_failed)
return future
def finish_transcription(self, *, status: str = "completed") -> Future[bool]:
"""Flush and finalize the current transcript exactly once."""
clean_status = str(status or "completed").strip().lower()
if clean_status not in {"completed", "partial", "failed"}:
raise ValueError("transcription status is invalid")
with self._lock:
if self._transcription_finish_future is not None:
return self._transcription_finish_future
if self._transcription_start_future is None or not self._transcription_session_id:
return _settled_future(False)
method = getattr(self.repository, "finish_call_transcription", None)
if not callable(method):
raise ValueError("video repository does not implement transcription finalization")
session_id = self._transcription_session_id
def operation() -> bool:
with self._lock:
active = self._transcription_active
record_id = self.call_record_id
expected_count = len(self._segment_futures)
if not active:
return False
if record_id is None:
raise ValueError("server did not return the current call_record_id")
_call_repository_method(
method,
{
"diagnosis_id": self.request.diagnosis_id,
"call_record_id": record_id,
"transcription_session_id": session_id,
"expected_segment_count": expected_count,
"status": clean_status,
},
)
with self._lock:
self._transcription_active = False
self.logger.info(
"video call transcription finalized",
extra={
"video_call": self.request.safe_log_context(),
"transcription_status": clean_status,
"segment_count": expected_count,
},
)
return True
self._transcription_finish_future = self._worker.submit(
"transcription-finish", operation
)
return self._transcription_finish_future
def end(self, reason: str) -> Future[bool]:
with self._lock:
if self._end_future is not None:
return self._end_future
if (
self._transcription_start_future is not None
and self._transcription_finish_future is None
):
self.finish_transcription(status="partial")
method = getattr(self.repository, "end_call", None)
def operation() -> bool:
+120 -2
View File
@@ -210,6 +210,7 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
call_error = Signal(str) # type: ignore[misc]
_start_completed = Signal(bool) # type: ignore[misc]
_screenshot_completed = Signal(bool, str) # type: ignore[misc]
_transcription_completed = Signal(str, str, str, bool, str) # type: ignore[misc]
def __init__(
self,
@@ -284,6 +285,7 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
self._start_completed.connect(self._on_lifecycle_started)
self._screenshot_completed.connect(self._on_screenshot_completed)
self._transcription_completed.connect(self._on_transcription_completed)
self.web_view.loadFinished.connect(self._on_load_finished)
self.web_view.setUrl(QUrl(self.location.url))
@@ -426,6 +428,21 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
str(message.get("message") or "截屏图片无效。")[:200],
)
return
if event == "transcription-start-request":
self._start_transcription(
str(message.get("sessionId") or ""),
str(message.get("language") or "zh-CN"),
)
return
if event == "transcription-segment":
self._save_transcript_segment(message)
return
if event == "transcription-stop":
self._finish_transcription(
str(message.get("sessionId") or ""),
str(message.get("status") or "completed"),
)
return
room_id = message.get("roomId", message.get("room_id"))
if room_id not in (None, ""):
self.lifecycle.bind_room(room_id)
@@ -434,8 +451,6 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
if event == "status":
status = str(message.get("status", "unknown"))[:80]
self.status_changed.emit(status)
if status == "idle" and not self.open_im:
self._close_from_companion("remote-idle")
elif event == "hangup":
status = str(message.get("status", "ended"))[:80]
self.call_ended.emit(status)
@@ -509,6 +524,109 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
f"window.doctorConsultation?.screenshotResult?.({state}, {payload});"
)
def _notify_transcription_completed(
self,
operation: str,
session_id: str,
segment_id: str,
future: Future[bool],
) -> None:
try:
succeeded = bool(future.result())
except Exception as error:
succeeded = False
message = str(error)[:200] or "录音文字保存失败。"
else:
message = {
"start": "录音文字存储已准备。",
"segment": "",
"stop": "本次面诊对话文字已保存。",
}.get(operation, "")
with suppress(RuntimeError):
self._transcription_completed.emit(
operation, session_id, segment_id, succeeded, message
)
def _on_transcription_completed(
self,
operation: str,
session_id: str,
segment_id: str,
succeeded: bool,
message: str,
) -> None:
payload = json.dumps(str(message)[:200], ensure_ascii=True)
state = "true" if succeeded else "false"
self._page.runJavaScript(
"window.doctorConsultation?.transcriptionResult?.("
f"{json.dumps(operation)}, {json.dumps(session_id)}, "
f"{json.dumps(segment_id)}, {state}, {payload});"
)
def _start_transcription(self, session_id: str, language: str) -> None:
try:
future = self.lifecycle.start_transcription(session_id, language=language)
except Exception as error:
self._on_transcription_completed(
"start", session_id, "", False, str(error)[:200]
)
return
future.add_done_callback(
lambda completed: self._notify_transcription_completed(
"start", session_id, "", completed
)
)
def _save_transcript_segment(self, message: Mapping[str, Any]) -> None:
session_id = str(message.get("sessionId") or "").strip()
segment = message.get("segment")
segment_id = (
str(segment.get("segment_id") or segment.get("segmentId") or "").strip()
if isinstance(segment, Mapping)
else ""
)
if not isinstance(segment, Mapping):
self._on_transcription_completed(
"segment", session_id, segment_id, False, "录音文字片段无效。"
)
return
if session_id != str(self.lifecycle.transcription_session_id or ""):
self._on_transcription_completed(
"segment", session_id, segment_id, False, "录音会话标识不匹配。"
)
return
try:
future = self.lifecycle.save_transcript_segment(segment)
except Exception as error:
self._on_transcription_completed(
"segment", session_id, segment_id, False, str(error)[:200]
)
return
future.add_done_callback(
lambda completed: self._notify_transcription_completed(
"segment", session_id, segment_id, completed
)
)
def _finish_transcription(self, session_id: str, status: str) -> None:
if session_id.strip() != str(self.lifecycle.transcription_session_id or ""):
self._on_transcription_completed(
"stop", session_id, "", False, "录音会话标识不匹配。"
)
return
try:
future = self.lifecycle.finish_transcription(status=status)
except Exception as error:
self._on_transcription_completed(
"stop", session_id, "", False, str(error)[:200]
)
return
future.add_done_callback(
lambda completed: self._notify_transcription_completed(
"stop", session_id, "", completed
)
)
def _close_from_companion(self, reason: str) -> None:
self._companion_ended = True
self._close_reason = reason