From e5fc3a4be206348fe0bcc5b50f4fae0267abc39a Mon Sep 17 00:00:00 2001 From: Vasco Schiavo <115561717+VascoSch92@users.noreply.github.com> Date: Tue, 18 Aug 2026 14:35:30 +0200 Subject: [PATCH] perf: batch StreamingDeltaEvents so the UI keeps up with fast models (#16164) Co-authored-by: Graham Neubig --- .../conversation-websocket-context.test.tsx | 120 +++++++++ __tests__/hooks/use-websocket.test.ts | 65 ++--- __tests__/stores/use-event-store.test.ts | 32 ++- .../utils/streaming-delta-batcher.test.ts | 239 ++++++++++++++++++ .../conversation-websocket-context.tsx | 52 ++++ src/hooks/use-websocket.ts | 11 +- src/stores/use-event-store.ts | 20 +- src/utils/streaming-delta-batcher.ts | 75 ++++++ 8 files changed, 554 insertions(+), 60 deletions(-) create mode 100644 __tests__/utils/streaming-delta-batcher.test.ts create mode 100644 src/utils/streaming-delta-batcher.ts diff --git a/__tests__/contexts/conversation-websocket-context.test.tsx b/__tests__/contexts/conversation-websocket-context.test.tsx index 3c0f39288b..3d544c311e 100644 --- a/__tests__/contexts/conversation-websocket-context.test.tsx +++ b/__tests__/contexts/conversation-websocket-context.test.tsx @@ -17,6 +17,7 @@ import { } from "#/api/conversation-metadata-store"; import type { AppConversation } from "#/api/conversation-service/agent-server-conversation-service.types"; import type { MessageEvent } from "#/types/agent-server/core"; +import { isStreamingDeltaEvent } from "#/types/agent-server/type-guards"; type CapturedWebSocketOptions = { onMessage?: (event: { data: string }) => void; @@ -680,6 +681,125 @@ describe("ConversationWebSocketProvider — conversation-scoped event store", () expect(eventIds()).toHaveLength(2); }); + const makeStreamingDelta = (id: string, content: string) => ({ + id, + timestamp: new Date().toISOString(), + source: "agent", + kind: "StreamingDeltaEvent", + content, + reasoning_content: null, + }); + + const makeAgentMessage = (id: string, text: string): MessageEvent => ({ + id, + timestamp: new Date(Date.now() + 1000).toISOString(), + source: "agent", + llm_message: { role: "assistant", content: [{ type: "text", text }] }, + activated_skills: [], + extended_content: [], + }); + + const renderProviderWithUrl = (conversationId: string) => + render( + + +
+ + , + ); + + it("buffers streaming deltas, then flushes them (reconciled) when the final message arrives", async () => { + renderProviderWithUrl("conv-stream"); + await waitFor(() => expect(wsCapture.mainOnMessage).not.toBeNull()); + await waitFor(() => expect(eventIds()).toEqual(["user-msg-conv-stream"])); + + // Deltas arrive: they are buffered by the batcher, NOT committed per token. + act(() => { + wsCapture.mainOnMessage!({ + data: JSON.stringify(makeStreamingDelta("d1", "I'll help")), + }); + wsCapture.mainOnMessage!({ + data: JSON.stringify(makeStreamingDelta("d2", " with that.")), + }); + }); + expect(eventIds()).toEqual(["user-msg-conv-stream"]); + + // The final agent message is a non-delta event: the handler flushes the + // buffered deltas first, so the message reconciles the streamed text in + // place instead of racing ahead of it. + act(() => { + wsCapture.mainOnMessage!({ + data: JSON.stringify( + makeAgentMessage("agent-final", "I'll help with that. Done."), + ), + }); + }); + + const { uiEvents, eventIds: ids } = useEventStore.getState(); + // One reconciled agent bubble: the canonical final message supersedes the + // flushed deltas, so the streamed text renders once and is never duplicated. + expect(uiEvents).toHaveLength(2); + const bubble = uiEvents[1] as MessageEvent; + expect(bubble.id).toBe("agent-final"); + expect(bubble.llm_message.content).toEqual([ + { type: "text", text: "I'll help with that. Done." }, + ]); + expect(uiEvents.some((event) => isStreamingDeltaEvent(event))).toBe(false); + // eventIds tracks the two durable events, never the deltas. + expect(ids.size).toBe(2); + }); + + it("discards buffered deltas from the previous conversation on switch", async () => { + const { rerender } = renderProviderWithUrl("conv-a"); + await waitFor(() => expect(wsCapture.mainOnMessage).not.toBeNull()); + await waitFor(() => expect(eventIds()).toEqual(["user-msg-conv-a"])); + + // Buffer deltas for A, then switch to B before they flush. + act(() => { + wsCapture.mainOnMessage!({ + data: JSON.stringify(makeStreamingDelta("a1", "STALE")), + }); + }); + rerender( + + +
+ + , + ); + await waitFor(() => expect(eventIds()).toEqual(["user-msg-conv-b"])); + + // B streams and finalizes. If the switch had NOT reset the batcher, A's + // "STALE" delta would still be buffered and merge into B's stream here. + act(() => { + wsCapture.mainOnMessage!({ + data: JSON.stringify(makeStreamingDelta("b1", "fresh")), + }); + wsCapture.mainOnMessage!({ + data: JSON.stringify(makeAgentMessage("agent-b", "fresh.")), + }); + }); + + const { uiEvents, events } = useEventStore.getState(); + expect(uiEvents).toHaveLength(2); + expect((uiEvents[1] as MessageEvent).llm_message.content).toEqual([ + { type: "text", text: "fresh." }, + ]); + // The committed delta carries B's text only — had A's buffer survived the + // switch it would have merged in ahead of it as "STALEfresh". + const committedDeltas = events.filter((event) => + isStreamingDeltaEvent(event), + ); + expect(committedDeltas.map((delta) => delta.content)).toEqual(["fresh"]); + expect(JSON.stringify(events)).not.toContain("STALE"); + }); + it("consumes the optimistic pending bubble when the echoed user message arrives via REST preload", async () => { // Arrange: a cloud start-task conversation left a "Sending…" bubble whose // content matches the first message the server has already persisted. With diff --git a/__tests__/hooks/use-websocket.test.ts b/__tests__/hooks/use-websocket.test.ts index 3f2f4a0de0..5f7f49ec47 100644 --- a/__tests__/hooks/use-websocket.test.ts +++ b/__tests__/hooks/use-websocket.test.ts @@ -56,18 +56,22 @@ describe("useWebSocket", () => { }; it("should establish a WebSocket connection", async () => { - const { result } = renderHook(() => useWebSocket("ws://acme.com/ws")); + const messages: string[] = []; + const { result } = renderHook(() => + useWebSocket("ws://acme.com/ws", { + onMessage: (event) => messages.push(event.data), + }), + ); // Initially should not be connected expect(result.current.isConnected).toBe(false); - expect(result.current.lastMessage).toBe(null); // Wait for connection to be established await waitForConnection(result); - // Should receive the welcome message from our mock + // Should deliver the welcome message from our mock via onMessage await waitFor(() => { - expect(result.current.lastMessage).toBe("Welcome to the WebSocket!"); + expect(messages).toContain("Welcome to the WebSocket!"); }); // Confirm that the WebSocket connection is established when the hook is used @@ -116,8 +120,11 @@ describe("useWebSocket", () => { vi.stubGlobal("WebSocket", MockWebSocket); try { + const messages: string[] = []; const { result, unmount } = renderHook(() => - useWebSocket("ws://acme.com/ws"), + useWebSocket("ws://acme.com/ws", { + onMessage: (event) => messages.push(event.data), + }), ); await waitForConnection(result); @@ -134,7 +141,10 @@ describe("useWebSocket", () => { ); }); - expect(result.current.lastMessage).toBe("third"); + // Every frame is delivered via onMessage, but the hook retains no raw + // message history of its own — not even the latest. + expect(messages).toEqual(["first", "second", "third"]); + expect("lastMessage" in result.current).toBe(false); expect("messages" in result.current).toBe(false); unmount(); @@ -144,32 +154,6 @@ describe("useWebSocket", () => { } }); - it.skip("should handle incoming messages correctly", async () => { - const { result } = renderHook(() => useWebSocket("ws://acme.com/ws")); - - // Wait for connection to be established - await waitFor(() => { - expect(result.current.isConnected).toBe(true); - }); - - // Should receive the welcome message from our mock - await waitFor(() => { - expect(result.current.lastMessage).toBe("Welcome to the WebSocket!"); - }); - - // Send another message from the mock server - wsLink.broadcast("Hello from server!"); - - await waitFor(() => { - expect(result.current.lastMessage).toBe("Hello from server!"); - }); - - // The hook intentionally keeps only the latest message; consumers that - // need durable history should store parsed events in their own domain - // store instead of retaining every raw websocket frame here. - expect("messages" in result.current).toBe(false); - }); - it("should handle connection errors gracefully", async () => { // Create a mock that will simulate an error const errorLink = ws.link("ws://error-test.com/ws"); @@ -474,23 +458,18 @@ describe("useWebSocket", () => { expect(result.current.isConnected).toBe(true); }); - // Should receive the welcome message from our mock - await waitFor(() => { - expect(result.current.lastMessage).toBe("Welcome to the WebSocket!"); - }); - // onMessage handler should have been called for the welcome message - expect(onMessageSpy).toHaveBeenCalledOnce(); + await waitFor(() => { + expect(onMessageSpy).toHaveBeenCalledOnce(); + }); // Send another message from the mock server wsLink.broadcast("Hello from server!"); - await waitFor(() => { - expect(result.current.lastMessage).toBe("Hello from server!"); - }); - // onMessage handler should have been called twice now - expect(onMessageSpy).toHaveBeenCalledTimes(2); + await waitFor(() => { + expect(onMessageSpy).toHaveBeenCalledTimes(2); + }); }); it("should call onError handler when WebSocket encounters an error", async () => { diff --git a/__tests__/stores/use-event-store.test.ts b/__tests__/stores/use-event-store.test.ts index 8517826ac4..09fe92c65c 100644 --- a/__tests__/stores/use-event-store.test.ts +++ b/__tests__/stores/use-event-store.test.ts @@ -178,8 +178,10 @@ describe("useEventStore", () => { content: "hello world", }, ]); - expect(result.current.eventIds.has("delta-1")).toBe(true); - expect(result.current.eventIds.has("delta-2")).toBe(true); + // Transient deltas are never tracked in `eventIds` — copying that Set once + // per token would otherwise be O(n^2). + expect(result.current.eventIds.has("delta-1")).toBe(false); + expect(result.current.eventIds.has("delta-2")).toBe(false); }); it("should compact streaming deltas during bulk add", () => { @@ -196,8 +198,30 @@ describe("useEventStore", () => { id: "delta-1", content: "hello world", }); - expect(result.current.eventIds.has("delta-1")).toBe(true); - expect(result.current.eventIds.has("delta-2")).toBe(true); + // Transient deltas are never tracked in `eventIds`. + expect(result.current.eventIds.has("delta-1")).toBe(false); + expect(result.current.eventIds.has("delta-2")).toBe(false); + }); + + it("should not grow eventIds with the raw streaming-delta count", () => { + const { result } = renderHook(() => useEventStore()); + + act(() => { + result.current.addEvent(mockUserMessageEvent); + for (let i = 0; i < 1000; i += 1) { + result.current.addEvent(makeStreamingDeltaEvent(`delta-${i}`, "x")); + } + }); + + // 1000 deltas collapse to a single event alongside the user message, and + // eventIds tracks only the durable user message — not the deltas. This is + // what keeps the per-token Set copy from going quadratic. + expect(result.current.events).toHaveLength(2); + expect(result.current.eventIds.size).toBe(1); + expect(result.current.eventIds.has(mockUserMessageEvent.id)).toBe(true); + expect( + (result.current.events[1] as StreamingDeltaEvent).content, + ).toHaveLength(1000); }); it("should not compact streaming deltas from different senders (#1656)", () => { diff --git a/__tests__/utils/streaming-delta-batcher.test.ts b/__tests__/utils/streaming-delta-batcher.test.ts new file mode 100644 index 0000000000..25f44f2668 --- /dev/null +++ b/__tests__/utils/streaming-delta-batcher.test.ts @@ -0,0 +1,239 @@ +import { describe, it, expect } from "vitest"; +import { + createStreamingDeltaBatcher, + DeltaFlushScheduler, +} from "#/utils/streaming-delta-batcher"; +import { useEventStore } from "#/stores/use-event-store"; +import { StreamingDeltaEvent } from "#/types/agent-server/core/events/streaming-delta-event"; +import { MessageEvent } from "#/types/agent-server/core"; +import { isStreamingDeltaEvent } from "#/types/agent-server/type-guards"; + +const makeDelta = ( + id: string, + content: string | null, + reasoning: string | null = null, +): StreamingDeltaEvent => ({ + id, + timestamp: "2024-03-01T00:00:00Z", + source: "agent", + kind: "StreamingDeltaEvent", + content, + reasoning_content: reasoning, +}); + +/** + * Deterministic stand-in for `requestAnimationFrame`: callbacks only run when + * the test explicitly `tick()`s a frame, so cadence is fully controlled. + */ +function manualScheduler() { + const callbacks = new Map void>(); + let nextHandle = 1; + const scheduler: DeltaFlushScheduler = { + schedule: (callback) => { + const handle = nextHandle; + nextHandle += 1; + callbacks.set(handle, callback); + return handle; + }, + cancel: (handle) => { + callbacks.delete(handle); + }, + }; + return { + scheduler, + pendingFrames: () => callbacks.size, + tick: () => { + const scheduled = [...callbacks.values()]; + callbacks.clear(); + scheduled.forEach((callback) => callback()); + }, + }; +} + +describe("createStreamingDeltaBatcher", () => { + it("coalesces adjacent deltas into a single commit per frame", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + batcher.enqueue(makeDelta("d1", "Hello")); + batcher.enqueue(makeDelta("d2", ", ")); + batcher.enqueue(makeDelta("d3", "world")); + + // Nothing commits until the frame fires, and three enqueues schedule only + // ONE frame (not one per delta). + expect(commits).toHaveLength(0); + expect(clock.pendingFrames()).toBe(1); + + clock.tick(); + + expect(commits).toHaveLength(1); + expect(commits[0].content).toBe("Hello, world"); + // The coalesced event keeps the first delta's identity. + expect(commits[0].id).toBe("d1"); + }); + + it("merges content and reasoning_content independently, in order", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + batcher.enqueue(makeDelta("d1", "ans", "think-")); + batcher.enqueue(makeDelta("d2", "wer", null)); + batcher.enqueue(makeDelta("d3", null, "more")); + clock.tick(); + + expect(commits).toHaveLength(1); + expect(commits[0].content).toBe("answer"); + expect(commits[0].reasoning_content).toBe("think-more"); + }); + + it("flush() commits synchronously and cancels the scheduled frame", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + batcher.enqueue(makeDelta("d1", "a")); + batcher.enqueue(makeDelta("d2", "b")); + batcher.flush(); + + expect(commits).toHaveLength(1); + expect(commits[0].content).toBe("ab"); + // The pending frame was cancelled, so ticking must not double-commit. + expect(clock.pendingFrames()).toBe(0); + clock.tick(); + expect(commits).toHaveLength(1); + }); + + it("flush() is a no-op when nothing is buffered", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + batcher.flush(); + expect(commits).toHaveLength(0); + }); + + it("reset() drops buffered deltas without committing", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + batcher.enqueue(makeDelta("d1", "lost")); + batcher.reset(); + clock.tick(); + + expect(commits).toHaveLength(0); + expect(clock.pendingFrames()).toBe(0); + }); + + it("preserves text byte-for-byte and order across thousands of 1-char deltas faster than 60Hz", () => { + const commits: StreamingDeltaEvent[] = []; + const clock = manualScheduler(); + const batcher = createStreamingDeltaBatcher( + (delta) => commits.push(delta), + clock.scheduler, + ); + + const total = 5000; + let expected = ""; + for (let i = 0; i < total; i += 1) { + const char = String.fromCharCode(97 + (i % 26)); + expected += char; + batcher.enqueue(makeDelta(`d${i}`, char)); + // A frame only every 100 deltas => deltas arrive far faster than frames. + if (i % 100 === 99) { + clock.tick(); + } + } + batcher.flush(); // boundary flush, as a non-delta event would trigger + + // Commits are bounded by frames, not by provider chunk count. + expect(commits.length).toBeLessThan(total); + expect(commits.length).toBeLessThanOrEqual(total / 100 + 1); + // Concatenating the per-frame batches reproduces the stream exactly (the + // store folds these into one accumulating event by position). + expect(commits.map((delta) => delta.content).join("")).toBe(expected); + }); +}); + +describe("createStreamingDeltaBatcher wired into the event store", () => { + const userMessage: MessageEvent = { + id: "user-1", + timestamp: "2024-02-01T00:00:00Z", + source: "user", + llm_message: { role: "user", content: [{ type: "text", text: "hi" }] }, + activated_skills: [], + extended_content: [], + }; + + it("coalesces deltas across frames, then reconciles into one bubble when the final message arrives", () => { + useEventStore.getState().clearEvents(); + const clock = manualScheduler(); + // Commit into the real store exactly as ConversationWebSocketProvider does. + const batcher = createStreamingDeltaBatcher( + (delta) => useEventStore.getState().addEvent(delta), + clock.scheduler, + ); + + useEventStore.getState().addEvent(userMessage); + + // Stream one char per delta, flushing a frame only every 5 chars, so deltas + // arrive faster than frames — the case where the UI used to fall behind. + const streamed = "I'll start working on that."; + [...streamed].forEach((char, i) => { + batcher.enqueue(makeDelta(`d${i}`, char)); + if (i % 5 === 4) { + clock.tick(); + } + }); + + // A non-delta event (the final agent message) arrives. The provider flushes + // buffered deltas first, so the durable message can never overtake its own + // streamed text. + batcher.flush(); + const finalMessage: MessageEvent = { + id: "agent-1", + timestamp: "2024-04-01T00:00:00Z", + source: "agent", + llm_message: { + role: "assistant", + content: [{ type: "text", text: "I'll start working on that. Done." }], + }, + activated_skills: [], + extended_content: [], + }; + useEventStore.getState().addEvent(finalMessage); + + const state = useEventStore.getState(); + // The user message plus a single reconciled agent bubble — the canonical + // final message supersedes the streamed deltas rather than duplicating them. + expect(state.uiEvents).toHaveLength(2); + const bubble = state.uiEvents[1] as MessageEvent; + expect(bubble.id).toBe("agent-1"); + expect(bubble.llm_message.content).toEqual([ + { type: "text", text: "I'll start working on that. Done." }, + ]); + // No provisional delta survives, so the streamed text renders exactly once. + expect(state.uiEvents.some((event) => isStreamingDeltaEvent(event))).toBe( + false, + ); + // eventIds tracks only the two durable events, never the 27 deltas. + expect(state.eventIds.size).toBe(2); + }); +}); diff --git a/src/contexts/conversation-websocket-context.tsx b/src/contexts/conversation-websocket-context.tsx index bae553c18e..7b7f5ed57c 100644 --- a/src/contexts/conversation-websocket-context.tsx +++ b/src/contexts/conversation-websocket-context.tsx @@ -38,8 +38,13 @@ import { isBrowserNavigateActionEvent, isSwitchLLMObservationEvent, isCanvasUIActionEvent, + isStreamingDeltaEvent, isLaunchChildConversationActionEvent, } from "#/types/agent-server/type-guards"; +import { + createStreamingDeltaBatcher, + StreamingDeltaBatcher, +} from "#/utils/streaming-delta-batcher"; import { handleCanvasUIAction } from "#/services/canvas-ui"; import { handleLaunchChildConversationAction } from "#/services/child-conversation-launch"; import { ConversationStateUpdateEventStats } from "#/types/agent-server/core/events/conversation-state-event"; @@ -154,6 +159,26 @@ export function ConversationWebSocketProvider({ const { appendInput, appendOutput } = useCommandStore(); const resetBrowserStore = useBrowserStore((state) => state.reset); + // Coalesce streaming deltas to ≤1 store commit/render per frame. + // Separate batchers keep the main and planning streams from ever merging. + const mainDeltaBatcherRef = useRef(null); + if (mainDeltaBatcherRef.current === null) { + mainDeltaBatcherRef.current = createStreamingDeltaBatcher((delta) => { + useEventStore.getState().addEvent(delta); + // A delta means connectivity recovered — mirror handleNonErrorEvent. + useErrorMessageStore.getState().clearConnectionError(); + }); + } + const planningDeltaBatcherRef = useRef(null); + if (planningDeltaBatcherRef.current === null) { + planningDeltaBatcherRef.current = createStreamingDeltaBatcher((delta) => { + useEventStore + .getState() + .addEvent({ ...delta, isFromPlanningAgent: true }); + useErrorMessageStore.getState().clearConnectionError(); + }); + } + // History loading state. // - Main conversation history is now loaded via REST (`useConversationHistory`), // so its loading state mirrors the REST query state (see below). @@ -498,6 +523,17 @@ export function ConversationWebSocketProvider({ latestPlanningFileEventRef.current = null; }, [conversationId]); + // Drop buffered deltas on conversation switch/unmount: the store is cleared on + // switch, so flushing them would leak into the next conversation. + useEffect(() => { + const mainBatcher = mainDeltaBatcherRef.current; + const planningBatcher = planningDeltaBatcherRef.current; + return () => { + mainBatcher?.reset(); + planningBatcher?.reset(); + }; + }, [conversationId]); + // Merged loading history state - true if either connection is still loading const isLoadingHistory = useMemo( () => isLoadingHistoryMain || isLoadingHistoryPlanning, @@ -515,6 +551,14 @@ export function ConversationWebSocketProvider({ // Use type guard to validate v1 event structure if (isAgentServerEvent(event)) { + // Buffer deltas; nothing else in this handler applies to them. + if (isStreamingDeltaEvent(event)) { + mainDeltaBatcherRef.current?.enqueue(event); + return; + } + // Flush buffered deltas before this event so it can't overtake them. + mainDeltaBatcherRef.current?.flush(); + // A reconnect replays the backlog from a stale anchor. The store // dedups by id, but the side-effects below aren't idempotent, so skip // them for replayed events (#1656). @@ -739,6 +783,14 @@ export function ConversationWebSocketProvider({ // Use type guard to validate v1 event structure if (isAgentServerEvent(event)) { + // Buffer deltas (the commit re-applies the planning flag). + if (isStreamingDeltaEvent(event)) { + planningDeltaBatcherRef.current?.enqueue(event); + return; + } + // Flush buffered deltas before this event so it can't overtake them. + planningDeltaBatcherRef.current?.flush(); + // Skip non-idempotent side-effects for replayed events, as in the // main handler (#1656). const isDuplicateEvent = useEventStore diff --git a/src/hooks/use-websocket.ts b/src/hooks/use-websocket.ts index dd63a5f530..66516207ea 100644 --- a/src/hooks/use-websocket.ts +++ b/src/hooks/use-websocket.ts @@ -14,12 +14,8 @@ export interface WebSocketHookOptions { }; } -export const useWebSocket = ( - url: string, - options?: WebSocketHookOptions, -) => { +export const useWebSocket = (url: string, options?: WebSocketHookOptions) => { const [isConnected, setIsConnected] = React.useState(false); - const [lastMessage, setLastMessage] = React.useState(null); const [error, setError] = React.useState(null); const [isReconnecting, setIsReconnecting] = React.useState(false); const wsRef = React.useRef(null); @@ -67,7 +63,9 @@ export const useWebSocket = ( }; ws.onmessage = (event) => { - setLastMessage(event.data); + // Deliberately no `lastMessage` state here: nothing reads it, and a + // React state write per frame re-renders this hook's owner on every + // streamed token. Consumers subscribe via `onMessage`. optionsRef.current?.onMessage?.(event); }; @@ -207,7 +205,6 @@ export const useWebSocket = ( return { isConnected, - lastMessage, error, socket: wsRef.current, sendMessage, diff --git a/src/stores/use-event-store.ts b/src/stores/use-event-store.ts index ca4132d173..cbfa0f800b 100644 --- a/src/stores/use-event-store.ts +++ b/src/stores/use-event-store.ts @@ -90,14 +90,19 @@ export interface EventState { } const appendEvent = (state: EventState, event: OHEvent): EventState => { - // Deduplicate: skip if event with same id already exists (O(1) lookup) const eventId = getEventId(event); - if (eventId !== undefined && state.eventIds.has(eventId)) { + // Transient deltas merge by position and are never persisted/resent, so skip + // id tracking for them — copying the growing `eventIds` Set per token would + // otherwise be O(n^2). + const isDelta = isStreamingDeltaEvent(event); + + // Deduplicate: skip if event with same id already exists (O(1) lookup) + if (!isDelta && eventId !== undefined && state.eventIds.has(eventId)) { return state; } const newEventIds = - eventId !== undefined + !isDelta && eventId !== undefined ? new Set(state.eventIds).add(eventId) : state.eventIds; @@ -105,7 +110,7 @@ const appendEvent = (state: EventState, event: OHEvent): EventState => { const lastEvent = state.events[lastEventIndex]; const shouldMergeStreamingDelta = lastEvent && - isStreamingDeltaEvent(event) && + isDelta && isStreamingDeltaEvent(lastEvent) && isSameStreamingSender(event, lastEvent); const events = [...state.events]; @@ -162,11 +167,14 @@ export const useEventStore = create()((set) => ({ for (const event of incoming) { const eventId = getEventId(event); - const isDuplicate = eventId !== undefined && eventIds.has(eventId); + // See `appendEvent`: transient deltas are not tracked in `eventIds`. + const isDelta = isStreamingDeltaEvent(event); + const isDuplicate = + !isDelta && eventId !== undefined && eventIds.has(eventId); if (!isDuplicate) { added = true; - if (eventId !== undefined) { + if (!isDelta && eventId !== undefined) { eventIds.add(eventId); } diff --git a/src/utils/streaming-delta-batcher.ts b/src/utils/streaming-delta-batcher.ts new file mode 100644 index 0000000000..f0231f8027 --- /dev/null +++ b/src/utils/streaming-delta-batcher.ts @@ -0,0 +1,75 @@ +import { StreamingDeltaEvent } from "#/types/agent-server/core/events/streaming-delta-event"; +import { mergeStreamingDeltaEvent } from "#/utils/handle-event-for-ui"; + +/** Schedules a single deferred callback (defaults to the animation frame). */ +export interface DeltaFlushScheduler { + schedule: (callback: () => void) => number; + cancel: (handle: number) => void; +} + +const defaultScheduler: DeltaFlushScheduler = + typeof requestAnimationFrame === "function" + ? { + schedule: (callback) => requestAnimationFrame(callback), + cancel: (handle) => cancelAnimationFrame(handle), + } + : { + schedule: (callback) => setTimeout(callback, 16) as unknown as number, + cancel: (handle) => clearTimeout(handle), + }; + +export interface StreamingDeltaBatcher { + /** Buffer a delta; a flush is scheduled for the next frame if not already. */ + enqueue: (event: StreamingDeltaEvent) => void; + /** Commit buffered deltas now. Call before any non-delta event. */ + flush: () => void; + /** Drop buffered deltas without committing. Call on unmount / conversation switch. */ + reset: () => void; +} + +/** + * Coalesces adjacent `StreamingDeltaEvent`s and commits them at most once per + * animation frame, so a fast model can't force a store commit + re-render per + * token. Callers MUST `flush()` before any non-delta event so a + * durable message/action can't render ahead of its own streamed text. + */ +export function createStreamingDeltaBatcher( + commit: (event: StreamingDeltaEvent) => void, + scheduler: DeltaFlushScheduler = defaultScheduler, +): StreamingDeltaBatcher { + let pending: StreamingDeltaEvent[] = []; + let frame: number | null = null; + + const cancelFrame = () => { + if (frame !== null) { + scheduler.cancel(frame); + frame = null; + } + }; + + const flush = () => { + cancelFrame(); + if (pending.length === 0) { + return; + } + const batch = pending; + pending = []; + commit( + batch.reduce((merged, delta) => mergeStreamingDeltaEvent(delta, merged)), + ); + }; + + return { + enqueue: (event) => { + pending.push(event); + if (frame === null) { + frame = scheduler.schedule(flush); + } + }, + flush, + reset: () => { + cancelFrame(); + pending = []; + }, + }; +}