feat(consult): type out streamed text at a capped pace and write the tail out on settle instead of one frame (T3, BUG-1075)
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
86d99115c6
commit
3f8b817261
@@ -145,7 +145,7 @@ export const ChatMessageRow = memo(function ChatMessageRow({
|
||||
<div className="consultation-thinking-report">
|
||||
<ConsultationRunTimeline
|
||||
rows={consultTimeline}
|
||||
live={showActivity && message.state !== "settled"}
|
||||
live={showActivity && message.state !== "settled" && !message.settling}
|
||||
vargaSentence={vargaSentence}
|
||||
/>
|
||||
{hasAnswer ? (
|
||||
|
||||
@@ -489,7 +489,9 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
if (pendingConsultation.current?.requestId !== requestId) return;
|
||||
pendingConsultation.current = null;
|
||||
if (consultationReplayStarted.current === requestId) consultationReplayStarted.current = null;
|
||||
setStreamingReply(null);
|
||||
// A reply still being written out at the typing pace keeps the display
|
||||
// (BUG-1075); its last paced frame clears it. Everything else is done now.
|
||||
setStreamingReply((current) => (current?.settling ? current : null));
|
||||
setPendingSessionId(null);
|
||||
setPendingRequestId(null);
|
||||
setConsultationPhase(null);
|
||||
@@ -766,6 +768,11 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
let streamedThinking = "";
|
||||
let timelineState = emptyConsultationTimeline();
|
||||
let currentActivity: AgentActivityView | undefined;
|
||||
// True once the stream has ended and the buffer is writing out what is
|
||||
// left at the typing pace (BUG-1075). From then on the display belongs to
|
||||
// this request only while the streaming reply still carries its id: a new
|
||||
// send replaces the reply, and the old buffer's frames fall through.
|
||||
let settleRequested = false;
|
||||
// Every stream event lands in `frames`; it commits at most once per animation
|
||||
// frame and releases text at a steady pace. Nothing below calls
|
||||
// setStreamingReply directly while the response body is being read.
|
||||
@@ -775,7 +782,7 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
const partialReply = parseAgentReply(frame.answer).text;
|
||||
latestPartialReply = partialReply;
|
||||
thinkingSections = applyThinkingSectionProgress(thinkingSections, partialReply);
|
||||
setStreamingReply({
|
||||
const next: StreamingReply = {
|
||||
sessionId,
|
||||
text: partialReply,
|
||||
responseKind,
|
||||
@@ -783,12 +790,33 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
thinkingSections: thinkingSections.length ? thinkingSections : undefined,
|
||||
timeline: timelineState.rows,
|
||||
activity: currentActivity,
|
||||
});
|
||||
...(settleRequested ? { settling: requestId } : {}),
|
||||
};
|
||||
if (settleRequested) {
|
||||
setStreamingReply((current) => (
|
||||
current?.settling !== requestId ? current : frame.settled ? null : next
|
||||
));
|
||||
} else {
|
||||
setStreamingReply(next);
|
||||
}
|
||||
if (partialReply && pendingConsultation.current?.requestId === requestId) {
|
||||
pendingConsultation.current = { ...pendingConsultation.current, partialReply };
|
||||
}
|
||||
},
|
||||
});
|
||||
// The stream has ended: write out the rest at the typing pace unless the
|
||||
// run was stopped or failed, in which case what arrived shows at once.
|
||||
const settleFrames = (immediate: boolean) => {
|
||||
if (immediate) {
|
||||
void frames.settle({ immediate: true });
|
||||
return;
|
||||
}
|
||||
settleRequested = true;
|
||||
setStreamingReply((current) => (
|
||||
current && current.sessionId === sessionId ? { ...current, settling: requestId } : current
|
||||
));
|
||||
void frames.settle();
|
||||
};
|
||||
try {
|
||||
const response = await fetch("/api/consult", {
|
||||
method: "POST",
|
||||
@@ -949,7 +977,7 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
parser.push(decoder.decode(value, { stream: true }));
|
||||
}
|
||||
parser.finish(decoder.decode());
|
||||
frames.settle();
|
||||
settleFrames(controller.signal.aborted || Boolean(truncatedFailure));
|
||||
if (truncatedFailure) {
|
||||
const reply = parseAgentReply(answer);
|
||||
if (!reply.text) throw new ConsultationResponseError(502, truncatedFailure.message);
|
||||
@@ -992,7 +1020,7 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
}
|
||||
answer += decoder.decode();
|
||||
frames.setAnswer(answer);
|
||||
frames.settle();
|
||||
settleFrames(controller.signal.aborted);
|
||||
}
|
||||
if (controller.signal.aborted) return Boolean(latestPartialReply);
|
||||
const reply = parseAgentReply(answer);
|
||||
@@ -1043,7 +1071,8 @@ export function useConsultationRun(params: ConsultationRunParams) {
|
||||
} catch (caught) {
|
||||
// Whatever arrived before the failure is what gets kept, not just the
|
||||
// part the pacing had released so far.
|
||||
frames.settle();
|
||||
settleRequested = false;
|
||||
void frames.settle({ immediate: true });
|
||||
const cancelled = controller.signal.aborted;
|
||||
const ownsInterface = pendingConsultation.current?.requestId === requestId;
|
||||
const partialReply = latestPartialReply;
|
||||
|
||||
@@ -67,6 +67,12 @@ export type ChatMessage = {
|
||||
export type ChatMessageView = ChatMessage & {
|
||||
readonly renderKey: string;
|
||||
readonly state: "settled" | "streaming" | "thinking";
|
||||
/**
|
||||
* The run has ended and the reply is stored in full; the text shown is the
|
||||
* part the client has written out so far (BUG-1075). The step timeline is
|
||||
* already settled, so it is not live while this is true.
|
||||
*/
|
||||
readonly settling?: boolean;
|
||||
readonly activity?: AgentActivityView;
|
||||
readonly activityTrace?: readonly AgentActivityTraceItem[];
|
||||
readonly timeline?: readonly ConsultationTimelineRow[];
|
||||
@@ -145,9 +151,29 @@ export function latestAssistantView(
|
||||
if (streaming) return { view: streaming, views: [...settled, streaming] };
|
||||
const last = settled.at(-1);
|
||||
if (!last || last.role !== "assistant") return undefined;
|
||||
const writing = settlingChatMessageView(last, loading, streamingText);
|
||||
if (writing) return { view: writing, views: [...settled.slice(0, -1), writing] };
|
||||
return { view: last, views: settled };
|
||||
}
|
||||
|
||||
/**
|
||||
* After the run has ended the stored reply is complete, but the client may
|
||||
* still be writing it out at the typing pace (BUG-1075): while the released
|
||||
* text is a strict prefix of the stored reply, the trailing row shows the
|
||||
* released part under the same render key, without a live timeline. Anything
|
||||
* else (no text, text already complete, text that is not a prefix) renders
|
||||
* the stored reply as is.
|
||||
*/
|
||||
export function settlingChatMessageView(
|
||||
last: ChatMessageView,
|
||||
loading: boolean,
|
||||
streamingText: string,
|
||||
): ChatMessageView | undefined {
|
||||
if (loading || !streamingText || last.role !== "assistant") return undefined;
|
||||
if (last.text === streamingText || !last.text.startsWith(streamingText)) return undefined;
|
||||
return { ...last, text: streamingText, state: "streaming", settling: true };
|
||||
}
|
||||
|
||||
export function chatMessageViews(
|
||||
messages: readonly ChatMessage[],
|
||||
loading: boolean,
|
||||
@@ -169,5 +195,8 @@ export function chatMessageViews(
|
||||
timeline,
|
||||
responseKind,
|
||||
);
|
||||
return streaming ? [...settled, streaming] : settled;
|
||||
if (streaming) return [...settled, streaming];
|
||||
const last = settled.at(-1);
|
||||
const writing = last ? settlingChatMessageView(last, loading, streamingText) : undefined;
|
||||
return writing ? [...settled.slice(0, -1), writing] : settled;
|
||||
}
|
||||
|
||||
@@ -106,6 +106,8 @@ export type StreamingReply = {
|
||||
responseKind?: "smalltalk";
|
||||
sessionId: string;
|
||||
text: string;
|
||||
/** Set once the run has ended and the reply is being written out at the typing pace (BUG-1075): the request that owns the display. */
|
||||
settling?: string;
|
||||
activity?: AgentActivityView;
|
||||
thinkingText?: string;
|
||||
thinkingSections?: PublicThinkingSection[];
|
||||
|
||||
@@ -7,39 +7,62 @@
|
||||
* 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.
|
||||
* 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.
|
||||
*
|
||||
* Pure release arithmetic lives in exported functions so the policy is
|
||||
* testable without a DOM; scheduling is injectable for the same reason.
|
||||
*/
|
||||
|
||||
export const STREAM_RELEASE_MIN_CHARS = 2;
|
||||
/**
|
||||
* Typing pace: at most this many characters per frame while the backlog is
|
||||
* ordinary (BUG-1075). Four a frame is about 240 characters a second at 60 fps,
|
||||
* above what a model produces on average, so the backlog does not grow; a
|
||||
* server-held lump (the 160-character answer release, a whole short follow-up)
|
||||
* reads as writing instead of appearing at once.
|
||||
*/
|
||||
export const STREAM_RELEASE_MAX_CHARS = 4;
|
||||
export const STREAM_RELEASE_CATCHUP_DIVISOR = 12;
|
||||
/** A backlog above this (reconnect, slow tab) catches up in about twelve frames instead of typing it out. */
|
||||
export const STREAM_RELEASE_CATCHUP_CHARS = 600;
|
||||
/** A paced settle writes out what is left within this many frames (1.5 s at 60 fps). */
|
||||
export const STREAM_SETTLE_MAX_FRAMES = 90;
|
||||
export const STREAM_HIDDEN_FLUSH_MS = 250;
|
||||
|
||||
/**
|
||||
* Characters to reveal on one frame. `backlogChars` is how much was waiting
|
||||
* when the newest text arrived: dividing that by twelve clears any burst in
|
||||
* about twelve frames, while the two-character floor keeps a slow model from
|
||||
* reading as stalled. Callers without a backlog figure pass the pending count.
|
||||
* when the newest text arrived. An ordinary backlog is typed out at up to
|
||||
* STREAM_RELEASE_MAX_CHARS a frame, with a two-character floor so a slow model
|
||||
* never reads as stalled; a backlog above STREAM_RELEASE_CATCHUP_CHARS is
|
||||
* divided by twelve so any burst clears in about twelve frames. Callers
|
||||
* without a backlog figure pass the pending count.
|
||||
*/
|
||||
export function streamReleaseCount(pendingChars: number, backlogChars = pendingChars): number {
|
||||
if (pendingChars <= 0) return 0;
|
||||
return Math.min(
|
||||
pendingChars,
|
||||
Math.max(STREAM_RELEASE_MIN_CHARS, Math.ceil(backlogChars / STREAM_RELEASE_CATCHUP_DIVISOR)),
|
||||
);
|
||||
const catchUp = Math.ceil(backlogChars / STREAM_RELEASE_CATCHUP_DIVISOR);
|
||||
const perFrame = backlogChars > STREAM_RELEASE_CATCHUP_CHARS
|
||||
? catchUp
|
||||
: Math.min(STREAM_RELEASE_MAX_CHARS, catchUp);
|
||||
return Math.min(pendingChars, Math.max(STREAM_RELEASE_MIN_CHARS, perFrame));
|
||||
}
|
||||
|
||||
/** Characters per frame that write out `pendingChars` within the settle budget, never slower than the typing pace. */
|
||||
export function settleReleaseCount(pendingChars: number): number {
|
||||
if (pendingChars <= 0) return 0;
|
||||
return Math.max(STREAM_RELEASE_MAX_CHARS, Math.ceil(pendingChars / STREAM_SETTLE_MAX_FRAMES));
|
||||
}
|
||||
|
||||
/** Advance a released prefix toward its target by one frame's worth of text. */
|
||||
export function advanceStreamRelease(released: string, target: string, backlogChars?: number): string {
|
||||
export function advanceStreamRelease(released: string, target: string, backlogChars?: number, perFrame?: number): string {
|
||||
if (!target.startsWith(released)) {
|
||||
// The target was replaced rather than extended: restart from its head.
|
||||
return target.slice(0, streamReleaseCount(target.length, backlogChars ?? target.length));
|
||||
return target.slice(0, perFrame ?? streamReleaseCount(target.length, backlogChars ?? target.length));
|
||||
}
|
||||
const pending = target.length - released.length;
|
||||
if (pending <= 0) return target;
|
||||
return target.slice(0, released.length + streamReleaseCount(pending, backlogChars ?? pending));
|
||||
return target.slice(0, released.length + (perFrame ?? streamReleaseCount(pending, backlogChars ?? pending)));
|
||||
}
|
||||
|
||||
export type StreamFrameSnapshot<Meta> = Readonly<{
|
||||
@@ -70,8 +93,15 @@ export type StreamFrameBuffer<Meta> = Readonly<{
|
||||
setMeta: (next: Meta | ((current: Meta) => Meta)) => void;
|
||||
/** Publish meta-only changes (timeline rows, activity) on the next frame. */
|
||||
touch: () => void;
|
||||
/** Release everything received and flush synchronously. */
|
||||
settle: () => 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
|
||||
* at once.
|
||||
*/
|
||||
settle: (options?: Readonly<{ immediate?: boolean }>) => Promise<void>;
|
||||
/** Drop everything, including scheduled work, without flushing. */
|
||||
reset: (meta?: Meta) => void;
|
||||
dispose: () => void;
|
||||
@@ -107,6 +137,10 @@ export function createStreamFrameBuffer<Meta>(
|
||||
let disposed = false;
|
||||
let frameHandle: number | null = null;
|
||||
let timeoutHandle: number | null = null;
|
||||
// A paced settle in flight: its per-frame count, and who to tell when the
|
||||
// last frame has flushed. dispose() during it defers until that frame.
|
||||
let settling: { perFrame: number; resolve: () => void } | null = null;
|
||||
let disposeWhenSettled = false;
|
||||
|
||||
const cancelScheduled = () => {
|
||||
if (frameHandle !== null) {
|
||||
@@ -128,6 +162,16 @@ export function createStreamFrameBuffer<Meta>(
|
||||
});
|
||||
};
|
||||
|
||||
const finishSettle = () => {
|
||||
const done = settling;
|
||||
settling = null;
|
||||
if (disposeWhenSettled) {
|
||||
disposeWhenSettled = false;
|
||||
disposed = true;
|
||||
}
|
||||
done?.resolve();
|
||||
};
|
||||
|
||||
const step = () => {
|
||||
frameHandle = null;
|
||||
timeoutHandle = null;
|
||||
@@ -135,6 +179,9 @@ export function createStreamFrameBuffer<Meta>(
|
||||
if (scheduler.hidden()) {
|
||||
releasedAnswer = targetAnswer;
|
||||
releasedThinking = targetThinking;
|
||||
} else if (settling) {
|
||||
releasedAnswer = advanceStreamRelease(releasedAnswer, targetAnswer, answerBacklog, settling.perFrame);
|
||||
releasedThinking = targetThinking;
|
||||
} else {
|
||||
releasedAnswer = advanceStreamRelease(releasedAnswer, targetAnswer, answerBacklog);
|
||||
releasedThinking = advanceStreamRelease(releasedThinking, targetThinking, thinkingBacklog);
|
||||
@@ -143,7 +190,11 @@ export function createStreamFrameBuffer<Meta>(
|
||||
if (releasedThinking === targetThinking) thinkingBacklog = 0;
|
||||
const caughtUp = releasedAnswer === targetAnswer && releasedThinking === targetThinking;
|
||||
emit(caughtUp);
|
||||
if (!caughtUp) schedule();
|
||||
if (!caughtUp) {
|
||||
schedule();
|
||||
return;
|
||||
}
|
||||
if (settling) finishSettle();
|
||||
};
|
||||
|
||||
const schedule = () => {
|
||||
@@ -177,17 +228,38 @@ export function createStreamFrameBuffer<Meta>(
|
||||
if (disposed) return;
|
||||
schedule();
|
||||
},
|
||||
settle() {
|
||||
if (disposed) return;
|
||||
cancelScheduled();
|
||||
releasedAnswer = targetAnswer;
|
||||
releasedThinking = targetThinking;
|
||||
answerBacklog = 0;
|
||||
thinkingBacklog = 0;
|
||||
emit(true);
|
||||
settle(settleOptions) {
|
||||
if (disposed) return Promise.resolve();
|
||||
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) {
|
||||
return new Promise<void>((resolve) => {
|
||||
const previous = settling!.resolve;
|
||||
settling!.resolve = () => { previous(); resolve(); };
|
||||
});
|
||||
}
|
||||
cancelScheduled();
|
||||
}
|
||||
const pending = pendingChars(releasedAnswer, targetAnswer);
|
||||
if (settleOptions?.immediate || scheduler.hidden() || pending <= 0) {
|
||||
cancelScheduled();
|
||||
releasedAnswer = targetAnswer;
|
||||
releasedThinking = targetThinking;
|
||||
answerBacklog = 0;
|
||||
thinkingBacklog = 0;
|
||||
emit(true);
|
||||
if (settling) finishSettle();
|
||||
return Promise.resolve();
|
||||
}
|
||||
return new Promise<void>((resolve) => {
|
||||
settling = { perFrame: settleReleaseCount(pending), resolve };
|
||||
schedule();
|
||||
});
|
||||
},
|
||||
reset(nextMeta) {
|
||||
cancelScheduled();
|
||||
if (settling) finishSettle();
|
||||
targetAnswer = "";
|
||||
targetThinking = "";
|
||||
releasedAnswer = "";
|
||||
@@ -197,6 +269,11 @@ export function createStreamFrameBuffer<Meta>(
|
||||
if (nextMeta !== undefined) meta = nextMeta;
|
||||
},
|
||||
dispose() {
|
||||
if (settling) {
|
||||
// Let the paced settle write out its last frames; it disposes itself.
|
||||
disposeWhenSettled = true;
|
||||
return;
|
||||
}
|
||||
disposed = true;
|
||||
cancelScheduled();
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user