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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8
This commit is contained in:
Jesse_Chen
2026-09-26 22:17:36 +08:00
co-authored by Claude Opus 5.5
parent 16f4600d5e
commit 1530a0dd63
6 changed files with 557 additions and 17 deletions
+18 -5
View File
@@ -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),
+6
View File
@@ -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.
+109 -12
View File
@@ -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<string, number>();
// 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<Uint8Array> | 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<Uint8Array> | 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<Uint8Array> | 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<Uint8Array> | 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<Uint8Array> | 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<Uint8Array>({
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 {}
+40
View File
@@ -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 }),
};
}
@@ -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<typeof createConsultationRuntimeState>;
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<void>((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<unknown>;
}
/** 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<ReadableStream<unknown>>;
continueAfterLength?: () => Promise<ReadableStream<unknown>>;
}) {
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");
});
@@ -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({