fix(rectification): stream-first typed turns, 10 s classifier cap, stage progress, run timings (BUG-1047)

- 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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8
This commit is contained in:
Jesse_Chen
2026-09-26 11:45:47 +08:00
co-authored by Claude Opus 5.5
parent fe76d54795
commit ab0f01a9cf
20 changed files with 1913 additions and 75 deletions
+110 -13
View File
@@ -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<Response | null> => {
let classifierDiagnostic: TurnIntentClassifierDiagnostic | null = null;
const computeImmediateResponse = async (): Promise<Response | null> => {
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<string, unknown>) => 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<string, unknown>);
}
};
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<Uint8Array>({
// 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<Uint8Array>({
async start(controller) {
let closed = false;
const send = (event: Record<string, unknown>) => {
@@ -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: {
@@ -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);
@@ -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<Record<PublicRecti
preparing_result: "正在准备结果…",
};
/**
* Stage lines for a typed answer (BUG-1047 D3). The first one shows the moment
* the reader sends; the server's `turn.progress` events move it forward. They
* live only on the live row and are never part of the saved reply.
*/
export const RECTIFICATION_TURN_PROGRESS_LABELS: Readonly<Record<RectificationTurnProgressStage, string>> = {
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 = "正在分析…";
@@ -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<V9Agen
|| !shouldAutoRetry(outcome.errorCode ?? "run_failed", signal, outcome.status)
|| attemptNumber === MAX_ATTEMPTS) break;
await emit({ type: "attempt.reset" });
resetTurnProgress();
reportTurnProgress("received");
}
const outcome = finalOutcome ?? {
@@ -665,9 +689,14 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise<V9Agen
const answerText = outcome.answerText;
let spokenAnswer = answerText;
let interviewIdle: Awaited<ReturnType<typeof persistNextInterviewIfIdle>> | 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<V9Agen
toolsUsed: outcome.toolsUsed,
errorCode: null,
previousFocusId,
interviewSettled,
};
async function streamAttempt(
@@ -730,6 +760,9 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise<V9Agen
attemptId: string,
previousErrorCode: string | null,
): Promise<AttemptOutcome> {
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<V9Agen
try {
for await (const chunk of result.fullStream) {
recordStepChunk(stepTimings, chunk, Date.now() - attemptStartedAt, isPublicRectificationToolName);
const progressStage = turnProgressForChunk(chunk as never);
if (progressStage) reportTurnProgress(progressStage);
const rawToolName = typeof chunk.payload?.toolName === "string" ? chunk.payload.toolName : "";
// Identical public tool-call + args are idempotent. Throwing
// `repeated_tool_call` (BUG-368 P0-3) aborted the turn after
@@ -1237,23 +1273,31 @@ export async function runV9AgentTurn(options: V9AgentRunOptions): Promise<V9Agen
} finally {
clearTimeout(timeout);
signal?.removeEventListener("abort", onAbort);
console.info(JSON.stringify({
scope: "RectificationRunDiagnostic",
const stepSummary = diagnosticStepsFromTimings(stepTimings);
const diagnostic: RectificationRunDiagnostic = {
runId: attemptId,
modelId: options.modelName,
finishReason: finishReason ?? "unknown",
inputTokens: null,
reasoningTokens: null,
outputTokens: null,
stepCount: toolsUsed.size,
toolCallCount: toolsUsed.size,
inputTokens: stepSummary.inputTokens,
reasoningTokens: stepSummary.reasoningTokens,
outputTokens: stepSummary.outputTokens,
stepCount: stepSummary.stepCount,
toolCallCount: stepSummary.toolCallCount,
distinctToolCount: toolsUsed.size,
readCasePayloadBytes: null,
elapsedMs: Date.now() - startedAt,
attemptNumber,
attemptStartMs: attemptStartedAt - startedAt,
attemptElapsedMs: Date.now() - attemptStartedAt,
steps: stepSummary.steps,
lastCompletedTool: lastCompletedPublicTool(toolTerminalStatus),
stateMutationCommitted: publicWriteToolCompleted(toolTerminalStatus) || hostFallbackUsed,
expectedWrite: options.expectedWrite ?? null,
collectIntent: options.collectIntent ?? null,
}));
classifier: options.classifierDiagnostic ?? null,
engineCalls: currentEngineCallTimings().slice(engineCallsBefore),
};
console.info(JSON.stringify({ scope: "RectificationRunDiagnostic", ...diagnostic }));
}
}
@@ -1634,7 +1634,14 @@ export async function persistNextInterviewIfIdle(input: {
askedTurnId?: string | null;
narrateAdopt?: AdoptNarrationWriter;
userStopped?: boolean;
}): Promise<{ persisted: boolean; choiceReady: boolean; hostNarration: string | null; terminalNote?: boolean }> {
}): 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 {
@@ -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<Record<string, unknown>> {
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<Record<string, unknown>> {
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<string, unknown>;
} 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<object, Map<string, VersionsMemo>>();
async function readNativeScoringIdentity(): Promise<LiveEngineScoringIdentity> {
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<LiveEngineScoringIdentity> {
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)
@@ -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<string, unknown> } } | 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<Record<string, string>> = {
answer_truncated: "模型输出达到上限,状态已记录。",
length: "模型输出达到上限,状态已记录。",
@@ -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<PublicErrorCode>([
"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<PublicRectificationTool>([
"rectification-record-evidence-batch",
"rectification-propose-evidence",
"rectification-confirm-evidence",
"rectification-revise-evidence",
]);
const RESCORING_TOOLS = new Set<PublicRectificationTool>([
"rectification-compare-candidates",
"rectification-read-diagnostics",
]);
const QUESTION_TOOLS = new Set<PublicRectificationTool>([
"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;
}
@@ -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<void> {
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);
@@ -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<TurnInstrumentationScope>();
/** 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<T>(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<T>(
input: Readonly<{ onProgress?: (stage: RectificationTurnProgressStage) => void; now?: () => number }>,
body: (instrumentation: TurnInstrumentation) => Promise<T>,
): Promise<T> {
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 ?? [])];
}
@@ -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<ClassifyTurnIntent>[1],
timeoutMs: number,
): Promise<RectificationTurnIntent | null> {
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<typeof setTimeout> | undefined;
const timeout = new Promise<never>((_, 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<TurnIntentClassifierResult> {
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(
@@ -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);
}
@@ -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;
}
@@ -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<string>(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 {