Files
zyt/app/tests/test_local_audio_queue.py
T
2026-08-20 17:47:14 +08:00

483 lines
16 KiB
Python

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()