/** * Frame-coalesced release of streamed agent output. * * Every network chunk used to become its own React commit, and each commit * re-parsed the whole partial answer. This buffer sits between the event * parser and `setState`: events mutate an accumulator, and at most one flush * 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. * * 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; export const STREAM_RELEASE_CATCHUP_DIVISOR = 12; 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. */ 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)), ); } /** Advance a released prefix toward its target by one frame's worth of text. */ export function advanceStreamRelease(released: string, target: string, backlogChars?: 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)); } const pending = target.length - released.length; if (pending <= 0) return target; return target.slice(0, released.length + streamReleaseCount(pending, backlogChars ?? pending)); } export type StreamFrameSnapshot = Readonly<{ answer: string; thinking: string; meta: Meta; /** True when this flush released everything that had arrived. */ settled: boolean; }>; export type StreamFrameScheduler = Readonly<{ requestFrame: (callback: () => void) => number; cancelFrame: (handle: number) => void; requestTimeout: (callback: () => void, delayMs: number) => number; cancelTimeout: (handle: number) => void; hidden: () => boolean; }>; export type StreamFrameBufferOptions = Readonly<{ initialMeta: Meta; flush: (snapshot: StreamFrameSnapshot) => void; scheduler?: StreamFrameScheduler; }>; export type StreamFrameBuffer = Readonly<{ setAnswer: (fullText: string) => void; setThinking: (fullText: string) => void; 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; /** Drop everything, including scheduled work, without flushing. */ reset: (meta?: Meta) => void; dispose: () => void; /** Text released so far, for callers that persist partial output. */ released: () => Readonly<{ answer: string; thinking: string }>; }>; function pendingChars(released: string, target: string): number { return target.startsWith(released) ? target.length - released.length : target.length; } function browserScheduler(): StreamFrameScheduler { return { requestFrame: (callback) => window.requestAnimationFrame(callback), cancelFrame: (handle) => window.cancelAnimationFrame(handle), requestTimeout: (callback, delayMs) => window.setTimeout(callback, delayMs), cancelTimeout: (handle) => window.clearTimeout(handle), hidden: () => typeof document !== "undefined" && document.hidden, }; } export function createStreamFrameBuffer( options: StreamFrameBufferOptions, ): StreamFrameBuffer { const scheduler = options.scheduler ?? browserScheduler(); let targetAnswer = ""; let targetThinking = ""; let releasedAnswer = ""; let releasedThinking = ""; let answerBacklog = 0; let thinkingBacklog = 0; let meta = options.initialMeta; let disposed = false; let frameHandle: number | null = null; let timeoutHandle: number | null = null; const cancelScheduled = () => { if (frameHandle !== null) { scheduler.cancelFrame(frameHandle); frameHandle = null; } if (timeoutHandle !== null) { scheduler.cancelTimeout(timeoutHandle); timeoutHandle = null; } }; const emit = (settled: boolean) => { options.flush({ answer: releasedAnswer, thinking: releasedThinking, meta, settled, }); }; const step = () => { frameHandle = null; timeoutHandle = null; if (disposed) return; if (scheduler.hidden()) { releasedAnswer = targetAnswer; releasedThinking = targetThinking; } else { releasedAnswer = advanceStreamRelease(releasedAnswer, targetAnswer, answerBacklog); releasedThinking = advanceStreamRelease(releasedThinking, targetThinking, thinkingBacklog); } if (releasedAnswer === targetAnswer) answerBacklog = 0; if (releasedThinking === targetThinking) thinkingBacklog = 0; const caughtUp = releasedAnswer === targetAnswer && releasedThinking === targetThinking; emit(caughtUp); if (!caughtUp) schedule(); }; const schedule = () => { if (disposed || frameHandle !== null || timeoutHandle !== null) return; if (scheduler.hidden()) { timeoutHandle = scheduler.requestTimeout(step, STREAM_HIDDEN_FLUSH_MS); } else { frameHandle = scheduler.requestFrame(step); } }; return { setAnswer(fullText) { if (disposed || fullText === targetAnswer) return; targetAnswer = fullText; answerBacklog = Math.max(answerBacklog, pendingChars(releasedAnswer, targetAnswer)); schedule(); }, setThinking(fullText) { if (disposed || fullText === targetThinking) return; targetThinking = fullText; thinkingBacklog = Math.max(thinkingBacklog, pendingChars(releasedThinking, targetThinking)); schedule(); }, setMeta(next) { if (disposed) return; meta = typeof next === "function" ? (next as (current: Meta) => Meta)(meta) : next; schedule(); }, touch() { if (disposed) return; schedule(); }, settle() { if (disposed) return; cancelScheduled(); releasedAnswer = targetAnswer; releasedThinking = targetThinking; answerBacklog = 0; thinkingBacklog = 0; emit(true); }, reset(nextMeta) { cancelScheduled(); targetAnswer = ""; targetThinking = ""; releasedAnswer = ""; releasedThinking = ""; answerBacklog = 0; thinkingBacklog = 0; if (nextMeta !== undefined) meta = nextMeta; }, dispose() { disposed = true; cancelScheduled(); }, released() { return { answer: releasedAnswer, thinking: releasedThinking }; }, }; }