Files
Jyotisha/frontend/src/lib/stream-agent-response.ts
T

1133 lines
46 KiB
TypeScript

import {
appendConsultationRuntimeStep,
CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
CONSULTATION_NATAL_CALC_TOOL_ID,
CONSULTATION_WINDOW_CALC_TOOL_ID,
type ConsultationRuntimeState,
} from "../mastra/consultation-tools.ts";
import {
agentExecutionReceiptSchema,
consultationAgentPublicEventSchema,
publicActivityPhaseSchema,
type AgentExecutionReceipt,
type ConsultationAgentPublicEvent,
} from "./consultation-agent-events.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";
import {
classifyPass4,
takeClosedSentences,
GENERAL_NO_BIRTH_TIME_REFUSAL,
PASS4_RETRY_HINT,
type Pass4Mode,
type Pass4Step,
} from "./timing-output-guard.ts";
import {
applyThinkingSectionProgress,
generalConsultationThinkingPlan,
thinkPlanSteps,
type PublicThinkingSection,
} from "./consultation-thinking-plan.ts";
type Chunk = { type?: string; payload?: Record<string, unknown>; data?: unknown; error?: unknown };
type ChunkStream = AsyncIterable<unknown> | ReadableStream<unknown>;
async function* readChunks(stream: ChunkStream): AsyncIterable<Chunk> {
const values = Symbol.asyncIterator in stream
? stream as AsyncIterable<unknown>
: (async function* () {
const reader = (stream as ReadableStream<unknown>).getReader();
try {
while (true) {
const { done, value } = await reader.read();
if (done) return;
yield value;
}
} finally {
reader.releaseLock();
}
})();
for await (const value of values) {
if (value && typeof value === "object") yield value as Chunk;
}
}
type Status = "ready" | "degraded" | "blocked";
type EventOptions = {
runId: string;
requestId: string;
toolStatus: () => Status;
receipt: () => AgentExecutionReceipt;
state?: ConsultationRuntimeState;
};
export function consultationPublicActivityEvent(
value: unknown,
): Extract<ConsultationAgentPublicEvent, { type: "activity" }> | null {
if (!value || typeof value !== "object") return null;
const data = value as { phase?: unknown; label?: unknown };
const phase = publicActivityPhaseSchema.safeParse(data.phase);
if (!phase.success || typeof data.label !== "string") return null;
return { type: "activity", phase: phase.data, label: data.label.slice(0, 120) };
}
/**
* The runtime only reveals how the model loop ended through the stream: one
* `step-finish` per model step, then a terminal `finish` carrying the reason
* the model stopped and the authoritative step list. Without this, a run that
* exhausted its step budget is indistinguishable from one that chose to stop,
* because progressive-disclosure reads never reach the public event stream.
*/
function finishTelemetry(chunk: Chunk) {
const payload = chunk.payload as {
stepResult?: { reason?: unknown };
output?: { steps?: unknown };
} | undefined;
const steps = payload?.output?.steps;
return {
reason: toAgentModelFinishReason(payload?.stepResult?.reason),
stepCount: Array.isArray(steps) ? steps.length : null,
};
}
function safeToolError(error: unknown) {
if (error instanceof DOMException && error.name === "AbortError") return "cancelled" as const;
if (error instanceof DOMException && error.name === "TimeoutError") return "timeout" as const;
return "calculation_failed" as const;
}
function isTimeoutOrAbort(error: unknown) {
return error instanceof Error && (error.name === "TimeoutError" || error.name === "AbortError");
}
type RunFailedCode = "runtime_contract_incomplete" | "empty_answer" | "answer_truncated" | "calculation_failed";
type ProviderStreamErrorCode = "thinking_tool_choice_unsupported" | "provider_error";
const SKILL_BINDING_FAILED = "skill_binding_failed";
const SKILL_BINDING_ABORT_STEP = "skill-binding-abort";
const RUNTIME_CONTRACT_INCOMPLETE_STEP = "runtime-contract-incomplete";
const CONTRACT_DEGRADED_STEP = "contract-degraded";
/** Server-owned; never generated by the model. Voice: 直接,不说法务腔. */
export const CONTRACT_DEGRADED_NOTE = "\n\n这次没跑完星盘计算,上面是模型直接写的,先看着。";
function providerErrorMessage(error: unknown): string {
return error instanceof Error ? error.message : String(error);
}
function classifyProviderStreamError(error: unknown): ProviderStreamErrorCode {
return providerErrorMessage(error).includes("Thinking mode does not support this tool_choice")
? "thinking_tool_choice_unsupported"
: "provider_error";
}
function chunkProviderError(chunk: Chunk): unknown {
if (chunk.error !== undefined) return chunk.error;
return chunk.payload?.error;
}
function runFailedCode(error: unknown, emitted: boolean): RunFailedCode {
if (error instanceof Error && error.message === "runtime_contract_incomplete") return "runtime_contract_incomplete";
if (error instanceof Error && error.message === SKILL_BINDING_FAILED) return "runtime_contract_incomplete";
if (error instanceof Error && error.message === "empty_answer") {
return emitted ? "answer_truncated" : "empty_answer";
}
if (error instanceof Error && error.message === "answer_truncated") return "answer_truncated";
if (
error instanceof Error
&& (error.message === "thinking_tool_choice_unsupported" || error.message === "provider_error")
) {
return "calculation_failed";
}
if (emitted && isTimeoutOrAbort(error)) return "answer_truncated";
return "calculation_failed";
}
function runFailedMessage(code: RunFailedCode) {
if (code === "runtime_contract_incomplete") return "Agent 未完成必要的方法与计算步骤,本次不会扣点。";
if (code === "empty_answer") return "计算已完成,但这次没有生成回答,本次不会扣点。请再发送一次。";
if (code === "answer_truncated") return "回答未完成,已保留现有内容;本次不会扣点。";
return "咨询暂时无法完成,本次不会扣点。";
}
/**
* Whether a tool result is really an input rejection Mastra resolved with.
*
* Mastra validates arguments against the tool's inputSchema before `execute`
* and reports a mismatch by *resolving* with an error envelope rather than
* throwing. Passing that through as `tool.completed` would tell the client a
* calculation finished while the tool body never ran. Matched on the envelope's
* shape, not its English message, so an upstream wording change cannot silently
* turn a rejection back into a success.
*/
function isToolInputRejection(result: unknown) {
if (!result || typeof result !== "object") return false;
const envelope = result as { error?: unknown; validationErrors?: unknown };
return envelope.error === true
&& typeof envelope.validationErrors === "object"
&& envelope.validationErrors !== null;
}
/**
* Record a tool failure the tool itself could not record.
*
* A call rejected against the tool's input schema never enters `execute`, so
* nothing in the tool runs: the run reported a `tool.failed` event to the client
* while its receipt showed no failed step and its budget counted no call. The
* stream is the only place that observes every failure, whichever side of the
* schema it came from.
*
* The tool still records the failures it can see, with the duration and cause it
* alone knows, so this only fills the gap: it appends when the stream has seen
* more tool errors than the state has failed tool steps. The error object here
* is provider-shaped and cannot be classified further, but the call not reaching
* `execute` is itself the diagnosis.
*/
function recordUnrecordedToolFailure(
state: ConsultationRuntimeState,
toolErrorsSeen: number,
durationMs: number,
) {
const recorded = state.steps.filter((step) => step.kind === "tool" && step.status === "failed").length;
if (recorded >= toolErrorsSeen) return;
appendConsultationRuntimeStep(state, {
kind: "tool",
name: "run-jyotish-consultation",
status: "failed",
durationMs,
failureCode: "tool_call_rejected",
});
}
/**
* The two events that used to report the model activating the skill. The server
* binds the method into the instructions now, so they are announced once at the
* start of a run instead of being read off a tool call that no longer happens.
* They stay in the public stream because they are what tells a waiting client
* that method is in hand, and the first of them is a run's first activity.
*/
const skillBoundEvents: readonly ConsultationAgentPublicEvent[] = [
{ type: "skill.started", name: "jyotish-vedic-astrology" },
{ type: "skill.completed", name: "jyotish-vedic-astrology" },
];
function thinkingSectionEvents(sections: readonly PublicThinkingSection[] | undefined): ConsultationAgentPublicEvent[] {
if (!sections?.length) return [];
return sections.map((section) => consultationAgentPublicEventSchema.parse({
type: "thinking.section",
...section,
}));
}
function mapChunk(
chunk: Chunk,
options: EventOptions,
startedAt: Map<string, number>,
toolErrors: { seen: number },
): ConsultationAgentPublicEvent[] {
const payload = chunk.payload ?? {};
if (chunk.type === "data-jyotish-activity") {
const event = consultationPublicActivityEvent(chunk.data);
return event ? [event] : [];
}
if (chunk.type === "tool-call") {
const toolName = payload.toolName;
const callId = typeof payload.toolCallId === "string" ? payload.toolCallId : "tool";
if (toolName === "run-jyotish-consultation") {
startedAt.set(callId, Date.now());
return [{ type: "tool.started", callId, tool: "run-jyotish-consultation", label: "正在计算个人星盘" }];
}
}
if (chunk.type === "tool-result") {
const toolName = payload.toolName;
const callId = typeof payload.toolCallId === "string" ? payload.toolCallId : "tool";
if (toolName === "run-jyotish-consultation") {
const durationMs = Math.max(0, Date.now() - (startedAt.get(callId) ?? Date.now()));
if (isToolInputRejection(payload.result)) {
toolErrors.seen += 1;
if (options.state) recordUnrecordedToolFailure(options.state, toolErrors.seen, durationMs);
return [{ type: "tool.failed", callId, tool: "run-jyotish-consultation", code: "calculation_failed" }];
}
return [{
type: "tool.completed", callId, tool: "run-jyotish-consultation", status: options.toolStatus(),
durationMs,
}];
}
}
if (chunk.type === "tool-error") {
const toolName = payload.toolName;
if (toolName === "run-jyotish-consultation") {
const callId = typeof payload.toolCallId === "string" ? payload.toolCallId : "tool";
toolErrors.seen += 1;
if (options.state) {
recordUnrecordedToolFailure(
options.state,
toolErrors.seen,
Math.max(0, Date.now() - (startedAt.get(callId) ?? Date.now())),
);
}
return [{
type: "tool.failed",
callId,
tool: "run-jyotish-consultation",
code: safeToolError(payload.error),
}];
}
}
return [];
}
export async function collectAgentPublicEvents(stream: ChunkStream | Iterable<Chunk>, options: EventOptions) {
const events: ConsultationAgentPublicEvent[] = [
{ type: "run.started", runId: options.runId, requestId: options.requestId },
...skillBoundEvents,
];
const startedAt = new Map<string, number>();
const toolErrors = { seen: 0 };
let planSent = false;
const flushPlan = () => {
if (planSent) return;
const planned = thinkingSectionEvents(options.state?.thinkingPlan);
const steps = thinkPlanSteps(options.state?.thinkingPlan ?? []);
if (!planned.length && !steps.length) return;
planSent = true;
if (steps.length) events.push({ type: "think.plan", steps });
events.push(...planned);
};
for await (const chunk of stream instanceof ReadableStream || Symbol.asyncIterator in stream ? readChunks(stream as ChunkStream) : stream) {
events.push(...mapChunk(chunk, options, startedAt, toolErrors));
flushPlan();
if (chunk.type === "reasoning-delta" && typeof chunk.payload?.text === "string") {
logTruncatedReasoning(options.requestId, chunk.payload.text);
}
if (chunk.type === "text-delta" && typeof chunk.payload?.text === "string") {
events.push({ type: "answer.delta", text: chunk.payload.text });
}
}
flushPlan();
events.push({ type: "run.completed", receipt: agentExecutionReceiptSchema.parse(options.receipt()) });
return events.map((event) => consultationAgentPublicEventSchema.parse(event));
}
type AgentStreamSource = ChunkStream | (() => ChunkStream | Promise<ChunkStream>);
async function resolveAgentStream(stream: AgentStreamSource): Promise<ChunkStream> {
return typeof stream === "function" ? await stream() : stream;
}
type StreamAgentResponseOptions = EventOptions & {
state: ConsultationRuntimeState;
stream: AgentStreamSource;
warmup?: (send: (event: ConsultationAgentPublicEvent) => void) => Promise<void>;
transformText?: (text: string) => string;
requireTool: boolean;
retry?: () => Promise<ChunkStream>;
/**
* Empty-answer fallback. It keeps the tools, so it can fetch this request's
* cached calculation again. `retryHint` carries the Pass 4 rewrite hint when
* every sentence of the first answer was rejected.
*/
retryForAnswer?: (retryHint?: string) => Promise<ChunkStream>;
/**
* Length continuation. `evidence` is the calculation result this run's
* model saw (the calculation tool's result), so the continuation writes
* with the same evidence instead of blind (BUG-1053).
*/
continueAfterLength?: (output: string, evidence?: unknown) => Promise<ChunkStream>;
/**
* The answer is the text of the loop's step(s) that end without a tool call
* (BUG-1053). Set for routes whose loop has tools: a step's text is held
* until it reads as the answer (a Markdown heading, or
* ANSWER_RELEASE_CHARS of text), and dropped when the step calls a tool, so
* narration around tool calls never reaches the answer.
*/
stepScopedAnswer?: boolean;
/**
* Called once, when the run contract first turns ready (the calculation
* result is in hand). The route hands the loop to the answer clock here.
*/
onAnswerPhase?: () => void;
pass4Mode?: Pass4Mode;
continueAfterDisconnect?: boolean;
headers?: HeadersInit;
onFirstActivity?: () => void | Promise<void>;
onFirstOutput?: () => void | Promise<void>;
/**
* Content moderation of the finished answer (2026-09-30). Returns the text
* that must replace the whole answer, or null to keep it. A replacement is
* sent as one `answer.delta` with `replace: true` and is what `onComplete`
* receives.
*/
moderateOutput?: (output: string) => Promise<string | null>;
onComplete?: (
output: string,
receipt: AgentExecutionReceipt,
thinkingText?: string,
thinkingSections?: PublicThinkingSection[],
) => void | Promise<void>;
onError?: (error: unknown, emitted: boolean, output: string) => void | Promise<void>;
onCancel?: (emitted: boolean) => void | Promise<void>;
sideEvent?: Promise<ConsultationAgentPublicEvent | null>;
};
// Failed attempts are retried by the model against the same request-scoped
// calculation cache, so only successful workflow executions may count against
// the single-calculation boundary. Gating on total attempts would make any
// transient failure permanently unrecoverable.
function contractReady(options: StreamAgentResponseOptions) {
return options.state.jyotishSkillBound
&& (!options.requireTool || (options.state.consultationToolCompleted && options.state.consultationToolSuccessCount === 1));
}
function isSkillBindingAbortError(error: unknown) {
if (!(error instanceof Error)) return false;
if (error.message === SKILL_BINDING_FAILED) return true;
return error.message.includes("not bound into the system prompt");
}
function isSkillBindingTripwire(chunk: Chunk) {
if (chunk.type !== "tripwire") return false;
const reason = typeof chunk.payload?.reason === "string" ? chunk.payload.reason : "";
const processorId = typeof chunk.payload?.processorId === "string" ? chunk.payload.processorId : "";
return processorId === "jyotish-skill-bound" || reason.includes("not bound into the system prompt");
}
function recordSkillBindingAbort(options: StreamAgentResponseOptions) {
if (options.state.steps.some((step) => step.name === SKILL_BINDING_ABORT_STEP)) return;
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: SKILL_BINDING_ABORT_STEP,
status: "failed",
failureCode: SKILL_BINDING_FAILED,
});
console.error("[consult-binding-error]", {
requestId: options.requestId,
code: SKILL_BINDING_FAILED,
});
}
/**
* 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));
}
/**
* How much of a step's text is held before it is released as the answer when
* no Markdown heading has appeared yet. Narration before a tool call is one
* short sentence ("我再读一下参考"); the natal opener is 3-6 sentences and at
* most 400 characters, so 160 characters is two or three sentences into it:
* about a second or two of delay before the first visible sentence.
*/
export const ANSWER_RELEASE_CHARS = 160;
function readsAsAnswer(text: string) {
return /(^|\n)#{1,3} \S/.test(text) || Array.from(text).length >= ANSWER_RELEASE_CHARS;
}
function isCalculationTool(toolName: unknown) {
return toolName === CONSULTATION_NATAL_CALC_TOOL_ID || toolName === CONSULTATION_WINDOW_CALC_TOOL_ID;
}
/**
* How many characters a step after an evidence lookup must repeat, from the
* start of the answer already put out in this attempt, before it counts as a
* restart and the repeat is dropped. Two answers that merely open alike
* ("你这盘…") diverge well before this.
*/
export const LOOKUP_RESTART_MATCH_CHARS = 40;
function stepFinishReason(chunk: Chunk) {
const payload = chunk.payload as { stepResult?: { reason?: unknown } } | undefined;
return payload?.stepResult?.reason;
}
export function streamAgentResponse(options: StreamAgentResponseOptions) {
const encoder = new TextEncoder();
let disconnected = false;
let settled = false;
let settling = false;
let emitted = false;
let firstActivity = false;
let firstOutput = false;
let fullOutput = "";
let uncontractedText = "";
let planSent = false;
let answerPhaseStarted = false;
// The calculation result the model saw, kept so a length continuation
// writes with the same evidence (BUG-1053). Never sent to the client.
let calculationEvidence: unknown;
// The one-shot evidence lookup's result, if the model used it; carried into
// a length continuation with the card (TASK-consult-evidence-card-20260927).
let lookupEvidence: unknown;
// Pass 4 buffers only the current open sentence. Closed sentences are
// classified and either sent whole or dropped whole. Whole-answer rewrite
// is allowed only before any answer.delta has gone out; after the first
// sentence is public, later rejects are dropped in place so the user never
// sees a flash-then-replace. That is the boundary between "don't flash
// twice" and "stream by sentence".
let pass4Buffer = "";
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: an attempt that
// wrote nothing (a contract retry, a loop that only called tools) says
// nothing about whether the answer 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;
void Promise.resolve(options.onFirstActivity?.()).catch(() => {});
}
if (!disconnected && controller) controller.enqueue(encoder.encode(`${JSON.stringify(consultationAgentPublicEventSchema.parse(event))}\n`));
};
const flushThinkingPlan = (controller: ReadableStreamDefaultController<Uint8Array> | undefined) => {
if (planSent) return;
if (!options.requireTool && !options.state.thinkingPlan?.length) {
options.state.thinkingPlan = generalConsultationThinkingPlan();
}
const plan = options.state.thinkingPlan ?? [];
const steps = thinkPlanSteps(plan);
const planned = thinkingSectionEvents(plan);
if (!steps.length && !planned.length) return;
planSent = true;
if (steps.length) send(controller, { type: "think.plan", steps });
for (const event of planned) send(controller, event);
};
const pendingAnswer = () => fullOutput + pass4Buffer;
function recordPass4Steps(steps: readonly Pass4Step[]) {
for (const step of steps) {
const name = `${step.action === "observe" ? "pass4-observe" : "pass4-reject"}:${step.kind}`;
if (options.state.steps.some((item) => item.name === name)) continue;
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name,
status: step.action === "observe" ? "completed" : "failed",
failureCode: step.kind,
});
}
}
async function releasePass4Sentences(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
text: string,
flush: boolean,
) {
if (!options.pass4Mode) return;
pass4Buffer += text;
const { closed, rest } = flush
? { closed: pass4Buffer ? [pass4Buffer] : [], rest: "" }
: takeClosedSentences(pass4Buffer);
pass4Buffer = rest;
for (const sentence of closed) {
if (!sentence) continue;
const steps = classifyPass4(sentence, options.pass4Mode);
recordPass4Steps(steps);
if (steps.some((step) => step.action === "reject")) continue;
if (!firstOutput && /\S/.test(sentence)) {
firstOutput = true;
await options.onFirstOutput?.();
}
send(controller, { type: "answer.delta", text: sentence });
fullOutput += sentence;
if (/\S/.test(sentence)) emitted = true;
}
}
/**
* 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);
}
function markAnswerPhase() {
if (answerPhaseStarted || !contractReady(options)) return;
answerPhaseStarted = true;
options.onAnswerPhase?.();
}
async function consumeAttempt(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
stream: ChunkStream,
attempt: { 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 = "";
markAnswerPhase();
const outcome: AttemptOutcome = {
finishReason: "missing",
aborted: false,
answerPhase: Boolean(attempt.answerPhase),
contributed: false,
};
// An answer-phase attempt that starts with the calculation already in hand
// can only re-fetch it from the request cache; announcing that as a new
// calculation would put a live "calculating" row back above the answer.
const refetchOnly = Boolean(attempt.answerPhase) && contractReady(options);
const visible = createVisibleTextTransformer(options.transformText ?? ((value) => value));
let held = "";
let composingSent = Boolean(attempt.suppressCompositionActivity);
// Per model step (BUG-1053): the answer is the text of the step that ends
// without a tool call. Held text is released once it reads as the answer
// or when its step ends on its own; a tool call in the step drops it.
let stepText = "";
let stepReleased = false;
let stepCalledTool = false;
// Whether the current step put answer text out, and how the last such
// step ended. Settlement judges the step that wrote the answer: Mastra's
// loop runs another model step after a step that ends on `other`,
// `unknown` or a bare `tool-calls`, and that step's `stop` must not make
// the cut answer read as finished (BUG-1051 carried into the loop).
let stepWrote = false;
let answerStepReason: AgentModelFinishReason | undefined;
// Evidence lookup during the answer (T5): if the model had already put
// answer text out in this attempt and then called the lookup, the next
// step may start the answer over. A verbatim restart of the text already
// out is dropped as it arrives, so nothing released is repeated; anything
// that diverges within LOOKUP_RESTART_MATCH_CHARS is kept as written.
let attemptText = "";
let restartCheck = false;
let restartPos = 0;
let restartHeld = "";
let restartConfirmed = false;
const resetStep = () => {
stepText = "";
stepReleased = false;
stepCalledTool = false;
};
const outputText = async (text: string) => {
// Text the model writes before the contract is ready is not the answer: it
// is the model narrating its own in-progress or failed tool calls. Holding
// it meant a later successful call released that narration as the entire
// visible answer, so a run where the model recovered read as a run where it
// explained itself instead of answering. Drop it from the live stream, but
// keep a copy so a still-red contract can degrade instead of discarding it.
if (!contractReady(options)) {
if (text) uncontractedText += text;
return;
}
uncontractedText = "";
attemptText += text;
held += text;
if (!held) return;
if (/\S/.test(held)) {
outcome.contributed = true;
stepWrote = true;
}
if (!composingSent) {
composingSent = true;
send(controller, { type: "activity", phase: "answer-composition", label: "正在组织回答" });
}
if (options.pass4Mode) {
await releasePass4Sentences(controller, held, false);
held = "";
return;
}
if (!firstOutput && /\S/.test(held)) {
firstOutput = true;
await options.onFirstOutput?.();
}
send(controller, { type: "answer.delta", text: held });
fullOutput += held;
if (/\S/.test(held)) emitted = true;
held = "";
};
const filterRestart = (text: string) => {
let index = 0;
while (index < text.length && restartCheck) {
const char = text[index]!;
if (restartPos === 0 && !restartConfirmed && /\s/.test(char)) {
restartHeld += char;
index += 1;
continue;
}
if (restartPos < attemptText.length && char === attemptText[restartPos]) {
restartPos += 1;
index += 1;
if (!restartConfirmed) {
restartHeld += char;
if (restartPos >= LOOKUP_RESTART_MATCH_CHARS) {
restartConfirmed = true;
restartHeld = "";
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: "answer-restart-dropped",
status: "completed",
});
}
}
if (restartPos >= attemptText.length) restartCheck = false;
continue;
}
restartCheck = false;
}
// Still repeating: everything so far is held (not yet a confirmed
// restart) or dropped (confirmed).
if (restartCheck) return "";
// The check ended inside this chunk. An unconfirmed match was a
// coincidence and goes out as written; a confirmed repeat stays dropped.
const kept = `${restartConfirmed ? "" : restartHeld}${text.slice(index)}`;
restartHeld = "";
return kept;
};
const acceptText = async (raw: string) => {
const text = restartCheck ? filterRestart(raw) : raw;
await acceptAnswerText(text);
};
const flushRestartHeld = async () => {
// Only a step that has started repeating ends the check here; the
// lookup's own step ends before the step that might restart begins.
if (!restartCheck || (restartPos === 0 && !restartHeld)) return;
restartCheck = false;
const held = restartConfirmed ? "" : restartHeld;
restartHeld = "";
if (held) await acceptAnswerText(held);
};
const acceptAnswerText = async (text: string) => {
if (!options.stepScopedAnswer || !contractReady(options) || stepReleased) {
await outputText(text);
return;
}
if (stepCalledTool || !text) return;
stepText += text;
if (!readsAsAnswer(stepText)) return;
stepReleased = true;
const released = stepText;
stepText = "";
await outputText(released);
};
// The step ended (or the stream did): its held text is the answer unless
// the step called a tool.
const settleStep = async (reason?: unknown) => {
await flushRestartHeld();
const pending = stepText;
const toolStep = stepCalledTool || reason === "tool-calls";
resetStep();
if (pending && !toolStep) await outputText(pending);
};
// A retry runs a second model loop under the same step budget, so the run
// total accumulates while the finish reason describes the latest attempt.
const stepCountBeforeAttempt = options.state.modelStepCount;
try {
for await (const chunk of readChunks(stream)) {
if (isSkillBindingTripwire(chunk)) {
recordSkillBindingAbort(options);
throw new Error(SKILL_BINDING_FAILED);
}
if (chunk.type === "error") {
const raw = chunkProviderError(chunk);
const code = classifyProviderStreamError(raw);
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: "model-stream-error",
status: "failed",
});
console.error("[consult-provider-error]", {
requestId: options.requestId,
code,
messageHead: providerErrorMessage(raw).slice(0, 200),
});
throw new Error(code);
}
for (const event of mapChunk(chunk, options, startedAt, toolErrors)) {
if (refetchOnly && (event.type === "tool.started" || event.type === "tool.completed")) continue;
send(controller, event);
}
if (
chunk.type === "tool-result"
&& isCalculationTool(chunk.payload?.toolName)
&& !isToolInputRejection(chunk.payload?.result)
) {
calculationEvidence = chunk.payload?.result;
}
if (
chunk.type === "tool-result"
&& chunk.payload?.toolName === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID
&& !isToolInputRejection(chunk.payload?.result)
) {
lookupEvidence = chunk.payload?.result;
}
markAnswerPhase();
flushThinkingPlan(controller);
if (chunk.type === "step-start") {
resetStep();
stepWrote = false;
}
if (chunk.type === "tool-call") {
// The visible-text transformer holds the step's open clause until a
// sentence boundary. It used to carry that clause into the next
// step's first text, so narration without closing punctuation
// ("我先排一下盘:") leaked into the answer (BUG-1059). Settle it with
// this step: answer text it continues goes out whole, anything else
// is narration and dropped with the step.
const openClause = visible.finish("");
if (openClause && (!contractReady(options) || stepReleased)) await acceptAnswerText(openClause);
// Whatever this step said before calling a tool is narration.
stepCalledTool = true;
stepText = "";
if (
chunk.payload?.toolName === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID
&& /\S/.test(attemptText)
) {
restartCheck = true;
restartPos = 0;
restartHeld = "";
restartConfirmed = false;
}
}
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;
recordAttemptAbort();
}
if (chunk.type === "step-finish") {
options.state.modelStepCount += 1;
const reason = stepFinishReason(chunk);
await settleStep(reason);
if (stepWrote) answerStepReason = toAgentModelFinishReason(reason);
stepWrote = false;
}
if (chunk.type === "finish") {
const finish = finishTelemetry(chunk);
options.state.modelFinishReason = finish.reason;
outcome.finishReason = answerStepReason ?? finish.reason;
if (finish.stepCount !== null) options.state.modelStepCount = stepCountBeforeAttempt + finish.stepCount;
}
if (chunk.type === "reasoning-delta" && typeof chunk.payload?.text === "string") {
logTruncatedReasoning(options.requestId, chunk.payload.text);
}
if (chunk.type === "text-delta" && typeof chunk.payload?.text === "string") {
await acceptText(visible.push(chunk.payload.text));
}
}
flushThinkingPlan(controller);
await acceptText(visible.finish(""));
await settleStep();
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) && !outcome.aborted) {
outcome.aborted = true;
recordAttemptAbort();
}
try {
await acceptText(visible.finish(""));
await settleStep();
} catch {}
lastAttempt = outcome;
if (outcome.contributed || outcome.answerPhase) answerTail = outcome;
throw error;
}
}
// An abort before the calculation result is a tool-phase abort; after it,
// the loop was writing the answer. `compose-abort` keeps its BUG-1051 name:
// "compose" now means the answer-writing step, not a separate stream.
function recordAttemptAbort() {
appendConsultationRuntimeStep(options.state, {
kind: "abort",
name: contractReady(options) ? "compose-abort" : "tool-abort",
status: "failed",
});
}
async function continueCurrentAnswer(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
heading?: string,
) {
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" });
send(controller, {
type: "activity",
phase: "answer-composition",
label: heading ? consultationWriteLabel(heading, true) : "正在组织回答",
});
const evidence = lookupEvidence !== undefined
&& calculationEvidence
&& typeof calculationEvidence === "object"
&& !Array.isArray(calculationEvidence)
? { ...calculationEvidence, evidence_lookup: lookupEvidence }
: calculationEvidence;
await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer(), evidence), {
suppressCompositionActivity: true,
answerPhase: true,
});
if (!/\S/.test(pendingAnswer())) throw new Error("empty_answer");
if (lastAttempt?.finishReason === "length" && pendingAnswer() === beforeContinue) {
throw new Error("answer_truncated");
}
}
/**
* Release the answer's last sentence and apply the general-mode refusal when
* Pass 4 rejected every sentence. A natal answer that Pass 4 emptied is
* rewritten by the caller through `retryForAnswer` with PASS4_RETRY_HINT.
*/
async function finishPass4(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
origin: string,
) {
if (!options.pass4Mode) return;
await releaseFinalSentence(controller);
const produced = () => fullOutput.slice(origin.length);
if (!/\S/.test(produced()) && options.pass4Mode === "general_no_birth_time" && hadRetryablePass4Reject()) {
if (!firstOutput) {
firstOutput = true;
await options.onFirstOutput?.();
}
send(controller, { type: "answer.delta", text: GENERAL_NO_BIRTH_TIME_REFUSAL });
fullOutput = origin + GENERAL_NO_BIRTH_TIME_REFUSAL;
emitted = true;
}
}
function hadRetryablePass4Reject() {
return options.state.steps.some((step) =>
step.name === "pass4-reject:guarantee"
|| step.name === "pass4-reject:personal-chart"
|| step.name === "pass4-reject:methodology"
);
}
async function deliverDegradedAnswer(
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 releaseFinalSentence(controller);
} else if (/\S/.test(uncontractedText)) {
if (!firstOutput) {
firstOutput = true;
await options.onFirstOutput?.();
}
send(controller, { type: "answer.delta", text: uncontractedText });
fullOutput += uncontractedText;
emitted = true;
}
if (!/\S/.test(fullOutput.slice(origin.length))) {
if (options.pass4Mode === "general_no_birth_time") {
if (!firstOutput) {
firstOutput = true;
await options.onFirstOutput?.();
}
send(controller, { type: "answer.delta", text: GENERAL_NO_BIRTH_TIME_REFUSAL });
fullOutput += GENERAL_NO_BIRTH_TIME_REFUSAL;
emitted = true;
} else {
return false;
}
}
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: CONTRACT_DEGRADED_STEP,
status: "failed",
});
send(controller, { type: "answer.delta", text: CONTRACT_DEGRADED_NOTE });
fullOutput += CONTRACT_DEGRADED_NOTE;
emitted = true;
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
? options.sideEvent.then((event) => {
if (event) send(controller, event);
return event;
}).catch(() => null)
: Promise.resolve(null);
void (async () => {
send(controller, { type: "run.started", runId: options.runId, requestId: options.requestId });
for (const event of skillBoundEvents) send(controller, event);
flushThinkingPlan(controller);
try {
let deliveredDegraded = false;
if (options.warmup) {
await options.warmup((event) => send(controller, event));
flushThinkingPlan(controller);
}
await consumeAttempt(controller, await resolveAgentStream(options.stream));
if (!contractReady(options) && options.retry) {
appendConsultationRuntimeStep(options.state, { kind: "validation", name: "runtime-contract-retry", status: "completed" });
send(controller, { type: "activity", phase: "loading-method", label: "正在补齐方法与计算步骤" });
await consumeAttempt(controller, await options.retry());
}
if (!contractReady(options)) {
const canDegrade = options.requireTool
&& options.state.consultationToolSuccessCount === 0
&& /\S/.test(uncontractedText);
if (canDegrade) {
deliveredDegraded = await deliverDegradedAnswer(controller);
}
if (!deliveredDegraded) {
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: RUNTIME_CONTRACT_INCOMPLETE_STEP,
status: "failed",
});
throw new Error("runtime_contract_incomplete");
}
}
if (!deliveredDegraded) {
// The loop's own final step wrote the answer with the calculation
// in its context (BUG-1053); there is no second, blind compose.
await continueCurrentAnswer(controller);
if (options.pass4Mode) await finishPass4(controller, "");
if (!/\S/.test(fullOutput) && options.retryForAnswer) {
appendConsultationRuntimeStep(options.state, { kind: "validation", name: "answer-retry", status: "completed" });
send(controller, { type: "activity", phase: "answer-composition", label: "正在组织回答" });
const retryOrigin = fullOutput;
pass4Buffer = "";
await consumeAttempt(
controller,
await options.retryForAnswer(
options.pass4Mode && hadRetryablePass4Reject() ? PASS4_RETRY_HINT : undefined,
),
{ answerPhase: true },
);
if (options.pass4Mode) await finishPass4(controller, retryOrigin);
}
}
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 replacement = options.moderateOutput ? await options.moderateOutput(fullOutput) : null;
if (replacement !== null) {
fullOutput = replacement;
send(controller, { type: "answer.delta", text: replacement, replace: true });
}
const receipt = agentExecutionReceiptSchema.parse(options.receipt());
const thinkingSections = applyThinkingSectionProgress(options.state.thinkingPlan ?? [], fullOutput);
// Public thinking text came from the removed interpret pass, which
// only ever published step ids; provider reasoning never reaches it.
await options.onComplete?.(
fullOutput,
receipt,
undefined,
thinkingSections.length > 0 ? thinkingSections : undefined,
);
settled = true;
settling = false;
await sideEvent;
send(controller, { type: "run.completed", receipt });
if (!disconnected) controller.close();
} catch (error) {
if (settled) return;
settled = true;
settling = false;
recordAnswerTelemetry();
try {
await options.onError?.(error, emitted, fullOutput);
} catch {}
const code = runFailedCode(error, emitted);
// Step durations, the step budget and the workflow route are the only
// evidence the caller has for why a run failed. Building the receipt
// must not be able to replace the failure event with a silent close.
let failureReceipt: AgentExecutionReceipt | undefined;
try {
failureReceipt = agentExecutionReceiptSchema.parse(options.receipt());
} catch {}
send(controller, {
type: "run.failed",
code,
message: runFailedMessage(code),
...(failureReceipt ? { receipt: failureReceipt } : {}),
});
if (!disconnected) controller.close();
}
})();
},
async cancel() {
if (settled) return;
if (settling || options.continueAfterDisconnect) {
disconnected = true;
return;
}
settled = true;
await options.onCancel?.(emitted);
},
});
return new Response(body, {
headers: {
"cache-control": "no-cache, no-transform",
"content-type": "application/x-ndjson; charset=utf-8",
"x-accel-buffering": "no",
"x-ayanam-mode": "mastra-agentic",
"x-ayanam-request-id": options.requestId,
...options.headers,
},
});
}