From ee2a8ed7d8fce1c6947cd9ef70c170bbe9d825d5 Mon Sep 17 00:00:00 2001 From: ChHsiching Date: Sun, 27 Sep 2026 18:25:55 +0800 Subject: [PATCH 1/4] fix: batch SSE events in use-convert to survive fast reasoning streams --- next/src/lib/__tests__/use-convert.test.ts | 273 +++++++++++++++++++++ next/src/lib/use-convert.ts | 84 ++++++- 2 files changed, 350 insertions(+), 7 deletions(-) create mode 100644 next/src/lib/__tests__/use-convert.test.ts diff --git a/next/src/lib/__tests__/use-convert.test.ts b/next/src/lib/__tests__/use-convert.test.ts new file mode 100644 index 00000000..8ab9c8cc --- /dev/null +++ b/next/src/lib/__tests__/use-convert.test.ts @@ -0,0 +1,273 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { act, createElement } from 'react'; +import { createRoot, type Root } from 'react-dom/client'; +import { createEventBatcher, useConvert } from '../use-convert'; +import { useStore } from '../store'; + +// act() refuses to run unless React knows it is inside a test environment +(globalThis as Record).IS_REACT_ACT_ENVIRONMENT = true; + +describe('createEventBatcher', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('buffers events until the 100ms window closes, preserving arrival order', async () => { + const seen: Array<[string, unknown]> = []; + const b = createEventBatcher((event, data) => seen.push([event, data])); + b.push('meta', { key: 'model' }); + b.push('stderr', { text: 'noise' }); + expect(seen).toHaveLength(0); + await vi.advanceTimersByTimeAsync(100); + expect(seen).toEqual([ + ['meta', { key: 'model' }], + ['stderr', { text: 'noise' }], + ]); + }); + + it('passes the first delta through immediately, draining the buffer first', () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('meta', { key: 'model' }); + b.push('meta', { key: 'session' }); + b.push('delta', { text: '

