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);
}