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
361 lines
17 KiB
TypeScript
361 lines
17 KiB
TypeScript
// 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");
|
|
});
|