The natal loop's step after run-jyotish-consultation saw the evidence, but its text was drained and a second, blind compose stream (history + question only, empty findings) wrote the user-visible answer. Remove compose, interpret and the drain; keep the loop's own final-step text. - stepScopedAnswer: per-step holding; text of a step that calls a tool is dropped, so narration around tool calls never reaches the answer - writing shape (opener + four headings) moves into the user turn - length continuation receives the calculation result; Pass 4 whole-answer reject retries through retryForAnswer with the rewrite hint - createConsultationRunClock: tools keep the 110s tool phase; the loop is handed to the 70s answer clock when the calculation result arrives - settlement judges the step that wrote the answer (BUG-1051 kept) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8
373 lines
19 KiB
TypeScript
373 lines
19 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.
|
||
//
|
||
// 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<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>;
|
||
}
|
||
|
||
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,
|
||
// 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");
|
||
});
|