From 0bc66590522731f446d6c2edca7c60f2d0517ee0 Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Mon, 28 Sep 2026 09:20:23 +0800 Subject: [PATCH] feat(consult): open the response stream before preparation, reservation and classification; rejections ride the stream (T2, BUG-1074) Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_0199rbQDTsUbCVw84wc8BTFe --- frontend/src/app/api/consult/route.ts | 33 +++ frontend/src/components/chat-message-row.tsx | 1 + frontend/src/hooks/use-consultation-run.ts | 16 +- frontend/src/lib/consultation-agent-events.ts | 22 +- frontend/src/lib/stream-first-response.ts | 154 ++++++++++++ ...t-first-frame-and-pacing-20260928.test.tsx | 225 ++++++++++++++++++ .../consultation-agentic-runtime.test.ts | 6 +- frontend/tests/consultation-smalltalk.test.ts | 6 +- 8 files changed, 458 insertions(+), 5 deletions(-) create mode 100644 frontend/src/lib/stream-first-response.ts diff --git a/frontend/src/app/api/consult/route.ts b/frontend/src/app/api/consult/route.ts index be4e227b..2d579e53 100644 --- a/frontend/src/app/api/consult/route.ts +++ b/frontend/src/app/api/consult/route.ts @@ -40,6 +40,7 @@ import { createServerSupabaseClient } from "@/lib/supabase/server"; import { streamTextResponse } from "@/lib/stream-text-response"; import { classifyConsultationTurn, type SmalltalkUsage } from "@/lib/consultation-smalltalk"; import { streamSmalltalkResponse } from "@/lib/stream-smalltalk-response"; +import { streamFirstResponse } from "@/lib/stream-first-response"; import { consultationPublicActivityEvent, streamAgentResponse } from "@/lib/stream-agent-response"; import type { AgentExecutionReceipt, WorkflowReceipt } from "@/lib/consultation-agent-events"; import { @@ -283,6 +284,7 @@ function workflowStatus(status: string | undefined): "ready" | "degraded" | "blo } export async function POST(request: Request) { + const requestStartedAt = Date.now(); let supabase: Awaited>; let accounting: ReturnType; try { @@ -406,6 +408,16 @@ export async function POST(request: Request) { type SelectedModel = NonNullable>>; type ReservationResult = { success: boolean; credits: number | null; error_code: string | null }; type ModelSelection = Awaited>>; + + // Everything above answered with a real HTTP status: authentication, the + // request body, the session and its model. Everything below — subject and + // profile reads, the credit reservation, appending the question, the model + // classification, the agent itself — runs inside the response stream in the + // agentic runtime (BUG-1074): `runTurn` still returns the same Responses it + // always did, and `streamFirstResponse` relays them after a first frame that + // goes out before any of this work starts. The legacy runtime keeps the + // pre-stream order because its text/plain replies cannot ride an NDJSON stream. + const runTurn = async (): Promise => { let prepared: PreparedConsultationRoute; try { prepared = await prepareConsultationRoute({ @@ -1625,4 +1637,25 @@ export async function POST(request: Request) { { status: 503 }, ); } + }; + + if (!shouldUseAgenticRuntime(user)) return runTurn(); + return streamFirstResponse({ + requestId, + run: runTurn, + startedAt: requestStartedAt, + onFirstFrame: (elapsedMs) => { + logAgentObservability({ + requestId, + sessionId, + agentVersion: "consultation-stream-first-v1", + contractPhases: [{ phase: "first_byte", durationMs: elapsedMs, status: "completed" }], + billingSettlementResult: "not_applicable", + }); + }, + onUnhandled: (error) => { + const reason = error instanceof Error ? error.name : "UnknownError"; + console.error(`[consult] turn failed before a response request=${requestId} reason=${reason}`); + }, + }); } diff --git a/frontend/src/components/chat-message-row.tsx b/frontend/src/components/chat-message-row.tsx index 34536130..13286f32 100644 --- a/frontend/src/components/chat-message-row.tsx +++ b/frontend/src/components/chat-message-row.tsx @@ -66,6 +66,7 @@ export const ChatMessageRow = memo(function ChatMessageRow({ : "Jyotisha"; const activityState = message.activity ? ({ + received: "working", "loading-method": "searching", "chart-calculation": "solving", "evidence-validation": "working", diff --git a/frontend/src/hooks/use-consultation-run.ts b/frontend/src/hooks/use-consultation-run.ts index 51bade8c..4c42a73e 100644 --- a/frontend/src/hooks/use-consultation-run.ts +++ b/frontend/src/hooks/use-consultation-run.ts @@ -863,7 +863,9 @@ export function useConsultationRun(params: ConsultationRunParams) { if (!response.body) { throw new ConsultationResponseError(502, "浏览器未收到可读取的回答流"); } - responseKind = response.headers.get("x-jyotish-response-kind") === "smalltalk" ? "smalltalk" : undefined; + // Whether the reply was smalltalk is only known once the server has + // classified the turn, which now happens inside the stream (BUG-1074): + // the run.completed event says so; there is no response header for it. let techniqueTruth = response.headers.get("x-jyotish-technique-truth") ?? "unknown"; let workflowReceipt: AgentExecutionReceipt["workflow"] = { route: response.headers.get("x-jyotish-workflow-route") ?? "unknown", @@ -891,6 +893,9 @@ export function useConsultationRun(params: ConsultationRunParams) { label: CONSULTATION_CHART_CALCULATION_LABEL, completedTrail: activityCompletedTrail([CONSULTATION_DONE_SKILL_LABEL]), }; + } else if (event.type === "activity" && event.phase === "received") { + // The stream's first byte, before anything has been done: no trail. + activity = { phase: event.phase, label: event.label }; } else if (event.type === "activity") { activity = { phase: event.phase, @@ -957,6 +962,15 @@ export function useConsultationRun(params: ConsultationRunParams) { techniqueTruth = event.receipt.techniqueTruth ?? "unknown"; } } + if (event.type === "run.failed" && event.code === "request_rejected") { + // A rejection the route used to send as a JSON body with an HTTP + // status before the stream opened (BUG-1074): same status, same + // sentence, same code, so the handling below does not change. + const status = event.status ?? 503; + if (status === 401) window.location.assign("/login"); + if (status === 402) openAccountDialog("billing", { source: "insufficient-credits" }); + throw new ConsultationResponseError(status, payloadMessage({ message: event.message }, "服务暂时不可用"), event.reason); + } if (event.type === "run.failed") { if (event.code === "answer_truncated") { truncatedFailure = event; diff --git a/frontend/src/lib/consultation-agent-events.ts b/frontend/src/lib/consultation-agent-events.ts index c846aa42..0441b361 100644 --- a/frontend/src/lib/consultation-agent-events.ts +++ b/frontend/src/lib/consultation-agent-events.ts @@ -3,6 +3,10 @@ import { consultationDomainSchema, type ConsultationDomain } from "./consultatio import { publicThinkingSectionSchema } from "./consultation-thinking-plan.ts"; export const publicActivityPhaseSchema = z.enum([ + // The first byte of a consultation stream, before the turn is classified or + // the chart is read (BUG-1074). The client already shows the same queued + // row from the send frame, so this frame changes nothing on screen. + "received", "loading-method", "chart-calculation", "evidence-validation", @@ -142,10 +146,18 @@ const runCompletedSchema = z.object({ // allowlisted receipt a completed run does. It stays optional because the // receipt is built from live state that a hard failure may leave unparseable, // and losing the whole failure event would be worse than losing its receipt. +// `request_rejected` is a failure the route used to answer with a JSON body +// and an HTTP status before the stream existed (credits, plan, subject, +// session full…). Since the stream opens first (BUG-1074) it travels here +// instead: `status` keeps the HTTP semantics the client acts on (402 opens +// billing, 4xx means the reservation never committed), `reason` is the JSON +// body's `code` when it had one (session_full, pricing_configuration_unavailable). const runFailedSchema = z.object({ type: z.literal("run.failed"), - code: z.enum(["runtime_contract_incomplete", "calculation_failed", "empty_answer", "answer_truncated", "cancelled"]), + code: z.enum(["runtime_contract_incomplete", "calculation_failed", "empty_answer", "answer_truncated", "cancelled", "request_rejected"]), message: z.string().max(200), + status: z.number().int().min(400).max(599).optional(), + reason: z.string().max(80).optional(), receipt: agentExecutionReceiptSchema.optional(), }).strict(); @@ -159,6 +171,14 @@ export const consultationAgentPublicEventSchema = z.discriminatedUnion("type", [ && (event.responseKind === "smalltalk" ? event.receipt !== undefined : event.receipt === undefined)) { ctx.addIssue({ code: z.ZodIssueCode.custom, message: "completion_requires_receipt_or_smalltalk" }); } + if (event.type === "run.failed") { + if ((event.status !== undefined || event.reason !== undefined) && event.code !== "request_rejected") { + ctx.addIssue({ code: z.ZodIssueCode.custom, message: "status_and_reason_only_on_request_rejected" }); + } + if (event.code === "request_rejected" && event.status === undefined) { + ctx.addIssue({ code: z.ZodIssueCode.custom, message: "request_rejected_requires_status" }); + } + } }); export type ConsultationAgentPublicEvent = z.infer; diff --git a/frontend/src/lib/stream-first-response.ts b/frontend/src/lib/stream-first-response.ts new file mode 100644 index 00000000..258a0bb4 --- /dev/null +++ b/frontend/src/lib/stream-first-response.ts @@ -0,0 +1,154 @@ +import { + consultationAgentPublicEventSchema, + type ConsultationAgentPublicEvent, +} from "./consultation-agent-events.ts"; +import { CONSULTATION_RECEIVED_LABEL } from "./consultation-activity-labels.ts"; + +/** + * Stream-first consultation response (TASK-consult-first-frame-and-pacing-20260928 + * D2, BUG-1074; the same shape BUG-1047 gave the rectification turn). + * + * The route used to read the session, prepare the chart, reserve credits, + * append the question and classify the turn with a model call before it built + * a Response, so the browser received nothing for all of that. Here the + * Response exists first: its first byte is a deterministic `activity` frame, + * and the whole turn runs inside the stream. The inner work still produces the + * same Responses it always did — an NDJSON stream (consultation or smalltalk) + * or a JSON rejection with an HTTP status — and they are relayed as they are: + * NDJSON bytes pass through unchanged after the first frame; a JSON rejection + * becomes one `run.failed` event that carries the status, message and code the + * client used to read from the HTTP response. + * + * Only what has to answer with a real HTTP status stays in front of the + * stream: authentication, request-body validation and the session lookup. + */ + +/** The first byte of every consultation stream: the same sentence the client's queued row shows from the send frame. */ +export const CONSULTATION_RECEIVED_EVENT: ConsultationAgentPublicEvent = { + type: "activity", + phase: "received", + label: CONSULTATION_RECEIVED_LABEL, +}; + +export const STREAM_FIRST_FALLBACK_MESSAGE = "咨询服务暂时不可用,请稍后再试。"; + +export type StreamFirstOptions = Readonly<{ + requestId: string; + /** Written before anything else happens; defaults to the received frame. */ + firstFrame?: ConsultationAgentPublicEvent; + /** The turn. Its Response is relayed: NDJSON as bytes, JSON as one run.failed event. */ + run: () => Promise; + /** Called once the first frame has been written, with the milliseconds since `startedAt`. */ + onFirstFrame?: (elapsedMs: number) => void; + /** Called when `run` throws instead of returning a Response; the stream then ends with a generic failure. */ + onUnhandled?: (error: unknown) => void | Promise; + startedAt?: number; + now?: () => number; +}>; + +/** + * The JSON body the route answered with, as the failure event the stream + * carries instead. Message precedence is the one the client applied to the + * HTTP body (`payloadMessage`: recovery, then message, then error), so the + * sentence the user sees does not change; `code` rides along as `reason`. + */ +export function consultationRejectionEvent(status: number, payload: unknown): ConsultationAgentPublicEvent { + const body = payload && typeof payload === "object" ? payload as Record : {}; + const message = [body.recovery, body.message, body.error] + .find((value): value is string => typeof value === "string" && value.trim().length > 0) + ?? STREAM_FIRST_FALLBACK_MESSAGE; + const reason = typeof body.code === "string" && body.code ? body.code.slice(0, 80) : undefined; + const safeStatus = Number.isInteger(status) && status >= 400 && status <= 599 ? status : 503; + return { + type: "run.failed", + code: "request_rejected", + message: message.slice(0, 200), + status: safeStatus, + ...(reason ? { reason } : {}), + }; +} + +function isNdjson(response: Response): boolean { + return (response.headers.get("content-type") ?? "").includes("application/x-ndjson"); +} + +export function streamFirstResponse(options: StreamFirstOptions): Response { + const encoder = new TextEncoder(); + const now = options.now ?? (() => Date.now()); + const startedAt = options.startedAt ?? now(); + let disconnected = false; + let innerReader: ReadableStreamDefaultReader | null = null; + + const body = new ReadableStream({ + start(controller) { + const send = (event: ConsultationAgentPublicEvent) => { + if (disconnected) return; + const safe = consultationAgentPublicEventSchema.parse(event); + controller.enqueue(encoder.encode(`${JSON.stringify(safe)}\n`)); + }; + // The first byte goes out before the turn starts, synchronously. + send(options.firstFrame ?? CONSULTATION_RECEIVED_EVENT); + options.onFirstFrame?.(Math.max(0, now() - startedAt)); + + void (async () => { + try { + const inner = await options.run(); + if (isNdjson(inner) && inner.body) { + innerReader = inner.body.getReader(); + if (disconnected) { + // The client left while the turn was being prepared: hand the + // disconnect to the inner stream, which keeps its own settlement. + await innerReader.cancel().catch(() => {}); + } else { + for (;;) { + const { done, value } = await innerReader.read(); + if (done) break; + if (!disconnected && value) controller.enqueue(value); + } + } + } else { + const text = await inner.text().catch(() => ""); + let payload: unknown = null; + try { + payload = text ? JSON.parse(text) : null; + } catch { + payload = { message: text }; + } + send(consultationRejectionEvent(inner.status, payload)); + } + } catch (error) { + await Promise.resolve(options.onUnhandled?.(error)).catch(() => {}); + try { + send({ type: "run.failed", code: "calculation_failed", message: STREAM_FIRST_FALLBACK_MESSAGE }); + } catch { + // The stream is already gone; nothing to tell. + } + } finally { + if (!disconnected) { + try { + controller.close(); + } catch { + // Closed by the consumer first. + } + } + } + })(); + }, + async cancel() { + disconnected = true; + // The inner stream owns settlement and continues after a disconnect, + // exactly as it did when the client held it directly. + await innerReader?.cancel().catch(() => {}); + }, + }); + + return new Response(body, { + headers: { + "cache-control": "no-cache, no-transform", + "content-type": "application/x-ndjson; charset=utf-8", + "x-accel-buffering": "no", + "x-ayanam-mode": "mastra-agentic", + "x-ayanam-request-id": options.requestId, + }, + }); +} diff --git a/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx b/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx index 7e4f20af..1c018c01 100644 --- a/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx +++ b/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx @@ -74,3 +74,228 @@ test("the consultation hook paces the write-out only for a finished run; stop, f // No second animation system: the write-out is the same frame buffer the live stream uses. assert.equal((hook.match(/createStreamFrameBuffer { + const reader = response.body!.getReader(); + let text = ""; + const lines: string[] = []; + for (;;) { + const { done, value } = await reader.read(); + if (value) text += decoder.decode(value, { stream: true }); + const parts = text.split("\n"); + text = parts.pop() ?? ""; + lines.push(...parts.filter(Boolean)); + if (done || lines.length >= limit) break; + } + if (lines.length >= limit) await reader.cancel().catch(() => {}); + return lines; +} + +function ndjson(events: ConsultationAgentPublicEvent[], headers: Record = {}) { + const encoder = new TextEncoder(); + const body = new ReadableStream({ + start(controller) { + for (const event of events) controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`)); + controller.close(); + }, + }); + return new Response(body, { headers: { "content-type": "application/x-ndjson; charset=utf-8", ...headers } }); +} + +test("① the first frame is written before the turn has done anything, even while classification hangs", async () => { + let firstFrameMs = -1; + let ran = false; + const response = streamFirstResponse({ + requestId: "fake-request", + startedAt: 1_000, + now: () => 1_040, + onFirstFrame: (elapsed) => { firstFrameMs = elapsed; }, + run: () => new Promise(() => { ran = true; }), // never resolves: a hanging classification + }); + assert.equal(response.headers.get("content-type"), "application/x-ndjson; charset=utf-8"); + assert.equal(response.headers.get("x-ayanam-request-id"), "fake-request"); + const [first] = await readLines(response, 1); + const event = consultationAgentPublicEventSchema.parse(JSON.parse(first!)); + assert.deepEqual(event, CONSULTATION_RECEIVED_EVENT); + assert.equal(event.type === "activity" && event.label, CONSULTATION_RECEIVED_LABEL); + assert.equal(firstFrameMs, 40); + assert.equal(ran, true); +}); + +test("② a smalltalk turn completes inside the stream: free settlement once, answer, receipt-free completion", async () => { + let persisted = 0; + const response = streamFirstResponse({ + requestId: "fake-request", + run: async () => streamSmalltalkResponse({ + requestId: "fake-request", + reply: "你好,想聊什么都可以", + complete: async () => { persisted += 1; }, + onError: async () => assert.fail("no error path"), + }), + }); + const events = (await readLines(response)).map((line) => consultationAgentPublicEventSchema.parse(JSON.parse(line))); + assert.deepEqual(events.map((event) => event.type), ["activity", "answer.delta", "run.completed"]); + assert.equal(persisted, 1); + assert.deepEqual(events[2], { type: "run.completed", responseKind: "smalltalk" }); + // The header the client used to read is gone with the stream-first order; the event carries the kind. + assert.equal(response.headers.get("x-jyotish-response-kind"), null); +}); + +test("③ a JSON rejection becomes one run.failed request_rejected with the same status, sentence and code", async () => { + const insufficient = consultationRejectionEvent(402, { error: "咨询点数不足", message: "请先兑换咨询点数后再继续。" }); + assert.deepEqual(insufficient, { type: "run.failed", code: "request_rejected", status: 402, message: "请先兑换咨询点数后再继续。" }); + const full = consultationRejectionEvent(409, { error: "这段对话已写满", message: "这段对话已写满,开个新对话继续吧", code: "session_full" }); + assert.deepEqual(full, { type: "run.failed", code: "request_rejected", status: 409, message: "这段对话已写满,开个新对话继续吧", reason: "session_full" }); + // Same precedence the client applied to the HTTP body: recovery, then message, then error. + const generation = consultationRejectionEvent(503, { error: "暂时无法生成解读", message: "咨询服务暂时不可用,请稍后再试。", recovery: "稍后重试,或换一个模型继续。" }); + assert.equal(generation.type === "run.failed" && generation.message, "稍后重试,或换一个模型继续。"); + assert.equal(consultationRejectionEvent(200, {}).type === "run.failed" && consultationRejectionEvent(200, {}).status, 503); + assert.equal(consultationRejectionEvent(503, null).type === "run.failed" && consultationRejectionEvent(503, null).message, STREAM_FIRST_FALLBACK_MESSAGE); + for (const event of [insufficient, full, generation]) consultationAgentPublicEventSchema.parse(event); + + const response = streamFirstResponse({ + requestId: "fake-request", + run: async () => new Response(JSON.stringify({ error: "咨询点数不足", message: "请先兑换咨询点数后再继续。" }), { + status: 402, + headers: { "content-type": "application/json" }, + }), + }); + const events = (await readLines(response)).map((line) => consultationAgentPublicEventSchema.parse(JSON.parse(line))); + assert.deepEqual(events.map((event) => event.type), ["activity", "run.failed"]); + assert.deepEqual(events[1], insufficient); +}); + +test("the failure schema keeps status and reason to request_rejected, which must carry a status", () => { + assert.equal(consultationAgentPublicEventSchema.safeParse({ type: "run.failed", code: "request_rejected", message: "x" }).success, false); + assert.equal(consultationAgentPublicEventSchema.safeParse({ type: "run.failed", code: "calculation_failed", message: "x", status: 503 }).success, false); + assert.equal(consultationAgentPublicEventSchema.safeParse({ type: "run.failed", code: "empty_answer", message: "x", reason: "why" }).success, false); + assert.equal(consultationAgentPublicEventSchema.safeParse({ type: "run.failed", code: "request_rejected", message: "x", status: 409, reason: "session_full" }).success, true); + assert.equal(consultationAgentPublicEventSchema.safeParse({ type: "activity", phase: "received", label: CONSULTATION_RECEIVED_LABEL }).success, true); +}); + +test("④ a consultation stream passes through byte for byte after the first frame; the inner headers are not needed", async () => { + const inner: ConsultationAgentPublicEvent[] = [ + { type: "run.started", runId: "fake-request", requestId: "fake-request" }, + { type: "skill.started", name: "jyotish-vedic-astrology" }, + { type: "answer.delta", text: "不是。" }, + { type: "run.failed", code: "answer_truncated", message: "回答没写完" }, + ]; + const response = streamFirstResponse({ + requestId: "fake-request", + run: async () => ndjson(inner, { "x-jyotish-birth-time-mode": "verified_chart" }), + }); + const lines = await readLines(response); + assert.equal(lines.length, inner.length + 1); + assert.deepEqual(JSON.parse(lines[0]!), CONSULTATION_RECEIVED_EVENT); + assert.deepEqual(lines.slice(1), inner.map((event) => JSON.stringify(event))); + assert.equal(response.headers.get("x-jyotish-birth-time-mode"), null); +}); + +test("a client disconnect is handed to the inner stream, which keeps its own settlement", async () => { + let innerCancelled = false; + let release: (() => void) | undefined; + const innerBody = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`${JSON.stringify({ type: "run.started", runId: "r", requestId: "r" })}\n`)); + // Then hang until the test releases it. + void new Promise((resolve) => { release = resolve; }).then(() => { + try { + controller.close(); + } catch { + // Cancelled by the consumer first: the disconnect path under test. + } + }); + }, + cancel() { innerCancelled = true; }, + }); + const response = streamFirstResponse({ + requestId: "r", + run: async () => new Response(innerBody, { headers: { "content-type": "application/x-ndjson; charset=utf-8" } }), + }); + const reader = response.body!.getReader(); + await reader.read(); // the first frame + await reader.read(); // run.started, relayed + await reader.cancel(); + await new Promise((resolve) => setTimeout(resolve, 0)); + assert.equal(innerCancelled, true); + release?.(); +}); + +test("a turn that throws before producing any response ends the stream with a generic failure and is logged", async () => { + const seen: unknown[] = []; + const response = streamFirstResponse({ + requestId: "fake-request", + onUnhandled: (error) => { seen.push(error); }, + run: async () => { throw new Error("session_full"); }, + }); + const events = (await readLines(response)).map((line) => consultationAgentPublicEventSchema.parse(JSON.parse(line))); + assert.deepEqual(events.map((event) => event.type), ["activity", "run.failed"]); + assert.deepEqual(events[1], { type: "run.failed", code: "calculation_failed", message: STREAM_FIRST_FALLBACK_MESSAGE }); + assert.equal(seen.length, 1); +}); + +test("the consult route opens the stream before subject reads, reservation, question append and classification", () => { + const route = source("../src/app/api/consult/route.ts"); + const openIndex = route.indexOf("const runTurn = async (): Promise =>"); + assert.ok(openIndex > 0); + // Stays in front of the stream: auth, body, session existence and type, the session model. + for (const marker of [ + "status: 401", + '{ error: "出生资料或问题格式不正确", details: parsed.error.flatten() }', + 'error: "咨询会话不存在"', + 'code: "session_not_consultation"', + 'error: "会话模型已经变化"', + "if (blocksPromptExtraction(userControlledPrompt))", + ]) { + assert.ok(route.indexOf(marker) < openIndex, `${marker} must answer with an HTTP status`); + } + // Inside the stream: the slow part, and every rejection it can raise. + for (const marker of [ + "prepareConsultationRoute({", + 'accounting.rpc("reserve_consultation_usage"', + 'accounting.rpc("append_consultation_question"', + "await classifyConsultationTurn({", + 'code: "session_full"', + "insufficient ? 402 : 503", + 'error: "暂时无法生成解读"', + ]) { + assert.ok(route.indexOf(marker) > openIndex, `${marker} must run inside the stream`); + } + assert.match(route, /if \(!shouldUseAgenticRuntime\(user\)\) return runTurn\(\);/); + assert.match(route, /return streamFirstResponse\(\{\s*requestId,\s*run: runTurn,\s*startedAt: requestStartedAt,/); + assert.match(route, /phase: "first_byte", durationMs: elapsedMs, status: "completed"/); + // The smalltalk branch is untouched: it still settles for free and never runs the agent. + const smalltalk = route.slice(route.indexOf('if (turn.kind === "smalltalk")'), route.indexOf("const expectedTitle")); + assert.match(smalltalk, /complete_consultation_free/); + assert.doesNotMatch(smalltalk, /streamAgentResponse\(/); +}); + +test("the client reads smalltalk and rejections from events, with the same status semantics as before", () => { + const hook = source("../src/hooks/use-consultation-run.ts"); + assert.doesNotMatch(hook, /x-jyotish-response-kind/); + assert.match(hook, /if \(event\.responseKind === "smalltalk"\) \{\s*responseKind = "smalltalk";/); + assert.match(hook, /event\.type === "run\.failed" && event\.code === "request_rejected"/); + assert.match(hook, /if \(status === 402\) openAccountDialog\("billing", \{ source: "insufficient-credits" \}\);/); + assert.match(hook, /throw new ConsultationResponseError\(status, payloadMessage\(\{ message: event\.message \}, "服务暂时不可用"\), event\.reason\);/); + // The rollback / notice handling keyed on status and code is unchanged. + assert.match(hook, /caught\.code === "session_full"/); + assert.match(hook, /const reserveDidNotCommit = caught\.status === 400\s*\|\| caught\.status === 401\s*\|\| caught\.status === 402\s*\|\| caught\.status === 409;/); + // The first frame carries no completed trail. + assert.match(hook, /event\.type === "activity" && event\.phase === "received"/); +}); diff --git a/frontend/tests/consultation-agentic-runtime.test.ts b/frontend/tests/consultation-agentic-runtime.test.ts index eb3a80ca..3497787e 100644 --- a/frontend/tests/consultation-agentic-runtime.test.ts +++ b/frontend/tests/consultation-agentic-runtime.test.ts @@ -32,7 +32,7 @@ import { ConsultationWorkflowError, consultationWorkflowFailureCode, } from "../src/mastra/consultation-workflow.ts"; -import { agentExecutionReceiptSchema } from "../src/lib/consultation-agent-events.ts"; +import { agentExecutionReceiptSchema, type PublicActivityPhase } from "../src/lib/consultation-agent-events.ts"; import { consultationDomainIds, consultationDomainPlanValues } from "../src/lib/consultation-domain-registry.ts"; import { OPENER_THINKING_TITLE, @@ -1582,7 +1582,9 @@ function makeWindowCtx(options: { fetchWindowChart?: () => Promise void) { +// 原值: phase 手写四值联合 新值: 用 schema 导出的 PublicActivityPhase(多了 "received") +// 原因: TASK-consult-first-frame-and-pacing-20260928 D2 / BUG-1074:流的首帧是 phase "received" 的 activity 事件 +async function writerToSend(send: (event: { type: "activity"; phase: PublicActivityPhase; label: string }) => void) { return { custom: async (chunk: unknown) => { if (!chunk || typeof chunk !== "object") return; diff --git a/frontend/tests/consultation-smalltalk.test.ts b/frontend/tests/consultation-smalltalk.test.ts index 43fdbaa4..6049d493 100644 --- a/frontend/tests/consultation-smalltalk.test.ts +++ b/frontend/tests/consultation-smalltalk.test.ts @@ -213,7 +213,11 @@ test("disconnect does not cancel the server-owned free completion", async () => test("smalltalk views never reconstruct an invented settled execution timeline", () => { assert.equal(settledChatMessageViews([{ role: "assistant", text: "你好", responseKind: "smalltalk" }])[0]?.timeline, undefined); assert.equal(streamingChatMessageView([{ role: "user", text: "你好" }], true, "你好", undefined, undefined, undefined, [], "smalltalk")?.responseKind, "smalltalk"); - assert.match(source("../src/components/chat-message-row.tsx"), /const quiet = smalltalk \|\| awaitingClassification/); + // 原值: /const quiet = smalltalk \|\| awaitingClassification/(分类回来前整行静默) + // 新值: /const quiet = smalltalk;/(只有寒暄回复静默;分类前显示排队行) + // 原因: TASK-consult-first-frame-and-pacing-20260928 D1 / BUG-1074 + assert.match(source("../src/components/chat-message-row.tsx"), /const quiet = smalltalk;/); + assert.doesNotMatch(source("../src/components/chat-message-row.tsx"), /awaitingClassification/); assert.match(source("../src/lib/home-cloud-sync.ts"), /stored.responseKind === "smalltalk"/); }); test("route gates all special entries before a single selected-model classifier and keeps tool contracts", () => {