fix(consult): paced settle is opt-in for the consultation reply; one rejection path for HTTP bodies and stream events (T2/T3 follow-up)
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
0bc6659052
commit
891e9c62da
@@ -808,14 +808,14 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
// run was stopped or failed, in which case what arrived shows at once.
|
||||
const settleFrames = (immediate: boolean) => {
|
||||
if (immediate) {
|
||||
void frames.settle({ immediate: true });
|
||||
void frames.settle();
|
||||
return;
|
||||
}
|
||||
settleRequested = true;
|
||||
setStreamingReply((current) => (
|
||||
current && current.sessionId === sessionId ? { ...current, settling: requestId } : current
|
||||
));
|
||||
void frames.settle();
|
||||
void frames.settle({ paced: true });
|
||||
};
|
||||
try {
|
||||
const response = await fetch("/api/consult", {
|
||||
@@ -849,16 +849,18 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
}),
|
||||
signal: controller.signal,
|
||||
});
|
||||
// One rejection path for the HTTP body a pre-stream check answers with
|
||||
// and for the run.failed request_rejected event the stream carries once
|
||||
// it is open (BUG-1074): same status semantics either way.
|
||||
const rejectRun = (status: number, message: string, code?: string): never => {
|
||||
if (status === 401) window.location.assign("/login");
|
||||
if (status === 402) openAccountDialog("billing", { source: "insufficient-credits" });
|
||||
throw new ConsultationResponseError(status, message, code);
|
||||
};
|
||||
if (!response.ok) {
|
||||
const contentType = response.headers.get("content-type") ?? "";
|
||||
const errorPayload = contentType.includes("application/json") ? await response.json() : { message: await response.text() };
|
||||
if (response.status === 401) window.location.assign("/login");
|
||||
if (response.status === 402) openAccountDialog("billing", { source: "insufficient-credits" });
|
||||
throw new ConsultationResponseError(
|
||||
response.status,
|
||||
payloadMessage(errorPayload, "服务暂时不可用"),
|
||||
payloadCode(errorPayload),
|
||||
);
|
||||
rejectRun(response.status, payloadMessage(errorPayload, "服务暂时不可用"), payloadCode(errorPayload));
|
||||
}
|
||||
if (!response.body) {
|
||||
throw new ConsultationResponseError(502, "浏览器未收到可读取的回答流");
|
||||
@@ -966,10 +968,7 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
// 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);
|
||||
rejectRun(event.status ?? 503, payloadMessage({ message: event.message }, "服务暂时不可用"), event.reason);
|
||||
}
|
||||
if (event.type === "run.failed") {
|
||||
if (event.code === "answer_truncated") {
|
||||
@@ -1086,7 +1085,7 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
// Whatever arrived before the failure is what gets kept, not just the
|
||||
// part the pacing had released so far.
|
||||
settleRequested = false;
|
||||
void frames.settle({ immediate: true });
|
||||
void frames.settle();
|
||||
const cancelled = controller.signal.aborted;
|
||||
const ownsInterface = pendingConsultation.current?.requestId === requestId;
|
||||
const partialReply = latestPartialReply;
|
||||
|
||||
@@ -7,9 +7,10 @@
|
||||
* happens per animation frame. Answer and thinking text are released at a
|
||||
* steady per-frame pace so a burst of chunks reads as flowing text instead of
|
||||
* a jump, while a large backlog (reconnect, slow tab) catches up in roughly a
|
||||
* dozen frames. When the stream ends, what is still unreleased is written out
|
||||
* at the same pace within a short budget instead of appearing in one frame
|
||||
* (BUG-1075); only a stop, a failure or a hidden document releases at once.
|
||||
* dozen frames. When a consultation stream ends, what is still unreleased is
|
||||
* written out at the same pace within a short budget instead of appearing in
|
||||
* one frame (a paced settle, BUG-1075); a stop, a failure, a hidden document
|
||||
* and every other caller release at once.
|
||||
*
|
||||
* Pure release arithmetic lives in exported functions so the policy is
|
||||
* testable without a DOM; scheduling is injectable for the same reason.
|
||||
@@ -94,14 +95,15 @@ export type StreamFrameBuffer<Meta> = Readonly<{
|
||||
/** Publish meta-only changes (timeline rows, activity) on the next frame. */
|
||||
touch: () => void;
|
||||
/**
|
||||
* The stream has ended. By default what is still unreleased is written out
|
||||
* at the typing pace within STREAM_SETTLE_MAX_FRAMES, and the last frame
|
||||
* flushes with `settled: true`; the promise resolves on that frame. With
|
||||
* `immediate` (stop, failure, disconnect) everything is released in one
|
||||
* synchronous flush, as before BUG-1075. A hidden document always releases
|
||||
* The stream has ended. By default everything unreleased goes out in one
|
||||
* synchronous flush (stop, failure, disconnect, and surfaces that merge a
|
||||
* snapshot right after). With `paced` (the consultation reply, BUG-1075)
|
||||
* what is left is written out at the typing pace within
|
||||
* STREAM_SETTLE_MAX_FRAMES, only the last frame flushes with `settled: true`,
|
||||
* and the promise resolves on that frame. A hidden document always releases
|
||||
* at once.
|
||||
*/
|
||||
settle: (options?: Readonly<{ immediate?: boolean }>) => Promise<void>;
|
||||
settle: (options?: Readonly<{ paced?: boolean }>) => Promise<void>;
|
||||
/** Drop everything, including scheduled work, without flushing. */
|
||||
reset: (meta?: Meta) => void;
|
||||
dispose: () => void;
|
||||
@@ -233,7 +235,7 @@ export function createStreamFrameBuffer<Meta>(
|
||||
if (settling) {
|
||||
// A second settle while one is writing out: an immediate one takes
|
||||
// over and flushes now; a paced one just waits for the first.
|
||||
if (!settleOptions?.immediate) {
|
||||
if (settleOptions?.paced) {
|
||||
return new Promise<void>((resolve) => {
|
||||
const previous = settling!.resolve;
|
||||
settling!.resolve = () => { previous(); resolve(); };
|
||||
@@ -242,7 +244,7 @@ export function createStreamFrameBuffer<Meta>(
|
||||
cancelScheduled();
|
||||
}
|
||||
const pending = pendingChars(releasedAnswer, targetAnswer);
|
||||
if (settleOptions?.immediate || scheduler.hidden() || pending <= 0) {
|
||||
if (!settleOptions?.paced || scheduler.hidden() || pending <= 0) {
|
||||
cancelScheduled();
|
||||
releasedAnswer = targetAnswer;
|
||||
releasedThinking = targetThinking;
|
||||
|
||||
Reference in New Issue
Block a user