from __future__ import annotations import sqlite3 import threading from pathlib import Path from typing import Any import pytest from PySide6.QtTest import QTest from PySide6.QtWidgets import QApplication from doctor_workstation.services.local_audio_queue import ( LocalAudioQueueStore, LocalAudioUploadManager, ) from doctor_workstation.ui.dialogs.local_audio_queue import ( LocalAudioQueueDialog, _display_time, ) def _application() -> QApplication: return QApplication.instance() or QApplication([]) def _ready_record( store: LocalAudioQueueStore, *, diagnosis_id: int, call_record_id: int, session_id: str, room_id: str = "", ) -> int: record = store.begin_recording( session_id=session_id, diagnosis_id=diagnosis_id, mime_type="audio/webm", call_record_id=call_record_id, room_id=room_id, ) payload = b"\x1aE\xdf\xa3" + (b"local-call-audio" * 128) record.file_path.write_bytes(payload) finalized = store.finalize_recording(record.id, size_bytes=len(payload)) return finalized.id class _ConcurrentRepository: def __init__(self, expected: int) -> None: self.expected = expected self.lock = threading.Lock() self.release = threading.Event() self.all_started = threading.Event() self.active = 0 self.maximum_active = 0 self.calls: list[dict[str, Any]] = [] def upload_call_recording(self, **payload: Any) -> dict[str, str]: path = Path(payload["path"]) assert path.is_file() with self.lock: self.active += 1 self.maximum_active = max(self.maximum_active, self.active) self.calls.append(payload) if self.active >= self.expected: self.all_started.set() try: assert self.release.wait(5), "concurrent uploads did not receive release" return {"file_url": f"cos://recordings/{path.name}"} finally: with self.lock: self.active -= 1 class _RetryRepository: def __init__(self) -> None: self.calls = 0 self.call_records: dict[int, list[dict[str, Any]]] = {} self.list_calls: list[int] = [] def upload_call_recording(self, **payload: Any) -> dict[str, str]: self.calls += 1 if self.calls == 1: raise RuntimeError("COS 暂时不可用") return {"file_url": f"cos://recordings/{Path(payload['path']).name}"} def list_call_records(self, diagnosis_id: int) -> list[dict[str, Any]]: self.list_calls.append(int(diagnosis_id)) return list(self.call_records.get(int(diagnosis_id), [])) def test_existing_queue_schema_adds_room_id_without_losing_rows( tmp_path: Path, ) -> None: root = tmp_path / "old-audio-queue" root.mkdir() database = root / "queue.sqlite3" with sqlite3.connect(database) as connection: connection.execute( """ CREATE TABLE local_audio_uploads ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL UNIQUE, diagnosis_id INTEGER NOT NULL, call_record_id INTEGER, mime_type TEXT NOT NULL, file_path TEXT NOT NULL, size_bytes INTEGER NOT NULL DEFAULT 0, status TEXT NOT NULL, error_text TEXT NOT NULL DEFAULT '', uploaded_url TEXT NOT NULL DEFAULT '', attempts INTEGER NOT NULL DEFAULT 0, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, uploaded_at TEXT NOT NULL DEFAULT '' ) """ ) connection.execute( """ INSERT INTO local_audio_uploads ( session_id, diagnosis_id, call_record_id, mime_type, file_path, size_bytes, status, error_text, uploaded_url, attempts, created_at, updated_at, uploaded_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( "legacy-uploaded-session", 8169, 901, "audio/webm", str(root / "legacy.webm"), 2048, "uploaded", "", "cos://recordings/legacy.webm", 2, "2026-08-20T01:00:00+00:00", "2026-08-20T01:02:00+00:00", "2026-08-20T01:02:00+00:00", ), ) store = LocalAudioQueueStore(root) record = store.list_records()[0] with sqlite3.connect(database) as connection: columns = { str(row[1]) for row in connection.execute( "PRAGMA table_info(local_audio_uploads)" ).fetchall() } assert "room_id" in columns assert record.room_id == "" assert record.status == "uploaded" assert record.uploaded_url == "cos://recordings/legacy.webm" assert record.attempts == 2 assert LocalAudioQueueStore(root).require(record.id).room_id == "" def test_local_audio_identity_binding_is_idempotent_and_rejects_conflicts( tmp_path: Path, ) -> None: store = LocalAudioQueueStore(tmp_path / "audio-queue") record = store.begin_recording( session_id="identity-binding-session", diagnosis_id=8169, mime_type="audio/webm", call_record_id=901, room_id="00123456", ) assert record.call_record_id == 901 assert record.room_id == "00123456" assert store.bind_identity( record.id, call_record_id=901, room_id="00123456", ).room_id == "00123456" assert store.bind_identity(record.id, call_record_id=901).room_id == "00123456" with pytest.raises(RuntimeError, match="房间号发生冲突"): store.bind_identity(record.id, call_record_id=901, room_id="99887766") with pytest.raises(RuntimeError, match="通话记录 ID 或房间号发生冲突"): store.bind_identity(record.id, call_record_id=902, room_id="00123456") assert store.require(record.id).call_record_id == 901 assert store.require(record.id).room_id == "00123456" def test_local_audio_queue_uploads_three_files_concurrently(tmp_path: Path) -> None: store = LocalAudioQueueStore(tmp_path / "audio-queue") record_ids = [ _ready_record( store, diagnosis_id=8169, call_record_id=900 + index, session_id=f"session-{index}", ) for index in range(3) ] repository = _ConcurrentRepository(expected=3) manager = LocalAudioUploadManager(repository, store, max_workers=3) futures = [manager.submit(record_id) for record_id in record_ids] try: assert repository.all_started.wait(5) assert repository.maximum_active == 3 finally: repository.release.set() assert [future.result(timeout=5) for future in futures] == [True, True, True] records = [store.require(record_id) for record_id in record_ids] assert all(record.status == "uploaded" for record in records) assert all(record.exists for record in records) assert all(record.uploaded_url.startswith("cos://recordings/") for record in records) assert {call["call_record_id"] for call in repository.calls} == {900, 901, 902} def test_failed_local_audio_is_kept_and_can_be_retried(tmp_path: Path) -> None: store = LocalAudioQueueStore(tmp_path / "audio-queue") record_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="retry-session", ) repository = _RetryRepository() manager = LocalAudioUploadManager(repository, store, max_workers=1) assert manager.submit(record_id).result(timeout=5) is False failed = store.require(record_id) assert failed.status == "failed" assert failed.exists assert failed.attempts == 1 assert "COS 暂时不可用" in failed.error_text store.retry(record_id) assert manager.submit(record_id).result(timeout=5) is True uploaded = store.require(record_id) assert uploaded.status == "uploaded" assert uploaded.exists assert uploaded.attempts == 2 assert uploaded.uploaded_url.startswith("cos://recordings/") def test_local_audio_manager_notifies_when_an_upload_reaches_uploaded( tmp_path: Path, ) -> None: store = LocalAudioQueueStore(tmp_path / "audio-queue") repository = _RetryRepository() repository.calls = 1 manager = LocalAudioUploadManager(repository, store, max_workers=1) record_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="upload-notification-session", ) uploads: list[tuple[int, int, str]] = [] def record_upload(record: Any) -> None: uploads.append( (record.diagnosis_id, int(record.call_record_id or 0), record.uploaded_url) ) manager.add_upload_listener(record_upload) assert manager.submit(record_id).result(timeout=5) is True assert manager.submit(record_id).result(timeout=5) is True manager.remove_upload_listener(record_upload) assert uploads == [ ( 8169, 901, store.require(record_id).uploaded_url, ) ] def test_local_audio_manager_does_not_report_success_without_uploaded_url( tmp_path: Path, ) -> None: class _IncompleteRepository: @staticmethod def upload_call_recording(**_payload: Any) -> dict[str, bool]: return {"completed": True} store = LocalAudioQueueStore(tmp_path / "audio-queue") record_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="missing-url-session", ) manager = LocalAudioUploadManager(_IncompleteRepository(), store, max_workers=1) uploads: list[Any] = [] manager.add_upload_listener(uploads.append) assert manager.submit(record_id).result(timeout=5) is False record = store.require(record_id) assert record.status == "failed" assert "文件地址" in record.error_text assert uploads == [] def test_local_audio_dialog_displays_utc_recording_time_in_business_timezone() -> None: assert _display_time("2026-08-20T01:51:00+00:00") == ( "2026-08-20 09:51:00" ) def test_local_audio_dialog_lists_status_and_retry_controls(tmp_path: Path) -> None: application = _application() store = LocalAudioQueueStore(tmp_path / "audio-queue") repository = _RetryRepository() manager = LocalAudioUploadManager(repository, store, max_workers=1) failed_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="failed-session", room_id="67534825", ) store.update_status(failed_id, "failed", "等待医生重试") uploaded_id = _ready_record( store, diagnosis_id=8169, call_record_id=902, session_id="uploaded-session", room_id="1692231119", ) store.mark_uploaded(uploaded_id, "cos://recordings/uploaded.webm") dialog = LocalAudioQueueDialog( repository, 8169, store=store, manager=manager, ) dialog.show() application.processEvents() try: assert dialog.objectName() == "LocalAudioQueueDialog" assert dialog.table.columnCount() == 8 assert dialog.table.rowCount() == 2 assert dialog.summary_failed.text() == "失败 1" assert dialog.summary_uploaded.text() == "已上传 1" assert dialog.retry_failed_button.isEnabled() assert dialog.table.horizontalHeaderItem(1).text() == "通话记录 ID" assert dialog.table.horizontalHeaderItem(2).text() == "房间号" assert dialog.table.horizontalHeaderItem(5).text() == "上传状态" assert dialog.table.horizontalHeaderItem(6).text() == "失败原因" assert dialog.table.columnWidth(dialog._room_column) == 180 rooms = { dialog.table.item(row, 2).text() for row in range(dialog.table.rowCount()) } assert rooms == {"67534825", "1692231119"} statuses = { dialog.table.item(row, 5).text() for row in range(dialog.table.rowCount()) } assert statuses == {"上传失败", "已上传"} finally: dialog.close() application.processEvents() def test_global_local_audio_dialog_lists_all_diagnoses_and_outcomes( tmp_path: Path, ) -> None: application = _application() store = LocalAudioQueueStore(tmp_path / "audio-queue") repository = _RetryRepository() manager = LocalAudioUploadManager(repository, store, max_workers=1) failed_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="global-failed-session", room_id="67534825", ) store.update_status(failed_id, "failed", "等待医生重试") uploaded_id = _ready_record( store, diagnosis_id=9001, call_record_id=902, session_id="global-uploaded-session", room_id="407179477", ) store.mark_uploaded(uploaded_id, "cos://recordings/global-uploaded.webm") dialog = LocalAudioQueueDialog( repository, None, store=store, manager=manager, ) dialog.show() application.processEvents() try: assert dialog.title_label.text() == "本机录音上传管理" assert dialog.table.columnCount() == 9 assert dialog.table.rowCount() == 2 assert dialog.table.horizontalHeaderItem(1).text() == "诊单 ID" assert dialog.table.horizontalHeaderItem(2).text() == "通话记录 ID" assert dialog.table.horizontalHeaderItem(3).text() == "房间号" assert dialog.table.horizontalHeaderItem(6).text() == "上传状态" diagnosis_ids = { dialog.table.item(row, 1).text() for row in range(dialog.table.rowCount()) } assert diagnosis_ids == {"8169", "9001"} statuses = { dialog.table.item(row, 6).text() for row in range(dialog.table.rowCount()) } assert statuses == {"上传失败", "已上传"} rooms = { dialog.table.item(row, 3).text() for row in range(dialog.table.rowCount()) } assert rooms == {"67534825", "407179477"} assert dialog.summary_failed.text() == "失败 1" assert dialog.summary_uploaded.text() == "已上传 1" finally: dialog.close() application.processEvents() def test_dialog_fetches_and_persists_historical_room_ids_once_per_diagnosis( tmp_path: Path, ) -> None: application = _application() store = LocalAudioQueueStore(tmp_path / "audio-queue") repository = _RetryRepository() repository.call_records[8169] = [ {"id": 902, "room_id": "1692231119"}, {"id": 901, "room_id": "67534825"}, ] manager = LocalAudioUploadManager(repository, store, max_workers=1) first_id = _ready_record( store, diagnosis_id=8169, call_record_id=901, session_id="legacy-room-first", ) second_id = _ready_record( store, diagnosis_id=8169, call_record_id=902, session_id="legacy-room-second", ) dialog = LocalAudioQueueDialog( repository, 8169, store=store, manager=manager, ) dialog.show() try: for _ in range(200): application.processEvents() if store.require(first_id).room_id and store.require(second_id).room_id: break QTest.qWait(10) assert store.require(first_id).room_id == "67534825" assert store.require(second_id).room_id == "1692231119" assert repository.list_calls == [8169] for _ in range(3): dialog.refresh_records() application.processEvents() assert repository.list_calls == [8169] assert { dialog.table.item(row, 2).text() for row in range(dialog.table.rowCount()) } == {"67534825", "1692231119"} finally: dialog.close() application.processEvents()