diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md index 34f1997f..706bcf4a 100644 --- a/frontend/DESIGN.md +++ b/frontend/DESIGN.md @@ -2,6 +2,10 @@ This file adapts the full visual analysis in `CLAUDE_DESIGN.md` to the shipped Jyotisha application. `CLAUDE_DESIGN.md` remains the upstream reference; this file is the implementation contract. +## 普通咨询一遍成文(2026-09-27,BUG-1053) + +本命咨询不再有单独的「写结论」阶段:算完盘后,同一个模型循环接着就写回答。时间线的变化只有两处,没有新组件、没有新动效。一是不再出现「正在整理判断依据」「正在写结论」这两个阶段事件,也不再发 `think.step`;计划里的思考行在第一段正文出现时收口(沿用 `completeLiveThink`),写作行的进行中句统一为「正在组织回答」。二是第一段正文要等模型写出第一个二级标题或满 160 字才出现(为了把工具前后的过程说明挡在正文外),开场句会整段一起到达,之后照常按帧放出。截断、停止、失败三种收口与提示不变(BUG-1051)。 + ## 校正开场与步骤回执(2026-09-26,BUG-1049 / BUG-1050) 开场一轮只有两段字:正文两句讲做法(要把出生时间缩小到更准的范围、现在先在哪段时间里找;你说几件大事和大概年月,我拿去和星盘对照),下面是问题块里的题干(六个例子 + 一个带年月的回答示例)。题干只出现一次:正文里和题干逐句相同或几乎相同的句子才会被去掉;只是开头几个字一样、内容不同的句子保留(以前按前 12 个字判定,把正文里的例子句删掉了)。 diff --git a/frontend/src/app/api/consult/route.ts b/frontend/src/app/api/consult/route.ts index 9bbb1037..be2c4739 100644 --- a/frontend/src/app/api/consult/route.ts +++ b/frontend/src/app/api/consult/route.ts @@ -42,13 +42,18 @@ import { classifyConsultationTurn, type SmalltalkUsage } from "@/lib/consultatio import { streamSmalltalkResponse } from "@/lib/stream-smalltalk-response"; import { consultationPublicActivityEvent, streamAgentResponse } from "@/lib/stream-agent-response"; import type { AgentExecutionReceipt, WorkflowReceipt } from "@/lib/consultation-agent-events"; -import { consultationComposePrompt, consultationContinuePrompt, natalConsultationThinkingPlan, type PublicThinkingSection } from "@/lib/consultation-thinking-plan"; +import { + consultationContinueMessages, + dailyAnswerShapeInstruction, + natalAnswerShapeInstruction, + natalConsultationThinkingPlan, + type PublicThinkingSection, +} from "@/lib/consultation-thinking-plan"; import { AGENT_MAX_STEPS, AGENT_TIMEOUT_MS, - AGENT_SLICE_MAX_STEPS, - CONSULTATION_COMPOSE_TIMEOUT_MS, - createConsultationAnswerClock, + CONSULTATION_ANSWER_TIMEOUT_MS, + createConsultationRunClock, consultationContinueGenerationSettings, consultationGenerationSettings, consultationNatalPrepareStep, @@ -108,9 +113,10 @@ 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: 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. +// stay under: the tool phase (AGENT_TIMEOUT_MS) plus the answer phase's own +// clock (CONSULTATION_ANSWER_TIMEOUT_MS), 180s, plus setup and settlement +// (BUG-1051, BUG-1053). Self-hosted `node server.js` does not enforce it; it +// documents the ceiling. const chatRequestMetadataSchema = z.object({ requestId: z.string().uuid(), @@ -969,9 +975,11 @@ export async function POST(request: Request) { const natalToolInstruction = pinsConsultationDomains(consultEntrypoint) ? "如需新的个人星盘结论,必须调用服务器绑定的排盘工具。调用时不要填写 domains,沿用服务器已选定的主题。" : "如需新的个人星盘结论,必须调用服务器绑定的排盘工具。"; + // The loop's own step after the calculation writes the answer (BUG-1053), + // so the answer shape the removed compose prompt used to carry is here. const natalInstruction = consultEntrypoint === "daily_starlanguage" - ? `${natalToolInstruction}按三节写:今日趋势、适合推进 / 需要避开、一个行动(把技法审计表和「探索性日提示,不是确定预测」放进最后一节)。不要复述内部 JSON 字段。` - : `${natalToolInstruction}事业/财富/婚恋/家庭先给口语开场,再按四个标题写结论:先回答你的问题、盘里支持这个判断的地方、时间怎么看、这周可以做的一件事。不要把统一参数或技法审计表写进正文。不要复述内部 JSON 字段。`; + ? `${natalToolInstruction}${dailyAnswerShapeInstruction()}` + : `${natalToolInstruction}${natalAnswerShapeInstruction()}`; const adoptedRangeNote = consultationMode === "verified_chart" && prepared.serverChart?.truth.birthTimeStatus === "accepted" && prepared.serverChart.toolInput.candidate_range @@ -1009,11 +1017,21 @@ export async function POST(request: Request) { }; let baseMessages = consultationBaseMessages(false); let windowPacketMessage: string | null = null; - const agentAbortSignal = AbortSignal.timeout(AGENT_TIMEOUT_MS); + // One run, two clocks (BUG-1053, BUG-1051). Tools and the window precompute + // keep the tool phase's AGENT_TIMEOUT_MS deadline. The model loop runs on + // that deadline until the calculation result is in hand, then on the answer + // clock (CONSULTATION_ANSWER_TIMEOUT_MS), because its next step writes the + // answer. Continuation and answer retries share the same answer clock. + const runClock = createConsultationRunClock({ + toolPhaseMs: AGENT_TIMEOUT_MS, + answerMs: CONSULTATION_ANSWER_TIMEOUT_MS, + answerReady: () => state.consultationToolCompleted, + }); + const agentAbortSignal = runClock.toolSignal; const streamOptions = { runId: requestId, maxSteps: AGENT_MAX_STEPS, - abortSignal: agentAbortSignal, + abortSignal: runClock.loopSignal, hooks, ...consultationGenerationSettings(selectedModel.model), }; @@ -1025,11 +1043,10 @@ 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); + const answerPhaseSignal = runClock.answerSignal; + const startAnswerPhase = () => { + answerPhaseSignal(); + }; async function streamWithOverflowRetry( agent: { stream: ( @@ -1079,20 +1096,16 @@ export async function POST(request: Request) { const result = await streamWithOverflowRetry(agent); // This mode has no calculation to require and no chart method to bind, // so there is no contract for a retry to repair. - const retryForAnswer = async () => { + const retryForAnswer = async (retryHint?: string) => { const retried = await agent.stream([ ...baseMessages, - { role: "user" as const, content: "上一轮没有输出任何回答文本。请直接给出这个问题的回答,不要只说明过程。" }, + { role: "user" as const, content: `上一轮没有输出任何回答文本。请直接给出这个问题的回答,不要只说明过程。${retryHint ? `\n${retryHint}` : ""}` }, ], { ...streamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; const continueAfterLength = async (output: string) => { - const continued = await agent.stream([ - ...baseMessages, - { role: "assistant" as const, content: output }, - { role: "user" as const, content: consultationContinuePrompt(output) }, - ], { + const continued = await agent.stream(consultationContinueMessages(baseMessages, output), { ...streamOptions, abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), @@ -1175,23 +1188,26 @@ export async function POST(request: Request) { usages.push(retried.totalUsage); return retried.fullStream; }; - const retryForAnswer = async () => { + const retryForAnswer = async (retryHint?: string) => { const retried = await agent.stream([ ...baseMessages, { role: "user" as const, - content: "服务器窗口计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。", + content: `服务器窗口计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。${retryHint ? `\n${retryHint}` : ""}`, }, ], { ...streamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; - const continueAfterLength = async (output: string) => { - const continued = await agent.stream([ - ...baseMessages, - { role: "assistant" as const, content: output }, - { role: "user" as const, content: consultationContinuePrompt(output) }, - ], { + // A continuation is a new stream and the agent keeps nothing between + // streams (BUG-1053). The precomputed packet is already in baseMessages; + // a packet the model fetched itself is not, so it travels along. + const continueAfterLength = async (output: string, evidence?: unknown) => { + const continued = await agent.stream(consultationContinueMessages( + baseMessages, + output, + windowPacketMessage ? undefined : evidence, + ), { ...streamOptions, abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), @@ -1252,6 +1268,8 @@ export async function POST(request: Request) { retry, retryForAnswer, continueAfterLength, + stepScopedAnswer: true, + onAnswerPhase: startAnswerPhase, continueAfterDisconnect: true, pass4Mode: consultationMode, toolStatus: () => workflowStatus(state.workflowReceipt?.status), @@ -1308,23 +1326,22 @@ export async function POST(request: Request) { // The tool caches this request's calculation, so this attempt gets the same // evidence back without paying for it twice; keeping the tools available is // what puts that evidence in front of the model at all. - const retryForAnswer = async () => { + const retryForAnswer = async (retryHint?: string) => { const retried = await agent.stream([ ...baseMessages, { role: "user" as const, - content: "服务器计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。", + content: `服务器计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答;不要只描述过程或工具调用。${retryHint ? `\n${retryHint}` : ""}`, }, ], { ...natalStreamOptions, abortSignal: answerPhaseSignal() }); usages.push(retried.totalUsage); return retried.fullStream; }; - const continueAfterLength = async (output: string) => { - const continued = await agent.stream([ - ...baseMessages, - { role: "assistant" as const, content: output }, - { role: "user" as const, content: consultationContinuePrompt(output) }, - ], { + // A continuation is a new stream and the agent keeps nothing between + // streams: without the calculation result the loop saw, it would continue + // the answer blind (BUG-1053). + const continueAfterLength = async (output: string, evidence?: unknown) => { + const continued = await agent.stream(consultationContinueMessages(baseMessages, output, evidence), { ...streamOptions, abortSignal: answerPhaseSignal(), ...consultationContinueGenerationSettings(selectedModel.model), @@ -1332,23 +1349,6 @@ export async function POST(request: Request) { usages.push(continued.totalUsage); return continued.fullStream; }; - const interpretFindings = async () => ( - (state.thinkingPlan ?? []).map((section) => ({ id: section.id })) - ); - const composeAnswer = async (_: readonly { id: string; text?: string }[] = [], retryHint?: string) => { - const composed = await agent.stream([ - ...baseMessages, - { role: "user" as const, content: `${consultationComposePrompt()}${retryHint ? `\n${retryHint}` : ""}` }, - ], { - ...streamOptions, - abortSignal: answerPhaseSignal(), - maxSteps: AGENT_SLICE_MAX_STEPS, - toolChoice: "none", - ...consultationContinueGenerationSettings(selectedModel.model), - }); - usages.push(composed.totalUsage); - return composed.fullStream; - }; const executionReceipt = (): AgentExecutionReceipt => ({ runId: requestId, runtime: "mastra-agentic", @@ -1376,8 +1376,8 @@ export async function POST(request: Request) { retry, retryForAnswer, continueAfterLength, - interpretFindings, - composeAnswer, + stepScopedAnswer: true, + onAnswerPhase: startAnswerPhase, pass4Mode: consultationMode, continueAfterDisconnect: true, toolStatus: () => workflowStatus(state.workflowReceipt?.status), diff --git a/frontend/src/lib/consultation-thinking-plan.ts b/frontend/src/lib/consultation-thinking-plan.ts index 37eab9d0..3c455a72 100644 --- a/frontend/src/lib/consultation-thinking-plan.ts +++ b/frontend/src/lib/consultation-thinking-plan.ts @@ -330,40 +330,56 @@ export function consultationContinuePrompt(output: string): string { ].join("\n"); } -export type ConsultationSectionPromptReason = "write" | "empty-retry"; +/** + * The line that grounds the answer in this turn's calculation. Since BUG-1053 + * the loop's own step after the calculation result writes the answer; there is + * no separate compose stream (its prompt said 「服务器计算已经完成」 while it + * never saw the result), so the writing instructions travel with the user turn. + */ +const ANSWER_FROM_EVIDENCE = "拿到本轮计算结果后直接写给用户的回答,只用结果里的盘面事实下判断;调用工具前后都不要写过程说明。"; -export function consultationComposePrompt( - findings: ReadonlyArray<{ id: string; text?: string }> = [], -): string { - const grounded = findings - .filter((item) => item.text?.trim()) - .map((item) => `- ${item.id}: ${item.text}`) - .join("\n"); +/** Natal answer shape (VOICE §7): heading-free opener, then the four headings in one pass. */ +export function natalAnswerShapeInstruction(): string { return [ - "服务器计算已经完成。不要再调用排盘工具,不要重算。", - "只写结论散文。开场不要标题,然后一次写完这些二级标题,不要拆成多轮:", - `## ${REPORT_HEADING.question}`, - `## ${REPORT_HEADING.support}`, - `## ${REPORT_HEADING.timing}`, - `## ${REPORT_HEADING.action}`, - "不要写「统一参数与原始结构」,不要写技法审计表,不要写思考过程清单,不要复述内部 JSON。", - grounded ? `判断依据(写正文时引用,不要原样粘贴成列表):\n${grounded}` : "", - ].filter(Boolean).join("\n"); + ANSWER_FROM_EVIDENCE, + `开场不要标题,先用口语回答问题;然后一次写完这四个二级标题:## ${REPORT_HEADING.question}、## ${REPORT_HEADING.support}、## ${REPORT_HEADING.timing}、## ${REPORT_HEADING.action}。`, + "不要写「统一参数与原始结构」,不要写技法审计表,不要写思考过程清单,不要复述内部 JSON 字段。", + ].join(""); } -export function consultationSectionPrompt( - heading: string, - priorOutput: string, - options?: { reason?: ConsultationSectionPromptReason }, -): string { - const title = heading.trim() || REPORT_HEADING.question; +/** Daily entrypoint answer shape: three sections, audit table and boundary line in the last. */ +export function dailyAnswerShapeInstruction(): string { return [ - ...(options?.reason === "empty-retry" ? ["上一段没有输出正文,请直接写这一节。"] : []), - "服务器计算已经完成。不要再调用排盘工具,不要重算,不要读取其他二级标题。", - `只写这一个二级标题及其正文:## ${title}`, - "不要写其他 ## 标题,不要复述已经写出的段落,不要写思考过程清单。", - priorOutput.trim() - ? `已经写出的上文(冻结,勿重复):\n${priorOutput.slice(-4000)}` - : "这是正文的第一节。", + ANSWER_FROM_EVIDENCE, + `按三节写:${DAILY_HEADING.trend}、${DAILY_HEADING.actAvoid}、${DAILY_HEADING.action}(把技法审计表和「探索性日提示,不是确定预测」放进最后一节)。`, + "不要复述内部 JSON 字段。", + ].join(""); +} + +/** + * The calculation result handed to a length continuation (BUG-1053): the + * continuation is a new stream and the agent keeps no memory between streams, + * so without this it would continue the answer blind. + */ +export function consultationContinueEvidencePrompt(evidence: unknown): string { + return [ + "本轮服务器计算结果(同一请求已完成,不要重算)。续写时只用这里的盘面事实:", + JSON.stringify(evidence), ].join("\n"); } + +/** + * The messages a length continuation streams with: the run's own messages, the + * calculation result when the loop fetched one, the answer so far, and the + * continue instruction. One builder for every consultation route. + */ +export function consultationContinueMessages(base: readonly M[], output: string, evidence?: unknown) { + return [ + ...base, + ...(evidence === undefined + ? [] + : [{ role: "user" as const, content: consultationContinueEvidencePrompt(evidence) }]), + { role: "assistant" as const, content: output }, + { role: "user" as const, content: consultationContinuePrompt(output) }, + ]; +} diff --git a/frontend/src/lib/stream-agent-response.ts b/frontend/src/lib/stream-agent-response.ts index 1b6addc5..b3bdf48f 100644 --- a/frontend/src/lib/stream-agent-response.ts +++ b/frontend/src/lib/stream-agent-response.ts @@ -1,5 +1,7 @@ import { appendConsultationRuntimeStep, + CONSULTATION_NATAL_CALC_TOOL_ID, + CONSULTATION_WINDOW_CALC_TOOL_ID, type ConsultationRuntimeState, } from "../mastra/consultation-tools.ts"; import { @@ -13,7 +15,6 @@ import { toAgentModelFinishReason, type AgentModelFinishReason } from "./agent-o import { createVisibleTextTransformer } from "./stream-text-response.ts"; import { consultationWriteLabel } from "./consultation-activity-labels.ts"; import { logTruncatedReasoning } from "./consultation-budget.ts"; -import { acceptThinkStepText } from "./think-step-gate.ts"; import { classifyPass4, takeClosedSentences, @@ -308,11 +309,6 @@ export async function collectAgentPublicEvents(stream: ChunkStream | Iterable consultationAgentPublicEventSchema.parse(event)); } -export type ThinkFinding = Readonly<{ - id: string; - text?: string; -}>; - type AgentStreamSource = ChunkStream | (() => ChunkStream | Promise); async function resolveAgentStream(stream: AgentStreamSource): Promise { @@ -326,10 +322,31 @@ type StreamAgentResponseOptions = EventOptions & { transformText?: (text: string) => string; requireTool: boolean; retry?: () => Promise; - retryForAnswer?: () => Promise; - continueAfterLength?: (output: string) => Promise; - composeAnswer?: (findings: readonly ThinkFinding[], retryHint?: string) => Promise; - interpretFindings?: () => Promise; + /** + * 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; + /** + * 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; + /** + * 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; @@ -412,6 +429,28 @@ 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; +} + +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; @@ -422,8 +461,11 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { let firstOutput = false; let fullOutput = ""; let uncontractedText = ""; - let thinkingText = ""; 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; // 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 @@ -435,8 +477,9 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { // 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. + // 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 | undefined, event: ConsultationAgentPublicEvent) => { @@ -518,24 +561,52 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { await releasePass4Sentences(controller, "", true); } + function markAnswerPhase() { + if (answerPhaseStarted || !contractReady(options)) return; + answerPhaseStarted = true; + options.onAnswerPhase?.(); + } + async function consumeAttempt( controller: ReadableStreamDefaultController | undefined, stream: ChunkStream, - attempt: { drainSpoken?: boolean; suppressCompositionActivity?: boolean; answerPhase?: boolean } = {}, + 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); - const drainingSpoken = () => Boolean(attempt.drainSpoken) && contractReady(options); + // 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; + 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 @@ -543,14 +614,17 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { // 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) || drainingSpoken()) { - if (!contractReady(options) && text) uncontractedText += text; + if (!contractReady(options)) { + if (text) uncontractedText += text; return; } uncontractedText = ""; held += text; if (!held) return; - if (/\S/.test(held)) outcome.contributed = true; + if (/\S/.test(held)) { + outcome.contributed = true; + stepWrote = true; + } if (!composingSent) { composingSent = true; send(controller, { type: "activity", phase: "answer-composition", label: "正在组织回答" }); @@ -569,6 +643,27 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { if (/\S/.test(held)) emitted = true; held = ""; }; + const acceptText = 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) => { + 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; @@ -593,34 +688,57 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { }); throw new Error(code); } - for (const event of mapChunk(chunk, options, startedAt, toolErrors)) send(controller, event); + 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; + } + markAnswerPhase(); flushThinkingPlan(controller); + if (chunk.type === "step-start") { + resetStep(); + stepWrote = false; + } + if (chunk.type === "tool-call") { + // Whatever this step said before calling a tool is narration. + stepCalledTool = true; + stepText = ""; + } 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", - }); + 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 === "step-finish") options.state.modelStepCount += 1; if (chunk.type === "finish") { const finish = finishTelemetry(chunk); options.state.modelFinishReason = finish.reason; - outcome.finishReason = 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 outputText(visible.push(chunk.payload.text)); + await acceptText(visible.push(chunk.payload.text)); } } flushThinkingPlan(controller); - await outputText(visible.finish("")); + await acceptText(visible.finish("")); + await settleStep(); lastAttempt = outcome; if (outcome.contributed || outcome.answerPhase) answerTail = outcome; } catch (error) { @@ -630,21 +748,29 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } if (isTimeoutOrAbort(error) && !outcome.aborted) { outcome.aborted = true; - appendConsultationRuntimeStep(options.state, { - kind: "abort", - name: attempt.drainSpoken ? "tool-abort" : "compose-abort", - status: "failed", - }); + recordAttemptAbort(); } + try { + await acceptText(visible.finish("")); + await settleStep(); + } catch {} lastAttempt = outcome; if (outcome.contributed || outcome.answerPhase) answerTail = outcome; - try { - await outputText(visible.finish("")); - } catch {} 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 | undefined, heading?: string, @@ -658,7 +784,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { phase: "answer-composition", label: heading ? consultationWriteLabel(heading, true) : "正在组织回答", }); - await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer()), { + await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer(), calculationEvidence), { suppressCompositionActivity: true, answerPhase: true, }); @@ -668,67 +794,19 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } } - async function publishFindings( - controller: ReadableStreamDefaultController | undefined, - ): Promise { - flushThinkingPlan(controller); - const plan = options.state.thinkingPlan ?? []; - const steps = thinkPlanSteps(plan); - let findings: readonly ThinkFinding[] = []; - if (options.interpretFindings) { - send(controller, { - type: "phase.started", - phase: "interpret", - label: "正在整理判断依据", - }); - const started = Date.now(); - findings = await options.interpretFindings(); - send(controller, { - type: "phase.completed", - phase: "interpret", - durationMs: Math.max(0, Date.now() - started), - }); - } - const byId = new Map(findings.map((item) => [item.id, item])); - for (const step of steps) { - send(controller, { type: "think.step", id: step.id, status: "running" }); - const accepted = acceptThinkStepText(byId.get(step.id)?.text ?? ""); - if (accepted) { - thinkingText += accepted; - send(controller, { type: "think.step", id: step.id, status: "done", text: accepted }); - } else { - send(controller, { type: "think.step", id: step.id, status: "done" }); - } - } - return findings; - } - + /** + * 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 | undefined, origin: string, - findings: readonly ThinkFinding[], - allowComposeRetry = true, ) { if (!options.pass4Mode) return; await releaseFinalSentence(controller); const produced = () => fullOutput.slice(origin.length); - const hadRetryableReject = options.state.steps.some((step) => - step.name === "pass4-reject:guarantee" - || step.name === "pass4-reject:personal-chart" - || step.name === "pass4-reject:methodology" - ); - if (allowComposeRetry && !/\S/.test(produced()) && hadRetryableReject && options.composeAnswer) { - fullOutput = origin; - pass4Buffer = ""; - await consumeAttempt( - controller, - await options.composeAnswer(findings, PASS4_RETRY_HINT), - { suppressCompositionActivity: true, answerPhase: true }, - ); - await continueCurrentAnswer(controller); - await releaseFinalSentence(controller); - } - if (!/\S/.test(produced()) && options.pass4Mode === "general_no_birth_time" && hadRetryableReject) { + if (!/\S/.test(produced()) && options.pass4Mode === "general_no_birth_time" && hadRetryablePass4Reject()) { if (!firstOutput) { firstOutput = true; await options.onFirstOutput?.(); @@ -739,36 +817,12 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } } - async function composeOnce( - controller: ReadableStreamDefaultController | undefined, - findings: readonly ThinkFinding[], - ) { - if (!options.composeAnswer) return false; - send(controller, { - type: "activity", - phase: "answer-composition", - label: consultationWriteLabel("回答", true), - }); - send(controller, { - type: "phase.started", - phase: "compose", - label: "正在写结论", - }); - const started = Date.now(); - const origin = fullOutput; - await consumeAttempt( - controller, - await options.composeAnswer(findings), - { suppressCompositionActivity: true, answerPhase: 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" ); - await continueCurrentAnswer(controller); - await finishPass4(controller, origin, findings); - send(controller, { - type: "phase.completed", - phase: "compose", - durationMs: Math.max(0, Date.now() - started), - }); - return true; } async function deliverDegradedAnswer( @@ -841,15 +895,11 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { await options.warmup((event) => send(controller, event)); flushThinkingPlan(controller); } - await consumeAttempt(controller, await resolveAgentStream(options.stream), { - drainSpoken: Boolean(options.composeAnswer), - }); + 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(), { - drainSpoken: Boolean(options.composeAnswer), - }); + await consumeAttempt(controller, await options.retry()); } if (!contractReady(options)) { const canDegrade = options.requireTool @@ -868,18 +918,23 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { } } if (!deliveredDegraded) { - const findings = await publishFindings(controller); - const composed = await composeOnce(controller, findings); - if (!composed) { - await continueCurrentAnswer(controller); - if (options.pass4Mode) await finishPass4(controller, "", findings); - } + // 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; - await consumeAttempt(controller, await options.retryForAnswer(), { answerPhase: true }); - if (options.pass4Mode) await finishPass4(controller, retryOrigin, findings, false); + 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"); @@ -898,10 +953,12 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { settling = 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, - thinkingText.trim() || undefined, + undefined, thinkingSections.length > 0 ? thinkingSections : undefined, ); settled = true; diff --git a/frontend/src/mastra/consultation-tools.ts b/frontend/src/mastra/consultation-tools.ts index ce4f282f..9b326343 100644 --- a/frontend/src/mastra/consultation-tools.ts +++ b/frontend/src/mastra/consultation-tools.ts @@ -55,36 +55,83 @@ export { AGENT_MAX_OUTPUT_TOKENS as CONSULTATION_MAX_OUTPUT_TOKENS } from "../li // calculation inside the 110s budget—so it is measured, not guessed. 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). + * The answer phase runs on its own clock, started when the answer phase + * starts. It used to share AGENT_TIMEOUT_MS with the tool loop, so a long loop + * left the answer a few seconds and Mastra's abort cut it mid-sentence + * (BUG-1051; BUG-944 had already asked for one budget per stream). + * + * Since BUG-1053 there is no separate compose stream: the tool loop's own step + * after the calculation result writes the answer, with that result in its + * context. The phase therefore starts when the calculation result is in hand + * (the run contract turns ready), not when a second stream opens, and it + * covers the rest of the loop plus any length continuation or empty-answer + * retry. * * 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. + * therefore holds about 6,300-7,000 output tokens, three to four times a + * typical four-heading answer; at half that throughput it still holds about + * 3,000 tokens. The answer step keeps provider thinking on (it is the loop's + * step), so thinking tokens now come out of the same 70s; the loop's step + * after the tool result used to think and write inside the 110s as well, and + * its text was thrown away. The calculation must finish inside the 110s tool + * phase, so the worst case is 110 + 70 = 180s, the three minutes the product + * accepted. The route's maxDuration must stay above the sum. */ -export const CONSULTATION_COMPOSE_TIMEOUT_MS = 70_000; +export const CONSULTATION_ANSWER_TIMEOUT_MS = 70_000; + +export type ConsultationRunClock = Readonly<{ + /** The tool phase's deadline. Tools and precompute run under it, unchanged. */ + toolSignal: AbortSignal; + /** + * The model loop's signal: the tool-phase deadline until the answer phase + * starts, then the answer clock. A loop that is writing the answer is never + * cut by the tool phase's timer. + */ + loopSignal: AbortSignal; + /** + * The answer clock. The first call starts it; every later call (loop hand + * over, continuation, retries) returns the same signal, so the phase as a + * whole is bounded. + */ + answerSignal: () => AbortSignal; +}>; /** - * 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. + * One consultation run's clocks (BUG-1053). `answerReady` covers the race + * where the calculation settles just before the tool-phase timer fires but + * the stream consumer has not seen the result yet: the loop is then handed to + * the answer clock instead of being cut. */ -export function createConsultationAnswerClock(timeoutMs = CONSULTATION_COMPOSE_TIMEOUT_MS) { - let signal: AbortSignal | null = null; - return () => { - signal ??= AbortSignal.timeout(timeoutMs); - return signal; +export function createConsultationRunClock(options: { + toolPhaseMs?: number; + answerMs?: number; + answerReady?: () => boolean; +} = {}): ConsultationRunClock { + const toolPhaseMs = options.toolPhaseMs ?? AGENT_TIMEOUT_MS; + const answerMs = options.answerMs ?? CONSULTATION_ANSWER_TIMEOUT_MS; + const toolSignal = AbortSignal.timeout(toolPhaseMs); + const loop = new AbortController(); + let answer: AbortSignal | null = null; + const answerSignal = () => { + if (!answer) { + const started = AbortSignal.timeout(answerMs); + answer = started; + started.addEventListener("abort", () => loop.abort(started.reason), { once: true }); + } + return answer; }; + toolSignal.addEventListener("abort", () => { + if (answer) return; + if (options.answerReady?.()) { + answerSignal(); + return; + } + loop.abort(toolSignal.reason); + }, { once: true }); + return Object.freeze({ toolSignal, loopSignal: loop.signal, answerSignal }); } export const CONSULTATION_NATAL_CALC_TOOL_ID = "run-jyotish-consultation"; export const CONSULTATION_WINDOW_CALC_TOOL_ID = "run-jyotish-window-consultation"; diff --git a/frontend/tests/application-billing-contract.test.ts b/frontend/tests/application-billing-contract.test.ts index fadcf0f0..468ccdf5 100644 --- a/frontend/tests/application-billing-contract.test.ts +++ b/frontend/tests/application-billing-contract.test.ts @@ -188,7 +188,10 @@ test("standard consultation awaits real usage before durable response settlement // Window and natal each have a contract retry plus an answer retry; general // has only an answer retry. Every one of those model calls spends tokens. assert.equal(consultRoute.match(/usages\.push\(retried\.totalUsage\)/g)?.length, 5); - assert.equal(consultRoute.match(/const retryForAnswer = async \(\) => \{/g)?.length, 3); + // 原值: /const retryForAnswer = async \(\) => \{/ 三处 + // 新值: /const retryForAnswer = async \(retryHint\?: string\) => \{/ 三处 + // 原因: BUG-1053 删除 compose 后,Pass 4 整篇被拒的重写提示改由 answer retry 携带;三处回答重试仍各计一次 usage + assert.equal(consultRoute.match(/const retryForAnswer = async \(retryHint\?: string\) => \{/g)?.length, 3); }); test("standard consultation forwards its stable reservation request as the usage event key", async () => { diff --git a/frontend/tests/consult-answer-truncation-20260926.test.ts b/frontend/tests/consult-answer-truncation-20260926.test.ts index 39cf965e..203b77cc 100644 --- a/frontend/tests/consult-answer-truncation-20260926.test.ts +++ b/frontend/tests/consult-answer-truncation-20260926.test.ts @@ -5,6 +5,15 @@ // 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"; @@ -16,10 +25,10 @@ import { natalConsultationThinkingPlan, REPORT_HEADING } from "../src/lib/consul import { streamAgentResponse } from "../src/lib/stream-agent-response.ts"; import { AGENT_TIMEOUT_MS, - CONSULTATION_COMPOSE_TIMEOUT_MS, + CONSULTATION_ANSWER_TIMEOUT_MS, consultationModelStepTelemetry, consultationStepBudgetReceipt, - createConsultationAnswerClock, + createConsultationRunClock, createConsultationRuntimeState, publicConsultationRuntimeSteps, } from "../src/mastra/consultation-tools.ts"; @@ -149,12 +158,6 @@ async function realStream( return result.fullStream as ReadableStream; } -/** 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>; continueAfterLength?: () => Promise>; @@ -167,13 +170,14 @@ async function runNatal(input: { runId: "run", requestId: "req", state, - stream: toolLoop(), + // 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, - 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"); @@ -224,28 +228,35 @@ test("a shared timeout firing mid-answer ends as answer_truncated, not a charged }); 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); + // 原值: 工具循环的 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(loopSignal.aborted, true); - - const answerClock = createConsultationAnswerClock(2_000); + assert.equal(clock.toolSignal.aborted, true); const run = await runNatal({ - compose: () => realStream([...pieces(CLOSED), DANGLING, REST], "stop", { delayMs: 5, abortSignal: answerClock() }), + 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: 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 }), + // 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(shared.completed, null); - assert.notEqual(shared.terminal[0]?.type, "run.completed"); + 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 () => { - const clock = createConsultationAnswerClock(30); + // 原值: 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"); @@ -336,25 +347,26 @@ test("the observability record carries the compose ending and answer length, nev }); 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_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\(\),/); + 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) ?? []; + 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) ?? []; + 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 tool loop keeps its own 110s signal. - assert.match(route, /const agentAbortSignal = AbortSignal\.timeout\(AGENT_TIMEOUT_MS\)/); + // 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_COMPOSE_TIMEOUT_MS < maxDuration * 1000); - assert.ok(AGENT_TIMEOUT_MS + CONSULTATION_COMPOSE_TIMEOUT_MS <= 180_000, "product accepted about three minutes"); + 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"); }); diff --git a/frontend/tests/consult-single-pass-answer-20260927.test.ts b/frontend/tests/consult-single-pass-answer-20260927.test.ts new file mode 100644 index 00000000..dbc972e0 --- /dev/null +++ b/frontend/tests/consult-single-pass-answer-20260927.test.ts @@ -0,0 +1,440 @@ +// BUG-1053: the user-visible consultation answer was written without the chart. +// +// The natal loop's step after `run-jyotish-consultation` saw the evidence, but +// its text was drained and a second `agent.stream` (compose) wrote the answer +// from history + the question only. These tests drive the real personal Agent +// (`getJyotishAgent`, real skill binding, real calculation tool, real +// prepareStep) over a prompt-recording fake model, with the calculation fed by +// the golden engine capture in fixtures/ (AGENTS §7.4). What each model call +// was shown is the assertion: the call that writes the answer must have the +// tool result in its prompt, and there must be no second, blind call. +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; + +import { getJyotishAgent } from "../src/mastra/index.ts"; +import { createNdjsonParser } from "../src/lib/consultation-agent-events.ts"; +import { + consultationContinueMessages, + natalAnswerShapeInstruction, + REPORT_HEADING, +} from "../src/lib/consultation-thinking-plan.ts"; +import { streamAgentResponse } from "../src/lib/stream-agent-response.ts"; +import { PASS4_RETRY_HINT } from "../src/lib/timing-output-guard.ts"; +import { + AGENT_MAX_STEPS, + AGENT_TIMEOUT_MS, + CONSULTATION_ANSWER_TIMEOUT_MS, + consultationNatalPrepareStep, + consultationStepBudgetReceipt, + createConsultationAgentContext, + createConsultationRunClock, + createConsultationRuntimeState, + publicConsultationRuntimeSteps, +} from "../src/mastra/consultation-tools.ts"; + +const golden = JSON.parse(readFileSync( + new URL("./fixtures/consultation-workflow-report-blocked-repairs-golden.json", import.meta.url), + "utf8", +)) as { smoke_birth: Record; themes: Record> }; +const careerWorkflow = golden.themes.career!; +// A fact only the calculation carries: the golden chart's Moon longitude. +const GOLDEN_MOON_DEGREE = "3.2259"; + +// Fictional capture identity from the fixture itself, not a real person. +const serverChart = { + name: "测试", + toolInput: { + year: 1993, month: 6, day: 15, hour: 10, minute: 30, city: "Handan", lat: 36.42, lon: 114.21, tz: 8, + ayanamsa: "raman" as const, declared_accuracy: "15min" as const, time_source: "family_vague", + }, + truth: { + birthDate: "1993-06-15", reportedBirthTime: "10:30", activeBirthTime: null, + selectedTimeKind: "reported" as const, birthTimeSource: "reported", birthTimeStatus: "reported", + placeLabel: "Handan", placeCodes: { countryCode: "CN", provinceCode: null, cityCode: null, districtCode: null }, + placeId: null, placeType: "city", placeProvider: "profile", latitude: 36.42, longitude: 114.21, + timezoneId: "Asia/Shanghai", timezoneSource: "profile", timezoneOffset: 8, + }, +}; + +const PRE_TOOL = "我先排一下盘,稍等。"; +const BETWEEN_TOOLS = "我再核对一次计算结果。"; +const OPENER = "你这盘外面看着稳,底下其实一直在换跑道:表面求安定,底下要自己说了算。"; +const ANSWER = [ + OPENER, + "", + `## ${REPORT_HEADING.question}`, + "能换,但先换做法,再换岗位。", + "", + `## ${REPORT_HEADING.support}`, + "月亮落在白羊座九宫,想法来得快。", + "", + `## ${REPORT_HEADING.timing}`, + "这半年先试,不急着定。", + "", + `## ${REPORT_HEADING.action}`, + "- **这周**:约一位同行聊一次。", + "", +].join("\n"); + +type Part = { text?: string; tool?: { name: string; input: Record } }; +type Turn = { parts: Part[]; finish: string; delayMs?: number }; + +/** A LanguageModelV2 that plays `turns` in order and records every prompt. */ +function scriptedModel(turns: Turn[]) { + const prompts: unknown[][] = []; + let call = 0; + const model = { + specificationVersion: "v2", + provider: "fake", + modelId: "fake-single-pass", + supportedUrls: {}, + async doGenerate() { + throw new Error("not used"); + }, + async doStream(options: { prompt: unknown[]; abortSignal?: AbortSignal }) { + prompts.push(options.prompt); + const turn = turns[call] ?? { parts: [], finish: "stop" }; + call += 1; + const signal = options.abortSignal; + const stream = new ReadableStream({ + async start(controller) { + controller.enqueue({ type: "stream-start", warnings: [] }); + let textOpen = false; + for (const [index, part] of turn.parts.entries()) { + if (turn.delayMs) { + try { + await new Promise((resolve, reject) => { + if (signal?.aborted) return reject(signal.reason); + const timer = setTimeout(resolve, turn.delayMs); + signal?.addEventListener("abort", () => { + clearTimeout(timer); + reject(signal.reason); + }, { once: true }); + }); + } catch (error) { + controller.error(error); + return; + } + } + if (part.text !== undefined) { + if (!textOpen) { + controller.enqueue({ type: "text-start", id: `t${call}` }); + textOpen = true; + } + controller.enqueue({ type: "text-delta", id: `t${call}`, delta: part.text }); + } + if (part.tool) { + if (textOpen) { + controller.enqueue({ type: "text-end", id: `t${call}` }); + textOpen = false; + } + controller.enqueue({ + type: "tool-call", + toolCallId: `call-${call}-${index}`, + toolName: part.tool.name, + input: JSON.stringify(part.tool.input), + }); + } + } + if (textOpen) controller.enqueue({ type: "text-end", id: `t${call}` }); + controller.enqueue({ + type: "finish", + finishReason: turn.finish, + usage: { inputTokens: 10, outputTokens: 10, totalTokens: 20 }, + }); + controller.close(); + }, + }); + return { stream }; + }, + }; + return { model, prompts, calls: () => call }; +} + +const calcCall = (text?: string): Turn => ({ + parts: [ + ...(text ? [{ text }] : []), + { tool: { name: "run-jyotish-consultation", input: { question: "我的事业能换方向吗", domains: ["career"] } } }, + ], + finish: "tool-calls", +}); + +function pieces(text: string, size = 12) { + return (text.match(new RegExp(`[\\s\\S]{1,${size}}`, "g")) ?? []).map((value) => ({ text: value })); +} + +let seq = 0; + +/** + * The natal route's wiring, minus HTTP, auth and billing: the same Agent, + * prepareStep, run clock, stream options, step-scoped answer, answer-phase + * hand-over, continuation builder and answer retry as `runAgenticConsultation`. + */ +async function runNatal(turns: Turn[], options: { toolPhaseMs?: number; answerMs?: number } = {}) { + seq += 1; + const { model, prompts, calls } = scriptedModel(turns); + const state = createConsultationRuntimeState({ plannedSteps: AGENT_MAX_STEPS }); + const clock = createConsultationRunClock({ + toolPhaseMs: options.toolPhaseMs ?? AGENT_TIMEOUT_MS, + answerMs: options.answerMs ?? CONSULTATION_ANSWER_TIMEOUT_MS, + answerReady: () => state.consultationToolCompleted, + }); + let workflowRuns = 0; + const agentContext = createConsultationAgentContext({ + userId: "u", sessionId: "s", requestId: `single-pass-${seq}`, consultationMode: "verified_chart", + theme: "career", serverChart: serverChart as never, abortSignal: clock.toolSignal, state, + runWorkflow: async () => { + workflowRuns += 1; + return structuredClone(careerWorkflow) as never; + }, + }); + const agent = getJyotishAgent({ id: `fake-${seq}`, model } as never, agentContext); + const baseMessages = [{ role: "user" as const, content: `${natalAnswerShapeInstruction()}\n问题:我的事业能换方向吗` }]; + const streamOptions = { runId: `run-${seq}`, maxSteps: AGENT_MAX_STEPS, abortSignal: clock.loopSignal }; + const natalStreamOptions = { ...streamOptions, prepareStep: consultationNatalPrepareStep }; + const continuations: unknown[] = []; + const retryHints: Array = []; + let completed: string | null = null; + let charges = 0; + let errored: unknown = null; + const result = await agent.stream(baseMessages as never, natalStreamOptions as never); + const response = streamAgentResponse({ + runId: `run-${seq}`, + requestId: `req-${seq}`, + state, + stream: result.fullStream as ReadableStream, + requireTool: true, + stepScopedAnswer: true, + onAnswerPhase: () => { clock.answerSignal(); }, + pass4Mode: "verified_chart", + retryForAnswer: async (retryHint) => { + retryHints.push(retryHint); + const retried = await agent.stream([ + ...baseMessages, + { role: "user" as const, content: `服务器计算已经完成,但上一轮没有输出任何回答文本。请重新取回本次计算结果,然后直接给出回答。${retryHint ? `\n${retryHint}` : ""}` }, + ] as never, { ...natalStreamOptions, abortSignal: clock.answerSignal() } as never); + return retried.fullStream as ReadableStream; + }, + continueAfterLength: async (output, evidence) => { + continuations.push(evidence); + const continued = await agent.stream( + consultationContinueMessages(baseMessages, output, evidence) as never, + { ...streamOptions, abortSignal: clock.answerSignal() } as never, + ); + return continued.fullStream as ReadableStream; + }, + toolStatus: () => "ready", + receipt: () => ({ + runId: `run-${seq}`, + runtime: "mastra-agentic", + skill: { name: "jyotish-vedic-astrology", loaded: true, referenceReads: 0, methodologySections: 0 }, + steps: publicConsultationRuntimeSteps(state), + stepBudget: consultationStepBudgetReceipt(state), + workflow: state.workflowReceipt ?? { route: "pending", status: "blocked", preciseTiming: "blocked", missingLayers: [] }, + techniqueTruth: "unknown", + }) as never, + onComplete: (output) => { completed = output; charges += 1; }, + onError: (error) => { errored = error; }, + }); + const events: Array<{ type: string; code?: string; text?: string; receipt?: unknown }> = []; + const parser = createNdjsonParser((event) => events.push(event as never)); + 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, prompts, calls: calls(), continuations, retryHints, workflowRuns, + completed: completed as string | null, charges, errored, + }; +} + +/** Index of the first prompt that already contains the calculation result. */ +function firstPromptWithToolResult(prompts: unknown[][]) { + return prompts.findIndex((prompt) => prompt.some((message) => { + const value = message as { role?: string; content?: unknown }; + return value.role === "tool" && JSON.stringify(value.content).includes("run-jyotish-consultation"); + })); +} + +test("the model call that writes the answer has the calculation result in its prompt, and no blind call follows", async () => { + const run = await runNatal([calcCall(PRE_TOOL), { parts: pieces(ANSWER), finish: "stop" }]); + + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.workflowRuns, 1); + // Exactly two model calls: the one that called the tool and the one after it. + // The pre-fix route opened a third (compose) whose prompt had no tool result. + assert.equal(run.calls, 2); + const withResult = firstPromptWithToolResult(run.prompts); + assert.equal(withResult, 1, "the call after the tool result sees it"); + assert.equal(run.prompts.length - withResult, 1, "exactly one model call after the tool result"); + assert.ok(JSON.stringify(run.prompts[1]).includes(GOLDEN_MOON_DEGREE), "the golden chart fact reaches the writer"); + // What that call wrote is what the user got, with the four headings intact. + assert.equal(run.answer, ANSWER); + assert.equal(run.completed, ANSWER); + for (const heading of Object.values(REPORT_HEADING)) assert.match(run.answer, new RegExp(`^## ${heading}$`, "m")); + assert.ok(run.answer.startsWith(OPENER), "the opener has no heading"); + // Pre-tool narration is not the answer. + assert.equal(run.answer.includes(PRE_TOOL), false); +}); + +test("the writing instructions travel with the user turn the loop sees", async () => { + const run = await runNatal([calcCall(), { parts: pieces(ANSWER), finish: "stop" }]); + const writerPrompt = JSON.stringify(run.prompts[1]); + for (const heading of Object.values(REPORT_HEADING)) assert.ok(writerPrompt.includes(`## ${heading}`)); + assert.match(writerPrompt, /开场不要标题/); + // The removed compose prompt claimed a finished calculation it never saw. + assert.doesNotMatch(writerPrompt, /服务器计算已经完成。不要再调用排盘工具/); +}); + +test("narration in a step that calls a tool never reaches the answer, even after the calculation", async () => { + const run = await runNatal([ + calcCall(PRE_TOOL), + // After the calculation the model narrates and calls the tool again (a + // request-cache hit), then writes the answer. + calcCall(BETWEEN_TOOLS), + { parts: pieces(ANSWER), finish: "stop" }, + ]); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.workflowRuns, 1, "the second call is served from the request cache"); + assert.equal(run.answer, ANSWER); + assert.equal(run.answer.includes(PRE_TOOL), false); + assert.equal(run.answer.includes(BETWEEN_TOOLS), false); + assert.equal(run.charges, 1); +}); + +test("a normal stop completes and charges exactly once", async () => { + const run = await runNatal([calcCall(), { parts: pieces(ANSWER), finish: "stop" }]); + assert.equal(run.charges, 1); + assert.equal(run.errored, null); + assert.equal(run.state.composeFinishReason, "stop", "the answer-writing step ended on its own"); + assert.equal(run.state.composeAborted, false); + assert.equal(run.state.steps.some((step) => step.kind === "abort"), false); +}); + +test("length continues with the calculation result in the continuation prompt", async () => { + const cut = ANSWER.indexOf(`## ${REPORT_HEADING.timing}`); + const head = ANSWER.slice(0, cut); + const tail = ANSWER.slice(cut); + const run = await runNatal([ + calcCall(), + { parts: pieces(head), finish: "length" }, + { parts: pieces(tail), finish: "stop" }, + ]); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.continuations.length, 1); + assert.ok(run.continuations[0], "the continuation receives the calculation result"); + assert.equal(run.calls, 3); + const continuationPrompt = JSON.stringify(run.prompts[2]); + assert.ok(continuationPrompt.includes(GOLDEN_MOON_DEGREE), "the continuation is not blind"); + assert.ok(continuationPrompt.includes(REPORT_HEADING.support), "and it sees the answer so far"); + assert.equal(run.completed, ANSWER); + assert.equal(run.charges, 1); +}); + +for (const reason of ["content-filter", "tool-calls", "other", "unknown"] as const) { + test(`an answer step that ends on ${reason} with visible text is truncated and not charged`, async () => { + const run = await runNatal([calcCall(), { parts: pieces(ANSWER.slice(0, 120)), finish: reason }]); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.charges, 0); + assert.equal(run.continuations.length, 0, "only length is continued"); + assert.ok(run.state.steps.some((step) => step.name === "answer-truncated" && step.status === "failed")); + }); +} + +test("the answer clock cutting the answer step mid-sentence is truncated, recorded and not charged", async () => { + const run = await runNatal( + [calcCall(), { parts: pieces(ANSWER, 6), finish: "stop", delayMs: 25 }], + { answerMs: 250 }, + ); + assert.deepEqual(run.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]); + assert.equal(run.charges, 0); + assert.ok(run.answer.length > 0, "what was written stays"); + assert.ok(ANSWER.startsWith(run.answer), "and nothing but the answer was written"); + 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")); + assert.equal(run.state.modelFinishReason, "tripwire"); + assert.equal(run.state.composeAborted, true); +}); + +test("the tool phase's deadline does not starve the answer step", async () => { + // The answer step takes about 0.5s; the tool phase's clock is 150ms. The + // loop is handed to the answer clock when the calculation result arrives, + // so the tool phase expiring mid-answer does not cut it. + const run = await runNatal( + [calcCall(), { parts: pieces(ANSWER, 6), finish: "stop", delayMs: 12 }], + { toolPhaseMs: 150, answerMs: 5_000 }, + ); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.completed, ANSWER); +}); + +test("without the hand-over the tool phase's deadline would cut the answer (control)", async () => { + const clock = createConsultationRunClock({ toolPhaseMs: 30, answerMs: 5_000 }); + await new Promise((resolve) => setTimeout(resolve, 60)); + assert.equal(clock.loopSignal.aborted, true, "a loop never handed over ends with the tool phase"); + + const handed = createConsultationRunClock({ toolPhaseMs: 30, answerMs: 5_000 }); + handed.answerSignal(); + await new Promise((resolve) => setTimeout(resolve, 60)); + assert.equal(handed.toolSignal.aborted, true, "tools still stop at the tool-phase deadline"); + assert.equal(handed.loopSignal.aborted, false, "the answer-writing loop is not cut by it"); + + // The calculation settled just before the tool-phase timer fired, but the + // consumer has not seen the result yet: the loop is handed over, not cut. + const raced = createConsultationRunClock({ toolPhaseMs: 30, answerMs: 60, answerReady: () => true }); + await new Promise((resolve) => setTimeout(resolve, 45)); + assert.equal(raced.loopSignal.aborted, false); + await new Promise((resolve) => setTimeout(resolve, 80)); + assert.equal(raced.loopSignal.aborted, true, "the answer clock still bounds it"); +}); + +test("an empty answer falls back to the answer retry, which keeps the tools and hits the request cache", async () => { + const run = await runNatal([ + calcCall(), + { parts: [], finish: "stop" }, + calcCall(), + { parts: pieces(ANSWER), finish: "stop" }, + ]); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); + assert.equal(run.retryHints.length, 1); + assert.equal(run.retryHints[0], undefined, "no Pass 4 hint when nothing was rejected"); + assert.equal(run.workflowRuns, 1); + assert.ok(run.state.steps.some((step) => step.name === "answer-retry")); + assert.equal(run.answer, ANSWER); + // The retry re-fetch is not announced as a second calculation. + assert.equal(run.events.filter((event) => event.type === "tool.started").length, 1); +}); + +test("an answer Pass 4 rejected whole is retried with the rewrite hint", async () => { + const run = await runNatal([ + calcCall(), + { parts: [{ text: "我保证你一定会升职。" }], finish: "stop" }, + { parts: pieces(ANSWER), finish: "stop" }, + ]); + assert.equal(run.answer.includes("一定会升职"), false); + assert.deepEqual(run.retryHints, [PASS4_RETRY_HINT]); + assert.equal(run.answer, ANSWER); + assert.deepEqual(run.terminal.map((event) => event.type), ["run.completed"]); +}); + +test("the consult route has no separate compose stream and wires the single-pass answer", () => { + const route = readFileSync(new URL("../src/app/api/consult/route.ts", import.meta.url), "utf8"); + const stream = readFileSync(new URL("../src/lib/stream-agent-response.ts", import.meta.url), "utf8"); + assert.doesNotMatch(route, /composeAnswer|interpretFindings|consultationComposePrompt|toolChoice: "none"/); + assert.doesNotMatch(stream, /composeAnswer|interpretFindings|drainSpoken|publishFindings/); + const natal = route.slice(route.indexOf("if (!prepared.serverChart)")); + assert.match(natal, /stepScopedAnswer: true,/); + assert.match(natal, /onAnswerPhase: startAnswerPhase,/); + assert.match(natal, /consultationContinueMessages\(baseMessages, output, evidence\)/); + const window = route.slice(route.indexOf("if (shouldRunDeclaredWindowWorkflow(consultationMode)) {"), route.indexOf("if (!prepared.serverChart)")); + assert.match(window, /stepScopedAnswer: true,/); + assert.match(window, /windowPacketMessage \? undefined : evidence/); + assert.match(route, /natalAnswerShapeInstruction\(\)/); + // One run clock: tools on the tool phase, the loop handed to the answer clock. + assert.match(route, /const agentAbortSignal = runClock\.toolSignal;/); + assert.match(route, /abortSignal: runClock\.loopSignal,/); + assert.match(route, /const answerPhaseSignal = runClock\.answerSignal;/); + 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"); +}); diff --git a/frontend/tests/consultation-agentic-runtime.test.ts b/frontend/tests/consultation-agentic-runtime.test.ts index 20b59900..557f9926 100644 --- a/frontend/tests/consultation-agentic-runtime.test.ts +++ b/frontend/tests/consultation-agentic-runtime.test.ts @@ -1465,6 +1465,9 @@ test("degraded delivery keeps only the last attempt body (BUG-960)", async () => }); test("degraded delivery does not start a compose pass (BUG-961)", async () => { + // 原值: 传入 composeAnswer,断言降级交付后 composeCalls === 0 + // 新值: compose 已删除;改为传入 retryForAnswer,断言降级交付后不再发起任何写回答的流 + // 原因: BUG-1053 删除单独的 compose 流;「降级正文不得再被第二遍写作覆盖」这一性质不变 const state = createConsultationRuntimeState(); let composeCalls = 0; async function* chunks() { @@ -1475,7 +1478,7 @@ test("degraded delivery does not start a compose pass (BUG-961)", async () => { runId: "run", requestId: "req", state, stream: chunks(), requireTool: true, pass4Mode: "verified_chart", toolStatus: () => "blocked", receipt: () => ({ ...receipt(state), steps: publicConsultationRuntimeSteps(state) }), - composeAnswer: async () => { + retryForAnswer: async () => { composeCalls += 1; async function* composed() { yield { type: "text-delta", payload: { text: "不该出现的 compose。" } }; @@ -2053,53 +2056,57 @@ test("a calculated run publishes thinking.section without tool ids", async () => assert.doesNotMatch(JSON.stringify(sections), /run-jyotish/); }); -test("composeAnswer drains leftover first-stream text and writes one body", async () => { +test("the loop's own text after the calculation is the answer, written once (BUG-1053)", async () => { + // 原值: 测试名「composeAnswer drains leftover first-stream text and writes one body」: + // 主循环在工具结果后写的正文被丢弃(drainSpoken),由 compose 另写一遍 + // 新值: 主循环在工具结果后写的正文就是回答,只写一次,不再有第二个流 + // 原因: BUG-1053 实证被丢弃的那一步才看得到盘面证据,compose 看不到; + // 产品 2026-09-27 决定删除 compose,保留主循环的正文。思考通道断言不变 const state = toolOnlyRunState(); state.thinkingPlan = natalConsultationThinkingPlan({ domains: ["career", "wealth"] }); - async function* leftover() { + async function* loop() { yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; - yield { type: "text-delta", payload: { text: "## 事业\n整篇都写了" } }; + yield { type: "reasoning-delta", payload: { text: "The user asked about career." } }; + yield { type: "text-delta", payload: { text: `## ${REPORT_HEADING.question}\n外松内紧。\n` } }; yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; } - let composed = 0; + let retries = 0; const response = streamAgentResponse({ - runId: "run", requestId: "req", state, stream: leftover(), requireTool: true, + runId: "run", requestId: "req", state, stream: loop(), requireTool: true, stepScopedAnswer: true, toolStatus: () => "ready", receipt: () => receipt(state), - composeAnswer: async () => { - composed += 1; - async function* body() { - yield { type: "reasoning-delta", payload: { text: "The user asked about career." } }; - yield { type: "text-delta", payload: { text: `## ${REPORT_HEADING.question}\n外松内紧。\n` } }; - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; - } - return body(); + retryForAnswer: async () => { + retries += 1; + throw new Error("a written answer must not be retried"); }, }); const events: unknown[] = []; const parser = createNdjsonParser((event) => events.push(event)); parser.finish(await response.text()); - assert.equal(composed, 1); + assert.equal(retries, 0); const answer = events .filter((event): event is { type: string; text: string } => (event as { type?: string }).type === "answer.delta") .map((event) => event.text) .join(""); - assert.doesNotMatch(answer, /整篇都写了/); - assert.match(answer, new RegExp(`## ${REPORT_HEADING.question}`)); + assert.equal(answer, `## ${REPORT_HEADING.question}\n外松内紧。\n`); assert.equal(events.some((event) => (event as { type?: string }).type === "think.plan"), true); assert.equal(events.filter((event) => (event as { type?: string }).type === "thinking.delta").length, 0); + assert.doesNotMatch(JSON.stringify(events), /The user asked/); + assert.equal(events.filter((event) => (event as { type?: string }).type === "run.completed").length, 1); }); -test("empty composeAnswer falls through to answer-retry once", async () => { +test("an empty answer step falls through to answer-retry once", async () => { + // 原值: 测试名「empty composeAnswer falls through to answer-retry once」:compose 返回空,再走 answer-retry + // 新值: 主循环在计算后没写正文,直接走 answer-retry 一次 + // 原因: BUG-1053 删除 compose;空回答的兜底仍是带工具的 answer-retry(只调用一次) const state = toolOnlyRunState(); state.thinkingPlan = dailyConsultationThinkingPlan(); let answerRetries = 0; - let composed = 0; async function* first() { yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; } const response = streamAgentResponse({ - runId: "run", requestId: "req", state, stream: first(), requireTool: true, + runId: "run", requestId: "req", state, stream: first(), requireTool: true, stepScopedAnswer: true, toolStatus: () => "ready", receipt: () => receipt(state), retryForAnswer: async () => { answerRetries += 1; @@ -2109,13 +2116,6 @@ test("empty composeAnswer falls through to answer-retry once", async () => { } return body(); }, - composeAnswer: async () => { - composed += 1; - async function* empty() { - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; - } - return empty(); - }, }); const events: unknown[] = []; const parser = createNdjsonParser((event) => events.push(event)); @@ -2124,7 +2124,6 @@ test("empty composeAnswer falls through to answer-retry once", async () => { .filter((event): event is { type: string; text: string } => (event as { type?: string }).type === "answer.delta") .map((event) => event.text) .join(""); - assert.equal(composed, 1); assert.equal(answerRetries, 1); assert.equal(state.steps.some((step) => step.name === "section-empty-retry"), false); assert.equal(state.steps.some((step) => step.name === "answer-retry"), true); @@ -2132,31 +2131,33 @@ test("empty composeAnswer falls through to answer-retry once", async () => { assert.equal(events.filter((event) => (event as { type?: string }).type === "run.completed").length, 1); }); -test("composeAnswer that stays empty still asks for a full answer", async () => { +test("an answer step that stays empty still asks for a full answer", async () => { + // 原值: 测试名「composeAnswer that stays empty still asks for a full answer」:compose 空 → 最后一步是 answer-retry + // 新值: 主循环的计算后一步空 → 最后一步是 answer-retry,且重试不带 Pass 4 提示 + // 原因: BUG-1053 删除 compose;空回答必须仍然要一份完整回答 const state = toolOnlyRunState(); state.thinkingPlan = dailyConsultationThinkingPlan(); let answerRetries = 0; + const hints: Array = []; async function* first() { yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; + yield { type: "step-finish", payload: { stepResult: { reason: "tool-calls" } } }; + yield { type: "text-delta", payload: { text: " " } }; + yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}, {}] } } }; } const response = streamAgentResponse({ - runId: "run", requestId: "req", state, stream: first(), requireTool: true, + runId: "run", requestId: "req", state, stream: first(), requireTool: true, stepScopedAnswer: true, + pass4Mode: "verified_chart", toolStatus: () => "ready", receipt: () => receipt(state), - retryForAnswer: async () => { + retryForAnswer: async (hint) => { answerRetries += 1; + hints.push(hint); async function* body() { yield { type: "text-delta", payload: { text: `## ${DAILY_HEADING.trend}\n补写。\n` } }; yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; } return body(); }, - composeAnswer: async () => { - async function* empty() { - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; - } - return empty(); - }, }); const events: unknown[] = []; const parser = createNdjsonParser((event) => events.push(event)); @@ -2166,32 +2167,34 @@ test("composeAnswer that stays empty still asks for a full answer", async () => .map((event) => event.text) .join(""); assert.equal(answerRetries, 1); + assert.deepEqual(hints, [undefined]); assert.equal(state.steps.some((step) => step.name === "section-empty-retry"), false); assert.equal(state.steps.at(-1)?.name, "answer-retry"); assert.match(answer, /补写/); }); -test("composeAnswer length continue finishes the same body", async () => { +test("a length-cut answer step continues the same body with the calculation result", async () => { + // 原值: 测试名「composeAnswer length continue finishes the same body」:compose 停在 length,续写补完 + // 新值: 主循环写回答的那一步停在 length,续写补完;续写拿到本轮计算结果(工具结果) + // 原因: BUG-1053 删除 compose;续写是新的流、Agent 不跨流记忆,必须把证据交给它 const state = toolOnlyRunState(); state.thinkingPlan = natalConsultationThinkingPlan({ domains: ["career"] }); - async function* first() { - yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; + const evidence = { evidence_contract: { must_use_layers: ["D10"] } }; + async function* pinched() { + yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: evidence } }; + yield { type: "step-finish", payload: { stepResult: { reason: "tool-calls" } } }; + yield { type: "text-delta", payload: { text: `## ${REPORT_HEADING.question}\n岁差` } }; + yield { type: "step-finish", payload: { stepResult: { reason: "length" } } }; + yield { type: "finish", payload: { stepResult: { reason: "length" }, output: { usage: {}, steps: [{}, {}] } } }; } let continues = 0; const response = streamAgentResponse({ - runId: "run", requestId: "req", state, stream: first(), requireTool: true, + runId: "run", requestId: "req", state, stream: pinched(), requireTool: true, stepScopedAnswer: true, toolStatus: () => "ready", receipt: () => receipt(state), - composeAnswer: async () => { - async function* pinched() { - yield { type: "text-delta", payload: { text: `## ${REPORT_HEADING.question}\n岁差` } }; - yield { type: "finish", payload: { stepResult: { reason: "length" }, output: { usage: {}, steps: [{}] } } }; - } - return pinched(); - }, - continueAfterLength: async (output) => { + continueAfterLength: async (output, received) => { continues += 1; assert.match(output, new RegExp(`## ${REPORT_HEADING.question}`)); + assert.deepEqual(received, evidence); async function* rest() { yield { type: "text-delta", payload: { text: " Lahiri。\n" } }; yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; @@ -2240,35 +2243,34 @@ test("pass4 holds verified dates and records pass4-observe without rewriting", a assert.equal(state.steps.some((step) => step.name === "pass4-observe:exact-timing"), true); }); -test("pass4 retries compose once on guarantee then drops leftover clauses", async () => { - // 原值:整段 hold,compose 两次后一次发出 - // 新值:按句放行,保证句从不出现在任何 answer.delta;已发出过句子不再整篇重写 - // 原因:BUG-950,「不闪两次」只约束已发出的不撤回。 +test("pass4 drops a guarantee clause in place without a second writing pass", async () => { + // 原值: 测试名「pass4 retries compose once on guarantee then drops leftover clauses」: + // 正文由 compose 写,断言 compose 只调用一次、保证句不外发 + // 新值: 正文由主循环写;已发出过句子时保证句就地丢弃,不再发起第二遍(answer-retry 调用 0 次) + // 原因: BUG-1053 删除 compose;BUG-950「已发出的不撤回、不闪两次」的性质原样保留 + // (history: 原值:整段 hold,compose 两次后一次发出 → 新值:按句放行 → 原因:BUG-950) const state = toolOnlyRunState(); - async function* first() { + async function* loop() { yield { type: "tool-result", payload: { toolCallId: "tool-1", toolName: "run-jyotish-consultation", result: {} } }; + yield { type: "text-delta", payload: { text: "方向可以推进。" } }; + yield { type: "text-delta", payload: { text: "我保证你一定会升职。" } }; + yield { type: "text-delta", payload: { text: "第三句照常。" } }; yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; } - let composed = 0; + let retried = 0; const response = streamAgentResponse({ - runId: "run", requestId: "req", state, stream: first(), requireTool: true, + runId: "run", requestId: "req", state, stream: loop(), requireTool: true, pass4Mode: "verified_chart", toolStatus: () => "ready", receipt: () => receipt(state), - composeAnswer: async () => { - composed += 1; - async function* body() { - yield { type: "text-delta", payload: { text: "方向可以推进。" } }; - yield { type: "text-delta", payload: { text: "我保证你一定会升职。" } }; - yield { type: "text-delta", payload: { text: "第三句照常。" } }; - yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}] } } }; - } - return body(); + retryForAnswer: async () => { + retried += 1; + throw new Error("a partly released answer must not be rewritten"); }, }); const events: unknown[] = []; const parser = createNdjsonParser((event) => events.push(event)); parser.finish(await response.text()); - assert.equal(composed, 1); + assert.equal(retried, 0); const answers = events .filter((event): event is { type: string; text: string } => (event as { type?: string }).type === "answer.delta") .map((event) => event.text); diff --git a/frontend/tests/consultation-stream-recovery.test.ts b/frontend/tests/consultation-stream-recovery.test.ts index b5ce10c6..f79586de 100644 --- a/frontend/tests/consultation-stream-recovery.test.ts +++ b/frontend/tests/consultation-stream-recovery.test.ts @@ -62,8 +62,15 @@ test("Agentic failures always refund and detached execution uses a server-owned // so the route imports it rather than restating it. assert.match(consultRoute, /import \{\n AGENT_MAX_STEPS,\n AGENT_TIMEOUT_MS,[\s\S]*?\} from "@\/mastra\/consultation-tools";/); assert.doesNotMatch(consultRoute, /const AGENT_TIMEOUT_MS =/); - assert.match(agentic, /const agentAbortSignal = AbortSignal\.timeout\(AGENT_TIMEOUT_MS\)/); - assert.equal(agentic.match(/abortSignal: agentAbortSignal/g)?.length, 3); + // 原值: const agentAbortSignal = AbortSignal.timeout(AGENT_TIMEOUT_MS),且 `abortSignal: agentAbortSignal` 三处 + // (共用 streamOptions + 两个 agentContext) + // 新值: agentAbortSignal 来自 run clock 的工具阶段(toolPhaseMs: AGENT_TIMEOUT_MS),两个 agentContext 仍用它; + // 共用 streamOptions 改用 runClock.loopSignal(拿到计算结果后交给答案时钟) + // 原因: BUG-1053 写回答发生在主循环里,循环不能再被工具阶段的 110s 掐断;超时仍由服务端持有,不用 request.signal + assert.match(agentic, /const runClock = createConsultationRunClock\(\{\s+toolPhaseMs: AGENT_TIMEOUT_MS,/); + assert.match(agentic, /const agentAbortSignal = runClock\.toolSignal;/); + assert.equal(agentic.match(/abortSignal: agentAbortSignal/g)?.length, 2); + assert.equal(agentic.match(/abortSignal: runClock\.loopSignal/g)?.length, 1); assert.doesNotMatch(agentic, /abortSignal: request\.signal/); assert.equal( agentic.match(/onError: \(error\) => settleRun\(\s*cancel,[\s\S]*?toAgentObservabilityErrorCode\(error\),\s*\)/g)?.length, diff --git a/frontend/tests/consultation-thinking-plan.test.ts b/frontend/tests/consultation-thinking-plan.test.ts index 8e0fef46..f2432aa4 100644 --- a/frontend/tests/consultation-thinking-plan.test.ts +++ b/frontend/tests/consultation-thinking-plan.test.ts @@ -3,12 +3,13 @@ import test from "node:test"; import { applyThinkingSectionProgress, - consultationComposePrompt, + consultationContinueMessages, consultationContinuePrompt, consultationReportHeadings, - consultationSectionPrompt, consultationSpokenHeadingRule, + dailyAnswerShapeInstruction, dailyConsultationThinkingPlan, + natalAnswerShapeInstruction, DAILY_HEADING, natalConsultationThinkingPlan, REPORT_HEADING, @@ -118,10 +119,18 @@ test("think plan flattens natal sections into v2 steps", () => { }); test("section prompt asks for one heading and forbids another calculation", () => { - const prompt = consultationSectionPrompt("事业", `## ${REPORT_HEADING.question}\n外松内紧。\n`); - assert.match(prompt, /只写这一个二级标题及其正文:## 事业/); - assert.match(prompt, /不要再调用排盘工具/); - assert.match(consultationComposePrompt(), /不要写「统一参数与原始结构」/); + // 原值: 锁 consultationSectionPrompt(只写一节、不要再调用排盘工具)与 consultationComposePrompt 禁写统一参数 + // 新值: 两个提示词都已删除;改锁搬进用户轮的本命写作要求:只用本轮计算结果、开场无标题、 + // 一次写完四个标题、禁写统一参数与技法审计表 + // 原因: BUG-1053 删除单独的 compose 流(分段写作早在 BUG-944 已取消,consultationSectionPrompt 是死代码); + // 写作形态要求改由主循环看到的用户轮携带。测试名保留,避免名单无说明地消失 + const natal = natalAnswerShapeInstruction(); + assert.match(natal, /只用结果里的盘面事实/); + assert.match(natal, /开场不要标题/); + for (const heading of Object.values(REPORT_HEADING)) assert.ok(natal.includes(`## ${heading}`)); + assert.match(natal, /不要写「统一参数与原始结构」/); + assert.match(natal, /不要写技法审计表/); + assert.doesNotMatch(natal, /不要再调用排盘工具/, "the loop still has to call the tool before writing"); }); test("daily thinking plan is three sections without a close heading", () => { @@ -149,7 +158,18 @@ test("daily thinking plan is three sections without a close heading", () => { }); test("empty-retry section prompt asks to write the heading directly", () => { - const prompt = consultationSectionPrompt("今日趋势", "", { reason: "empty-retry" }); - assert.match(prompt, /上一段没有输出正文,请直接写这一节/); - assert.match(prompt, /只写这一个二级标题及其正文:## 今日趋势/); + // 原值: 锁 consultationSectionPrompt 的 empty-retry 分支(「上一段没有输出正文,请直接写这一节」) + // 新值: 该提示词已删除;改锁今日入口的写作要求(三节标题、审计表与边界句进最后一节)和 + // 续写消息:带上本轮计算结果,再接已写正文与续写要求 + // 原因: BUG-1053 删除 compose 提示词链;空回答走 answer-retry,续写不得看不到证据。测试名保留 + const daily = dailyAnswerShapeInstruction(); + assert.match(daily, new RegExp(`${DAILY_HEADING.trend}、${DAILY_HEADING.actAvoid.replace(/[/]/g, "\\/")}、${DAILY_HEADING.action}`)); + assert.match(daily, /探索性日提示,不是确定预测/); + const base = [{ role: "user" as const, content: "问题" }]; + const withEvidence = consultationContinueMessages(base, `## ${REPORT_HEADING.question}\n岁差`, { moon: 3.2259 }); + assert.equal(withEvidence.length, 4); + assert.match(String(withEvidence[1]?.content), /本轮服务器计算结果[\s\S]*3\.2259/); + assert.equal(withEvidence[2]?.role, "assistant"); + assert.equal(withEvidence[3]?.content, consultationContinuePrompt(`## ${REPORT_HEADING.question}\n岁差`)); + assert.equal(consultationContinueMessages(base, "x").length, 3, "no evidence, no evidence message"); }); diff --git a/frontend/tests/consultation-workflow-contract.test.ts b/frontend/tests/consultation-workflow-contract.test.ts index bd4c5359..fe3b0198 100644 --- a/frontend/tests/consultation-workflow-contract.test.ts +++ b/frontend/tests/consultation-workflow-contract.test.ts @@ -174,12 +174,19 @@ test("the model step budget and the wall-clock budget are declared as one pair", // The pair now lives beside the domain cap it funds: the cap is derived from // the wall clock, so a change to one that forgets the other is impossible. assert.match(tools, /export const AGENT_MAX_STEPS = 8;\nexport const AGENT_TIMEOUT_MS = 110_000;/); - assert.match(tools, /export const AGENT_SLICE_MAX_STEPS = 1;/); + // 原值: 锁 AGENT_SLICE_MAX_STEPS = 1 与路由里的 maxSteps: AGENT_SLICE_MAX_STEPS + // 新值: 锁该常数已删除、路由不再出现单步写作流 + // 原因: BUG-1053 删除 compose(它是唯一的单步流),写回答发生在 AGENT_MAX_STEPS 的主循环里 + assert.doesNotMatch(tools, /AGENT_SLICE_MAX_STEPS/); assert.match(tools, /MAX_CONSULTATION_DOMAINS = Math\.max\(\s*1,\s*Math\.floor\(CONSULTATION_DOMAIN_WALL_CLOCK_MS \/ CONSULTATION_DOMAIN_DURATION_MS\),\s*\)/); assert.doesNotMatch(route, /const AGENT_(MAX_STEPS|TIMEOUT_MS) =/); assert.match(route, /maxSteps: AGENT_MAX_STEPS,/); - assert.match(route, /maxSteps: AGENT_SLICE_MAX_STEPS,/); - assert.match(route, /AbortSignal\.timeout\(AGENT_TIMEOUT_MS\)/); + assert.doesNotMatch(route, /AGENT_SLICE_MAX_STEPS/); + // 原值: assert.match(route, /AbortSignal\.timeout\(AGENT_TIMEOUT_MS\)/) + // 新值: 工具阶段的 AGENT_TIMEOUT_MS 交给 run clock(createConsultationRunClock 内部 AbortSignal.timeout) + // 原因: BUG-1053 主循环在拿到计算结果后改走答案时钟;工具仍在 110s 的 toolSignal 下 + assert.match(route, /toolPhaseMs: AGENT_TIMEOUT_MS,/); + assert.match(tools, /const toolSignal = AbortSignal\.timeout\(toolPhaseMs\);/); assert.match(route, /createConsultationRuntimeState\(\{ plannedSteps: AGENT_MAX_STEPS \}\)/); assert.doesNotMatch(route, /maxSteps: \d/); assert.doesNotMatch(route, /AbortSignal\.timeout\(\d/); @@ -202,11 +209,16 @@ test("consult streams reserve an answer budget and keep provider thinking on a s assert.match(route, /\.\.\.consultationGenerationSettings\(selectedModel\.model\)/); assert.match(route, /consultationContinueGenerationSettings\(selectedModel\.model\)/); assert.match(route, /consultationContinueGenerationSettings\(selectedModel\.model\)/); - assert.match(route, /composeAnswer,/); + // 原值: 锁 composeAnswer 接线、恰好一个 composeAnswer、续写签名 (output: string) 三处、 + // 以及 compose 的 toolChoice: "none" + // 新值: 锁路由没有 composeAnswer、没有 toolChoice: "none"(主循环开 thinking,BUG-937 不得用非 auto); + // 续写仍是三处,本命与申报时段两处接收 evidence + // 原因: BUG-1053 删除 compose;续写必须拿到计算结果 + assert.doesNotMatch(route, /composeAnswer/); assert.doesNotMatch(route, /maxOutputTokens:\s*\d/); - assert.equal(route.match(/const continueAfterLength = async \(output: string\) => \{/g)?.length, 3); - assert.equal(route.match(/const composeAnswer = async \(/g)?.length, 1); - assert.match(route, /toolChoice: "none"/); + assert.equal(route.match(/const continueAfterLength = async \(output: string(, evidence\?: unknown)?\) => \{/g)?.length, 3); + assert.equal(route.match(/const continueAfterLength = async \(output: string, evidence\?: unknown\) => \{/g)?.length, 2); + assert.doesNotMatch(route, /toolChoice: "none"/); assert.match(route, /entrypoint: consultEntrypoint/); assert.match(route, /entrypoint: parsed\.data\.entrypoint/); assert.match(tools, /export function consultationNatalPrepareStep/); @@ -235,22 +247,32 @@ test("consult streams reserve an answer budget and keep provider thinking on a s ); assert.match(natalRetry, /], natalStreamOptions\)/); assert.doesNotMatch(natalRetry, /], streamOptions\)/); - const composeBlock = route.slice( - route.indexOf("const composeAnswer = async ("), - route.indexOf("const executionReceipt = (): AgentExecutionReceipt => ({", route.indexOf("const composeAnswer = async (")), + // 原值: 锁 compose 块用 streamOptions、带 retryHint、不用 natalStreamOptions + // 新值: Pass 4 的 retryHint 改由本命 retryForAnswer 携带;该重试保留工具(natalStreamOptions,命中同请求缓存) + // 原因: BUG-1053 删除 compose;Pass 4 整篇被拒后的重写改走带工具的 answer-retry + const natalAnswerRetry = route.slice( + route.indexOf("const retryForAnswer = async (retryHint?: string) => {", route.indexOf("if (!prepared.serverChart)")), + route.indexOf("const continueAfterLength = async", route.indexOf("if (!prepared.serverChart)")), ); - assert.match(composeBlock, /\.\.\.streamOptions,/); - assert.match(composeBlock, /retryHint/); - assert.doesNotMatch(composeBlock, /natalStreamOptions/); + assert.match(natalAnswerRetry, /retryHint/); + assert.match(natalAnswerRetry, /\.\.\.natalStreamOptions, abortSignal: answerPhaseSignal\(\)/); assert.doesNotMatch(stream, /section-empty-retry/); assert.match(tools, /dailyConsultationThinkingPlan/); assert.match(tools, /pinsConsultationDomains/); assert.match(stream, /thinking\.section/); assert.match(stream, /think\.plan/); - assert.match(stream, /think\.step/); + // 原值: assert.match(stream, /think\.step/); + // 新值: 咨询流不再发 think.step;计划行由第一条 answer.delta / run.* 收口(timeline completeLiveThink) + // 原因: BUG-1053 删除 interpret/compose 阶段。think.step 只由 publishFindings 在 compose 前发出, + // 而 interpretFindings 只回 id、文本恒空,事件只是把行从 running 翻到 done;事件 schema 与客户端 reducer 保留兼容 + assert.doesNotMatch(stream, /publishFindings|interpretFindings/); assert.match(stream, /continueAfterLength/); - assert.match(stream, /composeAnswer/); - assert.match(stream, /drainSpoken/); + // 原值: assert.match(stream, /composeAnswer/); assert.match(stream, /drainSpoken/); + // 新值: 流里没有 composeAnswer / drainSpoken;回答按步取(stepScopedAnswer),续写带 calculationEvidence + // 原因: BUG-1053 删除 compose 与「丢弃主循环正文」 + assert.doesNotMatch(stream, /composeAnswer|drainSpoken/); + assert.match(stream, /stepScopedAnswer/); + assert.match(stream, /continueAfterLength\(pendingAnswer\(\), calculationEvidence\)/); assert.match(stream, /answer-continue/); });