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 = [];
+ },
+ };
+}