From 3f8b8172616cfca86caf4c30c94c9df4e79e07a0 Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Mon, 28 Sep 2026 09:14:46 +0800 Subject: [PATCH] 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 Claude-Session: https://claude.ai/code/session_0199rbQDTsUbCVw84wc8BTFe --- frontend/DESIGN.md | 1 + frontend/src/components/chat-message-row.tsx | 2 +- frontend/src/hooks/use-consultation-run.ts | 41 +++++- frontend/src/lib/chat-message-view.ts | 31 ++++- frontend/src/lib/home-types.ts | 2 + frontend/src/lib/stream-frame-buffer.ts | 121 ++++++++++++++---- ...t-first-frame-and-pacing-20260928.test.tsx | 76 +++++++++++ frontend/tests/stream-frame-buffer.test.ts | 104 ++++++++++++++- 8 files changed, 346 insertions(+), 32 deletions(-) create mode 100644 frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md index b9a1daf2..f5017bbd 100644 --- a/frontend/DESIGN.md +++ b/frontend/DESIGN.md @@ -753,6 +753,7 @@ or user IDs. | 行内 / 局部等待 | 出生地解析、两个会话面时间线的 live 步、兜底活动面板的 live 行、个人报告列表与详情 | `InlineSpinner`(`inline-spin`) | 0.8s linear | `animation: none`,收成静止圆点,不要半圈圆弧 | | 盘位骨架(仅星盘页) | D1 / 非 D1 分盘 / 西洋盘尚未返回 | `VedicChartSvg skeleton` / `WesternWheelSvg skeleton`;统一 `.chart-page-skeleton` | 2.4s ease-in-out,opacity 0.35 ↔ 0.6 | 停止动画,opacity 0.5 | | 流式生成中 | 引导语打字、时间线 summary 与 live 行的文案 | `onboarding-caret` / `agent-activity-shimmer` | 700ms steps / 1.6s linear | 保持现有全局降级 | +| 流式正文的放字节奏 | 普通对话与校正面的回答正文(`stream-frame-buffer`) | 每帧最多 4 字(≈ 240 字/秒),积压超过 600 字才按 1/12 追赶;流结束后剩余文字按同一节奏在 1.5 s 内写完,最后一帧才算 settled;停止 / 失败 / 断线 / 页面隐藏一帧全放(2026-09-28,BUG-1075) | 逐帧 | 不涉及 CSS 动画,无需降级 | Agent 的 live 标记只有 `InlineSpinner` 一种。曾经并存的 canvas 小球(`thinking-orbs`)已移除,不得再引入第二种 live 标记。 diff --git a/frontend/src/components/chat-message-row.tsx b/frontend/src/components/chat-message-row.tsx index f204a120..34536130 100644 --- a/frontend/src/components/chat-message-row.tsx +++ b/frontend/src/components/chat-message-row.tsx @@ -145,7 +145,7 @@ export const ChatMessageRow = memo(function ChatMessageRow({
{hasAnswer ? ( diff --git a/frontend/src/hooks/use-consultation-run.ts b/frontend/src/hooks/use-consultation-run.ts index fbe1f48a..51bade8c 100644 --- a/frontend/src/hooks/use-consultation-run.ts +++ b/frontend/src/hooks/use-consultation-run.ts @@ -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; diff --git a/frontend/src/lib/chat-message-view.ts b/frontend/src/lib/chat-message-view.ts index 9dab5f3b..1a7b3ddf 100644 --- a/frontend/src/lib/chat-message-view.ts +++ b/frontend/src/lib/chat-message-view.ts @@ -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; } diff --git a/frontend/src/lib/home-types.ts b/frontend/src/lib/home-types.ts index d31d2ff3..9fa9512b 100644 --- a/frontend/src/lib/home-types.ts +++ b/frontend/src/lib/home-types.ts @@ -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[]; diff --git a/frontend/src/lib/stream-frame-buffer.ts b/frontend/src/lib/stream-frame-buffer.ts index 92410aef..5c0942e4 100644 --- a/frontend/src/lib/stream-frame-buffer.ts +++ b/frontend/src/lib/stream-frame-buffer.ts @@ -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 = Readonly<{ @@ -70,8 +93,15 @@ export type StreamFrameBuffer = 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; /** Drop everything, including scheduled work, without flushing. */ reset: (meta?: Meta) => void; dispose: () => void; @@ -107,6 +137,10 @@ export function createStreamFrameBuffer( 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( }); }; + 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( 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( 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( 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((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((resolve) => { + settling = { perFrame: settleReleaseCount(pending), resolve }; + schedule(); + }); }, reset(nextMeta) { cancelScheduled(); + if (settling) finishSettle(); targetAnswer = ""; targetThinking = ""; releasedAnswer = ""; @@ -197,6 +269,11 @@ export function createStreamFrameBuffer( 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(); }, diff --git a/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx b/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx new file mode 100644 index 00000000..7e4f20af --- /dev/null +++ b/frontend/tests/consult-first-frame-and-pacing-20260928.test.tsx @@ -0,0 +1,76 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; +import React from "react"; +import { renderToStaticMarkup } from "react-dom/server"; + +import { ChatMessageRow } from "../src/components/chat-message-row.tsx"; +import { + chatMessageViews, + latestAssistantView, + settledChatMessageViews, + settlingChatMessageView, +} from "../src/lib/chat-message-view.ts"; + +const source = (path: string) => readFileSync(new URL(path, import.meta.url), "utf8"); + +const stored = [ + { role: "user" as const, text: "我妈是不是不太在意我" }, + { + role: "assistant" as const, + text: "不是。4 宫主水星逆行落 10 宫,她把心力投在你的前途上——不是不在意,是不会用你想要的方式在意。", + techniqueTruth: "unknown", + }, +]; + +// T3 / BUG-1075: after the run has ended the stored reply is complete, but the +// client keeps writing it out at the typing pace. The trailing row shows the +// released prefix under the same render key, and its timeline is not live. +test("a settled reply still being written out renders the released prefix, same render key, timeline not live", () => { + const prefix = "不是。4 宫主水星逆行落 10 宫,"; + const latest = latestAssistantView(stored, false, prefix); + assert.ok(latest); + assert.equal(latest.view.text, prefix); + assert.equal(latest.view.state, "streaming"); + assert.equal(latest.view.settling, true); + assert.equal(latest.view.renderKey, settledChatMessageViews(stored)[1]!.renderKey); + assert.equal(latest.views.length, 2); + assert.equal(latest.views[1]!.text, prefix); + + const views = chatMessageViews(stored, false, prefix); + assert.equal(views.at(-1)!.text, prefix); + assert.equal(views.at(-1)!.settling, true); + + const html = renderToStaticMarkup(); + assert.match(html, /4 宫主水星逆行落 10 宫,/); + assert.doesNotMatch(html, /她把心力投在你的前途上/); + // The step timeline is the settled one: no live row, no spinner, no 正在分析. + assert.doesNotMatch(html, /inline-spinner|正在分析|收到,正在看你的问题/); +}); + +test("the write-out override applies only to a strict prefix of the stored reply while nothing is loading", () => { + const last = settledChatMessageViews(stored)[1]!; + assert.equal(settlingChatMessageView(last, true, "不是。"), undefined, "still loading: the streaming view owns the row"); + assert.equal(settlingChatMessageView(last, false, ""), undefined, "nothing released: stored reply as is"); + assert.equal(settlingChatMessageView(last, false, last.text), undefined, "fully written out: stored reply as is"); + assert.equal(settlingChatMessageView(last, false, "另一条回答"), undefined, "not a prefix (a replaced reply): stored reply as is"); + assert.equal(latestAssistantView(stored, false, "另一条回答")?.view.text, last.text); + const user = settledChatMessageViews([{ role: "user", text: "问" }])[0]!; + assert.equal(settlingChatMessageView(user, false, "问"), undefined); +}); + +test("the consultation hook paces the write-out only for a finished run; stop, failure and truncation release at once", () => { + const hook = source("../src/hooks/use-consultation-run.ts"); + // The stream ended normally: paced unless the user stopped it or the server cut it. + assert.match(hook, /settleFrames\(controller\.signal\.aborted \|\| Boolean\(truncatedFailure\)\);/); + assert.match(hook, /settleFrames\(controller\.signal\.aborted\);/); + // Failure path keeps what arrived, at once. + assert.match(hook, /settleRequested = false;\s*void frames\.settle\(\{ immediate: true \}\);/); + // The paced write-out owns the display through the streaming reply's `settling` id; + // the interface completes without waiting for it and without clearing it. + assert.match(hook, /setStreamingReply\(\(current\) => \(current\?\.settling \? current : null\)\);/); + assert.match(hook, /current\?\.settling !== requestId \? current : frame\.settled \? null : next/); + assert.match(hook, /\{ \.\.\.current, settling: requestId \}/); + // No second animation system: the write-out is the same frame buffer the live stream uses. + assert.equal((hook.match(/createStreamFrameBuffer { +test("many events collapse into one flush per frame and a paced settle writes the rest out within the budget", () => { const fake = fakeScheduler(); const flushes: StreamFrameSnapshot[] = []; const buffer = createStreamFrameBuffer({ @@ -112,11 +115,108 @@ test("many events collapse into one flush per frame and settle releases everythi assert.equal(flushes.at(-1)!.meta.length, 200); assert.ok(flushes.at(-1)!.answer.length < 200, "pacing is still behind the network"); - buffer.settle(); + // 原值: buffer.settle() 同步一帧放完剩余文字(settled: true 立即到) + // 新值: 默认 settle 按打字节奏在 STREAM_SETTLE_MAX_FRAMES 内写完,最后一帧才 settled: true;immediate 才一帧放完 + // 原因: TASK-consult-first-frame-and-pacing-20260928 D3 / BUG-1075:流尾一帧全放让短回答「一下全出来」 + const pendingAtSettle = answer.length - flushes.at(-1)!.answer.length; + let resolved = false; + void buffer.settle().then(() => { resolved = true; }); + assert.equal(flushes.at(-1)!.settled, false, "a paced settle does not flush synchronously"); + let settleFrames = 0; + while (fake.scheduledFrames > 0 && settleFrames < 200) { + fake.tick(); + settleFrames += 1; + } + assert.ok(settleFrames >= Math.ceil(pendingAtSettle / STREAM_RELEASE_MAX_CHARS) - 1, `wrote out ${pendingAtSettle} in ${settleFrames} frames`); + assert.ok(settleFrames <= STREAM_SETTLE_MAX_FRAMES, `took ${settleFrames} frames`); assert.equal(flushes.at(-1)!.answer, answer); assert.equal(flushes.at(-1)!.settled, true); + assert.equal(flushes.at(-2)!.settled, false, "only the last paced frame is settled"); assert.equal(fake.scheduledFrames, 0); assert.equal(buffer.released().answer, answer); + return Promise.resolve().then(() => assert.equal(resolved, true)); +}); + +test("typing pace: an ordinary backlog is capped at four characters a frame, a large one still catches up", () => { + // The server releases the natal opener as one 160-character lump (ANSWER_RELEASE_CHARS). + assert.equal(streamReleaseCount(160), STREAM_RELEASE_MAX_CHARS); + assert.equal(streamReleaseCount(STREAM_RELEASE_CATCHUP_CHARS), STREAM_RELEASE_MAX_CHARS); + assert.equal(streamReleaseCount(STREAM_RELEASE_CATCHUP_CHARS + 1), Math.ceil((STREAM_RELEASE_CATCHUP_CHARS + 1) / STREAM_RELEASE_CATCHUP_DIVISOR)); + let released = ""; + const lump = "字".repeat(160); + let frames = 0; + while (released !== lump && frames < 1_000) { + released = advanceStreamRelease(released, lump, lump.length); + frames += 1; + } + assert.equal(frames, 160 / STREAM_RELEASE_MAX_CHARS, "a 160-character lump is typed out over 40 frames, not 12"); +}); + +test("a paced settle writes 300 leftover characters within the budget; immediate, stop-style settle is one flush", () => { + const fake = fakeScheduler(); + const flushes: StreamFrameSnapshot[] = []; + const buffer = createStreamFrameBuffer({ + initialMeta: null, + scheduler: fake.scheduler, + flush: (snapshot) => flushes.push(snapshot), + }); + const text = "字".repeat(300); + buffer.setAnswer(text); + fake.tick(); + const shownBefore = flushes.at(-1)!.answer.length; + assert.ok(shownBefore < 300); + + void buffer.settle(); + let frames = 0; + while (fake.scheduledFrames > 0 && frames < 500) { + fake.tick(); + frames += 1; + } + assert.ok(frames <= STREAM_SETTLE_MAX_FRAMES, `took ${frames} frames`); + assert.ok(frames > 1, "not one frame"); + assert.equal(flushes.at(-1)!.answer, text); + assert.equal(flushes.at(-1)!.settled, true); + for (let index = 1; index < flushes.length; index += 1) { + assert.ok(flushes[index]!.answer.length >= flushes[index - 1]!.answer.length, "monotonic"); + } + + const immediate = createStreamFrameBuffer({ + initialMeta: null, + scheduler: fake.scheduler, + flush: (snapshot) => flushes.push(snapshot), + }); + immediate.setAnswer(text); + const before = flushes.length; + void immediate.settle({ immediate: true }); + assert.equal(flushes.length, before + 1); + assert.equal(flushes.at(-1)!.answer, text); + assert.equal(flushes.at(-1)!.settled, true); + assert.equal(fake.scheduledFrames, 0); +}); + +test("dispose during a paced settle lets it finish, then silences the buffer", () => { + const fake = fakeScheduler(); + const flushes: StreamFrameSnapshot[] = []; + const buffer = createStreamFrameBuffer({ + initialMeta: null, + scheduler: fake.scheduler, + flush: (snapshot) => flushes.push(snapshot), + }); + const text = "字".repeat(40); + buffer.setAnswer(text); + fake.tick(); + void buffer.settle(); + buffer.dispose(); + let frames = 0; + while (fake.scheduledFrames > 0 && frames < 100) { + fake.tick(); + frames += 1; + } + assert.equal(flushes.at(-1)!.answer, text); + assert.equal(flushes.at(-1)!.settled, true); + buffer.setAnswer("不再发布"); + buffer.touch(); + assert.equal(fake.tick(), 0); }); test("thinking text is paced separately from the answer and meta-only touches still flush", () => {