diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-attempt.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-attempt.ts new file mode 100644 index 00000000..4e8f070d --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-attempt.ts @@ -0,0 +1,634 @@ +/** + * One attempt of a rectification agent turn: build the agent, verify the + * bound Skill, stream the model with the read-case-first tool gate, publish + * tool activity and phases, stream the spoken answer, apply the host + * fallbacks, check the evidence write and the Case invariants, and log the + * run diagnostic. Moved out of runV9AgentTurn's nested `streamAttempt` + * unchanged (TASK-rectification-code-split-20260926); the retry constraint + * reads the `previousErrorCode` argument, which the loop passes as the + * previous attempt's error (the nested version read the same value from the + * enclosing `lastAttemptError`). + */ +import { resolveRectificationStepBudget } from "./step-budget"; +import { insertV9SkillRunReceipt, loadV9CaseDossier, type V9CaseDossier } from "./tool-service"; +import { RECTIFICATION_AGENT_TOOLS } from "./public-receipt"; +import { agentGenerationSettings, promptCacheUsage } from "../../agent-generation-settings.ts"; +import { toAgentModelFinishReason } from "../../agent-observability.ts"; +import { decideFromDossier } from "./decision-from-dossier"; +import { parseAgentChoiceCopy, isPersistedFocusId } from "./choice-card"; +import { + RECTIFICATION_USER_COPY, + isAcceptableOpeningBody, + openingRangeFromCandidateRange, + openingSpokenBody, + withCompareFailedRetryNotice, + withRangeChangedAfterEvidence, + withRescoreSkippedNotice, +} from "../user-copy"; +import { attemptTimeoutForRemainingBudget, canStartRetryAttempt } from "../../rectification-run-budget.ts"; +import { stripQuestionSentences, stripVerbalWindowChange } from "./collect-prompt"; +import { focusSpokenPrompt } from "./turn-question"; +import { previousInferenceFromReceipt } from "../core/compose-receipt.ts"; +import { + activityChangedFromTool, + mapStreamChunkToActivity, + mapStreamChunkToPhase, + streamToolNames, + isPublicRectificationToolName, + turnProgressForChunk, + type PublicStreamEvent, +} from "./stream-mapping"; +import { + diagnosticStepsFromTimings, + mapModelFinishToErrorCode, + recordStepChunk, + type RectificationRunDiagnostic, + type RectificationStepTiming, +} from "./run-diagnostic"; +import { currentEngineCallTimings, reportTurnProgress } from "./turn-instrumentation.ts"; +import { + batchResultFromToolChunk, + batchRescoreFailed, + composeHostFallbackNarration, + lastCompletedPublicTool, + publicWriteToolCompleted, + answerClaimsEvidenceRecorded, + turnExpectsEvidenceWrite, +} from "./host-fallback"; +import { applyStepAnswerChunk, createStepAnswerState, flushStepAnswerOnStreamFinish } from "./step-answer"; +import { + type AttemptOutcome, + MAX_ATTEMPTS, + streamFinishReason, + publicToolCallKey, + v9TurnReceipts, +} from "./agent-run-support"; +import { buildAgentMessages } from "./agent-run-messages"; +import type { V9TimedTurn } from "./agent-run-prepare"; + +export async function streamV9Attempt( + turn: V9TimedTurn, + attemptNumber: number, + attemptId: string, + previousErrorCode: string | null, +): Promise { + const { + options, + userId, + caseId, + action, + accounting, + buildAgent, + emit, + signal, + skillName, + dossier, + skillPackage, + birthTimeClue, + turnId, + startedAt, + minRetryAttemptMs, + attemptCapMs, + runDeadlineAt, + } = turn; + const { failedAttempt, persistCommittedPhase } = v9TurnReceipts(turn); + const attemptStartedAt = Date.now(); + const stepTimings: RectificationStepTiming[] = []; + const engineCallsBefore = currentEngineCallTimings().length; + const agent = await buildAgent(turnId, skillPackage, attemptId); + let frameworkSkill: unknown = null; + try { + frameworkSkill = await (agent as unknown as { getSkill(name: string): Promise }).getSkill(skillName); + } catch { + frameworkSkill = null; + } + if (!frameworkSkill) { + return failedAttempt(attemptId, "skill_not_loaded"); + } + + const rawSkillInstructions = (frameworkSkill as { instructions?: unknown }).instructions; + const skillInstructions = typeof rawSkillInstructions === "string" + ? rawSkillInstructions.trim() + : ""; + if (!skillInstructions) { + return failedAttempt(attemptId, "skill_not_loaded"); + } + + const messages = buildAgentMessages( + options, + attemptNumber, + dossier, + skillInstructions, + birthTimeClue, + previousErrorCode, + ); + const maxSteps = resolveRectificationStepBudget(action); + const abortController = new AbortController(); + const onAbort = () => abortController.abort(); + signal?.addEventListener("abort", onAbort, { once: true }); + let timedOut = false; + const timeout = setTimeout(() => { + timedOut = true; + abortController.abort(); + }, attemptTimeoutForRemainingBudget(runDeadlineAt - Date.now(), attemptCapMs)); + + let skillBound = true; + let caseLoaded = false; + let intentClassified = false; + let streamFailed = false; + let finished = false; + let finishReason: ReturnType | null = null; + let answerText = ""; + const answerDeltas: string[] = []; + const phases: string[] = []; + const toolsUsed = new Set(); + const events: PublicStreamEvent[] = []; + const toolTerminalStatus = new Map(); + let batchToolResult: unknown = null; + let hostFallbackUsed = false; + const emittedKeys = new Set(); + const emittedActivities = new Set(); + const repeatedCalls = new Map(); + let phaseSequence = 0; + const rangeBeforeCompare = previousInferenceFromReceipt( + dossier.latestResult?.decisionReceipt ?? null, + )?.credible_range ?? null; + + const recordPhase = async (phase: string, tool: string | null = null) => { + if ( + phase === "answer.delta" + || phase === "thinking.delta" + || phase === "activity.changed" + || phase === "choice.applied" + || phase === "attempt.reset" + || emittedKeys.has(`${phase}:${tool ?? ""}`) + ) return; + emittedKeys.add(`${phase}:${tool ?? ""}`); + phases.push(phase); + phaseSequence += 1; + await persistCommittedPhase(phase, tool, attemptId, phaseSequence); + }; + + const publish = async (event: PublicStreamEvent) => { + events.push(event); + await emit(event); + }; + + try { + await recordPhase("run.started"); + await insertV9SkillRunReceipt( + accounting, + userId, + caseId, + turnId, + attemptId, + "turn", + skillPackage, + ); + await recordPhase("skill.bound"); + await publish({ type: "skill.bound" }); + emittedKeys.add("event:skill.bound::"); + + const generation = agentGenerationSettings(options.generationModel, { + thinking: "enabled", + answerTokens: 8_192, + thinkingTokens: 8_192, + }); + const result = await (agent as unknown as { + stream( + messages: unknown[], + streamOptions: { + maxSteps: number; + abortSignal: AbortSignal; + modelSettings?: { maxOutputTokens?: number }; + providerOptions?: Record; + prepareStep: (input: { stepNumber: number }) => { + activeTools: string[]; + toolChoice: "auto"; + }; + }, + ): Promise<{ + fullStream: AsyncIterable<{ + type: string; + payload?: { + toolName?: unknown; + text?: unknown; + args?: unknown; + error?: unknown; + stepResult?: { reason?: unknown }; + reason?: unknown; + }; + object?: unknown; + }>; + totalUsage?: Promise>; + }>; + }).stream(messages, { + maxSteps, + abortSignal: abortController.signal, + ...generation, + // Thinking-mode providers reject named/required tool_choice. Restrict + // the first step to read-case and keep tool_choice auto; the runner + // still refuses any other public tool before case.loaded. + prepareStep: ({ stepNumber }) => stepNumber === 0 + ? { + activeTools: ["rectification-read-case"], + toolChoice: "auto", + } + : { + activeTools: [...RECTIFICATION_AGENT_TOOLS], + toolChoice: "auto", + }, + }); + + const stepAnswer = createStepAnswerState(); + + let spokenRaw = ""; + let visibleEmitted = ""; + + const emitVisibleSpoken = async (visible: string) => { + if (!caseLoaded) return; + if (visible === visibleEmitted) { + answerText = visible; + return; + } + if (visible.startsWith(visibleEmitted)) { + const growth = visible.slice(visibleEmitted.length); + if (!growth) return; + answerText = visible; + answerDeltas.push(growth); + visibleEmitted = visible; + await emit({ type: "answer.delta", text: growth }); + return; + } + answerText = visible; + answerDeltas.push(visible); + visibleEmitted = visible; + await emit({ type: "answer.delta", text: visible, replace: true }); + }; + + const publishSpokenStep = async (pieces: readonly string[], live = false) => { + const joined = pieces.join(""); + if (!joined) return; + if (!caseLoaded) return; + const spoken = live ? joined : joined.trim(); + if (!spoken) return; + spokenRaw += spoken; + await emitVisibleSpoken(spokenRaw); + }; + + const retractSpoken = async () => { + const hadVisible = Boolean(visibleEmitted || answerText); + spokenRaw = ""; + visibleEmitted = ""; + answerText = ""; + answerDeltas.length = 0; + if (hadVisible && caseLoaded) { + await emit({ type: "answer.delta", text: "", replace: true }); + } + }; + + const applyHostFallback = async (): Promise => { + if (answerText.trim()) return false; + if (toolTerminalStatus.get("rectification-record-evidence-batch") !== "completed") { + return false; + } + const spoken = composeHostFallbackNarration(batchToolResult ?? {}); + if (!spoken) return false; + hostFallbackUsed = true; + await recordPhase("answer.host_fallback"); + await publish({ type: "answer.host_fallback" }); + await emitVisibleSpoken(spoken); + return true; + }; + + 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 + // evidence / diagnostics / compare had already committed, because + // the model often re-issued compare with the same caseId. Bound + // loops with maxSteps / timeout instead; do not attempt.reset. + let skipDuplicateToolCallReceipt = false; + if (chunk.type === "tool-call") { + if (rawToolName === "rectification-read-case" && !skillBound) { + throw new Error("skill_not_bound"); + } + if (isPublicRectificationToolName(rawToolName) && rawToolName !== "rectification-read-case" && !caseLoaded) { + throw new Error("case_not_loaded"); + } + if (isPublicRectificationToolName(rawToolName)) { + const key = publicToolCallKey(rawToolName, chunk.payload); + const count = (repeatedCalls.get(key) ?? 0) + 1; + repeatedCalls.set(key, count); + skipDuplicateToolCallReceipt = count > 1; + } + } + + const stepEffect = applyStepAnswerChunk( + stepAnswer, + chunk, + isPublicRectificationToolName, + ); + if (stepEffect.kind === "live") await publishSpokenStep([stepEffect.text], true); + if (stepEffect.kind === "publish") await publishSpokenStep(stepEffect.pieces); + if (stepEffect.kind === "retract") await retractSpoken(); + + const activityEvent = skipDuplicateToolCallReceipt + ? null + : mapStreamChunkToActivity(chunk as never); + if (activityEvent) { + await publish(activityEvent); + if (activityEvent.status === "started") { + const changed = activityChangedFromTool(activityEvent.tool); + if (!emittedActivities.has(changed.activity)) { + emittedActivities.add(changed.activity); + await publish(changed); + } + } + if (activityEvent.status === "failed" || activityEvent.status === "completed") { + toolTerminalStatus.set(activityEvent.tool, activityEvent.status); + } + } + const capturedBatch = batchResultFromToolChunk(chunk as never); + if (capturedBatch != null) batchToolResult = capturedBatch; + const phaseEvent = skipDuplicateToolCallReceipt + ? null + : mapStreamChunkToPhase(chunk as never); + if (phaseEvent) { + if (phaseEvent.type === "skill.bound" && !skillBound) { + skillBound = true; + await insertV9SkillRunReceipt( + accounting, + userId, + caseId, + turnId, + attemptId, + "turn", + skillPackage, + ); + } + if (phaseEvent.type === "case.loaded" && !skillBound) { + throw new Error("skill_not_bound"); + } + await recordPhase(phaseEvent.type, phaseEvent.tool ?? null); + const key = `${phaseEvent.type}:${phaseEvent.tool ?? ""}:${(phaseEvent.methods ?? []).join(",")}`; + if (!emittedKeys.has(`event:${key}`)) { + emittedKeys.add(`event:${key}`); + await publish(phaseEvent); + } + if (phaseEvent.type === "case.loaded") { + caseLoaded = true; + if (!intentClassified) { + intentClassified = true; + await recordPhase("intent.classified"); + await publish({ type: "intent.classified" }); + } + } + } + for (const toolName of streamToolNames(chunk as never)) toolsUsed.add(toolName); + if (chunk.type === "error" || chunk.type === "abort") streamFailed = true; + if (chunk.type === "finish") { + finished = true; + finishReason = streamFinishReason(chunk); + } + } + } catch (error) { + if (!timedOut && !abortController.signal.aborted) throw error; + streamFailed = true; + } + + if (finished && !streamFailed) { + const flushReason = finishReason === "length" + ? "length" + : finishReason === "stop" || finishReason === "unknown" || finishReason === null + ? "stop" + : finishReason; + const flushed = flushStepAnswerOnStreamFinish(stepAnswer, flushReason); + if (flushed.kind === "publish") await publishSpokenStep(flushed.pieces); + } + + if (toolTerminalStatus.get("rectification-compare-candidates") === "failed" && !batchRescoreFailed(batchToolResult)) { + await emitVisibleSpoken(withCompareFailedRetryNotice(answerText)); + } + + const completeAttempt = async (settleBilling = true): Promise => { + let inputTokens = 0; + let outputTokens = 0; + let cache: ReturnType = null; + try { + const raw = await (result.totalUsage ?? Promise.resolve({ inputTokens: 0, outputTokens: 0 })); + inputTokens = Math.max(0, Math.trunc(typeof raw.inputTokens === "number" ? raw.inputTokens : 0)); + outputTokens = Math.max(0, Math.trunc(typeof raw.outputTokens === "number" ? raw.outputTokens : 0)); + cache = promptCacheUsage(raw); + } catch { + // Timeout/abort can leave provider usage unread. + } + await recordPhase("answer.composed"); + await publish({ type: "answer.composed" }); + return { + ok: true, + status: "completed", + errorCode: null, + usage: { inputTokens, outputTokens, ...(cache ? { cache } : {}) }, + answerText, + answerDeltas, + phases, + toolsUsed: [...toolsUsed], + events, + skillBound, + caseLoaded, + attemptId, + settleBilling, + }; + }; + + if (!skillBound) return failedAttempt(attemptId, "skill_not_loaded"); + if (!caseLoaded) return failedAttempt(attemptId, "case_not_loaded"); + const mapped = mapModelFinishToErrorCode({ + finishReason, + aborted: abortController.signal.aborted, + timedOut, + answerText, + stepCount: toolsUsed.size, + maxSteps, + }); + if (mapped === "run_timeout") return failedAttempt(attemptId, "run_timeout"); + if (streamFailed || abortController.signal.aborted) return failedAttempt(attemptId, mapped ?? "stream_aborted"); + if (!finished) return failedAttempt(attemptId, mapped ?? "stream_unfinished"); + if (mapped === "answer_truncated") { + if (!await applyHostFallback()) { + return { + ok: false, + status: "failed", + errorCode: "answer_truncated", + usage: { inputTokens: 0, outputTokens: 0 }, + answerText, + answerDeltas, + phases, + toolsUsed: [...toolsUsed], + events, + skillBound, + caseLoaded, + attemptId, + }; + } + } else if (mapped === "max_steps" || mapped === "provider_error") { + if (mapped !== "max_steps" || !await applyHostFallback()) { + return failedAttempt(attemptId, mapped); + } + } else if (!answerText.trim() && !await applyHostFallback()) { + const couldRetry = !publicWriteToolCompleted(toolTerminalStatus) && attemptNumber < MAX_ATTEMPTS; + const hasRetryBudget = canStartRetryAttempt(runDeadlineAt - Date.now(), minRetryAttemptMs); + if (couldRetry && hasRetryBudget) { + return { + ok: false, + status: "retryable", + errorCode: "empty_stream", + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [...phases], + toolsUsed: [...toolsUsed], + events, + skillBound, + caseLoaded, + attemptId, + }; + } + if (couldRetry && !hasRetryBudget) { + hostFallbackUsed = true; + await recordPhase("answer.host_fallback"); + await publish({ type: "answer.host_fallback" }); + await emitVisibleSpoken( + turnExpectsEvidenceWrite(action, options.expectedWrite) + ? RECTIFICATION_USER_COPY.evidenceNotRecorded + : RECTIFICATION_USER_COPY.hostNarrationFallback, + ); + return completeAttempt(false); + } + return { + ok: false, + status: "failed", + errorCode: "empty_stream", + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [...phases], + toolsUsed: [...toolsUsed], + events, + skillBound, + caseLoaded, + attemptId, + }; + } + const needsWrite = turnExpectsEvidenceWrite(action, options.expectedWrite) + || (action !== "opening" && action !== "read_only" && answerClaimsEvidenceRecorded(answerText)); + if (needsWrite && !publicWriteToolCompleted(toolTerminalStatus) && !hostFallbackUsed) { + await retractSpoken(); + if ( + attemptNumber < MAX_ATTEMPTS + && canStartRetryAttempt(runDeadlineAt - Date.now(), minRetryAttemptMs) + ) { + return { + ok: false, + status: "retryable", + errorCode: "evidence_not_written", + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [...phases], + toolsUsed: [...toolsUsed], + events, + skillBound, + caseLoaded, + attemptId, + }; + } + hostFallbackUsed = true; + await recordPhase("answer.host_fallback"); + await publish({ type: "answer.host_fallback" }); + await emitVisibleSpoken(RECTIFICATION_USER_COPY.evidenceNotRecorded); + return completeAttempt(false); + } + let latestDossier: V9CaseDossier; + try { + latestDossier = await loadV9CaseDossier(accounting, userId, caseId); + } catch { + return failedAttempt(attemptId, "state_invariant_failed"); + } + const decision = decideFromDossier(latestDossier); + if (decision.nextAction === "ask_candidate_discriminator") { + const openFocus = latestDossier.conversationSummary.activeFocus; + if (!( + openFocus + && isPersistedFocusId(openFocus.id) + && parseAgentChoiceCopy(openFocus.expectedAnswerSchema) + )) { + return failedAttempt(attemptId, "state_invariant_failed"); + } + } + const askedFocus = latestDossier.conversationSummary.activeFocus; + if (askedFocus?.askedTurnId === turnId) { + const stem = focusSpokenPrompt(askedFocus.expectedAnswerSchema); + if (stem) { + const stripped = stripQuestionSentences(answerText, stem); + const next = stripped || RECTIFICATION_USER_COPY.collectHandoff; + if (next !== answerText) answerText = next; + } + } + if (action === "opening") { + const openingRange = openingRangeFromCandidateRange(latestDossier.case.candidateRange); + if (!isAcceptableOpeningBody(answerText)) { + answerText = openingSpokenBody(openingRange); + } + } + const rangeAfterEvidence = previousInferenceFromReceipt( + latestDossier.latestResult?.decisionReceipt ?? null, + )?.credible_range ?? null; + if (batchRescoreFailed(batchToolResult)) { + answerText = withRescoreSkippedNotice(answerText); + } else { + answerText = withRangeChangedAfterEvidence( + answerText, + rangeBeforeCompare, + rangeAfterEvidence, + ); + } + answerText = stripVerbalWindowChange(answerText) + || RECTIFICATION_USER_COPY.declaredWindowLockedReply; + if (answerText !== visibleEmitted) await emitVisibleSpoken(answerText); + return completeAttempt(); + } finally { + clearTimeout(timeout); + signal?.removeEventListener("abort", onAbort); + const stepSummary = diagnosticStepsFromTimings(stepTimings); + const diagnostic: RectificationRunDiagnostic = { + runId: attemptId, + modelId: options.modelName, + finishReason: finishReason ?? "unknown", + 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 })); + } +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts new file mode 100644 index 00000000..63746ed5 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-finish.ts @@ -0,0 +1,197 @@ +/** + * After the last attempt of a rectification agent turn: release or settle + * billing, write the terminal receipts, persist the next interview (with the + * exhaustion gate and the non-terminal repair), trim the spoken answer for the + * interview, finalize the turn and emit `run.completed` / `run.failed`. + * Moved out of runV9AgentTurn unchanged (TASK-rectification-code-split-20260926). + */ +import { finalizeV10RunAttempt } from "./tool-service"; +import { + ensureNonTerminalTurnExit, + persistExhaustionGateTurn, + persistNextInterviewIfIdle, +} from "./answer-choice"; +import { trimSpokenTurnForInterview, composeIdleGapIntoSpoken } from "./collect-prompt"; +import { userFacingRunFailure } from "./run-diagnostic"; +import { reportTurnProgress } from "./turn-instrumentation.ts"; +import { + type AttemptOutcome, + safeErrorCode, + v9TurnReceipts, +} from "./agent-run-support"; +import type { V9AgentRunResult } from "./agent-run"; +import type { V9TimedTurn } from "./agent-run-prepare"; + +export async function finishV9AgentTurn( + turn: V9TimedTurn, + outcome: AttemptOutcome, +): Promise { + const { + userId, + caseId, + action, + accounting, + billing, + emit, + previousFocusId, + resultId, + turnId, + startedAt, + } = turn; + const { persistCommittedPhase, finalizeTurn } = v9TurnReceipts(turn); + + if (!outcome.ok) { + await billing.release(); + await finalizeTurn(outcome.status, null, outcome.attemptId, null); + await emit({ + type: "run.failed", + code: outcome.errorCode ?? "run_failed", + recoverable: outcome.status === "retryable", + message: userFacingRunFailure(outcome.errorCode), + }); + return { + ok: false, + turnId, + turnStatus: outcome.status, + skillLoaded: false, + answerText: "", + phases: outcome.phases, + toolsUsed: outcome.toolsUsed, + errorCode: outcome.errorCode, + previousFocusId, + }; + } + + const skipBilling = outcome.settleBilling === false; + if (!skipBilling) { + const durationMs = Date.now() - startedAt; + const completed = await billing.complete({ ...outcome.usage, durationMs }); + if (!completed) { + await persistCommittedPhase("run.failed", null, outcome.attemptId, outcome.phases.length + 1); + await finalizeV10RunAttempt( + accounting, + userId, + caseId, + turnId, + outcome.attemptId, + "retryable", + "usage_settlement_failed", + outcome.usage, + ); + await finalizeTurn("retryable", null, outcome.attemptId, null); + await emit({ + type: "run.failed", + code: "run_failed", + recoverable: true, + message: userFacingRunFailure("run_failed"), + }); + return { + ok: false, + turnId, + turnStatus: "retryable", + skillLoaded: false, + answerText: "", + phases: [], + toolsUsed: [], + errorCode: "usage_settlement_failed", + previousFocusId, + }; + } + + await persistCommittedPhase( + "billing.settled", + null, + outcome.attemptId, + outcome.phases.length + 1, + true, + ); + } else { + await billing.release(); + } + + await persistCommittedPhase( + "run.completed", + null, + outcome.attemptId, + outcome.phases.length + (skipBilling ? 1 : 2), + true, + ); + await finalizeV10RunAttempt( + accounting, + userId, + caseId, + turnId, + outcome.attemptId, + "completed", + null, + outcome.usage, + ); + const answerText = outcome.answerText; + let spokenAnswer = answerText; + let interviewIdle: Awaited> | 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({ + accounting, + userId, + caseId, + askedTurnId: turnId, + hostNarration: interviewIdle.hostNarration, + resultId, + }); + } catch (error) { + console.warn( + `[rectification-v9] persist exhaustion gate before collect attach failed case=${caseId} reason=${safeErrorCode(error)}`, + ); + } + } + } catch (error) { + console.warn( + `[rectification-v9] persist interview before collect attach failed case=${caseId} reason=${safeErrorCode(error)}`, + ); + try { + interviewIdle = await ensureNonTerminalTurnExit({ accounting, userId, caseId }); + } catch (repairError) { + console.warn( + `[rectification-v9] nonterminal exit after idle failure failed case=${caseId} reason=${safeErrorCode(repairError)}`, + ); + } + } + } + if (action === "evidence") { + spokenAnswer = trimSpokenTurnForInterview(answerText, interviewIdle?.terminalNote === true); + if (interviewIdle?.terminalNote && interviewIdle.hostNarration) { + spokenAnswer = composeIdleGapIntoSpoken(spokenAnswer, interviewIdle.hostNarration); + } + if (spokenAnswer !== answerText) { + await emit({ type: "answer.delta", text: spokenAnswer, replace: true }); + } + } + await finalizeTurn("completed", spokenAnswer, outcome.attemptId, outcome.attemptId, true); + + if (!skipBilling) await emit({ type: "billing.settled" }); + await emit({ type: "run.completed", turnId }); + + return { + ok: true, + turnId, + turnStatus: "completed", + skillLoaded: outcome.skillBound, + answerText: spokenAnswer, + phases: skipBilling + ? [...outcome.phases, "run.completed"] + : [...outcome.phases, "billing.settled", "run.completed"], + toolsUsed: outcome.toolsUsed, + errorCode: null, + previousFocusId, + interviewSettled, + }; +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-messages.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-messages.ts new file mode 100644 index 00000000..f1801070 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-messages.ts @@ -0,0 +1,74 @@ +/** + * The two messages a rectification agent attempt sends the model: the bound + * Skill bootstrap (with the retry constraint on a second attempt) and the user + * turn (the opening brief, or the typed message), both behind the server time + * and Case id lines. Moved out of agent-run.ts unchanged + * (TASK-rectification-code-split-20260926); agent-run.ts re-exports both + * builders. + */ +import type { V9CaseDossier } from "./tool-service"; +import { cachedSystemMessage } from "../../agent-generation-settings.ts"; +import { OPENING_COLLECT_DOMAINS } from "../user-copy"; +import { retryConstraintForAttempt } from "./host-fallback"; +import type { V9AgentRunOptions } from "./agent-run"; + +function clockWindow(range: { start_time?: string | null; end_time?: string | null } | null | undefined): string | null { + const start = range?.start_time?.trim().slice(0, 5) || ""; + const end = range?.end_time?.trim().slice(0, 5) || ""; + return start && end ? `${start}–${end}` : null; +} + +export function buildOpeningBrief(dossier: V9CaseDossier, birthTimeClue?: string | null): string { + const confirmed = dossier.evidence.filter((item) => item.status === "confirmed"); + const pending = dossier.evidence.filter((item) => item.status === "draft" || item.status === "pending_confirmation"); + const domains = [...new Set(confirmed.map((item) => item.domain))].slice(0, 6); + const range = dossier.case.candidateRange; + const window = clockWindow(range); + const uncertaintyType = window + ? `当前搜索窗口 ${window},来自用户在资料里声明的不确定档` + : "用户的出生时间精度仍需通过经历证据核对"; + const clue = typeof birthTimeClue === "string" && birthTimeClue.trim() + ? birthTimeClue.trim() + : ""; + return [ + "【服务端 opening brief】", + `Case 状态:${dossier.case.status}。`, + `当前搜索窗口:${window ?? "尚未锁定"}。来源:intake 声明的不确定档。`, + `出生时间不确定类型:${uncertaintyType}。`, + `已有证据摘要:已确认 ${confirmed.length} 条,待澄清或待确认 ${pending.length} 条${domains.length ? `;已覆盖 ${domains.join("、")}` : ""}。`, + ...(clue + ? [`家人或本人关于出生时段的线索(仅旁白建议,不得改搜索窗口):${clue}`] + : []), + `做法要点:一句当前窗口与核对做法;一句「最后给区间和代表分钟,不给精确到秒」;一句「想到几件说几件,有大概年月就行」并点出${OPENING_COLLECT_DOMAINS.join("、")}。一条消息可以报多件,想到几件说几件。不得写具体年份,不得要求先准备材料。不要提问。先用 rectification-set-focus 的 spokenPrompt 写出当前采集题,题干写成「先说你最容易想起的一两件,年月大概就行」。`, + ].join("\n"); +} + +export function buildAgentMessages( + options: V9AgentRunOptions, + attempt: number, + dossier: V9CaseDossier, + skillInstructions: string, + birthTimeClue: string | null = null, + previousErrorCode: string | null = null, +): unknown[] { + const timeContext = options.timeContext + ?? `服务端当前时间(权威):${new Date().toISOString()}。涉及“现在、今天、今年、未来几个月”等相对时间时,以此为准。`; + const caseContext = `【服务端 Case ID】${options.caseId}。所有 rectification 工具调用的 caseId 必须原样使用此值。`; + const bootstrapContent = [ + "【服务器已绑定当前 Case 的精确 Skill】运行器已在本 attempt 内加载并核验下列指令;不要重复调用 skill。第一步必须调用 rectification-read-case。", + skillInstructions, + ...(attempt > 1 ? [retryConstraintForAttempt(previousErrorCode)] : []), + ].join("\n\n"); + const bootstrap = cachedSystemMessage(bootstrapContent, options.generationModel) + ?? { role: "system" as const, content: bootstrapContent }; + if (options.action === "opening") { + return [bootstrap, { + role: "user", + content: [timeContext, caseContext, buildOpeningBrief(dossier, birthTimeClue)].join("\n"), + }]; + } + return [bootstrap, { + role: "user", + content: [timeContext, caseContext, options.message ?? ""].join("\n"), + }]; +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-prepare.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-prepare.ts new file mode 100644 index 00000000..e6b3f52d --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-prepare.ts @@ -0,0 +1,274 @@ +/** + * Before any attempt of a rectification agent turn: date-reliability answer, + * session check, the year-entry host re-ask, bound Skill identity, the + * already-delivered guard, the opening birth-time clue, billing reserve, and + * the pending turn row (with an idempotent replay of a finished one). Moved + * out of runV9AgentTurn unchanged (TASK-rectification-code-split-20260926). + * A finished answer comes back as the run result; otherwise the prepared + * turn the attempts run on. + */ +import type { RectificationAgentAction } from "./step-budget"; +import { + loadV9CaseDossier, + loadV9CaseSkillIdentity, + loadV9CaseCompute, + RectificationToolServiceError, + persistV9DeterministicTurn, + resolveV10ConversationFocus, + setV9EvidenceDateReliability, + type RectificationRpcClient, + type V9CaseDossier, +} from "./tool-service"; +import { RECTIFICATION_SKILL_NAME } from "./case-status"; +import { classifyDateReliabilityUtterance, isDateReliabilitySchema } from "./date-reliability.ts"; +import { shouldHostReaskYearEntry } from "./year-entry-host.ts"; +import { alreadyDelivered } from "./delivery-turn-guard"; +import { isPersistedFocusId } from "./choice-card"; +import { resolveExactSkillPackage, type ResolvedSkillPackageIdentity } from "../../skill-package-registry.ts"; +import { defaultMessageOrigin, isRectificationMessageOrigin, messageContentHash } from "./message-origin"; +import { + rpcOf, +} from "./agent-run-support"; +import type { V9AgentRunOptions, V9AgentRunResult, V9RunBilling } from "./agent-run"; + +/** Everything the attempts, the retry loop and the finish read from preparation. */ +export type V9PreparedTurn = Readonly<{ + options: V9AgentRunOptions; + userId: string; + caseId: string; + action: RectificationAgentAction; + message: string | null; + accounting: RectificationRpcClient; + buildAgent: V9AgentRunOptions["buildAgent"]; + billing: V9RunBilling; + emit: V9AgentRunOptions["emit"]; + signal: AbortSignal | undefined; + skillName: string; + dossier: V9CaseDossier; + previousFocusId: string | null; + resultId: string; + skillPackage: ResolvedSkillPackageIdentity; + birthTimeClue: string | null; + turnId: string; +}>; + +/** The prepared turn plus the run clock runV9AgentTurn starts before the first attempt. */ +export type V9TimedTurn = V9PreparedTurn & Readonly<{ + startedAt: number; + runBudgetMs: number; + minRetryAttemptMs: number; + attemptCapMs: number; + runDeadlineAt: number; +}>; + +export async function prepareV9AgentTurn( + options: V9AgentRunOptions, +): Promise { + const { + userId, caseId, sessionId, action, message, modelName, + accounting, buildAgent, billing, emit, signal, + } = options; + const skillName = options.skillName ?? RECTIFICATION_SKILL_NAME; + + let dossier = await loadV9CaseDossier(accounting, userId, caseId); + const reliabilityFocus = dossier.conversationSummary.activeFocus; + const reliabilitySchema = reliabilityFocus?.expectedAnswerSchema ?? null; + if (message && isDateReliabilitySchema(reliabilitySchema)) { + const classified = classifyDateReliabilityUtterance(message); + const evidenceId = String(reliabilitySchema.target_evidence_id ?? ""); + const focusId = reliabilityFocus?.id ?? ""; + try { + if (classified && evidenceId) { + await setV9EvidenceDateReliability(accounting, userId, caseId, evidenceId, classified); + if (isPersistedFocusId(focusId)) { + await resolveV10ConversationFocus(accounting, userId, caseId, { + focusId, + status: "resolved", + evidenceId, + }); + } + } else if (isPersistedFocusId(focusId)) { + await resolveV10ConversationFocus(accounting, userId, caseId, { + focusId, + status: "skipped", + }); + } + dossier = await loadV9CaseDossier(accounting, userId, caseId); + } catch (error) { + console.warn( + `[rectification-v9] date reliability write deferred case=${caseId} reason=${ + error instanceof Error ? error.message : String(error) + }`, + ); + } + } + const previousFocusId = dossier.conversationSummary.activeFocus?.id ?? null; + if (dossier.case.sessionId !== sessionId) { + throw new RectificationToolServiceError("agentic_rectification_case_session_mismatch"); + } + // Date reliability (above) runs first. Year-stage typed collect with a + // yearless short reply does not enter the model (BUG-910). BUG-916: this host + // answer runs after the session check, so a message from another session is + // rejected instead of written as a deterministic turn; the phase reuses the + // existing `answer.host_fallback` allowlist entry. + const yearReask = shouldHostReaskYearEntry(dossier.conversationSummary.activeFocus, message); + if (yearReask && message) { + const persisted = await persistV9DeterministicTurn(accounting, userId, caseId, { + requestId: options.requestId, + userMessage: message, + assistantMessage: yearReask, + }); + await emit({ type: "answer.delta", text: yearReask, replace: true }); + return { + ok: true, + turnId: persisted.turnId, + turnStatus: "completed", + skillLoaded: true, + answerText: yearReask, + phases: ["answer.host_fallback"], + toolsUsed: [], + errorCode: null, + previousFocusId, + }; + } + const boundIdentity = await loadV9CaseSkillIdentity(accounting, userId, caseId); + if ((options.skillName && options.skillName !== boundIdentity.name) + || (options.skillVersion && options.skillVersion !== boundIdentity.version) + || dossier.case.skillName !== boundIdentity.name + || dossier.case.skillVersion !== boundIdentity.version) { + throw new RectificationToolServiceError("agentic_rectification_skill_identity_mismatch"); + } + const resultId = dossier.latestResult?.resultId ?? ""; + if (alreadyDelivered({ + caseId, + resultId, + hasUserMessage: Boolean(message?.trim()), + })) { + return { + ok: true, + turnId: "", + turnStatus: "completed", + skillLoaded: true, + answerText: "", + phases: [], + toolsUsed: [], + errorCode: "already_delivered", + previousFocusId, + }; + } + const skillPackage = resolveExactSkillPackage( + boundIdentity.name, + boundIdentity.version, + boundIdentity.sha256, + ); + if (skillPackage.sourceCommit !== boundIdentity.sourceCommit) { + throw new RectificationToolServiceError("agentic_rectification_skill_identity_mismatch"); + } + let birthTimeClue: string | null = null; + if (action === "opening") { + try { + const compute = await loadV9CaseCompute(accounting, userId, caseId); + const raw = compute.baselineBirthSnapshot.birth_time_clue; + birthTimeClue = typeof raw === "string" && raw.trim() ? raw.trim() : null; + } catch { + birthTimeClue = null; + } + } + + const reserve = await billing.reserve(); + if (!reserve.success) { + throw new RectificationToolServiceError(reserve.reason ?? "billing_denied"); + } + + let turnRow: unknown; + try { + turnRow = await rpcOf(accounting, "append_agentic_rectification_turn", { + p_user_id: userId, + p_case_id: caseId, + p_user_message: message, + p_assistant_message: null, + p_model_name: modelName, + p_model_version: null, + p_status: "pending", + p_request_id: options.requestId, + }); + } catch (error) { + await billing.release(); + throw error; + } + const turnRecord = turnRow && typeof turnRow === "object" + ? turnRow as Record + : null; + const turnId = typeof turnRecord?.turn_id === "string" ? turnRecord.turn_id : ""; + if (!turnId) { + await billing.release(); + throw new RectificationToolServiceError("agentic_rectification_turn_incomplete"); + } + + try { + await rpcOf(accounting, "record_agentic_rectification_turn_origin", { + p_user_id: userId, + p_case_id: caseId, + p_turn_id: turnId, + p_origin: isRectificationMessageOrigin(options.messageOrigin) + ? options.messageOrigin + : defaultMessageOrigin(options.action), + p_client_action_id: options.clientActionId ?? options.requestId, + p_content_hash: messageContentHash(message), + }); + } catch { + // Origin is audit metadata; a missing RPC must not fail the turn. + } + + const shouldExecute = turnRecord?.should_execute === undefined + ? true + : turnRecord.should_execute === true; + const existingStatus = typeof turnRecord?.status === "string" ? turnRecord.status : "pending"; + const existingAnswer = typeof turnRecord?.assistant_message === "string" + ? turnRecord.assistant_message + : ""; + if (!shouldExecute) { + if (existingStatus === "completed" && existingAnswer.trim()) { + await emit({ type: "run.started" }); + await emit({ type: "answer.delta", text: existingAnswer }); + await emit({ type: "run.completed", turnId }); + return { + ok: true, + turnId, + turnStatus: "completed", + skillLoaded: typeof turnRecord?.successful_attempt_id === "string", + answerText: existingAnswer, + phases: ["run.completed"], + toolsUsed: [], + errorCode: null, + previousFocusId, + }; + } + if (existingStatus !== "pending") await billing.release(); + throw new RectificationToolServiceError( + existingStatus === "pending" + ? "agentic_rectification_turn_in_progress" + : "agentic_rectification_turn_already_finalized", + ); + } + + return { + options, + userId, + caseId, + action, + message, + accounting, + buildAgent, + billing, + emit, + signal, + skillName, + dossier, + previousFocusId, + resultId, + skillPackage, + birthTimeClue, + turnId, + }; +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-retry.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-retry.ts new file mode 100644 index 00000000..b5e47889 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-retry.ts @@ -0,0 +1,125 @@ +/** + * The attempt loop of a rectification agent turn: claim an attempt, stream + * it, record a failed attempt, and retry once (`attempt.reset` first) while + * the error is retryable and the run budget allows. Moved out of + * runV9AgentTurn unchanged (TASK-rectification-code-split-20260926). + */ +import { createV10RunAttempt, finalizeV10RunAttempt, RectificationToolServiceError } from "./tool-service"; +import { canStartRetryAttempt } from "../../rectification-run-budget.ts"; +import { reportTurnProgress, resetTurnProgress } from "./turn-instrumentation.ts"; +import { + type AttemptOutcome, + MAX_ATTEMPTS, + safeErrorCode, + isRetryableError, + shouldAutoRetry, + v9TurnReceipts, +} from "./agent-run-support"; +import { streamV9Attempt } from "./agent-run-attempt"; +import type { V9TimedTurn } from "./agent-run-prepare"; + +export async function runV9AttemptsWithRetry(turn: V9TimedTurn): Promise { + const { + userId, + caseId, + accounting, + billing, + emit, + signal, + turnId, + minRetryAttemptMs, + runDeadlineAt, + } = turn; + const { persistCommittedPhase } = v9TurnReceipts(turn); + const streamAttempt = ( + attemptNumber: number, + attemptId: string, + previousErrorCode: string | null, + ) => streamV9Attempt(turn, attemptNumber, attemptId, previousErrorCode); + + let lastAttemptError: string | null = null; + let finalOutcome: AttemptOutcome | null = null; + for (let attemptNumber = 1; attemptNumber <= MAX_ATTEMPTS; attemptNumber += 1) { + const remainingMs = runDeadlineAt - Date.now(); + if (attemptNumber > 1 && !canStartRetryAttempt(remainingMs, minRetryAttemptMs)) break; + const claim = await createV10RunAttempt( + accounting, + userId, + caseId, + turnId, + attemptNumber, + ); + const { attemptId } = claim; + if (!claim.shouldExecute) { + await billing.release(); + throw new RectificationToolServiceError( + claim.alreadyInProgress + ? "agentic_rectification_attempt_in_progress" + : "agentic_rectification_attempt_already_finalized", + ); + } + let outcome: AttemptOutcome; + try { + outcome = await streamAttempt(attemptNumber, attemptId, lastAttemptError); + } catch (error) { + const errorCode = safeErrorCode(error); + const reason = error instanceof Error ? error.message.slice(0, 180) : "UnknownError"; + console.error( + `[rectification-v10] attempt failed case=${caseId} turn=${turnId} attempt=${attemptId} code=${errorCode} reason=${reason}`, + ); + outcome = { + ok: false, + status: isRetryableError(errorCode) ? "retryable" : "failed", + errorCode, + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [], + toolsUsed: [], + events: [], + skillBound: false, + caseLoaded: false, + attemptId, + }; + } + if (!outcome.ok) { + await persistCommittedPhase("run.failed", null, attemptId, 1_000_000 + attemptNumber); + await finalizeV10RunAttempt( + accounting, + userId, + caseId, + turnId, + attemptId, + outcome.status, + outcome.errorCode, + outcome.usage, + ); + } + finalOutcome = outcome; + if (!outcome.ok) lastAttemptError = outcome.errorCode; + if (outcome.ok + || outcome.status === "failed" + || !shouldAutoRetry(outcome.errorCode ?? "run_failed", signal, outcome.status) + || attemptNumber === MAX_ATTEMPTS) break; + await emit({ type: "attempt.reset" }); + resetTurnProgress(); + reportTurnProgress("received"); + } + + const outcome = finalOutcome ?? { + ok: false, + status: "failed" as const, + errorCode: "run_failed", + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [], + toolsUsed: [], + events: [], + skillBound: false, + caseLoaded: false, + attemptId: "", + }; + + return outcome; +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run-support.ts b/frontend/src/lib/rectification-agentic/v9/agent-run-support.ts new file mode 100644 index 00000000..d292ea82 --- /dev/null +++ b/frontend/src/lib/rectification-agentic/v9/agent-run-support.ts @@ -0,0 +1,207 @@ +/** + * Shared pieces of the V10 rectification turn runner: attempt outcome types, + * retry classification, RPC unwrapping, and the turn-scoped receipt writers + * (committed phases, turn finalize). Moved out of agent-run.ts unchanged + * (TASK-rectification-code-split-20260926). + */ +import { insertV9RunPhase, RectificationToolServiceError, type RectificationRpcClient } from "./tool-service"; +import { promptCacheUsage } from "../../agent-generation-settings.ts"; +import { toAgentModelFinishReason } from "../../agent-observability.ts"; +import type { PublicStreamEvent } from "./stream-mapping"; + +export type AttemptStatus = "completed" | "failed" | "retryable"; +export type Usage = Readonly<{ inputTokens: number; outputTokens: number; cache?: ReturnType }>; +export type AttemptOutcome = Readonly<{ + ok: boolean; + status: AttemptStatus; + errorCode: string | null; + usage: Usage; + answerText: string; + answerDeltas: readonly string[]; + phases: readonly string[]; + toolsUsed: readonly string[]; + events: readonly PublicStreamEvent[]; + skillBound: boolean; + caseLoaded: boolean; + attemptId: string; + settleBilling?: boolean; +}>; + +export const MAX_ATTEMPTS = 2; +export const RETRYABLE_ERROR_CODES = new Set([ + "stream_aborted", + "stream_unfinished", + "skill_not_loaded", + "skill_not_bound", + "case_not_loaded", +]); + +export function streamFinishReason(chunk: { + type: string; + payload?: { stepResult?: { reason?: unknown }; reason?: unknown }; +}): ReturnType | null { + if (chunk.type !== "finish") return null; + const raw = chunk.payload?.stepResult?.reason ?? chunk.payload?.reason; + if (typeof raw !== "string" || !raw.trim()) return null; + return toAgentModelFinishReason(raw); +} + +export function first(value: unknown): unknown { + if (Array.isArray(value)) return value[0] ?? null; + if (value && typeof value === "object" && "value" in value) { + return (value as { value?: unknown }).value; + } + return value; +} + +/** + * Identity for a public tool-call chunk. Mastra may put the model input on + * `args`, `input`, or omit it; missing input collapses to `{}` so a second + * call of the same tool name still looks identical. + */ +export function publicToolCallKey(toolName: string, payload: unknown): string { + if (!payload || typeof payload !== "object") return `${toolName}:{}`; + const record = payload as Record; + const args = record.args ?? record.input ?? record.toolArgs ?? {}; + try { + return `${toolName}:${JSON.stringify(args)}`; + } catch { + return `${toolName}:{}`; + } +} + +export async function rpcOf( + accounting: RectificationRpcClient, + fn: string, + args: Record, +): Promise { + const { data, error } = await accounting.rpc(fn, args); + if (error) throw new RectificationToolServiceError(error.message); + return first(data); +} + +export function safeErrorCode(error: unknown): string { + const message = error instanceof Error ? error.message : String(error); + for (const code of [ + "empty_stream", + "evidence_not_written", + "stream_aborted", + "stream_unfinished", + "skill_not_loaded", + "skill_not_bound", + "case_not_loaded", + "repeated_tool_call", + "focus_persistence_failed", + "answer_truncated", + "run_timeout", + "max_steps", + "provider_error", + ]) { + if (message.includes(code)) return code; + } + if (message.includes("agentic_rectification_case_terminal")) return "case_terminal"; + if (message.includes("agentic_rectification_case_not_found")) return "case_not_found"; + if (message.includes("agentic_rectification_case_session_mismatch")) return "case_session_mismatch"; + if (message.includes("Thinking mode does not support this tool_choice")) { + return "thinking_tool_choice_unsupported"; + } + return "run_failed"; +} + +export function isRetryableError(errorCode: string): boolean { + return RETRYABLE_ERROR_CODES.has(errorCode); +} + +export function shouldAutoRetry( + errorCode: string, + signal?: AbortSignal, + status?: AttemptStatus, +): boolean { + if (signal?.aborted) return false; + if (errorCode === "empty_stream") return status === "retryable"; + if (errorCode === "evidence_not_written") return status === "retryable"; + return isRetryableError(errorCode); +} + +/** The Case turn the receipt writers below belong to. */ +export type V9TurnScope = Readonly<{ + accounting: RectificationRpcClient; + userId: string; + caseId: string; + turnId: string; +}>; + +/** + * `failedAttempt`, `persistCommittedPhase` and `finalizeTurn` as they were + * nested inside runV9AgentTurn, bound to one turn so every phase of the + * runner calls them with the same arguments as before. + */ +export function v9TurnReceipts(scope: V9TurnScope) { + const { accounting, userId, caseId, turnId } = scope; + return { failedAttempt, persistCommittedPhase, finalizeTurn }; + + function failedAttempt(attemptId: string, errorCode: string): AttemptOutcome { + return { + ok: false, + status: isRetryableError(errorCode) ? "retryable" : "failed", + errorCode, + usage: { inputTokens: 0, outputTokens: 0 }, + answerText: "", + answerDeltas: [], + phases: [], + toolsUsed: [], + events: [], + skillBound: false, + caseLoaded: false, + attemptId, + }; + } + + async function persistCommittedPhase( + phase: string, + toolName: string | null, + attemptId: string, + sequence: number, + strict = false, + ) { + try { + await insertV9RunPhase( + accounting, + userId, + caseId, + turnId, + phase, + toolName, + sequence, + attemptId, + ); + } catch (error) { + if (strict) throw error; + // Non-terminal activity receipts remain best effort. The completion + // receipts above are strict because they are part of business truth. + } + } + + async function finalizeTurn( + status: AttemptStatus, + assistantText: string | null, + attemptId: string, + successfulAttemptId: string | null, + strict = false, + ) { + try { + await rpcOf(accounting, "finalize_agentic_rectification_turn", { + p_user_id: userId, + p_case_id: caseId, + p_turn_id: turnId, + p_attempt_id: attemptId, + p_status: status, + p_assistant_message: status === "completed" ? assistantText : null, + p_successful_attempt_id: successfulAttemptId, + }); + } catch (error) { + if (strict) throw error; + console.warn(`[rectification-v10] turn finalize failed turn=${turnId} status=${status} reason=${safeErrorCode(error)}`); + } + } +} diff --git a/frontend/src/lib/rectification-agentic/v9/agent-run.ts b/frontend/src/lib/rectification-agentic/v9/agent-run.ts index f1798057..5afb1cba 100644 --- a/frontend/src/lib/rectification-agentic/v9/agent-run.ts +++ b/frontend/src/lib/rectification-agentic/v9/agent-run.ts @@ -9,103 +9,34 @@ * `attempt.reset` first so the client discards the abandoned attempt. * Durable receipts, billing, and settled history still come only from the * successful attempt. + * + * Phases live in their own modules (TASK-rectification-code-split-20260926): + * preparation → agent-run-prepare.ts, one attempt → agent-run-attempt.ts, + * the retry loop → agent-run-retry.ts, the finish → agent-run-finish.ts, + * shared outcome types and receipt writers → agent-run-support.ts, the model + * messages → agent-run-messages.ts. This file keeps the public types and the + * order of the phases. */ import type { Agent } from "@mastra/core/agent"; -import { type RectificationAgentAction, resolveRectificationStepBudget } from "./step-budget"; -import { - createV10RunAttempt, - finalizeV10RunAttempt, - insertV9RunPhase, - insertV9SkillRunReceipt, - loadV9CaseDossier, - loadV9CaseSkillIdentity, - loadV9CaseCompute, - RectificationToolServiceError, - persistV9DeterministicTurn, - resolveV10ConversationFocus, - setV9EvidenceDateReliability, - type RectificationRpcClient, - type V9CaseDossier, -} from "./tool-service"; +import type { RectificationAgentAction } from "./step-budget"; +import type { RectificationRpcClient } from "./tool-service"; import { RECTIFICATION_SKILL_NAME, RECTIFICATION_SKILL_VERSION } from "./case-status"; -import { RECTIFICATION_AGENT_TOOLS } from "./public-receipt"; -import { agentGenerationSettings, cachedSystemMessage, promptCacheUsage } from "../../agent-generation-settings.ts"; -import { toAgentModelFinishReason } from "../../agent-observability.ts"; -import { classifyDateReliabilityUtterance, isDateReliabilitySchema } from "./date-reliability.ts"; -import { shouldHostReaskYearEntry } from "./year-entry-host.ts"; -import { decideFromDossier } from "./decision-from-dossier"; -import { ensureNonTerminalTurnExit, persistExhaustionGateTurn, persistNextInterviewIfIdle } from "./answer-choice"; -import { alreadyDelivered } from "./delivery-turn-guard"; -import { parseAgentChoiceCopy, isPersistedFocusId } from "./choice-card"; -import { - RECTIFICATION_USER_COPY, - OPENING_COLLECT_DOMAINS, - isAcceptableOpeningBody, - openingRangeFromCandidateRange, - openingSpokenBody, - withCompareFailedRetryNotice, - withRangeChangedAfterEvidence, - withRescoreSkippedNotice, -} from "../user-copy"; +import type { promptCacheUsage } from "../../agent-generation-settings.ts"; import { RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS, RECTIFICATION_AGENT_ROUTE_MAX_DURATION_S, RECTIFICATION_MIN_RETRY_ATTEMPT_MS, RECTIFICATION_RUN_BUDGET_MS, - attemptTimeoutForRemainingBudget, - canStartRetryAttempt, } from "../../rectification-run-budget.ts"; -import { stripQuestionSentences, stripVerbalWindowChange, trimSpokenTurnForInterview, composeIdleGapIntoSpoken } from "./collect-prompt"; -import { focusSpokenPrompt } from "./turn-question"; -import { previousInferenceFromReceipt } from "../core/compose-receipt.ts"; -import { - resolveExactSkillPackage, - type ResolvedSkillPackageIdentity, -} from "../../skill-package-registry.ts"; -import { - activityChangedFromTool, - mapStreamChunkToActivity, - mapStreamChunkToPhase, - streamToolNames, - isPublicRectificationToolName, - turnProgressForChunk, - type PublicStreamEvent, -} from "./stream-mapping"; -import { - diagnosticStepsFromTimings, - mapModelFinishToErrorCode, - recordStepChunk, - userFacingRunFailure, - type RectificationRunDiagnostic, - type RectificationStepTiming, -} from "./run-diagnostic"; -import { - currentEngineCallTimings, - reportTurnProgress, - resetTurnProgress, -} from "./turn-instrumentation.ts"; +import type { ResolvedSkillPackageIdentity } from "../../skill-package-registry.ts"; +import type { PublicStreamEvent } from "./stream-mapping"; import type { TurnIntentClassifierDiagnostic } from "./turn-intent-classifier"; -import { - batchResultFromToolChunk, - batchRescoreFailed, - composeHostFallbackNarration, - lastCompletedPublicTool, - publicWriteToolCompleted, - answerClaimsEvidenceRecorded, - retryConstraintForAttempt, - turnExpectsEvidenceWrite, -} from "./host-fallback"; -import { - applyStepAnswerChunk, - createStepAnswerState, - flushStepAnswerOnStreamFinish, -} from "./step-answer"; -import { - defaultMessageOrigin, - isRectificationMessageOrigin, - messageContentHash, - type RectificationMessageOrigin, -} from "./message-origin"; +import type { RectificationMessageOrigin } from "./message-origin"; +import { prepareV9AgentTurn, type V9TimedTurn } from "./agent-run-prepare"; +import { runV9AttemptsWithRetry } from "./agent-run-retry"; +import { finishV9AgentTurn } from "./agent-run-finish"; + +export { buildAgentMessages, buildOpeningBrief } from "./agent-run-messages"; export type V9RunBilling = Readonly<{ reserve(): Promise<{ success: boolean; reason?: string; status: number }>; @@ -175,1226 +106,30 @@ export type V9AgentRunResult = Readonly<{ interviewSettled?: boolean; }>; -type AttemptStatus = "completed" | "failed" | "retryable"; -type Usage = Readonly<{ inputTokens: number; outputTokens: number; cache?: ReturnType }>; -type AttemptOutcome = Readonly<{ - ok: boolean; - status: AttemptStatus; - errorCode: string | null; - usage: Usage; - answerText: string; - answerDeltas: readonly string[]; - phases: readonly string[]; - toolsUsed: readonly string[]; - events: readonly PublicStreamEvent[]; - skillBound: boolean; - caseLoaded: boolean; - attemptId: string; - settleBilling?: boolean; -}>; - -const MAX_ATTEMPTS = 2; -const RETRYABLE_ERROR_CODES = new Set([ - "stream_aborted", - "stream_unfinished", - "skill_not_loaded", - "skill_not_bound", - "case_not_loaded", -]); - -function streamFinishReason(chunk: { - type: string; - payload?: { stepResult?: { reason?: unknown }; reason?: unknown }; -}): ReturnType | null { - if (chunk.type !== "finish") return null; - const raw = chunk.payload?.stepResult?.reason ?? chunk.payload?.reason; - if (typeof raw !== "string" || !raw.trim()) return null; - return toAgentModelFinishReason(raw); -} - -function first(value: unknown): unknown { - if (Array.isArray(value)) return value[0] ?? null; - if (value && typeof value === "object" && "value" in value) { - return (value as { value?: unknown }).value; - } - return value; -} - -/** - * Identity for a public tool-call chunk. Mastra may put the model input on - * `args`, `input`, or omit it; missing input collapses to `{}` so a second - * call of the same tool name still looks identical. - */ -function publicToolCallKey(toolName: string, payload: unknown): string { - if (!payload || typeof payload !== "object") return `${toolName}:{}`; - const record = payload as Record; - const args = record.args ?? record.input ?? record.toolArgs ?? {}; - try { - return `${toolName}:${JSON.stringify(args)}`; - } catch { - return `${toolName}:{}`; - } -} - -async function rpcOf( - accounting: RectificationRpcClient, - fn: string, - args: Record, -): Promise { - const { data, error } = await accounting.rpc(fn, args); - if (error) throw new RectificationToolServiceError(error.message); - return first(data); -} - -function safeErrorCode(error: unknown): string { - const message = error instanceof Error ? error.message : String(error); - for (const code of [ - "empty_stream", - "evidence_not_written", - "stream_aborted", - "stream_unfinished", - "skill_not_loaded", - "skill_not_bound", - "case_not_loaded", - "repeated_tool_call", - "focus_persistence_failed", - "answer_truncated", - "run_timeout", - "max_steps", - "provider_error", - ]) { - if (message.includes(code)) return code; - } - if (message.includes("agentic_rectification_case_terminal")) return "case_terminal"; - if (message.includes("agentic_rectification_case_not_found")) return "case_not_found"; - if (message.includes("agentic_rectification_case_session_mismatch")) return "case_session_mismatch"; - if (message.includes("Thinking mode does not support this tool_choice")) { - return "thinking_tool_choice_unsupported"; - } - return "run_failed"; -} - -function isRetryableError(errorCode: string): boolean { - return RETRYABLE_ERROR_CODES.has(errorCode); -} - -function shouldAutoRetry( - errorCode: string, - signal?: AbortSignal, - status?: AttemptStatus, -): boolean { - if (signal?.aborted) return false; - if (errorCode === "empty_stream") return status === "retryable"; - if (errorCode === "evidence_not_written") return status === "retryable"; - return isRetryableError(errorCode); -} - -function clockWindow(range: { start_time?: string | null; end_time?: string | null } | null | undefined): string | null { - const start = range?.start_time?.trim().slice(0, 5) || ""; - const end = range?.end_time?.trim().slice(0, 5) || ""; - return start && end ? `${start}–${end}` : null; -} - -export function buildOpeningBrief(dossier: V9CaseDossier, birthTimeClue?: string | null): string { - const confirmed = dossier.evidence.filter((item) => item.status === "confirmed"); - const pending = dossier.evidence.filter((item) => item.status === "draft" || item.status === "pending_confirmation"); - const domains = [...new Set(confirmed.map((item) => item.domain))].slice(0, 6); - const range = dossier.case.candidateRange; - const window = clockWindow(range); - const uncertaintyType = window - ? `当前搜索窗口 ${window},来自用户在资料里声明的不确定档` - : "用户的出生时间精度仍需通过经历证据核对"; - const clue = typeof birthTimeClue === "string" && birthTimeClue.trim() - ? birthTimeClue.trim() - : ""; - return [ - "【服务端 opening brief】", - `Case 状态:${dossier.case.status}。`, - `当前搜索窗口:${window ?? "尚未锁定"}。来源:intake 声明的不确定档。`, - `出生时间不确定类型:${uncertaintyType}。`, - `已有证据摘要:已确认 ${confirmed.length} 条,待澄清或待确认 ${pending.length} 条${domains.length ? `;已覆盖 ${domains.join("、")}` : ""}。`, - ...(clue - ? [`家人或本人关于出生时段的线索(仅旁白建议,不得改搜索窗口):${clue}`] - : []), - `做法要点:一句当前窗口与核对做法;一句「最后给区间和代表分钟,不给精确到秒」;一句「想到几件说几件,有大概年月就行」并点出${OPENING_COLLECT_DOMAINS.join("、")}。一条消息可以报多件,想到几件说几件。不得写具体年份,不得要求先准备材料。不要提问。先用 rectification-set-focus 的 spokenPrompt 写出当前采集题,题干写成「先说你最容易想起的一两件,年月大概就行」。`, - ].join("\n"); -} - export async function runV9AgentTurn(options: V9AgentRunOptions): Promise { - const { - userId, caseId, sessionId, action, message, modelName, - accounting, buildAgent, billing, emit, signal, - } = options; - const skillName = options.skillName ?? RECTIFICATION_SKILL_NAME; - - let dossier = await loadV9CaseDossier(accounting, userId, caseId); - const reliabilityFocus = dossier.conversationSummary.activeFocus; - const reliabilitySchema = reliabilityFocus?.expectedAnswerSchema ?? null; - if (message && isDateReliabilitySchema(reliabilitySchema)) { - const classified = classifyDateReliabilityUtterance(message); - const evidenceId = String(reliabilitySchema.target_evidence_id ?? ""); - const focusId = reliabilityFocus?.id ?? ""; - try { - if (classified && evidenceId) { - await setV9EvidenceDateReliability(accounting, userId, caseId, evidenceId, classified); - if (isPersistedFocusId(focusId)) { - await resolveV10ConversationFocus(accounting, userId, caseId, { - focusId, - status: "resolved", - evidenceId, - }); - } - } else if (isPersistedFocusId(focusId)) { - await resolveV10ConversationFocus(accounting, userId, caseId, { - focusId, - status: "skipped", - }); - } - dossier = await loadV9CaseDossier(accounting, userId, caseId); - } catch (error) { - console.warn( - `[rectification-v9] date reliability write deferred case=${caseId} reason=${ - error instanceof Error ? error.message : String(error) - }`, - ); - } - } - const previousFocusId = dossier.conversationSummary.activeFocus?.id ?? null; - if (dossier.case.sessionId !== sessionId) { - throw new RectificationToolServiceError("agentic_rectification_case_session_mismatch"); - } - // Date reliability (above) runs first. Year-stage typed collect with a - // yearless short reply does not enter the model (BUG-910). BUG-916: this host - // answer runs after the session check, so a message from another session is - // rejected instead of written as a deterministic turn; the phase reuses the - // existing `answer.host_fallback` allowlist entry. - const yearReask = shouldHostReaskYearEntry(dossier.conversationSummary.activeFocus, message); - if (yearReask && message) { - const persisted = await persistV9DeterministicTurn(accounting, userId, caseId, { - requestId: options.requestId, - userMessage: message, - assistantMessage: yearReask, - }); - await emit({ type: "answer.delta", text: yearReask, replace: true }); - return { - ok: true, - turnId: persisted.turnId, - turnStatus: "completed", - skillLoaded: true, - answerText: yearReask, - phases: ["answer.host_fallback"], - toolsUsed: [], - errorCode: null, - previousFocusId, - }; - } - const boundIdentity = await loadV9CaseSkillIdentity(accounting, userId, caseId); - if ((options.skillName && options.skillName !== boundIdentity.name) - || (options.skillVersion && options.skillVersion !== boundIdentity.version) - || dossier.case.skillName !== boundIdentity.name - || dossier.case.skillVersion !== boundIdentity.version) { - throw new RectificationToolServiceError("agentic_rectification_skill_identity_mismatch"); - } - const resultId = dossier.latestResult?.resultId ?? ""; - if (alreadyDelivered({ - caseId, - resultId, - hasUserMessage: Boolean(message?.trim()), - })) { - return { - ok: true, - turnId: "", - turnStatus: "completed", - skillLoaded: true, - answerText: "", - phases: [], - toolsUsed: [], - errorCode: "already_delivered", - previousFocusId, - }; - } - const skillPackage = resolveExactSkillPackage( - boundIdentity.name, - boundIdentity.version, - boundIdentity.sha256, - ); - if (skillPackage.sourceCommit !== boundIdentity.sourceCommit) { - throw new RectificationToolServiceError("agentic_rectification_skill_identity_mismatch"); - } - let birthTimeClue: string | null = null; - if (action === "opening") { - try { - const compute = await loadV9CaseCompute(accounting, userId, caseId); - const raw = compute.baselineBirthSnapshot.birth_time_clue; - birthTimeClue = typeof raw === "string" && raw.trim() ? raw.trim() : null; - } catch { - birthTimeClue = null; - } - } - - const reserve = await billing.reserve(); - if (!reserve.success) { - throw new RectificationToolServiceError(reserve.reason ?? "billing_denied"); - } - - let turnRow: unknown; - try { - turnRow = await rpcOf(accounting, "append_agentic_rectification_turn", { - p_user_id: userId, - p_case_id: caseId, - p_user_message: message, - p_assistant_message: null, - p_model_name: modelName, - p_model_version: null, - p_status: "pending", - p_request_id: options.requestId, - }); - } catch (error) { - await billing.release(); - throw error; - } - const turnRecord = turnRow && typeof turnRow === "object" - ? turnRow as Record - : null; - const turnId = typeof turnRecord?.turn_id === "string" ? turnRecord.turn_id : ""; - if (!turnId) { - await billing.release(); - throw new RectificationToolServiceError("agentic_rectification_turn_incomplete"); - } - - try { - await rpcOf(accounting, "record_agentic_rectification_turn_origin", { - p_user_id: userId, - p_case_id: caseId, - p_turn_id: turnId, - p_origin: isRectificationMessageOrigin(options.messageOrigin) - ? options.messageOrigin - : defaultMessageOrigin(options.action), - p_client_action_id: options.clientActionId ?? options.requestId, - p_content_hash: messageContentHash(message), - }); - } catch { - // Origin is audit metadata; a missing RPC must not fail the turn. - } - - const shouldExecute = turnRecord?.should_execute === undefined - ? true - : turnRecord.should_execute === true; - const existingStatus = typeof turnRecord?.status === "string" ? turnRecord.status : "pending"; - const existingAnswer = typeof turnRecord?.assistant_message === "string" - ? turnRecord.assistant_message - : ""; - if (!shouldExecute) { - if (existingStatus === "completed" && existingAnswer.trim()) { - await emit({ type: "run.started" }); - await emit({ type: "answer.delta", text: existingAnswer }); - await emit({ type: "run.completed", turnId }); - return { - ok: true, - turnId, - turnStatus: "completed", - skillLoaded: typeof turnRecord?.successful_attempt_id === "string", - answerText: existingAnswer, - phases: ["run.completed"], - toolsUsed: [], - errorCode: null, - previousFocusId, - }; - } - if (existingStatus !== "pending") await billing.release(); - throw new RectificationToolServiceError( - existingStatus === "pending" - ? "agentic_rectification_turn_in_progress" - : "agentic_rectification_turn_already_finalized", - ); - } + // Preparation either answers the turn outright (year re-ask, already + // delivered, idempotent replay) or hands back the pending turn to run. + const prepared = await prepareV9AgentTurn(options); + if ("ok" in prepared) return prepared; const startedAt = Date.now(); const runBudgetMs = options.runBudgetMs ?? RECTIFICATION_RUN_BUDGET_MS; const minRetryAttemptMs = options.minRetryAttemptMs ?? RECTIFICATION_MIN_RETRY_ATTEMPT_MS; const attemptCapMs = options.attemptTimeoutMs ?? RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS; const runDeadlineAt = startedAt + runBudgetMs; + const turn: V9TimedTurn = { + ...prepared, + startedAt, + runBudgetMs, + minRetryAttemptMs, + attemptCapMs, + runDeadlineAt, + }; + const { emit } = turn; await emit({ type: "run.started" }); - let lastAttemptError: string | null = null; - let finalOutcome: AttemptOutcome | null = null; - for (let attemptNumber = 1; attemptNumber <= MAX_ATTEMPTS; attemptNumber += 1) { - const remainingMs = runDeadlineAt - Date.now(); - if (attemptNumber > 1 && !canStartRetryAttempt(remainingMs, minRetryAttemptMs)) break; - const claim = await createV10RunAttempt( - accounting, - userId, - caseId, - turnId, - attemptNumber, - ); - const { attemptId } = claim; - if (!claim.shouldExecute) { - await billing.release(); - throw new RectificationToolServiceError( - claim.alreadyInProgress - ? "agentic_rectification_attempt_in_progress" - : "agentic_rectification_attempt_already_finalized", - ); - } - let outcome: AttemptOutcome; - try { - outcome = await streamAttempt(attemptNumber, attemptId, lastAttemptError); - } catch (error) { - const errorCode = safeErrorCode(error); - const reason = error instanceof Error ? error.message.slice(0, 180) : "UnknownError"; - console.error( - `[rectification-v10] attempt failed case=${caseId} turn=${turnId} attempt=${attemptId} code=${errorCode} reason=${reason}`, - ); - outcome = { - ok: false, - status: isRetryableError(errorCode) ? "retryable" : "failed", - errorCode, - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [], - toolsUsed: [], - events: [], - skillBound: false, - caseLoaded: false, - attemptId, - }; - } - if (!outcome.ok) { - await persistCommittedPhase("run.failed", null, attemptId, 1_000_000 + attemptNumber); - await finalizeV10RunAttempt( - accounting, - userId, - caseId, - turnId, - attemptId, - outcome.status, - outcome.errorCode, - outcome.usage, - ); - } - finalOutcome = outcome; - if (!outcome.ok) lastAttemptError = outcome.errorCode; - if (outcome.ok - || outcome.status === "failed" - || !shouldAutoRetry(outcome.errorCode ?? "run_failed", signal, outcome.status) - || attemptNumber === MAX_ATTEMPTS) break; - await emit({ type: "attempt.reset" }); - resetTurnProgress(); - reportTurnProgress("received"); - } - - const outcome = finalOutcome ?? { - ok: false, - status: "failed" as const, - errorCode: "run_failed", - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [], - toolsUsed: [], - events: [], - skillBound: false, - caseLoaded: false, - attemptId: "", - }; - - if (!outcome.ok) { - await billing.release(); - await finalizeTurn(outcome.status, null, outcome.attemptId, null); - await emit({ - type: "run.failed", - code: outcome.errorCode ?? "run_failed", - recoverable: outcome.status === "retryable", - message: userFacingRunFailure(outcome.errorCode), - }); - return { - ok: false, - turnId, - turnStatus: outcome.status, - skillLoaded: false, - answerText: "", - phases: outcome.phases, - toolsUsed: outcome.toolsUsed, - errorCode: outcome.errorCode, - previousFocusId, - }; - } - - const skipBilling = outcome.settleBilling === false; - if (!skipBilling) { - const durationMs = Date.now() - startedAt; - const completed = await billing.complete({ ...outcome.usage, durationMs }); - if (!completed) { - await persistCommittedPhase("run.failed", null, outcome.attemptId, outcome.phases.length + 1); - await finalizeV10RunAttempt( - accounting, - userId, - caseId, - turnId, - outcome.attemptId, - "retryable", - "usage_settlement_failed", - outcome.usage, - ); - await finalizeTurn("retryable", null, outcome.attemptId, null); - await emit({ - type: "run.failed", - code: "run_failed", - recoverable: true, - message: userFacingRunFailure("run_failed"), - }); - return { - ok: false, - turnId, - turnStatus: "retryable", - skillLoaded: false, - answerText: "", - phases: [], - toolsUsed: [], - errorCode: "usage_settlement_failed", - previousFocusId, - }; - } - - await persistCommittedPhase( - "billing.settled", - null, - outcome.attemptId, - outcome.phases.length + 1, - true, - ); - } else { - await billing.release(); - } - - await persistCommittedPhase( - "run.completed", - null, - outcome.attemptId, - outcome.phases.length + (skipBilling ? 1 : 2), - true, - ); - await finalizeV10RunAttempt( - accounting, - userId, - caseId, - turnId, - outcome.attemptId, - "completed", - null, - outcome.usage, - ); - const answerText = outcome.answerText; - let spokenAnswer = answerText; - let interviewIdle: Awaited> | 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({ - accounting, - userId, - caseId, - askedTurnId: turnId, - hostNarration: interviewIdle.hostNarration, - resultId, - }); - } catch (error) { - console.warn( - `[rectification-v9] persist exhaustion gate before collect attach failed case=${caseId} reason=${safeErrorCode(error)}`, - ); - } - } - } catch (error) { - console.warn( - `[rectification-v9] persist interview before collect attach failed case=${caseId} reason=${safeErrorCode(error)}`, - ); - try { - interviewIdle = await ensureNonTerminalTurnExit({ accounting, userId, caseId }); - } catch (repairError) { - console.warn( - `[rectification-v9] nonterminal exit after idle failure failed case=${caseId} reason=${safeErrorCode(repairError)}`, - ); - } - } - } - if (action === "evidence") { - spokenAnswer = trimSpokenTurnForInterview(answerText, interviewIdle?.terminalNote === true); - if (interviewIdle?.terminalNote && interviewIdle.hostNarration) { - spokenAnswer = composeIdleGapIntoSpoken(spokenAnswer, interviewIdle.hostNarration); - } - if (spokenAnswer !== answerText) { - await emit({ type: "answer.delta", text: spokenAnswer, replace: true }); - } - } - await finalizeTurn("completed", spokenAnswer, outcome.attemptId, outcome.attemptId, true); - - if (!skipBilling) await emit({ type: "billing.settled" }); - await emit({ type: "run.completed", turnId }); - - return { - ok: true, - turnId, - turnStatus: "completed", - skillLoaded: outcome.skillBound, - answerText: spokenAnswer, - phases: skipBilling - ? [...outcome.phases, "run.completed"] - : [...outcome.phases, "billing.settled", "run.completed"], - toolsUsed: outcome.toolsUsed, - errorCode: null, - previousFocusId, - interviewSettled, - }; - - async function streamAttempt( - attemptNumber: number, - attemptId: string, - previousErrorCode: string | null, - ): Promise { - const attemptStartedAt = Date.now(); - const stepTimings: RectificationStepTiming[] = []; - const engineCallsBefore = currentEngineCallTimings().length; - const agent = await buildAgent(turnId, skillPackage, attemptId); - let frameworkSkill: unknown = null; - try { - frameworkSkill = await (agent as unknown as { getSkill(name: string): Promise }).getSkill(skillName); - } catch { - frameworkSkill = null; - } - if (!frameworkSkill) { - return failedAttempt(attemptId, "skill_not_loaded"); - } - - const rawSkillInstructions = (frameworkSkill as { instructions?: unknown }).instructions; - const skillInstructions = typeof rawSkillInstructions === "string" - ? rawSkillInstructions.trim() - : ""; - if (!skillInstructions) { - return failedAttempt(attemptId, "skill_not_loaded"); - } - - const messages = buildAgentMessages( - options, - attemptNumber, - dossier, - skillInstructions, - birthTimeClue, - lastAttemptError, - ); - const maxSteps = resolveRectificationStepBudget(action); - const abortController = new AbortController(); - const onAbort = () => abortController.abort(); - signal?.addEventListener("abort", onAbort, { once: true }); - let timedOut = false; - const timeout = setTimeout(() => { - timedOut = true; - abortController.abort(); - }, attemptTimeoutForRemainingBudget(runDeadlineAt - Date.now(), attemptCapMs)); - - let skillBound = true; - let caseLoaded = false; - let intentClassified = false; - let streamFailed = false; - let finished = false; - let finishReason: ReturnType | null = null; - let answerText = ""; - const answerDeltas: string[] = []; - const phases: string[] = []; - const toolsUsed = new Set(); - const events: PublicStreamEvent[] = []; - const toolTerminalStatus = new Map(); - let batchToolResult: unknown = null; - let hostFallbackUsed = false; - const emittedKeys = new Set(); - const emittedActivities = new Set(); - const repeatedCalls = new Map(); - let phaseSequence = 0; - const rangeBeforeCompare = previousInferenceFromReceipt( - dossier.latestResult?.decisionReceipt ?? null, - )?.credible_range ?? null; - - const recordPhase = async (phase: string, tool: string | null = null) => { - if ( - phase === "answer.delta" - || phase === "thinking.delta" - || phase === "activity.changed" - || phase === "choice.applied" - || phase === "attempt.reset" - || emittedKeys.has(`${phase}:${tool ?? ""}`) - ) return; - emittedKeys.add(`${phase}:${tool ?? ""}`); - phases.push(phase); - phaseSequence += 1; - await persistCommittedPhase(phase, tool, attemptId, phaseSequence); - }; - - const publish = async (event: PublicStreamEvent) => { - events.push(event); - await emit(event); - }; - - try { - await recordPhase("run.started"); - await insertV9SkillRunReceipt( - accounting, - userId, - caseId, - turnId, - attemptId, - "turn", - skillPackage, - ); - await recordPhase("skill.bound"); - await publish({ type: "skill.bound" }); - emittedKeys.add("event:skill.bound::"); - - const generation = agentGenerationSettings(options.generationModel, { - thinking: "enabled", - answerTokens: 8_192, - thinkingTokens: 8_192, - }); - const result = await (agent as unknown as { - stream( - messages: unknown[], - streamOptions: { - maxSteps: number; - abortSignal: AbortSignal; - modelSettings?: { maxOutputTokens?: number }; - providerOptions?: Record; - prepareStep: (input: { stepNumber: number }) => { - activeTools: string[]; - toolChoice: "auto"; - }; - }, - ): Promise<{ - fullStream: AsyncIterable<{ - type: string; - payload?: { - toolName?: unknown; - text?: unknown; - args?: unknown; - error?: unknown; - stepResult?: { reason?: unknown }; - reason?: unknown; - }; - object?: unknown; - }>; - totalUsage?: Promise>; - }>; - }).stream(messages, { - maxSteps, - abortSignal: abortController.signal, - ...generation, - // Thinking-mode providers reject named/required tool_choice. Restrict - // the first step to read-case and keep tool_choice auto; the runner - // still refuses any other public tool before case.loaded. - prepareStep: ({ stepNumber }) => stepNumber === 0 - ? { - activeTools: ["rectification-read-case"], - toolChoice: "auto", - } - : { - activeTools: [...RECTIFICATION_AGENT_TOOLS], - toolChoice: "auto", - }, - }); - - const stepAnswer = createStepAnswerState(); - - let spokenRaw = ""; - let visibleEmitted = ""; - - const emitVisibleSpoken = async (visible: string) => { - if (!caseLoaded) return; - if (visible === visibleEmitted) { - answerText = visible; - return; - } - if (visible.startsWith(visibleEmitted)) { - const growth = visible.slice(visibleEmitted.length); - if (!growth) return; - answerText = visible; - answerDeltas.push(growth); - visibleEmitted = visible; - await emit({ type: "answer.delta", text: growth }); - return; - } - answerText = visible; - answerDeltas.push(visible); - visibleEmitted = visible; - await emit({ type: "answer.delta", text: visible, replace: true }); - }; - - const publishSpokenStep = async (pieces: readonly string[], live = false) => { - const joined = pieces.join(""); - if (!joined) return; - if (!caseLoaded) return; - const spoken = live ? joined : joined.trim(); - if (!spoken) return; - spokenRaw += spoken; - await emitVisibleSpoken(spokenRaw); - }; - - const retractSpoken = async () => { - const hadVisible = Boolean(visibleEmitted || answerText); - spokenRaw = ""; - visibleEmitted = ""; - answerText = ""; - answerDeltas.length = 0; - if (hadVisible && caseLoaded) { - await emit({ type: "answer.delta", text: "", replace: true }); - } - }; - - const applyHostFallback = async (): Promise => { - if (answerText.trim()) return false; - if (toolTerminalStatus.get("rectification-record-evidence-batch") !== "completed") { - return false; - } - const spoken = composeHostFallbackNarration(batchToolResult ?? {}); - if (!spoken) return false; - hostFallbackUsed = true; - await recordPhase("answer.host_fallback"); - await publish({ type: "answer.host_fallback" }); - await emitVisibleSpoken(spoken); - return true; - }; - - 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 - // evidence / diagnostics / compare had already committed, because - // the model often re-issued compare with the same caseId. Bound - // loops with maxSteps / timeout instead; do not attempt.reset. - let skipDuplicateToolCallReceipt = false; - if (chunk.type === "tool-call") { - if (rawToolName === "rectification-read-case" && !skillBound) { - throw new Error("skill_not_bound"); - } - if (isPublicRectificationToolName(rawToolName) && rawToolName !== "rectification-read-case" && !caseLoaded) { - throw new Error("case_not_loaded"); - } - if (isPublicRectificationToolName(rawToolName)) { - const key = publicToolCallKey(rawToolName, chunk.payload); - const count = (repeatedCalls.get(key) ?? 0) + 1; - repeatedCalls.set(key, count); - skipDuplicateToolCallReceipt = count > 1; - } - } - - const stepEffect = applyStepAnswerChunk( - stepAnswer, - chunk, - isPublicRectificationToolName, - ); - if (stepEffect.kind === "live") await publishSpokenStep([stepEffect.text], true); - if (stepEffect.kind === "publish") await publishSpokenStep(stepEffect.pieces); - if (stepEffect.kind === "retract") await retractSpoken(); - - const activityEvent = skipDuplicateToolCallReceipt - ? null - : mapStreamChunkToActivity(chunk as never); - if (activityEvent) { - await publish(activityEvent); - if (activityEvent.status === "started") { - const changed = activityChangedFromTool(activityEvent.tool); - if (!emittedActivities.has(changed.activity)) { - emittedActivities.add(changed.activity); - await publish(changed); - } - } - if (activityEvent.status === "failed" || activityEvent.status === "completed") { - toolTerminalStatus.set(activityEvent.tool, activityEvent.status); - } - } - const capturedBatch = batchResultFromToolChunk(chunk as never); - if (capturedBatch != null) batchToolResult = capturedBatch; - const phaseEvent = skipDuplicateToolCallReceipt - ? null - : mapStreamChunkToPhase(chunk as never); - if (phaseEvent) { - if (phaseEvent.type === "skill.bound" && !skillBound) { - skillBound = true; - await insertV9SkillRunReceipt( - accounting, - userId, - caseId, - turnId, - attemptId, - "turn", - skillPackage, - ); - } - if (phaseEvent.type === "case.loaded" && !skillBound) { - throw new Error("skill_not_bound"); - } - await recordPhase(phaseEvent.type, phaseEvent.tool ?? null); - const key = `${phaseEvent.type}:${phaseEvent.tool ?? ""}:${(phaseEvent.methods ?? []).join(",")}`; - if (!emittedKeys.has(`event:${key}`)) { - emittedKeys.add(`event:${key}`); - await publish(phaseEvent); - } - if (phaseEvent.type === "case.loaded") { - caseLoaded = true; - if (!intentClassified) { - intentClassified = true; - await recordPhase("intent.classified"); - await publish({ type: "intent.classified" }); - } - } - } - for (const toolName of streamToolNames(chunk as never)) toolsUsed.add(toolName); - if (chunk.type === "error" || chunk.type === "abort") streamFailed = true; - if (chunk.type === "finish") { - finished = true; - finishReason = streamFinishReason(chunk); - } - } - } catch (error) { - if (!timedOut && !abortController.signal.aborted) throw error; - streamFailed = true; - } - - if (finished && !streamFailed) { - const flushReason = finishReason === "length" - ? "length" - : finishReason === "stop" || finishReason === "unknown" || finishReason === null - ? "stop" - : finishReason; - const flushed = flushStepAnswerOnStreamFinish(stepAnswer, flushReason); - if (flushed.kind === "publish") await publishSpokenStep(flushed.pieces); - } - - if (toolTerminalStatus.get("rectification-compare-candidates") === "failed" && !batchRescoreFailed(batchToolResult)) { - await emitVisibleSpoken(withCompareFailedRetryNotice(answerText)); - } - - const completeAttempt = async (settleBilling = true): Promise => { - let inputTokens = 0; - let outputTokens = 0; - let cache: ReturnType = null; - try { - const raw = await (result.totalUsage ?? Promise.resolve({ inputTokens: 0, outputTokens: 0 })); - inputTokens = Math.max(0, Math.trunc(typeof raw.inputTokens === "number" ? raw.inputTokens : 0)); - outputTokens = Math.max(0, Math.trunc(typeof raw.outputTokens === "number" ? raw.outputTokens : 0)); - cache = promptCacheUsage(raw); - } catch { - // Timeout/abort can leave provider usage unread. - } - await recordPhase("answer.composed"); - await publish({ type: "answer.composed" }); - return { - ok: true, - status: "completed", - errorCode: null, - usage: { inputTokens, outputTokens, ...(cache ? { cache } : {}) }, - answerText, - answerDeltas, - phases, - toolsUsed: [...toolsUsed], - events, - skillBound, - caseLoaded, - attemptId, - settleBilling, - }; - }; - - if (!skillBound) return failedAttempt(attemptId, "skill_not_loaded"); - if (!caseLoaded) return failedAttempt(attemptId, "case_not_loaded"); - const mapped = mapModelFinishToErrorCode({ - finishReason, - aborted: abortController.signal.aborted, - timedOut, - answerText, - stepCount: toolsUsed.size, - maxSteps, - }); - if (mapped === "run_timeout") return failedAttempt(attemptId, "run_timeout"); - if (streamFailed || abortController.signal.aborted) return failedAttempt(attemptId, mapped ?? "stream_aborted"); - if (!finished) return failedAttempt(attemptId, mapped ?? "stream_unfinished"); - if (mapped === "answer_truncated") { - if (!await applyHostFallback()) { - return { - ok: false, - status: "failed", - errorCode: "answer_truncated", - usage: { inputTokens: 0, outputTokens: 0 }, - answerText, - answerDeltas, - phases, - toolsUsed: [...toolsUsed], - events, - skillBound, - caseLoaded, - attemptId, - }; - } - } else if (mapped === "max_steps" || mapped === "provider_error") { - if (mapped !== "max_steps" || !await applyHostFallback()) { - return failedAttempt(attemptId, mapped); - } - } else if (!answerText.trim() && !await applyHostFallback()) { - const couldRetry = !publicWriteToolCompleted(toolTerminalStatus) && attemptNumber < MAX_ATTEMPTS; - const hasRetryBudget = canStartRetryAttempt(runDeadlineAt - Date.now(), minRetryAttemptMs); - if (couldRetry && hasRetryBudget) { - return { - ok: false, - status: "retryable", - errorCode: "empty_stream", - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [...phases], - toolsUsed: [...toolsUsed], - events, - skillBound, - caseLoaded, - attemptId, - }; - } - if (couldRetry && !hasRetryBudget) { - hostFallbackUsed = true; - await recordPhase("answer.host_fallback"); - await publish({ type: "answer.host_fallback" }); - await emitVisibleSpoken( - turnExpectsEvidenceWrite(action, options.expectedWrite) - ? RECTIFICATION_USER_COPY.evidenceNotRecorded - : RECTIFICATION_USER_COPY.hostNarrationFallback, - ); - return completeAttempt(false); - } - return { - ok: false, - status: "failed", - errorCode: "empty_stream", - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [...phases], - toolsUsed: [...toolsUsed], - events, - skillBound, - caseLoaded, - attemptId, - }; - } - const needsWrite = turnExpectsEvidenceWrite(action, options.expectedWrite) - || (action !== "opening" && action !== "read_only" && answerClaimsEvidenceRecorded(answerText)); - if (needsWrite && !publicWriteToolCompleted(toolTerminalStatus) && !hostFallbackUsed) { - await retractSpoken(); - if ( - attemptNumber < MAX_ATTEMPTS - && canStartRetryAttempt(runDeadlineAt - Date.now(), minRetryAttemptMs) - ) { - return { - ok: false, - status: "retryable", - errorCode: "evidence_not_written", - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [...phases], - toolsUsed: [...toolsUsed], - events, - skillBound, - caseLoaded, - attemptId, - }; - } - hostFallbackUsed = true; - await recordPhase("answer.host_fallback"); - await publish({ type: "answer.host_fallback" }); - await emitVisibleSpoken(RECTIFICATION_USER_COPY.evidenceNotRecorded); - return completeAttempt(false); - } - let latestDossier: V9CaseDossier; - try { - latestDossier = await loadV9CaseDossier(accounting, userId, caseId); - } catch { - return failedAttempt(attemptId, "state_invariant_failed"); - } - const decision = decideFromDossier(latestDossier); - if (decision.nextAction === "ask_candidate_discriminator") { - const openFocus = latestDossier.conversationSummary.activeFocus; - if (!( - openFocus - && isPersistedFocusId(openFocus.id) - && parseAgentChoiceCopy(openFocus.expectedAnswerSchema) - )) { - return failedAttempt(attemptId, "state_invariant_failed"); - } - } - const askedFocus = latestDossier.conversationSummary.activeFocus; - if (askedFocus?.askedTurnId === turnId) { - const stem = focusSpokenPrompt(askedFocus.expectedAnswerSchema); - if (stem) { - const stripped = stripQuestionSentences(answerText, stem); - const next = stripped || RECTIFICATION_USER_COPY.collectHandoff; - if (next !== answerText) answerText = next; - } - } - if (action === "opening") { - const openingRange = openingRangeFromCandidateRange(latestDossier.case.candidateRange); - if (!isAcceptableOpeningBody(answerText)) { - answerText = openingSpokenBody(openingRange); - } - } - const rangeAfterEvidence = previousInferenceFromReceipt( - latestDossier.latestResult?.decisionReceipt ?? null, - )?.credible_range ?? null; - if (batchRescoreFailed(batchToolResult)) { - answerText = withRescoreSkippedNotice(answerText); - } else { - answerText = withRangeChangedAfterEvidence( - answerText, - rangeBeforeCompare, - rangeAfterEvidence, - ); - } - answerText = stripVerbalWindowChange(answerText) - || RECTIFICATION_USER_COPY.declaredWindowLockedReply; - if (answerText !== visibleEmitted) await emitVisibleSpoken(answerText); - return completeAttempt(); - } finally { - clearTimeout(timeout); - signal?.removeEventListener("abort", onAbort); - const stepSummary = diagnosticStepsFromTimings(stepTimings); - const diagnostic: RectificationRunDiagnostic = { - runId: attemptId, - modelId: options.modelName, - finishReason: finishReason ?? "unknown", - 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 })); - } - } - - function failedAttempt(attemptId: string, errorCode: string): AttemptOutcome { - return { - ok: false, - status: isRetryableError(errorCode) ? "retryable" : "failed", - errorCode, - usage: { inputTokens: 0, outputTokens: 0 }, - answerText: "", - answerDeltas: [], - phases: [], - toolsUsed: [], - events: [], - skillBound: false, - caseLoaded: false, - attemptId, - }; - } - - async function persistCommittedPhase( - phase: string, - toolName: string | null, - attemptId: string, - sequence: number, - strict = false, - ) { - try { - await insertV9RunPhase( - accounting, - userId, - caseId, - turnId, - phase, - toolName, - sequence, - attemptId, - ); - } catch (error) { - if (strict) throw error; - // Non-terminal activity receipts remain best effort. The completion - // receipts above are strict because they are part of business truth. - } - } - - async function finalizeTurn( - status: AttemptStatus, - assistantText: string | null, - attemptId: string, - successfulAttemptId: string | null, - strict = false, - ) { - try { - await rpcOf(accounting, "finalize_agentic_rectification_turn", { - p_user_id: userId, - p_case_id: caseId, - p_turn_id: turnId, - p_attempt_id: attemptId, - p_status: status, - p_assistant_message: status === "completed" ? assistantText : null, - p_successful_attempt_id: successfulAttemptId, - }); - } catch (error) { - if (strict) throw error; - console.warn(`[rectification-v10] turn finalize failed turn=${turnId} status=${status} reason=${safeErrorCode(error)}`); - } - } -} - -export function buildAgentMessages( - options: V9AgentRunOptions, - attempt: number, - dossier: V9CaseDossier, - skillInstructions: string, - birthTimeClue: string | null = null, - previousErrorCode: string | null = null, -): unknown[] { - const timeContext = options.timeContext - ?? `服务端当前时间(权威):${new Date().toISOString()}。涉及“现在、今天、今年、未来几个月”等相对时间时,以此为准。`; - const caseContext = `【服务端 Case ID】${options.caseId}。所有 rectification 工具调用的 caseId 必须原样使用此值。`; - const bootstrapContent = [ - "【服务器已绑定当前 Case 的精确 Skill】运行器已在本 attempt 内加载并核验下列指令;不要重复调用 skill。第一步必须调用 rectification-read-case。", - skillInstructions, - ...(attempt > 1 ? [retryConstraintForAttempt(previousErrorCode)] : []), - ].join("\n\n"); - const bootstrap = cachedSystemMessage(bootstrapContent, options.generationModel) - ?? { role: "system" as const, content: bootstrapContent }; - if (options.action === "opening") { - return [bootstrap, { - role: "user", - content: [timeContext, caseContext, buildOpeningBrief(dossier, birthTimeClue)].join("\n"), - }]; - } - return [bootstrap, { - role: "user", - content: [timeContext, caseContext, options.message ?? ""].join("\n"), - }]; + const outcome = await runV9AttemptsWithRetry(turn); + return finishV9AgentTurn(turn, outcome); } export { RECTIFICATION_SKILL_NAME, RECTIFICATION_SKILL_VERSION }; diff --git a/frontend/tests/agent-voice-copy-contract.test.ts b/frontend/tests/agent-voice-copy-contract.test.ts index 6a4452a4..72c6b846 100644 --- a/frontend/tests/agent-voice-copy-contract.test.ts +++ b/frontend/tests/agent-voice-copy-contract.test.ts @@ -30,6 +30,7 @@ import { attachQuestionsToTurns } from "../src/lib/rectification-agentic/v9/turn import { CASE_ID, FOCUS_ID, TURN_ID } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; /** Task 0 frozen “before” snapshot: accident case f83d9b42, origin/staging @ 8617eb56. */ export const ACCIDENT_CASE_BEFORE_COPY = { @@ -253,7 +254,8 @@ test("settled assistant body with a focus has no question-mark sentences", () => assert.doesNotMatch(withRescoreSkippedNotice("记下了:2016 年 3 月入学。范围在收窄。"), /收窄/); assert.doesNotMatch(RECTIFICATION_USER_COPY.evidenceNotRecorded, /系统|模型/); assert.doesNotMatch(RECTIFICATION_USER_COPY.collectHandoff, /请回答下面的问题/); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.match(agentRun, /stripQuestionSentences/); assert.match(agentRun, /withRangeChangedAfterEvidence/); }); diff --git a/frontend/tests/rectification-agent-run-surface.ts b/frontend/tests/rectification-agent-run-surface.ts new file mode 100644 index 00000000..8e12249b --- /dev/null +++ b/frontend/tests/rectification-agent-run-surface.ts @@ -0,0 +1,29 @@ +import { readFileSync } from "node:fs"; + +/** + * The V10 rectification turn runner as source text: agent-run.ts (public + * types and the order of the phases) plus the phase modules it was split into + * (TASK-rectification-code-split-20260926), in the order the code ran inside + * the old runV9AgentTurn: preparation, the retry loop, the finish, one + * attempt, then the shared pieces and the model messages. A whole-source + * contract that held on agent-run.ts before the split holds on this + * concatenation after it. Slices that no longer sit in one file are + * rewritten per test. + */ +export const rectificationAgentRunFiles = [ + "../src/lib/rectification-agentic/v9/agent-run.ts", + "../src/lib/rectification-agentic/v9/agent-run-prepare.ts", + "../src/lib/rectification-agentic/v9/agent-run-retry.ts", + "../src/lib/rectification-agentic/v9/agent-run-finish.ts", + "../src/lib/rectification-agentic/v9/agent-run-attempt.ts", + "../src/lib/rectification-agentic/v9/agent-run-support.ts", + "../src/lib/rectification-agentic/v9/agent-run-messages.ts", +] as const; + +export function readRectificationAgentRunFile(relativePath: (typeof rectificationAgentRunFiles)[number]): string { + return readFileSync(new URL(relativePath, import.meta.url), "utf8"); +} + +export const rectificationAgentRunSurface = rectificationAgentRunFiles + .map((relativePath) => readRectificationAgentRunFile(relativePath)) + .join("\n"); diff --git a/frontend/tests/rectification-agentic-entry.test.ts b/frontend/tests/rectification-agentic-entry.test.ts index fb01f8b1..cd9e5853 100644 --- a/frontend/tests/rectification-agentic-entry.test.ts +++ b/frontend/tests/rectification-agentic-entry.test.ts @@ -17,6 +17,7 @@ import { deriveRectificationChatView } from "../src/lib/rectification-chat-view. import { completedReceiptFromPersisted } from "../src/lib/rectification-chat-messages.ts"; import { RectificationQuestionGapNotices } from "../src/components/rectification-question-gap-notices.tsx"; import { rectificationAgentRouteSurface } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const component = readFileSync( new URL("../src/components/conversational-birth-time-rectification.tsx", import.meta.url), @@ -288,10 +289,8 @@ test("candidate acceptance is non-billable, mutually exclusive, and continues th }); test("usage completes or releases without hiding settlement failures", () => { - const run = readFileSync( - new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), - "utf8", - ); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const run = rectificationAgentRunSurface; assert.match(run, /billing\.complete\(/); assert.match(run, /billing\.release\(/); assert.match(run, /usage_settlement_failed/); diff --git a/frontend/tests/rectification-answer-choice.test.ts b/frontend/tests/rectification-answer-choice.test.ts index 914f2893..28d064f5 100644 --- a/frontend/tests/rectification-answer-choice.test.ts +++ b/frontend/tests/rectification-answer-choice.test.ts @@ -55,6 +55,7 @@ import { } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface, readRectificationAgentRouteFile } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const ACTION_ID = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa"; const QUESTION_ID = "question-1"; @@ -1508,7 +1509,8 @@ test("the public agent route treats structured choice as a non-model command", ( }); test("rectification attempt timeout stays under the agent route budget", () => { - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const budget = readFileSync(new URL("../src/lib/rectification-run-budget.ts", import.meta.url), "utf8"); const regenerate = readFileSync( new URL("../src/app/api/rectification/cases/[caseId]/turns/[turnId]/regenerate/route.ts", import.meta.url), diff --git a/frontend/tests/rectification-collect-prompt.test.ts b/frontend/tests/rectification-collect-prompt.test.ts index b21c2ad2..4f170148 100644 --- a/frontend/tests/rectification-collect-prompt.test.ts +++ b/frontend/tests/rectification-collect-prompt.test.ts @@ -6,6 +6,7 @@ import { composeCollectSpokenAssistantText, detachCollectSpokenAssistantText, st import { attachQuestionsToTurns } from "../src/lib/rectification-agentic/v9/turn-question.ts"; import { GENERIC_COLLECT_QUESTION, RECTIFICATION_USER_COPY, USER_COLLECT_QUESTION } from "../src/lib/rectification-agentic/user-copy.ts"; import { CASE_ID, FOCUS_ID, TURN_ID } from "./rectification-v9-test-support.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; test("composeCollectSpokenAssistantText joins by exact prompt identity", () => { const greeting = "你好,我是生时校正助手。"; @@ -89,7 +90,8 @@ test("composeCollectSpokenAssistantText drops a near-duplicate restatement of th }); test("runtime no longer composes the stem into assistant_message", () => { - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const attach = readFileSync(new URL("../src/lib/rectification-agentic/v9/turn-question.ts", import.meta.url), "utf8"); assert.doesNotMatch(agentRun, /composeCollectSpokenAssistantText/); assert.match(attach, /stripQuestionSentences/); diff --git a/frontend/tests/rectification-delivery-ui-simplify-20260908.test.ts b/frontend/tests/rectification-delivery-ui-simplify-20260908.test.ts index f08c0cc0..638db467 100644 --- a/frontend/tests/rectification-delivery-ui-simplify-20260908.test.ts +++ b/frontend/tests/rectification-delivery-ui-simplify-20260908.test.ts @@ -17,6 +17,7 @@ import { import { RectificationVerificationReport } from "../src/components/rectification-verification-report.tsx"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; test("delivery body is three sentences and omits the eight-method report", () => { const text = deliveryTurnNarration({ @@ -74,7 +75,8 @@ test("already_delivered blocks a second no-message run for the same result withi }); test("agent-run returns already_delivered before billing and the route skips a second turn", () => { - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; // 原值: 读 app/api/rectification/agent/route.ts 单文件。新值: rectificationAgentRouteSurface(route.ts + 拆出的 agent-route-*.ts,见 rectification-agent-route-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 const route = rectificationAgentRouteSurface; const exit = readFileSync(new URL("../src/lib/rectification-agentic/v9/turn-exit.ts", import.meta.url), "utf8"); @@ -94,7 +96,8 @@ test("agent-run returns already_delivered before billing and the route skips a s }); test("agent-run trims evidence after persistNextInterviewIfIdle so delivery can keep three sentences", () => { - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const persistIdx = agentRun.lastIndexOf("persistNextInterviewIfIdle({"); const trimIdx = agentRun.lastIndexOf("trimSpokenTurnForInterview("); const replaceIdx = agentRun.lastIndexOf('type: "answer.delta", text: spokenAnswer, replace: true'); diff --git a/frontend/tests/rectification-growth-contract.test.ts b/frontend/tests/rectification-growth-contract.test.ts index 80ffaadf..4836e9c6 100644 --- a/frontend/tests/rectification-growth-contract.test.ts +++ b/frontend/tests/rectification-growth-contract.test.ts @@ -122,3 +122,41 @@ test("rectification agent POST stays an assembly of branch handlers", () => { assert.match(post, new RegExp(`\\b${handler}\\(`), `POST must delegate to ${handler}`); } }); + +// --------------------------------------------------------------------------- +// T3 · lib/rectification-agentic/v9/agent-run.ts +// Measured 2026-09-26 after the runner split. Before (origin/staging 12cbe6f8): +// file 1400 lines, runV9AgentTurn 1046 lines with four nested functions +// (streamAttempt, failedAttempt, persistCommittedPhase, finalizeTurn). After: +// file 135, runV9AgentTurn 25 lines, no nested functions; the phases live in +// agent-run-{prepare,retry,attempt,finish,support,messages}.ts. +const RUN_PATH = "../src/lib/rectification-agentic/v9/agent-run.ts"; +const RUN_LINE_BASELINE = 135; +const RUN_LINE_CAP = RUN_LINE_BASELINE + 150; +/** D2 of the task: runV9AgentTurn stays at or under 300 lines. */ +const RUN_FUNCTION_LINE_CAP = 300; +const RUN_FUNCTION_MARKER = "export async function runV9AgentTurn("; + +function functionBody(source: string, marker: string): string { + const start = source.indexOf(marker); + assert.notEqual(start, -1, `missing marker: ${marker}`); + const end = source.indexOf("\n}\n", start); + assert.notEqual(end, -1, `unterminated function: ${marker}`); + return source.slice(start, end + 2); +} + +test("agent-run.ts line count stays within the coarse guardrail", () => { + const n = lineCount(read(RUN_PATH)); + assert.ok(n <= RUN_LINE_CAP, `agent-run.ts has ${n} lines; cap is ${RUN_LINE_CAP} (${RUN_LINE_BASELINE} baseline + 150).`); +}); + +test("runV9AgentTurn stays the order of its phases", () => { + const body = functionBody(read(RUN_PATH), RUN_FUNCTION_MARKER); + const n = lineCount(body); + assert.ok(n <= RUN_FUNCTION_LINE_CAP, `runV9AgentTurn has ${n} lines; cap is ${RUN_FUNCTION_LINE_CAP}. Phase logic belongs in agent-run-*.ts.`); + assert.doesNotMatch(body, /^\s+(?:async )?function\s/m, "runV9AgentTurn must not grow nested functions again."); + const prepare = body.indexOf("prepareV9AgentTurn("); + const attempts = body.indexOf("runV9AttemptsWithRetry("); + const finish = body.indexOf("finishV9AgentTurn("); + assert.ok(prepare >= 0 && attempts > prepare && finish > attempts, "prepare → attempts → finish"); +}); diff --git a/frontend/tests/rectification-latency-20260926.test.ts b/frontend/tests/rectification-latency-20260926.test.ts index b4bb9dde..d138c1ad 100644 --- a/frontend/tests/rectification-latency-20260926.test.ts +++ b/frontend/tests/rectification-latency-20260926.test.ts @@ -60,6 +60,7 @@ import { finalizeSuccessfulTurnExit } from "../src/lib/rectification-agentic/v9/ import { RECTIFICATION_TURN_PROGRESS_LABELS } from "../src/lib/rectification-activity-labels.ts"; import { rectificationInitialLiveLabel } from "../src/lib/rectification-surface-state.ts"; import { rectificationAgentRouteSurface } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const dummyModel = { id: "test-model" } as ResolvedLanguageModel; const PROGRESS_LINES = Object.values(RECTIFICATION_TURN_PROGRESS_LABELS); @@ -527,7 +528,8 @@ test("D5: the exit gate skips its idle re-read only when told the run settled th }); 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"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; // 原值: 读 app/api/rectification/agent/route.ts 单文件。新值: rectificationAgentRouteSurface(route.ts + 拆出的 agent-route-*.ts,见 rectification-agent-route-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 const route = rectificationAgentRouteSurface; assert.match(agentRun, /interviewSettled = interviewIdle\.focusActive === true;/); diff --git a/frontend/tests/rectification-occupation-coverage-exit.test.ts b/frontend/tests/rectification-occupation-coverage-exit.test.ts index 20efe457..05cacef1 100644 --- a/frontend/tests/rectification-occupation-coverage-exit.test.ts +++ b/frontend/tests/rectification-occupation-coverage-exit.test.ts @@ -1,5 +1,4 @@ import assert from "node:assert/strict"; -import { readFileSync } from "node:fs"; import test from "node:test"; import { OPEN_ENGINE_CAPABILITY_CEILING } from "./rectification-v9-test-support.ts"; @@ -34,6 +33,7 @@ import { } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface, rectificationTypedMessageFastPath } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const STYLE_OPTIONS = [ { label: "明确发生且时间吻合", answer_class: "yes" as const }, @@ -828,7 +828,8 @@ test("idle persist on the covered-domain accident delivers adopt and writes no f }); test("agent-run still writes the exhaustion gate when idle persist returns terminalNote", () => { - const agent = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agent = rectificationAgentRunSurface; // 原值: 读 rectification-agentic-chat.tsx 单文件。新值: rectificationChatSurface(容器 + 拆出的 hook / 参数函数 / 块,见 rectification-chat-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改为调用函数。 const chat = rectificationChatSurface; assert.match(agent, /interviewIdle\.terminalNote && interviewIdle\.hostNarration/); diff --git a/frontend/tests/rectification-other-collect-fallback-20260908.test.ts b/frontend/tests/rectification-other-collect-fallback-20260908.test.ts index 3d97ce33..dcc50423 100644 --- a/frontend/tests/rectification-other-collect-fallback-20260908.test.ts +++ b/frontend/tests/rectification-other-collect-fallback-20260908.test.ts @@ -31,6 +31,7 @@ import { receiptHandlers, } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; type ExecutableTool = { execute(input: unknown): Promise; @@ -325,7 +326,8 @@ test("orphan collect:other focus is skipped so idle persist can deliver", async }); test("agent-run writes the exhaustion gate when idle persist returns terminalNote", () => { - const agent = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agent = rectificationAgentRunSurface; // 原值: 读 rectification-agentic-chat.tsx 单文件。新值: rectificationChatSurface(容器 + 拆出的 hook / 参数函数 / 块,见 rectification-chat-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改为调用函数。 const chat = rectificationChatSurface; assert.match(agent, /interviewIdle\.terminalNote && interviewIdle\.hostNarration/); diff --git a/frontend/tests/rectification-probe-replay-loss-20260908.test.ts b/frontend/tests/rectification-probe-replay-loss-20260908.test.ts index 5c749a2a..0600502a 100644 --- a/frontend/tests/rectification-probe-replay-loss-20260908.test.ts +++ b/frontend/tests/rectification-probe-replay-loss-20260908.test.ts @@ -1,5 +1,4 @@ import assert from "node:assert/strict"; -import { readFileSync } from "node:fs"; import test from "node:test"; import { @@ -23,6 +22,7 @@ import { import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { deriveRectificationChatView } from "../src/lib/rectification-chat-view.ts"; import { chatViewInput, offerableCandidateResult, settledAssistant } from "./rectification-chat-run-fixtures.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const S1_TIMES = [ "04:50", "04:51", "04:52", "04:53", "04:54", "04:55", "04:56", "04:57", "05:13", @@ -286,10 +286,8 @@ test("BUG-588 evidence narration names a widened range after compare", () => { ); assert.equal(rangeChangedAfterEvidence(["04:50", "04:57"], ["04:50", "04:57"]), null); - const agentRun = readFileSync( - new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), - "utf8", - ); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const askedFocusAt = agentRun.indexOf("const askedFocus = latestDossier.conversationSummary.activeFocus"); const settlement = askedFocusAt >= 0 ? agentRun.slice(askedFocusAt) : ""; const stripAt = settlement.indexOf("stripQuestionSentences"); diff --git a/frontend/tests/rectification-skipped-health-deadend-20260909.test.ts b/frontend/tests/rectification-skipped-health-deadend-20260909.test.ts index 1a288637..61741d2f 100644 --- a/frontend/tests/rectification-skipped-health-deadend-20260909.test.ts +++ b/frontend/tests/rectification-skipped-health-deadend-20260909.test.ts @@ -23,6 +23,7 @@ import { type MethodFollowupEvidence, } from "../src/lib/rectification-agentic/v9/method-followup.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const SEVEN_WITHOUT_HEALTH: MethodFollowupEvidence[] = [ { id: "e-edu", status: "confirmed", domain: "education", datePrecision: "month", occurredFrom: "2016-09-01", occurredTo: null, eventKind: "education_start" }, @@ -149,7 +150,8 @@ test("unavailable gap button posts repair-exit and says 接着问", () => { new URL("../src/app/api/rectification/cases/[caseId]/repair-exit/route.ts", import.meta.url), "utf8", ); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.equal(RECTIFICATION_QUESTION_RELOAD_LABEL, "接着问"); assert.match(surface, /RECTIFICATION_QUESTION_RELOAD_LABEL = "接着问"/); assert.match(surface, /暂时接不上,请新建一次校正/); diff --git a/frontend/tests/rectification-spoken-answer.test.ts b/frontend/tests/rectification-spoken-answer.test.ts index 79fe2ecc..c21c1b5b 100644 --- a/frontend/tests/rectification-spoken-answer.test.ts +++ b/frontend/tests/rectification-spoken-answer.test.ts @@ -1,14 +1,12 @@ import assert from "node:assert/strict"; -import { readFileSync } from "node:fs"; import test from "node:test"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; // 原值: 读 rectification-agentic-chat.tsx 单文件。新值: rectificationChatSurface(容器 + 拆出的 hook / 参数函数 / 块,见 rectification-chat-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改为调用函数。 const chat = rectificationChatSurface; -const agentRun = readFileSync( - new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), - "utf8", -); +// 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 +const agentRun = rectificationAgentRunSurface; test("live and persisted assistant text stay verbatim after spoken-answer removal", () => { assert.match(chat, /const raw = failed \? "" : turn\.text \?\? ""/); diff --git a/frontend/tests/rectification-spoken-collect.test.ts b/frontend/tests/rectification-spoken-collect.test.ts index 4a13ae7a..fb369542 100644 --- a/frontend/tests/rectification-spoken-collect.test.ts +++ b/frontend/tests/rectification-spoken-collect.test.ts @@ -39,6 +39,7 @@ import type { DiscriminatingEventProbe } from "../src/lib/rectification-agentic/ import { collectFocusCloseStatus } from "../src/lib/rectification-agentic/v9/turn-intent-classifier.ts"; import { rectificationChatSurface, readRectificationChatFile } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface, rectificationTypedMessageFastPath } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const RELATIONSHIP_PROMPT = USER_COLLECT_QUESTION.relationship; @@ -162,7 +163,8 @@ test("collect_spoken stem lives on turn.question inside the same assistant artic const chat = rectificationChatSurface; const row = readFileSync(new URL("../src/components/chat-message-row.tsx", import.meta.url), "utf8"); const styles = readFileSync(new URL("../src/app/globals.css", import.meta.url), "utf8"); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.match(chat, /currentQuestion\?\.kind === "collect_spoken"/); assert.match(chat, /const collectSpokenPrompt =/); assert.match(chat, /parseTurnQuestion\(turn\.question\)/); @@ -482,7 +484,8 @@ test("agent route keeps question ownership in the server Case projection", () => ); assert.match(successExit, /await finalizeSuccessfulTurnExit/); assert.doesNotMatch(successExit, /send\(\{\s*type:\s*"error"/); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.doesNotMatch(agentRun, /collectSpokenPromptForNewFocus/); assert.doesNotMatch(agentRun, /composeCollectSpokenAssistantText/); assert.match(agentRun, /askedTurnId: turnId/); @@ -578,7 +581,8 @@ test("after one or two dated events the next collect is an invite, not a domain assert.match(tools, /recordV10EvidenceBatch/); assert.match(tools, /result\.acceptedCount > 0[\s\S]*autoRescoreAfterEvidenceChange/); const agent = readFileSync(new URL("../src/mastra/agentic-rectification.ts", import.meta.url), "utf8"); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.doesNotMatch(agent, /还能想起别的吗/); assert.doesNotMatch(agentRun, /还能想起别的吗/); }); diff --git a/frontend/tests/rectification-stale-compare-fix-20260907.test.ts b/frontend/tests/rectification-stale-compare-fix-20260907.test.ts index 9ab04d9a..8c587e03 100644 --- a/frontend/tests/rectification-stale-compare-fix-20260907.test.ts +++ b/frontend/tests/rectification-stale-compare-fix-20260907.test.ts @@ -37,6 +37,7 @@ import { fakeAccounting, receiptHandlers, } from "./rectification-v9-test-support.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const ACTION_ID = "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa"; const QUESTION_ID = "question-1"; @@ -219,7 +220,8 @@ test("compare failure copy and receipt detail stay user-visible without PII", () }); assert.equal(detail?.safe_error_code, "engine_request_failed"); assert.match(String(detail?.engine_message), /asked_probe_keys/); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.match(agentRun, /withCompareFailedRetryNotice/); assert.match(agentRun, /rectification-compare-candidates/); const tools = readFileSync(new URL("../src/mastra/rectification-v9-tools.ts", import.meta.url), "utf8"); diff --git a/frontend/tests/rectification-unwritten-evidence.test.ts b/frontend/tests/rectification-unwritten-evidence.test.ts index 0445d2a3..d099cad9 100644 --- a/frontend/tests/rectification-unwritten-evidence.test.ts +++ b/frontend/tests/rectification-unwritten-evidence.test.ts @@ -32,6 +32,7 @@ import { } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; import { rectificationAgentRouteSurface, readRectificationAgentRouteFile } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const REJECTION = { error: true, @@ -164,10 +165,8 @@ test("retry bootstrap uses the unwritten-evidence constraint only for that error test("route passes classifier expectedWrite into the agent runner", () => { // 原值: 读 app/api/rectification/agent/route.ts 单文件。新值: rectificationAgentRouteSurface(route.ts + 拆出的 agent-route-*.ts,见 rectification-agent-route-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 const route = rectificationAgentRouteSurface; - const agentRun = readFileSync( - new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), - "utf8", - ); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; assert.match(route, /classifyTurnIntentWithRetry/); // 原值: route 含 `expectedWrite,`(runV9AgentTurn 参数)与 `let expectedWrite: "evidence" | "none" | "unknown" = "none"`。 // 新值: 分类结果放在 turnState 上:类型 `expectedWrite: "evidence" | "none" | "unknown"`,route.ts 初值 `expectedWrite: "none"`, diff --git a/frontend/tests/rectification-v9-stream.test.ts b/frontend/tests/rectification-v9-stream.test.ts index bf2f7611..21ef4e72 100644 --- a/frontend/tests/rectification-v9-stream.test.ts +++ b/frontend/tests/rectification-v9-stream.test.ts @@ -31,6 +31,7 @@ import { RECTIFICATION_SKILL_NAME } from "../src/lib/rectification-agentic/v9/ca import { messageContentHash } from "../src/lib/rectification-agentic/v9/message-origin.ts"; import { safeToolErrorCode } from "../src/lib/rectification-agentic/v9/tool-service.ts"; import { rectificationAgentRouteSurface } from "./rectification-agent-route-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; type StreamChunk = { type: string; @@ -1868,7 +1869,8 @@ test("whole-run budget is less than both route maxDuration values", () => { assert.ok(RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS <= RECTIFICATION_RUN_BUDGET_MS); assert.ok(RECTIFICATION_MIN_RETRY_ATTEMPT_MS > 0); const budget = readFileSync(new URL("../src/lib/rectification-run-budget.ts", import.meta.url), "utf8"); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const labels = readFileSync(new URL("../src/lib/rectification-activity-labels.ts", import.meta.url), "utf8"); assert.match(budget, /export const RECTIFICATION_AGENT_ATTEMPT_TIMEOUT_MS = 210_000/); assert.doesNotMatch(agentRun, /210_000/); diff --git a/frontend/tests/rectification-year-focus-overlay-20260916.test.ts b/frontend/tests/rectification-year-focus-overlay-20260916.test.ts index 931fc7db..25d01b0d 100644 --- a/frontend/tests/rectification-year-focus-overlay-20260916.test.ts +++ b/frontend/tests/rectification-year-focus-overlay-20260916.test.ts @@ -25,6 +25,7 @@ import { fakeAccounting, } from "./rectification-v9-test-support.ts"; import { rectificationChatSurface } from "./rectification-chat-surface.ts"; +import { rectificationAgentRunSurface } from "./rectification-agent-run-surface.ts"; const YEAR_Q = "collect:targeted:health_pressure:year"; const SERVER_PROMPT = TARGETED_YEAR_PROMPT; @@ -135,7 +136,8 @@ test("yearless short replies are held by the host pre-step (BUG-910)", () => { }; assert.equal(shouldHostReaskYearEntry(focus, "是"), SERVER_PROMPT); assert.equal(shouldHostReaskYearEntry(focus, "2019年3月做过手术"), null); - const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); + // 原值: 读 lib/rectification-agentic/v9/agent-run.ts 单文件。新值: rectificationAgentRunSurface(agent-run.ts + 拆出的 agent-run-*.ts,见 rectification-agent-run-surface.ts)。原因: TASK-rectification-code-split-20260926 只搬不改,整文件断言跟着代码走;切片断言另行改写。 + const agentRun = rectificationAgentRunSurface; const reliability = agentRun.indexOf("classifyDateReliabilityUtterance"); const yearHost = agentRun.indexOf("shouldHostReaskYearEntry"); const reserve = agentRun.indexOf("billing.reserve");