// 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. // // BUG-1053 conversion note (applies to the whole file) // 原值: runNatal 先喂一段「已排空的工具循环」,再由 composeAnswer 返回写回答的流; // 时钟用 createConsultationAnswerClock / CONSULTATION_COMPOSE_TIMEOUT_MS // 新值: 写回答的流就是本命主循环本身(stream: 真实 Mastra 流,stepScopedAnswer), // 时钟用 createConsultationRunClock / CONSULTATION_ANSWER_TIMEOUT_MS; // 测试名保持不变,其中的 "compose" 指写回答的那一步 // 原因: 产品 2026-09-27 决定删除单独的 compose 流(它看不到计算结果); // BUG-1051 的每条断言(截断、不扣点、abort 步、半句不外发、观测字段)原样保留 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_ANSWER_TIMEOUT_MS, consultationModelStepTelemetry, consultationStepBudgetReceipt, createConsultationRunClock, 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. // 原值: 夹具标题用 REPORT_HEADING.question(先回答你的问题) // 新值: 改用 REPORT_HEADING.support;这里只需要「任意一个二级标题」 // 原因: TASK-consult-answer-the-question-20260927 D8 / BUG-1073 删除了「先回答你的问题」这一节 const CLOSED = `开场段落。\n- 见面比线上聊管用\n\n## ${REPORT_HEADING.support}\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; } 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, // The calculation is already in hand (toolReadyState); the loop's step // that writes the answer is the stream itself (BUG-1053). stream: () => input.compose(), requireTool: true, stepScopedAnswer: true, pass4Mode: "verified_chart", toolStatus: () => "ready", receipt: () => receipt(state) as never, 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 () => { // 原值: 工具循环的 signal 已过期,compose 用 createConsultationAnswerClock 另起的时钟照常完成 // 新值: 同一个 run clock 里工具阶段的计时已到,但循环已交给答案时钟,写回答照常完成 // 原因: BUG-1053 删掉单独的 compose 流后,写回答发生在主循环里, // 「写回答不被工具阶段的闸刀掐断」改由循环 signal 的交接来保证 const clock = createConsultationRunClock({ toolPhaseMs: 5, answerMs: 2_000 }); clock.answerSignal(); await new Promise((resolve) => setTimeout(resolve, 20)); assert.equal(clock.toolSignal.aborted, true); const run = await runNatal({ compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop", { delayMs: 5, abortSignal: clock.loopSignal }), }); assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); assert.equal(run.completed, `${CLOSED}${DANGLING}${REST}`); // Control: a loop never handed to the answer clock ends with the tool phase. const shared = createConsultationRunClock({ toolPhaseMs: 5, answerMs: 2_000 }); await new Promise((resolve) => setTimeout(resolve, 20)); const cut = await runNatal({ compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop", { delayMs: 5, abortSignal: shared.loopSignal }), }); assert.equal(cut.completed, null); assert.notEqual(cut.terminal[0]?.type, "run.completed"); }); test("the answer clock starts on first use and is shared by the answer phase", async () => { // 原值: createConsultationAnswerClock(30) // 新值: createConsultationRunClock({ answerMs: 30 }).answerSignal // 原因: BUG-1053 把答案时钟并进 run clock;首用才起算、全阶段共用的断言不变 const clock = createConsultationRunClock({ answerMs: 30 }).answerSignal; 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", () => { // 原值: 锁 CONSULTATION_COMPOSE_TIMEOUT_MS、createConsultationAnswerClock 与 composeAnswer 块用答案时钟 // 新值: 锁 CONSULTATION_ANSWER_TIMEOUT_MS、run clock(工具用 toolSignal、循环用 loopSignal、 // 拿到计算结果即交给答案时钟),续写与回答重试仍用答案时钟 // 原因: BUG-1053 删除 compose 流;「写回答有自己的时钟、不与工具阶段共用闸刀」这一性质不变 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_ANSWER_TIMEOUT_MS = 70_000;/); assert.match(route, /const runClock = createConsultationRunClock\(\{\s+toolPhaseMs: AGENT_TIMEOUT_MS,\s+answerMs: CONSULTATION_ANSWER_TIMEOUT_MS,/); assert.match(route, /const answerPhaseSignal = runClock\.answerSignal;/); assert.equal(route.match(/onAnswerPhase: startAnswerPhase,/g)?.length, 2, "natal and window hand the loop over"); // 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 \(retryHint\?: string\) => \{[\s\S]*?\n\s+\};/g) ?? []; assert.equal(retries.length, 3); for (const block of retries) assert.match(block, /abortSignal: answerPhaseSignal\(\)/); // The tools keep the tool phase's own 110s deadline. assert.match(route, /const agentAbortSignal = runClock\.toolSignal;/); const maxDuration = Number(route.match(/export const maxDuration = (\d+);/)?.[1]); assert.ok(AGENT_TIMEOUT_MS + CONSULTATION_ANSWER_TIMEOUT_MS < maxDuration * 1000); assert.ok(AGENT_TIMEOUT_MS + CONSULTATION_ANSWER_TIMEOUT_MS <= 180_000, "product accepted about three minutes"); });