From 5706e9a6a51227ebf798df46fabcbabf83f91ab2 Mon Sep 17 00:00:00 2001 From: Your Name Date: Fri, 9 Oct 2026 14:13:57 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../doctor_workstation/services/repository.py | 22 +- .../ui/dialogs/ai_consult.py | 40 ++- app/tests/test_ai_consult_ui.py | 63 ++++ app/tests/test_repository_parity.py | 64 +++- .../controller/tcm/DiagnosisController.php | 56 ++- .../adminapi/logic/tcm/DiagnosisAiLogic.php | 336 +++++++++++------- .../adminapi/service/AssistantSseTiming.php | 89 +++++ .../service/ClinicalKnowledgeReference.php | 9 + server/app/common/service/DifyChatService.php | 4 +- .../service/TcmAncientBooksReference.php | 18 +- server/knowledge/TCM_ANCIENT_BOOKS.md | 2 +- ...DiagnosisAiAssistantStreamContractTest.php | 2 +- .../tests/DiagnosisAiAssistantTimingTest.php | 73 ++++ .../DifyChatAncientBooksReuseContractTest.php | 102 ++++++ server/tests/TcmAncientBooksReferenceTest.php | 14 + 15 files changed, 732 insertions(+), 162 deletions(-) create mode 100644 server/app/adminapi/service/AssistantSseTiming.php create mode 100644 server/tests/DiagnosisAiAssistantTimingTest.php create mode 100644 server/tests/DifyChatAncientBooksReuseContractTest.php diff --git a/app/src/doctor_workstation/services/repository.py b/app/src/doctor_workstation/services/repository.py index 9010fc53e..90bb3c1fa 100644 --- a/app/src/doctor_workstation/services/repository.py +++ b/app/src/doctor_workstation/services/repository.py @@ -19,7 +19,6 @@ from doctor_workstation.core.errors import ( ApiBusinessError, ApiHttpError, ApiProtocolError, - ApiTimeoutError, ApiTransportError, AuthenticationExpiredError, ) @@ -1413,11 +1412,11 @@ class RemoteDoctorRepository: task: str = "custom", cancelled: Callable[[], bool] | None = None, ) -> Iterator[dict[str, Any]]: - """Stream a diagnosis answer, with one legacy fallback before first content.""" + """Stream a diagnosis answer; use the legacy route only when streaming is absent.""" clean_prompt, clean_task = _diagnosis_ai_request(diagnosis_id, prompt, task) body = {"id": diagnosis_id, "prompt": clean_prompt, "task": clean_task} - received_delta = False + received_event = False received_done = False try: for raw_event in self.client.post_event_stream( @@ -1428,13 +1427,12 @@ class RemoteDoctorRepository: ): if cancelled is not None and cancelled(): return + received_event = True event = _normalise_diagnosis_ai_event(raw_event) if event is None: continue kind = event["event"] - if kind == "delta": - received_delta = True - elif kind == "done": + if kind == "done": received_done = True yield event if kind == "done": @@ -1443,12 +1441,16 @@ class RemoteDoctorRepository: return if not received_done: raise ApiProtocolError("AI assistant stream ended before done") - except (ApiHttpError, ApiProtocolError, ApiTimeoutError, ApiTransportError): - if received_delta or (cancelled is not None and cancelled()): + except ApiHttpError as error: + if ( + error.status_code not in {404, 405} + or received_event + or (cancelled is not None and cancelled()) + ): raise - # Older deployments do not expose the stream route. Submit exactly one - # request through the confirmed non-streaming endpoint in that case. + # Only a missing stream route on an older deployment justifies another + # full request. Timeouts, transport failures and malformed streams do not. result = self.analyze_diagnosis_ai( diagnosis_id, clean_prompt, diff --git a/app/src/doctor_workstation/ui/dialogs/ai_consult.py b/app/src/doctor_workstation/ui/dialogs/ai_consult.py index 4442063b2..0fdb6f2bf 100644 --- a/app/src/doctor_workstation/ui/dialogs/ai_consult.py +++ b/app/src/doctor_workstation/ui/dialogs/ai_consult.py @@ -3146,6 +3146,8 @@ class _AiStreamWorker(QRunnable): @Slot() def run(self) -> None: try: + if self.is_cancelled(): + return stream = getattr(self.repository, "stream_diagnosis_ai", None) if callable(stream): events = stream( @@ -3169,6 +3171,8 @@ class _AiStreamWorker(QRunnable): {"event": "delta", "text": answer, "fallback": True}, {**payload, "event": "done", "fallback": True}, ) + if self.is_cancelled(): + return for event in events: if self.is_cancelled(): return @@ -3207,13 +3211,18 @@ class _ClickCard(QFrame): AI_CONTEXT_MAX_CHARS = 320 AI_PROMPT_LIMIT = 500 -#: 首个增量到达前显示的占位文案,避免留下一块没有任何说明的空白气泡。 -AI_STREAM_PENDING_TEXT = "正在生成…" +#: 等待服务端整理诊单资料时显示的占位文案。 +AI_STREAM_PENDING_TEXT = "正在整理诊单资料…" + +#: 服务端确认资料已准备好、尚未收到回答时显示的占位文案。 +AI_STREAM_WAITING_TEXT = "资料已整理,正在等待 AI 回复…" #: 连接结束但一个字都没收到时的兜底文案。 AI_STREAM_SILENT_TEXT = ( "AI 助手没有返回内容,可能是服务端未响应或连接中断,请稍后重试。" ) +AI_STREAM_TOTAL_TIMEOUT_MS = 180_000 +AI_STREAM_TIMEOUT_TEXT = "等待 AI 回复超时,已停止本次等待,请稍后重试。" AI_CONTEXT_SEPARATOR = "\n\n— 医生提问 —\n" @@ -3534,6 +3543,9 @@ class AiConsultDialog(QDialog): self._flush_timer.setSingleShot(True) self._flush_timer.setInterval(40) self._flush_timer.timeout.connect(self._flush_stream_chunks) + self._stream_timeout_timer = QTimer(self) + self._stream_timeout_timer.setSingleShot(True) + self._stream_timeout_timer.timeout.connect(self._stream_timed_out) self.setObjectName("AiConsultDialog") self.setWindowTitle("问诊详情") # QDialog 默认只带关闭按钮,医生无法把这个信息密度很高的窗口放大到整屏。 @@ -6028,6 +6040,7 @@ class AiConsultDialog(QDialog): worker.signals.finished.connect( lambda: self._stream_finished(generation, stream_generation, worker) ) + self._stream_timeout_timer.start(AI_STREAM_TOTAL_TIMEOUT_MS) QThreadPool.globalInstance().start(worker) QTimer.singleShot(0, self._scroll_chat_to_bottom) @@ -6046,6 +6059,8 @@ class AiConsultDialog(QDialog): kind = str(payload.get("event") or "").lower() if kind == "start": self._stream_meta.update(payload) + if self._stream_bubble is not None and not self._stream_text and not self._pending_chunks: + self._stream_bubble.set_payload(AI_STREAM_WAITING_TEXT) return if kind == "delta": chunk = payload.get("text") @@ -6058,6 +6073,7 @@ class AiConsultDialog(QDialog): return self._stream_meta.update(payload) self._stream_completed = True + self._stream_timeout_timer.stop() if not self._stream_text and not self._pending_chunks: answer = first_value(payload, "answer", "content", default="") if answer not in (None, ""): @@ -6094,6 +6110,7 @@ class AiConsultDialog(QDialog): ) -> None: if not self._stream_is_current(generation, stream_generation): return + self._stream_timeout_timer.stop() self._flush_stream_chunks() message = friendly_error(error) if self._stream_text: @@ -6112,10 +6129,11 @@ class AiConsultDialog(QDialog): ) -> None: if not self._stream_is_current(generation, stream_generation): return + self._stream_timeout_timer.stop() if self._stream_worker is worker: self._stream_worker = None - # 连接结束却既没有 done 也没有报错时(服务端静默断开),占位气泡会永远 - # 停在“正在生成…”。这里补一条明确说明,而不是留一块空白卡片。 + # 连接结束却既没有 done 也没有报错时(服务端静默断开), + # 将阶段占位文案替换成明确说明。 if ( not self._stream_completed and not self._stream_text @@ -6127,9 +6145,23 @@ class AiConsultDialog(QDialog): self._asking = False self.send_button.setEnabled(True) + def _stream_timed_out(self) -> None: + if not self._asking or self._stream_completed: + return + self._flush_stream_chunks() + bubble = self._stream_bubble + text = self._stream_text + self._cancel_stream() + if bubble is not None: + bubble.set_payload( + f"{text}\n\n> {AI_STREAM_TIMEOUT_TEXT}" if text else AI_STREAM_TIMEOUT_TEXT + ) + bubble.set_time_text(datetime.now().strftime("%H:%M")) + def _cancel_stream(self) -> None: self._stream_generation += 1 self._flush_timer.stop() + self._stream_timeout_timer.stop() if self._stream_worker is not None: self._stream_worker.cancel() self._stream_worker = None diff --git a/app/tests/test_ai_consult_ui.py b/app/tests/test_ai_consult_ui.py index 1cb99a54e..ee1fa8bcc 100644 --- a/app/tests/test_ai_consult_ui.py +++ b/app/tests/test_ai_consult_ui.py @@ -1491,6 +1491,55 @@ def test_pending_and_silently_closed_streams_never_show_a_blank_bubble( dialog.close() +def test_stalled_stream_exits_pending_state_and_ignores_late_events( + application: QApplication, + monkeypatch: pytest.MonkeyPatch, +) -> None: + dialog = _silent_dialog(monkeypatch) + dialog.show() + application.processEvents() + + dialog._ask("请总结当前病情") + generation, stream_generation = dialog._generation, dialog._stream_generation + bubble = dialog._stream_bubble + worker = dialog._stream_worker + assert bubble is not None and worker is not None + assert dialog._stream_timeout_timer.isActive() + + dialog._stream_timed_out() + assert worker.is_cancelled() + assert not dialog._stream_timeout_timer.isActive() + assert not dialog._asking + assert dialog.send_button.isEnabled() + assert bubble._raw_payload == ai_consult_module.AI_STREAM_TIMEOUT_TEXT + + calls: list[int] = [] + monkeypatch.setattr( + dialog.repository, + "stream_diagnosis_ai", + lambda *args, **kwargs: calls.append(1), + raising=False, + ) + worker.run() + assert calls == [] + dialog._stream_event(generation, stream_generation, {"event": "delta", "text": "迟到回复"}) + dialog._stream_finished(generation, stream_generation, worker) + assert bubble._raw_payload == ai_consult_module.AI_STREAM_TIMEOUT_TEXT + + dialog._ask("请补充说明") + second_bubble = dialog._stream_bubble + assert second_bubble is not None + dialog._stream_event( + dialog._generation, + dialog._stream_generation, + {"event": "delta", "text": "已收到部分分析"}, + ) + dialog._stream_timed_out() + assert second_bubble._raw_payload.startswith("已收到部分分析\n\n>") + assert ai_consult_module.AI_STREAM_TIMEOUT_TEXT in second_bubble._raw_payload + dialog.close() + + def test_answered_stream_replaces_the_pending_placeholder( application: QApplication, monkeypatch: pytest.MonkeyPatch, @@ -1501,15 +1550,29 @@ def test_answered_stream_replaces_the_pending_placeholder( dialog._ask("请总结当前病情") generation, stream_generation = dialog._generation, dialog._stream_generation + assert dialog._stream_bubble is not None + assert dialog._stream_bubble._raw_payload == ai_consult_module.AI_STREAM_PENDING_TEXT + dialog._stream_event( + generation, + stream_generation, + {"event": "start", "source_summary": {"patient_name": "敏感姓名"}}, + ) + assert dialog._stream_bubble._raw_payload == ai_consult_module.AI_STREAM_WAITING_TEXT + assert "敏感姓名" not in dialog._stream_bubble._raw_payload dialog._stream_event(generation, stream_generation, {"event": "delta", "text": "证候:"}) + dialog._flush_timer.stop() + dialog._flush_stream_chunks() + assert dialog._stream_bubble._raw_payload == "证候:" dialog._stream_event(generation, stream_generation, {"event": "delta", "text": "脾肾两虚"}) dialog._stream_event(generation, stream_generation, {"event": "done", "model_label": "千问"}) dialog._stream_finished(generation, stream_generation, dialog._stream_worker) application.processEvents() assert dialog._stream_text == "证候:脾肾两虚" + assert not dialog._stream_timeout_timer.isActive() texts = _bubble_texts(dialog) assert not any(ai_consult_module.AI_STREAM_PENDING_TEXT in text for text in texts) + assert not any(ai_consult_module.AI_STREAM_WAITING_TEXT in text for text in texts) assert not any(ai_consult_module.AI_STREAM_SILENT_TEXT in text for text in texts) dialog.close() diff --git a/app/tests/test_repository_parity.py b/app/tests/test_repository_parity.py index ae10ea09a..a6bcf3e5b 100644 --- a/app/tests/test_repository_parity.py +++ b/app/tests/test_repository_parity.py @@ -8,7 +8,13 @@ from typing import Any import pytest -from doctor_workstation.core.errors import ApiBusinessError, ApiHttpError, ApiProtocolError +from doctor_workstation.core.errors import ( + ApiBusinessError, + ApiHttpError, + ApiProtocolError, + ApiTimeoutError, + ApiTransportError, +) from doctor_workstation.core.models import Appointment, Consultation, PageResult, Prescription from doctor_workstation.services.mock_repository import DemoDoctorRepository from doctor_workstation.services.repository import ( @@ -774,10 +780,11 @@ def test_remote_diagnosis_ai_stream_normalises_chunks_in_order() -> None: assert client.post_calls == [] -def test_remote_diagnosis_ai_stream_falls_back_once_but_not_for_error_event() -> None: +@pytest.mark.parametrize("status_code", [404, 405]) +def test_remote_diagnosis_ai_stream_falls_back_for_missing_route(status_code: int) -> None: class MissingStreamClient(RecordingClient): def post_event_stream(self, *args: Any, **kwargs: Any): - raise ApiHttpError("missing", status_code=404) + raise ApiHttpError("missing", status_code=status_code) missing_client = MissingStreamClient() events = list( @@ -789,6 +796,57 @@ def test_remote_diagnosis_ai_stream_falls_back_once_but_not_for_error_event() -> "tcm.diagnosis/aiAssistant" ] + +@pytest.mark.parametrize( + "stream_error", + [ + ApiHttpError("server error", status_code=500), + ApiTimeoutError("timed out"), + ApiTransportError("disconnected"), + ApiProtocolError("invalid event stream"), + ], +) +def test_remote_diagnosis_ai_stream_failures_do_not_issue_full_request( + stream_error: Exception, +) -> None: + class FailedStreamClient(RecordingClient): + def post_event_stream(self, *args: Any, **kwargs: Any): + raise stream_error + + client = FailedStreamClient() + with pytest.raises(type(stream_error)): + list(RemoteDoctorRepository(client).stream_diagnosis_ai(501, "请分析")) + assert client.post_calls == [] + + +@pytest.mark.parametrize("after_start", [False, True]) +def test_remote_diagnosis_ai_stream_incomplete_response_does_not_fall_back( + after_start: bool, +) -> None: + class IncompleteStreamClient(RecordingClient): + def post_event_stream(self, *args: Any, **kwargs: Any): + if after_start: + yield {"event": "start", "data": {"model_key": "qwen"}} + + client = IncompleteStreamClient() + with pytest.raises(ApiProtocolError, match="ended before done"): + list(RemoteDoctorRepository(client).stream_diagnosis_ai(501, "请分析")) + assert client.post_calls == [] + + +def test_remote_diagnosis_ai_stream_does_not_fall_back_after_start() -> None: + class FailedAfterStartClient(RecordingClient): + def post_event_stream(self, *args: Any, **kwargs: Any): + yield {"event": "start", "data": {"model_key": "qwen"}} + raise ApiHttpError("missing", status_code=404) + + client = FailedAfterStartClient() + with pytest.raises(ApiHttpError): + list(RemoteDoctorRepository(client).stream_diagnosis_ai(501, "请分析")) + assert client.post_calls == [] + + +def test_remote_diagnosis_ai_stream_error_event_does_not_fall_back() -> None: class ErrorStreamClient(RecordingClient): def post_event_stream(self, *args: Any, **kwargs: Any): yield {"event": "error", "data": {"message": "模型繁忙"}} diff --git a/server/app/adminapi/controller/tcm/DiagnosisController.php b/server/app/adminapi/controller/tcm/DiagnosisController.php index bd0279f0e..f16fe653b 100755 --- a/server/app/adminapi/controller/tcm/DiagnosisController.php +++ b/server/app/adminapi/controller/tcm/DiagnosisController.php @@ -22,6 +22,7 @@ use app\adminapi\logic\tcm\DiagnosisLogic; use app\adminapi\logic\tcm\PatientAiReportLogic; use app\adminapi\logic\tcm\TrackingNoteLogic; use app\adminapi\service\AssistantSseProtocol; +use app\adminapi\service\AssistantSseTiming; use app\adminapi\validate\tcm\DiagnosisValidate; use app\common\model\Order; use app\common\model\WechatChatRecord; @@ -930,6 +931,7 @@ class DiagnosisController extends BaseAdminController */ public function aiAssistantStream() { + $timing = new AssistantSseTiming(); // 登录由全局中间件完成;请求校验、旧助手权限与 DataScope 必须全部 // 在任何 SSE header / start 事件之前完成,失败时仍返回标准 JSON。 $params = (new DiagnosisValidate())->post()->goCheck('aiAssistant'); @@ -938,17 +940,18 @@ class DiagnosisController extends BaseAdminController (string) $params['task'], (string) ($params['prompt'] ?? ''), $this->adminId, - $this->adminInfo + $this->adminInfo, + $timing ); if ($prepared === null) { return $this->fail(DiagnosisAiLogic::getError()); } - $this->runAssistantSse($prepared); + $this->runAssistantSse($prepared, $timing); } /** @param array $prepared */ - private function runAssistantSse(array $prepared): void + private function runAssistantSse(array $prepared, AssistantSseTiming $timing): void { while (ob_get_level() > 0) { ob_end_clean(); @@ -984,7 +987,7 @@ class DiagnosisController extends BaseAdminController return !connection_aborted(); }; - $emit('start', [ + $startEmitted = $emit('start', [ 'task' => (string) ($prepared['task'] ?? ''), 'model_key' => (string) ($prepared['profile'] ?? ''), 'diagnosis_id' => (int) ($prepared['diagnosis_id'] ?? 0), @@ -995,14 +998,37 @@ class DiagnosisController extends BaseAdminController : [], 'message' => '已连接,正在生成…', ]); + $sseStartedAt = hrtime(true); + $firstDeltaMs = null; + $timing->log('sse_start', ['status' => $startEmitted ? 'ok' : 'disconnected']); try { $result = DiagnosisAiLogic::streamPreparedAssistant( $prepared, - static fn (string $delta): bool => $emit('delta', ['text' => $delta]), - static fn (): bool => connection_aborted() === 1 + static function (string $delta) use ($emit, $timing, $sseStartedAt, &$firstDeltaMs): bool { + $observedAtMs = $firstDeltaMs === null && $delta !== '' + ? AssistantSseTiming::elapsedMs($sseStartedAt) + : null; + $accepted = $emit('delta', ['text' => $delta]); + if ($accepted && $observedAtMs !== null) { + $firstDeltaMs = $observedAtMs; + $timing->log('first_delta', [ + 'duration_ms' => $firstDeltaMs, + 'delta_bytes' => strlen($delta), + ]); + } + return $accepted; + }, + static fn (): bool => connection_aborted() === 1, + $timing ); if (connection_aborted()) { + $timing->log('sse_done', [ + 'duration_ms' => AssistantSseTiming::elapsedMs($sseStartedAt), + 'first_delta_observed' => $firstDeltaMs !== null, + 'first_delta_ms' => $firstDeltaMs, + 'status' => 'disconnected', + ]); exit; } if ($result === null) { @@ -1013,14 +1039,22 @@ class DiagnosisController extends BaseAdminController } else { $emit('done', $result); } + $timing->log('sse_done', [ + 'duration_ms' => AssistantSseTiming::elapsedMs($sseStartedAt), + 'first_delta_observed' => $firstDeltaMs !== null, + 'first_delta_ms' => $firstDeltaMs, + 'status' => $result === null ? 'error' : 'ok', + ]); } catch (\Throwable $e) { \think\facade\Log::warning('diagnosis ai assistant sse failed ' . json_encode([ - 'diagnosis_id' => (int) ($prepared['diagnosis_id'] ?? 0), - 'profile' => (string) ($prepared['profile'] ?? ''), - 'task' => (string) ($prepared['task'] ?? ''), - 'admin_id' => (int) ($prepared['admin_id'] ?? 0), - 'exception_class' => get_class($e), + 'trace_id' => $timing->id(), ], JSON_UNESCAPED_SLASHES | JSON_INVALID_UTF8_SUBSTITUTE)); + $timing->log('sse_done', [ + 'duration_ms' => AssistantSseTiming::elapsedMs($sseStartedAt), + 'first_delta_observed' => $firstDeltaMs !== null, + 'first_delta_ms' => $firstDeltaMs, + 'status' => 'exception', + ]); $emit('error', [ 'code' => 'AI_ASSISTANT_FAILED', 'message' => 'AI 助手暂时不可用,请稍后重试', diff --git a/server/app/adminapi/logic/tcm/DiagnosisAiLogic.php b/server/app/adminapi/logic/tcm/DiagnosisAiLogic.php index 8b2f5922c..011aa1113 100644 --- a/server/app/adminapi/logic/tcm/DiagnosisAiLogic.php +++ b/server/app/adminapi/logic/tcm/DiagnosisAiLogic.php @@ -12,6 +12,7 @@ use app\common\model\tcm\DiagnosisAiReport; use app\common\service\DifyChatService; use app\common\service\NihaixiaClinicalSkill; use app\common\service\TcmAncientBooksReference; +use app\adminapi\service\AssistantSseTiming; use think\facade\Db; use think\facade\Log; @@ -359,79 +360,117 @@ class DiagnosisAiLogic extends BaseLogic string $task, string $prompt, int $adminId, - array $adminInfo + array $adminInfo, + ?AssistantSseTiming $timing = null ): ?array { - self::$assistantErrorCode = 'AI_ASSISTANT_FAILED'; - $task = strtolower(trim($task)); - if (!isset(self::ASSISTANT_TASKS[$task])) { - self::setError('不支持的 AI 助手任务'); - return null; - } - if ($task === 'prescription_generate' - && !self::hasPermission($adminId, $adminInfo, self::PERMISSION_GENERATE_PRESCRIPTION)) { - self::setError('权限不足,无法使用 AI 生成处方'); - return null; - } + $prepareStartedAt = hrtime(true); + $aggregateMs = 0; + $compactionMs = 0; + $sourceBytes = 0; + $queryBytes = 0; + $attachmentCount = 0; + $compactionAttempted = false; + $status = 'error'; + try { + self::$assistantErrorCode = 'AI_ASSISTANT_FAILED'; + $task = strtolower(trim($task)); + if (!isset(self::ASSISTANT_TASKS[$task])) { + self::setError('不支持的 AI 助手任务'); + return null; + } + if ($task === 'prescription_generate' + && !self::hasPermission($adminId, $adminInfo, self::PERMISSION_GENERATE_PRESCRIPTION)) { + self::setError('权限不足,无法使用 AI 生成处方'); + return null; + } - $diagnosis = self::loadAuthorizedDiagnosis( - $diagnosisId, - $adminId, - $adminInfo, - self::PERMISSION_ASSISTANT, - '权限不足,无法使用诊单 AI 助手' - ); - if ($diagnosis === null) { - return null; - } + $diagnosis = self::loadAuthorizedDiagnosis( + $diagnosisId, + $adminId, + $adminInfo, + self::PERMISSION_ASSISTANT, + '权限不足,无法使用诊单 AI 助手' + ); + if ($diagnosis === null) { + return null; + } - $prompt = self::cleanText($prompt, self::MAX_ASSISTANT_PROMPT_LENGTH, true); - if ($task === 'custom' && $prompt === '') { - self::setError('请输入要咨询的问题'); - return null; - } + $prompt = self::cleanText($prompt, self::MAX_ASSISTANT_PROMPT_LENGTH, true); + if ($task === 'custom' && $prompt === '') { + self::setError('请输入要咨询的问题'); + return null; + } - $context = self::buildCaseContext($diagnosis, $adminId, $adminInfo); - if ($context['case_lines'] === []) { - self::setError('患者纵向资料为空或聚合失败,无法使用 AI 助手'); - return null; - } + $aggregateStartedAt = hrtime(true); + $context = self::buildCaseContext($diagnosis, $adminId, $adminInfo); + $aggregateMs = AssistantSseTiming::elapsedMs($aggregateStartedAt); + $sourceBytes = strlen((string) ($context['case_text'] ?? '')); + $attachmentCount = count(is_array($context['files'] ?? null) ? $context['files'] : []); + if ($context['case_lines'] === []) { + self::setError('患者纵向资料为空或聚合失败,无法使用 AI 助手'); + return null; + } - $profile = self::selectAssistantProfile($task, $prompt); - $modelConfig = self::modelConfigs()[$profile] ?? []; - $model = trim((string) ($modelConfig['name'] ?? '')); - $modelLabel = trim((string) ($modelConfig['label'] ?? $profile)); - if ($model === '') { - self::setError('AI 模型服务尚未完整配置'); - return null; - } + $profile = self::selectAssistantProfile($task, $prompt); + $modelConfig = self::modelConfigs()[$profile] ?? []; + $model = trim((string) ($modelConfig['name'] ?? '')); + $modelLabel = trim((string) ($modelConfig['label'] ?? $profile)); + if ($model === '') { + self::setError('AI 模型服务尚未完整配置'); + return null; + } - if (!self::fitContextForPrompt($context, $profile)) { - return null; - } + $compactionAttempted = $sourceBytes > self::MAX_PROMPT_SOURCE_BYTES; + $compactionStartedAt = hrtime(true); + if (!self::fitContextForPrompt($context, $profile, $timing)) { + $compactionMs = $compactionAttempted ? AssistantSseTiming::elapsedMs($compactionStartedAt) : 0; + return null; + } + $compactionMs = $compactionAttempted ? AssistantSseTiming::elapsedMs($compactionStartedAt) : 0; - return [ - 'diagnosis_id' => $diagnosisId, - 'profile' => $profile, - 'model_name' => $model, - 'model_label' => $modelLabel, - 'task' => $task, - 'inputs' => self::buildUpstreamInputs( + $query = self::buildAssistantPrompt($context, $task, $prompt); + $inputs = self::buildUpstreamInputs( $context, '病例问诊助手', self::ASSISTANT_PROMPT_VERSION - ), - 'query' => self::buildAssistantPrompt($context, $task, $prompt), - 'knowledge_source' => NihaixiaClinicalSkill::knowledgeSource(), - 'user' => 'admin-diagnosis-assistant-' . $adminId, - 'admin_id' => $adminId, - 'files' => is_array($context['files'] ?? null) ? $context['files'] : [], - 'context_scope' => (string) ($context['context_scope'] ?? 'patient_longitudinal'), - 'context_version' => (string) ($context['context_version'] ?? self::ASSISTANT_PROMPT_VERSION), - 'source_summary' => is_array($context['source_summary'] ?? null) ? $context['source_summary'] : [], - 'source_diagnosis_ids' => is_array($context['source_diagnosis_ids'] ?? null) - ? $context['source_diagnosis_ids'] - : [], - ]; + ); + $knowledgeSource = NihaixiaClinicalSkill::knowledgeSource(); + $queryBytes = strlen($query); + $status = 'ok'; + + return [ + 'diagnosis_id' => $diagnosisId, + 'profile' => $profile, + 'model_name' => $model, + 'model_label' => $modelLabel, + 'task' => $task, + 'inputs' => $inputs, + 'query' => $query, + 'knowledge_source' => $knowledgeSource, + 'user' => 'admin-diagnosis-assistant-' . $adminId, + 'admin_id' => $adminId, + 'files' => is_array($context['files'] ?? null) ? $context['files'] : [], + 'context_scope' => (string) ($context['context_scope'] ?? 'patient_longitudinal'), + 'context_version' => (string) ($context['context_version'] ?? self::ASSISTANT_PROMPT_VERSION), + 'source_summary' => is_array($context['source_summary'] ?? null) ? $context['source_summary'] : [], + 'source_diagnosis_ids' => is_array($context['source_diagnosis_ids'] ?? null) + ? $context['source_diagnosis_ids'] + : [], + ]; + } finally { + if ($timing !== null) { + $timing->log('prepare', [ + 'duration_ms' => AssistantSseTiming::elapsedMs($prepareStartedAt), + 'aggregate_ms' => $aggregateMs, + 'compaction_ms' => $compactionMs, + 'source_bytes' => $sourceBytes, + 'query_bytes' => $queryBytes, + 'attachment_count' => $attachmentCount, + 'compaction_attempted' => $compactionAttempted, + 'status' => $status, + ]); + } + } } /** @@ -443,74 +482,42 @@ class DiagnosisAiLogic extends BaseLogic public static function streamPreparedAssistant( array $prepared, callable $onDelta, - ?callable $shouldAbort = null + ?callable $shouldAbort = null, + ?AssistantSseTiming $timing = null ): ?array { $diagnosisId = (int) ($prepared['diagnosis_id'] ?? 0); $profile = (string) ($prepared['profile'] ?? ''); $adminId = (int) ($prepared['admin_id'] ?? 0); $deliveredDelta = false; - $forwardDelta = static function (string $delta) use (&$deliveredDelta, $onDelta) { + $deltaCount = 0; + $deltaBytes = 0; + $retryAttempted = false; + $retryMs = 0; + $status = 'error'; + $upstreamStartedAt = hrtime(true); + $forwardDelta = static function (string $delta) use (&$deliveredDelta, &$deltaCount, &$deltaBytes, $onDelta) { $accepted = $onDelta($delta); if ($accepted !== false) { $deliveredDelta = true; + ++$deltaCount; + $deltaBytes += strlen($delta); } return $accepted; }; try { - $result = DifyChatService::streamChat( - $profile, - is_array($prepared['inputs'] ?? null) ? $prepared['inputs'] : [], - (string) ($prepared['query'] ?? ''), - (string) ($prepared['user'] ?? ''), - $forwardDelta, - $shouldAbort, - is_array($prepared['files'] ?? null) ? $prepared['files'] : [] - ); - } catch (\Throwable $e) { - self::logAssistantFailure( - $diagnosisId, - $profile, - $adminId, - $e, - (string) ($prepared['task'] ?? '') - ); - self::$assistantErrorCode = 'UPSTREAM_UNAVAILABLE'; - self::setError('AI 助手暂时不可用,请稍后重试'); - return null; - } - - if (empty($result['ok'])) { - self::logAssistantUpstreamError( - $diagnosisId, - $profile, - $adminId, - (string) ($prepared['task'] ?? ''), - is_array($result) ? $result : [] - ); - - // Some Dify-compatible gateways accept blocking chat but reject or - // incompletely terminate streaming responses. Before any delta has - // reached the doctor it is safe to make one blocking compatibility - // attempt; after a delta, retrying could duplicate clinical text. - $streamErrorCode = strtoupper(trim((string) ($result['error_code'] ?? ''))); - if ( - !$deliveredDelta - && in_array( - $streamErrorCode, - ['UPSTREAM_REJECTED', 'INCOMPLETE_RESPONSE', 'EMPTY_RESPONSE'], - true - ) - ) { - try { - $result = DifyChatService::chat( - $profile, - is_array($prepared['inputs'] ?? null) ? $prepared['inputs'] : [], - (string) ($prepared['query'] ?? ''), - (string) ($prepared['user'] ?? ''), - is_array($prepared['files'] ?? null) ? $prepared['files'] : [] - ); - } catch (\Throwable $e) { + try { + $result = DifyChatService::streamChat( + $profile, + is_array($prepared['inputs'] ?? null) ? $prepared['inputs'] : [], + (string) ($prepared['query'] ?? ''), + (string) ($prepared['user'] ?? ''), + $forwardDelta, + $shouldAbort, + is_array($prepared['files'] ?? null) ? $prepared['files'] : [] + ); + } catch (\Throwable $e) { + if ($timing === null) { self::logAssistantFailure( $diagnosisId, $profile, @@ -519,7 +526,14 @@ class DiagnosisAiLogic extends BaseLogic (string) ($prepared['task'] ?? '') ); } - if (empty($result['ok'])) { + $status = 'exception'; + self::$assistantErrorCode = 'UPSTREAM_UNAVAILABLE'; + self::setError('AI 助手暂时不可用,请稍后重试'); + return null; + } + + if (empty($result['ok'])) { + if ($timing === null) { self::logAssistantUpstreamError( $diagnosisId, $profile, @@ -528,10 +542,72 @@ class DiagnosisAiLogic extends BaseLogic is_array($result) ? $result : [] ); } + + // Some Dify-compatible gateways accept blocking chat but reject or + // incompletely terminate streaming responses. Before any delta has + // reached the doctor it is safe to make one blocking compatibility + // attempt; after a delta, retrying could duplicate clinical text. + $streamErrorCode = strtoupper(trim((string) ($result['error_code'] ?? ''))); + if ( + !$deliveredDelta + && in_array( + $streamErrorCode, + ['UPSTREAM_REJECTED', 'INCOMPLETE_RESPONSE', 'EMPTY_RESPONSE'], + true + ) + ) { + $retryAttempted = true; + $retryStartedAt = hrtime(true); + try { + $result = DifyChatService::chat( + $profile, + is_array($prepared['inputs'] ?? null) ? $prepared['inputs'] : [], + (string) ($prepared['query'] ?? ''), + (string) ($prepared['user'] ?? ''), + is_array($prepared['files'] ?? null) ? $prepared['files'] : [] + ); + } catch (\Throwable $e) { + if ($timing === null) { + self::logAssistantFailure( + $diagnosisId, + $profile, + $adminId, + $e, + (string) ($prepared['task'] ?? '') + ); + } + } + $retryMs = AssistantSseTiming::elapsedMs($retryStartedAt); + if (empty($result['ok'])) { + if ($timing === null) { + self::logAssistantUpstreamError( + $diagnosisId, + $profile, + $adminId, + (string) ($prepared['task'] ?? ''), + is_array($result) ? $result : [] + ); + } + } + } + } + + $formatted = self::formatAssistantResult($prepared, $result); + $status = $formatted === null ? 'error' : 'ok'; + return $formatted; + } finally { + if ($timing !== null) { + $timing->log('upstream_done', [ + 'duration_ms' => AssistantSseTiming::elapsedMs($upstreamStartedAt), + 'retry_ms' => $retryMs, + 'retry_attempted' => $retryAttempted, + 'delta_count' => $deltaCount, + 'delta_bytes' => $deltaBytes, + 'status' => $status, + 'error_code' => $status !== 'ok' ? self::getAssistantErrorCode() : '', + ]); } } - - return self::formatAssistantResult($prepared, $result); } /** @@ -1477,7 +1553,11 @@ PROMPT; * * @param array $context */ - private static function fitContextForPrompt(array &$context, string $profile): bool + private static function fitContextForPrompt( + array &$context, + string $profile, + ?AssistantSseTiming $timing = null + ): bool { $caseText = (string) ($context['case_text'] ?? ''); if ($caseText === '' || strlen($caseText) <= self::MAX_PROMPT_SOURCE_BYTES) { @@ -1487,12 +1567,12 @@ PROMPT; try { $compacted = PatientAiReportLogic::compactSourceForPrompt($profile, $caseText); } catch (\Throwable $e) { - Log::warning('diagnosis ai context compaction failed', [ - 'diagnosis_id' => (int) ($context['diagnosis_id'] ?? 0), - 'profile' => $profile, - 'source_bytes' => strlen($caseText), - 'exception_class' => get_class($e), - ]); + if ($timing === null) { + Log::warning('diagnosis ai context compaction failed', [ + 'source_bytes' => strlen($caseText), + 'exception_class' => get_class($e), + ]); + } self::setError('患者纵向资料过大,AI 分片读取失败,请稍后重试'); return false; } diff --git a/server/app/adminapi/service/AssistantSseTiming.php b/server/app/adminapi/service/AssistantSseTiming.php new file mode 100644 index 000000000..3586ed470 --- /dev/null +++ b/server/app/adminapi/service/AssistantSseTiming.php @@ -0,0 +1,89 @@ +traceId = bin2hex(random_bytes(16)); + } + + public function id(): string + { + return $this->traceId; + } + + public static function elapsedMs(int $startedAt): int + { + return max(0, (int) floor((hrtime(true) - $startedAt) / 1000000)); + } + + /** + * Only a fixed set of numeric and boolean metadata reaches the log. In particular, + * arbitrary prompt text, patient identifiers, URLs and credentials are discarded. + * + * @param array $metrics + * @return array + */ + public function event(string $phase, array $metrics = []): array + { + $event = [ + 'trace_id' => $this->traceId, + 'phase' => in_array($phase, self::PHASES, true) ? $phase : 'upstream_done', + ]; + foreach (self::NUMBER_FIELDS as $field) { + if (isset($metrics[$field]) && is_int($metrics[$field])) { + $event[$field] = max(0, $metrics[$field]); + } + } + foreach (self::BOOLEAN_FIELDS as $field) { + if (isset($metrics[$field]) && is_bool($metrics[$field])) { + $event[$field] = $metrics[$field]; + } + } + if (isset($metrics['status']) && in_array($metrics['status'], ['ok', 'error', 'exception', 'disconnected'], true)) { + $event['status'] = $metrics['status']; + } + if (isset($metrics['error_code']) && in_array($metrics['error_code'], self::ERROR_CODES, true)) { + $event['error_code'] = $metrics['error_code']; + } + return $event; + } + + /** @param array $metrics */ + public function log(string $phase, array $metrics = []): void + { + try { + Log::info('diagnosis ai assistant timing ' . json_encode($this->event($phase, $metrics))); + } catch (\Throwable $ignored) { + // Telemetry must never interrupt the clinical request. + } + } +} diff --git a/server/app/common/service/ClinicalKnowledgeReference.php b/server/app/common/service/ClinicalKnowledgeReference.php index eedc40a63..56b1ed7eb 100644 --- a/server/app/common/service/ClinicalKnowledgeReference.php +++ b/server/app/common/service/ClinicalKnowledgeReference.php @@ -12,6 +12,15 @@ final class ClinicalKnowledgeReference return TcmAncientBooksReference::augmentQuery(NihaixiaClinicalSkill::augmentQuery($query)); } + /** @param list> $ancientBookSources */ + public static function augmentQueryWithAncientBookSources(string $query, array $ancientBookSources): string + { + return TcmAncientBooksReference::augmentQueryWithSources( + NihaixiaClinicalSkill::augmentQuery($query), + $ancientBookSources + ); + } + /** @param array> $messages * @return array> */ diff --git a/server/app/common/service/DifyChatService.php b/server/app/common/service/DifyChatService.php index f9324f231..63073872e 100644 --- a/server/app/common/service/DifyChatService.php +++ b/server/app/common/service/DifyChatService.php @@ -79,7 +79,7 @@ class DifyChatService $ancientBookSources = ClinicalKnowledgeReference::ancientBookSources($query); try { - $query = ClinicalKnowledgeReference::augmentQuery($query); + $query = ClinicalKnowledgeReference::augmentQueryWithAncientBookSources($query, $ancientBookSources); } catch (\RuntimeException $error) { return self::error('SKILL_UNAVAILABLE', '临床 AI 参考资料不可用'); } @@ -239,7 +239,7 @@ class DifyChatService $ancientBookSources = ClinicalKnowledgeReference::ancientBookSources($query); try { - $query = ClinicalKnowledgeReference::augmentQuery($query); + $query = ClinicalKnowledgeReference::augmentQueryWithAncientBookSources($query, $ancientBookSources); } catch (\RuntimeException $error) { return self::error('SKILL_UNAVAILABLE', '临床 AI 参考资料不可用'); } diff --git a/server/app/common/service/TcmAncientBooksReference.php b/server/app/common/service/TcmAncientBooksReference.php index 14250f9dc..87690acd8 100644 --- a/server/app/common/service/TcmAncientBooksReference.php +++ b/server/app/common/service/TcmAncientBooksReference.php @@ -163,7 +163,16 @@ final class TcmAncientBooksReference if (str_contains($query, self::MARKER)) { return $query; } - $reference = self::referenceForText($query); + return self::augmentQueryWithSources($query, self::sourcesForText($query)); + } + + /** @param list> $sources */ + public static function augmentQueryWithSources(string $query, array $sources): string + { + if (str_contains($query, self::MARKER)) { + return $query; + } + $reference = self::referenceForSources($sources); return $reference === '' ? $query : $reference . "\n\n" . $query; } @@ -193,7 +202,12 @@ final class TcmAncientBooksReference public static function referenceForText(string $text): string { - $sources = self::sourcesForText($text); + return self::referenceForSources(self::sourcesForText($text)); + } + + /** @param list> $sources */ + public static function referenceForSources(array $sources): string + { if ($sources === []) { return ''; } diff --git a/server/knowledge/TCM_ANCIENT_BOOKS.md b/server/knowledge/TCM_ANCIENT_BOOKS.md index 55a8945bc..47dd087ac 100644 --- a/server/knowledge/TCM_ANCIENT_BOOKS.md +++ b/server/knowledge/TCM_ANCIENT_BOOKS.md @@ -15,7 +15,7 @@ php server/scripts/install_tcm_ancient_books.php ```bash mv server/knowledge/tcm-ancient-books server/knowledge/tcm-ancient-books.submodule-backup git pull -php server/scripts/install_tcm_ancient_books.php + ``` 确认新目录已有 701 个 TXT 且索引成功后,旧备份目录即可清理。 diff --git a/server/tests/DiagnosisAiAssistantStreamContractTest.php b/server/tests/DiagnosisAiAssistantStreamContractTest.php index 2ab7755c0..a41c9d871 100644 --- a/server/tests/DiagnosisAiAssistantStreamContractTest.php +++ b/server/tests/DiagnosisAiAssistantStreamContractTest.php @@ -71,7 +71,7 @@ assistantStreamExpect( $actionStart = strpos($controller, 'public function aiAssistantStream()'); $checkAt = strpos($controller, "goCheck('aiAssistant')", $actionStart); $prepareAt = strpos($controller, 'DiagnosisAiLogic::prepareAssistant(', $actionStart); -$runAt = strpos($controller, '$this->runAssistantSse($prepared)', $actionStart); +$runAt = strpos($controller, '$this->runAssistantSse($prepared, $timing)', $actionStart); $headerAt = strpos($controller, "header('Content-Type: text/event-stream; charset=utf-8')", $actionStart); assistantStreamExpect( $actionStart !== false && $checkAt > $actionStart && $prepareAt > $checkAt && $runAt > $prepareAt && $headerAt > $runAt, diff --git a/server/tests/DiagnosisAiAssistantTimingTest.php b/server/tests/DiagnosisAiAssistantTimingTest.php new file mode 100644 index 000000000..827fd7393 --- /dev/null +++ b/server/tests/DiagnosisAiAssistantTimingTest.php @@ -0,0 +1,73 @@ +event('prepare', [ + 'duration_ms' => 120, + 'aggregate_ms' => 75, + 'compaction_ms' => 0, + 'source_bytes' => 1000, + 'query_bytes' => 1100, + 'attachment_count' => 2, + 'compaction_attempted' => false, + 'status' => 'ok', + 'patient_name' => '患者姓名', + 'query' => '患者正文', + 'url' => 'https://private.example/path', + 'api_key' => 'secret', +]); +$second = $one->event('upstream_done', [ + 'duration_ms' => 4300, + 'retry_ms' => 1200, + 'retry_attempted' => true, + 'delta_count' => 4, + 'delta_bytes' => 24, + 'status' => 'error', + 'error_code' => 'UPSTREAM_TIMEOUT', +]); + +timingExpect(preg_match('/^[a-f0-9]{32}$/', $first['trace_id']) === 1, 'trace id is random hex only'); +timingExpect($first['trace_id'] !== $two->event('prepare')['trace_id'], 'requests use distinct trace ids'); +timingExpect($first['trace_id'] === $second['trace_id'], 'all phases share one trace id'); +timingExpect( + $first['aggregate_ms'] === 75 && $first['attachment_count'] === 2 && $first['compaction_attempted'] === false, + 'prepare phase retains useful numeric metadata' +); +timingExpect( + $second['retry_attempted'] === true && $second['retry_ms'] === 1200 + && $second['error_code'] === 'UPSTREAM_TIMEOUT', + 'upstream phase retains retry, duration and safe status' +); +$encoded = json_encode($first, JSON_UNESCAPED_UNICODE); +foreach (['患者姓名', '患者正文', 'private.example', 'secret'] as $sensitive) { + timingExpect(!str_contains($encoded, $sensitive), 'timing log excludes sensitive payload'); +} +timingExpect(!isset($one->event('prepare', ['duration_ms' => '患者正文'])['duration_ms']), 'numeric fields reject text'); +timingExpect(!isset($one->event('upstream_done', ['error_code' => 'KEY=secret'])['error_code']), 'error code is validated'); +timingExpect(!isset($one->event('upstream_done', ['error_code' => 'API_KEY_SECRET'])['error_code']), 'unknown machine codes are not logged'); + +$startedAt = hrtime(true); +timingExpect(AssistantSseTiming::elapsedMs($startedAt) >= 0, 'monotonic elapsed time is non-negative'); + +$controller = file_get_contents(dirname(__DIR__) . '/app/adminapi/controller/tcm/DiagnosisController.php'); +$logic = file_get_contents(dirname(__DIR__) . '/app/adminapi/logic/tcm/DiagnosisAiLogic.php'); +timingExpect(is_string($controller) && is_string($logic), 'instrumented source is readable'); +timingExpect(str_contains($controller, "log('first_delta'") && str_contains($controller, "log('sse_done'"), 'SSE first delta and completion are timed'); +timingExpect(str_contains($logic, "log('prepare'") && str_contains($logic, "log('upstream_done'"), 'prepare and upstream are timed'); +timingExpect(str_contains($logic, 'compaction_attempted') && str_contains($logic, 'retry_attempted'), 'compression and compatibility retry are recorded'); + +echo "Diagnosis AI assistant timing: OK\n"; diff --git a/server/tests/DifyChatAncientBooksReuseContractTest.php b/server/tests/DifyChatAncientBooksReuseContractTest.php new file mode 100644 index 000000000..bd79bd69d --- /dev/null +++ b/server/tests/DifyChatAncientBooksReuseContractTest.php @@ -0,0 +1,102 @@ +options = $options; + $GLOBALS['ancientReuseRequests'][] = json_decode($options[CURLOPT_POSTFIELDS], true); + return true; + } + + function curl_exec(\stdClass $handle): string + { + $options = $handle->options; + if (isset($options[CURLOPT_WRITEFUNCTION])) { + $options[CURLOPT_HEADERFUNCTION]($handle, "HTTP/1.1 200 OK\r\n"); + $options[CURLOPT_WRITEFUNCTION]($handle, + "data: {\"event\":\"message\",\"answer\":\"offline\"}\n\n" + . "data: {\"event\":\"message_end\"}\n\n"); + return ''; + } + return '{"answer":"offline"}'; + } + + function curl_errno(\stdClass $handle): int { return 0; } + function curl_getinfo(\stdClass $handle, int $option): int { return 200; } + function curl_close(\stdClass $handle): void {} +} + +namespace { + require dirname(__DIR__) . '/vendor/autoload.php'; + + use app\common\service\ClinicalKnowledgeReference; + use app\common\service\DifyChatService; + use app\common\service\TcmAncientBooksReference; + + function ancientReuseExpect(bool $condition, string $message): void + { + if (!$condition) { + fwrite(STDERR, "FAIL: {$message}\n"); + exit(1); + } + } + + $indexProperty = new ReflectionProperty(TcmAncientBooksReference::class, 'index'); + $indexProperty->setValue(null, [ + 'commit' => TcmAncientBooksReference::COMMIT, + 'index_version' => TcmAncientBooksReference::INDEX_VERSION, + 'terms' => [ + '消渴' => [[ + 'title' => '甲书', 'chapter' => '消渴篇', 'path' => '甲书.txt', + 'line' => 42, 'excerpt' => '消渴原文。', 'score' => 15, + ]], + ], + 'chapter_terms' => [ + '痛风' => [[ + 'title' => '乙书', 'chapter' => '痛风篇', 'path' => '乙书.txt', + 'line' => 88, 'excerpt' => '痛风原文。', 'score' => 12, + ]], + ], + ]); + $ancientReuseConfig = [ + 'enable' => true, 'base_url' => 'https://ai.example.test/v1/chat-messages', + 'timeout' => 30, 'models' => ['qwen' => ['name' => 'offline', 'api_key' => 'test-key']], + ]; + + foreach (['患者糖尿病,请分析。', '患者痛风,请分析。', '患者有未收录的症状。'] as $case) { + $expectedSources = ClinicalKnowledgeReference::ancientBookSources($case); + $expectedPrompt = ClinicalKnowledgeReference::augmentQuery($case); + foreach (['blocking', 'streaming'] as $mode) { + $ancientReuseRequests = []; + if ($mode === 'blocking') { + $result = DifyChatService::chat('qwen', [], $case, 'offline-user'); + } else { + $deltas = []; + $result = DifyChatService::streamChat('qwen', [], $case, 'offline-user', + static function (string $delta) use (&$deltas): void { $deltas[] = $delta; }); + ancientReuseExpect($deltas === ['offline'], 'streaming forwards the upstream delta'); + } + ancientReuseExpect($result['ok'] === true, "$mode request succeeds"); + ancientReuseExpect(count($ancientReuseRequests) === 1, "$mode sends one request"); + ancientReuseExpect($ancientReuseRequests[0]['query'] === $expectedPrompt, + "$mode prompt remains byte-for-byte identical"); + ancientReuseExpect($result['ancient_book_sources'] === $expectedSources, + "$mode returns exactly the selected citations"); + } + } + + $indexProperty->setValue(null, null); + echo "Dify ancient-books reuse contract passed.\n"; +} diff --git a/server/tests/TcmAncientBooksReferenceTest.php b/server/tests/TcmAncientBooksReferenceTest.php index 551f28488..4d17150a8 100644 --- a/server/tests/TcmAncientBooksReferenceTest.php +++ b/server/tests/TcmAncientBooksReferenceTest.php @@ -52,6 +52,20 @@ ancientExpect(str_contains($augmented, NihaixiaClinicalSkill::VERSION), 'existin ancientExpect(str_ends_with($augmented, $query), 'original patient question remains intact'); ancientExpect(ClinicalKnowledgeReference::augmentQuery($augmented) === $augmented, 'retry does not duplicate references'); +foreach ([$query, '患者出现痛风症状。', '患者有未收录的症状。', '阶段=text 患者糖尿病', '阶段=final 患者糖尿病'] as $case) { + $selected = ClinicalKnowledgeReference::ancientBookSources($case); + ancientExpect( + ClinicalKnowledgeReference::augmentQueryWithAncientBookSources($case, $selected) + === ClinicalKnowledgeReference::augmentQuery($case), + 'selected sources preserve the exact legacy prompt for each clinical case' + ); +} +ancientExpect( + ClinicalKnowledgeReference::augmentQueryWithAncientBookSources('患者有未收录的症状。', []) + === ClinicalKnowledgeReference::augmentQuery('患者有未收录的症状。'), + 'an unrelated request receives no preceding patient citations' +); + $messages = [['role' => 'user', 'content' => $query]]; $wire = ClinicalKnowledgeReference::augmentMessages($messages); ancientExpect(count($wire) === 2 && str_contains($wire[0]['content'], '测试古籍'), 'chat receives reference in system message');