From e3c59c00cdc35ce55673d808d0119de090bf77d2 Mon Sep 17 00:00:00 2001 From: Danielle Maywood Date: Wed, 1 Apr 2026 16:38:16 +0100 Subject: [PATCH] fix(site): clear stream state atomically with durable message commit (#23924) --- .../ChatConversation/chatStore.test.tsx | 285 ++++++++++++++++++ .../ChatConversation/useChatStore.ts | 59 ++-- 2 files changed, 318 insertions(+), 26 deletions(-) diff --git a/site/src/pages/AgentsPage/components/ChatConversation/chatStore.test.tsx b/site/src/pages/AgentsPage/components/ChatConversation/chatStore.test.tsx index e9c7273b26..ace1f20e95 100644 --- a/site/src/pages/AgentsPage/components/ChatConversation/chatStore.test.tsx +++ b/site/src/pages/AgentsPage/components/ChatConversation/chatStore.test.tsx @@ -38,6 +38,7 @@ import type { OneWayMessageEvent } from "#/utils/OneWayWebSocket"; import { selectChatStatus, selectIsAwaitingFirstStreamChunk, + selectMessagesByID, selectOrderedMessageIDs, selectQueuedMessages, selectReconnectState, @@ -3677,3 +3678,287 @@ describe("updateSidebarChat via stream events", () => { }); }); }); + +describe("stream-to-durable transition (Bug 1)", () => { + it("does not render both stream state and durable message after assistant message commits", async () => { + immediateAnimationFrame(); + + const chatID = "chat-b1-overlap"; + const userMsg = makeMessage(chatID, 1, "user", "hello"); + const mockSocket = createMockSocket(); + mockWatchChatReturn(mockSocket); + + const queryClient = createTestQueryClient(); + const wrapper = ({ children }: PropsWithChildren) => ( + {children} + ); + + const { result } = renderHook( + () => { + const { store } = useChatStore({ + chatID, + chatMessages: [userMsg], + chatRecord: makeChat(chatID), + chatMessagesData: { + messages: [userMsg], + queued_messages: [], + has_more: false, + }, + chatQueuedMessages: [], + setChatErrorReason: vi.fn(), + clearChatErrorReason: vi.fn(), + }); + return { + streamState: useChatSelector(store, selectStreamState), + orderedIDs: useChatSelector(store, selectOrderedMessageIDs), + messagesByID: useChatSelector(store, selectMessagesByID), + }; + }, + { wrapper }, + ); + + await waitFor(() => { + expect(watchChat).toHaveBeenCalledWith(chatID, 1); + }); + + // Build up streaming content. + act(() => { + mockSocket.emitData({ + type: "message_part", + chat_id: chatID, + message_part: { + role: "assistant", + part: { type: "text", text: "response" }, + }, + }); + }); + + await waitFor(() => { + expect(result.current.streamState?.blocks).toEqual([ + { type: "response", text: "response" }, + ]); + }); + + // Commit the assistant message as durable. With the old + // code, streamState stayed non-null here because + // clearStreamState was deferred to a rAF. Both the + // durable message and the stream content coexisted, + // causing duplicate rendering. + act(() => { + mockSocket.emitData({ + type: "message", + chat_id: chatID, + message: makeMessage(chatID, 2, "assistant", "response"), + }); + }); + + // The durable message must be present AND streamState + // must be null in the same snapshot. + await waitFor(() => { + expect(result.current.orderedIDs).toContain(2); + expect(result.current.messagesByID.get(2)?.role).toBe("assistant"); + expect(result.current.streamState).toBeNull(); + }); + }); + + it("no snapshot ever has both durable assistant and stream state", async () => { + immediateAnimationFrame(); + + const chatID = "chat-b1-atomic"; + const userMsg = makeMessage(chatID, 1, "user", "hi"); + const mockSocket = createMockSocket(); + mockWatchChatReturn(mockSocket); + + const queryClient = createTestQueryClient(); + const wrapper = ({ children }: PropsWithChildren) => ( + {children} + ); + + // Track every snapshot emitted to subscribers. + const snapshots: Array<{ + hasStream: boolean; + hasDurableAssistant: boolean; + }> = []; + + const { result } = renderHook( + () => { + const { store } = useChatStore({ + chatID, + chatMessages: [userMsg], + chatRecord: makeChat(chatID), + chatMessagesData: { + messages: [userMsg], + queued_messages: [], + has_more: false, + }, + chatQueuedMessages: [], + setChatErrorReason: vi.fn(), + clearChatErrorReason: vi.fn(), + }); + const streamState = useChatSelector(store, selectStreamState); + const messagesByID = useChatSelector(store, selectMessagesByID); + const hasDurableAssistant = Array.from(messagesByID.values()).some( + (m) => m.role === "assistant", + ); + + snapshots.push({ + hasStream: streamState !== null, + hasDurableAssistant, + }); + + return { streamState, messagesByID }; + }, + { wrapper }, + ); + + await waitFor(() => { + expect(watchChat).toHaveBeenCalledWith(chatID, 1); + }); + + act(() => { + mockSocket.emitData({ + type: "message_part", + chat_id: chatID, + message_part: { + role: "assistant", + part: { type: "text", text: "hello" }, + }, + }); + }); + + await waitFor(() => { + expect(result.current.streamState).not.toBeNull(); + }); + + // Clear snapshot history before the critical transition. + snapshots.length = 0; + + act(() => { + mockSocket.emitData({ + type: "message", + chat_id: chatID, + message: makeMessage(chatID, 2, "assistant", "hello"), + }); + }); + + await waitFor(() => { + expect(result.current.messagesByID.has(2)).toBe(true); + }); + + // No snapshot should ever have BOTH a durable assistant + // message AND non-null stream state. + const overlapping = snapshots.filter( + (s) => s.hasStream && s.hasDurableAssistant, + ); + expect(overlapping).toEqual([]); + }); +}); + +describe("partsBuf cleanup on reconnect (Bug 2)", () => { + it("discards stale buffered parts when the socket reconnects", async () => { + immediateAnimationFrame(); + + const chatID = "chat-b2-reconnect"; + const userMsg = makeMessage(chatID, 1, "user", "test"); + const mockSocket1 = createMockSocket(); + mockWatchChatReturnOnce(mockSocket1); + + const queryClient = createTestQueryClient(); + const wrapper = ({ children }: PropsWithChildren) => ( + {children} + ); + + const { result } = renderHook( + () => { + const { store } = useChatStore({ + chatID, + chatMessages: [userMsg], + chatRecord: makeChat(chatID), + chatMessagesData: { + messages: [userMsg], + queued_messages: [], + has_more: false, + }, + chatQueuedMessages: [], + setChatErrorReason: vi.fn(), + clearChatErrorReason: vi.fn(), + }); + return { + streamState: useChatSelector(store, selectStreamState), + }; + }, + { wrapper }, + ); + + await waitFor(() => { + expect(watchChat).toHaveBeenCalledWith(chatID, 1); + }); + + // Stream a message_part on the first socket. + act(() => { + mockSocket1.emitData({ + type: "message_part", + chat_id: chatID, + message_part: { + role: "assistant", + part: { type: "text", text: "stale content" }, + }, + }); + }); + + await waitFor(() => { + expect(result.current.streamState?.blocks).toEqual([ + { type: "response", text: "stale content" }, + ]); + }); + + // Disconnect. The reconnecting websocket utility + // schedules a reconnect after a 1s delay. + act(() => { + mockSocket1.emitError(); + }); + + // Prepare the second socket and wait for the reconnect + // timer to fire (real timers, ~1s). + const mockSocket2 = createMockSocket(); + mockWatchChatReturnOnce(mockSocket2); + + await waitFor( + () => { + expect(watchChat).toHaveBeenCalledTimes(2); + }, + { timeout: 3_000 }, + ); + + // Open the new socket. This should clear stale state + // including any buffered parts from socket1. + act(() => { + mockSocket2.emitOpen(); + }); + + await waitFor(() => { + expect(result.current.streamState).toBeNull(); + }); + + // Stream fresh content on the new socket. + act(() => { + mockSocket2.emitData({ + type: "message_part", + chat_id: chatID, + message_part: { + role: "assistant", + part: { type: "text", text: "fresh content" }, + }, + }); + }); + + // The stream should show only the new content, not a + // mix of stale + fresh. With the old code, stale parts + // from socket1 could leak into the new stream. + await waitFor(() => { + expect(result.current.streamState?.blocks).toEqual([ + { type: "response", text: "fresh content" }, + ]); + }); + }); +}); diff --git a/site/src/pages/AgentsPage/components/ChatConversation/useChatStore.ts b/site/src/pages/AgentsPage/components/ChatConversation/useChatStore.ts index 6721aa67e0..63db88d184 100644 --- a/site/src/pages/AgentsPage/components/ChatConversation/useChatStore.ts +++ b/site/src/pages/AgentsPage/components/ChatConversation/useChatStore.ts @@ -94,7 +94,6 @@ export const useChatStore = ( const queryClient = useQueryClient(); const [store] = useState(createChatStore); - const streamResetFrameRef = useRef(null); const queuedMessagesHydratedChatIDRef = useRef(null); // Tracks whether the WebSocket has delivered a queue_update for the // current chat. When true, the stream is the authoritative source @@ -240,22 +239,6 @@ export const useChatStore = ( }); }; - const cancelScheduledStreamReset = () => { - if (streamResetFrameRef.current === null) { - return; - } - window.cancelAnimationFrame(streamResetFrameRef.current); - streamResetFrameRef.current = null; - }; - - const scheduleStreamReset = () => { - cancelScheduledStreamReset(); - streamResetFrameRef.current = window.requestAnimationFrame(() => { - store.clearStreamState(); - streamResetFrameRef.current = null; - }); - }; - const updateChatQueuedMessages = ( queuedMessages: readonly TypesGen.ChatQueuedMessage[] | undefined, ) => { @@ -336,7 +319,6 @@ export const useChatStore = ( }); }; - cancelScheduledStreamReset(); store.resetTransientState(); activeChatIDRef.current = chatID ?? null; @@ -366,7 +348,6 @@ export const useChatStore = ( if (partsFlushTimer !== null || partsBuf.length === 0) { return; } - cancelScheduledStreamReset(); partsFlushTimer = setTimeout(() => { partsFlushTimer = null; if (disposed || activeChatIDRef.current !== chatID) { @@ -391,7 +372,6 @@ export const useChatStore = ( clearTimeout(partsFlushTimer); partsFlushTimer = null; } - cancelScheduledStreamReset(); const parts = partsBuf.splice(0); if (activeChatIDRef.current !== chatID || !shouldApplyMessagePart()) { return; @@ -452,7 +432,6 @@ export const useChatStore = ( const part = streamEvent.message_part?.part; if (part) { store.clearRetryState(); - cancelScheduledStreamReset(); partsBuf.push(part); } continue; @@ -568,7 +547,7 @@ export const useChatStore = ( } } - // Schedule a rAF-coalesced flush for any remaining + // Schedule a coalesced flush for any remaining // parts. If parts were already flushed by a // non-message_part event above, this is a no-op. schedulePartsFlush(); @@ -579,10 +558,33 @@ export const useChatStore = ( store.upsertDurableMessages(pendingMessages); upsertCacheMessages(pendingMessages); } + + // Clear stream state atomically with the durable + // message commit so subscribers never see a + // snapshot where both the committed message and + // the streaming output coexist. Previously this + // was deferred to a requestAnimationFrame, which + // left a window where ConversationTimeline and + // LiveStreamTail rendered the same content. + if (needsStreamReset) { + store.clearStreamState(); + // If more message_part events arrived in this + // batch after the durable message, they belong + // to the next turn. Apply them immediately so + // the stream transitions from the old turn to + // the new one without a flash. + if (partsBuf.length > 0) { + if (partsFlushTimer !== null) { + clearTimeout(partsFlushTimer); + partsFlushTimer = null; + } + const nextParts = partsBuf.splice(0); + if (shouldApplyMessagePart()) { + store.applyMessageParts(nextParts); + } + } + } }); - if (needsStreamReset) { - scheduleStreamReset(); - } }; const disposeSocket = createReconnectingWebSocket({ connect() { @@ -599,6 +601,12 @@ export const useChatStore = ( // partial output or failures do not leak into the new // stream. store.resetTransportReplayState(); + // Drain any message parts buffered from the + // previous socket. Without this, a pending + // flush timer could fire after reconnect and + // apply stale parts from the old connection + // into the fresh stream state. + discardBufferedParts(); }, onDisconnect( reconnectState: import("#/utils/reconnectingWebSocket").ReconnectSchedule, @@ -616,7 +624,6 @@ export const useChatStore = ( return () => { disposed = true; disposeSocket(); - cancelScheduledStreamReset(); if (partsFlushTimer !== null) { clearTimeout(partsFlushTimer); }