' }); + expect(seen).toEqual(['meta', 'meta', 'delta']); + }); + + it('buffers deltas after the first one', async () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('delta', { text: 'a' }); + b.push('delta', { text: 'b' }); + expect(seen).toEqual(['delta']); + await vi.advanceTimersByTimeAsync(100); + expect(seen).toEqual(['delta', 'delta']); + }); + + it('passes done and error through immediately, draining the buffer first', () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('meta', { key: 'model' }); + b.push('done', { code: 0 }); + expect(seen).toEqual(['meta', 'done']); + b.push('meta', { key: 'cost_usd' }); + b.push('error', { message: 'boom' }); + expect(seen).toEqual(['meta', 'done', 'meta', 'error']); + }); + + it('does not re-deliver drained events; a post-terminal event opens a fresh window', async () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('meta', { key: 'a' }); + b.push('done', { code: 0 }); + await vi.advanceTimersByTimeAsync(1000); + expect(seen).toEqual(['meta', 'done']); + b.push('meta', { key: 'b' }); + expect(seen).toEqual(['meta', 'done']); + await vi.advanceTimersByTimeAsync(100); + expect(seen).toEqual(['meta', 'done', 'meta']); + }); + + it('flush() drains synchronously and cancels the pending timer', async () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('meta', { key: 'a' }); + b.flush(); + expect(seen).toEqual(['meta']); + await vi.advanceTimersByTimeAsync(1000); + expect(seen).toEqual(['meta']); + b.push('meta', { key: 'b' }); + expect(seen).toEqual(['meta']); + await vi.advanceTimersByTimeAsync(100); + expect(seen).toEqual(['meta', 'meta']); + }); + + it('dispose() only clears the timer and never delivers', async () => { + const seen: string[] = []; + const b = createEventBatcher((event) => seen.push(event)); + b.push('meta', { key: 'a' }); + b.dispose(); + await vi.advanceTimersByTimeAsync(1000); + expect(seen).toHaveLength(0); + }); +}); + +describe('useConvert().run', () => { + let root: Root | null = null; + let container: HTMLElement | null = null; + let api: ReturnType | null = null; + + // useConvert is a hook — render a null harness once per test to capture it + function Harness() { + api = useConvert(); + return null; + } + + function sseFrame(event: string, data: unknown) { + return `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`; + } + + type ReadResult = { value?: Uint8Array; done: boolean }; + + // Duck-typed fetch: returns a Response-like whose body reader yields the + // given SSE frames one read() at a time. With hangUntilAbort the reader + // never reports done — it rejects with an AbortError once the request + // signal aborts, like a network read stalled mid-stream. + function stubStreamFetch(frames: string[], opts?: { hangUntilAbort?: boolean }) { + const enc = new TextEncoder(); + let i = 0; + let signal: AbortSignal | undefined; + let hanging = false; + const read = (): Promise => { + if (i < frames.length) { + return Promise.resolve({ value: enc.encode(frames[i++]), done: false }); + } + if (!opts?.hangUntilAbort) { + return Promise.resolve({ value: undefined, done: true }); + } + hanging = true; + return new Promise((_resolve, reject) => { + const abort = () => { + const err = new Error('aborted'); + err.name = 'AbortError'; + reject(err); + }; + if (signal?.aborted) abort(); + else signal?.addEventListener('abort', abort, { once: true }); + }); + }; + const res = { ok: true, body: { getReader: () => ({ read }) } }; + vi.stubGlobal( + 'fetch', + vi.fn(async (_url: unknown, init?: { signal?: AbortSignal }) => { + signal = init?.signal; + return res; + }), + ); + return { isHanging: () => hanging }; + } + + async function renderHarness() { + await act(async () => { + root!.render(createElement(Harness)); + }); + } + + function runReq(taskId: string) { + return { + taskId, + agent: 'test-agent', + templateId: 'article-magazine', + content: 'hello world', + }; + } + + function taskOf(taskId: string) { + const task = useStore.getState().tasks.find((t) => t.id === taskId); + expect(task).toBeDefined(); + return task!; + } + + beforeEach(() => { + localStorage.clear(); + container = document.createElement('div'); + document.body.appendChild(container); + root = createRoot(container); + }); + + afterEach(() => { + vi.unstubAllGlobals(); + if (root) { + act(() => root!.unmount()); + root = null; + } + container?.remove(); + container = null; + api = null; + }); + + it('marks the task errored with a truncation log when the stream ends without done/error', async () => { + const taskId = useStore.getState().newTask({ name: 'no-terminal' }); + stubStreamFetch([ + sseFrame('start', { bin: '/usr/bin/agent', promptBytes: 12 }), + sseFrame('delta', { text: '

partial' }), + sseFrame('meta', { key: 'model', value: 'test-model' }), + ]); + await renderHarness(); + + await act(async () => { + await api!.run(runReq(taskId)); + }); + + const task = taskOf(taskId); + expect(task.status).toBe('error'); + const truncation = task.log.find((l) => l.text === '连接中断,未收到结束事件'); + expect(truncation?.kind).toBe('error'); + // streamed content still landed (delta pass-through + final flush) + expect(task.html).toBe('

partial'); + expect(task.stats.model).toBe('test-model'); + // no terminal → the diff-edit baseline must stay untouched + expect(task.baseHtml).toBeUndefined(); + expect(task.log.some((l) => l.kind === 'done')).toBe(false); + }); + + it('finishes done and commits the diff-edit baseline when the stream ends with done', async () => { + const taskId = useStore.getState().newTask({ name: 'happy-path' }); + stubStreamFetch([ + sseFrame('start', { bin: '/usr/bin/agent', promptBytes: 12 }), + sseFrame('delta', { text: '

hello' }), + sseFrame('meta', { key: 'model', value: 'test-model' }), + sseFrame('done', { code: 0 }), + ]); + await renderHarness(); + + await act(async () => { + await api!.run(runReq(taskId)); + }); + + const task = taskOf(taskId); + expect(task.status).toBe('done'); + expect(task.log.some((l) => l.kind === 'done' && l.text.includes('agent 进程退出'))).toBe(true); + expect(task.html).toBe('

