fix(site): clear stream state atomically with durable message commit (#23924)

This commit is contained in:
Danielle Maywood
2026-04-01 16:38:16 +01:00
committed by GitHub
parent ba734f8b10
commit e3c59c00cd
2 changed files with 318 additions and 26 deletions
@@ -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) => (
<QueryClientProvider client={queryClient}>{children}</QueryClientProvider>
);
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) => (
<QueryClientProvider client={queryClient}>{children}</QueryClientProvider>
);
// 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) => (
<QueryClientProvider client={queryClient}>{children}</QueryClientProvider>
);
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" },
]);
});
});
});
@@ -94,7 +94,6 @@ export const useChatStore = (
const queryClient = useQueryClient();
const [store] = useState(createChatStore);
const streamResetFrameRef = useRef<number | null>(null);
const queuedMessagesHydratedChatIDRef = useRef<string | null>(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);
}