Files
Jyotisha/frontend/tests/consultation-stream-recovery.test.ts
T
Jesse_ChenandClaude Opus 5.5 eef0cb7486 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
2026-09-27 01:01:38 +08:00

185 lines
12 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import assert from "node:assert/strict";
import { readFileSync } from "node:fs";
import test from "node:test";
const read = (path: string) => readFileSync(new URL(`../${path}`, import.meta.url), "utf8");
const consultRoute = read("src/app/api/consult/route.ts");
const statusRoute = read("src/app/api/consult/status/route.ts");
const migration = read("supabase/migrations/20260808030000_consultation_stream_recovery.sql");
const statusReadMigration = read("supabase/migrations/20260811010000_consultation_status_service_role_read.sql");
test("reserves usage and binds the owned consultation session atomically", () => {
assert.match(migration, /add column if not exists session_id uuid references public\.chat_sessions\(id\) on delete set null/i);
assert.match(migration, /reserve_consultation_usage\([\s\S]*p_session_id uuid[\s\S]*from public\.chat_sessions as session[\s\S]*session\.id = p_session_id[\s\S]*session\.user_id = p_user_id[\s\S]*session\.session_type = 'consultation'/i);
assert.match(migration, /from public\.authorize_usage[\s\S]*insert into public\.consultation_requests\(user_id, request_id, session_id, status\)[\s\S]*values \(p_user_id, btrim\(p_request_id\), p_session_id, 'reserved'\)/i);
assert.match(migration, /reserve_consultation_usage\(uuid, text, uuid, text, integer\)/i);
assert.match(consultRoute, /rpc\("reserve_consultation_usage"[\s\S]*p_session_id: sessionId/);
assert.doesNotMatch(consultRoute, /from\("consultation_requests"\)[\s\S]*\.update\(\{ session_id: sessionId/);
});
test("persists transformed assistant metadata before atomically settling usage", () => {
assert.equal(consultRoute.match(/continueAfterDisconnect: true/g)?.length, 5);
assert.equal(consultRoute.match(/onComplete: \(rawTransformedText\) => settle\(\(\) => completeResponse\(/g)?.length, 2);
assert.match(
consultRoute,
/parseAgentReply\([\s\S]*?rawTransformedText,[\s\S]*?createConsultationReplyMetadata\(\{ question: visibleQuestion \}\),[\s\S]*?\)/,
);
assert.match(consultRoute, /role: "assistant" as const,[\s\S]*thinkingText: persistedThinking[\s\S]*techniqueTruth,[\s\S]*workflowReceipt/);
const append = migration.indexOf("set messages = session.messages || jsonb_build_array(p_response_message)");
const store = migration.indexOf("set response_message = p_response_message");
const settle = migration.indexOf("from public.complete_usage");
const complete = migration.indexOf("set status = 'completed'");
assert.ok(append >= 0 && append < store && store < settle && settle < complete);
});
test("persists partial transformed output when the upstream stream errors", () => {
assert.doesNotMatch(consultRoute, /completeInterrupted|completeUsage\(accounting, userId, requestId/);
assert.equal(
consultRoute.match(/onError: \(_error, emitted, output: string\) => settleErrored\(emitted, output\)/g)?.length,
2,
);
assert.equal(consultRoute.match(/const settleErrored = \(emitted: boolean, output: string\) => settle\(/g)?.length, 2);
assert.equal(
consultRoute.match(/emitted\s*\n\s*\? \(\) => completeResponse\(\n\s*output,\n\s*result\.totalUsage,/g)?.length,
2,
);
assert.equal(consultRoute.match(/onCancel: \(\) => settle\(cancel\)/g)?.length, 2);
assert.equal(
consultRoute.match(/onCancel: \(\) => settleRun\(cancel, "cancelled"\)/g)?.length,
3,
);
});
test("Agentic failures always refund and detached execution uses a server-owned timeout", () => {
const agentic = consultRoute.slice(
consultRoute.indexOf("async function runAgenticConsultation("),
// Former value: `const { history } = parsed.data` marked the end of the
// agentic function. History is now copied from stored session messages.
consultRoute.indexOf(" try {\n const history = storedHistory;"),
);
// The value lives beside the model step budget and the domain cap it funds,
// 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 =/);
// 原值: 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,
3,
);
const onErrorBlocks = agentic.match(/onError:[\s\S]*?onCancel:/g) ?? [];
assert.equal(onErrorBlocks.length, 3);
for (const block of onErrorBlocks) assert.doesNotMatch(block, /completeResponse|completed_partial/);
});
test("Agentic settlement telemetry uses the actual cancellation result", () => {
const cancellation = consultRoute.slice(
consultRoute.indexOf("async function cancel("),
consultRoute.indexOf("const usageStartedAt"),
);
const settlement = consultRoute.slice(
consultRoute.indexOf("let settlement:"),
consultRoute.indexOf("const baseMessages"),
);
assert.match(cancellation, /return "cancelled";/);
assert.match(cancellation, /catch \(error\) \{[\s\S]*return "failed";/);
assert.match(settlement, /const actualSettlementResult = await settleResult\(action\)/);
assert.match(settlement, /settlementTelemetryOutcome\(actualSettlementResult, errorCode\)/);
assert.doesNotMatch(settlement, /logRun\(errorCode, "cancelled"\)/);
});
test("best-effort cancels a failed or uncertain durable completion before rethrowing", () => {
const completion = consultRoute.slice(
consultRoute.indexOf("async function completeResponse("),
consultRoute.indexOf("let settlement:"),
);
assert.match(completion, /try \{[\s\S]*accounting\.rpc\("complete_consultation_response"/);
assert.match(completion, /catch \(error\) \{[\s\S]*await cancel\(\);[\s\S]*throw error;/);
assert.ok(completion.indexOf("await cancel();") > completion.indexOf('accounting.rpc("complete_consultation_response"'));
});
test("serializes explicit cancellation against response persistence and charging", () => {
const cancelFunction = migration.slice(
migration.indexOf("create or replace function public.cancel_consultation_credit"),
migration.indexOf("create or replace function public.complete_consultation_response"),
);
assert.match(cancelFunction, /pg_advisory_xact_lock\(hashtextextended\(p_user_id::text \|\| ':' \|\| btrim\(p_request_id\), 0\)\)/);
assert.match(cancelFunction, /from public\.consultation_requests[\s\S]*for update/);
assert.match(cancelFunction, /if v_request\.status = 'completed' then[\s\S]*'request_completed'/i);
assert.match(cancelFunction, /if v_request\.status = 'cancelled' then[\s\S]*select true/i);
assert.match(cancelFunction, /if v_request\.status <> 'reserved' then[\s\S]*from public\.release_usage[\s\S]*set status = 'cancelled'/i);
assert.match(cancelFunction, /grant execute on function public\.cancel_consultation_credit\(uuid, text\)[\s\S]*to service_role/);
assert.match(migration, /if v_request\.status = 'cancelled' then[\s\S]*'request_cancelled'/i);
assert.match(consultRoute, /cancel_consultation_credit/);
assert.match(consultRoute, /error_code !== "request_cancelled"/);
});
test("status endpoint supports one global reserved lookup and strict bound polling", () => {
assert.ok(statusRoute.indexOf("supabase.auth.getUser()") < statusRoute.indexOf("createAdminSupabaseClient()"));
assert.match(statusRoute, /activeLookup = sessionIdValue === null && requestIdValue === null/);
assert.match(statusRoute, /activeLookup[\s\S]*\.eq\("status", "reserved"\)\.order\("created_at", \{ ascending: false \}\)\.limit\(1\)/);
assert.match(statusRoute, /\.eq\("user_id", user\.id\)[\s\S]*\.eq\("session_id", sessionId!\.data\)[\s\S]*\.eq\("request_id", requestId!\.data\)/);
assert.match(statusRoute, /requestId: statusData\.request_id,[\s\S]*sessionId: statusData\.session_id,[\s\S]*status: statusData\.status/);
assert.doesNotMatch(statusRoute, /String\(data\.response_message\)|responseMessage:\s*data\.response_message as string/);
});
test("status polling grants consultation request reads only to the server role", () => {
assert.match(statusReadMigration, /grant select on table public\.consultation_requests to service_role/i);
assert.doesNotMatch(statusReadMigration, /grant select[\s\S]*to (anon|authenticated)/i);
});
test("detached completion and cancellation use a bounded retry ceiling", () => {
assert.match(consultRoute, /const detachedSettlementAttempts = 3/);
assert.match(consultRoute, /ponytail: Staging MVP ceiling—without a queue\/worker/);
assert.match(consultRoute, /for \(let attempt = 1; attempt <= detachedSettlementAttempts; attempt \+= 1\)/);
assert.match(consultRoute, /setTimeout\(resolve, attempt \* 150\)/);
assert.match(consultRoute, /async function cancel\(\)[\s\S]*retryDetachedSettlement\(async \(\) =>[\s\S]*rpc\("cancel_consultation_credit"/);
assert.match(consultRoute, /async function completeResponse\([\s\S]*retryDetachedSettlement\(async \(\) =>[\s\S]*rpc\("complete_consultation_response"/);
});
test("status expires stale reservations after fifteen minutes but leaves fresh reservations active", () => {
assert.match(statusRoute, /const reservedLeaseMs = 15 \* 60 \* 1000/);
assert.match(statusRoute, /ponytail: Staging MVP ceiling—without a queue\/worker/);
assert.match(statusRoute, /Number\.isFinite\(timestamp\) && now - timestamp >= reservedLeaseMs/);
const leaseBlock = statusRoute.slice(
statusRoute.indexOf('if (statusData.status === "reserved" && reservationLeaseExpired(statusData.updated_at))'),
statusRoute.indexOf("return NextResponse.json({\n requestId: statusData.request_id"),
);
assert.match(leaseBlock, /runCreditRpc\([\s\S]*"cancel_consultation_credit"/);
assert.match(leaseBlock, /\.eq\("session_id", statusData\.session_id\)[\s\S]*\.eq\("request_id", statusData\.request_id\)/);
assert.doesNotMatch(statusRoute.slice(0, statusRoute.indexOf('if (statusData.status === "reserved"')), /cancel_consultation_credit/);
});
test("an empty active consultation lookup is 204 and a missing bound request stays 404", () => {
const activeEmpty = statusRoute.slice(
statusRoute.lastIndexOf("if (activeLookup)"),
statusRoute.indexOf("咨询请求不存在"),
);
assert.match(activeEmpty, /!data \|\| typeof data\.session_id !== "string"/);
assert.match(activeEmpty, /new NextResponse\(null, \{ status: 204 \}\)/);
assert.doesNotMatch(activeEmpty, /咨询请求不存在|status: 404/);
assert.match(statusRoute, /return NextResponse\.json\(\{ error: "咨询请求不存在" \}, \{ status: 404 \}\)/);
const cloud = read("src/lib/home-cloud-sync.ts");
const activeFn = cloud.slice(
cloud.indexOf("export async function fetchActiveConsultationStatus"),
cloud.indexOf("export async function patchSessionModel"),
);
const boundFn = cloud.slice(
cloud.indexOf("export async function fetchConsultationStatus"),
cloud.indexOf("export async function fetchActiveConsultationStatus"),
);
const statusGuard = activeFn.indexOf("if (response.status === 204 || response.status === 404) return null;");
assert.ok(statusGuard >= 0 && statusGuard < activeFn.indexOf("response.json()"));
assert.doesNotMatch(boundFn, /status === 204/);
});