diff --git a/webui/src/hooks/useNanobotStream.ts b/webui/src/hooks/useNanobotStream.ts index 9b1d7ac2f..dc6f0cc01 100644 --- a/webui/src/hooks/useNanobotStream.ts +++ b/webui/src/hooks/useNanobotStream.ts @@ -643,6 +643,28 @@ export function useNanobotStream( return () => document.removeEventListener("visibilitychange", flushOnReturn); }, [flushPendingStreamEvents]); + useEffect(() => { + if (!chatId) return; + return client.onRunStatus((runChatId, startedAt) => { + if (runChatId !== chatId) return; + if (startedAt !== null) { + setRunStartedAt(startedAt); + setIsStreaming(true); + return; + } + flushPendingStreamEvents(); + buffer.current = null; + activeAssistantRef.current = null; + closedAssistantStreamIdsRef.current.clear(); + clearActivitySegment(); + setMessages((prev) => prev.map((message) => ( + message.isStreaming ? { ...message, isStreaming: false } : message + ))); + setRunStartedAt(null); + setIsStreaming(false); + }); + }, [chatId, client, clearActivitySegment, flushPendingStreamEvents]); + // Reset local state when switching chats. Do not reset on every // ``initialMessages`` update: a brand-new chat can receive an empty/404 // history response after the optimistic first message has already rendered. diff --git a/webui/src/tests/thread-shell.test.tsx b/webui/src/tests/thread-shell.test.tsx index 7a16aef9b..ed75437c3 100644 --- a/webui/src/tests/thread-shell.test.tsx +++ b/webui/src/tests/thread-shell.test.tsx @@ -21,6 +21,7 @@ function makeClient() { (modelName: string | null, modelPreset?: string | null) => void >(); const sessionUpdateHandlers = new Set<(chatId: string, scope?: string) => void>(); + const runStatusHandlers = new Set<(chatId: string, startedAt: number | null) => void>(); const runStartedAtByChatId = new Map(); const runGenerationByChatId = new Map(); const latestRunTurnIdByChatId = new Map(); @@ -98,6 +99,13 @@ function makeClient() { statusHandlers.delete(handler); }; }, + onRunStatus: (handler: (chatId: string, startedAt: number | null) => void) => { + runStatusHandlers.add(handler); + for (const [chatId, startedAt] of runStartedAtByChatId) handler(chatId, startedAt); + return () => { + runStatusHandlers.delete(handler); + }; + }, onRuntimeModelUpdate: ( handler: (modelName: string | null, modelPreset?: string | null) => void, ) => { @@ -157,11 +165,13 @@ function makeClient() { ) { advanceRunGeneration(chatId, ev.turn_id); runStartedAtByChatId.set(chatId, ev.started_at); + for (const h of runStatusHandlers) h(chatId, ev.started_at); } else if ( (ev.event === "goal_status" && ev.status === "idle") || ev.event === "turn_end" ) { runStartedAtByChatId.delete(chatId); + for (const h of runStatusHandlers) h(chatId, null); } if (ev.event === "goal_state") { goalStateByChatId.set(chatId, ev.goal_state); diff --git a/webui/src/tests/useNanobotStream.test.tsx b/webui/src/tests/useNanobotStream.test.tsx index 85c4056f2..c3a2d3931 100644 --- a/webui/src/tests/useNanobotStream.test.tsx +++ b/webui/src/tests/useNanobotStream.test.tsx @@ -70,6 +70,7 @@ function normalizeProjection(messages: UIMessage[]): Array void>>(); const statusHandlers = new Set<(status: ConnectionStatus) => void>(); + const runStatusHandlers = new Set<(chatId: string, startedAt: number | null) => void>(); const errorHandlers = new Set<(error: StreamError) => void>(); const runStartedAtByChatId = new Map(); const unsettledRunByChatId = new Map(); @@ -111,6 +112,11 @@ function fakeClient() { handler(status); return () => statusHandlers.delete(handler); }, + onRunStatus(handler: (chatId: string, startedAt: number | null) => void) { + runStatusHandlers.add(handler); + for (const [chatId, startedAt] of runStartedAtByChatId) handler(chatId, startedAt); + return () => runStatusHandlers.delete(handler); + }, onError(handler: (error: StreamError) => void) { errorHandlers.add(handler); return () => errorHandlers.delete(handler); @@ -154,6 +160,11 @@ function fakeClient() { status = nextStatus; statusHandlers.forEach((handler) => handler(status)); }, + emitRunStatus(chatId: string, startedAt: number | null) { + if (startedAt === null) runStartedAtByChatId.delete(chatId); + else runStartedAtByChatId.set(chatId, startedAt); + runStatusHandlers.forEach((handler) => handler(chatId, startedAt)); + }, emitError(error: StreamError) { errorHandlers.forEach((handler) => handler(error)); }, @@ -327,6 +338,39 @@ describe("useNanobotStream", () => { }); }); + it("clears stale stream state when the transport resets a run", async () => { + const fake = fakeClient(); + const { result } = renderHook( + () => useNanobotStream("chat-reconnect-reset", EMPTY_MESSAGES), + { wrapper: wrap(fake.client) }, + ); + + act(() => { + fake.emit("chat-reconnect-reset", { + event: "goal_status", + chat_id: "chat-reconnect-reset", + status: "running", + started_at: 1_700, + }); + fake.emit("chat-reconnect-reset", { + event: "delta", + chat_id: "chat-reconnect-reset", + text: "partial", + }); + }); + await flushStreamFrame(); + expect(result.current.isStreaming).toBe(true); + + act(() => fake.emitRunStatus("chat-reconnect-reset", null)); + + expect(result.current.runStartedAt).toBeNull(); + expect(result.current.isStreaming).toBe(false); + expect(result.current.messages[0]).toMatchObject({ + content: "partial", + isStreaming: false, + }); + }); + it("flushes pending delta text before turn_end finalizes the turn", () => { const fake = fakeClient(); const { result } = renderHook(() => useNanobotStream("chat-flush", EMPTY_MESSAGES), {