From 924737fa6f3e4aaef86e0081de3fa2a14b81b6de Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 5 Oct 2026 23:41:11 -0700 Subject: [PATCH 1/4] fix(mothership): re-attach at once when the user returns to a recovered stream Once the user left a running chat and came back, the return recovery owned the stream for the rest of the turn, and every later online/visible/pageshow event just joined it. So a network drop after that waited out whatever the recovery was doing: a tail that went silent held the stream until the 45s idle timeout, and a tail that failed slept out a reconnect backoff of up to 30s. The same drop on a stream the send still owned re-attached immediately, because the return signal supersedes the send's reader. On the local QA stack the stream resumed 48s after the network came back. A return signal now supersedes an in-flight recovery the same way, and the new recovery re-attaches from the cursor. --- .../home/hooks/use-chat.dom.test.tsx | 90 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 13 ++- 2 files changed, 99 insertions(+), 4 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index affe8a1c15c..a7fcb27b1e9 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -1136,6 +1136,96 @@ describe('useChat remount send recovery', () => { } }) + /** + * After the user leaves and returns once, the return recovery owns the stream for + * the rest of the turn. When the network then drops, its tail either goes silent + * (the socket stalls) or fails into the reconnect backoff, which grows to 30s. + * Coming back online must re-attach at once, as it does while the send still owns + * the stream, instead of waiting out the idle timeout or the backoff. + */ + it.each(['stalled', 'failed'] as const)( + 're-attaches at once when the network returns to a return recovery whose tail %s', + async (tailOutcome) => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) + try { + let online = true + let backOnline = false + let tailOpenedAfterReturn = false + let failedReconnects = 0 + const openTails: ReadableStreamDefaultController[] = [] + const history: MothershipChatHistory = { + id: `chat-recovery-${tailOutcome}`, + mode: 'agent', + title: 'Recovery', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => + Promise.resolve({ + chat: { ...history, activeStreamId: state.postBodies[0]?.userMessageId ?? null }, + }) + ) + state.postBehavior = 'accept' + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init) + if (!online) { + failedReconnects++ + throw new TypeError('Failed to fetch') + } + if (url.includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + if (backOnline) tailOpenedAfterReturn = true + return new Response( + new ReadableStream({ + start: (controller) => void openTails.push(controller), + }), + { headers: { 'Content-Type': 'text/event-stream' } } + ) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Keep going while I am away') + }) + await act(async () => vi.advanceTimersByTimeAsync(100)) + await act(async () => { + window.dispatchEvent(new Event('pageshow')) + await vi.advanceTimersByTimeAsync(100) + }) + expect(openTails.length).toBeGreaterThan(0) + + online = false + if (tailOutcome === 'failed') { + await act(async () => { + for (const tail of openTails.splice(0)) tail.error(new TypeError('network error')) + await vi.advanceTimersByTimeAsync(0) + }) + for (let second = 0; second < 120 && failedReconnects < 6; second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + expect(failedReconnects).toBeGreaterThanOrEqual(6) + } else { + await act(async () => vi.advanceTimersByTimeAsync(20_000)) + } + + online = true + backOnline = true + await act(async () => { + window.dispatchEvent(new Event('online')) + await vi.advanceTimersByTimeAsync(500) + }) + + expect(tailOpenedAfterReturn).toBe(true) + expect(getResult().isSending).toBe(true) + expect(state.postBodies).toHaveLength(1) + } finally { + vi.useRealTimers() + } + } + ) + it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) try { diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index feb4fde0e4a..ae24143017d 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -2982,12 +2982,17 @@ export function useChat( if (!chatId) return const subjectKey = buildRecoverySubjectKey(startingChatId, startingSelectedChatId) + /* A return signal supersedes a recovery already in flight, as it supersedes the + send's own reader: that recovery's tail may have gone silent while the tab was + away or offline, or it may be sleeping out a reconnect backoff, and either + would hold the stream for up to the idle timeout or the backoff. */ const existingRecovery = activeStreamReturnRecoveryRef.current - if (existingRecovery?.subjectKey === subjectKey) { - return existingRecovery.promise - } if (existingRecovery) { - existingRecovery.controller.abort('replaced_by_new_recovery_subject') + existingRecovery.controller.abort( + existingRecovery.subjectKey === subjectKey + ? 'superseded_by_return' + : 'replaced_by_new_recovery_subject' + ) activeStreamReturnRecoveryRef.current = null } From b62ebb9f93506c60999607f3946c030904c498b5 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 09:23:37 -0700 Subject: [PATCH 2/4] fix(mothership): finish a turn that ended while its reader was silent A return recovery used the chat history's `activeStreamId` to decide what to re-attach to. When the turn had ended while this surface was not listening (its reader stalled, or a later return superseded the recovery that held it), the history listed no running turn, so recovery returned without touching the stream this surface still showed as running, and the chat stayed on Stop until an idle timeout or backoff happened to run into the terminal state. When the history lists no running turn but this surface is still sending, recovery now resolves that stream: its terminal status replays the remaining events and finalizes. --- .../home/hooks/use-chat.dom.test.tsx | 58 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 9 ++- 2 files changed, 66 insertions(+), 1 deletion(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index a7fcb27b1e9..4a678d1e89c 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -1226,6 +1226,64 @@ describe('useChat remount send recovery', () => { } ) + /** + * The turn ends on the server while this surface's reader is silent (a stalled + * socket, or a recovery a later return superseded), so it never sees `complete`. + * The next return reads a chat with no running turn; it must resolve the stream + * it still shows as running instead of leaving the chat stuck on Stop. + */ + it('finishes a turn that ended while its reader was silent when the user returns', async () => { + let turnRunning = true + const history: MothershipChatHistory = { + id: 'chat-ended-while-silent', + mode: 'agent', + title: 'Ended while silent', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => + Promise.resolve({ + chat: { + ...history, + activeStreamId: turnRunning ? (state.postBodies[0]?.userMessageId ?? null) : null, + }, + }) + ) + state.postBehavior = 'accept' + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init) + if (url.includes('batch=true')) { + return Response.json({ + success: true, + events: [], + status: turnRunning ? 'streaming' : 'complete', + }) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Finish while I am away') + }) + await act(async () => { + window.dispatchEvent(new Event('pageshow')) + await sleep(100) + }) + expect(getResult().isSending).toBe(true) + + turnRunning = false + await act(async () => { + window.dispatchEvent(new Event('online')) + }) + await waitFor(() => !getResult().isSending) + + expect(state.postBodies).toHaveLength(1) + }) + it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) try { diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index ae24143017d..8ccf271896c 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -3010,8 +3010,15 @@ export function useChat( const fallbackStreamId = streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? cached?.activeStreamId const loadedStream = await getActiveStreamIdForChat(chatId, recoveryController.signal) + /* The chat no longer lists a running turn, but this surface is still showing + one: it ended while nothing here was listening (its reader went silent, or a + recovery it superseded was attached). Resolve that stream instead of + leaving it running: its terminal state replays the rest and finalizes. */ + const locallyRunningStreamId = sendingRef.current + ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) + : undefined const streamId = loadedStream.loaded - ? (loadedStream.streamId ?? undefined) + ? (loadedStream.streamId ?? locallyRunningStreamId) : fallbackStreamId if ( !isSameRecoverySubject() || From 66c195d1b272e2fc78b2f790935184d4b8da9bb6 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 10:28:01 -0700 Subject: [PATCH 3/4] fix(mothership): leave a send waiting for admission alone on return A send shows as running before its POST is admitted, but the chat cannot list it yet. Resolving that stream on a return event read it as ended and aborted the POST. Recovery now resolves only a locally running stream whose POST was admitted. --- .../home/hooks/use-chat.dom.test.tsx | 38 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 11 ++++-- 2 files changed, 45 insertions(+), 4 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 4a678d1e89c..2052307933f 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -1284,6 +1284,44 @@ describe('useChat remount send recovery', () => { expect(state.postBodies).toHaveLength(1) }) + /** + * Before its POST is admitted a send shows as running, but the chat cannot list + * it yet. A return event in that window must leave the POST alone. + */ + it('does not abort a send still waiting for admission when the user returns', async () => { + const history: MothershipChatHistory = { + id: 'chat-pending-admission', + mode: 'agent', + title: 'Pending admission', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let postSignal: AbortSignal | undefined + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) === '/api/mothership/chat' && init?.method === 'POST') { + state.postBodies.push(JSON.parse(String(init.body))) + postSignal = init.signal ?? undefined + return new Promise(() => {}) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Still being admitted') + }) + await waitFor(() => postSignal !== undefined) + + await act(async () => { + window.dispatchEvent(new Event('online')) + await sleep(200) + }) + + expect(postSignal?.aborted).toBe(false) + expect(getResult().isSending).toBe(true) + }) + it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) try { diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 8ccf271896c..9373264f5ca 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -3013,10 +3013,13 @@ export function useChat( /* The chat no longer lists a running turn, but this surface is still showing one: it ended while nothing here was listening (its reader went silent, or a recovery it superseded was attached). Resolve that stream instead of - leaving it running: its terminal state replays the rest and finalizes. */ - const locallyRunningStreamId = sendingRef.current - ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) - : undefined + leaving it running: its terminal state replays the rest and finalizes. A + send still waiting for its POST to be admitted is not such a stream: the + chat cannot list it yet, and recovering it would abort that POST. */ + const locallyRunningStreamId = + sendingRef.current && !pendingChatAdmissionRef.current + ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) + : undefined const streamId = loadedStream.loaded ? (loadedStream.streamId ?? locallyRunningStreamId) : fallbackStreamId From 8fc5dae39d7c12cea8882d4127030af7ccd9481d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 11:05:43 -0700 Subject: [PATCH 4/4] fix(mothership): finish an admitted turn whose POST never answered on return Recovery left a send alone while its POST had not answered, so a turn the server admitted and finished, whose answer never reached the client, kept the chat on Stop with nothing to clear it. The send now counts as admitted once the loaded chat holds its message, and recovery resolves its stream as any other. --- .../home/hooks/use-chat.dom.test.tsx | 56 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 14 ++++- 2 files changed, 67 insertions(+), 3 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 2052307933f..23ca9f64f10 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -1322,6 +1322,62 @@ describe('useChat remount send recovery', () => { expect(getResult().isSending).toBe(true) }) + /** + * The server admitted the send and finished its turn, but the POST's answer + * never arrived. Once the chat holds the message, a return resolves the turn + * rather than leaving the chat on Stop behind a POST that will not answer. + */ + it('finishes an admitted turn whose POST never answered when the user returns', async () => { + let admitted = false + const history: MothershipChatHistory = { + id: 'chat-admitted-unanswered', + mode: 'agent', + title: 'Admitted, unanswered', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => { + const userMessageId = state.postBodies[0]?.userMessageId + return Promise.resolve({ + chat: { + ...history, + messages: + admitted && userMessageId + ? [ + { id: userMessageId, role: 'user', content: 'Answer lost' }, + { id: 'saved-answer', role: 'assistant', content: 'Done.' }, + ] + : [], + }, + }) + }) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url === '/api/mothership/chat' && init?.method === 'POST') { + state.postBodies.push(JSON.parse(String(init.body))) + return new Promise(() => {}) + } + if (url.includes('/api/mothership/chat/stream')) { + return Response.json({ success: true, events: [], status: 'complete' }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Answer lost') + }) + await waitFor(() => state.postBodies.length === 1) + + admitted = true + await act(async () => { + window.dispatchEvent(new Event('online')) + }) + await waitFor(() => !getResult().isSending) + + expect(state.postBodies).toHaveLength(1) + }) + it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) try { diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 9373264f5ca..8d28f28ee8a 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -3014,10 +3014,18 @@ export function useChat( one: it ended while nothing here was listening (its reader went silent, or a recovery it superseded was attached). Resolve that stream instead of leaving it running: its terminal state replays the rest and finalizes. A - send still waiting for its POST to be admitted is not such a stream: the - chat cannot list it yet, and recovering it would abort that POST. */ + send whose POST has not answered yet is such a stream only once the loaded + chat holds its message (the server admitted it, and the answer is lost); + until then it may still be on its way, and recovering it would abort it. */ + const pendingAdmission = pendingChatAdmissionRef.current + const admitted = + !pendingAdmission || + (loadedStream.loaded && + queryClient + .getQueryData(mothershipChatKeys.detail(chatId)) + ?.messages.some((message) => message.id === pendingAdmission.userMessageId) === true) const locallyRunningStreamId = - sendingRef.current && !pendingChatAdmissionRef.current + sendingRef.current && admitted ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) : undefined const streamId = loadedStream.loaded