更新bug
This commit is contained in:
@@ -15,6 +15,7 @@ import threading
|
||||
from collections.abc import Callable, Mapping
|
||||
from concurrent.futures import Future
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any, TypeVar
|
||||
|
||||
from .launcher import VideoCallRequest
|
||||
@@ -110,6 +111,29 @@ def _extract_call_record_id(result: Any) -> int | str | None:
|
||||
return None
|
||||
|
||||
|
||||
def _extract_cloud_recording_outcome(result: Any) -> tuple[bool, str] | None:
|
||||
"""Read the cloud mixed-video recording outcome returned by room binding."""
|
||||
|
||||
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))
|
||||
recording = mapping.get("cloud_recording", mapping.get("cloudRecording"))
|
||||
if isinstance(recording, Mapping):
|
||||
started = bool(recording.get("started"))
|
||||
message = str(recording.get("message") or "").strip()[:200]
|
||||
return started, message
|
||||
for key in ("data", "result"):
|
||||
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()
|
||||
@@ -239,6 +263,7 @@ class OrderedCallLifecycle:
|
||||
self._claimed_room_id: str | None = None
|
||||
self._start_future: Future[bool] | None = None
|
||||
self._bind_future: Future[bool] | None = None
|
||||
self._local_recording_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
|
||||
@@ -256,6 +281,13 @@ class OrderedCallLifecycle:
|
||||
def transcription_session_id(self) -> str | None:
|
||||
return self._transcription_session_id
|
||||
|
||||
@property
|
||||
def current_room_id(self) -> str | None:
|
||||
"""Return the room observed for this call cycle, including an active bind."""
|
||||
|
||||
with self._lock:
|
||||
return self.bound_room_id or self._claimed_room_id
|
||||
|
||||
def start(self) -> Future[bool]:
|
||||
with self._lock:
|
||||
if self._start_future is not None:
|
||||
@@ -318,18 +350,31 @@ class OrderedCallLifecycle:
|
||||
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")
|
||||
if not callable(method):
|
||||
self.logger.warning(
|
||||
"video repository does not implement bind_call_room",
|
||||
extra={"video_call": self.request.safe_log_context()},
|
||||
)
|
||||
return False
|
||||
_call_repository_method(
|
||||
result = _call_repository_method(
|
||||
method,
|
||||
{"diagnosis_id": self.request.diagnosis_id, "room_id": cleaned},
|
||||
{
|
||||
"diagnosis_id": self.request.diagnosis_id,
|
||||
"room_id": cleaned,
|
||||
"call_record_id": record_id,
|
||||
},
|
||||
)
|
||||
recording_outcome = _extract_cloud_recording_outcome(result)
|
||||
if recording_outcome is not None and not recording_outcome[0]:
|
||||
raise RuntimeError(
|
||||
recording_outcome[1]
|
||||
or "automatic Tencent cloud mixed-video recording did not start"
|
||||
)
|
||||
with self._lock:
|
||||
self.bound_room_id = cleaned
|
||||
self.logger.info(
|
||||
@@ -338,8 +383,85 @@ class OrderedCallLifecycle:
|
||||
)
|
||||
return True
|
||||
|
||||
self._bind_future = self._worker.submit("bind", operation)
|
||||
return self._bind_future
|
||||
future = self._worker.submit("bind", operation)
|
||||
self._bind_future = future
|
||||
|
||||
def release_failed_claim(completed: Future[bool]) -> None:
|
||||
try:
|
||||
succeeded = bool(completed.result())
|
||||
except Exception:
|
||||
succeeded = False
|
||||
if succeeded:
|
||||
return
|
||||
# A transient bridge/API failure must not permanently pin the
|
||||
# room to a failed Future. The companion retries the exact
|
||||
# same room after the host acknowledgement, so release only
|
||||
# this failed claim while preserving successful bindings.
|
||||
with self._lock:
|
||||
if self._bind_future is completed and self.bound_room_id is None:
|
||||
self._bind_future = None
|
||||
self._claimed_room_id = None
|
||||
|
||||
future.add_done_callback(release_failed_claim)
|
||||
return future
|
||||
|
||||
def save_local_audio_recording(
|
||||
self,
|
||||
path: str | Path,
|
||||
*,
|
||||
mime_type: str = "audio/webm",
|
||||
) -> Future[bool]:
|
||||
"""Upload one locally mixed audio file to this exact call record.
|
||||
|
||||
The operation shares the lifecycle FIFO, so a caller that queues this
|
||||
before :meth:`end` is guaranteed to attach the COS object before the
|
||||
call record is marked ended.
|
||||
"""
|
||||
|
||||
recording_path = Path(path).expanduser().resolve()
|
||||
clean_mime = str(mime_type or "audio/webm").strip()[:120] or "audio/webm"
|
||||
if not recording_path.is_file() or recording_path.stat().st_size <= 0:
|
||||
raise ValueError("local audio recording is empty or missing")
|
||||
with self._lock:
|
||||
if self._end_future is not None:
|
||||
raise RuntimeError("video call has already ended")
|
||||
if self._local_recording_future is not None:
|
||||
return self._local_recording_future
|
||||
if self._start_future is None:
|
||||
self.start()
|
||||
method = getattr(self.repository, "upload_call_recording", None)
|
||||
if not callable(method):
|
||||
raise ValueError("video repository does not implement local recording upload")
|
||||
|
||||
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")
|
||||
result = _call_repository_method(
|
||||
method,
|
||||
{
|
||||
"path": recording_path,
|
||||
"diagnosis_id": self.request.diagnosis_id,
|
||||
"call_record_id": record_id,
|
||||
"mime_type": clean_mime,
|
||||
},
|
||||
)
|
||||
if result is False:
|
||||
raise RuntimeError("local audio recording upload failed")
|
||||
self.logger.info(
|
||||
"local call audio uploaded and attached",
|
||||
extra={"video_call": self.request.safe_log_context()},
|
||||
)
|
||||
return True
|
||||
|
||||
self._local_recording_future = self._worker.submit(
|
||||
"local-audio-upload", operation
|
||||
)
|
||||
return self._local_recording_future
|
||||
|
||||
def save_screenshot(self, content: bytes, filename: str) -> Future[str]:
|
||||
"""Upload one video frame and append it to the diagnosis doctor notes."""
|
||||
@@ -554,15 +676,21 @@ class OrderedCallLifecycle:
|
||||
def operation() -> bool:
|
||||
with self._lock:
|
||||
started = self.started
|
||||
record_id = self.call_record_id
|
||||
if not started:
|
||||
with self._lock:
|
||||
self.ended = True
|
||||
return False
|
||||
if not callable(method):
|
||||
raise ValueError("video repository does not implement end_call")
|
||||
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},
|
||||
{
|
||||
"diagnosis_id": self.request.diagnosis_id,
|
||||
"call_record_id": record_id,
|
||||
},
|
||||
)
|
||||
with self._lock:
|
||||
self.ended = True
|
||||
|
||||
@@ -11,6 +11,9 @@ import base64
|
||||
import binascii
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import sqlite3
|
||||
import sys
|
||||
from collections.abc import Callable, Mapping
|
||||
from concurrent.futures import Future
|
||||
@@ -20,6 +23,11 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
from urllib.parse import parse_qsl, urlsplit
|
||||
|
||||
from ..services.local_audio_queue import (
|
||||
LocalAudioQueueStore,
|
||||
LocalAudioUploadManager,
|
||||
get_local_audio_upload_manager,
|
||||
)
|
||||
from .launcher import (
|
||||
VideoCallRequest,
|
||||
VideoTicketError,
|
||||
@@ -29,7 +37,7 @@ from .lifecycle import OrderedCallLifecycle
|
||||
from .security import TrustedDocumentError, TrustedDocumentPolicy
|
||||
|
||||
try: # Optional by design: core-only builds must still import this module.
|
||||
from PySide6.QtCore import QObject, Qt, QUrl, Signal, Slot
|
||||
from PySide6.QtCore import QObject, Qt, QTimer, QUrl, Signal, Slot
|
||||
from PySide6.QtWebChannel import QWebChannel
|
||||
from PySide6.QtWebEngineCore import (
|
||||
QWebEnginePage,
|
||||
@@ -39,7 +47,7 @@ try: # Optional by design: core-only builds must still import this module.
|
||||
from PySide6.QtWebEngineWidgets import QWebEngineView
|
||||
from PySide6.QtWidgets import QApplication, QMainWindow
|
||||
except (ImportError, OSError) as _qt_import_error: # pragma: no cover - no Qt runtime.
|
||||
QObject = Qt = QUrl = Signal = Slot = None # type: ignore[assignment]
|
||||
QObject = QTimer = Qt = QUrl = Signal = Slot = None # type: ignore[assignment]
|
||||
QWebChannel = QWebEnginePage = QWebEngineProfile = None # type: ignore[assignment]
|
||||
QWebEngineSettings = QWebEngineView = None # type: ignore[assignment]
|
||||
QApplication = QMainWindow = None # type: ignore[assignment]
|
||||
@@ -63,6 +71,18 @@ class CompanionLocation:
|
||||
is_local: bool
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class _LocalAudioCapture:
|
||||
record_id: int
|
||||
session_id: str
|
||||
mime_type: str
|
||||
path: Path
|
||||
handle: Any
|
||||
lifecycle: OrderedCallLifecycle
|
||||
next_sequence: int = 0
|
||||
bytes_written: int = 0
|
||||
|
||||
|
||||
def _validate_remote_url(value: str) -> str:
|
||||
parsed = urlsplit(value)
|
||||
if parsed.scheme.lower() != "https" or not parsed.hostname:
|
||||
@@ -204,13 +224,67 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
}
|
||||
)
|
||||
|
||||
@Slot(str, str) # type: ignore[misc]
|
||||
def startLocalAudioRecording( # noqa: N802 - Qt bridge API
|
||||
self, session_id: str, mime_type: str
|
||||
) -> None:
|
||||
self._callback(
|
||||
{
|
||||
"source": "doctor-call",
|
||||
"event": "local-audio-start",
|
||||
"sessionId": session_id,
|
||||
"mimeType": mime_type,
|
||||
}
|
||||
)
|
||||
|
||||
@Slot(str, int, str) # type: ignore[misc]
|
||||
def appendLocalAudioChunk( # noqa: N802 - Qt bridge API
|
||||
self, session_id: str, sequence: int, encoded: str
|
||||
) -> None:
|
||||
self._callback(
|
||||
{
|
||||
"source": "doctor-call",
|
||||
"event": "local-audio-chunk",
|
||||
"sessionId": session_id,
|
||||
"sequence": sequence,
|
||||
"data": encoded,
|
||||
}
|
||||
)
|
||||
|
||||
@Slot(str, int) # type: ignore[misc]
|
||||
def finishLocalAudioRecording( # noqa: N802 - Qt bridge API
|
||||
self, session_id: str, total_bytes: int
|
||||
) -> None:
|
||||
self._callback(
|
||||
{
|
||||
"source": "doctor-call",
|
||||
"event": "local-audio-finish",
|
||||
"sessionId": session_id,
|
||||
"totalBytes": total_bytes,
|
||||
}
|
||||
)
|
||||
|
||||
@Slot(str) # type: ignore[misc]
|
||||
def abortLocalAudioRecording(self, session_id: str) -> None: # noqa: N802
|
||||
self._callback(
|
||||
{
|
||||
"source": "doctor-call",
|
||||
"event": "local-audio-abort",
|
||||
"sessionId": session_id,
|
||||
}
|
||||
)
|
||||
|
||||
class _EmbeddedVideoWindow(QMainWindow): # type: ignore[misc, valid-type]
|
||||
status_changed = Signal(str) # type: ignore[misc]
|
||||
call_ended = Signal(str) # type: ignore[misc]
|
||||
call_error = Signal(str) # type: ignore[misc]
|
||||
_start_completed = Signal(bool) # type: ignore[misc]
|
||||
_room_completed = Signal(str, bool, str) # type: ignore[misc]
|
||||
_screenshot_completed = Signal(bool, str) # type: ignore[misc]
|
||||
_transcription_completed = Signal(str, str, str, bool, str) # type: ignore[misc]
|
||||
_local_recording_completed = Signal( # type: ignore[misc]
|
||||
str, str, int, bool, str
|
||||
)
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
@@ -247,8 +321,15 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self._close_reason = "window-closed"
|
||||
self._call_cycle_closed = False
|
||||
self._start_requested = False
|
||||
self._shutdown_requested = False
|
||||
self._local_audio_capture: _LocalAudioCapture | None = None
|
||||
self._local_audio_store: LocalAudioQueueStore | None = None
|
||||
self._local_audio_uploads: LocalAudioUploadManager | None = None
|
||||
self._legacy_grants: list[tuple[Any, Any]] = []
|
||||
self._permission_grants: list[Any] = []
|
||||
self._shutdown_timer = QTimer(self)
|
||||
self._shutdown_timer.setSingleShot(True)
|
||||
self._shutdown_timer.timeout.connect(self._force_requested_shutdown)
|
||||
|
||||
self.setWindowTitle(
|
||||
f"与 {self.patient_name} IM 问诊" if self.open_im else "视频面诊"
|
||||
@@ -284,8 +365,10 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self._connect_permissions()
|
||||
|
||||
self._start_completed.connect(self._on_lifecycle_started)
|
||||
self._room_completed.connect(self._on_room_completed)
|
||||
self._screenshot_completed.connect(self._on_screenshot_completed)
|
||||
self._transcription_completed.connect(self._on_transcription_completed)
|
||||
self._local_recording_completed.connect(self._on_local_recording_completed)
|
||||
self.web_view.loadFinished.connect(self._on_load_finished)
|
||||
self.web_view.setUrl(QUrl(self.location.url))
|
||||
|
||||
@@ -428,6 +511,28 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
str(message.get("message") or "截屏图片无效。")[:200],
|
||||
)
|
||||
return
|
||||
if event == "local-audio-start":
|
||||
self._start_local_audio_recording(
|
||||
str(message.get("sessionId") or ""),
|
||||
str(message.get("mimeType") or "audio/webm"),
|
||||
)
|
||||
return
|
||||
if event == "local-audio-chunk":
|
||||
self._append_local_audio_chunk(
|
||||
str(message.get("sessionId") or ""),
|
||||
message.get("sequence"),
|
||||
str(message.get("data") or ""),
|
||||
)
|
||||
return
|
||||
if event == "local-audio-finish":
|
||||
self._finish_local_audio_recording(
|
||||
str(message.get("sessionId") or ""),
|
||||
message.get("totalBytes"),
|
||||
)
|
||||
return
|
||||
if event == "local-audio-abort":
|
||||
self._abort_local_audio_recording(str(message.get("sessionId") or ""))
|
||||
return
|
||||
if event == "transcription-start-request":
|
||||
self._start_transcription(
|
||||
str(message.get("sessionId") or ""),
|
||||
@@ -445,7 +550,13 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
return
|
||||
room_id = message.get("roomId", message.get("room_id"))
|
||||
if room_id not in (None, ""):
|
||||
self.lifecycle.bind_room(room_id)
|
||||
clean_room_id = str(room_id).strip()
|
||||
future = self.lifecycle.bind_room(clean_room_id)
|
||||
future.add_done_callback(
|
||||
lambda completed, current_room_id=clean_room_id: (
|
||||
self._notify_room_completed(current_room_id, completed)
|
||||
)
|
||||
)
|
||||
if event == "room":
|
||||
return
|
||||
if event == "status":
|
||||
@@ -457,7 +568,7 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self.lifecycle.end(f"companion-{status}")
|
||||
self._call_cycle_closed = True
|
||||
self._start_requested = False
|
||||
if not self.open_im:
|
||||
if not self.open_im or self._shutdown_requested:
|
||||
self._close_from_companion("companion-hangup")
|
||||
elif event == "error":
|
||||
message_text = str(message.get("message", "视频通话错误"))[:400]
|
||||
@@ -466,9 +577,343 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self.lifecycle.end("companion-error")
|
||||
self._call_cycle_closed = True
|
||||
self._start_requested = False
|
||||
if not self.open_im:
|
||||
if not self.open_im or self._shutdown_requested:
|
||||
self._close_from_companion("companion-error")
|
||||
|
||||
def _notify_room_completed(self, room_id: str, future: Future[bool]) -> None:
|
||||
try:
|
||||
succeeded = bool(future.result())
|
||||
except Exception as error:
|
||||
succeeded = False
|
||||
message = str(error)[:200] or "腾讯云混流视频录制未启动。"
|
||||
else:
|
||||
message = (
|
||||
"腾讯云混流视频已启动;本机录音将在结束后另行上传 COS。"
|
||||
if succeeded
|
||||
else "通话房间尚未绑定,云端视频和本机录音无法关联通话记录。"
|
||||
)
|
||||
with suppress(RuntimeError):
|
||||
self._room_completed.emit(room_id, succeeded, message)
|
||||
|
||||
def _on_room_completed(
|
||||
self,
|
||||
room_id: str,
|
||||
succeeded: bool,
|
||||
message: str,
|
||||
) -> None:
|
||||
if self._closing:
|
||||
return
|
||||
room_payload = json.dumps(str(room_id)[:160], ensure_ascii=True)
|
||||
payload = json.dumps(str(message)[:200], ensure_ascii=True)
|
||||
state = "true" if succeeded else "false"
|
||||
self._page.runJavaScript(
|
||||
"window.doctorConsultation?.roomBindingResult?.("
|
||||
f"{room_payload}, {state}, {payload});"
|
||||
)
|
||||
|
||||
def _emit_local_recording_result(
|
||||
self,
|
||||
operation: str,
|
||||
session_id: str,
|
||||
sequence: int,
|
||||
succeeded: bool,
|
||||
message: str,
|
||||
) -> None:
|
||||
with suppress(RuntimeError):
|
||||
self._local_recording_completed.emit(
|
||||
operation,
|
||||
session_id,
|
||||
sequence,
|
||||
succeeded,
|
||||
str(message)[:200],
|
||||
)
|
||||
|
||||
def _on_local_recording_completed(
|
||||
self,
|
||||
operation: str,
|
||||
session_id: str,
|
||||
sequence: int,
|
||||
succeeded: bool,
|
||||
message: str,
|
||||
) -> None:
|
||||
if self._closing or self._released:
|
||||
return
|
||||
self._page.runJavaScript(
|
||||
"window.doctorConsultation?.localRecordingResult?.("
|
||||
f"{json.dumps(operation)}, {json.dumps(session_id)}, {sequence}, "
|
||||
f"{'true' if succeeded else 'false'}, "
|
||||
f"{json.dumps(str(message)[:200], ensure_ascii=True)});"
|
||||
)
|
||||
|
||||
def _start_local_audio_recording(self, session_id: str, mime_type: str) -> None:
|
||||
cleaned = str(session_id or "").strip()
|
||||
clean_mime = str(mime_type or "audio/webm").strip().lower()[:120]
|
||||
if not re.fullmatch(r"[A-Za-z0-9_-]{12,64}", cleaned):
|
||||
self._emit_local_recording_result(
|
||||
"start", cleaned, -1, False, "本地录音会话标识无效。"
|
||||
)
|
||||
return
|
||||
if not clean_mime.startswith(("audio/webm", "audio/ogg")):
|
||||
self._emit_local_recording_result(
|
||||
"start", cleaned, -1, False, "当前浏览器录音格式不受支持。"
|
||||
)
|
||||
return
|
||||
if self._local_audio_capture is not None:
|
||||
existing = self._local_audio_capture.session_id == cleaned
|
||||
self._emit_local_recording_result(
|
||||
"start",
|
||||
cleaned,
|
||||
-1,
|
||||
existing,
|
||||
"本地录音已启动。" if existing else "已有另一条本地录音正在进行。",
|
||||
)
|
||||
return
|
||||
try:
|
||||
lifecycle = self.lifecycle
|
||||
raw_call_record_id = lifecycle.call_record_id
|
||||
call_record_id = (
|
||||
int(raw_call_record_id) if raw_call_record_id not in (None, "") else None
|
||||
)
|
||||
store, uploads = get_local_audio_upload_manager(
|
||||
lifecycle.repository,
|
||||
self._local_audio_store,
|
||||
)
|
||||
self._local_audio_store = store
|
||||
self._local_audio_uploads = uploads
|
||||
record = store.begin_recording(
|
||||
session_id=cleaned,
|
||||
diagnosis_id=self.request.diagnosis_id,
|
||||
mime_type=clean_mime,
|
||||
call_record_id=call_record_id,
|
||||
room_id=lifecycle.current_room_id or "",
|
||||
)
|
||||
# The handle intentionally remains open across WebChannel chunks.
|
||||
handle = record.file_path.open("w+b")
|
||||
except (OSError, RuntimeError, sqlite3.Error) as error:
|
||||
self._emit_local_recording_result(
|
||||
"start", cleaned, -1, False, str(error)[:200]
|
||||
)
|
||||
return
|
||||
self._local_audio_capture = _LocalAudioCapture(
|
||||
record_id=record.id,
|
||||
session_id=cleaned,
|
||||
mime_type=clean_mime,
|
||||
path=record.file_path,
|
||||
handle=handle,
|
||||
lifecycle=lifecycle,
|
||||
)
|
||||
self._emit_local_recording_result(
|
||||
"start", cleaned, -1, True, "本机语音录音已启动。"
|
||||
)
|
||||
|
||||
def _append_local_audio_chunk(
|
||||
self,
|
||||
session_id: str,
|
||||
sequence_value: Any,
|
||||
encoded: str,
|
||||
) -> None:
|
||||
capture = self._local_audio_capture
|
||||
try:
|
||||
sequence = int(sequence_value)
|
||||
except (TypeError, ValueError):
|
||||
sequence = -1
|
||||
if capture is None or session_id != capture.session_id:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音会话标识不匹配。"
|
||||
)
|
||||
return
|
||||
if sequence != capture.next_sequence:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音分片顺序不连续。"
|
||||
)
|
||||
return
|
||||
if not encoded or len(encoded) > 16_384:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音分片过大或为空。"
|
||||
)
|
||||
return
|
||||
try:
|
||||
content = base64.b64decode(encoded, validate=True)
|
||||
except (ValueError, binascii.Error):
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音分片解析失败。"
|
||||
)
|
||||
return
|
||||
if not content or len(content) > 12 * 1024:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音分片大小无效。"
|
||||
)
|
||||
return
|
||||
if capture.bytes_written + len(content) > 512 * 1024 * 1024:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, "本地录音超过 512 MB 限制。"
|
||||
)
|
||||
self._abort_local_audio_recording(session_id)
|
||||
return
|
||||
try:
|
||||
capture.handle.write(content)
|
||||
except OSError as error:
|
||||
self._emit_local_recording_result(
|
||||
"chunk", session_id, sequence, False, str(error)[:200]
|
||||
)
|
||||
self._abort_local_audio_recording(session_id)
|
||||
return
|
||||
capture.bytes_written += len(content)
|
||||
capture.next_sequence += 1
|
||||
|
||||
def _finish_local_audio_recording(
|
||||
self, session_id: str, total_bytes_value: Any
|
||||
) -> None:
|
||||
capture = self._local_audio_capture
|
||||
try:
|
||||
total_bytes = int(total_bytes_value)
|
||||
except (TypeError, ValueError):
|
||||
total_bytes = -1
|
||||
if capture is None or session_id != capture.session_id:
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, "本地录音会话标识不匹配。"
|
||||
)
|
||||
return
|
||||
self._local_audio_capture = None
|
||||
try:
|
||||
capture.handle.flush()
|
||||
os.fsync(capture.handle.fileno())
|
||||
capture.handle.close()
|
||||
except OSError as error:
|
||||
self._mark_local_audio_invalid(capture.record_id, str(error))
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, str(error)[:200]
|
||||
)
|
||||
return
|
||||
if total_bytes != capture.bytes_written or total_bytes <= 0:
|
||||
self._mark_local_audio_invalid(
|
||||
capture.record_id, "本地录音文件不完整。"
|
||||
)
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, "本地录音文件不完整。"
|
||||
)
|
||||
return
|
||||
if capture.bytes_written < 1024:
|
||||
self._mark_local_audio_invalid(
|
||||
capture.record_id, "本地录音文件为空或只有容器信息。"
|
||||
)
|
||||
self._emit_local_recording_result(
|
||||
"finish",
|
||||
session_id,
|
||||
-1,
|
||||
False,
|
||||
"本地录音文件为空或只有容器信息,已阻止上传。",
|
||||
)
|
||||
return
|
||||
try:
|
||||
with capture.path.open("rb") as recording:
|
||||
signature = recording.read(4)
|
||||
except OSError as error:
|
||||
self._mark_local_audio_invalid(capture.record_id, str(error))
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, str(error)[:200]
|
||||
)
|
||||
return
|
||||
valid_signature = (
|
||||
capture.mime_type.startswith("audio/webm")
|
||||
and signature == b"\x1aE\xdf\xa3"
|
||||
) or (
|
||||
capture.mime_type.startswith("audio/ogg") and signature == b"OggS"
|
||||
)
|
||||
if not valid_signature:
|
||||
self._mark_local_audio_invalid(
|
||||
capture.record_id, "本地录音格式校验失败。"
|
||||
)
|
||||
self._emit_local_recording_result(
|
||||
"finish",
|
||||
session_id,
|
||||
-1,
|
||||
False,
|
||||
"本地录音格式校验失败,已阻止上传无效文件。",
|
||||
)
|
||||
return
|
||||
store = self._local_audio_store
|
||||
uploads = self._local_audio_uploads
|
||||
if store is None or uploads is None:
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, "本机录音队列尚未初始化。"
|
||||
)
|
||||
return
|
||||
try:
|
||||
store.finalize_recording(
|
||||
capture.record_id,
|
||||
size_bytes=capture.bytes_written,
|
||||
)
|
||||
except (OSError, RuntimeError, sqlite3.Error) as error:
|
||||
self._mark_local_audio_invalid(capture.record_id, str(error))
|
||||
self._emit_local_recording_result(
|
||||
"finish", session_id, -1, False, str(error)[:200]
|
||||
)
|
||||
return
|
||||
|
||||
lifecycle = capture.lifecycle
|
||||
|
||||
def enqueue_upload(start_result: Future[bool] | None = None) -> None:
|
||||
try:
|
||||
if start_result is not None and not bool(start_result.result()):
|
||||
raise RuntimeError("通话记录创建失败,录音已保存在本机,可稍后重试。")
|
||||
raw_call_record_id = lifecycle.call_record_id
|
||||
call_record_id = int(raw_call_record_id or 0)
|
||||
if call_record_id <= 0:
|
||||
raise RuntimeError("未取得通话记录编号,录音已保存在本机,可稍后重试。")
|
||||
store.bind_identity(
|
||||
capture.record_id,
|
||||
call_record_id=call_record_id,
|
||||
room_id=lifecycle.current_room_id or "",
|
||||
)
|
||||
uploads.submit(capture.record_id)
|
||||
except Exception as error:
|
||||
with suppress(Exception):
|
||||
store.update_status(
|
||||
capture.record_id,
|
||||
"failed",
|
||||
str(error)[:1000],
|
||||
)
|
||||
|
||||
if lifecycle.call_record_id:
|
||||
enqueue_upload()
|
||||
else:
|
||||
try:
|
||||
lifecycle.start().add_done_callback(enqueue_upload)
|
||||
except Exception as error:
|
||||
store.update_status(
|
||||
capture.record_id,
|
||||
"failed",
|
||||
str(error)[:1000] or "通话记录创建失败。",
|
||||
)
|
||||
self._emit_local_recording_result(
|
||||
"finish",
|
||||
session_id,
|
||||
-1,
|
||||
True,
|
||||
"本地录音已保存,正在后台上传 COS。",
|
||||
)
|
||||
|
||||
def _mark_local_audio_invalid(self, record_id: int, message: str) -> None:
|
||||
store = self._local_audio_store
|
||||
if store is None:
|
||||
return
|
||||
with suppress(Exception):
|
||||
store.mark_invalid(record_id, message)
|
||||
|
||||
def _abort_local_audio_recording(self, session_id: str) -> None:
|
||||
capture = self._local_audio_capture
|
||||
if capture is None or (session_id and capture.session_id != session_id):
|
||||
return
|
||||
self._local_audio_capture = None
|
||||
with suppress(OSError):
|
||||
capture.handle.flush()
|
||||
capture.handle.close()
|
||||
self._mark_local_audio_invalid(
|
||||
capture.record_id,
|
||||
"本次本地录音未正常结束,文件已保留以便排查。",
|
||||
)
|
||||
|
||||
def _start_call_cycle(self) -> None:
|
||||
if self._closing or self._start_requested:
|
||||
return
|
||||
@@ -633,7 +1078,49 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self.close()
|
||||
|
||||
def hangup(self) -> None:
|
||||
self._close_reason = "desktop-hangup"
|
||||
self._request_companion_shutdown("desktop-hangup")
|
||||
|
||||
def _request_companion_shutdown(self, reason: str) -> None:
|
||||
"""Let MediaRecorder finish and upload before WebEngine is destroyed."""
|
||||
|
||||
self._close_reason = reason
|
||||
if self._closing or self._released:
|
||||
return
|
||||
should_wait_for_companion = (
|
||||
self._injected
|
||||
and self._start_requested
|
||||
and not self._call_cycle_closed
|
||||
and not self._companion_ended
|
||||
)
|
||||
if not should_wait_for_companion:
|
||||
self.close()
|
||||
return
|
||||
if self._shutdown_requested:
|
||||
return
|
||||
self._shutdown_requested = True
|
||||
# doctorConsultation.close() stops MediaRecorder, drains every queued
|
||||
# WebChannel chunk, waits for the Qt/COS finish acknowledgement, and
|
||||
# only then emits hangup. Keeping _closing false here is essential:
|
||||
# bridge callbacks are deliberately rejected once final destruction
|
||||
# begins.
|
||||
self._page.runJavaScript(
|
||||
"void window.doctorConsultation?.close?.().catch(() => undefined)"
|
||||
)
|
||||
self._shutdown_timer.start(190_000)
|
||||
|
||||
def _force_requested_shutdown(self) -> None:
|
||||
"""Bound a failed companion shutdown without racing queued uploads."""
|
||||
|
||||
if not self._shutdown_requested or self._closing or self._released:
|
||||
return
|
||||
if self._start_requested and not self._call_cycle_closed:
|
||||
# OrderedCallLifecycle places end after any upload that already
|
||||
# reached the Qt bridge.
|
||||
self.lifecycle.end(f"{self._close_reason}-timeout")
|
||||
self._call_cycle_closed = True
|
||||
self._start_requested = False
|
||||
self._abort_local_audio_recording("")
|
||||
self._companion_ended = True
|
||||
self.close()
|
||||
|
||||
def _begin_shutdown(self) -> None:
|
||||
@@ -641,12 +1128,10 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
return
|
||||
self._closing = True
|
||||
self._media_active = False
|
||||
if self._injected and not self._companion_ended and not self._released:
|
||||
self._page.runJavaScript(
|
||||
"void window.doctorConsultation?.close?.().catch(() => undefined)"
|
||||
)
|
||||
self._shutdown_timer.stop()
|
||||
if self._start_requested and not self._call_cycle_closed:
|
||||
self.lifecycle.end(self._close_reason)
|
||||
self._abort_local_audio_recording("")
|
||||
self._release_webengine()
|
||||
|
||||
def wait_for_lifecycles(self, timeout: float) -> bool:
|
||||
@@ -698,6 +1183,15 @@ if WEBENGINE_AVAILABLE: # pragma: no cover - GUI behavior needs an integration
|
||||
self._profile.deleteLater()
|
||||
|
||||
def closeEvent(self, event: Any) -> None:
|
||||
if (
|
||||
self._injected
|
||||
and self._start_requested
|
||||
and not self._call_cycle_closed
|
||||
and not self._companion_ended
|
||||
):
|
||||
event.ignore()
|
||||
self._request_companion_shutdown(self._close_reason)
|
||||
return
|
||||
self._begin_shutdown()
|
||||
event.accept()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user