feat(consult): one read-only evidence lookup per turn; settle open clauses at tool calls (BUG-1059)

read-consultation-evidence returns one closed-enum section (a formal varga,
research/extended vargas, a Western layer, yogas, Ashtakavarga, Shadbala,
transits, Chara Dasha, arudha, karakas, KP, gulika, kakshya, mahadashas,
thematic evidence) from this request's finished calculation, never
recalculates, answers unavailable on a cache miss and refuses a second call.
The receipt records the step and the write row shows 「正在多看一眼:…」.
A lookup after answer text went out keeps the released text whole: a verbatim
restart is dropped as it arrives (40-char confirmation), a continuation is
kept, and settlement still reads the step that wrote the answer; the lookup
runs on the answer clock without resetting it. A length continuation carries
the lookup result with the card. BUG-1059: the visible-text transformer's open
clause is settled at each tool call, so unpunctuated narration no longer
leaks into the answer.

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 02:54:11 +08:00
co-authored by Claude Opus 5.5
parent 147ebc1789
commit cb3ee55837
9 changed files with 873 additions and 7 deletions
+15
View File
@@ -14201,3 +14201,18 @@
- 相关记录:BUG-287(同一条 allowlist 曾经把分盘和审计表整段挡住;这次是嵌套键还没放行,不是「没算」复发)
- 复发自:无
- 修复版本:`codex/consult-evidence-card-20260927`(本地提交,未推送)
## BUG-1059 | 调工具那一步没写完的半句过程说明,会粘到下一步正文开头
- 状态:resolved(本地修复 `codex/consult-evidence-card-20260927`,待部署 staging;部署后补门禁 run 号)
- 首次发现 / 最近更新:2026-09-27 / 2026-09-27
- 影响面:`frontend/src/lib/stream-agent-response.ts` `consumeAttempt`(按步取答 `stepScopedAnswer`);本命与申报时段两条带工具的咨询路径。
- 现象:数据卡实现(T5 补取工具)写测试时发现:模型在调用工具前写了一句不以句号 / 问号 / 换行收尾的过程说明(如「我先排一下盘:」「再看一眼 D60 分盘,」),这半句不会在那一步被丢掉,而是被原样接在下一步回答的第一句前面放给用户。用同一套真实 Agent + 假模型在修复前复现:回答以「我先排一下盘:」开头。
- 触发条件:带工具的咨询轮,模型在调工具那一步写的过程说明最后一句没有句末标点。以句号收尾的过程说明(BUG-1053 测试用的样例)不受影响,所以 BUG-1053 的回归没拦住。
- 根因:`createVisibleTextTransformer` 按句子边界放字,没收尾的半句留在它的缓冲里;这个缓冲按整次 attempt 共享、跨步保留。BUG-1053 的按步取答只丢掉「这一步里已经交给取答逻辑的文字」,调工具时没有把转换器里压着的半句一并结算,于是它在下一步第一次 `push` 时被冲出来,算作下一步(写回答那一步)的正文。
- 修复:遇到 `tool-call` 块时先 `visible.finish("")` 结算这一步压着的半句:计算结果到手前的交给原来的「未签约文字」缓冲(降级材料不变);这一步已经放出过正文的(回答中途补取)当作正文续上,保证已放出的句子不被截断;其余都是过程说明,随这一步丢掉。
- 验证:`frontend/tests/consult-evidence-lookup-20260927.test.ts`「BUG-1059: narration without closing punctuation before a tool call never leaks into the answer」(真实 `getJyotishAgent` + 记录提示词的假模型 + 公开 AA 盘真实引擎 golden),修复前失败、修复后通过;同文件「an open clause of released answer text before a lookup goes out whole」锁住另一面(回答中途补取时,已放出正文的半句照样完整放出)。BUG-1053 的 `consult-single-pass-answer-20260927.test.ts` 全部仍通过。
- 防复发:按步取答的回归必须包含「过程说明不以句末标点收尾」的样例;任何在步边界丢弃文字的逻辑,都要连同 `createVisibleTextTransformer` 的缓冲一起结算。
- 相关记录:BUG-1053(按步取答,本条补它漏掉的转换器缓冲)、BUG-1051(结算只认写回答那一步的 `stop`,不变)。
- 复发自:无(BUG-1053 引入按步取答时留下的缺口)。
- 修复版本:`codex/consult-evidence-card-20260927`(本地提交,未推送)
+4
View File
@@ -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,TASK-consult-evidence-card)
写回答的模型只读数据卡;界面上唯一新增的是写作行在模型补取一段卡外数据时的进行中句:「正在多看一眼:D60 分盘…」(分盘写「Dxx 分盘」,其余用中文名:格局明细、八分法(Ashtakavarga)、六力(Shadbala)、太阳回归等,见 `evidenceLookupSectionLabel`)。它复用写作行(`answer-composition` 活动),不加新行、新组件或动效;第一段正文出来时照常换成写作句。每轮最多一次。补取发生在已放出正文之后时,界面上的正文不会被截断,也不会重复出现:模型若把已写的开头原样重写一遍,重复部分在服务端丢掉(记 `answer-restart-dropped` 校验步)。一步里调工具前没写完的半句过程说明不会再粘到下一步的正文开头(BUG-1059)。
## 普通咨询一遍成文(2026-09-27,BUG-1053)
本命咨询不再有单独的「写结论」阶段:算完盘后,同一个模型循环接着就写回答。时间线的变化只有两处,没有新组件、没有新动效。一是不再出现「正在整理判断依据」「正在写结论」这两个阶段事件,也不再发 `think.step`;计划里的思考行在第一段正文出现时收口(沿用 `completeLiveThink`),写作行的进行中句统一为「正在组织回答」。二是第一段正文要等模型写出第一个二级标题或满 160 字才出现(为了把工具前后的过程说明挡在正文外),开场句会整段一起到达,之后照常按帧放出。截断、停止、失败三种收口与提示不变(BUG-1051)。
+5
View File
@@ -108,6 +108,11 @@ Jyotisha 的可见文案是产品的一部分。正确性红线(真实性、
普通报告只留下人能读的正文、结论、行动建议和必要限制。不写字段名、状态码、评分、权重或执行账本。说不到日期就写「这次说不到具体哪一天。」依据没补上就写「有些依据还没补上,相关说法不能当成确定预测。」判断没闭合就写「有些判断还没闭合,不能写成确定结论。」不出现 `technique_truth`、`workflow_route` 这类键。
## 普通咨询的领域名与「多看一眼」
- 领域中文名:父母、子女与原有的事业、关系、财富、身心压力、学习、迁居、家庭、年运、时运、综合并列;「家庭」只用于家里整体(家庭氛围、家里的事),问父母就说父母,问孩子就说子女。
- 写作行里补取卡外数据的进行中句:「正在多看一眼:<名称>…」,例如「正在多看一眼:D60 分盘…」「正在多看一眼:太阳回归…」「正在多看一眼:格局明细…」。不写「补取」「卡外」「数据卡」这些内部词,不写秒数、不写「正在加载」。只在活动行里,不进正文与历史。
## 寒暄只回一句
普通对话里的纯打招呼、道谢、告别,只回一句简体白话,最多 20 字,不以句号结尾;不套开场形状,不写星盘或运势主张,不用客服套话、星月比喻或一串追问。本命、申报时段、无出生分钟三种模式共用。含咨询、解释前文、纠错、抱怨或不确定意图时仍按咨询处理。
@@ -14,3 +14,47 @@ export function consultationWriteLabel(heading: string, live: boolean): string {
const title = heading.trim() || "回答";
return live ? `正在写${title}…` : `写${title}`;
}
const WESTERN_LAYER_LABELS: Readonly<Record<string, string>> = {
natal: "西洋本命",
transits: "西洋行运",
solar_return: "太阳回归",
secondary_progressions: "次限推运",
solar_arc_directions: "太阳弧",
converse_secondary_progressions: "逆推次限",
converse_solar_arc_directions: "逆推太阳弧",
midpoints: "中点",
lunar_return: "月亮回归",
transit_duration_scan: "行运时长扫描",
parans: "共升共落",
};
const LOOKUP_SECTION_LABELS: Readonly<Record<string, string>> = {
"varga:research_dn": "研究用分盘",
"varga:extended": "D81 / D108 / D144 分盘",
yogas: "格局明细",
ashtakavarga: "八分法(Ashtakavarga)",
shadbala: "六力(Shadbala)",
transits: "行运触发",
chara_dasha: "Chara 大运",
arudha_padas: "映点(Arudha)",
chara_karakas: "七个代表星(Chara Karaka)",
kp_cusps: "KP 宫头",
gulika: "Gulika",
kakshya: "Kakshya",
vimshottari_mahadashas: "全部大运起止",
domain_thematic_evidence: "这个主题的证据明细",
};
/** Plain words for one evidence-lookup section (TASK-consult-evidence-card-20260927 T5). */
export function evidenceLookupSectionLabel(section: string): string {
if (section.startsWith("varga:D")) return `${section.slice("varga:".length)} 分盘`;
if (section.startsWith("western:")) return WESTERN_LAYER_LABELS[section.slice("western:".length)] ?? "西洋层";
return LOOKUP_SECTION_LABELS[section] ?? "一段盘面数据";
}
/** The activity row while the answer model looks up one section outside the card. */
export function evidenceLookupActivityLabel(section: string, live = true): string {
const label = evidenceLookupSectionLabel(section);
return live ? `正在多看一眼:${label}…` : `多看了一眼:${label}`;
}
+107 -2
View File
@@ -1,5 +1,6 @@
import {
appendConsultationRuntimeStep,
CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
CONSULTATION_NATAL_CALC_TOOL_ID,
CONSULTATION_WINDOW_CALC_TOOL_ID,
type ConsultationRuntimeState,
@@ -446,6 +447,14 @@ function isCalculationTool(toolName: unknown) {
return toolName === CONSULTATION_NATAL_CALC_TOOL_ID || toolName === CONSULTATION_WINDOW_CALC_TOOL_ID;
}
/**
* How many characters a step after an evidence lookup must repeat, from the
* start of the answer already put out in this attempt, before it counts as a
* restart and the repeat is dropped. Two answers that merely open alike
* ("你这盘…") diverge well before this.
*/
export const LOOKUP_RESTART_MATCH_CHARS = 40;
function stepFinishReason(chunk: Chunk) {
const payload = chunk.payload as { stepResult?: { reason?: unknown } } | undefined;
return payload?.stepResult?.reason;
@@ -466,6 +475,9 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
// The calculation result the model saw, kept so a length continuation
// writes with the same evidence (BUG-1053). Never sent to the client.
let calculationEvidence: unknown;
// The one-shot evidence lookup's result, if the model used it; carried into
// a length continuation with the card (TASK-consult-evidence-card-20260927).
let lookupEvidence: unknown;
// Pass 4 buffers only the current open sentence. Closed sentences are
// classified and either sent whole or dropped whole. Whole-answer rewrite
// is allowed only before any answer.delta has gone out; after the first
@@ -602,6 +614,16 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
// the cut answer read as finished (BUG-1051 carried into the loop).
let stepWrote = false;
let answerStepReason: AgentModelFinishReason | undefined;
// Evidence lookup during the answer (T5): if the model had already put
// answer text out in this attempt and then called the lookup, the next
// step may start the answer over. A verbatim restart of the text already
// out is dropped as it arrives, so nothing released is repeated; anything
// that diverges within LOOKUP_RESTART_MATCH_CHARS is kept as written.
let attemptText = "";
let restartCheck = false;
let restartPos = 0;
let restartHeld = "";
let restartConfirmed = false;
const resetStep = () => {
stepText = "";
stepReleased = false;
@@ -619,6 +641,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
return;
}
uncontractedText = "";
attemptText += text;
held += text;
if (!held) return;
if (/\S/.test(held)) {
@@ -643,7 +666,58 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
if (/\S/.test(held)) emitted = true;
held = "";
};
const acceptText = async (text: string) => {
const filterRestart = (text: string) => {
let index = 0;
while (index < text.length && restartCheck) {
const char = text[index]!;
if (restartPos === 0 && !restartConfirmed && /\s/.test(char)) {
restartHeld += char;
index += 1;
continue;
}
if (restartPos < attemptText.length && char === attemptText[restartPos]) {
restartPos += 1;
index += 1;
if (!restartConfirmed) {
restartHeld += char;
if (restartPos >= LOOKUP_RESTART_MATCH_CHARS) {
restartConfirmed = true;
restartHeld = "";
appendConsultationRuntimeStep(options.state, {
kind: "validation",
name: "answer-restart-dropped",
status: "completed",
});
}
}
if (restartPos >= attemptText.length) restartCheck = false;
continue;
}
restartCheck = false;
}
// Still repeating: everything so far is held (not yet a confirmed
// restart) or dropped (confirmed).
if (restartCheck) return "";
// The check ended inside this chunk. An unconfirmed match was a
// coincidence and goes out as written; a confirmed repeat stays dropped.
const kept = `${restartConfirmed ? "" : restartHeld}${text.slice(index)}`;
restartHeld = "";
return kept;
};
const acceptText = async (raw: string) => {
const text = restartCheck ? filterRestart(raw) : raw;
await acceptAnswerText(text);
};
const flushRestartHeld = async () => {
// Only a step that has started repeating ends the check here; the
// lookup's own step ends before the step that might restart begins.
if (!restartCheck || (restartPos === 0 && !restartHeld)) return;
restartCheck = false;
const held = restartConfirmed ? "" : restartHeld;
restartHeld = "";
if (held) await acceptAnswerText(held);
};
const acceptAnswerText = async (text: string) => {
if (!options.stepScopedAnswer || !contractReady(options) || stepReleased) {
await outputText(text);
return;
@@ -659,6 +733,7 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
// The step ended (or the stream did): its held text is the answer unless
// the step called a tool.
const settleStep = async (reason?: unknown) => {
await flushRestartHeld();
const pending = stepText;
const toolStep = stepCalledTool || reason === "tool-calls";
resetStep();
@@ -699,6 +774,13 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
) {
calculationEvidence = chunk.payload?.result;
}
if (
chunk.type === "tool-result"
&& chunk.payload?.toolName === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID
&& !isToolInputRejection(chunk.payload?.result)
) {
lookupEvidence = chunk.payload?.result;
}
markAnswerPhase();
flushThinkingPlan(controller);
if (chunk.type === "step-start") {
@@ -706,9 +788,26 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
stepWrote = false;
}
if (chunk.type === "tool-call") {
// The visible-text transformer holds the step's open clause until a
// sentence boundary. It used to carry that clause into the next
// step's first text, so narration without closing punctuation
// ("我先排一下盘:") leaked into the answer (BUG-1059). Settle it with
// this step: answer text it continues goes out whole, anything else
// is narration and dropped with the step.
const openClause = visible.finish("");
if (openClause && (!contractReady(options) || stepReleased)) await acceptAnswerText(openClause);
// Whatever this step said before calling a tool is narration.
stepCalledTool = true;
stepText = "";
if (
chunk.payload?.toolName === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID
&& /\S/.test(attemptText)
) {
restartCheck = true;
restartPos = 0;
restartHeld = "";
restartConfirmed = false;
}
}
if (chunk.type === "abort" && !outcome.aborted) {
// Mastra's own abort chunk: the signal fired and the stream is about
@@ -784,7 +883,13 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) {
phase: "answer-composition",
label: heading ? consultationWriteLabel(heading, true) : "正在组织回答",
});
await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer(), calculationEvidence), {
const evidence = lookupEvidence !== undefined
&& calculationEvidence
&& typeof calculationEvidence === "object"
&& !Array.isArray(calculationEvidence)
? { ...calculationEvidence, evidence_lookup: lookupEvidence }
: calculationEvidence;
await consumeAttempt(controller, await options.continueAfterLength(pendingAnswer(), evidence), {
suppressCompositionActivity: true,
answerPhase: true,
});
+153 -3
View File
@@ -16,8 +16,13 @@ import type { TechniqueAuditRow, WorkflowReceipt } from "../lib/consultation-age
import { normalizeTechniqueAuditRows } from "../lib/consultation-technique-audit.ts";
import type { AgentModelFinishReason } from "../lib/agent-observability.ts";
import { agentGenerationSettings, AGENT_SLICE_ANSWER_OUTPUT_TOKENS, AGENT_SLICE_THINKING_OUTPUT_TOKENS } from "../lib/agent-generation-settings.ts";
import { chartCalculationProgressLabel } from "../lib/consultation-activity-labels.ts";
import { buildEvidenceCard, type EvidenceCard } from "../lib/consultation-evidence-card.ts";
import { chartCalculationProgressLabel, evidenceLookupActivityLabel } from "../lib/consultation-activity-labels.ts";
import {
buildEvidenceCard,
EVIDENCE_LOOKUP_SECTIONS,
type EvidenceCard,
type EvidenceLookupSection,
} from "../lib/consultation-evidence-card.ts";
import { pinsConsultationDomains, type ConsultationEntrypoint } from "../lib/consultation-entrypoint.ts";
import {
dailyConsultationThinkingPlan,
@@ -135,6 +140,13 @@ export function createConsultationRunClock(options: {
return Object.freeze({ toolSignal, loopSignal: loop.signal, answerSignal });
}
export const CONSULTATION_NATAL_CALC_TOOL_ID = "run-jyotish-consultation";
/**
* The one-shot, read-only lookup of a section outside the evidence card
* (TASK-consult-evidence-card-20260927 D8). It reads this request's finished
* calculation only; it never calculates.
*/
export const CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID = "read-consultation-evidence";
export const MAX_EVIDENCE_LOOKUPS_PER_TURN = 1;
export const CONSULTATION_WINDOW_CALC_TOOL_ID = "run-jyotish-window-consultation";
/**
@@ -276,6 +288,8 @@ export type ConsultationRuntimeState = {
evidenceCardChars?: number;
/** Characters of the whole tool result the model saw. */
modelVisibleChars?: number;
/** Calls to the one-shot evidence lookup this turn, including refused ones. */
evidenceLookupCallCount: number;
};
export function createConsultationRuntimeState(options: { plannedSteps?: number; reservedValidationSteps?: number } = {}): ConsultationRuntimeState {
@@ -294,6 +308,7 @@ export function createConsultationRuntimeState(options: { plannedSteps?: number;
stepBudget: { planned, reservedValidation, total: planned + reservedValidation },
stepsTruncated: false,
modelStepCount: 0,
evidenceLookupCallCount: 0,
};
// Binding is the run's first step and it costs no model step, so it is
// recorded here rather than observed from the stream. It carries no duration
@@ -845,8 +860,91 @@ export function toModelEvidenceView(full: FullDomainPlanContext, card: EvidenceC
export type ConsultationModelEvidenceView = ReturnType<typeof toModelEvidenceView>;
type EvidenceLookupCache = {
executions: readonly DomainExecution[];
};
function lookupRecord(value: unknown): Record<string, unknown> {
return value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : {};
}
function packetEvidence(execution: DomainExecution, category: string) {
const card = execution.modelOutput.claim_cards.find((item) => item.category === category);
return lookupRecord(card?.evidence);
}
function pickDefined(source: Record<string, unknown>, keys: readonly string[]) {
const out: Record<string, unknown> = {};
for (const key of keys) if (source[key] !== undefined) out[key] = source[key];
return Object.keys(out).length ? out : undefined;
}
/**
* One section from this request's calculation, copied from the model
* projection or the engine context exactly as they hold it. `undefined` means
* the calculation did not produce it; nothing is recomputed.
*/
export function readEvidenceLookupSection(
cache: EvidenceLookupCache,
section: EvidenceLookupSection,
): unknown {
const first = cache.executions[0];
if (!first) return undefined;
const natal = packetEvidence(first, "natal_foundation");
const timing = packetEvidence(first, "timing");
const spectrum = lookupRecord(natal.varga_spectrum);
const western = lookupRecord(timing.western_spectrum);
const layers = lookupRecord(first.context.local_layers);
if (section.startsWith("varga:D")) return lookupRecord(spectrum.formal)[section.slice("varga:".length)];
if (section === "varga:research_dn") return spectrum.research_dn;
if (section === "varga:extended") return spectrum.extended;
if (section === "western:natal") return pickDefined(western, ["zodiac", "house_system", "natal", "boundary"]);
if (section.startsWith("western:")) return lookupRecord(western.techniques)[section.slice("western:".length)];
switch (section) {
case "yogas":
case "arudha_padas":
case "kp_cusps":
case "gulika":
case "kakshya":
return natal[section];
case "transits":
case "chara_dasha":
return timing[section];
case "vimshottari_mahadashas":
return timing.dasha;
case "ashtakavarga":
return pickDefined(lookupRecord(layers.ashtakavarga), ["sav", "bav", "house_scores", "strongest_signs", "weakest_signs"]);
case "shadbala":
return pickDefined(
{ summary: first.context.chart.shadbala, boundary: layers.shadbala_boundary },
["summary", "boundary"],
);
case "chara_karakas": {
const table = lookupRecord(lookupRecord(layers.jaimini).chara_karakas);
return pickDefined(table, ["AK", "AmK", "BK", "MK", "PK", "GK", "DK"]);
}
case "domain_thematic_evidence": {
const byDomain: Record<string, unknown> = {};
for (const execution of cache.executions) {
const evidence = packetEvidence(execution, "domain");
if (Object.keys(evidence).length) byDomain[execution.domain] = evidence;
}
return Object.keys(byDomain).length ? byDomain : undefined;
}
default:
return undefined;
}
}
const evidenceLookupInputSchema = z.object({
section: z.enum(EVIDENCE_LOOKUP_SECTIONS),
}).strict();
export function createConsultationTools(ctx: ConsultationAgentContext) {
let calculation: Promise<ConsultationModelEvidenceView> | null = null;
// Request-scoped: filled when this request's calculation finishes, read by
// the lookup tool, dropped with the request.
let lookupCache: EvidenceLookupCache | null = null;
const consultationTool = createTool({
id: "run-jyotish-consultation",
description: `Run one server-validated plan of personal Jyotish consultation domains. Send only question and domains. The single ordered domains array is the only way to select domains: list every domain the question needs, in priority order, up to ${MAX_CONSULTATION_DOMAIN_PLAN_VALUES}. Do not drop a relevant domain to shorten the plan. Use only the ids enumerated in the schema; workflow or checklist names from the skill's methodology are not domain ids. Questions about one's parents (父母) use parents and questions about one's children (子女) use children; family is only for the household as a whole (家里、家庭氛围). Domains execute one after another and each costs about ${Math.round(CONSULTATION_DOMAIN_DURATION_MS / 1000)}s of the run's wall clock; the server executes as many as that clock can pay for (about ${MAX_CONSULTATION_DOMAINS}) and returns the rest in omitted_domains. Birth data is server-bound and must never be supplied. The result always carries one top-level answer contract—status, evidence_contract, claim_cards, rectification—which for several domains is the most restrictive merge of the executed ones, with each domain's own contract in consultations. claim_cards is the evidence card: the base chart (ascendant, house signs, placements, functional benefics/malefics, the running Vimshottari and Narayana periods) and a section per domain; evidence_card lists what else this calculation holds. One calculation is executed per request and reused, so repeating the call with different parameters cannot change the result.`,
@@ -939,6 +1037,7 @@ export function createConsultationTools(ctx: ConsultationAgentContext) {
ctx.state.consultationToolSuccessCount += 1;
appendConsultationRuntimeStep(ctx.state, { kind: "tool", name: "run-jyotish-consultation", status: "completed", durationMs: ctx.state.consultationToolDurationMs });
const modelContext = toModelDomainPlanContext(executions, omittedDomains);
lookupCache = { executions };
const card = buildEvidenceCard(executions.map((execution) => ({
domain: execution.domain,
packet: execution.modelOutput,
@@ -986,7 +1085,58 @@ export function createConsultationTools(ctx: ConsultationAgentContext) {
}
},
});
return { "run-jyotish-consultation": consultationTool };
const lookupTool = createTool({
id: CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
description: `Read one section of this request's finished chart calculation that is not on the evidence card, for example a varga the card does not carry (varga:D60), a Western layer (western:solar_return), or the yoga details. It never calculates: it returns what run-jyotish-consultation already computed, or status unavailable. At most ${MAX_EVIDENCE_LOOKUPS_PER_TURN} call per turn; a second call is refused. Call it before writing any answer text, and only when the question needs that section; the answer contract and card come from run-jyotish-consultation.`,
inputSchema: evidenceLookupInputSchema,
execute: async (input, context) => {
const startedAt = Date.now();
ctx.state.evidenceLookupCallCount += 1;
const note = "If you had already started writing the answer, continue from where you stopped; never repeat text already written.";
if (ctx.state.evidenceLookupCallCount > MAX_EVIDENCE_LOOKUPS_PER_TURN) {
if (ctx.state.evidenceLookupCallCount === MAX_EVIDENCE_LOOKUPS_PER_TURN + 1) {
appendConsultationRuntimeStep(ctx.state, {
kind: "tool",
name: CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
status: "failed",
durationMs: 0,
failureCode: "lookup_limit_reached",
});
}
return { status: "refused" as const, reason: "lookup_limit_reached", section: input.section, note };
}
if (!lookupCache) {
appendConsultationRuntimeStep(ctx.state, {
kind: "tool",
name: CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
status: "failed",
durationMs: Math.max(0, Date.now() - startedAt),
failureCode: "calculation_not_cached",
});
return { status: "unavailable" as const, reason: "calculation_not_in_request_cache", section: input.section, note };
}
await (context as { writer?: { custom?: (value: unknown) => unknown } } | undefined)?.writer?.custom?.({
type: "data-jyotish-activity",
data: { phase: "answer-composition", label: evidenceLookupActivityLabel(input.section) },
});
const data = readEvidenceLookupSection(lookupCache, input.section);
const found = data !== undefined && data !== null;
appendConsultationRuntimeStep(ctx.state, {
kind: "tool",
name: CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
status: found ? "completed" : "failed",
durationMs: Math.max(0, Date.now() - startedAt),
...(found ? {} : { failureCode: "section_not_computed" }),
});
return found
? { status: "ok" as const, section: input.section, data, note }
: { status: "unavailable" as const, reason: "section_not_computed", section: input.section, note };
},
});
return {
"run-jyotish-consultation": consultationTool,
[CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID]: lookupTool,
};
}
async function executeWindowConsultation(
@@ -0,0 +1,316 @@
// TASK-consult-evidence-card-20260927 red line 4 and T5.
//
// Red line 4 (BUG-1053 stays fixed with the card): the model call that writes
// the answer has the evidence card in its prompt; narration before a tool call
// is never answer text; a non-`stop` answer step is `answer_truncated` and not
// charged; a length continuation carries the card.
//
// T5 (D8): one read-only lookup of a section outside the card, from this
// request's calculation only, at most once per turn. A lookup during the
// answer phase must not cut or repeat released text and must not reset the
// 110 s tool / 70 s answer clocks.
//
// Real `getJyotishAgent` + prompt-recording fake model; the calculation is the
// golden engine capture of a public AA chart (AGENTS §7.4).
import assert from "node:assert/strict";
import { readFileSync } from "node:fs";
import test from "node:test";
import { REPORT_HEADING } from "../src/lib/consultation-thinking-plan.ts";
import { LOOKUP_RESTART_MATCH_CHARS } from "../src/lib/stream-agent-response.ts";
import { evidenceLookupActivityLabel } from "../src/lib/consultation-activity-labels.ts";
import {
CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID,
consultationNatalPrepareStep,
createConsultationRuntimeState,
createConsultationTools,
MAX_EVIDENCE_LOOKUPS_PER_TURN,
} from "../src/mastra/consultation-tools.ts";
import {
consultationWorkflowResponseSchema,
toAgentConsultationContext,
toModelOutput,
} from "../src/mastra/consultation-workflow.ts";
import {
firstPromptWithToolResult,
pieces,
publicServerChart,
runNatalAgent,
type Turn,
} from "./consult-natal-agent-test-support.ts";
type Json = Record<string, unknown>;
const golden = JSON.parse(readFileSync(
new URL("./fixtures/consult-evidence-card-golden.json", import.meta.url),
"utf8",
)) as { charts: Array<{ id: string; workflow: Json }> };
const workflow = golden.charts.find((chart) => chart.id === "steve_jobs")!.workflow;
const modules = (workflow.chart as Json).modules as Json;
// Facts only the calculation carries, read from the engine payload.
const PD_START = String(((modules.dasha_sub_periods as Json).pratyantar_dasha_timeline as Json & { current: Json }).current.start);
const NARAYANA_SIGN = String(((modules.narayana_dasha as Json).current_dasha as Json & { md: Json }).md.sign);
const D60_LAGNA = String((((modules.varga_spectrum as Json).formal as Json).D60 as Json).lagna);
const PRE_TOOL = "我先排一下盘,稍等。";
const LOOKUP_NARRATION = "我再看一眼 D60。";
const OPENER = "你和父母这条线,表面是各过各的,底下一直在较劲:你要自己说了算,他们要你稳。这不是谁对谁错,是同一件事的两种护法,所以叫「隔空守护」。";
const ANSWER = [
OPENER,
"",
`## ${REPORT_HEADING.question}`,
"关系能缓,但先换说话方式,再谈大事。",
"",
`## ${REPORT_HEADING.support}`,
"四宫的主星落在十宫,家里的事常被你当成要办成的事来处理。",
"",
`## ${REPORT_HEADING.timing}`,
"这几个月先试着每周打一次电话,不急着谈结论。",
"",
`## ${REPORT_HEADING.action}`,
"- **这周**:约他们吃一顿饭,只聊近况。",
"",
].join("\n");
// The part of the answer written before a mid-answer lookup: past the first
// heading, so it has already been released to the client.
const HEAD_CHARS = 120;
const calcCall = (text?: string): Turn => ({
parts: [
...(text ? [{ text }] : []),
{ tool: { name: "run-jyotish-consultation", input: { question: "我和父母关系如何", domains: ["parents"] } } },
],
finish: "tool-calls",
});
const lookupCall = (section: string, text?: string): Turn => ({
parts: [
...(text ? pieces(text) : []),
{ tool: { name: CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID, input: { section } } },
],
finish: "tool-calls",
});
const run = (turns: Turn[], options: { toolPhaseMs?: number; answerMs?: number } = {}) => runNatalAgent(turns, { workflow, ...options });
function toolResultText(prompt: unknown[] | undefined, toolName: string) {
return JSON.stringify((prompt ?? []).filter((message) => {
const value = message as { role?: string; content?: unknown };
return value.role === "tool" && JSON.stringify(value.content).includes(toolName);
}));
}
test("the answer-writing call sees the evidence card, not the full result (red line 4)", async () => {
const result = await run([calcCall(PRE_TOOL), { parts: pieces(ANSWER), finish: "stop" }]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.calls, 2, "no second, blind writer");
const writer = firstPromptWithToolResult(result.prompts, "run-jyotish-consultation");
assert.equal(writer, 1);
const card = toolResultText(result.prompts[writer], "run-jyotish-consultation");
assert.ok(card.includes("evidence-card-v1"), "the card reaches the writer");
assert.ok(card.includes(PD_START), "with the engine's pratyantar start");
assert.ok(card.includes(NARAYANA_SIGN), "and the running Narayana sign");
assert.ok(card.includes("can_answer_direction"), "and the answer contract");
for (const absent of ["technique_audit_table", "western_spectrum", "research_dn"]) {
assert.doesNotMatch(card, new RegExp(`${absent}\\\\*"\\s*:`), `${absent} stays server-side`);
}
assert.equal(result.answer, ANSWER);
assert.equal(result.answer.includes(PRE_TOOL), false, "pre-tool narration is not the answer");
assert.equal(result.charges, 1);
});
test("a non-stop answer step with the card is still truncated and not charged (red line 4)", async () => {
const result = await run([calcCall(), { parts: pieces(ANSWER.slice(0, 200)), finish: "content-filter" }]);
assert.deepEqual(result.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]);
assert.equal(result.charges, 0);
});
test("a length continuation carries the card (red line 4)", async () => {
const cut = ANSWER.indexOf(`## ${REPORT_HEADING.timing}`);
const result = await run([
calcCall(),
{ parts: pieces(ANSWER.slice(0, cut)), finish: "length" },
{ parts: pieces(ANSWER.slice(cut)), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.continuations.length, 1);
const continuation = JSON.stringify(result.prompts[2]);
assert.ok(continuation.includes("evidence-card-v1"));
assert.ok(continuation.includes(PD_START));
assert.equal(result.completed, ANSWER);
});
test("step 0 still exposes only the calculation tool; later steps may use the lookup", () => {
assert.deepEqual(consultationNatalPrepareStep({ stepNumber: 0 }).activeTools, ["run-jyotish-consultation"]);
assert.equal("activeTools" in consultationNatalPrepareStep({ stepNumber: 1 }), false);
assert.equal(MAX_EVIDENCE_LOOKUPS_PER_TURN, 1);
});
test("a lookup before the answer returns the section from the request cache and the answer is written once", async () => {
const result = await run([
calcCall(),
lookupCall("varga:D60", LOOKUP_NARRATION),
{ parts: pieces(ANSWER), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.workflowRuns, 1, "the lookup never recalculates");
assert.equal(result.answer, ANSWER);
assert.equal(result.answer.includes(LOOKUP_NARRATION), false, "narration around the lookup is not answer text");
const writer = firstPromptWithToolResult(result.prompts, CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID);
assert.equal(writer, 2);
const lookup = toolResultText(result.prompts[writer], CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID);
assert.ok(lookup.includes(D60_LAGNA), "the D60 lagna the engine computed");
assert.match(lookup, /status\\*"\s*:\s*\\*"ok/);
// Receipt step and the live activity row in plain words.
const step = result.state.steps.find((item) => item.name === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID);
assert.equal(step?.kind, "tool");
assert.equal(step?.status, "completed");
assert.ok(result.events.some((event) => event.type === "activity" && event.label === evidenceLookupActivityLabel("varga:D60")));
assert.equal(evidenceLookupActivityLabel("varga:D60"), "正在多看一眼:D60 分盘…");
assert.equal(result.state.evidenceLookupCallCount, 1);
assert.equal(result.charges, 1);
});
test("a second lookup in the same turn is refused", async () => {
const result = await run([
calcCall(),
lookupCall("varga:D60"),
lookupCall("western:solar_return"),
{ parts: pieces(ANSWER), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
const second = toolResultText(result.prompts[3], CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID);
assert.ok(second.includes("lookup_limit_reached"));
assert.equal(result.state.evidenceLookupCallCount, 2);
assert.deepEqual(
result.state.steps.filter((item) => item.name === CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID).map((item) => item.status),
["completed", "failed"],
);
assert.equal(result.answer, ANSWER);
});
test("a lookup with no calculation in the request cache returns unavailable and calculates nothing", async () => {
let workflowRuns = 0;
const state = createConsultationRuntimeState();
const tools = createConsultationTools({
userId: "u", sessionId: "s", requestId: "lookup-miss", consultationMode: "verified_chart",
serverChart: publicServerChart as never, state,
runWorkflow: async () => {
workflowRuns += 1;
return structuredClone(workflow) as never;
},
});
const miss = await tools[CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID].execute!({ section: "varga:D60" } as never, {} as never) as Json;
assert.equal(miss.status, "unavailable");
assert.equal(miss.reason, "calculation_not_in_request_cache");
assert.equal(workflowRuns, 0);
assert.equal(state.steps.at(-1)?.name, CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID);
assert.equal(state.steps.at(-1)?.status, "failed");
});
test("the lookup returns exactly the projected section it names", async () => {
const state = createConsultationRuntimeState();
const tools = createConsultationTools({
userId: "u", sessionId: "s", requestId: "lookup-hit", consultationMode: "verified_chart",
serverChart: publicServerChart as never, state,
runWorkflow: async () => structuredClone(workflow) as never,
});
await tools["run-jyotish-consultation"].execute!({ question: "q", domains: ["parents"] } as never, { writer: { custom: async () => {} } } as never);
const hit = await tools[CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID].execute!({ section: "varga:D60" } as never, { writer: { custom: async () => {} } } as never) as Json;
const packet = toModelOutput(toAgentConsultationContext(consultationWorkflowResponseSchema.parse(structuredClone(workflow))));
const natal = packet.claim_cards.find((card) => card.category === "natal_foundation")!.evidence as Json;
assert.equal(hit.status, "ok");
assert.deepEqual(hit.data, ((natal.varga_spectrum as Json).formal as Json).D60);
});
test("a lookup after answer text went out: a verbatim restart is dropped, the answer is whole and settles on stop", async () => {
const head = ANSWER.slice(0, HEAD_CHARS);
const result = await run([
calcCall(),
lookupCall("yogas", head),
// The model starts the whole answer over after the lookup.
{ parts: pieces(ANSWER), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.answer, ANSWER, "released text is neither cut nor repeated");
assert.equal(result.completed, ANSWER);
assert.equal(result.charges, 1);
assert.ok(result.state.steps.some((step) => step.name === "answer-restart-dropped"));
assert.equal(result.state.composeFinishReason, "stop");
});
test("a lookup after answer text went out: a continuation is kept as written", async () => {
const head = ANSWER.slice(0, HEAD_CHARS);
const result = await run([
calcCall(),
lookupCall("yogas", head),
{ parts: pieces(ANSWER.slice(HEAD_CHARS)), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.answer, ANSWER);
assert.equal(result.state.steps.some((step) => step.name === "answer-restart-dropped"), false);
});
test("a step that only opens like the released text is not mistaken for a restart", async () => {
const head = ANSWER.slice(0, HEAD_CHARS);
const shared = OPENER.slice(0, LOOKUP_RESTART_MATCH_CHARS - 10);
const tail = `${shared}——补一句:格局明细里没有新的东西。`;
const result = await run([
calcCall(),
lookupCall("yogas", head),
{ parts: pieces(tail), finish: "stop" },
]);
assert.equal(result.answer, `${head}${tail}`);
assert.equal(result.state.steps.some((step) => step.name === "answer-restart-dropped"), false);
});
test("a lookup does not reset the answer clock: the same 70 s signal bounds the whole answer phase", async () => {
const result = await run([
calcCall(),
{ ...lookupCall("varga:D60"), delayMs: 40 },
{ parts: pieces(ANSWER, 6), finish: "stop", delayMs: 25 },
], { answerMs: 300 });
assert.deepEqual(result.terminal.map((event) => `${event.type}:${event.code ?? ""}`), ["run.failed:answer_truncated"]);
assert.equal(result.charges, 0);
assert.equal(result.answerSignals.length, 1, "the answer clock started once");
assert.ok(result.state.steps.some((step) => step.kind === "abort" && step.name === "compose-abort"));
});
test("a lookup in the answer phase is not cut by the tool phase's deadline", async () => {
const result = await run([
calcCall(),
{ ...lookupCall("varga:D60"), delayMs: 30 },
{ parts: pieces(ANSWER, 6), finish: "stop", delayMs: 12 },
], { toolPhaseMs: 150, answerMs: 5_000 });
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.clock.toolSignal.aborted, true, "the tool phase did expire");
assert.equal(result.answer, ANSWER);
});
test("BUG-1059: narration without closing punctuation before a tool call never leaks into the answer", async () => {
// The visible-text transformer held the open clause ("…:") across the tool
// call and emitted it glued to the next step's first sentence.
const openCalc = "我先排一下盘:";
const openLookup = "再看一眼 D60 分盘,";
const result = await run([
calcCall(openCalc),
lookupCall("varga:D60", openLookup),
{ parts: pieces(ANSWER), finish: "stop" },
]);
assert.deepEqual(result.terminal.map((event) => event.type), ["run.completed"]);
assert.equal(result.answer, ANSWER);
assert.equal(result.answer.includes("排一下盘"), false);
assert.equal(result.answer.includes("再看一眼"), false);
});
test("BUG-1059: an open clause of released answer text before a lookup goes out whole, not cut", async () => {
// The head ends mid-sentence; that fragment is answer text, not narration.
const head = ANSWER.slice(0, HEAD_CHARS);
assert.doesNotMatch(head.slice(-1), /[。!?.!?\n]/, "the head really ends mid-clause");
const result = await run([
calcCall(),
lookupCall("yogas", head),
{ parts: pieces(ANSWER.slice(HEAD_CHARS)), finish: "stop" },
]);
assert.equal(result.answer, ANSWER);
});
@@ -0,0 +1,220 @@
// Shared harness for natal-route tests that drive the real personal Agent
// (`getJyotishAgent`, real skill binding, real calculation and lookup tools,
// real prepareStep, real run clock) over a prompt-recording fake model.
// Same wiring as consult-single-pass-answer-20260927.test.ts (BUG-1053); the
// workflow the calculation tool receives is a golden engine capture.
import { getJyotishAgent } from "../src/mastra/index.ts";
import { createNdjsonParser } from "../src/lib/consultation-agent-events.ts";
import {
consultationContinueMessages,
natalAnswerShapeInstruction,
} from "../src/lib/consultation-thinking-plan.ts";
import { streamAgentResponse } from "../src/lib/stream-agent-response.ts";
import type { ConsultationDomain } from "../src/lib/consultation-domain-registry.ts";
import {
AGENT_MAX_STEPS,
AGENT_TIMEOUT_MS,
CONSULTATION_ANSWER_TIMEOUT_MS,
consultationNatalPrepareStep,
consultationStepBudgetReceipt,
createConsultationAgentContext,
createConsultationRunClock,
createConsultationRuntimeState,
publicConsultationRuntimeSteps,
} from "../src/mastra/consultation-tools.ts";
export type Part = { text?: string; tool?: { name: string; input: Record<string, unknown> } };
export type Turn = { parts: Part[]; finish: string; delayMs?: number };
/** A LanguageModelV2 that plays `turns` in order and records every prompt. */
export function scriptedModel(turns: Turn[]) {
const prompts: unknown[][] = [];
let call = 0;
const model = {
specificationVersion: "v2",
provider: "fake",
modelId: "fake-natal-agent",
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<void>((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 };
}
export function pieces(text: string, size = 12) {
return (text.match(new RegExp(`[\\s\\S]{1,${size}}`, "g")) ?? []).map((value) => ({ text: value }));
}
export const publicServerChart = {
name: "public",
toolInput: {
year: 1955, month: 2, day: 24, hour: 19, minute: 15, city: "San Francisco", lat: 37.77, lon: -122.42, tz: -8,
ayanamsa: "raman" as const, declared_accuracy: "minute" as const, time_source: "aa_rated",
},
truth: {
birthDate: "1955-02-24", reportedBirthTime: "19:15", activeBirthTime: null,
selectedTimeKind: "reported" as const, birthTimeSource: "reported", birthTimeStatus: "reported",
placeLabel: "San Francisco", placeCodes: { countryCode: "US", provinceCode: null, cityCode: null, districtCode: null },
placeId: null, placeType: "city", placeProvider: "profile", latitude: 37.77, longitude: -122.42,
timezoneId: "America/Los_Angeles", timezoneSource: "profile", timezoneOffset: -8,
},
};
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`.
*/
export async function runNatalAgent(
turns: Turn[],
options: {
workflow: Record<string, unknown>;
theme?: ConsultationDomain;
question?: string;
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 answerSignals: AbortSignal[] = [];
const agentContext = createConsultationAgentContext({
userId: "u", sessionId: "s", requestId: `natal-agent-${seq}`, consultationMode: "verified_chart",
theme: options.theme ?? "parents", serverChart: publicServerChart as never, abortSignal: clock.toolSignal, state,
runWorkflow: async () => {
workflowRuns += 1;
return structuredClone(options.workflow) as never;
},
});
const agent = getJyotishAgent({ id: `fake-${seq}`, model } as never, agentContext);
const question = options.question ?? "我和父母关系如何";
const baseMessages = [{ role: "user" as const, content: `${natalAnswerShapeInstruction()}\n问题:${question}` }];
const streamOptions = { runId: `run-${seq}`, maxSteps: AGENT_MAX_STEPS, abortSignal: clock.loopSignal };
const natalStreamOptions = { ...streamOptions, prepareStep: consultationNatalPrepareStep };
const continuations: unknown[] = [];
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<unknown>,
requireTool: true,
stepScopedAnswer: true,
onAnswerPhase: () => { answerSignals.push(clock.answerSignal()); },
pass4Mode: "verified_chart",
retryForAnswer: async (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<unknown>;
},
continueAfterLength: async (output, evidence) => {
continuations.push(evidence);
answerSignals.push(clock.answerSignal());
const continued = await agent.stream(
consultationContinueMessages(baseMessages, output, evidence) as never,
{ ...streamOptions, abortSignal: clock.answerSignal() } as never,
);
return continued.fullStream as ReadableStream<unknown>;
},
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; label?: string; phase?: 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, workflowRuns, clock, answerSignals,
completed: completed as string | null, charges, errored,
};
}
/** Index of the first prompt that already contains a tool result for `toolName`. */
export function firstPromptWithToolResult(prompts: unknown[][], toolName: string) {
return prompts.findIndex((prompt) => prompt.some((message) => {
const value = message as { role?: string; content?: unknown };
return value.role === "tool" && JSON.stringify(value.content).includes(toolName);
}));
}
@@ -275,7 +275,11 @@ test("consult streams reserve an answer budget and keep provider thinking on a s
// 原因: BUG-1053 删除 compose 与「丢弃主循环正文」
assert.doesNotMatch(stream, /composeAnswer|drainSpoken/);
assert.match(stream, /stepScopedAnswer/);
assert.match(stream, /continueAfterLength\(pendingAnswer\(\), calculationEvidence\)/);
// 原值: /continueAfterLength\(pendingAnswer\(\), calculationEvidence\)/
// 新值: 续写带 evidence = 计算结果(数据卡),模型本轮用过补取时再并上 evidence_lookup
// 原因: TASK-consult-evidence-card-20260927 T5,续写不能丢掉补取到的那一段
assert.match(stream, /continueAfterLength\(pendingAnswer\(\), evidence\)/);
assert.match(stream, /\{ \.\.\.calculationEvidence, evidence_lookup: lookupEvidence \}/);
assert.match(stream, /answer-continue/);
});
@@ -299,7 +303,10 @@ test("personal consultation lets the Agent invoke the server-bound workflow tool
assert.doesNotMatch(agenticBranch, /await runConsultationWorkflow/);
assert.doesNotMatch(agenticBranch, /JSON\.stringify\(toolInput\)/);
assert.match(tools, /\(ctx\.runWorkflow \?\? runConsultationWorkflow\)\(toolInput, \{/);
assert.match(tools, /return \{ "run-jyotish-consultation": consultationTool \};/);
// 原值: /return \{ "run-jyotish-consultation": consultationTool \};/
// 新值: 同一个工厂再返回只读补取工具(同请求缓存,每轮一次)
// 原因: TASK-consult-evidence-card-20260927 T5(D8)
assert.match(tools, /return \{\s*"run-jyotish-consultation": consultationTool,\s*\[CONSULTATION_EVIDENCE_LOOKUP_TOOL_ID\]: lookupTool,\s*\};/);
// 原值:transformText 读 workflowReceipt.preciseTiming === "allowed" 再套恒等壳
// 新值:Pass 4 只看 consultationMode,不再读 preciseTiming 开关去挖日期
// 原因:BUG-948,恒等壳下线;日期观察按模式,不按分钟敏感开关。