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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0199rbQDTsUbCVw84wc8BTFe
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
3f8b817261
commit
0bc6659052
@@ -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<typeof consultationAgentPublicEventSchema>;
|
||||
|
||||
|
||||
@@ -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<Response>;
|
||||
/** 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<void>;
|
||||
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<string, unknown> : {};
|
||||
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<Uint8Array> | null = null;
|
||||
|
||||
const body = new ReadableStream<Uint8Array>({
|
||||
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,
|
||||
},
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user