mirror of
https://github.com/OpenHands/OpenHands.git
synced 2026-10-06 15:03:43 +08:00
perf: batch StreamingDeltaEvents so the UI keeps up with fast models (#16164)
Co-authored-by: Graham Neubig <neubig@gmail.com>
This commit is contained in:
co-authored by
Graham Neubig
parent
96b6aab34d
commit
e5fc3a4be2
@@ -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(
|
||||
<QueryClientProvider client={queryClient}>
|
||||
<ConversationWebSocketProvider
|
||||
conversationId={conversationId}
|
||||
conversationUrl="http://localhost/api"
|
||||
>
|
||||
<div />
|
||||
</ConversationWebSocketProvider>
|
||||
</QueryClientProvider>,
|
||||
);
|
||||
|
||||
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(
|
||||
<QueryClientProvider client={queryClient}>
|
||||
<ConversationWebSocketProvider
|
||||
conversationId="conv-b"
|
||||
conversationUrl="http://localhost/api"
|
||||
>
|
||||
<div />
|
||||
</ConversationWebSocketProvider>
|
||||
</QueryClientProvider>,
|
||||
);
|
||||
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
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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)", () => {
|
||||
|
||||
@@ -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<number, () => 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);
|
||||
});
|
||||
});
|
||||
@@ -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<StreamingDeltaBatcher | null>(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<StreamingDeltaBatcher | null>(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
|
||||
|
||||
@@ -14,12 +14,8 @@ export interface WebSocketHookOptions {
|
||||
};
|
||||
}
|
||||
|
||||
export const useWebSocket = <T = string>(
|
||||
url: string,
|
||||
options?: WebSocketHookOptions,
|
||||
) => {
|
||||
export const useWebSocket = (url: string, options?: WebSocketHookOptions) => {
|
||||
const [isConnected, setIsConnected] = React.useState(false);
|
||||
const [lastMessage, setLastMessage] = React.useState<T | null>(null);
|
||||
const [error, setError] = React.useState<Error | null>(null);
|
||||
const [isReconnecting, setIsReconnecting] = React.useState(false);
|
||||
const wsRef = React.useRef<WebSocket | null>(null);
|
||||
@@ -67,7 +63,9 @@ export const useWebSocket = <T = string>(
|
||||
};
|
||||
|
||||
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 = <T = string>(
|
||||
|
||||
return {
|
||||
isConnected,
|
||||
lastMessage,
|
||||
error,
|
||||
socket: wsRef.current,
|
||||
sendMessage,
|
||||
|
||||
@@ -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<EventState>()((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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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 = [];
|
||||
},
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user