import assert from "node:assert/strict"; import test from "node:test"; import { STREAM_HIDDEN_FLUSH_MS, STREAM_RELEASE_CATCHUP_CHARS, STREAM_RELEASE_CATCHUP_DIVISOR, STREAM_RELEASE_MAX_CHARS, STREAM_RELEASE_MIN_CHARS, STREAM_SETTLE_MAX_FRAMES, advanceStreamRelease, createStreamFrameBuffer, streamReleaseCount, type StreamFrameScheduler, type StreamFrameSnapshot, } from "../src/lib/stream-frame-buffer.ts"; function fakeScheduler(hidden = () => false) { const frames: Array<() => void> = []; const timeouts: Array<{ callback: () => void; delayMs: number }> = []; let handle = 0; const scheduler: StreamFrameScheduler = { requestFrame(callback) { frames.push(callback); handle += 1; return handle; }, cancelFrame() { frames.length = 0; }, requestTimeout(callback, delayMs) { timeouts.push({ callback, delayMs }); handle += 1; return handle; }, cancelTimeout() { timeouts.length = 0; }, hidden, }; return { scheduler, tick() { const pending = frames.splice(0); for (const callback of pending) callback(); return pending.length; }, tickTimeouts() { const pending = timeouts.splice(0); for (const entry of pending) entry.callback(); return pending; }, get scheduledFrames() { return frames.length; }, }; } test("release count is at least two characters and catches up a backlog within about twelve frames", () => { assert.equal(streamReleaseCount(0), 0); assert.equal(streamReleaseCount(1), 1); assert.equal(streamReleaseCount(2), STREAM_RELEASE_MIN_CHARS); assert.equal(streamReleaseCount(5), STREAM_RELEASE_MIN_CHARS); assert.equal(streamReleaseCount(24), STREAM_RELEASE_MIN_CHARS); assert.equal(streamReleaseCount(25), 3); assert.equal(streamReleaseCount(1200), 1200 / STREAM_RELEASE_CATCHUP_DIVISOR); let released = ""; const target = "字".repeat(3_000); let frames = 0; while (released !== target && frames < 100) { released = advanceStreamRelease(released, target, target.length); frames += 1; } // A 3,000-character backlog that arrived at once clears in twelve frames (~200ms). assert.equal(frames, STREAM_RELEASE_CATCHUP_DIVISOR); assert.equal(released, target); // Without the backlog figure the pace still floors at two characters per frame. assert.equal(advanceStreamRelease("", "十二个字符十二个字符十二"), "十二"); }); test("a replaced target that no longer extends the released prefix jumps instead of stalling", () => { assert.equal(advanceStreamRelease("旧的回答", "新"), "新"); // A replacement restarts at the paced rate from the new head rather than showing stale text. assert.equal(advanceStreamRelease("abc", "abd"), "ab"); }); 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({ initialMeta: [], scheduler: fake.scheduler, flush: (snapshot) => flushes.push(snapshot), }); // 200 one-character tokens arrive four per frame across fifty frames, the way a // model streams Chinese text; each frame is allowed one commit. let answer = ""; let frameCount = 0; for (let index = 0; index < 200; index += 1) { answer += "字"; buffer.setAnswer(answer); buffer.setMeta((rows) => [...rows, `row-${index}`]); if (index % 4 === 3) frameCount += fake.tick(); } assert.equal(frameCount, 50); assert.equal(flushes.length, 50); assert.ok(flushes.length * 3 <= 200, "at most one commit per frame, not per token"); for (let index = 1; index < flushes.length; index += 1) { assert.ok(flushes[index]!.answer.length >= flushes[index - 1]!.answer.length); assert.ok(flushes[index]!.answer.length - flushes[index - 1]!.answer.length >= STREAM_RELEASE_MIN_CHARS); } assert.equal(flushes.at(-1)!.meta.length, 200); assert.ok(flushes.at(-1)!.answer.length < 200, "pacing is still behind the network"); // 原值: buffer.settle() 同步一帧放完剩余文字(settled: true 立即到) // 新值: settle({ paced: true }) 按打字节奏在 STREAM_SETTLE_MAX_FRAMES 内写完,最后一帧才 settled: true;无参 settle() 仍一帧放完(下一条测试) // 原因: 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({ paced: true }).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; the default 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({ paced: true }); 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(); 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({ paced: true }); 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", () => { const fake = fakeScheduler(); const flushes: StreamFrameSnapshot[] = []; const buffer = createStreamFrameBuffer({ initialMeta: null, scheduler: fake.scheduler, flush: (snapshot) => flushes.push(snapshot), }); buffer.setThinking("先看事业宫,再看大运。"); fake.tick(); assert.equal(flushes.length, 1); assert.equal(flushes[0]!.answer, ""); assert.ok(flushes[0]!.thinking.length >= STREAM_RELEASE_MIN_CHARS); assert.equal(flushes[0]!.settled, false); buffer.settle(); assert.equal(flushes.at(-1)!.thinking, "先看事业宫,再看大运。"); buffer.touch(); fake.tick(); assert.equal(flushes.length, 3); assert.equal(flushes.at(-1)!.settled, true); }); test("a burst that lands mid-stream is cleared within about twelve frames instead of trickling", () => { const fake = fakeScheduler(); const flushes: StreamFrameSnapshot[] = []; const buffer = createStreamFrameBuffer({ initialMeta: null, scheduler: fake.scheduler, flush: (snapshot) => flushes.push(snapshot), }); buffer.setAnswer("字".repeat(20)); fake.tick(); buffer.setAnswer("字".repeat(2_420)); let frames = 0; while (fake.scheduledFrames > 0 && frames < 100) { fake.tick(); frames += 1; } assert.equal(flushes.at(-1)!.answer.length, 2_420); assert.ok(frames <= STREAM_RELEASE_CATCHUP_DIVISOR + 1, `took ${frames} frames`); }); test("a hidden document falls back to a timeout and releases everything at once", () => { const fake = fakeScheduler(() => true); const flushes: StreamFrameSnapshot[] = []; const buffer = createStreamFrameBuffer({ initialMeta: null, scheduler: fake.scheduler, flush: (snapshot) => flushes.push(snapshot), }); buffer.setAnswer("字".repeat(500)); assert.equal(fake.scheduledFrames, 0); const fired = fake.tickTimeouts(); assert.equal(fired.length, 1); assert.equal(fired[0]!.delayMs, STREAM_HIDDEN_FLUSH_MS); assert.equal(flushes.length, 1); assert.equal(flushes[0]!.answer.length, 500); assert.equal(flushes[0]!.settled, true); }); test("reset drops received and released text plus scheduled work, and dispose silences the buffer", () => { const fake = fakeScheduler(); const flushes: StreamFrameSnapshot[] = []; const buffer = createStreamFrameBuffer({ initialMeta: 1, scheduler: fake.scheduler, flush: (snapshot) => flushes.push(snapshot), }); buffer.setAnswer("第一次尝试的正文"); fake.tick(); assert.equal(flushes.length, 1); buffer.reset(2); assert.equal(fake.scheduledFrames, 0); assert.deepEqual(buffer.released(), { answer: "", thinking: "" }); buffer.touch(); fake.tick(); assert.equal(flushes.at(-1)!.answer, ""); assert.equal(flushes.at(-1)!.meta, 2); buffer.dispose(); buffer.setAnswer("不再发布"); buffer.touch(); assert.equal(fake.tick(), 0); buffer.settle(); assert.equal(flushes.at(-1)!.answer, ""); });