From 1530a0dd631873bd97b61db6685e70053a813571 Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Sat, 26 Sep 2026 22:17:13 +0800 Subject: [PATCH] fix(consult): give compose its own clock and never settle a cut answer (BUG-1051) The tool loop and the answer-writing stream shared one 110s AbortSignal. Mastra 1.50 does not throw on abort: it emits an abort chunk and finish(tripwire) and closes normally, so a half-written answer reached onComplete, was charged and persisted as completed. - Compose, length continuation and answer retry run on a 70s answer clock started on first use (worst case 110s + 70s = 180s; maxDuration 240). - Settlement requires finish=stop from the stream that wrote the answer; abort/tripwire, content-filter, tool-calls, other/unknown/error or a missing finish with visible text ends as answer_truncated (cancel, no charge). The abort chunk records an abort runtime step; a cut stream no longer flushes its dangling Pass 4 sentence. length still continues. - [agent-observability] gains composeFinishReason, composeAborted and answerVisibleChars (enum/boolean/count only). - Regression tests use a real Mastra Agent over a fake model; the BUG-305 hand-thrown DOMException fixture is kept with a three-column note, and eleven fixtures gain the finish(stop) chunk real streams always carry. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8 --- frontend/src/app/api/consult/route.ts | 23 +- frontend/src/lib/agent-observability.ts | 6 + frontend/src/lib/stream-agent-response.ts | 121 +++++- frontend/src/mastra/consultation-tools.ts | 40 ++ ...consult-answer-truncation-20260926.test.ts | 360 ++++++++++++++++++ .../consultation-agentic-runtime.test.ts | 24 ++ 6 files changed, 557 insertions(+), 17 deletions(-) create mode 100644 frontend/tests/consult-answer-truncation-20260926.test.ts diff --git a/frontend/src/app/api/consult/route.ts b/frontend/src/app/api/consult/route.ts index 8d27cf4c..9bbb1037 100644 --- a/frontend/src/app/api/consult/route.ts +++ b/frontend/src/app/api/consult/route.ts @@ -47,6 +47,8 @@ import { AGENT_MAX_STEPS, AGENT_TIMEOUT_MS, AGENT_SLICE_MAX_STEPS, + CONSULTATION_COMPOSE_TIMEOUT_MS, + createConsultationAnswerClock, consultationContinueGenerationSettings, consultationGenerationSettings, consultationNatalPrepareStep, @@ -101,12 +103,14 @@ import { generateSessionTitle, shouldGenerateSessionTitle } from "@/lib/session- import { z } from "zod"; export const runtime = "nodejs"; -export const maxDuration = 120; +export const maxDuration = 240; // The step budget, the wall-clock budget and the domain cap all bound this same // run, so they are declared as one group in @/mastra/consultation-tools with the // reasoning that ties them together. maxDuration above is the ceiling they must -// stay under; raising it here without raising that is meaningless. +// stay under: the tool loop (AGENT_TIMEOUT_MS) plus the answer phase's own clock +// (CONSULTATION_COMPOSE_TIMEOUT_MS), 180s, plus setup and settlement (BUG-1051). +// Self-hosted `node server.js` does not enforce it; it documents the ceiling. const chatRequestMetadataSchema = z.object({ requestId: z.string().uuid(), @@ -1021,6 +1025,11 @@ export async function POST(request: Request) { ...streamOptions, prepareStep: consultationWindowPrepareStep, }; + // Writing the answer (compose, length continuation, answer retry) runs on + // its own clock, started when the first of those streams starts. Sharing + // the loop's 110s left compose seconds and Mastra's abort cut it mid-sentence + // (BUG-1051, BUG-944). The loop's tools keep agentAbortSignal. + const answerPhaseSignal = createConsultationAnswerClock(CONSULTATION_COMPOSE_TIMEOUT_MS); async function streamWithOverflowRetry( agent: { stream: ( @@ -1074,7 +1083,7 @@ export async function POST(request: Request) { const retried = await agent.stream([ ...baseMessages, { role: "user" as const, content: "上一轮没有输出任何回答文本。请直接给出这个问题的回答,不要只说明过程。" }, - ], streamOptions); + ], { ...streamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; @@ -1085,6 +1094,7 @@ export async function POST(request: Request) { { role: "user" as const, content: consultationContinuePrompt(output) }, ], { ...streamOptions, + abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), }); usages.push(continued.totalUsage); @@ -1172,7 +1182,7 @@ export async function POST(request: Request) { role: "user" as const, content: "服务器窗口计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。", }, - ], streamOptions); + ], { ...streamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; @@ -1183,6 +1193,7 @@ export async function POST(request: Request) { { role: "user" as const, content: consultationContinuePrompt(output) }, ], { ...streamOptions, + abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), }); usages.push(continued.totalUsage); @@ -1304,7 +1315,7 @@ export async function POST(request: Request) { role: "user" as const, content: "服务器计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。", }, - ], natalStreamOptions); + ], { ...natalStreamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; @@ -1315,6 +1326,7 @@ export async function POST(request: Request) { { role: "user" as const, content: consultationContinuePrompt(output) }, ], { ...streamOptions, + abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), }); usages.push(continued.totalUsage); @@ -1329,6 +1341,7 @@ export async function POST(request: Request) { { role: "user" as const, content: `${consultationComposePrompt()}${retryHint ? `\n${retryHint}` : ""}` }, ], { ...streamOptions, + abortSignal: answerPhaseSignal(), maxSteps: AGENT_SLICE_MAX_STEPS, toolChoice: "none", ...consultationContinueGenerationSettings(selectedModel.model), diff --git a/frontend/src/lib/agent-observability.ts b/frontend/src/lib/agent-observability.ts index 6c188435..4fcafa8c 100644 --- a/frontend/src/lib/agent-observability.ts +++ b/frontend/src/lib/agent-observability.ts @@ -132,6 +132,12 @@ export const agentObservabilityEventSchema = z.object({ // every attempt. Both are enum-like machine values, never provider text. modelFinishReason: z.enum(agentModelFinishReasons).optional(), modelStepCount: countSchema.optional(), + // How the stream that wrote the answer ended, whether Mastra's abort chunk + // arrived in it, and how many visible characters the answer had. Enum, + // boolean and count only: never answer text (BUG-1051). + composeFinishReason: z.enum([...agentModelFinishReasons, "missing"]).optional(), + composeAborted: z.boolean().optional(), + answerVisibleChars: countSchema.optional(), // How many reference documents the model opened after loading the skill, and how many strict-method // sections the server delivered with the evidence. Both are needed to read the other: zero reads is // only a gap in the answer's method if nothing was delivered either. diff --git a/frontend/src/lib/stream-agent-response.ts b/frontend/src/lib/stream-agent-response.ts index 0c5dffc5..1b6addc5 100644 --- a/frontend/src/lib/stream-agent-response.ts +++ b/frontend/src/lib/stream-agent-response.ts @@ -9,7 +9,7 @@ import { type AgentExecutionReceipt, type ConsultationAgentPublicEvent, } from "./consultation-agent-events.ts"; -import { toAgentModelFinishReason } from "./agent-observability.ts"; +import { toAgentModelFinishReason, type AgentModelFinishReason } from "./agent-observability.ts"; import { createVisibleTextTransformer } from "./stream-text-response.ts"; import { consultationWriteLabel } from "./consultation-activity-labels.ts"; import { logTruncatedReasoning } from "./consultation-budget.ts"; @@ -382,6 +382,32 @@ function recordSkillBindingAbort(options: StreamAgentResponseOptions) { }); } +/** + * How one model stream ended. Mastra 1.50 does not throw when the abort signal + * fires mid-answer: it enqueues `{ type: "abort" }`, then `finish` with reason + * `tripwire`, and closes the stream normally. A thrown TimeoutError, which the + * BUG-305 fixture hand-built, never reaches us for that case, so the verdict has + * to be read off the chunks. + */ +type AttemptOutcome = { + finishReason: AgentModelFinishReason | "missing"; + aborted: boolean; + /** The attempt is part of writing the answer (compose, continuation, retry). */ + answerPhase: boolean; + /** The attempt added visible text to the answer. */ + contributed: boolean; +}; + +/** + * Only `stop` means the model finished the answer. `length` is recoverable by + * continuation; everything else with visible text (abort/tripwire, + * content-filter, tool-calls under toolChoice none, other/unknown, a stream + * that closed with no finish chunk) is a cut answer (BUG-1051). + */ +function attemptCut(outcome: AttemptOutcome) { + return outcome.aborted || outcome.finishReason !== "stop"; +} + function sliceAddedVisibleText(before: string, after: string) { return after.length > before.length && /\S/.test(after.slice(before.length)); } @@ -408,6 +434,11 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { const startedAt = new Map(); // A retry reuses these counters so a failure in either attempt is recorded once. const toolErrors = { seen: 0 }; + // The latest attempt, and the latest attempt that wrote (or was asked to + // write) the answer. Settlement is judged on the second: a drained tool loop + // that ended on `tool-calls` says nothing about whether compose finished. + let lastAttempt: AttemptOutcome | null = null; + let answerTail: AttemptOutcome | null = null; const send = (controller: ReadableStreamDefaultController | undefined, event: ConsultationAgentPublicEvent) => { if (!firstActivity && (event.type === "skill.started" || event.type === "tool.started" || event.type === "activity")) { firstActivity = true; @@ -470,14 +501,37 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } } + /** + * Release the open Pass 4 sentence at the end of the answer, unless the stream + * that wrote it was cut. A cut stream's last fragment is not a sentence the + * model finished; flushing it made the half line look like a deliberate end. + * The truncation notice explains the missing rest instead. + */ + async function releaseFinalSentence(controller: ReadableStreamDefaultController | undefined) { + const writer = answerTail ?? lastAttempt; + // Callers run this after any length continuation, so a writer still + // ending on `length` here is cut as well. + if (writer && attemptCut(writer)) { + pass4Buffer = ""; + return; + } + await releasePass4Sentences(controller, "", true); + } + async function consumeAttempt( controller: ReadableStreamDefaultController | undefined, stream: ChunkStream, - attempt: { drainSpoken?: boolean; suppressCompositionActivity?: boolean } = {}, + attempt: { drainSpoken?: boolean; suppressCompositionActivity?: boolean; answerPhase?: boolean } = {}, ) { // Each attempt owns its own uncontracted buffer. Accumulating across the // contract retry delivered the first draft and the retry as one answer. uncontractedText = ""; + const outcome: AttemptOutcome = { + finishReason: "missing", + aborted: false, + answerPhase: Boolean(attempt.answerPhase), + contributed: false, + }; const visible = createVisibleTextTransformer(options.transformText ?? ((value) => value)); let held = ""; let composingSent = Boolean(attempt.suppressCompositionActivity); @@ -496,6 +550,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { uncontractedText = ""; held += text; if (!held) return; + if (/\S/.test(held)) outcome.contributed = true; if (!composingSent) { composingSent = true; send(controller, { type: "activity", phase: "answer-composition", label: "正在组织回答" }); @@ -540,10 +595,21 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } for (const event of mapChunk(chunk, options, startedAt, toolErrors)) send(controller, event); flushThinkingPlan(controller); + if (chunk.type === "abort" && !outcome.aborted) { + // Mastra's own abort chunk: the signal fired and the stream is about + // to close normally. Leave the same trace a thrown abort leaves. + outcome.aborted = true; + appendConsultationRuntimeStep(options.state, { + kind: "abort", + name: attempt.drainSpoken ? "tool-abort" : "compose-abort", + status: "failed", + }); + } if (chunk.type === "step-finish") options.state.modelStepCount += 1; if (chunk.type === "finish") { const finish = finishTelemetry(chunk); options.state.modelFinishReason = finish.reason; + outcome.finishReason = finish.reason; if (finish.stepCount !== null) options.state.modelStepCount = stepCountBeforeAttempt + finish.stepCount; } if (chunk.type === "reasoning-delta" && typeof chunk.payload?.text === "string") { @@ -555,18 +621,23 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } flushThinkingPlan(controller); await outputText(visible.finish("")); + lastAttempt = outcome; + if (outcome.contributed || outcome.answerPhase) answerTail = outcome; } catch (error) { if (isSkillBindingAbortError(error)) { recordSkillBindingAbort(options); throw new Error(SKILL_BINDING_FAILED); } - if (isTimeoutOrAbort(error)) { + if (isTimeoutOrAbort(error) && !outcome.aborted) { + outcome.aborted = true; appendConsultationRuntimeStep(options.state, { kind: "abort", - name: drainingSpoken() ? "tool-abort" : "compose-abort", + name: attempt.drainSpoken ? "tool-abort" : "compose-abort", status: "failed", }); } + lastAttempt = outcome; + if (outcome.contributed || outcome.answerPhase) answerTail = outcome; try { await outputText(visible.finish("")); } catch {} @@ -578,7 +649,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { controller: ReadableStreamDefaultController | undefined, heading?: string, ) { - if (options.state.modelFinishReason !== "length") return; + if (lastAttempt?.finishReason !== "length" || lastAttempt.aborted) return; if (!options.continueAfterLength) throw new Error("answer_truncated"); const beforeContinue = pendingAnswer(); appendConsultationRuntimeStep(options.state, { kind: "validation", name: "answer-continue", status: "completed" }); @@ -589,9 +660,10 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { }); await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer()), { suppressCompositionActivity: true, + answerPhase: true, }); if (!/\S/.test(pendingAnswer())) throw new Error("empty_answer"); - if (options.state.modelFinishReason === "length" && pendingAnswer() === beforeContinue) { + if (lastAttempt?.finishReason === "length" && pendingAnswer() === beforeContinue) { throw new Error("answer_truncated"); } } @@ -638,7 +710,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { allowComposeRetry = true, ) { if (!options.pass4Mode) return; - await releasePass4Sentences(controller, "", true); + await releaseFinalSentence(controller); const produced = () => fullOutput.slice(origin.length); const hadRetryableReject = options.state.steps.some((step) => step.name === "pass4-reject:guarantee" @@ -651,10 +723,10 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { await consumeAttempt( controller, await options.composeAnswer(findings, PASS4_RETRY_HINT), - { suppressCompositionActivity: true }, + { suppressCompositionActivity: true, answerPhase: true }, ); await continueCurrentAnswer(controller); - await releasePass4Sentences(controller, "", true); + await releaseFinalSentence(controller); } if (!/\S/.test(produced()) && options.pass4Mode === "general_no_birth_time" && hadRetryableReject) { if (!firstOutput) { @@ -687,7 +759,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { await consumeAttempt( controller, await options.composeAnswer(findings), - { suppressCompositionActivity: true }, + { suppressCompositionActivity: true, answerPhase: true }, ); await continueCurrentAnswer(controller); await finishPass4(controller, origin, findings); @@ -703,9 +775,12 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { controller: ReadableStreamDefaultController | undefined, ) { const origin = fullOutput; + // The degraded text is the latest attempt's own words, so that attempt is + // the one whose ending decides whether this answer was finished. + if (lastAttempt) answerTail = { ...lastAttempt, contributed: true }; if (options.pass4Mode) { await releasePass4Sentences(controller, uncontractedText, false); - await releasePass4Sentences(controller, "", true); + await releaseFinalSentence(controller); } else if (/\S/.test(uncontractedText)) { if (!firstOutput) { firstOutput = true; @@ -739,6 +814,15 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { return true; } + function recordAnswerTelemetry() { + const writer = answerTail ?? lastAttempt; + if (writer) { + options.state.composeFinishReason = writer.finishReason; + options.state.composeAborted = writer.aborted; + } + options.state.answerVisibleChars = Array.from(fullOutput.trim()).length; + } + const body = new ReadableStream({ start(controller) { const sideEvent = options.sideEvent @@ -794,11 +878,23 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { appendConsultationRuntimeStep(options.state, { kind: "validation", name: "answer-retry", status: "completed" }); send(controller, { type: "activity", phase: "answer-composition", label: "正在组织回答" }); const retryOrigin = fullOutput; - await consumeAttempt(controller, await options.retryForAnswer()); + await consumeAttempt(controller, await options.retryForAnswer(), { answerPhase: true }); if (options.pass4Mode) await finishPass4(controller, retryOrigin, findings, false); } } if (!/\S/.test(fullOutput)) throw new Error("empty_answer"); + recordAnswerTelemetry(); + if (answerTail && attemptCut(answerTail)) { + // Visible text, but the stream that wrote it did not stop on its + // own: not a finished answer, so no onComplete, no charge. + appendConsultationRuntimeStep(options.state, { + kind: "validation", + name: "answer-truncated", + status: "failed", + failureCode: answerTail.aborted ? "abort" : answerTail.finishReason, + }); + throw new Error("answer_truncated"); + } settling = true; const receipt = agentExecutionReceiptSchema.parse(options.receipt()); const thinkingSections = applyThinkingSectionProgress(options.state.thinkingPlan ?? [], fullOutput); @@ -817,6 +913,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { if (settled) return; settled = true; settling = false; + recordAnswerTelemetry(); try { await options.onError?.(error, emitted, fullOutput); } catch {} diff --git a/frontend/src/mastra/consultation-tools.ts b/frontend/src/mastra/consultation-tools.ts index 4b6e0dd2..ce4f282f 100644 --- a/frontend/src/mastra/consultation-tools.ts +++ b/frontend/src/mastra/consultation-tools.ts @@ -56,6 +56,36 @@ export { AGENT_MAX_OUTPUT_TOKENS as CONSULTATION_MAX_OUTPUT_TOKENS } from "../li export const AGENT_MAX_STEPS = 8; export const AGENT_TIMEOUT_MS = 110_000; export const AGENT_SLICE_MAX_STEPS = 1; +/** + * The answer-writing phase (compose, its length continuation, the Pass 4 + * compose retry and the empty-answer retry) runs on its own clock, started + * when that phase starts. It used to share AGENT_TIMEOUT_MS with the tool loop, + * so a long loop left compose a few seconds and Mastra's abort cut the answer + * mid-sentence (BUG-1051; BUG-944 had already asked for one budget per stream). + * + * 70s is measured, not guessed: staging writer calls on the default model + * produced 560-1240 output tokens in 5.7-13.6s end to end (about 90-100 tok/s + * including first-token latency, PROGRESS-report-writer-failure-20260902). 70s + * therefore holds about 6,300-7,000 visible tokens, most of the 8,192 answer + * budget BUG-305 reserved and three to four times a typical four-heading + * answer; at half that throughput it still holds about 3,000 tokens. With the + * 110s loop the worst case is 180s, the three minutes the product accepted. + * The route's maxDuration must stay above the sum. + */ +export const CONSULTATION_COMPOSE_TIMEOUT_MS = 70_000; + +/** + * The answer phase's clock. It starts on the first call, when the first + * answer-writing stream is about to open, and every later answer stream in the + * same run (continuation, retries) shares it, so the phase as a whole is bounded. + */ +export function createConsultationAnswerClock(timeoutMs = CONSULTATION_COMPOSE_TIMEOUT_MS) { + let signal: AbortSignal | null = null; + return () => { + signal ??= AbortSignal.timeout(timeoutMs); + return signal; + }; +} export const CONSULTATION_NATAL_CALC_TOOL_ID = "run-jyotish-consultation"; export const CONSULTATION_WINDOW_CALC_TOOL_ID = "run-jyotish-window-consultation"; @@ -182,6 +212,13 @@ export type ConsultationRuntimeState = { // strict and would reject them, so they never enter it. modelStepCount: number; modelFinishReason?: AgentModelFinishReason; + // How the stream that wrote the answer ended (compose, continuation or answer + // retry; the main loop when a path has no compose; the last stream when no + // answer was written). "missing" means the stream closed without a finish + // chunk. Observability only, never in the public receipt (BUG-1051). + composeFinishReason?: AgentModelFinishReason | "missing"; + composeAborted?: boolean; + answerVisibleChars?: number; }; export function createConsultationRuntimeState(options: { plannedSteps?: number; reservedValidationSteps?: number } = {}): ConsultationRuntimeState { @@ -223,6 +260,9 @@ export function consultationModelStepTelemetry(state: ConsultationRuntimeState) skillReferenceReads: state.skillReferenceReadCount, methodologySections: state.methodologySectionCount, ...(state.modelFinishReason === undefined ? {} : { modelFinishReason: state.modelFinishReason }), + ...(state.composeFinishReason === undefined ? {} : { composeFinishReason: state.composeFinishReason }), + ...(state.composeAborted === undefined ? {} : { composeAborted: state.composeAborted }), + ...(state.answerVisibleChars === undefined ? {} : { answerVisibleChars: state.answerVisibleChars }), }; } diff --git a/frontend/tests/consult-answer-truncation-20260926.test.ts b/frontend/tests/consult-answer-truncation-20260926.test.ts new file mode 100644 index 00000000..39cf965e --- /dev/null +++ b/frontend/tests/consult-answer-truncation-20260926.test.ts @@ -0,0 +1,360 @@ +// BUG-1051: a consultation answer cut mid-sentence was completed and charged. +// +// Every stream here comes from a real Mastra `Agent` over a fake language model, +// because the shape is the bug: Mastra 1.50 does not throw when the abort signal +// fires mid-answer. It enqueues `{ type: "abort" }`, then `finish` with reason +// `tripwire`, and closes normally. The BUG-305 fixture hand-threw a +// DOMException instead, so its test stayed green while production charged. +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; +import { Agent } from "@mastra/core/agent"; + +import { agentObservabilityEventSchema } from "../src/lib/agent-observability.ts"; +import { createNdjsonParser } from "../src/lib/consultation-agent-events.ts"; +import { natalConsultationThinkingPlan, REPORT_HEADING } from "../src/lib/consultation-thinking-plan.ts"; +import { streamAgentResponse } from "../src/lib/stream-agent-response.ts"; +import { + AGENT_TIMEOUT_MS, + CONSULTATION_COMPOSE_TIMEOUT_MS, + consultationModelStepTelemetry, + consultationStepBudgetReceipt, + createConsultationAnswerClock, + createConsultationRuntimeState, + publicConsultationRuntimeSteps, +} from "../src/mastra/consultation-tools.ts"; + +type RuntimeState = ReturnType; +type Event = { type: string; code?: string; text?: string; message?: string; receipt?: unknown }; + +// Fictional answer text in the incident's shape: an opening, one heading, one +// closed sentence, then a fragment the stream was cut after. +const CLOSED = `开场段落。\n- 见面比线上聊管用\n\n## ${REPORT_HEADING.question}\n能遇到,但「遇到」和「成」这半年不是一回事。\n\n`; +const DANGLING = "你这段盘"; +const REST = "里金星落在七宫,所以关系会先从熟人圈里冒出来。\n"; + +function toolReadyState() { + const state = createConsultationRuntimeState(); + state.jyotishSkillBound = true; + state.consultationToolCallCount = 1; + state.consultationToolSuccessCount = 1; + state.consultationToolCompleted = true; + state.workflowReceipt = { route: "marriage", status: "ready", preciseTiming: "blocked", missingLayers: [] }; + state.thinkingPlan = natalConsultationThinkingPlan({ domains: ["marriage"] }); + return state; +} + +function receipt(state: RuntimeState) { + return { + runId: "run", + runtime: "mastra-agentic" as const, + skill: { name: "jyotish-vedic-astrology" as const, loaded: true, referenceReads: 0, methodologySections: 0 }, + steps: publicConsultationRuntimeSteps(state), + stepBudget: consultationStepBudgetReceipt(state), + workflow: state.workflowReceipt!, + techniqueTruth: "unknown", + }; +} + +/** + * A fake LanguageModelV2 that streams `parts` with `delayMs` between them and + * ends with `finishReason`. It honours the abort signal the way a provider + * fetch does: the pending read rejects with the signal's reason. `onPart` + * lets a test fire the timer at an exact point instead of racing wall time. + */ +function fakeModel( + parts: readonly string[], + finishReason: string | null, + delayMs = 0, + onPart?: (index: number) => void, +) { + return { + specificationVersion: "v2", + provider: "fake", + modelId: "fake-writer", + supportedUrls: {}, + async doGenerate() { + throw new Error("not used"); + }, + async doStream(options: { abortSignal?: AbortSignal }) { + const signal = options.abortSignal; + const stream = new ReadableStream({ + async start(controller) { + controller.enqueue({ type: "stream-start", warnings: [] }); + controller.enqueue({ type: "text-start", id: "1" }); + for (const [index, part] of parts.entries()) { + if (delayMs > 0) { + try { + await new Promise((resolve, reject) => { + if (signal?.aborted) return reject(signal.reason); + const timer = setTimeout(resolve, delayMs); + signal?.addEventListener("abort", () => { + clearTimeout(timer); + reject(signal.reason); + }, { once: true }); + }); + } catch (error) { + controller.error(error); + return; + } + } + controller.enqueue({ type: "text-delta", id: "1", delta: part }); + onPart?.(index); + } + controller.enqueue({ type: "text-end", id: "1" }); + if (finishReason !== null) { + controller.enqueue({ + type: "finish", + finishReason, + usage: { inputTokens: 10, outputTokens: parts.length, totalTokens: 10 + parts.length }, + }); + } + controller.close(); + }, + }); + return { stream }; + }, + }; +} + +let agentSeq = 0; +async function realStream( + parts: readonly string[], + finishReason: string | null, + options: { delayMs?: number; abortSignal?: AbortSignal; timeoutAfterParts?: number } = {}, +) { + agentSeq += 1; + // `timeoutAfterParts` stands in for AbortSignal.timeout firing right after + // that many parts went out: same TimeoutError reason, no wall-clock race. + const timer = options.timeoutAfterParts === undefined ? null : new AbortController(); + const onPart = timer + ? (index: number) => { + if (index + 1 === options.timeoutAfterParts) { + setTimeout(() => timer.abort(new DOMException("The operation was aborted due to timeout", "TimeoutError")), 5); + } + } + : undefined; + const abortSignal = timer?.signal ?? options.abortSignal; + const agent = new Agent({ + id: `writer-${agentSeq}`, + name: `writer-${agentSeq}`, + model: fakeModel(parts, finishReason, options.delayMs ?? 0, onPart) as never, + instructions: "fictional writer", + } as never); + const result = await agent.stream([{ role: "user", content: "fictional question" }] as never, { + maxSteps: 1, + toolChoice: "none", + ...(abortSignal ? { abortSignal } : {}), + } as never); + return result.fullStream as ReadableStream; +} + +/** The drained tool loop, as the natal path sees it before compose. */ +async function* toolLoop() { + yield { type: "tool-result", payload: { toolCallId: "t1", toolName: "run-jyotish-consultation", result: {} } }; + yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}, {}, {}, {}, {}, {}, {}] } } }; +} + +async function runNatal(input: { + compose: () => Promise>; + continueAfterLength?: () => Promise>; +}) { + const state = toolReadyState(); + let completed: string | null = null; + let errored: { error: unknown; emitted: boolean } | null = null; + let continues = 0; + const response = streamAgentResponse({ + runId: "run", + requestId: "req", + state, + stream: toolLoop(), + requireTool: true, + pass4Mode: "verified_chart", + toolStatus: () => "ready", + receipt: () => receipt(state) as never, + interpretFindings: async () => (state.thinkingPlan ?? []).map((section) => ({ id: section.id })), + composeAnswer: async () => input.compose(), + continueAfterLength: async () => { + continues += 1; + if (!input.continueAfterLength) throw new Error("continuation not expected"); + return input.continueAfterLength(); + }, + onComplete: (output) => { completed = output; }, + onError: (error, emitted) => { errored = { error, emitted }; }, + }); + const events: Event[] = []; + const parser = createNdjsonParser((event) => events.push(event as Event)); + parser.finish(await response.text()); + const answer = events.filter((event) => event.type === "answer.delta").map((event) => event.text ?? "").join(""); + const terminal = events.filter((event) => event.type === "run.completed" || event.type === "run.failed"); + return { + state, events, answer, terminal, continues, + completed: completed as string | null, + errored: errored as { error: unknown; emitted: boolean } | null, + }; +} + +function pieces(text: string, size = 8) { + return text.match(new RegExp(`[\\s\\S]{1,${size}}`, "g")) ?? []; +} + +test("a shared timeout firing mid-answer ends as answer_truncated, not a charged completion", async () => { + // The incident shape: the signal was already mostly spent by the tool loop, + // so it fires while compose is still writing. Cut right after the fragment. + const before = [...pieces(CLOSED), DANGLING]; + const run = await runNatal({ + compose: () => realStream([...before, REST], "stop", { delayMs: 30, timeoutAfterParts: before.length }), + }); + + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.terminal[0]?.message, "回答未完成,已保留现有内容;本次不会扣点。"); + assert.equal(run.completed, null, "onComplete (charge + persist as completed) must not run"); + assert.ok(run.errored, "onError runs, which the route maps to the cancel settlement"); + assert.equal(run.errored?.emitted, true); + // The streamed text stays; the half sentence the cut left behind is not + // released as if it were the answer's last line. + assert.match(run.answer, /不是一回事。/); + assert.equal(run.answer.includes(DANGLING), false); + assert.equal(run.answer.includes(REST.trim()), false); + // The receipt no longer says every step completed. + const failure = run.terminal[0]?.receipt as { steps: Array<{ kind: string; name: string; status: string }> }; + assert.ok(failure.steps.some((step) => step.kind === "abort" && step.name === "compose-abort" && step.status === "failed")); + assert.ok(failure.steps.some((step) => step.name === "answer-truncated" && step.status === "failed")); + assert.equal(run.state.modelFinishReason, "tripwire"); +}); + +test("the answer clock is its own: the tool loop's timer expiring does not cut compose", async () => { + // The tool loop's signal is already spent when compose starts. + const loopSignal = AbortSignal.timeout(5); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(loopSignal.aborted, true); + + const answerClock = createConsultationAnswerClock(2_000); + const run = await runNatal({ + compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop", { delayMs: 5, abortSignal: answerClock() }), + }); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.completed, `${CLOSED}${DANGLING}${REST}`); + + // Control: the pre-fix wiring handed compose the loop's spent signal. + const shared = await runNatal({ + compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop", { delayMs: 5, abortSignal: loopSignal }), + }); + assert.equal(shared.completed, null); + assert.notEqual(shared.terminal[0]?.type, "run.completed"); +}); + +test("the answer clock starts on first use and is shared by the answer phase", async () => { + const clock = createConsultationAnswerClock(30); + await new Promise((resolve) => setTimeout(resolve, 60)); + const first = clock(); + assert.equal(first.aborted, false, "the clock must not run before the answer phase starts"); + assert.equal(clock(), first, "continuation and retries share one answer-phase signal"); + await new Promise((resolve) => setTimeout(resolve, 60)); + assert.equal(first.aborted, true); +}); + +for (const reason of ["content-filter", "tool-calls", "other", "unknown", "error"] as const) { + test(`compose that ends on finish reason ${reason} with visible text is truncated, not charged`, async () => { + const run = await runNatal({ compose: () => realStream([...pieces(CLOSED), DANGLING], reason) }); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.completed, null); + assert.equal(run.continues, 0, "only length is continued"); + assert.match(run.answer, /不是一回事。/); + assert.equal(run.answer.includes(DANGLING), false); + }); +} + +test("a provider stream that closes without a finish part is truncated (Mastra reports it as unknown)", async () => { + // Mastra 1.50 still emits its own `finish` when the provider sends none; the + // reason is undefined, which the closed vocabulary reads as `unknown`. + const run = await runNatal({ compose: () => realStream([...pieces(CLOSED), DANGLING], null) }); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.completed, null); + assert.equal(run.state.composeFinishReason, "unknown"); +}); + +test("a compose stream with no finish chunk at all is truncated", async () => { + // Not a Mastra shape we have observed; guards a transport that closes early. + async function* noFinish() { + for (const piece of pieces(CLOSED)) yield { type: "text-delta", payload: { text: piece } }; + yield { type: "text-delta", payload: { text: DANGLING } }; + } + const run = await runNatal({ compose: async () => noFinish() as never }); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.completed, null); + assert.equal(run.state.composeFinishReason, "missing"); +}); + +test("length still continues, and a continuation that stops completes and charges", async () => { + const run = await runNatal({ + compose: () => realStream([...pieces(CLOSED), DANGLING], "length"), + continueAfterLength: () => realStream(pieces(REST), "stop"), + }); + assert.equal(run.continues, 1); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.completed, `${CLOSED}${DANGLING}${REST}`); + assert.ok(run.state.steps.some((step) => step.name === "answer-continue")); +}); + +test("a continuation cut by the answer clock is truncated too", async () => { + const restPieces = pieces(REST, 4); + const run = await runNatal({ + compose: () => realStream([...pieces(CLOSED), DANGLING], "length"), + continueAfterLength: () => realStream(restPieces, "stop", { delayMs: 30, timeoutAfterParts: 2 }), + }); + assert.equal(run.continues, 1); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.completed, null); + assert.ok(run.state.steps.some((step) => step.kind === "abort")); +}); + +test("a normal stop still completes, charges once and keeps the whole answer", async () => { + const run = await runNatal({ compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop") }); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.completed, `${CLOSED}${DANGLING}${REST}`); + assert.equal(run.answer, `${CLOSED}${DANGLING}${REST}`); + assert.equal(run.errored, null); + assert.equal(run.state.steps.some((step) => step.kind === "abort"), false); + assert.equal(run.state.composeFinishReason, "stop"); + assert.equal(run.state.composeAborted, false); +}); + +test("the observability record carries the compose ending and answer length, never text", async () => { + const before = [...pieces(CLOSED), DANGLING]; + const run = await runNatal({ + compose: () => realStream([...before, REST], "stop", { delayMs: 30, timeoutAfterParts: before.length }), + }); + const telemetry = consultationModelStepTelemetry(run.state); + assert.equal(telemetry.composeFinishReason, "tripwire"); + assert.equal(telemetry.composeAborted, true); + assert.equal(telemetry.answerVisibleChars, Array.from(run.answer.trim()).length); + const parsed = agentObservabilityEventSchema.parse({ requestId: "req", ...telemetry, errorCode: "answer_truncated" }); + assert.doesNotMatch(JSON.stringify(parsed), /不是一回事|你这段盘/); + // BUG-305 rule: the public receipt never carries the model's finish reason. + assert.doesNotMatch(JSON.stringify(run.terminal[0]?.receipt), /tripwire|composeFinishReason|modelFinishReason|answerVisibleChars/); +}); + +test("the consult route gives the answer phase its own clock inside maxDuration", () => { + const route = readFileSync(new URL("../src/app/api/consult/route.ts", import.meta.url), "utf8"); + const tools = readFileSync(new URL("../src/mastra/consultation-tools.ts", import.meta.url), "utf8"); + assert.match(tools, /export const CONSULTATION_COMPOSE_TIMEOUT_MS = 70_000;/); + assert.match(route, /const answerPhaseSignal = createConsultationAnswerClock\(CONSULTATION_COMPOSE_TIMEOUT_MS\);/); + const composeBlock = route.slice( + route.indexOf("const composeAnswer = async ("), + route.indexOf("const executionReceipt = (): AgentExecutionReceipt => ({", route.indexOf("const composeAnswer = async (")), + ); + assert.match(composeBlock, /\.\.\.streamOptions,\n\s+abortSignal: answerPhaseSignal\(\),/); + // Every continuation and every empty-answer retry writes on the answer clock. + const continuations = route.match(/const continueAfterLength = async \(output: string\) => \{[\s\S]*?\n\s+\};/g) ?? []; + assert.equal(continuations.length, 3); + for (const block of continuations) assert.match(block, /abortSignal: answerPhaseSignal\(\)/); + const retries = route.match(/const retryForAnswer = async \(\) => \{[\s\S]*?\n\s+\};/g) ?? []; + assert.equal(retries.length, 3); + for (const block of retries) assert.match(block, /abortSignal: answerPhaseSignal\(\)/); + // The tool loop keeps its own 110s signal. + assert.match(route, /const agentAbortSignal = AbortSignal\.timeout\(AGENT_TIMEOUT_MS\)/); + const maxDuration = Number(route.match(/export const maxDuration = (\d+);/)?.[1]); + assert.ok(AGENT_TIMEOUT_MS + CONSULTATION_COMPOSE_TIMEOUT_MS < maxDuration * 1000); + assert.ok(AGENT_TIMEOUT_MS + CONSULTATION_COMPOSE_TIMEOUT_MS <= 180_000, "product accepted about three minutes"); +}); diff --git a/frontend/tests/consultation-agentic-runtime.test.ts b/frontend/tests/consultation-agentic-runtime.test.ts index eb7da82f..20b59900 100644 --- a/frontend/tests/consultation-agentic-runtime.test.ts +++ b/frontend/tests/consultation-agentic-runtime.test.ts @@ -66,6 +66,12 @@ const serverChart = { // test can only send what the model can send. The cast keeps the argument // checked against that shape without depending on Mastra's inferred type. type ModelConsultationToolInput = { question: string; domains?: string[] }; +// BUG-1051 fixture note (applies to every `yield STOP_FINISH` below) +// 原值: 这些生成器只 yield text-delta 就结束,没有 finish chunk +// 新值: 末尾补 Mastra 1.50 正常结束时一定会发的 finish(stop) +// 原因: 无 finish 的流现在按夹断处理(answer_truncated,不扣点);Mastra 真实流 +// 正常结束必有 finish,fixture 按 §7.4 用真实形状,断言本身一条未改 +const STOP_FINISH = { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; const modelInput = (input: ModelConsultationToolInput) => input as never; const rejectedByInputSchema = (input: { question: string; theme?: string; domains?: string[] }) => input as never; @@ -1124,6 +1130,7 @@ test("text written before the contract completes is dropped, not released later" state.workflowReceipt = { route: "career", status: "ready", preciseTiming: "blocked", missingLayers: [] }; yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; yield { type: "text-delta", payload: { text: "这是真正的回答。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1231,6 +1238,7 @@ test("a calculation that succeeds only after failed attempts still satisfies the state.workflowReceipt = { route: "career", status: "ready", preciseTiming: "blocked", missingLayers: [] }; yield { type: "tool-result", payload: { toolCallId: "tool-3", toolName: "run-jyotish-consultation", result: {} } }; yield { type: "text-delta", payload: { text: "事业方向的判断如下。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1281,6 +1289,7 @@ test("incomplete runtime contract with body is delivered degraded instead of dis let failed = 0; async function* chunks() { yield { type: "text-delta", payload: { text: "不能保存" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1342,6 +1351,7 @@ test("degraded delivery drops guarantee sentences through Pass 4 (BUG-959)", asy async function* chunks() { yield { type: "text-delta", payload: { text: "我保证你一定会升职。" } }; yield { type: "text-delta", payload: { text: "方向上可以推进。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1405,6 +1415,7 @@ test("degraded delivery uses the refusal when Pass 4 drops every general-mode se const state = createConsultationRuntimeState(); async function* chunks() { yield { type: "text-delta", payload: { text: "你的上升是巨蟹座。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1428,9 +1439,11 @@ test("degraded delivery keeps only the last attempt body (BUG-960)", async () => const state = createConsultationRuntimeState(); async function* first() { yield { type: "text-delta", payload: { text: "第一段。" } }; + yield STOP_FINISH; } async function* second() { yield { type: "text-delta", payload: { text: "第二段。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: first(), requireTool: true, @@ -1456,6 +1469,7 @@ test("degraded delivery does not start a compose pass (BUG-961)", async () => { let composeCalls = 0; async function* chunks() { yield { type: "text-delta", payload: { text: "方向上可以推进。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -1575,6 +1589,7 @@ test("window precompute greens the contract without a model tool call (BUG-957)" const { ctx, state } = makeWindowCtx(); async function* chunks() { yield { type: "text-delta", payload: { text: "方向上可以推进。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, requireTool: true, @@ -1641,6 +1656,7 @@ test("window precompute failure still degrades when the model writes without a t }); async function* chunks() { yield { type: "text-delta", payload: { text: "方向上可以推进。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, requireTool: true, @@ -1861,6 +1877,7 @@ test("a calculation the model never wrote up is asked again instead of apologise } async function* answerChunks() { yield { type: "text-delta", payload: { text: "事业方向的判断如下。" } }; + yield STOP_FINISH; } const response = streamAgentResponse({ runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, @@ -2369,6 +2386,11 @@ test("natal tool success stores a Chinese thinking plan", async () => { }); test("a timeout after partial visible text is the same truncation, not a successful answer", async () => { + // 原值: 本条手工 throw DOMException("TimeoutError") 代表「超时掐断半截」,是 BUG-305 唯一的超时回归 + // 新值: 保留,只覆盖「真的抛出 TimeoutError」的 catch 分支,并补断言 abort 运行步; + // Mastra 1.50 真实超时形状(abort chunk + finish(tripwire),流正常关闭、不抛错) + // 由 consult-answer-truncation-20260926.test.ts 用真实 Agent + 假模型覆盖 + // 原因: BUG-1051 手造形状不是 Mastra 的真实行为,本条一直绿而线上半截回答照样扣点(§7.4) const state = toolOnlyRunState(); const pinchedHeading = "**先看命盘结构(Lahiri岁差、均交点口径"; let completed = 0; @@ -2395,6 +2417,7 @@ test("a timeout after partial visible text is the same truncation, not a success assert.equal(answer, pinchedHeading); const failure = events.find((event) => (event as { type?: string }).type === "run.failed") as { code: string }; assert.equal(failure.code, "answer_truncated"); + assert.ok(state.steps.some((step) => step.kind === "abort" && step.status === "failed")); }); test("consult generation reserves spoken-answer tokens and enables a separate thinking channel", () => { @@ -2420,6 +2443,7 @@ test("provider reasoning stays off the spoken answer and off the public think ch yield { type: "reasoning-delta", payload: { text: "The proposedKind value was rejected" } }; yield { type: "reasoning-delta", payload: { text: "先看事业宫的结构。" } }; yield { type: "text-delta", payload: { text: "事业方向的判断如下。" } }; + yield STOP_FINISH; } let completedThinking: string | undefined; const response = streamAgentResponse({