hello'); + expect(task.stats.model).toBe('test-model'); + expect(task.stats.endedAt).toBeDefined(); + expect(task.baseHtml).toBe(task.html); + }); + + it('lands buffered events in order before the 已取消 log when cancelled mid-stream', async () => { + const taskId = useStore.getState().newTask({ name: 'cancel' }); + const probe = stubStreamFetch( + [ + sseFrame('meta', { key: 'model', value: 'test-model' }), + sseFrame('meta', { key: 'session', value: 's1' }), + ], + { hangUntilAbort: true }, + ); + await renderHarness(); + + let runP: Promise | undefined; + await act(async () => { + runP = api!.run(runReq(taskId)); + }); + await vi.waitFor(() => expect(probe.isHanging()).toBe(true)); + + act(() => api!.cancel(taskId)); + await act(async () => { + await runP; + }); + + const task = taskOf(taskId); + expect(task.status).toBe('idle'); + const texts = task.log.map((l) => l.text); + const modelIdx = texts.indexOf('model = test-model'); + const sessionIdx = texts.indexOf('session = s1'); + const cancelIdx = texts.indexOf('已取消'); + expect(modelIdx).toBeGreaterThanOrEqual(0); + expect(sessionIdx).toBeGreaterThan(modelIdx); + expect(cancelIdx).toBeGreaterThan(sessionIdx); + }); +}); diff --git a/next/src/lib/use-convert.ts b/next/src/lib/use-convert.ts index 061da5f9..923025ac 100644 --- a/next/src/lib/use-convert.ts +++ b/next/src/lib/use-convert.ts @@ -20,6 +20,55 @@ const DIFF_LOG_PREFIX = "🔁 diff-edit 模式"; // per-task abort controllers — multiple tasks can stream concurrently const controllers = new Map(); +// SSE events are buffered and handed to handleEvent in 100ms windows. +// Inside a flush the store updates run back-to-back with no await between +// them, so React 19 automatic batching turns the whole window into one +// render — the read loop's awaited reader.read() would otherwise give every +// single event its own full-list re-render, and fast reasoning streams peak +// around 240 events/s. The first delta, done and error bypass the window so +// TTFB and terminal state are not delayed; each bypass drains the buffer +// first to keep event order. +const BATCH_WINDOW_MS = 100; + +export function createEventBatcher(onEvent: (event: string, data: unknown) => void) { + const pending: Array<{ event: string; data: unknown }> = []; + let timer: ReturnType | null = null; + let sawFirstDelta = false; + + const drain = () => { + if (timer !== null) { + clearTimeout(timer); + timer = null; + } + const out = pending.splice(0, pending.length); + for (const item of out) onEvent(item.event, item.data); + }; + + return { + push(event: string, data: unknown) { + if (event === 'delta' && !sawFirstDelta) { + sawFirstDelta = true; + drain(); + onEvent(event, data); + return; + } + if (event === 'done' || event === 'error') { + drain(); + onEvent(event, data); + return; + } + pending.push({ event, data }); + // arm only on the empty→non-empty transition so long streams hold + // exactly one timer, not one per event + if (timer === null) timer = setTimeout(drain, BATCH_WINDOW_MS); + }, + flush: drain, + dispose() { + if (timer !== null) clearTimeout(timer); + }, + }; +} + export function useConvert() { const cancel = useCallback((taskId: string) => { const ctl = controllers.get(taskId); @@ -97,6 +146,12 @@ export function useConvert() { : `准备调用 ${req.agent}${useModel ? ` · 模型 ${useModel}` : ""} · 模板 ${req.templateId} · ${sizeNote}`, }); + let sawTerminal = false; + const batch = createEventBatcher((event, data) => { + if (event === 'done' || event === 'error') sawTerminal = true; + handleEvent(taskId, event, data, startedAt); + }); + try { const res = await fetch("/api/convert", { method: "POST", @@ -139,16 +194,30 @@ export function useConvert() { } catch { continue; } - handleEvent(taskId, event, data, startedAt); + batch.push(event, data); } } - const endedAt = Date.now(); - useStore.getState().patchStatsFor(taskId, { endedAt, durationMs: endedAt - startedAt }); - useStore.getState().setStatusFor(taskId, "done"); - // record the just-finished (content, html) as the new diff-edit baseline - // so the user's next edit goes through diff mode instead of full regen - useStore.getState().commitBaseFor(taskId); + batch.flush(); + if (sawTerminal) { + const endedAt = Date.now(); + useStore.getState().patchStatsFor(taskId, { endedAt, durationMs: endedAt - startedAt }); + useStore.getState().setStatusFor(taskId, "done"); + // record the just-finished (content, html) as the new diff-edit baseline + // so the user's next edit goes through diff mode instead of full regen + useStore.getState().commitBaseFor(taskId); + } else if (ctl.signal.aborted) { + // user cancelled mid-stream: cancel() already set the task idle and + // the catch arm below wrote the log line — nothing left to finish + } else { + // stream ended without a done/error event (connection dropped): + // surface the truncation and keep the previous diff-edit baseline + useStore.getState().pushLogFor(taskId, { kind: 'error', text: '连接中断,未收到结束事件' }); + useStore.getState().setStatusFor(taskId, 'error'); + } } catch (err) { + // deliver whatever was already buffered before deciding how to fail — + // those events arrived, they must not vanish + batch.flush(); if ((err as Error)?.name === "AbortError") { useStore.getState().pushLogFor(taskId, { kind: "info", text: "已取消" }); useStore.getState().setStatusFor(taskId, "idle"); @@ -160,6 +229,7 @@ export function useConvert() { }); useStore.getState().setStatusFor(taskId, "error"); } finally { + batch.dispose(); if (controllers.get(taskId) === ctl) controllers.delete(taskId); } }, From 8b172a02ddf02af5dbc3a924242ba10fc5e9f1fe Mon Sep 17 00:00:00 2001 From: ChHsiching Date: Mon, 28 Sep 2026 12:10:08 +0800 Subject: [PATCH 2/4] test: assert bounded re-render count for a high-rate meta burst --- next/src/lib/__tests__/use-convert.test.ts | 36 ++++++++++++++++++++-- 1 file changed, 34 insertions(+), 2 deletions(-) diff --git a/next/src/lib/__tests__/use-convert.test.ts b/next/src/lib/__tests__/use-convert.test.ts index 8ab9c8cc..8dd5afe6 100644 --- a/next/src/lib/__tests__/use-convert.test.ts +++ b/next/src/lib/__tests__/use-convert.test.ts @@ -100,11 +100,21 @@ describe('useConvert().run', () => { let root: Root | null = null; let container: HTMLElement | null = null; let api: ReturnType | null = null; + let renderCount = 0; + + // useConvert is a hook — render a null harness once per test to capture it. + // RenderProbe subscribes to the active task's log length, so renderCount + // tracks how many store batches React committed, not how many events + // arrived — the bounded quantity the batching fix exists to guarantee. + function RenderProbe() { + useStore((s) => s.tasks.find((t) => t.id === s.activeTaskId)?.log.length ?? 0); + renderCount++; + return null; + } - // useConvert is a hook — render a null harness once per test to capture it function Harness() { api = useConvert(); - return null; + return createElement(RenderProbe); } function sseFrame(event: string, data: unknown) { @@ -174,6 +184,7 @@ describe('useConvert().run', () => { beforeEach(() => { localStorage.clear(); + renderCount = 0; container = document.createElement('div'); document.body.appendChild(container); root = createRoot(container); @@ -270,4 +281,25 @@ describe('useConvert().run', () => { expect(sessionIdx).toBeGreaterThan(modelIdx); expect(cancelIdx).toBeGreaterThan(sessionIdx); }); + + it('delivers every event of a high-rate meta burst while keeping re-renders bounded', async () => { + const taskId = useStore.getState().newTask({ name: 'burst' }); + // one network chunk carrying 200 meta frames — the burst shape from the + // issue. Unfixed code re-renders the full task list once per event; the + // batcher must land all 200 in a handful of store batches. + const burst = Array.from({ length: 200 }, (_, i) => + sseFrame('meta', { key: 'model', value: `m${i}` }), + ).join(''); + stubStreamFetch([burst, sseFrame('done', { code: 0 })]); + await renderHarness(); + + await act(async () => { + await api!.run(runReq(taskId)); + }); + + const task = taskOf(taskId); + expect(task.status).toBe('done'); + expect(task.log.filter((l) => l.kind === 'meta')).toHaveLength(200); + expect(renderCount).toBeLessThanOrEqual(10); + }); }); From ac60e27a80c1070cfe7500a87bbb7429a89f9730 Mon Sep 17 00:00:00 2001 From: ChHsiching Date: Mon, 28 Sep 2026 13:05:59 +0800 Subject: [PATCH 3/4] fix: finish error-terminated convert streams as errored, not done --- next/src/lib/__tests__/use-convert.test.ts | 23 ++++++++++++++++++++++ next/src/lib/use-convert.ts | 14 ++++++++++--- 2 files changed, 34 insertions(+), 3 deletions(-) diff --git a/next/src/lib/__tests__/use-convert.test.ts b/next/src/lib/__tests__/use-convert.test.ts index 8dd5afe6..dd37666e 100644 --- a/next/src/lib/__tests__/use-convert.test.ts +++ b/next/src/lib/__tests__/use-convert.test.ts @@ -249,6 +249,29 @@ describe('useConvert().run', () => { expect(task.baseHtml).toBe(task.html); }); + it('marks the task errored without committing a baseline when the stream ends with error', async () => { + const taskId = useStore.getState().newTask({ name: 'error-terminal' }); + stubStreamFetch([ + sseFrame('start', { bin: '/usr/bin/agent', promptBytes: 12 }), + sseFrame('delta', { text: '

