From ab0f01a9cfb00ce92587b6288c18d418482c15fe Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Sat, 26 Sep 2026 11:45:41 +0800 Subject: [PATCH] fix(rectification): stream-first typed turns, 10 s classifier cap, stage progress, run timings (BUG-1047) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - D2: each intent-classifier attempt is capped at 10 s (same session model, thinking untouched); a hang takes the existing retry -> classifier_unavailable path. Success is timed too (RectificationClassifierDiagnostic). - D3: a typed message builds the NDJSON stream first; the first line is turn.progress "received", then the classifier and deterministic replies run inside the stream. Preflight rejections become turn.rejected (old status, code, message) and the client handles them like the old HTTP rejection. Stage lines 收到,正在对照你的档案… / 正在记下这件事… / 正在重新对照盘面… / 正在准备下一个问题… are driven by existing tool events and engine calls, are transient (live row only) and never persisted. VOICE / DESIGN updated. - D4: RectificationRunDiagnostic records per-step start/end, provider token usage incl. reasoning tokens, classifier timing and per-engine-call durations (AsyncLocalStorage scope per turn); RectificationTurnDiagnostic for deterministic turns. No user text, birth data or model text. - D5: /v5/versions memo (30 s, complete identities only, per transport); the exit gate skips its second persistNextInterviewIfIdle when the run's own call found the next focus already active (provably identical). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8 --- frontend/DESIGN.md | 7 +- frontend/docs/VOICE.md | 14 + .../src/app/api/rectification/agent/route.ts | 123 +++- .../components/rectification-agentic-chat.tsx | 61 +- .../src/lib/rectification-activity-labels.ts | 19 + .../lib/rectification-agentic/v9/agent-run.ts | 62 +- .../rectification-agentic/v9/answer-choice.ts | 11 +- .../rectification-agentic/v9/engine-client.ts | 59 +- .../v9/run-diagnostic.ts | 119 +++- .../v9/stream-mapping.ts | 85 +++ .../lib/rectification-agentic/v9/turn-exit.ts | 69 ++- .../v9/turn-instrumentation.ts | 132 +++++ .../v9/turn-intent-classifier.ts | 90 ++- .../rectification-agentic/v9/turn-progress.ts | 18 + .../src/lib/rectification-surface-state.ts | 3 + .../src/lib/rectification-timeline-adapter.ts | 28 +- .../rectification-latency-20260926.test.ts | 556 ++++++++++++++++++ .../rectification-latency-20260926.test.tsx | 345 +++++++++++ ...ctification-latency-route-20260926.test.ts | 182 ++++++ .../tests/rectification-surface-state.test.ts | 5 +- 20 files changed, 1913 insertions(+), 75 deletions(-) create mode 100644 frontend/src/lib/rectification-agentic/v9/turn-instrumentation.ts create mode 100644 frontend/src/lib/rectification-agentic/v9/turn-progress.ts create mode 100644 frontend/tests/rectification-latency-20260926.test.ts create mode 100644 frontend/tests/rectification-latency-20260926.test.tsx create mode 100644 frontend/tests/rectification-latency-route-20260926.test.ts diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md index 690e82e7..7476a980 100644 --- a/frontend/DESIGN.md +++ b/frontend/DESIGN.md @@ -2,6 +2,10 @@ This file adapts the full visual analysis in `CLAUDE_DESIGN.md` to the shipped Jyotisha application. `CLAUDE_DESIGN.md` remains the upstream reference; this file is the implementation contract. +## 校正打字回答的阶段进度句(2026-09-26,BUG-1047) + +生时校正里打字回答一句话,发送的同一帧,这条回复的时间线 live 行就是「收到,正在对照你的档案…」,不再先出现「正在处理…」再干等。服务端先建流、第一行就推这一句,再做意图分类;之后由真实动作推动 `turn.progress`:写入经历 →「正在记下这件事…」,引擎重算 →「正在重新对照盘面…」,经历写完或定下一问 →「正在准备下一个问题…」。阶段只往前走;工具步骤进行中那一行也显示当前阶段句,完成后仍写工具的完成名(「读取校正记录」等);步骤之间不再落回「正在分析…」。这仍是 §9 的「流式生成中」:同一个 `InlineSpinner` live 行,没有新组件、没有新动效;结算时连同 live 行一起消失,不进正文、不进历史。选择题点选、开场与只读续轮不变。 + ## 读者版年运与正文投影(2026-09-25) 年主、Muntha 星座、年度上升与返照时刻分别来自对应年度的真实计算字段,不从年主对象猜补。两种口径不一致时并列注明,不合成单一结论;本地结果不等于完成外部验证。姓名缺失不显示空行,五要素字段使用中文列名。 @@ -300,6 +304,7 @@ The birth-time rectification session is the consultation transcript plus a house | `question-gap`, dead card on the last message | the latest assistant message carries a tap question with no answer and no live card (its focus was superseded, or GET has no card): the card is not drawn at all, and the gap is the repair exit “没有拿到下一个问题。” + “接着问”. A greyed, unclickable A/B/C/D never appears next to the collect-wait placeholder | enabled, generic placeholder | | `question-gap`, delivered | no current question, or the current question is a dead choice card, and the session already delivered a range (`completed_with_range` / `provisional_range` / adopt outcomes, or `tied_first`); range card or range line plus an exit note. Never “没有拿到下一个问题。” | enabled | | `verified_idle` | one closing line `postAdoptVerifyDone` under the still-visible range card (same assistant column); no spinner, no reload | enabled | +| `typed-pending` | the live row reads “收到,正在对照你的档案…” from the frame the answer is sent, then follows `turn.progress`: “正在记下这件事…” → “正在重新对照盘面…” → “正在准备下一个问题…” (BUG-1047). A running tool step shows the same stage line; finished steps keep their done labels. None of it survives the settle | enabled (typing queues), stop visible | | `choice-pending` | the answered card (`data-selected` fill, a top row “正在记录…”) and the same live row from “正在记录本次选择…” through the follow-up turn | enabled (typing queues), stop visible | | `candidates` | one range-delivery card titled “目前范围 …(对照了 N 件经历)”, a caption under the title. The card appears once every collect line has been asked (refresh, guided windows, the skip retry, the seven targeted kinds) or the reader said there is nothing more to add. The precision gate (range ≤ 10 minutes, top-two gap > 3 points, no exact tie) is reported as `precision_gate_met` but never gates the card on its own: meeting it skips no remaining line, and failing it adds no marker to the card once there is nothing left to ask. No free-text invite. Optional “再答两道参考题微调排序” only when unused D9/D10 remain; if those were already asked and the top two are still within one point, a line “这两分钟按现有信息分不开,参考题已经用过.”; then up to three compare columns (highest posterior first; “更像这个” adopts); a closed “查看验证报告” fold. On a normal convergence, unused style questions are asked before this card. Exhausted and closed-ceiling exits still deliver the card if a style question cannot be rendered, once the precision gate or the “nothing more” stop is met. | enabled | | `guided-collect` | a named yes/no card then an entry card. A window whose boundary track carries a domain (D9 / D10 Narayana) asks “YYYY 年 M 到 M 月之间,有没有<那一类事>?”; every other window asks openly — “YYYY 年 M 到 M 月之间,有没有什么事,比如?” — and one window is asked once, whatever label it carries. The entry card is seven type chips (an open window defaults to the first still-open kind) plus a month-granularity date picker (day optional, years from birth year to this year); submitting writes “YYYY 年 M 月(D 日),<领域标签>方面有一件事”, never the question’s example list. Composer stays enabled. | enabled (typing still records), stop visible | @@ -718,7 +723,7 @@ or user IDs. Agent 的 live 标记只有 `InlineSpinner` 一种。曾经并存的 canvas 小球(`thinking-orbs`)已移除,不得再引入第二种 live 标记。 -校正面的所有等待复用行内等待:进入前的 hydration 在揭幕之前完成,进入后唯一的等待形态是时间线 live 行(含「正在准备下一个问题…」这一条独立 live 行)。区间交付卡只挂在最新那条采用旁白下面,不得留在更早的采集/区分题下。有未答的采集或选择题时卡仍在,「更像这个」置灰并写「先答完上面这道,再选时间」,不得整卡消失。卡上至多三列并排,相同性格句只写一次,点「更像这个」即采用该列分钟,按钮显示「正在采用…」或「已采用」。采用过程中整张卡留在原处,不得因 `busy` 卸掉。采用后前事核对结束走 `verified_idle`:一行收尾文案跟在卡片下面、与助手列对齐,没有 live 行、没有重载、没有采用状态条。卡片与右栏细则见 §11、§12。 +校正面的所有等待复用行内等待:进入前的 hydration 在揭幕之前完成,进入后唯一的等待形态是时间线 live 行(含「正在准备下一个问题…」这一条独立 live 行)。打字回答的 live 行从发出那一帧起就是阶段句「收到,正在对照你的档案…」,按服务端 `turn.progress` 换成「正在记下这件事…」「正在重新对照盘面…」「正在准备下一个问题…」,属于流式生成中(BUG-1047)。区间交付卡只挂在最新那条采用旁白下面,不得留在更早的采集/区分题下。有未答的采集或选择题时卡仍在,「更像这个」置灰并写「先答完上面这道,再选时间」,不得整卡消失。卡上至多三列并排,相同性格句只写一次,点「更像这个」即采用该列分钟,按钮显示「正在采用…」或「已采用」。采用过程中整张卡留在原处,不得因 `busy` 卸掉。采用后前事核对结束走 `verified_idle`:一行收尾文案跟在卡片下面、与助手列对齐,没有 live 行、没有重载、没有采用状态条。卡片与右栏细则见 §11、§12。 **每次打开网页只揭幕一次,客户端返回首页走暖快照(2026-09-26,BUG-1040)。** 整页加载(首次打开、刷新、登录后跳转)照旧放下面这一次加载屏。同一次打开里从星盘 / 星历 / 我的报告 / 星盘档案经侧栏「新建对话」、历史会话行、账户页脚或浏览器返回回到 `/`,不再出现加载环:第一次冷启动成功后,模型目录、校正入口摘要、会话分页游标和今日星语记进模块级暖快照(`lib/home-warm-snapshot.ts`,按账户隔离、只在内存、不落任何存储;账户、资料与会话列表本来就在布局层的会话列表 provider 里跨页存活)。首页重挂时若快照齐全,首帧即可用,落点按冷启动同一套规则同步算出(`?new=1` 新建本地空对话、`?c=` 打开该会话、登录返回存根与其人物范围同 BUG-1038),随后后台刷新模型目录、账户、后台回答恢复、入口摘要与今日星语,结果到了静默替换;消息没缓存的历史会话沿用首页内切换会话的留白方式。快照缺任何一项、落点需要去服务端查、或客户端导航没留下目标地址(Next 先渲染新页面、后写地址栏,`AppLink` 在点击时记下目标),一律回到冷启动,不半揭幕。账户变化、退出、任一 401 清空快照;今日星语跨日或换人按新键重取,卡片先用静态句。 diff --git a/frontend/docs/VOICE.md b/frontend/docs/VOICE.md index 8a96dfbf..2c4c1d62 100644 --- a/frontend/docs/VOICE.md +++ b/frontend/docs/VOICE.md @@ -20,6 +20,19 @@ Jyotisha 的可见文案是产品的一部分。正确性红线(真实性、 首页开场语下面可以有一行今日趋势(今日星语卡片的 trend,每日生成);还没生成时写「今天的星语还没写出来。」,没有出生分钟时不写。入口按钮下面平时不写字,只在用户需要动手时写一句:有没做完的校正写「上次那次校正还没完成,可以在历史对话里接着做。」;当前人物不是本人写「生时校正暂时只支持本人。」。不再写「不确定出生时间时,用记得住的经历一步步缩小范围」「上次已经校正完,可以拿最新资料再来一次」,也不要用「·」把两句不相干的提示拼成一行。 +## 生时校正打字回答的阶段进度句(2026-09-26,BUG-1047) + +打字回答发出的那一刻,活动行就写「收到,正在对照你的档案…」;随后只按服务端真实做的事换句,不另起新说法: + +| 阶段 | 句子 | 何时出现 | +| --- | --- | --- | +| 发出即 | 收到,正在对照你的档案… | 按下发送;服务端流的第一行也是它 | +| 记录经历 | 正在记下这件事… | 开始写入这条经历(或记录这次作答) | +| 引擎重算 | 正在重新对照盘面… | 引擎开始重新比较候选 | +| 准备下一问 | 正在准备下一个问题… | 经历写完、在定下一问 | + +这四句只在活动行里,是「流式生成中」的一部分:不进正文、不进历史、刷新后不出现;只往前走,不回头;不写秒数、不写「请稍候」,不承诺多久。 + ## 星盘空状态、0 年大运、星历切日期 这三句只描述眼前缺的数据,不解释运势,也不许写成「正在加载」: @@ -107,6 +120,7 @@ Jyotisha 的可见文案是产品的一部分。正确性红线(真实性、 | ## 统一参数与原始结构(用户刚问「我婚姻怎么样」) | 先用几句口语答婚姻方向,再进入统一参数与原始结构。 | 骨架仍在,开场先回人。 | | 生成中输入框变灰,回车没反应 | 生成中仍可打字;回车后显示「已排队,回答结束后发出」,可撤回 | 对标主流对话产品,焦点不丢。 | | 停止后出现红色告警条 | 已停止,已生成的内容保留;本次不会扣点。 | 停止是用户动作,不是出错。 | +| 打字回答发出后一直是「正在处理… / 正在分析…」,半分钟没有变化 | 发出即「收到,正在对照你的档案…」,随后按阶段换成「正在记下这件事…」「正在重新对照盘面…」「正在准备下一个问题…」。 | 发出后 300 ms 内要有一句确定的话;阶段由真实动作推动,只在活动行里,正文与历史里都没有(BUG-1047)。 | | 也可以再说一件你记得大概时间的事。 / 还差带月份的经历,领域不限。 | 训练门开后给区间卡,并写「如果还记得……范围还能再收一截」。材料不够时只写记下了哪几件、再来一件不是这类的、至少两个具体例子;输入框占位「再说一件带年月的事」。 | 训练门未开时永远给结果或精确缺口;不说「领域不限」「做不了」。交付后采集线关完的邀请见下一行。 | | 能问的都问完了 / 你要是还记得确切哪一天的事,不限领域,说出来我接着算 | 还有题可问时:「现在还剩 04:51–05:11 里 5 个候选,再对照几件经历会更准」,下面是系统点名的题。答「有」之后是口述「大概哪年几月?」,在输入框打字回答,不再出「哪一类事 + 发生年月」的卡。题问完了或用户说「没有了」,就出目前范围卡。 | 有题就继续引导;年月阶段直接打字。没题了照常给结果,不写「未达门槛」。不承诺「再补几件就能定到分钟」。 | | 2018 年 3 到 5 月之间,有没有入职、换工作或职责变重?(引擎其实不知道是哪一类) | 「2018 年 3 到 5 月之间,有没有什么事,比如搬家或开始长期住外地、开始认真关系、分手或结婚?」只有 D9 / D10 那两条轨道才点名领域。 | 边界只说得出「哪段时间」,说不出「哪一类事」。例子最多三个,只列用户还没拒答、也还没说过的领域。 | diff --git a/frontend/src/app/api/rectification/agent/route.ts b/frontend/src/app/api/rectification/agent/route.ts index 1f2fd4f4..79099d65 100644 --- a/frontend/src/app/api/rectification/agent/route.ts +++ b/frontend/src/app/api/rectification/agent/route.ts @@ -53,6 +53,8 @@ import { type RectificationRouteAction, } from "@/lib/rectification-agentic/v9/turn-exit"; import { logRectificationDeliveryTurn } from "@/lib/rectification-agentic/v9/delivery-turn-guard"; +import { createTurnInstrumentation, reportTurnProgress } from "@/lib/rectification-agentic/v9/turn-instrumentation"; +import type { TurnIntentClassifierDiagnostic } from "@/lib/rectification-agentic/v9/turn-intent-classifier"; export const runtime = "nodejs"; export const maxDuration = 240; @@ -147,6 +149,7 @@ async function rectificationBillingRequestId( * phases only. Client history can never override server history. */ export async function POST(request: Request) { + const turnStartedAt = Date.now(); let supabase; let accounting; try { @@ -309,7 +312,8 @@ export async function POST(request: Request) { let expectedWrite: "evidence" | "none" | "unknown" = "none"; let collectIntent: "classified" | "unclassified" | null = null; let writeClassified = false; - const immediateResponse = await (async (): Promise => { + let classifierDiagnostic: TurnIntentClassifierDiagnostic | null = null; + const computeImmediateResponse = async (): Promise => { if (isStructuredChoice) { const actionId = parsed.data.actionId; const focusId = parsed.data.focusId; @@ -459,6 +463,7 @@ export async function POST(request: Request) { const classified = intent.classified; expectedWrite = intent.expectedWrite; writeClassified = true; + classifierDiagnostic = intent.diagnostic ?? null; if (intent.outcome === "classifier_unavailable") { const narration = RECTIFICATION_USER_COPY.classifierUnavailableReply; const turn = await persistV9DeterministicTurn(accounting, userId, caseId, { @@ -487,6 +492,7 @@ export async function POST(request: Request) { } const continueToAgent = shouldContinueAgentForDatedEvent(classified); const previous = previousInferenceFromReceipt(dossier.latestResult?.decisionReceipt ?? null); + reportTurnProgress("recording"); const applied = await applyRectificationChoice(accounting, { userId, caseId, @@ -540,6 +546,7 @@ export async function POST(request: Request) { classified = intent.classified; expectedWrite = intent.expectedWrite; writeClassified = true; + classifierDiagnostic = intent.diagnostic ?? null; collectIntent = classified ? "classified" : "unclassified"; if (classified?.intent === "stop_rectification" || classified?.intent === "ask_about_result") { const previous = previousInferenceFromReceipt(dossier.latestResult?.decisionReceipt ?? null); @@ -564,6 +571,7 @@ export async function POST(request: Request) { const closeStatus = collectFocusCloseStatus(classified); if (closeStatus) { const continueToAgent = shouldContinueAgentForDatedEvent(classified); + reportTurnProgress("recording"); const applied = await applyCollectFocusDenial(accounting, { userId, caseId, @@ -647,6 +655,7 @@ export async function POST(request: Request) { signal: request.signal, }); const classified = intent.classified; + classifierDiagnostic = intent.diagnostic ?? null; if (intent.outcome === "classifier_unavailable") { const narration = RECTIFICATION_USER_COPY.classifierUnavailableReply; const turn = await persistV9DeterministicTurn(accounting, userId, caseId, { @@ -660,6 +669,7 @@ export async function POST(request: Request) { const optionId = optionIdForAnswerClass(persisted.focus, classified.answer_class); if (optionId) { const previous = previousInferenceFromReceipt(dossier.latestResult?.decisionReceipt ?? null); + reportTurnProgress("recording"); const applied = await applyRectificationChoice(accounting, { userId, caseId, @@ -720,7 +730,12 @@ export async function POST(request: Request) { } return null; - })(); + }; + // BUG-1047 D3: a typed message classifies and answers inside the stream, + // after the first progress line. Structured choices and opening/read-only + // keep their existing request/response shape. + const deferToStream = action === "message"; + const immediateResponse = deferToStream ? null : await computeImmediateResponse(); if (immediateResponse) { if (immediateResponse.status !== 200) return immediateResponse; const response = await awaitTurnExitBeforeResponse( @@ -747,15 +762,60 @@ export async function POST(request: Request) { { status: 409 }, ); } - if (action === "message" && !writeClassified) { - const extra = await classifyTurnIntentWithRetry(selectedModel, { - focus: null, - userMessage: parsed.data.message ?? "", - caseStatus, - signal: request.signal, - }); - expectedWrite = extra.expectedWrite; - } + const classifyUnfocusedMessage = async () => { + if (action === "message" && !writeClassified) { + const extra = await classifyTurnIntentWithRetry(selectedModel, { + focus: null, + userMessage: parsed.data.message ?? "", + caseStatus, + signal: request.signal, + }); + expectedWrite = extra.expectedWrite; + classifierDiagnostic = extra.diagnostic ?? null; + } + }; + /** + * Replay a deferred preflight Response on the open stream. A 200 reply + * (answer + run.completed) first passes the same awaited exit gate as the + * immediate path; any non-2xx becomes `turn.rejected` with the old status, + * code and message, which the client treats like the old HTTP rejection. + */ + const replayDeferredResponse = async ( + deferred: Response, + send: (event: Record) => void, + ) => { + if (deferred.status !== 200) { + const payload = await deferred.json().catch(() => null) as + | { code?: unknown; message?: unknown; error?: unknown } + | null; + const message = typeof payload?.message === "string" && payload.message + ? payload.message + : typeof payload?.error === "string" && payload.error + ? payload.error + : `请求失败(${deferred.status})`; + send({ + type: "turn.rejected", + httpStatus: deferred.status, + ...(typeof payload?.code === "string" ? { code: payload.code } : {}), + message, + }); + return; + } + const replayed = await awaitTurnExitBeforeResponse( + deferred, + () => finalizeSuccessfulTurnExit({ + accounting: accounting as never, + userId, + caseId, + action, + askedTurnId: deferred.headers.get("x-rectification-turn-id"), + }), + ); + for (const line of (await replayed.text()).split("\n")) { + if (!line.trim()) continue; + send(JSON.parse(line) as Record); + } + }; const requestTime = new Date(); const chinaTime = new Date(requestTime.getTime() + 8 * 60 * 60 * 1000) @@ -764,8 +824,30 @@ export async function POST(request: Request) { .slice(0, 19); const timeContext = `服务端当前时间(权威):${requestTime.toISOString()};中国标准时间(UTC+8):${chinaTime}。涉及“现在、今天、今年、未来几个月”等相对时间时,以此为准。`; + const logRectificationTurnDiagnostic = ( + path: "deterministic" | "agent", + instrumentation: { engineCalls(): readonly unknown[]; elapsedMs(): number }, + status: number, + ) => { + // Durations, counts and paths only: no message text, no birth data. + console.info(JSON.stringify({ + scope: "RectificationTurnDiagnostic", + requestId, + action, + path, + status, + totalMs: Date.now() - turnStartedAt, + streamMs: instrumentation.elapsedMs(), + classifier: classifierDiagnostic, + engineCalls: instrumentation.engineCalls(), + })); + }; + const encoder = new TextEncoder(); - const body = new ReadableStream({ + // One AsyncLocalStorage scope per turn: engine timings and stage progress + // from anywhere inside the stream land here (BUG-1047). + const instrumentation = createTurnInstrumentation(); + const body = instrumentation.run(() => new ReadableStream({ async start(controller) { let closed = false; const send = (event: Record) => { @@ -774,6 +856,7 @@ export async function POST(request: Request) { if (!safe) return; controller.enqueue(encoder.encode(`${JSON.stringify(safe)}\n`)); }; + instrumentation.setProgressSink((stage) => send({ type: "turn.progress", stage })); const billing: V9RunBilling = { async reserve() { @@ -840,6 +923,17 @@ export async function POST(request: Request) { }; try { + if (deferToStream) { + // First byte of the turn: the progress line, before the classifier. + instrumentation.advance("received"); + const deferred = await computeImmediateResponse(); + if (deferred) { + await replayDeferredResponse(deferred, send); + logRectificationTurnDiagnostic("deterministic", instrumentation, deferred.status); + return; + } + await classifyUnfocusedMessage(); + } if (action === "opening") { logRectificationDeliveryTurn({ trigger: "opening", @@ -881,6 +975,7 @@ export async function POST(request: Request) { generationModel: selectedModel.model, expectedWrite, collectIntent, + classifierDiagnostic, buildAgent: async (turnId, skillPackage, attemptId) => { let decision; try { @@ -926,6 +1021,7 @@ export async function POST(request: Request) { action, assistantBodyPresent: Boolean(result.answerText?.trim()), askedTurnId: result.turnId, + interviewSettled: result.interviewSettled === true, }); } catch (error) { const code = error instanceof RectificationToolServiceError ? error.code : "run_failed"; @@ -933,6 +1029,7 @@ export async function POST(request: Request) { } send({ type: "done", emitted: true }); } + logRectificationTurnDiagnostic("agent", instrumentation, result.ok ? 200 : 500); } catch (error) { const code = error instanceof RectificationToolServiceError ? error.code : "run_failed"; console.error(`[rectification-v9] run failed case=${caseId} code=${code}`); @@ -966,7 +1063,7 @@ export async function POST(request: Request) { } } }, - }); + })); return new Response(body, { headers: { diff --git a/frontend/src/components/rectification-agentic-chat.tsx b/frontend/src/components/rectification-agentic-chat.tsx index 1f88f193..1e6d0ab1 100644 --- a/frontend/src/components/rectification-agentic-chat.tsx +++ b/frontend/src/components/rectification-agentic-chat.tsx @@ -20,11 +20,15 @@ import { RECTIFICATION_ANALYZING_LIVE_LABEL, RECTIFICATION_TOOL_DONE_LABELS, RECTIFICATION_TOOL_PROGRESS_LABELS, + RECTIFICATION_TURN_PROGRESS_LABELS, activityTraceFromReceipt, rectificationCompletedTrail, rectificationLiveProgressLabel, rectificationToolActivityPhase, + rectificationTurnProgressPhase, } from "@/lib/rectification-activity-labels"; +import { isRectificationTurnProgressStage } from "@/lib/rectification-agentic/v9/turn-progress"; +import { withLiveStepLabel } from "@/lib/rectification-timeline-adapter"; import { createRectificationActivityReceiptState, receiptFromRectificationActivityState, @@ -815,9 +819,16 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { clientActionId: requestId, }), }); - if (!response.ok) { - const payload = await response.json().catch(() => null); - const message = payload?.message || payload?.error || `请求失败(${response.status})`; + // A turn the server turned down: the old non-2xx response, or the same + // facts as a `turn.rejected` event once the stream is already open + // (BUG-1047). Either way the live row goes and the typed mark returns. + const rejectTurn = ( + response: { status: number }, + payload: { code?: unknown; message?: unknown; error?: unknown } | null, + ) => { + const message = (typeof payload?.message === "string" && payload.message) + || (typeof payload?.error === "string" && payload.error) + || `请求失败(${response.status})`; setMessages((current) => withdrawTypedMark( current.filter((message) => message.renderKey !== assistantRenderKey), )); @@ -833,6 +844,10 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { } if (response.status === 401) setError("请先登录。"); else setError(message); + }; + if (!response.ok) { + const payload = await response.json().catch(() => null); + rejectTurn(response, payload); return; } if (!response.body) { @@ -849,7 +864,11 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { let completed = false; let streamFailed = false; let runFailedMessage = ""; - while (true) { + let rejected: { status: number; payload: { code?: unknown; message?: unknown } } | null = null; + // Stage line of a typed answer (BUG-1047 D3); null until the first + // `turn.progress`. It replaces the generic tool / "正在分析…" labels. + let stageLabel: string | null = null; + while (!rejected) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); @@ -868,6 +887,8 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { turnId?: unknown; code?: unknown; activity?: unknown; + stage?: unknown; + httpStatus?: unknown; }; try { event = JSON.parse(line) as typeof event; @@ -875,7 +896,19 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { continue; } if (typeof event.type !== "string") continue; - if (event.type === "answer.delta" && typeof event.text === "string") { + if (event.type === "turn.rejected" && typeof event.httpStatus === "number") { + rejected = { status: event.httpStatus, payload: { code: event.code, message: event.message } }; + break; + } + if (event.type === "turn.progress" && isRectificationTurnProgressStage(event.stage)) { + stageLabel = RECTIFICATION_TURN_PROGRESS_LABELS[event.stage]; + activityTrace = withLiveStepLabel(activityTrace, stageLabel); + currentActivity = nextActivityView(currentActivity, { + phase: rectificationTurnProgressPhase(event.stage), + label: rememberLiveActivity(stageLabel, null), + }); + frames.touch(); + } else if (event.type === "answer.delta" && typeof event.text === "string") { raw = event.replace === true ? event.text : raw + event.text; activityTrace = freezeLiveThink(activityTrace); currentActivity = nextActivityView(currentActivity, { @@ -887,7 +920,7 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { const activity = event.activity; currentActivity = nextActivityView(currentActivity, { phase: "evidence-validation", - label: rememberLiveActivity(RECTIFICATION_ACTIVITY_PROGRESS_LABELS[activity]), + label: rememberLiveActivity(stageLabel ?? RECTIFICATION_ACTIVITY_PROGRESS_LABELS[activity]), }); frames.touch(); } else if (event.type === "attempt.reset") { @@ -896,9 +929,10 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { activityReceiptState = createRectificationActivityReceiptState(); completedReceipt = receiptFromRectificationActivityState(activityReceiptState); completedTurnId = undefined; + if (stageLabel) stageLabel = RECTIFICATION_TURN_PROGRESS_LABELS.received; currentActivity = nextActivityView(undefined, { phase: "evidence-validation", - label: rememberLiveActivity("正在处理…", null), + label: rememberLiveActivity(stageLabel ?? "正在处理…", null), }); frames.reset(); setMessages((current) => current.map((message) => message.renderKey === assistantRenderKey @@ -932,11 +966,11 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { activityTrace = startActivityTraceStep( activityTrace, tool, - RECTIFICATION_TOOL_PROGRESS_LABELS[tool], + stageLabel ?? RECTIFICATION_TOOL_PROGRESS_LABELS[tool], ); currentActivity = nextActivityView(currentActivity, { phase: rectificationToolActivityPhase(tool), - label: rememberLiveActivity(RECTIFICATION_TOOL_PROGRESS_LABELS[tool], tool), + label: rememberLiveActivity(stageLabel ?? RECTIFICATION_TOOL_PROGRESS_LABELS[tool], tool), completedTrail: rectificationCompletedTrail(activityReceiptState.completedSteps), }); frames.touch(); @@ -958,7 +992,9 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { completedReceipt = receiptFromRectificationActivityState(activityReceiptState); currentActivity = nextActivityView(currentActivity, { phase: "evidence-validation", - label: rememberLiveActivity(RECTIFICATION_ANALYZING_LIVE_LABEL, null), + label: stageLabel + ? rememberLiveActivity(stageLabel, null) + : rememberLiveActivity(RECTIFICATION_ANALYZING_LIVE_LABEL, null), completedTrail: rectificationCompletedTrail(activityReceiptState.completedSteps), }); frames.touch(); @@ -966,6 +1002,11 @@ export function RectificationAgenticChat(props: RectificationAgenticChatProps) { } } frames.settle(); + if (rejected) { + await reader.cancel().catch(() => undefined); + rejectTurn({ status: rejected.status }, rejected.payload); + return; + } const parsed = completed && !streamFailed ? parseAgentReply(raw) : { text: "", title: undefined }; const succeeded = completed && !streamFailed && Boolean(parsed.text); diff --git a/frontend/src/lib/rectification-activity-labels.ts b/frontend/src/lib/rectification-activity-labels.ts index 6a353325..8310048a 100644 --- a/frontend/src/lib/rectification-activity-labels.ts +++ b/frontend/src/lib/rectification-activity-labels.ts @@ -6,6 +6,7 @@ import type { PublicRectificationActivity, PublicRectificationTool, } from "./rectification-agentic/v9/public-receipt.ts"; +import type { RectificationTurnProgressStage } from "./rectification-agentic/v9/turn-progress.ts"; import { RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS as ATTEMPT_TIMEOUT_MS } from "./rectification-run-budget.ts"; export { RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS } from "./rectification-run-budget.ts"; @@ -36,6 +37,24 @@ export const RECTIFICATION_ACTIVITY_PROGRESS_LABELS: Readonly> = { + received: "收到,正在对照你的档案…", + recording: "正在记下这件事…", + rescoring: "正在重新对照盘面…", + preparing_question: "正在准备下一个问题…", +}; + +export function rectificationTurnProgressPhase(stage: RectificationTurnProgressStage): PublicActivityPhase { + if (stage === "received") return "loading-method"; + if (stage === "rescoring") return "chart-calculation"; + return "evidence-validation"; +} + /** Live label between a completed tool and the next tool or spoken answer. */ export const RECTIFICATION_ANALYZING_LIVE_LABEL = "正在分析…"; diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run.ts b/frontend/src/lib/rectification-agentic/v9/agent-run.ts index 7fcd361a..f1798057 100644 --- a/frontend/src/lib/rectification-agentic/v9/agent-run.ts +++ b/frontend/src/lib/rectification-agentic/v9/agent-run.ts @@ -68,9 +68,23 @@ import { mapStreamChunkToPhase, streamToolNames, isPublicRectificationToolName, + turnProgressForChunk, type PublicStreamEvent, } from "./stream-mapping"; -import { mapModelFinishToErrorCode, userFacingRunFailure } from "./run-diagnostic"; +import { + diagnosticStepsFromTimings, + mapModelFinishToErrorCode, + recordStepChunk, + userFacingRunFailure, + type RectificationRunDiagnostic, + type RectificationStepTiming, +} from "./run-diagnostic"; +import { + currentEngineCallTimings, + reportTurnProgress, + resetTurnProgress, +} from "./turn-instrumentation.ts"; +import type { TurnIntentClassifierDiagnostic } from "./turn-intent-classifier"; import { batchResultFromToolChunk, batchRescoreFailed, @@ -132,6 +146,8 @@ export type V9AgentRunOptions = Readonly<{ minRetryAttemptMs?: number; expectedWrite?: "evidence" | "none" | "unknown"; collectIntent?: "classified" | "unclassified" | null; + /** Classifier timing from the route preflight, for the run diagnostic only. */ + classifierDiagnostic?: TurnIntentClassifierDiagnostic | null; }>; export { @@ -151,6 +167,12 @@ export type V9AgentRunResult = Readonly<{ toolsUsed: readonly string[]; errorCode: string | null; previousFocusId: string | null; + /** + * True when this run's own post-turn interview call found the next focus + * already active and linked it to this turn (BUG-1047 D5). The route's + * common exit gate then skips its second, identical re-read. + */ + interviewSettled?: boolean; }>; type AttemptStatus = "completed" | "failed" | "retryable"; @@ -559,6 +581,8 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise> | null = null; + let interviewSettled = false; if (action === "opening" || action === "evidence") { try { + reportTurnProgress("preparing_question"); interviewIdle = await persistNextInterviewIfIdle({ accounting, userId, caseId, askedTurnId: turnId }); + // Only the "focus already active" exit is provably idempotent: a second + // call re-reads the same focus and re-links the same turn (BUG-1047 D5). + interviewSettled = interviewIdle.focusActive === true; if (interviewIdle.terminalNote && interviewIdle.hostNarration && !turnId) { try { await persistExhaustionGateTurn({ @@ -723,6 +752,7 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise { + const attemptStartedAt = Date.now(); + const stepTimings: RectificationStepTiming[] = []; + const engineCallsBefore = currentEngineCallTimings().length; const agent = await buildAgent(turnId, skillPackage, attemptId); let frameworkSkill: unknown = null; try { @@ -938,6 +971,9 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise { +}): Promise<{ + persisted: boolean; + choiceReady: boolean; + hostNarration: string | null; + terminalNote?: boolean; + /** The next focus was already active; it was only (re)linked to this turn. */ + focusActive?: boolean; +}> { const rescored = await rescoreStaleMinuteSnapshotIfNeeded({ accounting: input.accounting, userId: input.userId, @@ -1716,7 +1723,7 @@ export async function persistNextInterviewIfIdle(input: { ); } } - return finishIdle({ persisted: false, choiceReady: false, hostNarration: null }); + return finishIdle({ persisted: false, choiceReady: false, hostNarration: null, focusActive: true }); } let birthDate: string | null = null; try { diff --git a/frontend/src/lib/rectification-agentic/v9/engine-client.ts b/frontend/src/lib/rectification-agentic/v9/engine-client.ts index dfc043d2..f4f1bb6c 100644 --- a/frontend/src/lib/rectification-agentic/v9/engine-client.ts +++ b/frontend/src/lib/rectification-agentic/v9/engine-client.ts @@ -26,6 +26,7 @@ import { questionContractVersionIsCompatible } from "./probe-question-contract"; import { resolveAyanamsa } from "../../ayanamsa.ts"; import { requestCandidateIntervals, readCandidatePosition, type CandidatePosition, type DatedCandidateRange } from "../core/candidate-window.ts"; import { parseBlockScanPayload, type BlockScanBlock } from "./block-scan.ts"; +import { recordEngineCallTiming, reportTurnProgress } from "./turn-instrumentation.ts"; export type EngineCallFailureKind = "busy" | "http_error" | "timeout" | "bad_payload"; @@ -290,8 +291,33 @@ function retryAfterFromHeaders(headers: { get?(name: string): string | null } | return parseEngineRetryAfterSeconds(headers.get("retry-after") ?? headers.get("Retry-After")); } +/** Engine paths that recompare candidates against the chart (BUG-1047 stage). */ +const RESCORING_ENGINE_PATHS = new Set([ + "/api/rectification/v5/score", + "/api/rectification/v5/block_scan", +]); + async function readEngineJson(path: string, init: RequestInit, timeoutMs: number): Promise> { + if (RESCORING_ENGINE_PATHS.has(path)) reportTurnProgress("rescoring"); const started = Date.now(); + try { + return await readEngineJsonTimed(path, init, timeoutMs, started); + } catch (error) { + recordEngineCallTiming({ + path, + ms: Date.now() - started, + outcome: error instanceof RectificationEngineError && error.kind ? error.kind : "http_error", + }); + throw error; + } +} + +async function readEngineJsonTimed( + path: string, + init: RequestInit, + timeoutMs: number, + started: number, +): Promise> { try { const response = await fetch(`${engineBase()}${path}`, { ...init, @@ -332,6 +358,7 @@ async function readEngineJson(path: string, init: RequestInit, timeoutMs: number path, }); } + recordEngineCallTiming({ path, ms: Date.now() - started, outcome: "ok" }); return data as Record; } catch (error) { if (error instanceof RectificationEngineError) throw error; @@ -452,11 +479,41 @@ export function cachedEngineScoreIsReusable( return scoringIdentityMatches(stored, live); } +/** + * In-process memo for `/v5/versions` when the deployment does not pin the + * scoring identity in env (BUG-1047 D5). One turn asked the engine for its + * versions several times. Only a complete identity pair is kept; failures and + * partial payloads are always re-read. The memo is scoped to the transport + * (`fetch` implementation) and engine base URL, never to a user or Case. + */ +export const RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS = 30_000; + +type VersionsMemo = { identity: LiveEngineScoringIdentity; expiresAt: number }; +const versionsMemo = new WeakMap>(); + +async function readNativeScoringIdentity(): Promise { + const transport = globalThis.fetch as unknown as object; + const base = engineBase(); + const now = Date.now(); + const memo = versionsMemo.get(transport)?.get(base); + if (memo && memo.expiresAt > now) return { ...memo.identity }; + const native = scoringIdentityFromEnginePayload(await getEngine("/api/rectification/v5/versions")); + if (scoringIdentityIsTrusted(native)) { + let byBase = versionsMemo.get(transport); + if (!byBase) { + byBase = new Map(); + versionsMemo.set(transport, byBase); + } + byBase.set(base, { identity: native, expiresAt: now + RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS }); + } + return native; +} + export async function readV9EngineScoringIdentity(): Promise { const fromEnv = liveEngineScoringIdentityFromEnv(); if (scoringIdentityIsTrusted(fromEnv)) return fromEnv; try { - const native = scoringIdentityFromEnginePayload(await getEngine("/api/rectification/v5/versions")); + const native = await readNativeScoringIdentity(); // Partial overrides cannot invent the missing half of a computation identity. // A conflicting partial override is not a coherent identity pair. if ((fromEnv.algorithmVersion && fromEnv.algorithmVersion !== native.algorithmVersion) diff --git a/frontend/src/lib/rectification-agentic/v9/run-diagnostic.ts b/frontend/src/lib/rectification-agentic/v9/run-diagnostic.ts index f69da5c7..88616030 100644 --- a/frontend/src/lib/rectification-agentic/v9/run-diagnostic.ts +++ b/frontend/src/lib/rectification-agentic/v9/run-diagnostic.ts @@ -19,23 +19,140 @@ export const RECTIFICATION_FINISH_REASONS = [ export type RectificationFinishReason = (typeof RECTIFICATION_FINISH_REASONS)[number]; +/** + * One model step inside an attempt (BUG-1047 D4). Times are milliseconds from + * the attempt start; token counts only when the provider reports them. Tool + * names are the public allowlisted ids. No text, arguments or results. + */ +export type RectificationStepTiming = { + index: number; + startMs: number; + endMs: number | null; + inputTokens: number | null; + outputTokens: number | null; + reasoningTokens: number | null; + tools: string[]; +}; + export type RectificationRunDiagnostic = Readonly<{ runId: string; modelId: string; - finishReason: RectificationFinishReason; + /** Provider finish reason as mapped by `toAgentModelFinishReason`, or "unknown". */ + finishReason: string; + /** Sum over model steps when the provider reports usage; otherwise null. */ inputTokens: number | null; reasoningTokens: number | null; outputTokens: number | null; + /** Model steps in this attempt (was: distinct tools used; BUG-1047). */ stepCount: number; + /** Public tool calls in this attempt, repeats included. */ toolCallCount: number; + distinctToolCount: number; readCasePayloadBytes: number | null; + /** Milliseconds since the run started (the whole turn, all attempts). */ elapsedMs: number; + attemptNumber: number; + /** Milliseconds from run start to this attempt's start. */ + attemptStartMs: number; + attemptElapsedMs: number; + steps: readonly RectificationStepTiming[]; lastCompletedTool: string | null; stateMutationCommitted: boolean; expectedWrite: "evidence" | "none" | "unknown" | null; collectIntent: "classified" | "unclassified" | null; + classifier: Readonly<{ + outcome: string; + attempts: number; + timedOutAttempts: number; + elapsedMs: number; + }> | null; + engineCalls: readonly Readonly<{ path: string; ms: number; outcome: string }>[]; }>; +function usageNumber(value: unknown): number | null { + return typeof value === "number" && Number.isFinite(value) && value >= 0 ? Math.trunc(value) : null; +} + +/** + * Fold one fullStream chunk into the step timeline. `step-start` opens a + * step, public `tool-call`s are attached to the open step, `step-finish` + * closes it with the provider's usage for that step. + */ +function newStep(steps: readonly RectificationStepTiming[], startMs: number): RectificationStepTiming { + return { + index: steps.length, + startMs, + endMs: null, + inputTokens: null, + outputTokens: null, + reasoningTokens: null, + tools: [], + }; +} + +export function recordStepChunk( + steps: RectificationStepTiming[], + chunk: Readonly<{ type: string; payload?: unknown }>, + elapsedMs: number, + isPublicTool: (name: string) => boolean, +): void { + const last = steps.length > 0 ? steps[steps.length - 1] : null; + const open = last && last.endMs === null ? last : null; + // A chunk that arrives outside step-start/step-finish opens an implicit + // step starting where the previous one ended. + const openOrCreate = () => { + if (open) return open; + const created = newStep(steps, last?.endMs ?? 0); + steps.push(created); + return created; + }; + if (chunk.type === "step-start") { + if (open) open.endMs = elapsedMs; + steps.push(newStep(steps, elapsedMs)); + return; + } + if (chunk.type === "tool-call") { + const payload = chunk.payload as { toolName?: unknown } | undefined; + const name = typeof payload?.toolName === "string" ? payload.toolName : ""; + if (name && isPublicTool(name)) openOrCreate().tools.push(name); + return; + } + if (chunk.type === "step-finish") { + const payload = chunk.payload as { output?: { usage?: Record } } | undefined; + const usage = payload?.output?.usage ?? {}; + const target = openOrCreate(); + target.endMs = elapsedMs; + target.inputTokens = usageNumber(usage.inputTokens); + target.outputTokens = usageNumber(usage.outputTokens); + target.reasoningTokens = usageNumber(usage.reasoningTokens); + } +} + +function sumOrNull(values: readonly (number | null)[]): number | null { + const present = values.filter((value): value is number => value !== null); + return present.length > 0 ? present.reduce((total, value) => total + value, 0) : null; +} + +/** Totals derived from the step timeline; null when no step reported usage. */ +export function diagnosticStepsFromTimings(steps: readonly RectificationStepTiming[]): Readonly<{ + steps: readonly RectificationStepTiming[]; + stepCount: number; + inputTokens: number | null; + outputTokens: number | null; + reasoningTokens: number | null; + toolCallCount: number; +}> { + const copy = steps.map((step) => ({ ...step, tools: [...step.tools] })); + return { + steps: copy, + stepCount: copy.length, + inputTokens: sumOrNull(copy.map((step) => step.inputTokens)), + outputTokens: sumOrNull(copy.map((step) => step.outputTokens)), + reasoningTokens: sumOrNull(copy.map((step) => step.reasoningTokens)), + toolCallCount: copy.reduce((total, step) => total + step.tools.length, 0), + }; +} + const USER_COPY: Readonly> = { answer_truncated: "模型输出达到上限,状态已记录。", length: "模型输出达到上限,状态已记录。", diff --git a/frontend/src/lib/rectification-agentic/v9/stream-mapping.ts b/frontend/src/lib/rectification-agentic/v9/stream-mapping.ts index a1820c9a..cb112b26 100644 --- a/frontend/src/lib/rectification-agentic/v9/stream-mapping.ts +++ b/frontend/src/lib/rectification-agentic/v9/stream-mapping.ts @@ -23,6 +23,10 @@ import { type PublicRectificationTool, } from "./public-receipt"; import { isToolInputRejection, toolResultFromChunk } from "./host-fallback"; +import { + isRectificationTurnProgressStage, + type RectificationTurnProgressStage, +} from "./turn-progress.ts"; export type PublicPhaseStreamEvent = Readonly<{ type: PublicRectificationPhase; @@ -58,7 +62,30 @@ export type PublicFailedStreamEvent = Readonly<{ message?: string; }>; +/** + * Transient live-row stage for one turn (BUG-1047 D3). Never persisted as a + * phase, never part of `assistant_message`. + */ +export type PublicTurnProgressEvent = Readonly<{ + type: "turn.progress"; + stage: RectificationTurnProgressStage; +}>; + +/** + * A message turn that the route turned down before any model or write ran + * to completion. It carries what a non-2xx JSON response used to carry, so + * the client handles it exactly like the old HTTP rejection (BUG-1047). + */ +export type PublicTurnRejectedEvent = Readonly<{ + type: "turn.rejected"; + httpStatus: number; + code?: string; + message: string; +}>; + export type PublicStreamEvent = + | PublicTurnProgressEvent + | PublicTurnRejectedEvent | PublicPhaseStreamEvent | RectificationActivityEvent | RectificationActivityChangedEvent @@ -278,6 +305,25 @@ export function safePublicEvent(value: unknown): PublicStreamEvent | null { origin?: unknown; }; if (event.type === "thinking.delta") return null; + if (event.type === "turn.progress") { + const stage = (value as { stage?: unknown }).stage; + return isRectificationTurnProgressStage(stage) ? { type: "turn.progress", stage } : null; + } + if (event.type === "turn.rejected") { + const httpStatus = (value as { httpStatus?: unknown }).httpStatus; + if ( + typeof httpStatus !== "number" || !Number.isInteger(httpStatus) + || httpStatus < 400 || httpStatus > 599 + ) return null; + if (typeof event.message !== "string") return null; + const code = typeof event.code === "string" && /^[a-z_]{1,80}$/.test(event.code) ? event.code : undefined; + return { + type: "turn.rejected", + httpStatus, + ...(code ? { code } : {}), + message: event.message.slice(0, 500), + }; + } if (event.type === "error") { const codes = new Set([ "billing_denied", @@ -360,3 +406,42 @@ export function safePublicEvent(value: unknown): PublicStreamEvent | null { export function activityChangedFromTool(tool: PublicRectificationTool): RectificationActivityChangedEvent { return { type: "activity.changed", activity: activityForRectificationTool(tool) }; } + +const RECORDING_TOOLS = new Set([ + "rectification-record-evidence-batch", + "rectification-propose-evidence", + "rectification-confirm-evidence", + "rectification-revise-evidence", +]); + +const RESCORING_TOOLS = new Set([ + "rectification-compare-candidates", + "rectification-read-diagnostics", +]); + +const QUESTION_TOOLS = new Set([ + "rectification-set-focus", + "rectification-resolve-focus", + "rectification-offer-candidates", + "rectification-stop-and-review", +]); + +/** + * Stage signal carried by an existing tool event (BUG-1047 D3). Writing an + * experience starts "recording"; a recompare starts "rescoring"; a finished + * evidence write or a focus/offer tool starts "preparing_question". The + * engine client adds "rescoring" when the batch write recompares. + */ +export function turnProgressForChunk(chunk: AgentChunkType): RectificationTurnProgressStage | null { + if (chunk.type !== "tool-call" && chunk.type !== "tool-result") return null; + const toolName = typeof chunk.payload?.toolName === "string" ? chunk.payload.toolName : ""; + if (!isPublicRectificationTool(toolName)) return null; + if (chunk.type === "tool-call") { + if (RECORDING_TOOLS.has(toolName)) return "recording"; + if (RESCORING_TOOLS.has(toolName)) return "rescoring"; + if (QUESTION_TOOLS.has(toolName)) return "preparing_question"; + return null; + } + if (isToolInputRejection(toolResultFromChunk(chunk))) return null; + return RECORDING_TOOLS.has(toolName) ? "preparing_question" : null; +} diff --git a/frontend/src/lib/rectification-agentic/v9/turn-exit.ts b/frontend/src/lib/rectification-agentic/v9/turn-exit.ts index df195428..833c2db9 100644 --- a/frontend/src/lib/rectification-agentic/v9/turn-exit.ts +++ b/frontend/src/lib/rectification-agentic/v9/turn-exit.ts @@ -5,6 +5,7 @@ import { } from "./answer-choice.ts"; import { logRectificationDeliveryTurn } from "./delivery-turn-guard.ts"; import type { RectificationRpcClient } from "./tool-service.ts"; +import { reportTurnProgress } from "./turn-instrumentation.ts"; export type RectificationRouteAction = | "opening" @@ -36,45 +37,55 @@ export async function finalizeSuccessfulTurnExit(input: { action: RectificationRouteAction; assistantBodyPresent?: boolean; askedTurnId?: string | null; + /** + * The Agent run already made the same post-turn interview call for this + * turn and found the next focus active (BUG-1047 D5). Re-reading it here + * would return the same "focus active" answer, so go straight to the + * nonterminal repair check. + */ + interviewSettled?: boolean; }): Promise { if (input.action === "read_only") { // Read-only requests must never mutate the interview or create a focus. return; } - try { - const idle = await persistNextInterviewIfIdle({ - accounting: input.accounting, - userId: input.userId, - caseId: input.caseId, - askedTurnId: input.askedTurnId ?? null, - }); - if (idle.terminalNote) { - logRectificationDeliveryTurn({ - trigger: "finalizeSuccessfulTurnExit", + reportTurnProgress("preparing_question"); + if (!input.interviewSettled) { + try { + const idle = await persistNextInterviewIfIdle({ + accounting: input.accounting, + userId: input.userId, caseId: input.caseId, - terminalNote: true, + askedTurnId: input.askedTurnId ?? null, }); - if (!input.askedTurnId && idle.hostNarration) { - try { - await persistExhaustionGateTurn({ - accounting: input.accounting, - userId: input.userId, - caseId: input.caseId, - askedTurnId: input.askedTurnId ?? null, - hostNarration: idle.hostNarration, - }); - } catch (error) { - console.warn( - `[rectification-v9] persist exhaustion gate after turn failed case=${input.caseId} reason=${error instanceof Error ? error.name : "Unknown"}`, - ); + if (idle.terminalNote) { + logRectificationDeliveryTurn({ + trigger: "finalizeSuccessfulTurnExit", + caseId: input.caseId, + terminalNote: true, + }); + if (!input.askedTurnId && idle.hostNarration) { + try { + await persistExhaustionGateTurn({ + accounting: input.accounting, + userId: input.userId, + caseId: input.caseId, + askedTurnId: input.askedTurnId ?? null, + hostNarration: idle.hostNarration, + }); + } catch (error) { + console.warn( + `[rectification-v9] persist exhaustion gate after turn failed case=${input.caseId} reason=${error instanceof Error ? error.name : "Unknown"}`, + ); + } } + return; } - return; + } catch (error) { + console.warn( + `[rectification-v9] persist next interview after turn failed case=${input.caseId} reason=${error instanceof Error ? error.name : "Unknown"}`, + ); } - } catch (error) { - console.warn( - `[rectification-v9] persist next interview after turn failed case=${input.caseId} reason=${error instanceof Error ? error.name : "Unknown"}`, - ); } try { await ensureNonTerminalTurnExit(input); diff --git a/frontend/src/lib/rectification-agentic/v9/turn-instrumentation.ts b/frontend/src/lib/rectification-agentic/v9/turn-instrumentation.ts new file mode 100644 index 00000000..18108f26 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/turn-instrumentation.ts @@ -0,0 +1,132 @@ +/** + * Per-request instrumentation for one rectification turn (BUG-1047). + * + * Two things ride on one AsyncLocalStorage scope that the agent route opens + * around a turn: + * + * - engine call timings: every Python engine request made anywhere inside the + * turn (route preflight, answer-choice, Mastra tools, idle interview) is + * appended with its path, duration and outcome. No request body, birth data + * or engine payload is kept. + * - stage progress: the route supplies a sink that turns stage changes into + * transient `turn.progress` stream events. The stages only move forward, so + * the live row never flips back and forth. + * + * Outside a scope (unit tests, other routes) both calls are no-ops. + */ +import { AsyncLocalStorage } from "node:async_hooks"; + +import { + RECTIFICATION_TURN_PROGRESS_STAGES, + type RectificationTurnProgressStage, +} from "./turn-progress.ts"; + +export { + RECTIFICATION_TURN_PROGRESS_STAGES, + isRectificationTurnProgressStage, + type RectificationTurnProgressStage, +} from "./turn-progress.ts"; + +export type EngineCallTiming = Readonly<{ + path: string; + ms: number; + outcome: "ok" | "busy" | "http_error" | "timeout" | "bad_payload"; +}>; + +type TurnInstrumentationScope = { + readonly startedAt: number; + readonly engineCalls: EngineCallTiming[]; + onProgress?: (stage: RectificationTurnProgressStage) => void; + stageRank: number; +}; + +const storage = new AsyncLocalStorage(); + +/** Only the path is kept; query strings never reach diagnostics. */ +function safeEnginePath(path: string): string { + return path.split("?")[0]?.slice(0, 80) ?? ""; +} + +export type TurnInstrumentation = Readonly<{ + /** Run `body` inside this turn's scope; async work started there inherits it. */ + run(body: () => T): T; + /** Where stage changes go; set once the stream's `send` exists. */ + setProgressSink(sink: ((stage: RectificationTurnProgressStage) => void) | null): void; + engineCalls(): readonly EngineCallTiming[]; + elapsedMs(): number; + /** Advance the stage for this turn; lower or equal stages are ignored. */ + advance(stage: RectificationTurnProgressStage): void; + /** A retried model attempt starts over from the first stage. */ + resetStage(): void; +}>; + +export function createTurnInstrumentation( + input: Readonly<{ onProgress?: (stage: RectificationTurnProgressStage) => void; now?: () => number }> = {}, +): TurnInstrumentation { + const now = input.now ?? Date.now; + const scope: TurnInstrumentationScope = { + startedAt: now(), + engineCalls: [], + onProgress: input.onProgress, + stageRank: -1, + }; + return { + run: (body) => storage.run(scope, body), + setProgressSink: (sink) => { + scope.onProgress = sink ?? undefined; + }, + engineCalls: () => [...scope.engineCalls], + elapsedMs: () => now() - scope.startedAt, + advance: (stage) => advanceScope(scope, stage), + resetStage: () => { + scope.stageRank = -1; + }, + }; +} + +export function runWithTurnInstrumentation( + input: Readonly<{ onProgress?: (stage: RectificationTurnProgressStage) => void; now?: () => number }>, + body: (instrumentation: TurnInstrumentation) => Promise, +): Promise { + const instrumentation = createTurnInstrumentation(input); + return instrumentation.run(() => body(instrumentation)); +} + +function advanceScope(scope: TurnInstrumentationScope, stage: RectificationTurnProgressStage) { + const rank = RECTIFICATION_TURN_PROGRESS_STAGES.indexOf(stage); + if (rank <= scope.stageRank) return; + scope.stageRank = rank; + try { + scope.onProgress?.(stage); + } catch { + // Progress is a live nicety; it never fails the turn. + } +} + +/** A retried model attempt starts the stage sequence over (BUG-1047). */ +export function resetTurnProgress(): void { + const scope = storage.getStore(); + if (scope) scope.stageRank = -1; +} + +/** Report a stage from deep inside the turn (engine client, tools). */ +export function reportTurnProgress(stage: RectificationTurnProgressStage): void { + const scope = storage.getStore(); + if (scope) advanceScope(scope, stage); +} + +export function recordEngineCallTiming(timing: EngineCallTiming): void { + const scope = storage.getStore(); + if (!scope) return; + if (scope.engineCalls.length >= 64) return; + scope.engineCalls.push({ + path: safeEnginePath(timing.path), + ms: Math.max(0, Math.round(timing.ms)), + outcome: timing.outcome, + }); +} + +/** Engine calls recorded so far in the current scope (empty outside one). */ +export function currentEngineCallTimings(): readonly EngineCallTiming[] { + return [...(storage.getStore()?.engineCalls ?? [])]; +} diff --git a/frontend/src/lib/rectification-agentic/v9/turn-intent-classifier.ts b/frontend/src/lib/rectification-agentic/v9/turn-intent-classifier.ts index 4a3b3c0f..3e84bf0e 100644 --- a/frontend/src/lib/rectification-agentic/v9/turn-intent-classifier.ts +++ b/frontend/src/lib/rectification-agentic/v9/turn-intent-classifier.ts @@ -55,12 +55,70 @@ export type ExpectedWriteSignal = "evidence" | "none" | "unknown"; export type TurnIntentClassifierOutcome = "classified" | "unclear" | "classifier_unavailable"; +/** + * Timing facts for one classification (BUG-1047 D4). Durations and counts + * only; never the user message or model output. + */ +export type TurnIntentClassifierDiagnostic = Readonly<{ + outcome: TurnIntentClassifierOutcome; + attempts: number; + timedOutAttempts: number; + elapsedMs: number; +}>; + export type TurnIntentClassifierResult = Readonly<{ classified: RectificationTurnIntent | null; expectedWrite: ExpectedWriteSignal; outcome: TurnIntentClassifierOutcome; + diagnostic?: TurnIntentClassifierDiagnostic; }>; +/** + * Per-attempt ceiling for the intent classifier (BUG-1047 D2). The session + * model and provider thinking stay as they are; an attempt that has not + * answered after 10 s counts as a miss and takes the existing retry, then the + * existing `classifier_unavailable` path. + */ +export const RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS = 10_000; + +class ClassifierAttemptTimeout extends Error { + constructor() { + super("rectification_classifier_attempt_timeout"); + this.name = "ClassifierAttemptTimeout"; + } +} + +async function classifyWithinTimeout( + classify: ClassifyTurnIntent, + model: ResolvedLanguageModel, + input: Parameters[1], + timeoutMs: number, +): Promise { + const controller = new AbortController(); + const outer = input.signal; + const onOuterAbort = () => controller.abort(outer?.reason); + if (outer?.aborted) controller.abort(outer.reason); + else outer?.addEventListener("abort", onOuterAbort, { once: true }); + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => { + const error = new ClassifierAttemptTimeout(); + controller.abort(error); + reject(error); + }, timeoutMs); + }); + try { + // The race guarantees the wait ends even if a provider ignores the signal. + return await Promise.race([ + classify(model, { ...input, signal: controller.signal }), + timeout, + ]); + } finally { + if (timer !== undefined) clearTimeout(timer); + outer?.removeEventListener("abort", onOuterAbort); + } +} + export function expectedWriteFromCollectIntent( classified: RectificationTurnIntent | null, focus?: ConversationFocus | null, @@ -98,27 +156,49 @@ export async function classifyTurnIntentWithRetry( signal?: AbortSignal; }, classify: ClassifyTurnIntent = classifyRectificationTurnIntent, + options: Readonly<{ attemptTimeoutMs?: number }> = {}, ): Promise { const started = Date.now(); + const attemptTimeoutMs = options.attemptTimeoutMs ?? RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS; + let attempts = 0; + let timedOutAttempts = 0; for (let attempt = 0; attempt < 2; attempt += 1) { + attempts += 1; try { - const classified = await classify(model, input); + const classified = await classifyWithinTimeout(classify, model, input, attemptTimeoutMs); if (classified) { + const outcome = turnIntentOutcome(classified); + const diagnostic: TurnIntentClassifierDiagnostic = { + outcome, + attempts, + timedOutAttempts, + elapsedMs: Date.now() - started, + }; + console.info(JSON.stringify({ scope: "RectificationClassifierDiagnostic", ...diagnostic })); return { classified, expectedWrite: expectedWriteFromCollectIntent(classified, input.focus), - outcome: turnIntentOutcome(classified), + outcome, + diagnostic, }; } - } catch { + } catch (error) { + if (error instanceof ClassifierAttemptTimeout) timedOutAttempts += 1; // One retry, then fail-open as unknown / classifier_unavailable. } } - console.warn("rectification_classifier_unavailable", { + const diagnostic: TurnIntentClassifierDiagnostic = { + outcome: "classifier_unavailable", + attempts, + timedOutAttempts, elapsedMs: Date.now() - started, + }; + console.warn("rectification_classifier_unavailable", { + elapsedMs: diagnostic.elapsedMs, attempts: 2, + timedOutAttempts, }); - return { classified: null, expectedWrite: "unknown", outcome: "classifier_unavailable" }; + return { classified: null, expectedWrite: "unknown", outcome: "classifier_unavailable", diagnostic }; } export function optionIdForAnswerClass( diff --git a/frontend/src/lib/rectification-agentic/v9/turn-progress.ts b/frontend/src/lib/rectification-agentic/v9/turn-progress.ts new file mode 100644 index 00000000..b2184db2 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/turn-progress.ts @@ -0,0 +1,18 @@ +/** + * Stages of one typed rectification turn, as shown on the live row + * (BUG-1047 D3). Pure module: safe for the browser bundle. The server side + * that emits them lives in `turn-instrumentation.ts`. + */ +export const RECTIFICATION_TURN_PROGRESS_STAGES = [ + "received", + "recording", + "rescoring", + "preparing_question", +] as const; + +export type RectificationTurnProgressStage = (typeof RECTIFICATION_TURN_PROGRESS_STAGES)[number]; + +export function isRectificationTurnProgressStage(value: unknown): value is RectificationTurnProgressStage { + return typeof value === "string" + && (RECTIFICATION_TURN_PROGRESS_STAGES as readonly string[]).includes(value); +} diff --git a/frontend/src/lib/rectification-surface-state.ts b/frontend/src/lib/rectification-surface-state.ts index 45833546..33d215e5 100644 --- a/frontend/src/lib/rectification-surface-state.ts +++ b/frontend/src/lib/rectification-surface-state.ts @@ -10,6 +10,7 @@ import type { PersistedRectificationTurn } from "../components/conversational-birth-time-rectification.tsx"; import { BOOTSTRAP_PREPARE_TIMEOUT_MS } from "./home-bootstrap.ts"; +import { RECTIFICATION_TURN_PROGRESS_LABELS } from "./rectification-activity-labels.ts"; import { sessionOutcomeAllowsDelivery } from "./rectification-agentic/core/rectification-decision.ts"; import { boardDeclaredTimeLine, @@ -75,6 +76,8 @@ export function rectificationInitialLiveLabel( ): string { if (action === "opening") return RECTIFICATION_OPENING_LIVE_LABEL; if (action === "read_only" && continuationLabel) return continuationLabel; + // A typed answer shows its first stage line the moment it is sent (BUG-1047 D3). + if (action === "message") return RECTIFICATION_TURN_PROGRESS_LABELS.received; return RECTIFICATION_MESSAGE_LIVE_LABEL; } diff --git a/frontend/src/lib/rectification-timeline-adapter.ts b/frontend/src/lib/rectification-timeline-adapter.ts index 0571f53d..e3bb1b92 100644 --- a/frontend/src/lib/rectification-timeline-adapter.ts +++ b/frontend/src/lib/rectification-timeline-adapter.ts @@ -1,7 +1,7 @@ import type { AgentActivityTraceItem } from "./agent-activity-trace.ts"; import type { AgentActivityView } from "./chat-message-view.ts"; import type { ConsultationTimelineKind, ConsultationTimelineRow } from "./consultation-run-timeline.ts"; -import { RECTIFICATION_ANALYZING_LIVE_LABEL } from "./rectification-activity-labels.ts"; +import { RECTIFICATION_ANALYZING_LIVE_LABEL, RECTIFICATION_TURN_PROGRESS_LABELS } from "./rectification-activity-labels.ts"; import type { CompletedActivityReceiptView } from "./rectification-activity-receipt.ts"; import { PUBLIC_RECTIFICATION_METHOD_LABELS } from "./rectification-varga-sentence.ts"; @@ -30,9 +30,31 @@ function methodChips(receipt: CompletedActivityReceiptView | undefined): string[ return chips; } -/** Progress labels are the only ones that name an action still under way. */ +const TURN_PROGRESS_LINES = new Set(Object.values(RECTIFICATION_TURN_PROGRESS_LABELS)); + +/** + * Progress labels are the only ones that name an action still under way: the + * 「正在…」 labels and the typed-answer stage lines (「收到,正在对照你的档案…」, + * BUG-1047). + */ function isProgressLabel(label: string): boolean { - return /^正在/.test(label.trim()); + const trimmed = label.trim(); + return /^正在/.test(trimmed) || TURN_PROGRESS_LINES.has(trimmed); +} + +/** + * The stage line of a typed answer names what the turn is doing, so it is also + * what the tool step currently under way shows (BUG-1047 D3). Finished steps + * keep their own labels. + */ +export function withLiveStepLabel( + trace: readonly AgentActivityTraceItem[], + label: string, +): readonly AgentActivityTraceItem[] { + if (!trace.some((item) => item.kind === "activity" && item.status === "live" && item.label !== label)) return trace; + return trace.map((item) => ( + item.kind === "activity" && item.status === "live" ? { ...item, label } : item + )); } function liveRowKind(phase: AgentActivityView["phase"] | undefined, label: string): ConsultationTimelineKind { diff --git a/frontend/tests/rectification-latency-20260926.test.ts b/frontend/tests/rectification-latency-20260926.test.ts new file mode 100644 index 00000000..7b1e3903 --- /dev/null +++ b/frontend/tests/rectification-latency-20260926.test.ts @@ -0,0 +1,556 @@ +/** + * BUG-1047 (TASK-rectification-latency-20260926): a typed rectification answer + * waited 25–90 s behind an unbounded classifier and four thinking steps with + * no visible progress, and the run diagnostic could not say where the time + * went. + * + * Covered here (Node 20-runnable, no module mocks): + * - D2: each classifier attempt is capped at 10 s and a hang takes the + * existing retry → `classifier_unavailable` path; success is timed too. + * - D3: stage progress is monotonic, reaches the stream as `turn.progress`, + * passes the public allowlist, and never lands in the saved reply or phases. + * - D4: step start/end, provider reasoning tokens, classifier timing and + * engine call durations reach `RectificationRunDiagnostic` without text. + * - D5: `/v5/versions` memo, and the skipped second idle interview call. + * The route-level "first byte before the classifier" test needs module mocks + * and lives in `rectification-latency-route-20260926.test.ts`. + */ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; + +import { + CASE_ID, + RECTIFICATION_SKILL_SHA256, + SESSION_ID, + TURN_ID, + USER_ID, + activeFocusFixture, + conversationSummaryFixture, + dossierFixture, + fakeAccounting, + receiptHandlers, +} from "./rectification-v9-test-support.ts"; +import type { ResolvedLanguageModel } from "../src/mastra/model.ts"; +import { + RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS, + classifyTurnIntentWithRetry, +} from "../src/lib/rectification-agentic/v9/turn-intent-classifier.ts"; +import { + createTurnInstrumentation, + currentEngineCallTimings, + reportTurnProgress, + runWithTurnInstrumentation, + type RectificationTurnProgressStage, +} from "../src/lib/rectification-agentic/v9/turn-instrumentation.ts"; +import { + RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS, + readV9EngineScoringIdentity, +} from "../src/lib/rectification-agentic/v9/engine-client.ts"; +import { safePublicEvent, turnProgressForChunk } from "../src/lib/rectification-agentic/v9/stream-mapping.ts"; +import { + diagnosticStepsFromTimings, + recordStepChunk, + type RectificationStepTiming, +} from "../src/lib/rectification-agentic/v9/run-diagnostic.ts"; +import { runV9AgentTurn, type V9AgentRunOptions } from "../src/lib/rectification-agentic/v9/agent-run.ts"; +import { RECTIFICATION_SKILL_NAME } from "../src/lib/rectification-agentic/v9/case-status.ts"; +import { persistNextInterviewIfIdle } from "../src/lib/rectification-agentic/v9/answer-choice.ts"; +import { finalizeSuccessfulTurnExit } from "../src/lib/rectification-agentic/v9/turn-exit.ts"; +import { RECTIFICATION_TURN_PROGRESS_LABELS } from "../src/lib/rectification-activity-labels.ts"; +import { rectificationInitialLiveLabel } from "../src/lib/rectification-surface-state.ts"; + +const dummyModel = { id: "test-model" } as ResolvedLanguageModel; +const PROGRESS_LINES = Object.values(RECTIFICATION_TURN_PROGRESS_LABELS); + +async function flush(times = 5) { + for (let i = 0; i < times; i += 1) await Promise.resolve(); +} + +// ---------------------------------------------------------------- D2 classifier + +test("D2: a hanging classifier attempt ends at 10 s, retries once, then takes classifier_unavailable", async (t) => { + t.mock.timers.enable({ apis: ["setTimeout"] }); + assert.equal(RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS, 10_000); + const signals: AbortSignal[] = []; + const warnings: unknown[][] = []; + const originalWarn = console.warn; + console.warn = (...args: unknown[]) => { warnings.push(args); }; + try { + let settled = false; + const pending = classifyTurnIntentWithRetry( + dummyModel, + { userMessage: "2016 年 3 月入学", caseStatus: "collecting_evidence" }, + (_model, input) => { + if (input.signal) signals.push(input.signal); + return new Promise(() => {}); // provider never answers + }, + ).then((result) => { settled = true; return result; }); + await flush(); + assert.equal(signals.length, 1); + t.mock.timers.tick(9_999); + await flush(); + assert.equal(settled, false, "no answer before 10 s"); + assert.equal(signals[0].aborted, false); + t.mock.timers.tick(1); + await flush(); + assert.equal(signals[0].aborted, true, "the first attempt is aborted at 10 s"); + assert.equal(signals.length, 2, "exactly one retry starts"); + assert.equal(settled, false); + t.mock.timers.tick(10_000); + const result = await pending; + assert.equal(result.outcome, "classifier_unavailable"); + assert.equal(result.classified, null); + assert.equal(result.expectedWrite, "unknown"); + assert.equal(result.diagnostic?.attempts, 2); + assert.equal(result.diagnostic?.timedOutAttempts, 2); + assert.equal(signals[1].aborted, true); + const line = JSON.stringify(warnings[0]); + assert.match(line, /rectification_classifier_unavailable/); + assert.match(line, /timedOutAttempts/); + assert.doesNotMatch(line, /2016|入学/); + } finally { + console.warn = originalWarn; + } +}); + +test("D2: a first attempt that times out and a second that answers is classified (existing retry)", async (t) => { + t.mock.timers.enable({ apis: ["setTimeout"] }); + let calls = 0; + const pending = classifyTurnIntentWithRetry( + dummyModel, + { userMessage: "没有", caseStatus: "collecting_evidence" }, + async () => { + calls += 1; + if (calls === 1) return new Promise(() => {}); + return { intent: "answer_current_focus", answer_class: "no" } as const; + }, + ); + await flush(); + t.mock.timers.tick(10_000); + const result = await pending; + assert.equal(calls, 2); + assert.equal(result.outcome, "classified"); + assert.equal(result.diagnostic?.timedOutAttempts, 1); +}); + +test("D2/D4: a successful classification logs its timing without the user text", async () => { + const lines: string[] = []; + const originalInfo = console.info; + console.info = (...args: unknown[]) => { lines.push(args.map(String).join(" ")); }; + try { + const result = await classifyTurnIntentWithRetry( + dummyModel, + { userMessage: "2019 年换了工作", caseStatus: "collecting_evidence" }, + async () => ({ intent: "provide_new_evidence", answer_class: null, has_new_dated_event: true }), + ); + assert.equal(result.outcome, "classified"); + assert.equal(result.diagnostic?.attempts, 1); + assert.equal(result.diagnostic?.timedOutAttempts, 0); + assert.equal(typeof result.diagnostic?.elapsedMs, "number"); + } finally { + console.info = originalInfo; + } + const line = lines.find((item) => item.includes("RectificationClassifierDiagnostic")); + assert.ok(line, "success is logged too"); + assert.doesNotMatch(line, /2019|换了工作|provide_new_evidence/); +}); + +test("D2 red line: no model swap, thinking untouched, route uses the default 10 s cap", () => { + const classifier = readFileSync(new URL("../src/lib/rectification-agentic/v9/turn-intent-classifier.ts", import.meta.url), "utf8"); + const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); + // The Agent is still built on the session model with no thinking override. + assert.match(classifier, /model: model\.model,/); + assert.doesNotMatch(classifier, /thinking:|thinkingTokens|providerOptions|modelSettings|agentGenerationSettings/); + assert.doesNotMatch(route, /attemptTimeoutMs/); + assert.equal((route.match(/await classifyTurnIntentWithRetry\(/g) ?? []).length, 4); +}); + +// ---------------------------------------------------------------- D3 progress + +test("D3: stages only move forward, reach the sink, and reset for a retried attempt", async () => { + const seen: RectificationTurnProgressStage[] = []; + await runWithTurnInstrumentation({ onProgress: (stage) => seen.push(stage) }, async (instrumentation) => { + instrumentation.advance("received"); + reportTurnProgress("recording"); + await Promise.resolve(); + reportTurnProgress("received"); // lower: ignored + reportTurnProgress("rescoring"); + reportTurnProgress("rescoring"); // same: ignored + instrumentation.resetStage(); + reportTurnProgress("received"); + reportTurnProgress("preparing_question"); + }); + assert.deepEqual(seen, ["received", "recording", "rescoring", "received", "preparing_question"]); + // Outside a turn scope the calls are silent no-ops. + reportTurnProgress("recording"); + assert.deepEqual(currentEngineCallTimings(), []); +}); + +test("D3: the scope follows async work started inside run(), including a ReadableStream start", async () => { + const seen: string[] = []; + const instrumentation = createTurnInstrumentation(); + const stream = instrumentation.run(() => new ReadableStream({ + async start(controller) { + instrumentation.setProgressSink((stage) => controller.enqueue(stage)); + instrumentation.advance("received"); + await new Promise((resolve) => setTimeout(resolve, 5)); + reportTurnProgress("rescoring"); // e.g. from the engine client, deep in the turn + controller.close(); + }, + })); + const reader = stream.getReader(); + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + seen.push(value); + } + assert.deepEqual(seen, ["received", "rescoring"]); +}); + +test("D3: turn.progress and turn.rejected pass the public allowlist, nothing else rides along", () => { + assert.deepEqual(safePublicEvent({ type: "turn.progress", stage: "recording", text: "x" }), { + type: "turn.progress", + stage: "recording", + }); + assert.equal(safePublicEvent({ type: "turn.progress", stage: "thinking" }), null); + assert.deepEqual( + safePublicEvent({ type: "turn.rejected", httpStatus: 409, code: "stale_question", message: "这道题已经过期", secret: 1 }), + { type: "turn.rejected", httpStatus: 409, code: "stale_question", message: "这道题已经过期" }, + ); + assert.equal(safePublicEvent({ type: "turn.rejected", httpStatus: 200, message: "ok" }), null); + assert.deepEqual( + safePublicEvent({ type: "turn.rejected", httpStatus: 500, code: "Bad Code!", message: "m" }), + { type: "turn.rejected", httpStatus: 500, message: "m" }, + ); +}); + +test("D3: tool events map to the four stages", () => { + const call = (toolName: string) => ({ type: "tool-call", payload: { toolName } }) as never; + const result = (toolName: string) => ({ type: "tool-result", payload: { toolName, result: { ok: true } } }) as never; + assert.equal(turnProgressForChunk(call("rectification-read-case")), null); + assert.equal(turnProgressForChunk(call("rectification-record-evidence-batch")), "recording"); + assert.equal(turnProgressForChunk(result("rectification-record-evidence-batch")), "preparing_question"); + assert.equal(turnProgressForChunk(call("rectification-compare-candidates")), "rescoring"); + assert.equal(turnProgressForChunk(call("rectification-set-focus")), "preparing_question"); + assert.equal(turnProgressForChunk(call("skill")), null); + assert.equal(turnProgressForChunk({ type: "text-delta", payload: { text: "x" } } as never), null); +}); + +test("D3 copy: the four lines, and a typed answer starts on the first one", () => { + assert.deepEqual(RECTIFICATION_TURN_PROGRESS_LABELS, { + received: "收到,正在对照你的档案…", + recording: "正在记下这件事…", + rescoring: "正在重新对照盘面…", + preparing_question: "正在准备下一个问题…", + }); + assert.equal(rectificationInitialLiveLabel("message"), RECTIFICATION_TURN_PROGRESS_LABELS.received); + const voice = readFileSync(new URL("../docs/VOICE.md", import.meta.url), "utf8"); + const design = readFileSync(new URL("../DESIGN.md", import.meta.url), "utf8"); + for (const line of PROGRESS_LINES) { + assert.ok(voice.includes(line), `VOICE lists ${line}`); + assert.ok(design.includes(line), `DESIGN lists ${line}`); + } +}); + +// ---------------------------------------------------------------- D4 run diagnostic + +test("D4: step timeline records start/end, public tools and provider usage per step", () => { + const steps: RectificationStepTiming[] = []; + const isPublic = (name: string) => name.startsWith("rectification-"); + recordStepChunk(steps, { type: "step-start" }, 10, isPublic); + recordStepChunk(steps, { type: "tool-call", payload: { toolName: "rectification-read-case", args: { caseId: CASE_ID } } }, 900, isPublic); + recordStepChunk(steps, { type: "step-finish", payload: { output: { usage: { inputTokens: 1200, outputTokens: 40, reasoningTokens: 700 } } } }, 1_000, isPublic); + recordStepChunk(steps, { type: "step-start" }, 1_050, isPublic); + recordStepChunk(steps, { type: "tool-call", payload: { toolName: "skill" } }, 1_060, isPublic); + recordStepChunk(steps, { type: "step-finish", payload: { output: { usage: { inputTokens: 1300, outputTokens: 20 } } } }, 4_000, isPublic); + const summary = diagnosticStepsFromTimings(steps); + assert.equal(summary.stepCount, 2); + assert.deepEqual(summary.steps.map((step) => [step.startMs, step.endMs, step.tools]), [ + [10, 1_000, ["rectification-read-case"]], + [1_050, 4_000, []], + ]); + assert.equal(summary.reasoningTokens, 700); + assert.equal(summary.inputTokens, 2_500); + assert.equal(summary.toolCallCount, 1); + assert.equal(diagnosticStepsFromTimings([]).reasoningTokens, null, "unreported stays null, not 0"); +}); + +function chunk(type: string, payload?: Record) { + return { type, ...(payload ? { payload } : {}) }; +} + +function fakeAgent(chunks: Array<{ type: string; payload?: Record }>) { + return { + stream: async () => ({ + fullStream: (async function* () { + for (const item of chunks) yield item; + })(), + totalUsage: Promise.resolve({ inputTokens: 10, outputTokens: 20 }), + }), + getSkill: async () => ({ name: RECTIFICATION_SKILL_NAME, instructions: "skill" }), + }; +} + +test("D3/D4: an Agent turn reports stages, a full diagnostic, and never persists a progress line", async () => { + const accounting = fakeAccounting({ + ...receiptHandlers, + get_agentic_rectification_case_dossier: () => dossierFixture(), + append_agentic_rectification_turn: () => ({ turn_id: TURN_ID }), + finalize_agentic_rectification_turn: () => ({ turn_id: TURN_ID, status: "completed", idempotent: false }), + }); + const emitted: Array<{ type: string }> = []; + const stages: RectificationTurnProgressStage[] = []; + const infos: string[] = []; + const originalInfo = console.info; + console.info = (...args: unknown[]) => { infos.push(args.map(String).join(" ")); }; + const options: V9AgentRunOptions = { + userId: USER_ID, + caseId: CASE_ID, + sessionId: SESSION_ID, + requestId: "aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee", + action: "read_only", + message: null, + modelName: "test-model", + accounting: accounting.client, + billing: { + reserve: async () => ({ success: true, status: 200 }), + complete: async () => true, + release: async () => true, + }, + emit: (event) => { emitted.push(event); }, + classifierDiagnostic: { outcome: "classified", attempts: 1, timedOutAttempts: 0, elapsedMs: 4_200 }, + buildAgent: async () => fakeAgent([ + chunk("start"), + chunk("step-start"), + chunk("tool-call", { toolName: "rectification-read-case", args: { caseId: CASE_ID } }), + chunk("tool-result", { toolName: "rectification-read-case" }), + chunk("step-finish", { output: { usage: { inputTokens: 900, outputTokens: 30, reasoningTokens: 512 } } }), + chunk("step-start"), + chunk("tool-call", { toolName: "rectification-record-evidence-batch", args: { caseId: CASE_ID } }), + chunk("tool-result", { toolName: "rectification-record-evidence-batch", result: { accepted: [] } }), + chunk("step-finish", { output: { usage: { inputTokens: 950, outputTokens: 40, reasoningTokens: 256 } } }), + chunk("step-start"), + chunk("text-delta", { text: "记下了。" }), + chunk("step-finish", { output: { usage: { inputTokens: 980, outputTokens: 12 } } }), + chunk("finish"), + ]) as never, + }; + let result; + try { + result = await runWithTurnInstrumentation( + { onProgress: (stage) => stages.push(stage) }, + () => runV9AgentTurn(options), + ); + } finally { + console.info = originalInfo; + } + assert.equal(result.ok, true, JSON.stringify(result)); + assert.deepEqual(stages, ["recording", "preparing_question"]); + // Progress never goes through the persisted phase/answer channel. + assert.equal(emitted.some((event) => event.type === "turn.progress"), false); + const persistedText = JSON.stringify(accounting.calls.map((call) => call.args)); + for (const line of PROGRESS_LINES) assert.ok(!persistedText.includes(line), `not persisted: ${line}`); + assert.equal(persistedText.includes("turn.progress"), false); + const finalize = accounting.calls.findLast((call) => call.fn === "finalize_agentic_rectification_turn"); + assert.equal(finalize?.args.p_assistant_message, "记下了。"); + + const line = infos.find((item) => item.includes("RectificationRunDiagnostic")); + assert.ok(line); + const diagnostic = JSON.parse(line); + assert.equal(diagnostic.stepCount, 3); + assert.equal(diagnostic.reasoningTokens, 768); + assert.equal(diagnostic.inputTokens, 2_830); + assert.equal(diagnostic.toolCallCount, 2); + assert.equal(diagnostic.distinctToolCount, 2); + assert.deepEqual(diagnostic.steps.map((step: { tools: string[] }) => step.tools), [ + ["rectification-read-case"], + ["rectification-record-evidence-batch"], + [], + ]); + for (const step of diagnostic.steps) { + assert.equal(typeof step.startMs, "number"); + assert.ok(step.endMs >= step.startMs); + } + assert.deepEqual(diagnostic.classifier, { outcome: "classified", attempts: 1, timedOutAttempts: 0, elapsedMs: 4_200 }); + assert.deepEqual(diagnostic.engineCalls, []); + assert.equal(diagnostic.attemptNumber, 1); + assert.doesNotMatch(line, /记下了|RECTIFICATION_SKILL|北京|1990/); + assert.ok(!line.includes(RECTIFICATION_SKILL_SHA256)); +}); + +// ---------------------------------------------------------------- D4 engine timings + D5 versions memo + +function isolateIdentityEnv(t: { after(fn: () => void): void }) { + const saved = { + algorithm: process.env.RECTIFICATION_ALGORITHM_VERSION, + policy: process.env.RECTIFICATION_DECISION_POLICY_VERSION, + engine: process.env.RECTIFICATION_ENGINE_VERSION, + }; + delete process.env.RECTIFICATION_ALGORITHM_VERSION; + delete process.env.RECTIFICATION_DECISION_POLICY_VERSION; + delete process.env.RECTIFICATION_ENGINE_VERSION; + t.after(() => { + for (const [key, value] of [ + ["RECTIFICATION_ALGORITHM_VERSION", saved.algorithm], + ["RECTIFICATION_DECISION_POLICY_VERSION", saved.policy], + ["RECTIFICATION_ENGINE_VERSION", saved.engine], + ] as const) { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; + } + }); +} + +const VERSIONS = { algorithm_version: "rectification-v5-matrix-scoring-8", decision_policy_version: "rectification-candidate-policy-v3" }; + +test("D5: /v5/versions is read once per TTL when env does not pin the identity; engine timing is recorded", async (t) => { + isolateIdentityEnv(t); + let fetches = 0; + t.mock.method(globalThis, "fetch", async (url: unknown) => { + fetches += 1; + assert.ok(String(url).endsWith("/api/rectification/v5/versions")); + return Response.json(VERSIONS); + }); + let now = 1_000_000; + t.mock.method(Date, "now", () => now); + const timings = await runWithTurnInstrumentation({}, async () => { + const first = await readV9EngineScoringIdentity(); + const second = await readV9EngineScoringIdentity(); + assert.deepEqual(first, { algorithmVersion: VERSIONS.algorithm_version, policyVersion: VERSIONS.decision_policy_version }); + assert.deepEqual(second, first); + assert.notEqual(second, first, "callers get their own object"); + return currentEngineCallTimings(); + }); + assert.equal(fetches, 1); + assert.deepEqual(timings.map((item) => [item.path, item.outcome]), [["/api/rectification/v5/versions", "ok"]]); + now += RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS - 1; + await readV9EngineScoringIdentity(); + assert.equal(fetches, 1, "still inside the TTL"); + now += 2; + await readV9EngineScoringIdentity(); + assert.equal(fetches, 2, "TTL expired: read again"); +}); + +test("D5: failures and incomplete version payloads are never memoised", async (t) => { + isolateIdentityEnv(t); + let fetches = 0; + t.mock.method(globalThis, "fetch", async () => { + fetches += 1; + return fetches <= 2 ? Response.json({ algorithm_version: VERSIONS.algorithm_version }) : Response.json(VERSIONS); + }); + assert.equal((await readV9EngineScoringIdentity()).policyVersion, null); + assert.equal((await readV9EngineScoringIdentity()).policyVersion, null); + assert.equal((await readV9EngineScoringIdentity()).policyVersion, VERSIONS.decision_policy_version); + assert.equal(fetches, 3); + t.mock.method(globalThis, "fetch", async () => { throw new Error("engine down"); }); + assert.deepEqual(await readV9EngineScoringIdentity(), { algorithmVersion: null, policyVersion: null }, + "a different transport never sees another transport's memo"); +}); + +// ---------------------------------------------------------------- D5 second idle call + +function focusDossier() { + return dossierFixture({ + conversationSummary: conversationSummaryFixture({ + activeFocus: activeFocusFixture({ + intent: "distinguish_candidates", + questionId: "d9:relationship:2023", + targetDomain: "relationship", + expectedAnswerSchema: { + prompt: "2023 年前后,你有没有一段认真开始或结束的关系?", + choice: { + prompt: "2023 年前后,你有没有一段认真开始或结束的关系?", + options: [ + { key: "A", label: "明确发生且时间吻合", answer_class: "yes" }, + { key: "B", label: "发生过但程度较弱", answer_class: "weak_yes" }, + { key: "C", label: "没有这回事", answer_class: "no" }, + { key: "D", label: "这段记不清楚", answer_class: "unsure" }, + ], + }, + }, + }), + }), + }); +} + +function focusAccounting() { + return fakeAccounting({ + ...receiptHandlers, + get_agentic_rectification_case_dossier: () => focusDossier(), + set_agentic_rectification_conversation_focus: (_fn, args) => ({ + focus: { ...activeFocusFixture({ intent: "distinguish_candidates", questionId: String(args.p_question_id) }), asked_turn_id: args.p_asked_turn_id }, + idempotent: true, + }), + }); +} + +test("D5: with the next focus already active, a second idle call is a pure repeat — so skipping it is output-identical", async () => { + const accounting = focusAccounting(); + const first = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); + const writesAfterFirst = accounting.calls.filter((call) => call.fn === "set_agentic_rectification_conversation_focus").map((call) => call.args); + const second = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); + const writesAfterSecond = accounting.calls.filter((call) => call.fn === "set_agentic_rectification_conversation_focus").map((call) => call.args); + assert.equal(first.focusActive, true); + assert.deepEqual(second, first); + // The only write the second call makes is the identical idempotent re-link. + assert.equal(writesAfterSecond.length, writesAfterFirst.length * 2); + assert.deepEqual(writesAfterSecond.slice(writesAfterFirst.length), writesAfterFirst); +}); + +test("D5: the exit gate skips its idle re-read only when told the run settled the interview", async () => { + // Both arms first make the Agent run's own post-turn call, as agent-run.ts does. + const run = async (interviewSettled: boolean) => { + const accounting = focusAccounting(); + const idle = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); + assert.equal(idle.focusActive, true); + await finalizeSuccessfulTurnExit({ + accounting: accounting.client, + userId: USER_ID, + caseId: CASE_ID, + action: "message", + askedTurnId: TURN_ID, + interviewSettled, + }); + return accounting.calls; + }; + const unsettled = await run(false); + const settled = await run(true); + const writes = (calls: typeof settled) => calls + .filter((call) => call.fn !== "get_agentic_rectification_case_dossier" && call.fn !== "get_agentic_rectification_case_compute") + .map((call) => JSON.stringify([call.fn, call.args])); + assert.ok(settled.length < unsettled.length, `${settled.length} < ${unsettled.length}`); + // Same distinct writes either way: the skipped call only re-linked the same focus. + assert.deepEqual([...new Set(writes(settled))].sort(), [...new Set(writes(unsettled))].sort()); +}); + +test("D5 contract: the Agent run reports interviewSettled only from the focus-active exit; the route forwards it", () => { + const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); + assert.match(agentRun, /interviewSettled = interviewIdle\.focusActive === true;/); + assert.match(route, /interviewSettled: result\.interviewSettled === true/); +}); + +// ---------------------------------------------------------------- D3 route shape (source) + +test("D3 contract: a typed message builds the stream first and classifies inside it", () => { + const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); + assert.match(route, /const deferToStream = action === "message";/); + assert.match(route, /const immediateResponse = deferToStream \? null : await computeImmediateResponse\(\);/); + const start = route.indexOf("async start(controller)"); + const inStream = route.slice(start); + const advance = inStream.indexOf('instrumentation.advance("received")'); + const preflight = inStream.indexOf("await computeImmediateResponse()"); + const unfocused = inStream.indexOf("await classifyUnfocusedMessage()"); + const agent = inStream.indexOf("const result = await runV9AgentTurn"); + assert.ok(start > 0 && advance > 0, "first progress line is inside the stream"); + assert.ok(advance < preflight && preflight < unfocused && unfocused < agent); + // Before the stream exists nothing runs the classifier for a message: the + // preflight is only invoked outside the stream when it is not deferred, and + // the unfocused classification is only invoked inside the stream. + assert.equal((route.match(/await computeImmediateResponse\(\)/g) ?? []).length, 2); + assert.equal((route.match(/await classifyUnfocusedMessage\(\)/g) ?? []).length, 1); + assert.equal((inStream.match(/await computeImmediateResponse\(\)/g) ?? []).length, 1); + // The instrumentation scope wraps the stream construction. + assert.match(route, /const body = instrumentation\.run\(\(\) => new ReadableStream\(\{/); +}); diff --git a/frontend/tests/rectification-latency-20260926.test.tsx b/frontend/tests/rectification-latency-20260926.test.tsx new file mode 100644 index 00000000..1619f7a1 --- /dev/null +++ b/frontend/tests/rectification-latency-20260926.test.tsx @@ -0,0 +1,345 @@ +/** + * BUG-1047 D3 (TASK-rectification-latency-20260926), live path in the real + * `RectificationAgenticChat`: a typed answer shows 「收到,正在对照你的档案…」 + * the moment it is sent — before the server has answered at all — and the + * server's `turn.progress` events move the live row through the four stage + * lines. Tool events no longer overwrite it with tool names or 「正在分析…」. + * When the reply settles, no progress line is left in the transcript. A + * `turn.rejected` event is handled exactly like the old non-2xx response. + * + * Harness copied from `rectification-dup-question-20260926.test.tsx`. + */ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { RectificationAgenticChat } from "../src/components/rectification-agentic-chat.tsx"; +import { attachQuestionsToTurns } from "../src/lib/rectification-agentic/v9/turn-question.ts"; +import type { ConversationFocus } from "../src/lib/rectification-agentic/v9/tool-service.ts"; +import { RECTIFICATION_TURN_PROGRESS_LABELS } from "../src/lib/rectification-activity-labels.ts"; +import { createClientLifecycleHarness } from "./react-client-lifecycle-test-support.ts"; + +const CASE_ID = "44444444-4444-4444-8444-444444444444"; +const SESSION_ID = "33333333-3333-4333-8333-333333333333"; +const FOCUS_1 = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaa1"; +const FOCUS_2 = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaa2"; +const STEM_1 = "2023 年前后,有没有换过工作或职责明显变化?"; +const STEM_2 = "2024 年 3 月前后,有没有收入明显变化、大笔支出或欠债?"; + +const OPTIONS = [ + { key: "A" as const, label: "明确发生且时间吻合", answer_class: "yes", role: "primary" }, + { key: "B" as const, label: "发生过但程度较弱", answer_class: "weak_yes", role: "primary" }, + { key: "C" as const, label: "没有这回事", answer_class: "no", role: "primary" }, + { key: "D" as const, label: "这段记不清楚", answer_class: "unsure", role: "primary" }, +]; + +function choiceSchema(prompt: string) { + return { + prompt, + probe_id: `probe:${prompt.slice(0, 4)}`, + choice: { + prompt, + option_a: OPTIONS[0].label, + option_b: OPTIONS[1].label, + option_c: OPTIONS[2].label, + option_d: OPTIONS[3].label, + options: OPTIONS.map(({ key, label, answer_class }) => ({ key, label, answer_class })), + }, + }; +} + +function focus(id: string, prompt: string, overrides: Partial = {}): ConversationFocus { + return { + id, + caseId: CASE_ID, + questionId: `distinguish:${id.slice(-1)}`, + intent: "distinguish_candidates", + targetEvidenceId: null, + targetDomain: "career", + targetKind: null, + expectedAnswerSchema: choiceSchema(prompt), + status: "active", + askedAt: "2026-09-26T08:00:00.000Z", + resolvedAt: null, + askedTurnId: null, + answerOption: null, + ...overrides, + }; +} + +function choiceCard(focusId: string, prompt: string, questionId: string) { + return { + question_id: questionId, + method_id: "candidate_discriminator", + prompt, + choice_mode: "A/B/C/D", + options: OPTIONS, + focus_id: focusId, + case_revision: 3, + probe_id: `probe:${prompt.slice(0, 4)}`, + }; +} + +type RawTurn = { id: string; role: "user" | "assistant"; text: string; status: string }; + +/** The GET route's turn assembly: the real `attachQuestionsToTurns`. */ +function snapshot(rawTurns: RawTurn[], focuses: ConversationFocus[], current: ConversationFocus | null) { + const turns = attachQuestionsToTurns(rawTurns, focuses); + const collect = current && current.expectedAnswerSchema.collect === true; + return { + case: { status: "collecting_evidence" }, + question_source: current ? "focus" : null, + current_question: current + ? { + kind: collect ? "collect_spoken" : "choice", + prompt: current.expectedAnswerSchema.prompt, + focus_id: current.id, + question_id: current.questionId, + } + : null, + choice_card: current && !collect + ? choiceCard(current.id, String(current.expectedAnswerSchema.prompt), current.questionId) + : null, + turns, + }; +} + +function json(body: unknown, status = 200): Response { + return new Response(JSON.stringify(body), { status, headers: { "content-type": "application/json" } }); +} + +type Route = (request: { url: string; method: string; body: Record | null }) => Response | Promise; + +async function mountChat(input: { + initialTurns: RawTurn[]; + initialFocuses: ConversationFocus[]; + initialCurrent: ConversationFocus | null; + route: Route; +}) { + const harness = createClientLifecycleHarness(); + const win = (globalThis as unknown as { window: Record }).window; + win.matchMedia = () => ({ matches: false, addEventListener() {}, removeEventListener() {} }); + win.setInterval = setInterval; + win.clearInterval = clearInterval; + // Host DOM has no layout: the scroll-anchor hook only needs these to exist. + const proto = (harness.container as unknown as { constructor: { prototype: Record } }).constructor.prototype; + Object.assign(proto, { + querySelector: () => null, + querySelectorAll: () => [], + getBoundingClientRect: () => ({ top: 0, bottom: 0, left: 0, right: 0, width: 0, height: 0 }), + compareDocumentPosition: () => 0, + scrollTo() {}, + scrollTop: 0, + scrollHeight: 0, + clientHeight: 0, + }); + const doc = (globalThis as unknown as { document: { createElement: (tag: string) => { style: object } } }).document; + const create = doc.createElement.bind(doc); + const styled = (element: T): T => { + Object.assign(element.style, { setProperty() {}, removeProperty() {}, getPropertyValue: () => "" }); + return element; + }; + doc.createElement = (tag: string) => styled(create(tag)); + styled(harness.container as unknown as { style: object }); + + const originalFetch = globalThis.fetch; + const requests: Array<{ url: string; method: string; body: Record | null }> = []; + globalThis.fetch = (async (resource: RequestInfo | URL, init?: RequestInit) => { + const request = { + url: String(resource), + method: init?.method ?? "GET", + body: typeof init?.body === "string" ? JSON.parse(init.body) as Record : null, + }; + requests.push(request); + return input.route(request); + }) as typeof fetch; + + const initial = snapshot(input.initialTurns, input.initialFocuses, input.initialCurrent); + await harness.render( + {}} + headerSlot={null} + />, + ); + await harness.idle(); + + type Host = ReturnType[number]; + const byTag = (tag: string) => harness.elements().filter((node) => node.tagName === tag); + const attr = (node: Host, name: string) => node.getAttribute(name); + const insideDisabledFieldset = (node: Host): boolean => { + for (let at: Host | null = node; at; at = at.parentNode as Host | null) { + if (at.tagName === "FIELDSET" && at.hasAttribute("disabled")) return true; + } + return false; + }; + // What a sighted reader sees: screen-reader-only nodes (the card's + // ``) are not a second visible stem. + const visibleText = (node: Host): string => { + if ((node.props.className as string | undefined)?.split(/\s+/).includes("sr-only")) return ""; + return node.textContent + node.childNodes.map((child) => visibleText(child as Host)).join(""); + }; + const optionButtons = () => byTag("BUTTON").filter((node) => (node.props.className as string | undefined)?.includes("birth-time-choice-option")); + return { + harness, + requests, + text: () => visibleText(harness.container as Host), + occurrences: (needle: string) => visibleText(harness.container as Host).split(needle).length - 1, + cards: () => byTag("SECTION").filter((node) => (node.props.className as string | undefined)?.includes("rectification-choice-card")), + clickableOptions: () => optionButtons().filter((node) => !insideDisabledFieldset(node)), + selectedOptions: () => optionButtons().filter((node) => attr(node, "data-selected") === "true"), + persistedBlocks: () => harness.elements().filter((node) => attr(node, "data-testid") === "persisted-question"), + async type(message: string) { + const textarea = byTag("TEXTAREA")[0]; + assert.ok(textarea, "composer textarea"); + await harness.event(textarea, "onChange", { target: { value: message } }); + const form = byTag("FORM")[0]; + assert.ok(form, "composer form"); + await harness.event(form, "onSubmit"); + for (let i = 0; i < 5; i += 1) await harness.idle(); + }, + async tap(label: string) { + const button = optionButtons().find((node) => node.text.includes(label) && !insideDisabledFieldset(node)); + assert.ok(button, `clickable option ${label}`); + await harness.event(button, "onClick"); + for (let i = 0; i < 5; i += 1) await harness.idle(); + }, + async close() { + globalThis.fetch = originalFetch; + await harness.close(); + }, + }; +} + +const OPENING: RawTurn = { id: "t0", role: "assistant", text: "眼下按 14:35–15:05 来核对。", status: "completed" }; +const LABELS = RECTIFICATION_TURN_PROGRESS_LABELS; +const ALL_LINES = Object.values(LABELS); + +function controllableStream() { + const encoder = new TextEncoder(); + let controller!: ReadableStreamDefaultController; + const body = new ReadableStream({ start(c) { controller = c; } }); + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + return { + respond: async () => { + await gate; + return new Response(body, { status: 200, headers: { "content-type": "application/x-ndjson; charset=utf-8" } }); + }, + release: () => release(), + push: (event: unknown) => controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)), + end: () => controller.close(), + }; +} + +async function settle(chat: { harness: { idle(): Promise } }) { + for (let i = 0; i < 8; i += 1) { + await chat.harness.idle(); + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +function shown(chat: { text(): string }) { + const text = chat.text(); + return ALL_LINES.filter((line) => text.includes(line)); +} + +test("D3 · the first stage line shows on send (before any server byte), then follows turn.progress; none survives the settle", async () => { + const q1 = focus(FOCUS_1, STEM_1, { askedTurnId: "t0" }); + const stream = controllableStream(); + const reply = "记下了:2023 年 3 月换了工作。"; + const chat = await mountChat({ + initialTurns: [OPENING], + initialFocuses: [q1], + initialCurrent: q1, + route: ({ method }) => { + if (method === "POST") return stream.respond(); + const q2 = focus(FOCUS_2, STEM_2, { askedTurnId: "t2" }); + return json(snapshot([ + OPENING, + { id: "t1", role: "user", text: "2023 年 3 月换了工作", status: "completed" }, + { id: "t2", role: "assistant", text: reply, status: "completed" }, + ], [{ ...q1, status: "resolved" as const }, q2], q2)); + }, + }); + try { + const sentAt = Date.now(); + await chat.type("2023 年 3 月换了工作"); + // The fetch has not resolved yet: this line is on screen with no server byte. + assert.deepEqual(shown(chat), [LABELS.received], chat.text()); + assert.ok(Date.now() - sentAt <= 300, "first progress line within 300 ms of send"); + assert.equal(chat.text().includes("正在处理…"), false); + + stream.release(); + stream.push({ type: "turn.progress", stage: "received" }); + await settle(chat); + assert.deepEqual(shown(chat), [LABELS.received]); + + stream.push({ type: "run.started" }); + stream.push({ type: "tool.activity", tool: "rectification-read-case", status: "started" }); + stream.push({ type: "tool.activity", tool: "rectification-read-case", status: "completed" }); + await settle(chat); + assert.deepEqual(shown(chat), [LABELS.received], "a tool name or 正在分析… does not replace the stage line"); + assert.equal(chat.text().includes("正在分析…"), false); + assert.equal(chat.text().includes("正在读取校正记录…"), false); + + stream.push({ type: "turn.progress", stage: "recording" }); + stream.push({ type: "tool.activity", tool: "rectification-record-evidence-batch", status: "started" }); + await settle(chat); + assert.deepEqual(shown(chat), [LABELS.recording]); + assert.equal(chat.text().includes("正在整理多条事件证据…"), false); + + stream.push({ type: "turn.progress", stage: "rescoring" }); + await settle(chat); + assert.deepEqual(shown(chat), [LABELS.rescoring]); + + stream.push({ type: "tool.activity", tool: "rectification-record-evidence-batch", status: "completed" }); + stream.push({ type: "turn.progress", stage: "preparing_question" }); + await settle(chat); + assert.deepEqual(shown(chat), [LABELS.preparing_question]); + + stream.push({ type: "answer.delta", text: reply }); + stream.push({ type: "run.completed", turnId: "t2" }); + stream.end(); + await settle(chat); + assert.equal(chat.harness.errors.length, 0, String(chat.harness.errors[0] ?? "")); + assert.deepEqual(shown(chat), [], `no progress line after settle:\n${chat.text()}`); + assert.equal(chat.occurrences(reply), 1); + } finally { + await chat.close(); + } +}); + +test("D3 · turn.rejected on the open stream behaves like the old HTTP rejection: row gone, typed mark withdrawn, message shown", async () => { + const q1 = focus(FOCUS_1, STEM_1, { askedTurnId: "t0" }); + const stream = controllableStream(); + const chat = await mountChat({ + initialTurns: [OPENING], + initialFocuses: [q1], + initialCurrent: q1, + route: ({ method }) => (method === "POST" ? stream.respond() : json(snapshot([OPENING], [q1], q1))), + }); + try { + await chat.type("换过"); + assert.deepEqual(shown(chat), [LABELS.received]); + stream.release(); + stream.push({ type: "turn.progress", stage: "received" }); + stream.push({ type: "turn.rejected", httpStatus: 409, code: "stale_question", message: "这道题已经过期,请回答当前问题" }); + await settle(chat); + assert.deepEqual(shown(chat), [], chat.text()); + assert.ok(chat.text().includes("这道题已经过期,请回答当前问题"), chat.text()); + assert.equal(chat.cards().length, 1); + assert.equal(chat.clickableOptions().length, 4, "the typed mark is withdrawn: the card is tappable again"); + assert.equal(chat.occurrences(STEM_1), 1); + } finally { + await chat.close(); + } +}); diff --git a/frontend/tests/rectification-latency-route-20260926.test.ts b/frontend/tests/rectification-latency-route-20260926.test.ts new file mode 100644 index 00000000..0f512b03 --- /dev/null +++ b/frontend/tests/rectification-latency-route-20260926.test.ts @@ -0,0 +1,182 @@ +import assert from "node:assert/strict"; +import { spawnSync } from "node:child_process"; +import { fileURLToPath } from "node:url"; +import test from "node:test"; + +// BUG-1047 D3, route level (TASK-rectification-latency-20260926). The real +// POST /api/rectification/agent handler with fake persistence: a typed answer +// to a pending choice gets its first stream byte — the `turn.progress` +// received line — while the intent classifier is still running, and the +// deterministic reply follows on the same stream after the awaited exit gate. +// The progress line never reaches the persisted assistant message. +// Uses node:test module mocks (Node >= 22.3, same as the other route tests). +test("a typed message streams turn.progress before the classifier answers, then the reply after the exit gate", () => { + const script = String.raw` + import assert from 'node:assert/strict'; + import { mock } from 'node:test'; + import { pathToFileURL } from 'node:url'; + import { CASE_ID, SESSION_ID, USER_ID, TURN_ID, FOCUS_ID, dossierFixture, + conversationSummaryFixture, activeFocusFixture, receiptHandlers } from './tests/rectification-v9-test-support.ts'; + import * as realClassifier from './src/lib/rectification-agentic/v9/turn-intent-classifier.ts'; + import * as realTurnExit from './src/lib/rectification-agentic/v9/turn-exit.ts'; + import { RECTIFICATION_USER_COPY } from './src/lib/rectification-agentic/user-copy.ts'; + import { RECTIFICATION_TURN_PROGRESS_LABELS } from './src/lib/rectification-activity-labels.ts'; + + const prompt = '2023 年前后,有没有换过工作或职责明显变化?'; + const focus = activeFocusFixture({ intent: 'distinguish_candidates', questionId: 'd10:career:2023', + expectedAnswerSchema: { prompt, probe_id: 'probe:career:2023', choice: { prompt, + option_a: '明确发生且时间吻合', option_b: '发生过但程度较弱', option_c: '没有这回事', option_d: '这段记不清楚', + options: [ + { key: 'A', label: '明确发生且时间吻合', answer_class: 'yes' }, + { key: 'B', label: '发生过但程度较弱', answer_class: 'weak_yes' }, + { key: 'C', label: '没有这回事', answer_class: 'no' }, + { key: 'D', label: '这段记不清楚', answer_class: 'unsure' }, + ] } } }); + const calls = []; + const order = []; + const accounting = { rpc: async (fn, args = {}) => { + calls.push({ fn, args }); + let data; + if (fn === 'get_agentic_rectification_case') data = { ...dossierFixture().case, session_id: SESSION_ID, status: 'collecting_evidence' }; + else if (fn === 'get_agentic_rectification_case_dossier') data = dossierFixture({ + conversationSummary: conversationSummaryFixture({ activeFocus: focus }) }); + else if (fn === 'append_agentic_rectification_turn') { order.push('append'); data = { turn_id: TURN_ID, idempotent: false }; } + else if (receiptHandlers[fn]) data = await receiptHandlers[fn](fn, args); + else throw new Error('Unexpected persistence RPC: ' + fn); + return { data: structuredClone(data), error: null }; + }, from: () => { throw new Error('no table access expected'); } }; + + let releaseClassifier; + const classifierGate = new Promise((resolve) => { releaseClassifier = resolve; }); + let classifierStarted = 0; + mock.module('server-only', { namedExports: {} }); + mock.module('@/lib/supabase/server', { namedExports: { createServerSupabaseClient: async () => ({ + auth: { getUser: async () => ({ data: { user: { id: USER_ID } }, error: null }) }, + from: () => ({ select() { return this; }, eq() { return this; }, + maybeSingle: async () => ({ data: { id: SESSION_ID, session_type: 'birth_time_rectification', + agentic_rectification_case_id: CASE_ID, model_id: 'synthetic', model_config_version: 1 }, error: null }) }) + }) } }); + mock.module('@/lib/supabase/admin', { namedExports: { createAdminSupabaseClient: () => accounting } }); + mock.module('@/lib/product-access', { namedExports: { isProductEnabled: async () => true } }); + mock.module('@/lib/feature-flags', { namedExports: { loadRuntimeFeatureFlags: async () => new Map([ + ['rectification_runtime_version', { enabled: true }] + ]) } }); + mock.module('@/lib/model-catalog', { namedExports: { resolveSessionLanguageModel: async () => ({ + id: 'synthetic', configVersion: 1, model: {} + }) } }); + mock.module('@/mastra/agentic-rectification', { namedExports: { getRectificationV9Agent: async () => { + throw new Error('the deterministic path must not build an Agent'); + } } }); + mock.module('@/lib/rectification-agentic/v9/turn-intent-classifier', { namedExports: { + ...realClassifier, + classifyTurnIntentWithRetry: async () => { + classifierStarted += 1; + await classifierGate; + order.push('classified'); + return { classified: null, expectedWrite: 'unknown', outcome: 'classifier_unavailable', + diagnostic: { outcome: 'classifier_unavailable', attempts: 2, timedOutAttempts: 2, elapsedMs: 20000 } }; + }, + } }); + mock.module('@/lib/rectification-agentic/v9/turn-exit', { namedExports: { + ...realTurnExit, + finalizeSuccessfulTurnExit: async (input) => { + order.push('exit:' + input.action + ':' + input.askedTurnId); + }, + } }); + + const { POST } = await import(pathToFileURL(process.cwd() + '/src/app/api/rectification/agent/route.ts').href); + const sentAt = Date.now(); + const response = await POST(new Request('https://example.invalid/api/rectification/agent', { + method: 'POST', headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ caseId: CASE_ID, sessionId: SESSION_ID, requestId: TURN_ID, action: 'message', + message: '换过,2023 年 3 月' }) + })); + assert.equal(response.status, 200); + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + const first = await reader.read(); + const firstMs = Date.now() - sentAt; + const firstLine = decoder.decode(first.value).trim().split('\n')[0]; + assert.deepEqual(JSON.parse(firstLine), { type: 'turn.progress', stage: 'received' }); + assert.ok(firstMs <= 300, 'first byte within 300 ms, was ' + firstMs); + // The classifier is running (or about to) and has not answered: nothing else streamed yet. + await new Promise((resolve) => setTimeout(resolve, 50)); + assert.equal(classifierStarted, 1); + assert.deepEqual(order, []); + releaseClassifier(); + let rest = decoder.decode(first.value).trim().split('\n').slice(1).join('\n'); + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + rest += decoder.decode(value, { stream: true }); + } + const events = rest.trim().split('\n').filter(Boolean).map((line) => JSON.parse(line)); + const types = events.map((event) => event.type); + assert.deepEqual(types.filter((type) => type !== 'turn.progress'), ['answer.delta', 'run.completed'], JSON.stringify(events)); + assert.equal(events.find((event) => event.type === 'answer.delta').text, RECTIFICATION_USER_COPY.classifierUnavailableReply); + assert.equal(events.find((event) => event.type === 'run.completed').turnId, TURN_ID); + // Classified → persisted → awaited exit gate → only then the reply bytes. + assert.deepEqual(order, ['classified', 'append', 'exit:message:' + TURN_ID]); + const append = calls.find((call) => call.fn === 'append_agentic_rectification_turn'); + assert.equal(append.args.p_assistant_message, RECTIFICATION_USER_COPY.classifierUnavailableReply); + for (const line of Object.values(RECTIFICATION_TURN_PROGRESS_LABELS)) { + assert.equal(JSON.stringify(calls).includes(line), false, 'progress line persisted: ' + line); + } + console.log(JSON.stringify({ firstMs, types: ['turn.progress', ...types] })); + `; + const result = spawnSync(process.execPath, ["--experimental-test-module-mocks", "--import", "tsx", "--input-type=module", "--eval", script], { + cwd: fileURLToPath(new URL("../", import.meta.url)), encoding: "utf8", timeout: 60_000, + }); + assert.equal(result.status, 0, result.stderr + result.stdout); + console.log(result.stdout.trim()); +}); + +test("a rejected typed message arrives as turn.rejected on the stream with the old status, code and message", () => { + const script = String.raw` + import assert from 'node:assert/strict'; + import { mock } from 'node:test'; + import { pathToFileURL } from 'node:url'; + import { CASE_ID, SESSION_ID, USER_ID, TURN_ID, dossierFixture, receiptHandlers } from './tests/rectification-v9-test-support.ts'; + const accounting = { rpc: async (fn, args = {}) => { + if (fn === 'get_agentic_rectification_case') return { data: { ...dossierFixture().case, session_id: SESSION_ID, status: 'collecting_evidence' }, error: null }; + if (fn === 'get_agentic_rectification_case_dossier') return { data: null, error: { message: 'agentic_rectification_case_terminal' } }; + if (receiptHandlers[fn]) return { data: await receiptHandlers[fn](fn, args), error: null }; + throw new Error('Unexpected persistence RPC: ' + fn); + } }; + mock.module('server-only', { namedExports: {} }); + mock.module('@/lib/supabase/server', { namedExports: { createServerSupabaseClient: async () => ({ + auth: { getUser: async () => ({ data: { user: { id: USER_ID } }, error: null }) }, + from: () => ({ select() { return this; }, eq() { return this; }, + maybeSingle: async () => ({ data: { id: SESSION_ID, session_type: 'birth_time_rectification', + agentic_rectification_case_id: CASE_ID, model_id: 'synthetic', model_config_version: 1 }, error: null }) }) + }) } }); + mock.module('@/lib/supabase/admin', { namedExports: { createAdminSupabaseClient: () => accounting } }); + mock.module('@/lib/product-access', { namedExports: { isProductEnabled: async () => true } }); + mock.module('@/lib/feature-flags', { namedExports: { loadRuntimeFeatureFlags: async () => new Map([ + ['rectification_runtime_version', { enabled: true }] + ]) } }); + mock.module('@/lib/model-catalog', { namedExports: { resolveSessionLanguageModel: async () => ({ + id: 'synthetic', configVersion: 1, model: {} + }) } }); + mock.module('@/mastra/agentic-rectification', { namedExports: { getRectificationV9Agent: async () => { + throw new Error('no Agent expected'); + } } }); + const { POST } = await import(pathToFileURL(process.cwd() + '/src/app/api/rectification/agent/route.ts').href); + const response = await POST(new Request('https://example.invalid/api/rectification/agent', { + method: 'POST', headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ caseId: CASE_ID, sessionId: SESSION_ID, requestId: TURN_ID, action: 'message', message: '换过' }) + })); + assert.equal(response.status, 200); + const events = (await response.text()).trim().split('\n').map((line) => JSON.parse(line)); + assert.deepEqual(events, [ + { type: 'turn.progress', stage: 'received' }, + { type: 'turn.rejected', httpStatus: 409, code: 'case_terminal', message: '该校正已结束,不能继续修改' }, + ]); + console.log(JSON.stringify(events.map((event) => event.type))); + `; + const result = spawnSync(process.execPath, ["--experimental-test-module-mocks", "--import", "tsx", "--input-type=module", "--eval", script], { + cwd: fileURLToPath(new URL("../", import.meta.url)), encoding: "utf8", timeout: 60_000, + }); + assert.equal(result.status, 0, result.stderr + result.stdout); + console.log(result.stdout.trim()); +}); diff --git a/frontend/tests/rectification-surface-state.test.ts b/frontend/tests/rectification-surface-state.test.ts index 2496362f..80fce58a 100644 --- a/frontend/tests/rectification-surface-state.test.ts +++ b/frontend/tests/rectification-surface-state.test.ts @@ -215,7 +215,10 @@ test("readonly range copy never says 收窄 or 才会变", () => { test("live-row labels follow the action that started the turn", () => { assert.equal(rectificationInitialLiveLabel("opening"), "正在读取你的出生资料,准备第一个问题…"); - assert.equal(rectificationInitialLiveLabel("message"), "正在处理…"); + // 原值: "正在处理…" + // 新值: "收到,正在对照你的档案…" + // 原因: BUG-1047 D3,打字回答发出即显示第一句阶段进度,替代一直不变的「正在处理…」 + assert.equal(rectificationInitialLiveLabel("message"), "收到,正在对照你的档案…"); assert.equal(rectificationInitialLiveLabel("read_only", "正在记录本次选择…"), "正在记录本次选择…"); assert.equal(rectificationInitialLiveLabel("read_only"), "正在处理…"); assert.equal(rectificationAdoptingLabel("04:53"), "正在采用 04:53…");