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..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 @@ -1136,6 +1136,248 @@ 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() + } + } + ) + + /** + * 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) + }) + + /** + * 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) + }) + + /** + * 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 feb4fde0e4a..8d28f28ee8a 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 } @@ -3005,8 +3010,26 @@ 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. A + 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 && admitted + ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) + : undefined const streamId = loadedStream.loaded - ? (loadedStream.streamId ?? undefined) + ? (loadedStream.streamId ?? locallyRunningStreamId) : fallbackStreamId if ( !isSameRecoverySubject() ||