From ffa064557258e323fcf4a0851964c2e25028fabc Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Tue, 21 Jul 2026 15:50:04 -0400 Subject: [PATCH] perf(tui): batch event propagation --- packages/tui/src/context/client.tsx | 12 ++- packages/tui/src/context/event-batcher.ts | 57 ++++++++++++ packages/tui/test/cli/tui/data.test.tsx | 57 ++++++++++++ packages/tui/test/cli/tui/use-event.test.tsx | 48 ++++++++-- .../tui/test/context/event-batcher.test.ts | 89 +++++++++++++++++++ 5 files changed, 256 insertions(+), 7 deletions(-) create mode 100644 packages/tui/src/context/event-batcher.ts create mode 100644 packages/tui/test/context/event-batcher.test.ts diff --git a/packages/tui/src/context/client.tsx b/packages/tui/src/context/client.tsx index 83182f241e6..027fb37558d 100644 --- a/packages/tui/src/context/client.tsx +++ b/packages/tui/src/context/client.tsx @@ -1,8 +1,9 @@ import type { OpenCodeClient, OpenCodeEvent } from "@opencode-ai/client" import { createGlobalEmitter } from "@solid-primitives/event-bus" -import { onCleanup, onMount } from "solid-js" +import { batch, onCleanup, onMount } from "solid-js" import { createStore } from "solid-js/store" import { errorMessage } from "../util/error" +import { createEventBatcher } from "./event-batcher" import { createSimpleContext } from "./helper" import { useLog } from "./log" @@ -61,6 +62,7 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( const cancel = () => request.abort(controller.signal.reason) const timeout = setTimeout(() => request.abort(new Error("Timed out connecting to server")), connectTimeout) controller.signal.addEventListener("abort", cancel, { once: true }) + let queued: ReturnType> | undefined const error = await (async () => { record(attempt === 0 ? "connecting" : "reconnecting", attempt) log.info("event stream connecting", { attempt }) @@ -79,6 +81,11 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( log.info("event stream connected") events.emit(first.value.type, first.value) setConnection({ status: "connected", attempt: 0, error: undefined }) + queued = createEventBatcher((pending) => { + batch(() => { + for (const event of pending) events.emit(event.type, event) + }) + }) while (!abort.signal.aborted && !controller.signal.aborted) { const event = await iterator.next() if (abort.signal.aborted || controller.signal.aborted) return undefined @@ -89,12 +96,13 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( aggregateID: event.value.durable.aggregateID, seq: event.value.durable.seq, }) - events.emit(event.value.type, event.value) + queued.add(event.value) } return undefined })() .catch((error) => error) .finally(() => { + queued?.end(abort.signal.aborted || controller.signal.aborted) request.abort() clearTimeout(timeout) controller.signal.removeEventListener("abort", cancel) diff --git a/packages/tui/src/context/event-batcher.ts b/packages/tui/src/context/event-batcher.ts new file mode 100644 index 00000000000..60cc0b6713e --- /dev/null +++ b/packages/tui/src/context/event-batcher.ts @@ -0,0 +1,57 @@ +const defaultInterval = 16 +const defaultLimit = 1_024 + +type Options = { + interval?: number + limit?: number + now?: () => number + schedule?: (callback: () => void, delay: number) => ReturnType + cancel?: (timer: ReturnType) => void +} + +export function createEventBatcher(onFlush: (events: T[]) => void, options: Options = {}) { + const interval = options.interval ?? defaultInterval + const limit = options.limit ?? defaultLimit + const now = options.now ?? Date.now + const schedule = options.schedule ?? setTimeout + const cancel = options.cancel ?? clearTimeout + let queue: T[] = [] + let timer: ReturnType | undefined + let last = 0 + let ended = false + + function flush() { + if (queue.length === 0) return + const pending = queue + queue = [] + timer = undefined + last = now() + onFlush(pending) + } + + return { + add(event: T) { + if (ended) return + queue.push(event) + if (queue.length >= limit) { + if (timer !== undefined) cancel(timer) + flush() + return + } + if (timer !== undefined) return + if (now() - last >= interval) { + flush() + return + } + timer = schedule(flush, interval) + }, + end(discard: boolean) { + if (ended) return + ended = true + if (timer !== undefined) cancel(timer) + timer = undefined + if (!discard) flush() + queue = [] + }, + } +} diff --git a/packages/tui/test/cli/tui/data.test.tsx b/packages/tui/test/cli/tui/data.test.tsx index d392704dc41..362ae69e052 100644 --- a/packages/tui/test/cli/tui/data.test.tsx +++ b/packages/tui/test/cli/tui/data.test.tsx @@ -792,6 +792,63 @@ test("completes exploration when a queued prompt is promoted", async () => { } }) +test("batches burst event projections into fewer reactive executions", async () => { + const events = createEventStream() + const calls = createFetch(undefined, events) + const sessionID = "session-event-burst" + let client!: ReturnType + let received = 0 + let executions = 0 + + function Probe() { + const data = useData() + client = useClient() + client.event.on("session.input.admitted", () => received++) + createEffect(() => { + data.session.message.list(sessionID).length + executions++ + }) + return + } + + const app = await testRender(() => ( + + + + + + + + + + )) + + try { + await wait(() => client.connection.status() === "connected") + const baseline = executions + for (let index = 0; index < 10; index++) { + emitEvent(events, { + id: `evt_input_${index}`, + created: index, + type: "session.input.admitted", + durable: durable(sessionID, index), + data: { + sessionID, + inputID: `message-${index}`, + input: { type: "user", data: { text: `${index}` }, delivery: "steer" }, + }, + }) + } + + await wait(() => received === 10) + await Bun.sleep(20) + expect(received).toBe(10) + expect(executions - baseline).toBe(2) + } finally { + app.renderer.destroy() + } +}) + test("classifies live tool rows independently of their call ID", async () => { const events = createEventStream() const sessionID = "session-tool-call-id" diff --git a/packages/tui/test/cli/tui/use-event.test.tsx b/packages/tui/test/cli/tui/use-event.test.tsx index c26d2c876b5..e58bc3eb546 100644 --- a/packages/tui/test/cli/tui/use-event.test.tsx +++ b/packages/tui/test/cli/tui/use-event.test.tsx @@ -51,14 +51,12 @@ function update(version: string): OpenCodeEvent { } } -async function mount( - reconnect?: (signal: AbortSignal) => Promise<{ api: OpenCodeClient }>, - log?: LogSink, -) { +async function mount(reconnect?: (signal: AbortSignal) => Promise<{ api: OpenCodeClient }>, log?: LogSink) { const events = createEventStream() const calls = createFetch(undefined, events) const seen: OpenCodeEvent[] = [] const workspaces: Array = [] + const handshakes: string[] = [] let client!: ReturnType let done!: () => void const ready = new Promise((resolve) => { @@ -76,24 +74,27 @@ async function mount( }} seen={seen} workspaces={workspaces} + handshakes={handshakes} /> )) await ready - return { app, events, emit: events.emit, client, seen, workspaces } + return { app, events, emit: (event: OpenCodeEvent) => events.emit(event), client, seen, workspaces, handshakes } } function Probe(props: { seen: OpenCodeEvent[] workspaces: Array + handshakes: string[] onReady: (ctx: { client: ReturnType }) => void }) { const client = useClient() const event = useEvent() onMount(() => { + client.event.on("server.connected", () => props.handshakes.push(client.connection.status())) event.subscribe((evt, { workspace }) => { props.seen.push(evt) props.workspaces.push(workspace) @@ -105,6 +106,35 @@ function Probe(props: { } describe("useEvent", () => { + test("dispatches server.connected immediately", async () => { + const { app, client, handshakes } = await mount() + + try { + await wait(() => client.connection.status() === "connected") + expect(handshakes).toEqual(["connecting"]) + } finally { + app.renderer.destroy() + } + }) + + test("delivers a burst exactly once and in order", async () => { + const { app, client, emit, seen } = await mount() + + try { + await wait(() => client.connection.status() === "connected") + for (const branch of ["one", "two", "three"]) emit(vcs(branch)) + await wait(() => seen.length === 3) + + expect(seen.map((item) => (item.type === "vcs.branch.updated" ? item.data.branch : item.type))).toEqual([ + "one", + "two", + "three", + ]) + } finally { + app.renderer.destroy() + } + }) + test("logs only durable events", async () => { const logs: Array<{ level: LogLevel; message: string; tags: Readonly> }> = [] const { app, emit, seen } = await mount(undefined, (level, message, tags) => { @@ -195,6 +225,9 @@ describe("useEvent", () => { await wait(() => client.connection.status() === "connected") // Reconnection only runs when the stream is down, never while connected. expect(attempts).toEqual([]) + events.emit(event(vcs("before-drop"), { directory: "/tmp/original" })) + await wait(() => seen.some((item) => item.type === "vcs.branch.updated" && item.data.branch === "before-drop")) + events.emit(event(vcs("at-drop"), { directory: "/tmp/original" })) events.disconnect() await wait(() => client.connection.status() === "connected" && attempts.length > 0) replacementEvents.emit(event(vcs("rediscovered"), { directory: "/tmp/rediscovered" })) @@ -202,6 +235,11 @@ describe("useEvent", () => { expect(client.api).toBe(replacement.api) expect(attempts).toEqual([1]) + expect(seen.map((item) => (item.type === "vcs.branch.updated" ? item.data.branch : item.type))).toEqual([ + "before-drop", + "at-drop", + "rediscovered", + ]) const history = client.connection.internal.history() expect(history.map((event) => [event.data.status, event.data.attempt])).toEqual([ ["connecting", 0], diff --git a/packages/tui/test/context/event-batcher.test.ts b/packages/tui/test/context/event-batcher.test.ts new file mode 100644 index 00000000000..aa044158eb3 --- /dev/null +++ b/packages/tui/test/context/event-batcher.test.ts @@ -0,0 +1,89 @@ +import { describe, expect, test } from "bun:test" +import { createEventBatcher } from "../../src/context/event-batcher" + +function clock() { + let time = 100 + const scheduled = new Map, { callback: () => void; at: number }>() + return { + now: () => time, + schedule(callback: () => void, delay: number) { + const timer = setTimeout(() => {}, 60_000) + scheduled.set(timer, { callback, at: time + delay }) + return timer + }, + cancel(timer: ReturnType) { + clearTimeout(timer) + scheduled.delete(timer) + }, + advance(delay: number) { + time += delay + for (const [timer, task] of scheduled) { + if (task.at > time) continue + clearTimeout(timer) + scheduled.delete(timer) + task.callback() + } + }, + pending() { + return scheduled.size + }, + } +} + +describe("createEventBatcher", () => { + test("preserves events in frame-bounded flushes", () => { + const time = clock() + const flushes: number[][] = [] + const batcher = createEventBatcher((events) => flushes.push(events), time) + + batcher.add(1) + time.advance(1) + batcher.add(2) + batcher.add(3) + + expect(flushes).toEqual([[1]]) + expect(time.pending()).toBe(1) + time.advance(15) + expect(flushes).toEqual([[1]]) + time.advance(1) + expect(flushes).toEqual([[1], [2, 3]]) + expect(flushes.flat()).toEqual([1, 2, 3]) + }) + + test("flushes a live generation and discards an obsolete generation", () => { + const time = clock() + const live: number[][] = [] + const active = createEventBatcher((events) => live.push(events), time) + active.add(1) + time.advance(1) + active.add(2) + active.end(false) + + const obsolete: number[][] = [] + const stale = createEventBatcher((events) => obsolete.push(events), time) + stale.add(3) + time.advance(1) + stale.add(4) + stale.end(true) + time.advance(16) + + expect(live).toEqual([[1], [2]]) + expect(obsolete).toEqual([[3]]) + expect(time.pending()).toBe(0) + }) + + test("caps a batch when timers cannot run", () => { + const time = clock() + const flushes: number[][] = [] + const batcher = createEventBatcher((events) => flushes.push(events), { ...time, limit: 3 }) + + batcher.add(1) + time.advance(1) + batcher.add(2) + batcher.add(3) + batcher.add(4) + + expect(flushes).toEqual([[1], [2, 3, 4]]) + expect(time.pending()).toBe(0) + }) +})