From dbd1a6431f272cc59970247a4530975dd0192957 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 13:30:22 -0700 Subject: [PATCH 1/9] fix(mothership): close the remaining ways a chat send is lost, resent hot, or polled after unmount - A deleted chat's queue could come back through `enqueue`: a direct send to a chat deleted while its POST was failing (unreachable or busy) was re-queued, invisible, and would go out if the chat were restored. A cleared chat now refuses every queue write until it is restored from Recently Deleted. - A busy refusal re-queued its message for immediate redispatch. With Redis down every send is refused as busy without naming a turn, so the message was resent on every refusal (or stalled in the queue). Busy re-queues now wait out a jittered, growing delay (`backoffWithJitter`, 1s to 30s) before the drain sends them again under the same id. - A send waiting on the server's chat lock behind another tab's turn was aborted by a return-triggered recovery, and its message was lost. Recovery now leaves any POST that has not answered alone unless the chat already holds its message. When the busy refusal then arrives, the chat attaches to the turn holding it explicitly: the history read it relies on can repeat one the surface skipped while the POST was pending, so nothing else re-ran to attach, and the queued message waited behind a turn the tab never showed. - A retry told "already sent" while the server's earlier attempt with that id was still in flight (it had opened no stream yet) reattached to the missing stream, read the 404 as a finished turn, and dropped the message. Seen in the browser harness when Redis came back during a busy-refusal retry. Such a refusal is now retried later like a busy one. - A Send-now message whose Stop handoff failed (for example because the user switched chats while the Stop settled) returned a plain failure that the dispatch epoch check discarded. It is now handed back held, so it stays in its own chat's queue for the user to send; a surface that unmounted keeps resuming it from its stored handoff as before. - The saved-turn re-read kept refetching an unobserved chat for up to two minutes after the chat view unmounted. It now stops with the view. --- .../home/hooks/use-chat.dom.test.tsx | 342 ++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 138 +++++-- apps/sim/hooks/queries/mothership-chats.ts | 3 + apps/sim/stores/mothership-queue/store.ts | 23 +- apps/sim/stores/mothership-queue/types.ts | 10 +- 5 files changed, 483 insertions(+), 33 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 c3ced0b8475..88fadb352d8 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 @@ -915,6 +915,59 @@ describe('useChat remount send recovery', () => { expect(state.abortBodies).toHaveLength(1) }) + /** + * "Send now" on a queued message stops the running turn first. If the user + * switches chats while that Stop settles, the message is not sent into the + * other chat; it must go back to its own chat's queue, not vanish. + */ + it('keeps a send-now message in its chat when the user switches chats during the Stop', async () => { + state.postBehavior = 'task' + const { getResult, navigate } = renderUseChatInChat('chat-a') + await act(async () => { + void getResult().sendMessage('Original request') + }) + await waitFor(() => state.postBodies.length === 1 && getResult().isSending) + let releaseStop = () => {} + const stopGate = new Promise((resolve) => { + releaseStop = resolve + }) + let stopRequested = false + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).includes('/api/copilot/chat/abort')) { + stopRequested = true + await stopGate + } + return fetchStub(input, init) + }) + state.postBehavior = 'hang' + await act(async () => { + void getResult().sendMessage('Use the latest report') + }) + await waitFor(() => useMothershipQueueStore.getState().queues['chat-a']?.length === 1) + await act(async () => { + void getResult().sendNow() + }) + await waitFor(() => stopRequested) + + navigate('chat-b', { + id: 'chat-b', + mode: 'agent', + title: 'Other chat', + messages: [], + activeStreamId: null, + resources: [], + }) + await act(async () => { + releaseStop() + await sleep(200) + }) + + expect(state.postBodies).toHaveLength(1) + expect( + (useMothershipQueueStore.getState().queues['chat-a'] ?? []).map((message) => message.content) + ).toEqual(['Use the latest report']) + }) + it('surfaces an explicit admission rejection without reconnecting or marking the user turn stopped', async () => { const fetch = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { if (String(input) === '/api/mothership/chat' && init?.method === 'POST') @@ -2510,6 +2563,214 @@ describe('useChat remount send recovery', () => { expect(useMothershipQueueStore.getState().queues[history.id]).toBeUndefined() }) + /** + * A direct send to a chat that is deleted while its POST is failing must not + * bring the chat's queue back: nothing would show the message, and it would + * go out by itself if the chat were ever restored. + */ + it.each(['unreachable', 'busy'] as const)( + 'does not recreate a deleted chat through a %s direct send', + async (outcome) => { + const history = idleHistory(`chat-deleted-direct-${outcome}`) + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let answerPost: (() => void) | 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))) + return new Promise((resolve, reject) => { + answerPost = () => + outcome === 'unreachable' + ? reject(new TypeError('Failed to fetch')) + : resolve( + Response.json( + { error: 'A response is already in progress for this chat.' }, + { status: 409 } + ) + ) + }) + } + if (String(input).includes('/api/mothership/chat/stream')) { + return Response.json({ error: 'Stream not found' }, { status: 404 }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Sent to a chat I then deleted') + }) + await waitFor(() => answerPost !== undefined) + + useMothershipQueueStore.getState().clearChat(history.id) + await act(async () => { + answerPost?.() + await sleep(300) + }) + + expect(useMothershipQueueStore.getState().queues[history.id]).toBeUndefined() + expect(state.postBodies).toHaveLength(1) + } + ) + + /** + * With Redis down the server refuses every send as busy without naming a + * turn, and nothing is running. The message must be retried on a growing + * delay, not resent as fast as each refusal comes back. + */ + it('backs off retrying a send the server keeps refusing as busy without a turn', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout', 'Date'] }) + try { + const history = idleHistory('chat-busy-without-redis') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + 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))) + /** The server waits on the chat lock before refusing. */ + await sleep(5_000) + return Response.json( + { error: 'A response is already in progress for this chat.' }, + { status: 409 } + ) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Refused while Redis is down') + await vi.advanceTimersByTimeAsync(50) + }) + for (let second = 0; second < 90; second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + + /** Back to back, a 5s refusal allows 18 attempts in 90s; backing off allows far fewer. */ + expect(state.postBodies.length).toBeGreaterThan(2) + expect(state.postBodies.length).toBeLessThanOrEqual(8) + expect(new Set(state.postBodies.map((body) => body.userMessageId)).size).toBe(1) + expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe( + 'Refused while Redis is down' + ) + } finally { + vi.useRealTimers() + } + }) + + /** + * A send can wait several seconds on the server's chat lock behind another + * tab's turn. Returning to the tab in that window must not abort it: when + * the refusal arrives, the chat shows that turn running and the message + * waits in the queue behind it. + */ + it('keeps a send waiting on the chat lock when the user returns to the tab', async () => { + const history: MothershipChatHistory = { + ...idleHistory('chat-lock-wait-return'), + activeStreamId: 'turn-from-another-tab', + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let answerPost: (() => void) | undefined + let answered = false + const tailsAfterRefusal: string[] = [] + 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((resolve, reject) => { + init.signal?.addEventListener('abort', () => reject(init.signal?.reason), { + once: true, + }) + answerPost = () => + resolve( + Response.json( + { + error: 'A response is already in progress for this chat.', + activeStreamId: 'turn-from-another-tab', + }, + { status: 409 } + ) + ) + }) + } + if (url.includes('/api/mothership/chat/stream')) { + if (url.includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + if (answered) { + tailsAfterRefusal.push( + new URL(url, 'http://localhost').searchParams.get('streamId') ?? '' + ) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, { + ...history, + activeStreamId: null, + }) + await act(async () => { + void getResult().sendMessage('Waiting on the lock') + }) + await waitFor(() => answerPost !== undefined) + + await act(async () => { + window.dispatchEvent(new Event('pageshow')) + await sleep(100) + }) + await act(async () => { + answered = true + answerPost?.() + await sleep(300) + }) + + const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] + expect(queued.map((message) => message.content)).toEqual(['Waiting on the lock']) + expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) + expect(tailsAfterRefusal).toContain('turn-from-another-tab') + expect(getResult().isSending).toBe(true) + }) + + /** + * A retry can be told "already sent" while the server's earlier attempt with + * that id is still in flight and has opened no stream. The message must be + * retried once that attempt settles, not read as a finished turn and dropped. + */ + it('retries a send deduplicated against an attempt that opened no stream', async () => { + const history = idleHistory('chat-deduped-without-stream') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let posts = 0 + 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))) + posts++ + if (posts === 1) { + return Response.json( + { + error: 'This message was already sent.', + activeStreamId: state.postBodies[0].userMessageId, + }, + { status: 409 } + ) + } + return emptySseResponse() + } + if (url.includes('/api/mothership/chat/stream') && posts === 1) { + return Response.json({ error: 'Stream not found' }, { status: 404 }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + + await act(async () => { + await getResult().sendMessage('Told it was already sent') + }) + await waitFor(() => state.postBodies.length === 2, 5_000) + + expect(state.postBodies[1].message).toBe('Told it was already sent') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + /** The `online` event can fire while no surface for the chat is mounted. */ it('sends a held message when its chat mounts after the network came back', async () => { const history = idleHistory('chat-held-while-away') @@ -2869,6 +3130,87 @@ describe('useChat remount send recovery', () => { } ) + /** + * The saved-turn re-read belongs to the chat view: once it unmounts, nothing + * renders that chat, so the re-read must stop instead of refetching it for + * minutes. + */ + it('stops re-reading the saved turn when the chat view unmounts', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) + try { + const chatId = 'chat-reread-unmount' + const history: MothershipChatHistory = { + id: chatId, + mode: 'agent', + title: 'Never saved', + messages: [], + activeStreamId: null, + resources: [], + } + let streamId: string | undefined + let completed = false + let detailReads = 0 + mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => { + if (contract.path !== '/api/mothership/chats/[chatId]') { + return Promise.resolve({ chats: [] }) + } + if (!completed || !streamId) return Promise.resolve({ chat: history }) + detailReads++ + return Promise.resolve({ + chat: { + ...history, + activeStreamId: streamId, + messages: [ + { id: streamId, role: 'user', content: 'Summarize the run' }, + { id: `live-assistant:${streamId}`, role: 'assistant', content: 'Done.' }, + ], + }, + }) + }) + let stream: ReadableStreamDefaultController | undefined + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') { + return fetchStub(input, init) + } + streamId = JSON.parse(String(init.body)).userMessageId + return new Response( + new ReadableStream({ + start(controller) { + stream = controller + }, + }), + { headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } } + ) + }) + const { getResult, unmount } = renderUseChatInChat(chatId, history) + await act(async () => { + void getResult().sendMessage('Summarize the run') + await vi.advanceTimersByTimeAsync(50) + }) + await act(async () => { + stream?.enqueue( + new TextEncoder().encode( + `data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, seq: 1, type: 'complete', payload: { status: 'complete' } })}\n\n` + ) + ) + completed = true + stream?.close() + await vi.advanceTimersByTimeAsync(1_000) + }) + expect(detailReads).toBeGreaterThan(0) + + unmount() + const readsAtUnmount = detailReads + for (let second = 0; second < 60; second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + + expect(detailReads).toBe(readsAtUnmount) + } finally { + vi.useRealTimers() + } + }) + describe.each([ { kind: 'browser action', 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 38cb2fa0db4..9c8cbd302f8 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -205,6 +205,18 @@ interface FinalizeOptions { streamTerminal?: boolean } +/** + * Why the hook re-attaches to the chat's running turn: the browser returned + * (`pageshow`, `visible`, `online`), an exhausted reconnect rechecks, or a send + * was refused because that turn holds the chat. + */ +type ActiveStreamRecoveryReason = + | 'pageshow' + | 'visible' + | 'online' + | 'exhausted_recheck' + | 'busy_refusal' + /** * A send handed back to the caller instead of rendered. `userMessageId` is what * a retry reuses so the server deduplicates the two attempts. An `unreachable` @@ -217,6 +229,10 @@ interface WithdrawnSendResult { userMessageId: string unreachable?: boolean networkReturned?: boolean + /** Refused because another turn held the chat; retried on a growing delay. */ + busy?: boolean + /** Not sent at all (its Stop handoff failed); kept queued for the user to send. */ + held?: boolean } /** @@ -318,6 +334,9 @@ const PERSISTED_TURN_REFETCH_BASE_MS = 250 const PERSISTED_TURN_REFETCH_MAX_DELAY_MS = 5_000 /** How long a finished turn's save is waited for; a slow save still lands well inside it. */ const PERSISTED_TURN_WAIT_MS = 120_000 +/** Pacing for re-sending a message the server refused because the chat was busy. */ +const BUSY_RETRY_BASE_MS = 1_000 +const BUSY_RETRY_MAX_MS = 30_000 const STOP_REQUEST_TIMEOUT_MS = 15_000 const DETACHED_CHAT_RETRY_BASE_MS = 1000 const DETACHED_CHAT_RETRY_MAX_MS = 30_000 @@ -694,6 +713,16 @@ export function getWorkflowCopilotUseChatOptions( } } +/** Queue fields for the `attempt`th busy refusal of a message: when it may be sent again. */ +function busyRetry(attempt: number): { busyRetries: number; notBefore: number } { + return { + busyRetries: attempt, + notBefore: + Date.now() + + backoffWithJitter(attempt, null, { baseMs: BUSY_RETRY_BASE_MS, maxMs: BUSY_RETRY_MAX_MS }), + } +} + export function useChat( owner: string | { organizationId: string }, initialChatId?: string, @@ -954,9 +983,9 @@ export function useChat( >(async () => ({ terminal: false })) const finalizeRef = useRef<(options?: FinalizeOptions) => void>(() => {}) const recoveringQueuedSendHandoffRef = useRef(null) - const recoverActiveStreamRef = useRef< - (reason: 'pageshow' | 'visible' | 'online' | 'exhausted_recheck') => Promise - >(async () => {}) + const recoverActiveStreamRef = useRef<(reason: ActiveStreamRecoveryReason) => Promise>( + async () => {} + ) const reconnectExhaustedRecheckTimerRef = useRef | null>(null) const abortControllerRef = useRef(null) @@ -1946,7 +1975,8 @@ export function useChat( * in-flight copy: the stream listed as active and the answer under its live id. * That copy matches the optimistic one, so nothing would read it again; re-read * until the saved turn is there. The first pass joins finalize's own read. The - * wait ends as soon as this chat moves on: another send, or another chat. + * wait ends as soon as this chat moves on (another send, or another chat) or + * the chat view unmounts. */ const awaitPersistedTurn = useCallback( async (chatId: string, streamId: string) => { @@ -1963,7 +1993,11 @@ export function useChat( }) ) } - if (locallyTerminalStreamIdRef.current !== streamId || chatIdRef.current !== chatId) + if ( + !surfaceMountedRef.current || + locallyTerminalStreamIdRef.current !== streamId || + chatIdRef.current !== chatId + ) return await queryClient.refetchQueries( { queryKey: mothershipChatKeys.detail(chatId), exact: true }, @@ -3047,7 +3081,7 @@ export function useChat( retryReconnectRef.current = retryReconnect const recoverActiveStreamFromRedis = useCallback( - async (reason: 'pageshow' | 'visible' | 'online' | 'exhausted_recheck'): Promise => { + async (reason: ActiveStreamRecoveryReason): Promise => { const startingChatId = chatIdRef.current const startingSelectedChatId = selectedChatIdRef.current const chatId = startingChatId ?? startingSelectedChatId @@ -3096,10 +3130,14 @@ export function useChat( queryClient .getQueryData(mothershipChatKeys.detail(chatId)) ?.messages.some((message) => message.id === pendingAdmission.userMessageId) === true) - const locallyRunningStreamId = - sendingRef.current && admitted - ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) - : undefined + /* A POST still on its way owns this surface whatever the chat lists (another + tab's turn, or nothing yet): recovering would abort it before the server + answers, and that answer (the stream, or a busy refusal that re-queues the + message) is what must decide what comes next. */ + if (!admitted) return + const locallyRunningStreamId = sendingRef.current + ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId) + : undefined const streamId = loadedStream.loaded ? (loadedStream.streamId ?? locallyRunningStreamId) : fallbackStreamId @@ -3824,7 +3862,9 @@ export function useChat( setTransportIdle() } setError(getErrorMessage(err, 'Failed to stop the previous response')) - return false + /* Nothing was sent. Hand the message back so it stays in its chat's queue + even if the user has switched chats since the Stop began. */ + return { userMessageId, held: true } } } @@ -3935,15 +3975,10 @@ export function useChat( setTransportIdle() } setError('Previous response is still shutting down; queued message was restored.') - return false + return { userMessageId, held: true } } - if (conflictStreamId !== userMessageId) { - /* Another turn holds the chat: one started in another tab, or one this - surface lost track of. This message was not admitted (the server - released its id), so it goes back to the queue under the same id. The - queue drains only while the chat is idle, so the chat's running turn is - read before the message is handed back: the chat then attaches to that - turn and the message goes out once, after it ends. */ + /** Withdraws this refused send so the queue retries it, under the same id, later. */ + const releaseRefusedSend = () => { rollbackOptimisticSend() if (streamGenRef.current === gen) { streamGenRef.current++ @@ -3952,6 +3987,15 @@ export function useChat( clearActiveTurn() setTransportIdle() } + } + if (conflictStreamId !== userMessageId) { + /* Another turn holds the chat: one started in another tab, or one this + surface lost track of. This message was not admitted (the server + released its id), so it goes back to the queue under the same id. The + queue drains only while the chat is idle, so the chat's running turn is + read before the message is handed back: the chat then attaches to that + turn and the message goes out once, after it ends. */ + releaseRefusedSend() if (requestChatId) { const busyChatId = requestChatId if (conflictStreamId) { @@ -3964,7 +4008,13 @@ export function useChat( .refetchQueries({ queryKey: mothershipChatKeys.detail(busyChatId), exact: true }) .catch(() => {}) } - return { userMessageId } + /* The history read above can repeat one this surface skipped while the POST + was pending, so nothing re-runs to attach to the turn it lists. Attach + explicitly; the message is retried after that turn ends. */ + if (pendingChatAdmissionRef.current === admission) + pendingChatAdmissionRef.current = null + void recoverActiveStreamRef.current('busy_refusal') + return { userMessageId, busy: true } } /* A send deduplicated against an earlier attempt comes back naming the chat that attempt opened. Adopting it here spares a chatless @@ -3980,6 +4030,22 @@ export function useChat( }) streamTargetChatId = conflictChatId } + /* "Already sent" with no stream for it means the earlier attempt is still + in flight on the server (or died before starting a turn), not that a turn + ran: reattaching would read the missing stream as finished and drop the + message. Retry it later like a busy refusal; the server's claim settles. */ + const dedupedStreamExists = await fetchStreamBatch( + conflictStreamId, + '0', + abortController.signal + ).then( + () => true, + (error: unknown) => !isStreamGoneError(error) + ) + if (!dedupedStreamExists) { + releaseRefusedSend() + return { userMessageId, busy: true } + } streamIdRef.current = conflictStreamId const succeeded = await retryReconnect({ streamId: conflictStreamId, @@ -4147,6 +4213,7 @@ export function useChat( finalize, resumeOrFinalize, retryReconnect, + fetchStreamBatch, clearActiveTurn, resetStreamingBuffers, resolveChatIdForStream, @@ -4284,7 +4351,11 @@ export function useChat( ? { assistantSearchLevel: options?.assistantSearchLevel } : {}), } - if (!result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + if ( + !result.unreachable && + !result.held && + activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ) { handOffWithdrawnSend(withdrawn) return } @@ -4301,6 +4372,8 @@ export function useChat( ...(result.unreachable && !result.networkReturned ? { retryRequired: true, heldUntilOnline: true } : {}), + ...(result.held ? { retryRequired: true } : {}), + ...(result.busy ? busyRetry(1) : {}), ...(result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), @@ -4861,7 +4934,7 @@ export function useChat( withdrawn?: WithdrawnSendResult ) => { const withdrawnUserMessageId = withdrawn?.userMessageId - const retriesOnItsOwn = withdrawn !== undefined && !withdrawn.unreachable + const retriesOnItsOwn = withdrawn !== undefined && !withdrawn.unreachable && !withdrawn.held const savedHandoff = readQueuedSendHandoffState() const retainedHandoff = savedHandoff?.id === msg.id @@ -4878,8 +4951,11 @@ export function useChat( return } /* A withdrawn send was never admitted, so it goes back to its chat's queue - even when the user has moved on since its dispatch started. */ - if (options.epoch !== queueDispatchEpochRef.current && !withdrawn) { + even when the user has moved on since its dispatch started. The exception + is a held Stop handoff whose surface unmounted: its stored handoff is the + recovery, and the next mount of its chat resumes the Stop and the send. */ + const epochMoved = options.epoch !== queueDispatchEpochRef.current + if (epochMoved && (!withdrawn || (withdrawn.held && !surfaceMountedRef.current))) { return } // If the user explicitly removed this message during dispatch, honor @@ -4913,6 +4989,7 @@ export function useChat( ...dispatched, ...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}), retryRequired: withdrawn?.unreachable ? !withdrawn.networkReturned : !retriesOnItsOwn, + ...(withdrawn?.busy ? busyRetry((dispatched.busyRetries ?? 0) + 1) : {}), ...(withdrawn?.unreachable && !withdrawn.networkReturned ? { heldUntilOnline: true } : {}), @@ -4995,6 +5072,8 @@ export function useChat( const activeChatKey = chatKeyRef.current const msg = queueState.queues[activeChatKey]?.[0] if (!msg || msg.retryRequired) continue + /** A busy refusal's retry waits out its delay; the drain effect wakes it. */ + if (msg.notBefore !== undefined && msg.notBefore > Date.now()) continue // Pause draining if the head is bound to the composer; dispatching now // would race the eventual submit. The next kick on edit-resolve resumes us. if (queueState.editing[activeChatKey] === msg.id) continue @@ -5174,9 +5253,18 @@ export function useChat( const chatHistoryReady = chatHistory !== undefined const remoteActiveStreamId = chatHistory?.activeStreamId ?? null const queueHeadHeld = messageQueue[0]?.retryRequired === true + const queueHeadNotBefore = messageQueue[0]?.notBefore + const [busyRetryWakeup, setBusyRetryWakeup] = useState(0) useEffect(() => { if (!scopeKey) return if (messageQueue.length === 0 || queueHeadHeld) return + if (queueHeadNotBefore !== undefined && queueHeadNotBefore > Date.now()) { + const timer = setTimeout( + () => setBusyRetryWakeup((wakeups) => wakeups + 1), + queueHeadNotBefore - Date.now() + ) + return () => clearTimeout(timer) + } if (sendingRef.current || pendingStopPromiseRef.current) return if (queueDispatchTaskRef.current) return if (resolvedChatId && !chatHistoryReady) return @@ -5188,6 +5276,8 @@ export function useChat( scopeKey, messageQueue.length, queueHeadHeld, + queueHeadNotBefore, + busyRetryWakeup, resolvedChatId, chatHistoryReady, remoteActiveStreamId, diff --git a/apps/sim/hooks/queries/mothership-chats.ts b/apps/sim/hooks/queries/mothership-chats.ts index df914c9220c..487257a7c3e 100644 --- a/apps/sim/hooks/queries/mothership-chats.ts +++ b/apps/sim/hooks/queries/mothership-chats.ts @@ -365,6 +365,9 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) { const queryClient = useQueryClient() return useMutation({ mutationFn: restoreChat, + onSuccess: (_data, chatId) => { + useMothershipQueueStore.getState().reopenChat(chatId) + }, onSettled: () => { queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) }) }, diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index 418ef9d8bd6..94869bb918d 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -67,13 +67,15 @@ export const useMothershipQueueStore = create()( ...initialState, enqueue: (chatKey, message) => - set((state) => ({ - cleared: omitKey(state.cleared, chatKey), - queues: setQueueForChat(state.queues, chatKey, [ - ...(state.queues[chatKey] ?? []), - message, - ]), - })), + set((state) => { + if (state.cleared[chatKey]) return state + return { + queues: setQueueForChat(state.queues, chatKey, [ + ...(state.queues[chatKey] ?? []), + message, + ]), + } + }), insertAt: (chatKey, index, message) => set((state) => { @@ -99,6 +101,8 @@ export const useMothershipQueueStore = create()( retryRequired: _retry, heldUntilOnline: _held, heldSurface: _surface, + busyRetries: _busyRetries, + notBefore: _notBefore, ...rest } = next[index] next[index] = { @@ -204,6 +208,11 @@ export const useMothershipQueueStore = create()( cleared: { ...state.cleared, [chatKey]: true }, })), + reopenChat: (chatKey) => + set((state) => + state.cleared[chatKey] ? { cleared: omitKey(state.cleared, chatKey) } : state + ), + reset: () => set(initialState), }), { diff --git a/apps/sim/stores/mothership-queue/types.ts b/apps/sim/stores/mothership-queue/types.ts index af41dea707c..ad797c01ea9 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -24,6 +24,10 @@ export type QueuedMothershipMessage = QueuedMessage & { * mount: the next chatless surface for the same owner and workflow adopts it. */ heldSurface?: string + /** Busy refusals so far; paces the next retry. */ + busyRetries?: number + /** Epoch ms before which a busy-refused message is not sent again. */ + notBefore?: number /** * Message id of a prior attempt at this send that an unmount cleanup * withdrew. Reused when the entry is dispatched so the server deduplicates @@ -49,8 +53,8 @@ export interface MothershipQueueState { queues: Record editing: Record /** - * Chats cleared this session (deleted). A late restore of a send dispatched - * before the clear does not recreate their queue; a new enqueue lifts it. + * Chats cleared this session (deleted). No write recreates their queue (a late + * restore, or a failed send handed back); restoring the chat lifts it. */ cleared: Record @@ -65,5 +69,7 @@ export interface MothershipQueueState { /** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */ adoptHeldSends: (toKey: string, surface: string) => void clearChat: (chatKey: string) => void + /** Lifts `cleared` for a chat restored from Recently Deleted. */ + reopenChat: (chatKey: string) => void reset: () => void } From b81fde766cfaf715cb5a4b4f46bba2dc38f23589 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 14:38:11 -0700 Subject: [PATCH 2/9] fix(mothership): keep deleted chats empty across tabs and back off busy first messages - migrate no longer moves a new-chat queue into a chat deleted meanwhile - another tab's delete clears and tombstones the chat's queue; a restore (published as created) or the chat loading lifts it - a busy refusal on the new-chat surface waits out the backoff in its own queue instead of being handed to the surface's listener and resent at once - a deduplicated send whose stream check resolves after the user moved on leaves Stop on the new turn --- .../home/hooks/use-chat.dom.test.tsx | 158 +++++++++++++++++- .../[workspaceId]/home/hooks/use-chat.ts | 17 +- .../hooks/use-mothership-chat-events.test.ts | 22 +++ apps/sim/hooks/use-mothership-chat-events.ts | 5 + .../sim/stores/mothership-queue/store.test.ts | 9 + apps/sim/stores/mothership-queue/store.ts | 3 +- 6 files changed, 207 insertions(+), 7 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 88fadb352d8..93f34e05b17 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 @@ -2621,9 +2621,14 @@ describe('useChat remount send recovery', () => { try { const history = idleHistory('chat-busy-without-redis') mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + const redis = { up: false, acceptedPosts: 0 } 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))) + if (redis.up) { + redis.acceptedPosts++ + return emptySseResponse() + } /** The server waits on the chat lock before refusing. */ await sleep(5_000) return Response.json( @@ -2645,10 +2650,15 @@ describe('useChat remount send recovery', () => { /** Back to back, a 5s refusal allows 18 attempts in 90s; backing off allows far fewer. */ expect(state.postBodies.length).toBeGreaterThan(2) expect(state.postBodies.length).toBeLessThanOrEqual(8) + + /** Kept through every refusal: it goes out once Redis is back, under the same id. */ + redis.up = true + for (let second = 0; second < 45; second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + expect(redis.acceptedPosts).toBe(1) + expect(state.postBodies.at(-1)?.message).toBe('Refused while Redis is down') expect(new Set(state.postBodies.map((body) => body.userMessageId)).size).toBe(1) - expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe( - 'Refused while Redis is down' - ) } finally { vi.useRealTimers() } @@ -2771,6 +2781,148 @@ describe('useChat remount send recovery', () => { expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) }) + /** + * A first message from the new-chat surface can be refused as busy with no + * turn to wait for (Redis down). It must wait out the same growing delay + * there, not be handed to the surface's own send listener and resent at once. + */ + it('backs off a busy refusal of a first message on the new-chat surface', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout', 'Date'] }) + try { + 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))) + /** Caps a hot loop so it fails the count below instead of starving the test. */ + if (state.postBodies.length > 20) return new Promise(() => {}) + return Response.json( + { error: 'A response is already in progress for this chat.' }, + { status: 409 } + ) + } + return fetchStub(input, init) + }) + const { getResult } = renderHomeLikeSurface() + await act(async () => { + void getResult().sendMessage('First message while Redis is down') + await vi.advanceTimersByTimeAsync(50) + }) + for (let second = 0; second < 10; second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + + /** 1s, 2s, 4s, 8s of backoff fit at most 4 attempts in 10s. */ + expect(state.postBodies.length).toBeGreaterThan(1) + expect(state.postBodies.length).toBeLessThanOrEqual(4) + expect(new Set(state.postBodies.map((body) => body.userMessageId)).size).toBe(1) + expect(allQueuedMessages().map((message) => message.content)).toEqual([ + 'First message while Redis is down', + ]) + } finally { + vi.useRealTimers() + } + }) + + /** + * The "already sent" check of a deduplicated send can resolve after the user + * has moved to another chat and started a turn there. Stop must still stop + * that new turn, not the one the stale check found. + */ + it('keeps Stop on the new turn when a deduplicated send is checked after a chat switch', async () => { + const history = idleHistory('chat-deduped-then-left') + const other = idleHistory('chat-new-turn-after-switch') + mockRequestJson.mockImplementation((_contract: AnyApiRouteContract, input: unknown) => + Promise.resolve({ + chat: JSON.stringify(input).includes(other.id) ? other : history, + }) + ) + let finishCheck: (() => void) | undefined + 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))) + if (state.postBodies.length === 1) { + return Response.json( + { + error: 'This message was already sent.', + activeStreamId: state.postBodies[0].userMessageId, + }, + { status: 409 } + ) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + if (url.includes('/api/mothership/chat/stream') && finishCheck === undefined) { + return new Promise((resolve) => { + finishCheck = () => + resolve(Response.json({ success: true, events: [], status: 'streaming' })) + }) + } + return fetchStub(input, init) + }) + const { getResult, navigate } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Already sent here') + }) + await waitFor(() => finishCheck !== undefined) + + navigate(other.id, other) + await act(async () => { + void getResult().sendMessage('A new turn in the other chat') + }) + await waitFor(() => state.postBodies.length === 2) + await act(async () => { + finishCheck?.() + await sleep(100) + }) + await act(async () => { + await getResult().stopGeneration() + }) + + expect(state.abortBodies.map((body) => body.streamId)).toContain( + state.postBodies[1].userMessageId + ) + expect(state.abortBodies.map((body) => body.streamId)).not.toContain( + state.postBodies[0].userMessageId + ) + }) + + /** + * A chat this tab saw deleted can be restored from another tab without this + * tab hearing of it. Once the chat loads, it exists, so a follow-up queued + * behind its running turn must be kept. + */ + it('queues a follow-up in a chat that loads after this tab saw it deleted', async () => { + const history: MothershipChatHistory = { + ...idleHistory('chat-restored-elsewhere'), + activeStreamId: 'turn-still-running', + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).includes('/api/mothership/chat/stream')) { + if (String(input).includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + useMothershipQueueStore.getState().clearChat(history.id) + const { getResult } = renderUseChatInChat(history.id, history) + await waitFor(() => getResult().isSending) + + await act(async () => { + await getResult().sendMessage('Follow-up after the restore') + }) + + expect( + useMothershipQueueStore.getState().queues[history.id]?.map((message) => message.content) + ).toEqual(['Follow-up after the restore']) + }) + /** The `online` event can fire while no surface for the chat is mounted. */ it('sends a held message when its chat mounts after the network came back', async () => { const history = idleHistory('chat-held-while-away') 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 9c8cbd302f8..e329fbeffd0 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -2023,6 +2023,8 @@ export function useChat( const activeStreamId = chatHistory.activeStreamId appliedChatHistoryKeyRef.current = hydrationKey + /** The server returned this chat, so it exists: a delete seen earlier no longer applies. */ + useMothershipQueueStore.getState().reopenChat(chatHistory.id) const mappedMessages = chatHistory.messages.map(toDisplayMessage) const shouldReconnectActiveStream = Boolean(activeStreamId) && @@ -4046,6 +4048,8 @@ export function useChat( releaseRefusedSend() return { userMessageId, busy: true } } + /** The user may have moved on (another chat, another send) during the check. */ + if (streamGenRef.current !== gen) return consumedByTranscript streamIdRef.current = conflictStreamId const succeeded = await retryReconnect({ streamId: conflictStreamId, @@ -4354,6 +4358,7 @@ export function useChat( if ( !result.unreachable && !result.held && + !result.busy && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ) { handOffWithdrawnSend(withdrawn) @@ -4374,7 +4379,7 @@ export function useChat( : {}), ...(result.held ? { retryRequired: true } : {}), ...(result.busy ? busyRetry(1) : {}), - ...(result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ...((result.unreachable || result.busy) && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), }) @@ -4968,7 +4973,12 @@ export function useChat( restore would strand this under the dead instance's key — hand it to the next surface instead. A chat-bound key is the stable chat id, so the queue itself is the durable retry. */ - if (withdrawn && retriesOnItsOwn && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + if ( + withdrawn && + retriesOnItsOwn && + !withdrawn.busy && + dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ) { clearQueuedSendHandoffState(msg.id) handOffWithdrawnSend({ content: dispatched.content, @@ -4993,7 +5003,8 @@ export function useChat( ...(withdrawn?.unreachable && !withdrawn.networkReturned ? { heldUntilOnline: true } : {}), - ...(withdrawn?.unreachable && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ...((withdrawn?.unreachable || withdrawn?.busy) && + dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), diff --git a/apps/sim/hooks/use-mothership-chat-events.test.ts b/apps/sim/hooks/use-mothership-chat-events.test.ts index 940e04fe113..c451810d378 100644 --- a/apps/sim/hooks/use-mothership-chat-events.test.ts +++ b/apps/sim/hooks/use-mothership-chat-events.test.ts @@ -15,6 +15,7 @@ import { handleMothershipChatStatusEvent, resyncMothershipChatCaches, } from '@/hooks/use-mothership-chat-events' +import { useMothershipQueueStore } from '@/stores/mothership-queue/store' describe('handleMothershipChatStatusEvent', () => { const queryClient = { @@ -136,6 +137,27 @@ describe('handleMothershipChatStatusEvent', () => { expect(suspendTerminalScope).toHaveBeenCalledWith('chat-1') }) + it('drops the queue of a chat deleted elsewhere and takes sends again once it is restored', () => { + useMothershipQueueStore.getState().reset() + const queued = { id: 'm1', content: 'follow-up' } + const publish = (type: 'deleted' | 'created') => + handleMothershipChatStatusEvent( + queryClient, + 'ws-1', + JSON.stringify({ chatId: 'chat-1', type, timestamp: Date.now() }) + ) + + useMothershipQueueStore.getState().enqueue('chat-1', queued) + publish('deleted') + expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined() + useMothershipQueueStore.getState().enqueue('chat-1', queued) + expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined() + + publish('created') + useMothershipQueueStore.getState().enqueue('chat-1', queued) + expect(useMothershipQueueStore.getState().queues['chat-1']?.map((m) => m.id)).toEqual(['m1']) + }) + it('keeps started task detail when a stale started stream is older than the active stream', () => { queryClient.getQueryData.mockReturnValue({ id: 'chat-1', diff --git a/apps/sim/hooks/use-mothership-chat-events.ts b/apps/sim/hooks/use-mothership-chat-events.ts index eca7b80ec88..887eae2d3d5 100644 --- a/apps/sim/hooks/use-mothership-chat-events.ts +++ b/apps/sim/hooks/use-mothership-chat-events.ts @@ -10,6 +10,7 @@ import { type MothershipChatOwner, mothershipChatKeys, } from '@/hooks/queries/mothership-chats' +import { useMothershipQueueStore } from '@/stores/mothership-queue/store' const logger = createLogger('MothershipChatEvents') @@ -119,8 +120,12 @@ export function handleMothershipChatStatusEvent( // mutation would leave pages and PTYs running indefinitely. void suspendDesktopChatScopes(payload.chatId) queryClient.removeQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) }) + /** This tab's queue for the chat goes too, and no later send may bring it back. */ + useMothershipQueueStore.getState().clearChat(payload.chatId) return } + /** A restore is published as `created`; the chat takes queued sends again. */ + if (payload.type === 'created') useMothershipQueueStore.getState().reopenChat(payload.chatId) if (payload.type === 'renamed') { /** * The lists invalidated above carry the title every surface renders; the diff --git a/apps/sim/stores/mothership-queue/store.test.ts b/apps/sim/stores/mothership-queue/store.test.ts index 66f0ac3a2e2..bb3c4601211 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -88,5 +88,14 @@ describe('useMothershipQueueStore', () => { ]) expect(useMothershipQueueStore.getState().queues['pending::abc']).toBeUndefined() }) + + it('does not move a new chat surface queue into a chat deleted meanwhile', () => { + useMothershipQueueStore.getState().enqueue('pending::abc', message('pending-1')) + useMothershipQueueStore.getState().clearChat('chat-X') + useMothershipQueueStore.getState().migrate('pending::abc', 'chat-X') + const state = useMothershipQueueStore.getState() + expect(state.queues['chat-X']).toBeUndefined() + expect(state.queues['pending::abc']).toBeUndefined() + }) }) }) diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index 94869bb918d..f69bba53f78 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -148,7 +148,8 @@ export const useMothershipQueueStore = create()( if (!fromQueue && fromEditing === undefined) return state const queues = omitKey(state.queues, fromKey) - if (fromQueue && fromQueue.length > 0) { + /** A chat deleted meanwhile takes nothing: its queue is gone with it. */ + if (fromQueue && fromQueue.length > 0 && !state.cleared[toKey]) { // Merge defensively in case a stale bucket survived in // sessionStorage. FIFO: existing first, then the resolved stream. const existing = state.queues[toKey] ?? [] From c704a62b8f010d422e8e2324771a444171bf70a7 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:03:55 -0700 Subject: [PATCH 3/9] fix(mothership): lift a chat's delete only when a later server read returns it The chat history effect also ran on history written locally, such as the rollback after a busy refusal, so a chat deleted in another tab could take queued sends again before any restore. The delete is now lifted only by a server read that began after it. --- .../home/hooks/use-chat.dom.test.tsx | 54 +++++++++++++++- .../[workspaceId]/home/hooks/use-chat.ts | 2 - .../mothership-chat-history-restore.test.ts | 62 +++++++++++++++++++ apps/sim/hooks/queries/mothership-chats.ts | 17 ++++- 4 files changed, 131 insertions(+), 4 deletions(-) create mode 100644 apps/sim/hooks/queries/mothership-chat-history-restore.test.ts 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 93f34e05b17..68ab1baf0fb 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 @@ -2601,6 +2601,8 @@ describe('useChat remount send recovery', () => { await waitFor(() => answerPost !== undefined) useMothershipQueueStore.getState().clearChat(history.id) + /** The server no longer returns a deleted chat. */ + mockRequestJson.mockImplementation(() => Promise.reject(new Error('Chat not found'))) await act(async () => { answerPost?.() await sleep(300) @@ -2611,6 +2613,55 @@ describe('useChat remount send recovery', () => { } ) + /** + * Another tab deletes the chat while this tab's send waits on the lock. The + * busy refusal then rewrites the chat's history locally; that is not the + * server returning the chat, so the delete must still hold. + */ + it('keeps a chat deleted in another tab empty when a pending send there is refused as busy', async () => { + const history = idleHistory('chat-deleted-in-other-tab') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let answerPost: (() => void) | 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))) + return new Promise((resolve) => { + answerPost = () => + resolve( + Response.json( + { error: 'A response is already in progress for this chat.' }, + { status: 409 } + ) + ) + }) + } + if (String(input).includes('/api/mothership/chat')) { + return Response.json({ error: 'Chat not found' }, { status: 404 }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + await act(async () => { + void getResult().sendMessage('Sent as another tab deleted the chat') + }) + await waitFor(() => answerPost !== undefined) + + mockRequestJson.mockImplementation(() => Promise.reject(new Error('Chat not found'))) + handleMothershipChatStatusEvent( + queryClient, + 'ws-1', + JSON.stringify({ chatId: history.id, type: 'deleted', timestamp: Date.now() }) + ) + await act(async () => { + answerPost?.() + await sleep(300) + }) + useMothershipQueueStore.getState().enqueue(history.id, { id: 'later', content: 'Later' }) + + expect(useMothershipQueueStore.getState().queues[history.id]).toBeUndefined() + expect(state.postBodies).toHaveLength(1) + }) + /** * With Redis down the server refuses every send as busy without naming a * turn, and nothing is running. The message must be retried on a growing @@ -2911,7 +2962,8 @@ describe('useChat remount send recovery', () => { return fetchStub(input, init) }) useMothershipQueueStore.getState().clearChat(history.id) - const { getResult } = renderUseChatInChat(history.id, history) + /** Loaded from the server, not seeded: only a server read confirms the chat exists. */ + const { getResult } = renderUseChatInChat(history.id) await waitFor(() => getResult().isSending) await act(async () => { 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 e329fbeffd0..8a44b5e11b5 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -2023,8 +2023,6 @@ export function useChat( const activeStreamId = chatHistory.activeStreamId appliedChatHistoryKeyRef.current = hydrationKey - /** The server returned this chat, so it exists: a delete seen earlier no longer applies. */ - useMothershipQueueStore.getState().reopenChat(chatHistory.id) const mappedMessages = chatHistory.messages.map(toDisplayMessage) const shouldReconnectActiveStream = Boolean(activeStreamId) && diff --git a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts new file mode 100644 index 00000000000..085620fe6c9 --- /dev/null +++ b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts @@ -0,0 +1,62 @@ +import { QueryClient } from '@tanstack/react-query' +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const { mockRequestJson } = vi.hoisted(() => ({ mockRequestJson: vi.fn() })) + +vi.mock('@/lib/api/client/request', async (importOriginal) => ({ + ...(await importOriginal()), + requestJson: mockRequestJson, +})) + +import { + type MothershipChatHistory, + mothershipChatHistoryQueryOptions, +} from '@/hooks/queries/mothership-chats' +import { useMothershipQueueStore } from '@/stores/mothership-queue/store' + +const history: MothershipChatHistory = { + id: 'chat-1', + mode: 'agent', + title: 'Restored', + messages: [], + activeStreamId: null, + resources: [], +} + +/** Whether the queue store takes a send for the chat, i.e. whether its delete still holds. */ +function takesSends(chatId: string): boolean { + useMothershipQueueStore.getState().enqueue(chatId, { id: 'probe', content: 'probe' }) + return useMothershipQueueStore.getState().queues[chatId] !== undefined +} + +describe('chat history read after a delete', () => { + beforeEach(() => { + useMothershipQueueStore.getState().reset() + mockRequestJson.mockReset() + }) + + it('reopens a chat this tab saw deleted once the server returns it again', async () => { + useMothershipQueueStore.getState().clearChat(history.id) + mockRequestJson.mockResolvedValue({ chat: history }) + + await new QueryClient().fetchQuery(mothershipChatHistoryQueryOptions(history.id)) + + expect(takesSends(history.id)).toBe(true) + }) + + it('keeps the delete when the read returning the chat began before it', async () => { + let answer!: () => void + mockRequestJson.mockReturnValue( + new Promise((resolve) => { + answer = () => resolve({ chat: history }) + }) + ) + + const read = new QueryClient().fetchQuery(mothershipChatHistoryQueryOptions(history.id)) + useMothershipQueueStore.getState().clearChat(history.id) + answer() + await read + + expect(takesSends(history.id)).toBe(false) + }) +}) diff --git a/apps/sim/hooks/queries/mothership-chats.ts b/apps/sim/hooks/queries/mothership-chats.ts index 487257a7c3e..0c1196c9792 100644 --- a/apps/sim/hooks/queries/mothership-chats.ts +++ b/apps/sim/hooks/queries/mothership-chats.ts @@ -309,10 +309,25 @@ export async function fetchMothershipChatHistory( return parseChatHistory(await copilotRes.json()) } +/** + * A chat this tab saw deleted that the server returns again was restored, so it + * takes queued sends again. Only a read that began after the delete counts: one + * already in flight can return the chat from before it. + */ +async function fetchChatHistoryConfirmingRestore( + chatId: string, + signal: AbortSignal +): Promise { + const deletedBeforeRead = Boolean(useMothershipQueueStore.getState().cleared[chatId]) + const history = await fetchMothershipChatHistory(chatId, signal) + if (deletedBeforeRead) useMothershipQueueStore.getState().reopenChat(chatId) + return history +} + export function mothershipChatHistoryQueryOptions(chatId: string | undefined) { return queryOptions({ queryKey: mothershipChatKeys.detail(chatId), - queryFn: chatId ? ({ signal }) => fetchMothershipChatHistory(chatId, signal) : skipToken, + queryFn: chatId ? ({ signal }) => fetchChatHistoryConfirmingRestore(chatId, signal) : skipToken, staleTime: MOTHERSHIP_CHAT_HISTORY_STALE_TIME, }) } From ee245ebc63ca2c2c77425a586769ba3d6a11c33e Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:11:51 -0700 Subject: [PATCH 4/9] fix(mothership): apply the restore check to every server read of a chat Recovery on a return to the tab reads the chat directly, not through the history query, so a restore it found left the chat refusing follow-ups. The check now lives in fetchMothershipChatHistory itself. --- .../home/hooks/use-chat.dom.test.tsx | 32 +++++++++++++++++++ apps/sim/hooks/queries/mothership-chats.ts | 17 +++++----- 2 files changed, 41 insertions(+), 8 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 68ab1baf0fb..85d7511244c 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 @@ -2613,6 +2613,38 @@ describe('useChat remount send recovery', () => { } ) + /** The same holds when the server read is the one a return to the tab makes. */ + it('queues a follow-up after a return to the tab finds a chat this tab saw deleted', async () => { + const history = idleHistory('chat-restored-while-away') + const running: MothershipChatHistory = { ...history, activeStreamId: 'turn-after-restore' } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: running })) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).includes('/api/mothership/chat/stream')) { + if (String(input).includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + useMothershipQueueStore.getState().clearChat(history.id) + + await act(async () => { + window.dispatchEvent(new Event('pageshow')) + }) + await waitFor(() => getResult().isSending) + await act(async () => { + await getResult().sendMessage('Follow-up after coming back') + }) + + expect( + useMothershipQueueStore.getState().queues[history.id]?.map((message) => message.content) + ).toEqual(['Follow-up after coming back']) + }) + /** * Another tab deletes the chat while this tab's send waits on the lock. The * busy refusal then rewrites the chat's history locally; that is not the diff --git a/apps/sim/hooks/queries/mothership-chats.ts b/apps/sim/hooks/queries/mothership-chats.ts index 0c1196c9792..055caba4451 100644 --- a/apps/sim/hooks/queries/mothership-chats.ts +++ b/apps/sim/hooks/queries/mothership-chats.ts @@ -281,7 +281,7 @@ export function useOrganizationMothershipChats( }) } -export async function fetchMothershipChatHistory( +async function readMothershipChatHistory( chatId: string, signal?: AbortSignal ): Promise { @@ -310,16 +310,17 @@ export async function fetchMothershipChatHistory( } /** - * A chat this tab saw deleted that the server returns again was restored, so it - * takes queued sends again. Only a read that began after the delete counts: one - * already in flight can return the chat from before it. + * Reads a chat from the server. A chat this tab saw deleted that the server + * returns again was restored, so it takes queued sends again. Only a read that + * began after the delete counts: one already in flight can return the chat from + * before it. */ -async function fetchChatHistoryConfirmingRestore( +export async function fetchMothershipChatHistory( chatId: string, - signal: AbortSignal + signal?: AbortSignal ): Promise { const deletedBeforeRead = Boolean(useMothershipQueueStore.getState().cleared[chatId]) - const history = await fetchMothershipChatHistory(chatId, signal) + const history = await readMothershipChatHistory(chatId, signal) if (deletedBeforeRead) useMothershipQueueStore.getState().reopenChat(chatId) return history } @@ -327,7 +328,7 @@ async function fetchChatHistoryConfirmingRestore( export function mothershipChatHistoryQueryOptions(chatId: string | undefined) { return queryOptions({ queryKey: mothershipChatKeys.detail(chatId), - queryFn: chatId ? ({ signal }) => fetchChatHistoryConfirmingRestore(chatId, signal) : skipToken, + queryFn: chatId ? ({ signal }) => fetchMothershipChatHistory(chatId, signal) : skipToken, staleTime: MOTHERSHIP_CHAT_HISTORY_STALE_TIME, }) } From aae0be60d1daeb2eaa3524eeed5d7d1e8f7993af Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:24:08 -0700 Subject: [PATCH 5/9] test(mothership): use the central request mock in the restore read test --- .../queries/mothership-chat-history-restore.test.ts | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts index 085620fe6c9..d48c3cda4b2 100644 --- a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts +++ b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts @@ -1,12 +1,11 @@ +import { + apiClientRequestMock, + apiClientRequestMockFns, +} from '@sim/testing/mocks/api-client-request.mock' import { QueryClient } from '@tanstack/react-query' import { beforeEach, describe, expect, it, vi } from 'vitest' -const { mockRequestJson } = vi.hoisted(() => ({ mockRequestJson: vi.fn() })) - -vi.mock('@/lib/api/client/request', async (importOriginal) => ({ - ...(await importOriginal()), - requestJson: mockRequestJson, -})) +vi.mock('@/lib/api/client/request', () => apiClientRequestMock) import { type MothershipChatHistory, @@ -14,6 +13,8 @@ import { } from '@/hooks/queries/mothership-chats' import { useMothershipQueueStore } from '@/stores/mothership-queue/store' +const mockRequestJson = apiClientRequestMockFns.mockRequestJson + const history: MothershipChatHistory = { id: 'chat-1', mode: 'agent', From 4f6030924720bf433a640fdc08a5742e20c9db86 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:53:21 -0700 Subject: [PATCH 6/9] fix(mothership): lift only the chat delete a read or restore saw Each delete now gets a token. A history read or a restore lifts the delete it saw when it began, so a slow answer cannot reopen a chat deleted again while it was in flight. --- .../mothership-chat-history-restore.test.ts | 80 ++++++++++++++++--- apps/sim/hooks/queries/mothership-chats.ts | 11 ++- .../sim/stores/mothership-queue/store.test.ts | 18 +++++ apps/sim/stores/mothership-queue/store.ts | 18 +++-- apps/sim/stores/mothership-queue/types.ts | 14 ++-- 5 files changed, 114 insertions(+), 27 deletions(-) diff --git a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts index d48c3cda4b2..e14c3394303 100644 --- a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts +++ b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts @@ -2,14 +2,16 @@ import { apiClientRequestMock, apiClientRequestMockFns, } from '@sim/testing/mocks/api-client-request.mock' -import { QueryClient } from '@tanstack/react-query' +import { reactQueryMock } from '@sim/testing/mocks/react-query.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/api/client/request', () => apiClientRequestMock) +vi.mock('@tanstack/react-query', () => reactQueryMock) import { + fetchMothershipChatHistory, type MothershipChatHistory, - mothershipChatHistoryQueryOptions, + useRestoreMothershipChat, } from '@/hooks/queries/mothership-chats' import { useMothershipQueueStore } from '@/stores/mothership-queue/store' @@ -30,7 +32,36 @@ function takesSends(chatId: string): boolean { return useMothershipQueueStore.getState().queues[chatId] !== undefined } -describe('chat history read after a delete', () => { +/** A server answer the test releases when it chooses. */ +function deferredAnswer(value: unknown): () => void { + let answer!: () => void + mockRequestJson.mockReturnValue( + new Promise((resolve) => { + answer = () => resolve(value) + }) + ) + return answer +} + +/** The options `useRestoreMothershipChat` hands to `useMutation` (the mock returns them). */ +interface RestoreMutation { + mutationFn: (chatId: string) => Promise + onMutate?: (chatId: string) => { deleteSeen?: number } + onSuccess: (data: undefined, chatId: string, context?: { deleteSeen?: number }) => void +} + +async function restore(chatId: string, whileInFlight: () => void = () => {}) { + const mutation = useRestoreMothershipChat() as unknown as RestoreMutation + const context = mutation.onMutate?.(chatId) + const answer = deferredAnswer({ success: true }) + const done = mutation.mutationFn(chatId) + whileInFlight() + answer() + await done + mutation.onSuccess(undefined, chatId, context) +} + +describe('lifting a chat delete', () => { beforeEach(() => { useMothershipQueueStore.getState().reset() mockRequestJson.mockReset() @@ -40,24 +71,49 @@ describe('chat history read after a delete', () => { useMothershipQueueStore.getState().clearChat(history.id) mockRequestJson.mockResolvedValue({ chat: history }) - await new QueryClient().fetchQuery(mothershipChatHistoryQueryOptions(history.id)) + await fetchMothershipChatHistory(history.id) expect(takesSends(history.id)).toBe(true) }) it('keeps the delete when the read returning the chat began before it', async () => { - let answer!: () => void - mockRequestJson.mockReturnValue( - new Promise((resolve) => { - answer = () => resolve({ chat: history }) - }) - ) - - const read = new QueryClient().fetchQuery(mothershipChatHistoryQueryOptions(history.id)) + const answer = deferredAnswer({ chat: history }) + const read = fetchMothershipChatHistory(history.id) + useMothershipQueueStore.getState().clearChat(history.id) + answer() + await read + + expect(takesSends(history.id)).toBe(false) + }) + + it('keeps a newer delete that lands while a read after an earlier one is in flight', async () => { + useMothershipQueueStore.getState().clearChat(history.id) + const answer = deferredAnswer({ chat: history }) + const read = fetchMothershipChatHistory(history.id) + useMothershipQueueStore.getState().reopenChat(history.id) useMothershipQueueStore.getState().clearChat(history.id) answer() await read expect(takesSends(history.id)).toBe(false) }) + + it('reopens a chat restored from Recently Deleted', async () => { + useMothershipQueueStore.getState().clearChat(history.id) + + await restore(history.id) + + expect(takesSends(history.id)).toBe(true) + }) + + it('keeps a delete that lands while the restore is in flight', async () => { + useMothershipQueueStore.getState().clearChat(history.id) + + await restore(history.id, () => { + useMothershipQueueStore.getState().reopenChat(history.id) + useMothershipQueueStore.getState().clearChat(history.id) + }) + + expect(takesSends(history.id)).toBe(false) + }) }) diff --git a/apps/sim/hooks/queries/mothership-chats.ts b/apps/sim/hooks/queries/mothership-chats.ts index 055caba4451..5bcdc5966f4 100644 --- a/apps/sim/hooks/queries/mothership-chats.ts +++ b/apps/sim/hooks/queries/mothership-chats.ts @@ -319,9 +319,9 @@ export async function fetchMothershipChatHistory( chatId: string, signal?: AbortSignal ): Promise { - const deletedBeforeRead = Boolean(useMothershipQueueStore.getState().cleared[chatId]) + const deleteSeen = useMothershipQueueStore.getState().cleared[chatId] const history = await readMothershipChatHistory(chatId, signal) - if (deletedBeforeRead) useMothershipQueueStore.getState().reopenChat(chatId) + if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen) return history } @@ -381,8 +381,11 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) { const queryClient = useQueryClient() return useMutation({ mutationFn: restoreChat, - onSuccess: (_data, chatId) => { - useMothershipQueueStore.getState().reopenChat(chatId) + /** The delete this restore undoes; one that lands while it is in flight stays. */ + onMutate: (chatId) => ({ deleteSeen: useMothershipQueueStore.getState().cleared[chatId] }), + onSuccess: (_data, chatId, context) => { + if (context?.deleteSeen === undefined) return + useMothershipQueueStore.getState().reopenChat(chatId, context.deleteSeen) }, onSettled: () => { queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) }) diff --git a/apps/sim/stores/mothership-queue/store.test.ts b/apps/sim/stores/mothership-queue/store.test.ts index bb3c4601211..145be699d68 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -89,6 +89,24 @@ describe('useMothershipQueueStore', () => { expect(useMothershipQueueStore.getState().queues['pending::abc']).toBeUndefined() }) + it('lifts only the delete a restore saw, never a later one', () => { + useMothershipQueueStore.getState().clearChat('chat-X') + const seen = useMothershipQueueStore.getState().cleared['chat-X'] + useMothershipQueueStore.getState().reopenChat('chat-X') + useMothershipQueueStore.getState().clearChat('chat-X') + + useMothershipQueueStore.getState().reopenChat('chat-X', seen) + useMothershipQueueStore.getState().enqueue('chat-X', message('after-stale-restore')) + expect(useMothershipQueueStore.getState().queues['chat-X']).toBeUndefined() + + const latest = useMothershipQueueStore.getState().cleared['chat-X'] + useMothershipQueueStore.getState().reopenChat('chat-X', latest) + useMothershipQueueStore.getState().enqueue('chat-X', message('after-restore')) + expect(useMothershipQueueStore.getState().queues['chat-X']?.map((m) => m.id)).toEqual([ + 'after-restore', + ]) + }) + it('does not move a new chat surface queue into a chat deleted meanwhile', () => { useMothershipQueueStore.getState().enqueue('pending::abc', message('pending-1')) useMothershipQueueStore.getState().clearChat('chat-X') diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index f69bba53f78..c8003340501 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -41,10 +41,13 @@ const sessionStorageAdapter = { }, } +/** Numbers each delete, so a restore or read can tell the delete it saw from a later one. */ +let deleteCount = 0 + const initialState = { queues: {} as Record, editing: {} as Record, - cleared: {} as Record, + cleared: {} as Record, } const omitKey = (record: Record, key: string): Record => { @@ -206,13 +209,16 @@ export const useMothershipQueueStore = create()( set((state) => ({ queues: omitKey(state.queues, chatKey), editing: omitKey(state.editing, chatKey), - cleared: { ...state.cleared, [chatKey]: true }, + cleared: { ...state.cleared, [chatKey]: ++deleteCount }, })), - reopenChat: (chatKey) => - set((state) => - state.cleared[chatKey] ? { cleared: omitKey(state.cleared, chatKey) } : state - ), + reopenChat: (chatKey, deleteToken) => + set((state) => { + const current = state.cleared[chatKey] + if (current === undefined) return state + if (deleteToken !== undefined && deleteToken !== current) return state + return { cleared: omitKey(state.cleared, chatKey) } + }), reset: () => set(initialState), }), diff --git a/apps/sim/stores/mothership-queue/types.ts b/apps/sim/stores/mothership-queue/types.ts index ad797c01ea9..c5102343fed 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -53,10 +53,11 @@ export interface MothershipQueueState { queues: Record editing: Record /** - * Chats cleared this session (deleted). No write recreates their queue (a late - * restore, or a failed send handed back); restoring the chat lifts it. + * Chats cleared this session (deleted), each with the token of its latest + * delete. No write recreates their queue (a late restore, or a failed send + * handed back); restoring the chat lifts it. */ - cleared: Record + cleared: Record enqueue: (chatKey: string, message: QueuedMothershipMessage) => void insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void @@ -69,7 +70,10 @@ export interface MothershipQueueState { /** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */ adoptHeldSends: (toKey: string, surface: string) => void clearChat: (chatKey: string) => void - /** Lifts `cleared` for a chat restored from Recently Deleted. */ - reopenChat: (chatKey: string) => void + /** + * Lifts `cleared` for a restored chat. Given the delete token an operation + * saw when it began, lifts only that delete, never one that landed after it. + */ + reopenChat: (chatKey: string, deleteToken?: number) => void reset: () => void } From cf6651d034ab4d1207f4b535dffdaa3aee0f2b5a Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:53:30 -0700 Subject: [PATCH 7/9] fix(mothership): keep a send refused as busy after the user switched chats A chat switch detaches the view but leaves a send waiting on the chat lock running. Its busy refusal was dropped as a stale answer, so the message was lost. A conflict is now handled after a switch too, re-queuing the message in its own chat. The busy branch attaches through recovery's own history read instead of reading the chat twice. --- .../home/hooks/use-chat.dom.test.tsx | 66 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 49 +++++++------- 2 files changed, 92 insertions(+), 23 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 85d7511244c..818c3b44de8 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 @@ -2528,6 +2528,72 @@ describe('useChat remount send recovery', () => { expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) }) + /** + * A send can wait on the chat lock while the user switches chats. The switch + * detaches the view but leaves the POST running, so the busy refusal that + * answers it must still put the message back in its own chat's queue. + */ + it.each(['direct', 'queued'] as const)( + 'keeps a %s send refused as busy after the user switched chats', + async (origin) => { + const history = idleHistory(`chat-busy-after-switch-${origin}`) + const other = idleHistory(`chat-switched-to-during-lock-${origin}`) + mockRequestJson.mockImplementation((_contract: AnyApiRouteContract, input: unknown) => + Promise.resolve({ + chat: JSON.stringify(input).includes(other.id) ? other : history, + }) + ) + let answerPost: (() => void) | 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))) + return new Promise((resolve) => { + answerPost = () => + resolve( + Response.json( + { + error: 'A response is already in progress for this chat.', + activeStreamId: 'turn-from-another-tab', + }, + { status: 409 } + ) + ) + }) + } + if (String(input).includes('/api/mothership/chat/stream')) { + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + if (origin === 'queued') { + useMothershipQueueStore + .getState() + .enqueue(history.id, { id: 'queued-on-lock', content: 'Waiting on the lock' }) + } + const { getResult, navigate } = renderUseChatInChat(history.id, history) + if (origin === 'direct') { + await act(async () => { + void getResult().sendMessage('Waiting on the lock') + }) + } + await waitFor(() => answerPost !== undefined) + + navigate(other.id, other) + await act(async () => { + answerPost?.() + await sleep(300) + }) + + const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] + expect(queued.map((message) => message.content)).toEqual(['Waiting on the lock']) + expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) + expect(useMothershipQueueStore.getState().queues[other.id]).toBeUndefined() + expect(state.postBodies).toHaveLength(1) + } + ) + /** A chat deleted while its queued send was failing must not get that send back. */ it('does not recreate the queue of a chat deleted while its queued send was failing', async () => { const history = idleHistory('chat-deleted-mid-dispatch') 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 8a44b5e11b5..d43c72e55f3 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -3930,7 +3930,10 @@ export function useChat( : undefined resolveAdmission?.(admittedChatId) if (pendingChatAdmissionRef.current === admission) pendingChatAdmissionRef.current = null - if (streamGenRef.current !== gen) { + /* The user moved on (another chat) while the POST was out. A conflict is still + handled below: it means the message was not admitted, and must be re-queued + in its own chat rather than read as a turn this view no longer shows. */ + if (streamGenRef.current !== gen && response.status !== 409) { await response.body?.cancel() return consumedByTranscript } @@ -3964,6 +3967,8 @@ export function useChat( turn names that turn, or nothing when its stream id is unreadable. */ const conflictStreamId = typeof errorData.activeStreamId === 'string' ? errorData.activeStreamId : undefined + /** Whether this view still shows the send; otherwise only its own chat changes. */ + const viewOnSend = streamGenRef.current === gen const supersededStreamId = queuedSendHandoff?.supersededStreamId ?? pendingStopStreamId if (supersededStreamId && conflictStreamId === supersededStreamId) { rollbackOptimisticSend() @@ -3974,7 +3979,8 @@ export function useChat( clearActiveTurn() setTransportIdle() } - setError('Previous response is still shutting down; queued message was restored.') + if (viewOnSend) + setError('Previous response is still shutting down; queued message was restored.') return { userMessageId, held: true } } /** Withdraws this refused send so the queue retries it, under the same id, later. */ @@ -3992,28 +3998,25 @@ export function useChat( /* Another turn holds the chat: one started in another tab, or one this surface lost track of. This message was not admitted (the server released its id), so it goes back to the queue under the same id. The - queue drains only while the chat is idle, so the chat's running turn is - read before the message is handed back: the chat then attaches to that - turn and the message goes out once, after it ends. */ + queue drains only while the chat is idle, so the chat records the turn + the refusal names before the message is handed back. */ releaseRefusedSend() - if (requestChatId) { - const busyChatId = requestChatId - if (conflictStreamId) { - upsertChatHistory(busyChatId, (current) => ({ - ...current, - activeStreamId: conflictStreamId, - })) - } - await queryClient - .refetchQueries({ queryKey: mothershipChatKeys.detail(busyChatId), exact: true }) - .catch(() => {}) + if (requestChatId && conflictStreamId) { + upsertChatHistory(requestChatId, (current) => ({ + ...current, + activeStreamId: conflictStreamId, + })) } - /* The history read above can repeat one this surface skipped while the POST - was pending, so nothing re-runs to attach to the turn it lists. Attach - explicitly; the message is retried after that turn ends. */ - if (pendingChatAdmissionRef.current === admission) - pendingChatAdmissionRef.current = null - void recoverActiveStreamRef.current('busy_refusal') + /* Attach to the turn that holds the chat: recovery reads the chat once and + shows its running turn, and the message is retried after that turn ends. + A view that moved to another chat leaves it to load fresh when reopened. */ + if (viewOnSend) void recoverActiveStreamRef.current('busy_refusal') + else if (requestChatId) + void queryClient.invalidateQueries({ + queryKey: mothershipChatKeys.detail(requestChatId), + exact: true, + refetchType: 'none', + }) return { userMessageId, busy: true } } /* A send deduplicated against an earlier attempt comes back naming @@ -4022,7 +4025,7 @@ export function useChat( chat before the reconnect below replays it. */ const conflictChatId = typeof errorData.chatId === 'string' ? errorData.chatId : undefined - if (conflictChatId && !streamTargetChatId) { + if (viewOnSend && conflictChatId && !streamTargetChatId) { adoptNewChatEffort(conflictChatId, false) adoptResolvedChatId(conflictChatId, { replaceHomeHistory: true, From cec0ebb863af56e315e728863e51a61f451b3c56 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 16:13:51 -0700 Subject: [PATCH 8/9] fix(mothership): check a deduplicated send's stream before adopting its chat On the new-chat surface, an "already sent" answer names the chat the earlier attempt opened. Adopting it before finding no stream moved the surface to that chat while the retried message was queued under the new-chat key, where nothing sent it. The stream is now checked first. --- .../home/hooks/use-chat.dom.test.tsx | 42 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 32 +++++++------- 2 files changed, 59 insertions(+), 15 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 818c3b44de8..4243879eb2d 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 @@ -3073,6 +3073,48 @@ describe('useChat remount send recovery', () => { ).toEqual(['Follow-up after the restore']) }) + /** + * On the new-chat surface the "already sent" answer also names the chat the + * earlier attempt opened. With no stream yet, the retry must still go out, + * not wait under the new-chat key after the surface moved to that chat. + */ + it('retries a first message deduplicated against an attempt that opened no stream', async () => { + const opened = idleHistory('chat-opened-by-earlier-attempt') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: opened })) + let posts = 0 + 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))) + posts++ + if (posts === 1) { + return Response.json( + { + error: 'This message was already sent.', + activeStreamId: state.postBodies[0].userMessageId, + chatId: opened.id, + }, + { status: 409 } + ) + } + return emptySseResponse() + } + if (url.includes('/api/mothership/chat/stream') && posts === 1) { + return Response.json({ error: 'Stream not found' }, { status: 404 }) + } + return fetchStub(input, init) + }) + const { getResult } = renderHomeLikeSurface() + + await act(async () => { + await getResult().sendMessage('First message, told it was already sent') + }) + await waitFor(() => state.postBodies.length === 2, 5_000) + + expect(state.postBodies[1].message).toBe('First message, told it was already sent') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + /** The `online` event can fire while no surface for the chat is mounted. */ it('sends a held message when its chat mounts after the network came back', async () => { const history = idleHistory('chat-held-while-away') 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 d43c72e55f3..e169d5c72b2 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -4019,24 +4019,12 @@ export function useChat( }) return { userMessageId, busy: true } } - /* A send deduplicated against an earlier attempt comes back naming - the chat that attempt opened. Adopting it here spares a chatless - surface the stream-to-chat lookup and puts the user in the right - chat before the reconnect below replays it. */ - const conflictChatId = - typeof errorData.chatId === 'string' ? errorData.chatId : undefined - if (viewOnSend && conflictChatId && !streamTargetChatId) { - adoptNewChatEffort(conflictChatId, false) - adoptResolvedChatId(conflictChatId, { - replaceHomeHistory: true, - invalidateList: true, - }) - streamTargetChatId = conflictChatId - } /* "Already sent" with no stream for it means the earlier attempt is still in flight on the server (or died before starting a turn), not that a turn ran: reattaching would read the missing stream as finished and drop the - message. Retry it later like a busy refusal; the server's claim settles. */ + message. Retry it later like a busy refusal; the server's claim settles. + This is checked before adopting the chat the answer names, so a retried + message stays under the key it was sent from. */ const dedupedStreamExists = await fetchStreamBatch( conflictStreamId, '0', @@ -4051,6 +4039,20 @@ export function useChat( } /** The user may have moved on (another chat, another send) during the check. */ if (streamGenRef.current !== gen) return consumedByTranscript + /* A send deduplicated against an earlier attempt comes back naming + the chat that attempt opened. Adopting it here spares a chatless + surface the stream-to-chat lookup and puts the user in the right + chat before the reconnect below replays it. */ + const conflictChatId = + typeof errorData.chatId === 'string' ? errorData.chatId : undefined + if (conflictChatId && !streamTargetChatId) { + adoptNewChatEffort(conflictChatId, false) + adoptResolvedChatId(conflictChatId, { + replaceHomeHistory: true, + invalidateList: true, + }) + streamTargetChatId = conflictChatId + } streamIdRef.current = conflictStreamId const succeeded = await retryReconnect({ streamId: conflictStreamId, From be44382caa8830ab950d71402d4c7c7d6fc5a690 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 16:34:39 -0700 Subject: [PATCH 9/9] fix(mothership): retry a deduplicated send whose stream lookup failed A lookup that failed with a network error or 5xx was read as proof that the earlier attempt's stream existed, so a send that never started was reattached and finalized instead of re-sent. Only a lookup this send aborted skips the retry now; the server deduplicates the retry by id. --- .../home/hooks/use-chat.dom.test.tsx | 37 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 7 +++- 2 files changed, 42 insertions(+), 2 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 4243879eb2d..2cffbbc4628 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 @@ -3073,6 +3073,43 @@ describe('useChat remount send recovery', () => { ).toEqual(['Follow-up after the restore']) }) + /** A stream lookup that fails for another reason does not prove a turn ran either. */ + it('retries a deduplicated send whose stream lookup failed', async () => { + const history = idleHistory('chat-deduped-lookup-failed') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let posts = 0 + 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))) + posts++ + if (posts === 1) { + return Response.json( + { + error: 'This message was already sent.', + activeStreamId: state.postBodies[0].userMessageId, + }, + { status: 409 } + ) + } + return emptySseResponse() + } + if (url.includes('/api/mothership/chat/stream') && posts === 1) { + return Response.json({ error: 'Internal error' }, { status: 500 }) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + + await act(async () => { + await getResult().sendMessage('Told it was already sent, lookup failed') + }) + await waitFor(() => state.postBodies.length === 2, 5_000) + + expect(state.postBodies[1].message).toBe('Told it was already sent, lookup failed') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + /** * On the new-chat surface the "already sent" answer also names the chat the * earlier attempt opened. With no stream yet, the retry must still go out, 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 e169d5c72b2..ff5db89ef70 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -4024,14 +4024,17 @@ export function useChat( ran: reattaching would read the missing stream as finished and drop the message. Retry it later like a busy refusal; the server's claim settles. This is checked before adopting the chat the answer names, so a retried - message stays under the key it was sent from. */ + message stays under the key it was sent from. A lookup that fails for + another reason proves nothing either way, so it is retried too: the server + deduplicates the retry by id. Only a lookup this send aborted (Stop, or the + user moving on) is not retried. */ const dedupedStreamExists = await fetchStreamBatch( conflictStreamId, '0', abortController.signal ).then( () => true, - (error: unknown) => !isStreamGoneError(error) + (error: unknown) => !isStreamGoneError(error) && abortController.signal.aborted ) if (!dedupedStreamExists) { releaseRefusedSend()