fix(consult): write the answer in the step that saw the chart (BUG-1053)

The natal loop's step after run-jyotish-consultation saw the evidence, but
its text was drained and a second, blind compose stream (history + question
only, empty findings) wrote the user-visible answer. Remove compose,
interpret and the drain; keep the loop's own final-step text.

- stepScopedAnswer: per-step holding; text of a step that calls a tool is
  dropped, so narration around tool calls never reaches the answer
- writing shape (opener + four headings) moves into the user turn
- length continuation receives the calculation result; Pass 4 whole-answer
  reject retries through retryForAnswer with the rewrite hint
- createConsultationRunClock: tools keep the 110s tool phase; the loop is
  handed to the 70s answer clock when the calculation result arrives
- settlement judges the step that wrote the answer (BUG-1051 kept)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8
This commit is contained in:
Jesse_Chen
2026-09-27 01:01:38 +08:00
co-authored by Claude Opus 5.5
parent 03c7c0a97e
commit eef0cb7486
12 changed files with 1005 additions and 375 deletions
+58 -58
View File
@@ -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),