partial' }), + sseFrame('error', { message: 'agent binary not found' }), + ]); + await renderHarness(); + + await act(async () => { + await api!.run(runReq(taskId)); + }); + + const task = taskOf(taskId); + expect(task.status).toBe('error'); + expect(task.log.some((l) => l.kind === 'error' && l.text === 'agent binary not found')).toBe(true); + // streamed data stays visible … + expect(task.html).toBe('

partial'); + // … but the partial HTML must not be committed as the diff-edit baseline + expect(task.baseHtml).toBeUndefined(); + expect(task.log.some((l) => l.kind === 'done')).toBe(false); + }); + it('lands buffered events in order before the 已取消 log when cancelled mid-stream', async () => { const taskId = useStore.getState().newTask({ name: 'cancel' }); const probe = stubStreamFetch( diff --git a/next/src/lib/use-convert.ts b/next/src/lib/use-convert.ts index 923025ac..57b91853 100644 --- a/next/src/lib/use-convert.ts +++ b/next/src/lib/use-convert.ts @@ -146,9 +146,11 @@ export function useConvert() { : `准备调用 ${req.agent}${useModel ? ` · 模型 ${useModel}` : ""} · 模板 ${req.templateId} · ${sizeNote}`, }); - let sawTerminal = false; + // which terminal event the stream carried, if any — done and error need + // different finishing (below): a failed run must not read as success + let terminalEvent: "done" | "error" | null = null; const batch = createEventBatcher((event, data) => { - if (event === 'done' || event === 'error') sawTerminal = true; + if (event === "done" || event === "error") terminalEvent = event; handleEvent(taskId, event, data, startedAt); }); @@ -198,13 +200,19 @@ export function useConvert() { } } batch.flush(); - if (sawTerminal) { + if (terminalEvent === "done") { const endedAt = Date.now(); useStore.getState().patchStatsFor(taskId, { endedAt, durationMs: endedAt - startedAt }); useStore.getState().setStatusFor(taskId, "done"); // record the just-finished (content, html) as the new diff-edit baseline // so the user's next edit goes through diff mode instead of full regen useStore.getState().commitBaseFor(taskId); + } else if (terminalEvent === "error") { + // the server emits exactly one error event and closes (unknown agent, + // missing binary, crashed child). Keep whatever streamed for + // debugging and fail the run — a partial document must not become + // the next diff-edit baseline, so commitBaseFor is skipped. + useStore.getState().setStatusFor(taskId, "error"); } else if (ctl.signal.aborted) { // user cancelled mid-stream: cancel() already set the task idle and // the catch arm below wrote the log line — nothing left to finish From 9c5038f203b5dd797b9cf0cd8a490f3d969655e1 Mon Sep 17 00:00:00 2001 From: ChHsiching Date: Mon, 28 Sep 2026 13:33:57 +0800 Subject: [PATCH 4/4] fix: make an observed error win over a trailing done frame --- next/src/lib/__tests__/use-convert.test.ts | 33 ++++++++++++++++++++++ next/src/lib/use-convert.ts | 7 ++++- 2 files changed, 39 insertions(+), 1 deletion(-) diff --git a/next/src/lib/__tests__/use-convert.test.ts b/next/src/lib/__tests__/use-convert.test.ts index dd37666e..daaf9208 100644 --- a/next/src/lib/__tests__/use-convert.test.ts +++ b/next/src/lib/__tests__/use-convert.test.ts @@ -272,6 +272,39 @@ describe('useConvert().run', () => { expect(task.log.some((l) => l.kind === 'done')).toBe(false); }); + it('keeps the previous baseline and error status when error is followed by done', async () => { + // openclaw's close handler emits error for an empty response or a JSON + // parse failure, then an unconditional done — the trailing done must not + // turn the failed run into a success (invoke.ts child close handler) + const taskId = useStore.getState().newTask({ name: 'error-then-done' }); + useStore.setState((st) => ({ + tasks: st.tasks.map((t) => + t.id === taskId ? { ...t, baseContent: 'old content', baseHtml: '

old base

' } : t, + ), + })); + stubStreamFetch([ + sseFrame('start', { bin: '/usr/bin/agent', promptBytes: 12 }), + sseFrame('delta', { text: '

partial' }), + sseFrame('error', { message: 'OpenClaw returned an empty assistant message' }), + sseFrame('done', { code: 0 }), + ]); + await renderHarness(); + + await act(async () => { + await api!.run(runReq(taskId)); + }); + + const task = taskOf(taskId); + expect(task.status).toBe('error'); + expect(task.log.some((l) => l.kind === 'error' && l.text === 'OpenClaw returned an empty assistant message')).toBe(true); + // the done frame still lands in the log … + expect(task.log.some((l) => l.kind === 'done' && l.text.includes('agent 进程退出'))).toBe(true); + // … streamed data stays visible, and the previous baseline survives + expect(task.html).toBe('

partial'); + expect(task.baseHtml).toBe('

old base

'); + expect(task.baseContent).toBe('old content'); + }); + it('lands buffered events in order before the 已取消 log when cancelled mid-stream', async () => { const taskId = useStore.getState().newTask({ name: 'cancel' }); const probe = stubStreamFetch( diff --git a/next/src/lib/use-convert.ts b/next/src/lib/use-convert.ts index 57b91853..995ac896 100644 --- a/next/src/lib/use-convert.ts +++ b/next/src/lib/use-convert.ts @@ -150,7 +150,12 @@ export function useConvert() { // different finishing (below): a failed run must not read as success let terminalEvent: "done" | "error" | null = null; const batch = createEventBatcher((event, data) => { - if (event === "done" || event === "error") terminalEvent = event; + if (event === "done" || event === "error") { + // error wins once seen: openclaw's close handler emits error for an + // empty response or a JSON parse failure and then an unconditional + // done — the trailing done must not mask the failure + if (event === "error" || terminalEvent === null) terminalEvent = event; + } handleEvent(taskId, event, data, startedAt); });