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, }, }); }