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..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 @@ -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') @@ -2475,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') @@ -2510,6 +2629,529 @@ 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) + /** The server no longer returns a deleted chat. */ + mockRequestJson.mockImplementation(() => Promise.reject(new Error('Chat not found'))) + await act(async () => { + answerPost?.() + await sleep(300) + }) + + expect(useMothershipQueueStore.getState().queues[history.id]).toBeUndefined() + expect(state.postBodies).toHaveLength(1) + } + ) + + /** 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 + * 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 + * 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 })) + 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( + { 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) + + /** 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) + } 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) + }) + + /** + * 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) + /** 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 () => { + await getResult().sendMessage('Follow-up after the restore') + }) + + expect( + useMothershipQueueStore.getState().queues[history.id]?.map((message) => message.content) + ).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, + * 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') @@ -2869,6 +3511,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..ff5db89ef70 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 } } } @@ -3890,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 } @@ -3924,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() @@ -3934,16 +3979,12 @@ export function useChat( clearActiveTurn() setTransportIdle() } - setError('Previous response is still shutting down; queued message was restored.') - return false + if (viewOnSend) + setError('Previous response is still shutting down; queued message was restored.') + 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,20 +3993,55 @@ export function useChat( clearActiveTurn() setTransportIdle() } - if (requestChatId) { - const busyChatId = requestChatId - if (conflictStreamId) { - upsertChatHistory(busyChatId, (current) => ({ - ...current, - activeStreamId: conflictStreamId, - })) - } - await queryClient - .refetchQueries({ queryKey: mothershipChatKeys.detail(busyChatId), exact: true }) - .catch(() => {}) + } + 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 records the turn + the refusal names before the message is handed back. */ + releaseRefusedSend() + if (requestChatId && conflictStreamId) { + upsertChatHistory(requestChatId, (current) => ({ + ...current, + activeStreamId: conflictStreamId, + })) } - return { userMessageId } + /* 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 } } + /* "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. + This is checked before adopting the chat the answer names, so a retried + 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) && abortController.signal.aborted + ) + if (!dedupedStreamExists) { + releaseRefusedSend() + return { userMessageId, busy: true } + } + /** 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 @@ -4147,6 +4223,7 @@ export function useChat( finalize, resumeOrFinalize, retryReconnect, + fetchStreamBatch, clearActiveTurn, resetStreamingBuffers, resolveChatIdForStream, @@ -4284,7 +4361,12 @@ export function useChat( ? { assistantSearchLevel: options?.assistantSearchLevel } : {}), } - if (!result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + if ( + !result.unreachable && + !result.held && + !result.busy && + activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ) { handOffWithdrawnSend(withdrawn) return } @@ -4301,7 +4383,9 @@ export function useChat( ...(result.unreachable && !result.networkReturned ? { retryRequired: true, heldUntilOnline: true } : {}), - ...(result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ...(result.held ? { retryRequired: true } : {}), + ...(result.busy ? busyRetry(1) : {}), + ...((result.unreachable || result.busy) && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), }) @@ -4861,7 +4945,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 +4962,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 @@ -4892,7 +4979,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, @@ -4913,10 +5005,12 @@ 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 } : {}), - ...(withdrawn?.unreachable && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ...((withdrawn?.unreachable || withdrawn?.busy) && + dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), @@ -4995,6 +5089,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 +5270,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 +5293,8 @@ export function useChat( scopeKey, messageQueue.length, queueHeadHeld, + queueHeadNotBefore, + busyRetryWakeup, resolvedChatId, chatHistoryReady, remoteActiveStreamId, 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..e14c3394303 --- /dev/null +++ b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts @@ -0,0 +1,119 @@ +import { + apiClientRequestMock, + apiClientRequestMockFns, +} from '@sim/testing/mocks/api-client-request.mock' +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, + useRestoreMothershipChat, +} from '@/hooks/queries/mothership-chats' +import { useMothershipQueueStore } from '@/stores/mothership-queue/store' + +const mockRequestJson = apiClientRequestMockFns.mockRequestJson + +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 +} + +/** 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() + }) + + 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 fetchMothershipChatHistory(history.id) + + expect(takesSends(history.id)).toBe(true) + }) + + it('keeps the delete when the read returning the chat began before it', async () => { + 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 df914c9220c..5bcdc5966f4 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 { @@ -309,6 +309,22 @@ export async function fetchMothershipChatHistory( return parseChatHistory(await copilotRes.json()) } +/** + * 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. + */ +export async function fetchMothershipChatHistory( + chatId: string, + signal?: AbortSignal +): Promise { + const deleteSeen = useMothershipQueueStore.getState().cleared[chatId] + const history = await readMothershipChatHistory(chatId, signal) + if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen) + return history +} + export function mothershipChatHistoryQueryOptions(chatId: string | undefined) { return queryOptions({ queryKey: mothershipChatKeys.detail(chatId), @@ -365,6 +381,12 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) { const queryClient = useQueryClient() return useMutation({ mutationFn: restoreChat, + /** 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/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..145be699d68 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -88,5 +88,32 @@ 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') + 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 418ef9d8bd6..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 => { @@ -67,13 +70,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 +104,8 @@ export const useMothershipQueueStore = create()( retryRequired: _retry, heldUntilOnline: _held, heldSurface: _surface, + busyRetries: _busyRetries, + notBefore: _notBefore, ...rest } = next[index] next[index] = { @@ -144,7 +151,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] ?? [] @@ -201,9 +209,17 @@ 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, 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 af41dea707c..c5102343fed 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,10 +53,11 @@ 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), 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 @@ -65,5 +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 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 }