diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/desktop-tool-lifetimes.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/desktop-tool-lifetimes.ts index 5189f62177f..9e2f49aeee1 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/desktop-tool-lifetimes.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/desktop-tool-lifetimes.ts @@ -17,7 +17,7 @@ const runningTurns = new Map() /** A running desktop tool's hold on its turn. */ interface DesktopToolLease { - /** Aborted only by the user's Stop of the turn. */ + /** Aborted only by the user's Stop of the turn, or by signing out. */ signal: AbortSignal /** Called once when the tool settles. */ release(): void @@ -25,8 +25,9 @@ interface DesktopToolLease { /** * Starts a desktop tool (a browser action, a local file read or import) for a turn. Only the - * user's Stop of that turn cancels it: replacing the stream reader, leaving the chat view, or - * stopping another chat's turn leaves it running to finish and report its own result. + * user's Stop of that turn, or signing out (`stopAllDesktopTools`), cancels it: replacing the + * stream reader, leaving the chat view, or stopping another chat's turn leaves it running to + * finish and report its own result. */ export function leaseDesktopTool(streamId: string): DesktopToolLease { let turn = runningTurns.get(streamId) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index affe8a1c15c..7a7556a950c 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 @@ -507,7 +507,7 @@ describe('useChat remount send recovery', () => { state.stopBodies = [] state.abortTraceparents = [] mockRequestJson.mockResolvedValue({ chats: [] }) - useMothershipQueueStore.setState({ queues: {}, editing: {} }) + useMothershipQueueStore.setState({ queues: {}, editing: {}, cleared: {} }) useExecutionStore.setState({ workflowExecutions: new Map() }) window.sessionStorage.clear() window.localStorage.clear() @@ -1921,6 +1921,376 @@ describe('useChat remount send recovery', () => { }) }) + describe('a send the server never admitted', () => { + const idleHistory = (id: string): MothershipChatHistory => ({ + id, + mode: 'agent', + title: 'Not admitted', + messages: [], + activeStreamId: null, + resources: [], + }) + + /** + * The POST fails at the network layer until the network is back, and the + * stream it would have opened does not exist. + */ + const network = { online: false, acceptedPosts: 0 } + function stubUnreachableSend() { + network.online = false + network.acceptedPosts = 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))) + if (network.online) { + network.acceptedPosts++ + return emptySseResponse() + } + throw new TypeError('Failed to fetch') + } + if (url.includes('/api/mothership/chat/stream')) { + return Response.json({ error: 'Stream not found' }, { status: 404 }) + } + return fetchStub(input, init) + }) + } + + it('holds a message sent while offline and sends it under the same id once back online', async () => { + const history = idleHistory('chat-offline-send') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + stubUnreachableSend() + const { getResult } = renderUseChatInChat(history.id, history) + + await act(async () => { + await getResult().sendMessage('Written while offline') + }) + await waitFor(() => !getResult().isSending) + + const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] + expect(queued.map((message) => message.content)).toEqual(['Written while offline']) + expect(queued[0].retryRequired).toBe(true) + expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) + expect(getResult().error).not.toBeNull() + expect(state.postBodies).toHaveLength(1) + + network.online = true + await act(async () => { + window.dispatchEvent(new Event('online')) + }) + await waitFor(() => state.postBodies.length === 2) + + expect(state.postBodies[1].message).toBe('Written while offline') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + + it('keeps a queued follow-up whose dispatch could not reach the server', async () => { + const history = idleHistory('chat-offline-queue') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + stubUnreachableSend() + useMothershipQueueStore + .getState() + .enqueue(history.id, { id: 'queued-follow-up', content: 'Queued before the drop' }) + renderUseChatInChat(history.id, history) + + await waitFor(() => state.postBodies.length === 1) + await waitFor( + () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true + ) + + const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] + expect(queued.map((message) => message.content)).toEqual(['Queued before the drop']) + expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) + expect(state.postBodies).toHaveLength(1) + }) + + /** + * Another tab's turn holds the chat (this one missed its start). The server + * refuses the send, naming that turn or, when its stream id is unreadable, + * nothing. The message must go out exactly once, under its id, after that turn + * ends: never rendered under the other turn's answer, never lost, never resent + * while the turn still runs. + */ + it.each([ + ['names', 'turn-from-another-tab'], + ['does not name', undefined], + ] as const)( + 'sends a message once, after the turn that held the chat ends, when the refusal %s it', + async (_names, refusalStreamId) => { + const history = idleHistory(`chat-busy-${refusalStreamId ?? 'unnamed'}`) + let otherTurnRunning = true + mockRequestJson.mockImplementation(() => + Promise.resolve({ + chat: { + ...history, + activeStreamId: otherTurnRunning ? 'turn-from-another-tab' : null, + }, + }) + ) + 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 (otherTurnRunning) { + return Response.json( + { + error: 'A response is already in progress for this chat.', + ...(refusalStreamId ? { activeStreamId: refusalStreamId } : {}), + }, + { status: 409 } + ) + } + return emptySseResponse() + } + if (url.includes('/api/mothership/chat/stream')) { + if (url.includes('batch=true')) { + return Response.json({ + success: true, + events: [], + status: otherTurnRunning ? 'streaming' : 'complete', + }) + } + return emptySseResponse() + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChatInChat(history.id, history) + + await act(async () => { + await getResult().sendMessage('Sent from the second tab') + }) + await act(async () => { + await sleep(1500) + }) + expect(state.postBodies).toHaveLength(1) + expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe( + 'Sent from the second tab' + ) + + otherTurnRunning = false + await waitFor(() => state.postBodies.length === 2, 5000) + await act(async () => { + await sleep(500) + }) + + expect(state.postBodies).toHaveLength(2) + expect(state.postBodies[1].message).toBe('Sent from the second tab') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + } + ) + + /** + * A chatless surface keys its queue by mount, so a held first message would + * be stranded by a reload or remount before the network returns. The next + * chatless mount of the same surface adopts it. + */ + it('carries a first message held offline over to the next new-chat surface', async () => { + stubUnreachableSend() + const first = renderUseChat() + await act(async () => { + await first.getResult().sendMessage('First message, sent offline') + }) + await waitFor(() => allQueuedMessages().some((message) => message.retryRequired === true)) + first.unmount() + + const second = renderUseChat() + await waitFor(() => + second + .getResult() + .messageQueue.some((message) => message.content === 'First message, sent offline') + ) + + network.online = true + await act(async () => { + window.dispatchEvent(new Event('online')) + }) + await waitFor(() => network.acceptedPosts === 1) + await act(async () => { + await sleep(300) + }) + + expect(network.acceptedPosts).toBe(1) + expect(new Set(state.postBodies.map((body) => body.userMessageId)).size).toBe(1) + expect(state.postBodies.at(-1)?.message).toBe('First message, sent offline') + }) + + /** + * Releasing held sends on mount must not bypass the queue's own rules: a + * follow-up queued behind a turn that is still running (here, restored after a + * reload) waits for that turn instead of being sent into a busy chat. + */ + it('keeps a queued follow-up waiting on mount while the chat is still running', async () => { + const history: MothershipChatHistory = { + ...idleHistory('chat-still-running'), + activeStreamId: 'turn-still-running', + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + 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 emptySseResponse() + } + if (url.includes('/api/mothership/chat/stream')) { + if (url.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() + .enqueue(history.id, { id: 'queued-before-reload', content: 'Queued before the reload' }) + renderUseChatInChat(history.id) + + await act(async () => { + await sleep(1000) + }) + + expect(state.postBodies).toHaveLength(0) + expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe( + 'Queued before the reload' + ) + }) + + /** + * The browser can come back online while the failing POST is still pending, + * so the release fires before the message is held. It must not then wait + * for a release that already happened. + */ + it('sends a message whose POST failed after the network had already returned', async () => { + const history = idleHistory('chat-online-mid-send') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + let failFirstPost: (() => 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 new Promise((_, reject) => { + failFirstPost = () => reject(new TypeError('Failed to fetch')) + }) + } + return emptySseResponse() + } + if (url.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 as the network came back') + }) + await waitFor(() => failFirstPost !== undefined) + + await act(async () => { + window.dispatchEvent(new Event('online')) + failFirstPost?.() + }) + await waitFor(() => state.postBodies.length === 2) + + expect(state.postBodies[1].message).toBe('Sent as the network came back') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + + /** Switching chats while a queued send is failing must not drop it from its own chat. */ + it('keeps a queued send that failed after the user switched chats', async () => { + const history = idleHistory('chat-left-mid-dispatch') + const other = idleHistory('chat-switched-to') + mockRequestJson.mockImplementation((_contract: AnyApiRouteContract, input: unknown) => + Promise.resolve({ + chat: JSON.stringify(input).includes(other.id) ? other : history, + }) + ) + let failPost: (() => 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((_, reject) => { + failPost = () => reject(new TypeError('Failed to fetch')) + }) + } + return fetchStub(input, init) + }) + useMothershipQueueStore + .getState() + .enqueue(history.id, { id: 'queued-then-left', content: 'Sent as I switched chats' }) + const { navigate } = renderUseChatInChat(history.id, history) + await waitFor(() => failPost !== undefined) + + navigate(other.id, other) + await act(async () => { + failPost?.() + await sleep(100) + }) + + const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] + expect(queued.map((message) => message.content)).toEqual(['Sent as I switched chats']) + expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId) + }) + + /** 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') + const other = idleHistory('chat-open-after-delete') + mockRequestJson.mockImplementation((_contract: AnyApiRouteContract, input: unknown) => + Promise.resolve({ + chat: JSON.stringify(input).includes(other.id) ? other : history, + }) + ) + let failPost: (() => 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((_, reject) => { + failPost = () => reject(new TypeError('Failed to fetch')) + }) + } + return fetchStub(input, init) + }) + useMothershipQueueStore + .getState() + .enqueue(history.id, { id: 'queued-then-deleted', content: 'In a chat I deleted' }) + const { navigate } = renderUseChatInChat(history.id, history) + await waitFor(() => failPost !== undefined) + + navigate(other.id, other) + useMothershipQueueStore.getState().clearChat(history.id) + await act(async () => { + failPost?.() + await sleep(100) + }) + + expect(useMothershipQueueStore.getState().queues[history.id]).toBeUndefined() + }) + + /** 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') + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + stubUnreachableSend() + const first = renderUseChatInChat(history.id, history) + await act(async () => { + await first.getResult().sendMessage('Held while I was elsewhere') + }) + await waitFor( + () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true + ) + first.unmount() + + network.online = true + renderUseChatInChat(history.id, history) + await waitFor(() => state.postBodies.length === 2) + + expect(state.postBodies[1].message).toBe('Held while I was elsewhere') + expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId) + }) + }) + /** * A withdrawn send belongs to the chat it was sent to. The cross-surface * lanes deliver to whatever chat is mounted next, so routing a chat-bound 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 41fe7b51c47..00cfd9841bc 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -205,13 +205,27 @@ interface FinalizeOptions { streamTerminal?: boolean } +/** + * 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` + * send is held in the queue until the browser is back online or the user sends + * it: dispatching it again at once would fail the same way. When the browser + * came back online while that send was failing (`networkReturned`), the release + * it would have waited for has already fired, so it is not held. + */ +interface WithdrawnSendResult { + userMessageId: string + unreachable?: boolean + networkReturned?: boolean +} + /** * `true` when the send owns the transcript (rendered, or handed to reconnect), * `false` when the caller should restore the queue entry, and the object form - * when an unmount cleanup withdrew it — `userMessageId` is what a retry reuses - * so the server deduplicates the two attempts. + * when the send was withdrawn: by an unmount cleanup, because it never reached + * the server, or because another turn held the chat. */ -type StartSendMessageResult = boolean | { userMessageId: string } +type StartSendMessageResult = boolean | WithdrawnSendResult interface StartSendMessageOptions { /** Awaited before dispatch. Defaults to the hook's in-flight stop, if any. */ @@ -718,6 +732,10 @@ export function useChat( const pendingStopModeRef = useRef(null) const workflowIdRef = useRef(options?.workflowId) workflowIdRef.current = options?.workflowId + /** Counts `online` events, so a send can tell the network returned while it was failing. */ + const onlineEventsRef = useRef(0) + /** Identifies this chatless surface across mounts, for the sends it holds. */ + const heldSendSurface = `${scopeKey}:${options?.workflowId ?? 'home'}` const onToolResultRef = useRef(options?.onToolResult) onToolResultRef.current = options?.onToolResult const onTitleUpdateRef = useRef(options?.onTitleUpdate) @@ -3409,6 +3427,7 @@ export function useChat( let consumedByTranscript = false let sendReachedServer = false + const onlineEventsAtSend = onlineEventsRef.current setError(null) setTransportStreaming() @@ -3823,10 +3842,10 @@ export function useChat( if (!response.ok) { const errorData = await response.json().catch(() => ({})) if (response.status === 409) { + /* A deduplicated send always names itself; a chat busy with another + turn names that turn, or nothing when its stream id is unreadable. */ const conflictStreamId = - typeof errorData.activeStreamId === 'string' - ? errorData.activeStreamId - : userMessageId + typeof errorData.activeStreamId === 'string' ? errorData.activeStreamId : undefined const supersededStreamId = queuedSendHandoff?.supersededStreamId ?? pendingStopStreamId if (supersededStreamId && conflictStreamId === supersededStreamId) { rollbackOptimisticSend() @@ -3840,6 +3859,35 @@ export function useChat( setError('Previous response is still shutting down; queued message was restored.') return false } + 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. */ + rollbackOptimisticSend() + if (streamGenRef.current === gen) { + streamGenRef.current++ + abortController.abort('send_conflict:chat_busy') + abortControllerRef.current = null + 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(() => {}) + } + return { userMessageId } + } /* 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 @@ -3959,6 +4007,31 @@ export function useChat( return consumedByTranscript } + if (!sendReachedServer) { + /* The POST got no answer, so nothing shows the server admitted it, and a + stream that was never opened reads as finished. Hand the message back + under the same id: if the server did admit it, the retry deduplicates + and the chat's own recovery shows the running turn. */ + rollbackOptimisticSend() + if (gen !== undefined && streamGenRef.current === gen) { + streamGenRef.current++ + admission?.controller.abort('send_unreachable:no_response') + abortControllerRef.current = null + clearActiveTurn() + setTransportIdle() + } + setError( + err instanceof TypeError + ? 'Message not sent: Sim could not be reached. It will send when you are back online.' + : getErrorMessage(err, 'Failed to send message') + ) + return { + userMessageId, + unreachable: true, + ...(onlineEventsRef.current !== onlineEventsAtSend ? { networkReturned: true } : {}), + } + } + const activeStreamId = streamIdRef.current if (activeStreamId && gen !== undefined && streamGenRef.current === gen) { const succeeded = await retryReconnect({ @@ -4116,11 +4189,12 @@ export function useChat( const result = await startSendMessage(message, fileAttachments, contexts, options) if (typeof result !== 'object') return - /* An unmount cleanup withdrew the send. A chat-bound key is the stable - chat id, so re-queueing under the key this was sent to is the durable - retry — and keeps the message in that chat rather than following the - user into whichever one they opened next. Only a chatless surface, - whose key dies with the mount, goes to the cross-surface lanes. */ + /* The send was withdrawn. A chat-bound key is the stable chat id, so + re-queueing under the key this was sent to is the durable retry, and + keeps the message in that chat rather than following the user into + whichever one they opened next. Only a send an unmount withdrew from a + chatless surface, whose key dies with the mount, goes to the + cross-surface lanes. */ const withdrawn = { content: message, fileAttachments, @@ -4132,30 +4206,34 @@ export function useChat( ? { assistantSearchLevel: options?.assistantSearchLevel } : {}), } - if (activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + if (!result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { handOffWithdrawnSend(withdrawn) return } - useMothershipQueueStore - .getState() - .enqueue( - activeChatKey, - createQueuedMessage( - message, - fileAttachments, - contexts, - result.userMessageId, - options?.requestMode, - options?.assistantSearch, - options?.assistantSearchLevel - ) - ) + useMothershipQueueStore.getState().enqueue(activeChatKey, { + ...createQueuedMessage( + message, + fileAttachments, + contexts, + result.userMessageId, + options?.requestMode, + options?.assistantSearch, + options?.assistantSearchLevel + ), + ...(result.unreachable && !result.networkReturned + ? { retryRequired: true, heldUntilOnline: true } + : {}), + ...(result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ? { heldSurface: heldSendSurface } + : {}), + }) }, [ workspaceId, createQueuedMessage, startSendMessage, handOffWithdrawnSend, + heldSendSurface, hasPendingChatAdmission, ] ) @@ -4702,9 +4780,10 @@ export function useChat( let dispatched = msg const restoreQueuedMessage = ( handoff?: QueuedSendHandoffSeed, - withdrawnUserMessageId?: string + withdrawn?: WithdrawnSendResult ) => { - const withdrawnByCleanup = withdrawnUserMessageId !== undefined + const withdrawnUserMessageId = withdrawn?.userMessageId + const retriesOnItsOwn = withdrawn !== undefined && !withdrawn.unreachable const savedHandoff = readQueuedSendHandoffState() const retainedHandoff = savedHandoff?.id === msg.id @@ -4720,7 +4799,9 @@ export function useChat( if (!removedFromQueue) { return } - if (options.epoch !== queueDispatchEpochRef.current && !withdrawnByCleanup) { + /* 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) { return } // If the user explicitly removed this message during dispatch, honor @@ -4733,7 +4814,7 @@ 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 (withdrawnByCleanup && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + if (withdrawn && retriesOnItsOwn && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { clearQueuedSendHandoffState(msg.id) handOffWithdrawnSend({ content: dispatched.content, @@ -4744,7 +4825,7 @@ export function useChat( ...(dispatched.assistantSearchLevel !== undefined ? { assistantSearchLevel: dispatched.assistantSearchLevel } : {}), - userMessageId: withdrawnUserMessageId, + userMessageId: withdrawn.userMessageId, }) return } @@ -4753,7 +4834,13 @@ export function useChat( useMothershipQueueStore.getState().insertAt(dispatchChatKey, originalIndex, { ...dispatched, ...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}), - retryRequired: !withdrawnByCleanup, + retryRequired: withdrawn?.unreachable ? !withdrawn.networkReturned : !retriesOnItsOwn, + ...(withdrawn?.unreachable && !withdrawn.networkReturned + ? { heldUntilOnline: true } + : {}), + ...(withdrawn?.unreachable && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + ? { heldSurface: heldSendSurface } + : {}), ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), }) } @@ -4797,7 +4884,7 @@ export function useChat( if (sendResult !== true) { restoreQueuedMessage( activeQueuedSendHandoff, - typeof sendResult === 'object' ? sendResult.userMessageId : undefined + typeof sendResult === 'object' ? sendResult : undefined ) } } catch { @@ -4808,7 +4895,7 @@ export function useChat( userRemovedDuringDispatch.delete(msg.id) } }, - [startSendMessage, handOffWithdrawnSend] + [startSendMessage, handOffWithdrawnSend, heldSendSurface] ) const runQueueDispatchLoop = useCallback(async () => { @@ -4962,6 +5049,29 @@ export function useChat( } }, []) + /** + * Sends held because the server could not be reached are released once the + * browser is online: on the `online` event, and on mount in case it fired while + * no chat surface was listening. The queue drain below sends a released head + * under its usual rules (history loaded, no running turn). A chatless surface + * first adopts what a dead mount of the same surface held, since that mount's + * queue key died with it. + */ + useEffect(() => { + if (typeof window === 'undefined') return + const releaseHeldSends = () => useMothershipQueueStore.getState().releaseHeldUntilOnline() + const handleOnline = () => { + onlineEventsRef.current++ + releaseHeldSends() + } + if (chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + useMothershipQueueStore.getState().adoptHeldSends(chatKey, heldSendSurface) + } + if (navigator.onLine) releaseHeldSends() + window.addEventListener('online', handleOnline) + return () => window.removeEventListener('online', handleOnline) + }, [chatKey, heldSendSurface]) + /** A recovered send already in history belongs to its accepted turn, even after Stop. */ useEffect(() => { if (!chatHistory || chatHistory.id !== chatKeyRef.current) return @@ -4985,9 +5095,10 @@ export function useChat( // `notifyTurnEnded`. Idempotent — the dispatch loop dedupes. const chatHistoryReady = chatHistory !== undefined const remoteActiveStreamId = chatHistory?.activeStreamId ?? null + const queueHeadHeld = messageQueue[0]?.retryRequired === true useEffect(() => { if (!scopeKey) return - if (messageQueue.length === 0) return + if (messageQueue.length === 0 || queueHeadHeld) return if (sendingRef.current || pendingStopPromiseRef.current) return if (queueDispatchTaskRef.current) return if (resolvedChatId && !chatHistoryReady) return @@ -4998,6 +5109,7 @@ export function useChat( organizationId, scopeKey, messageQueue.length, + queueHeadHeld, resolvedChatId, chatHistoryReady, remoteActiveStreamId, diff --git a/apps/sim/lib/mothership/tools/client/browser-tool-execution.test.ts b/apps/sim/lib/mothership/tools/client/browser-tool-execution.test.ts index 7c45358c7c8..432a9cb7c30 100644 --- a/apps/sim/lib/mothership/tools/client/browser-tool-execution.test.ts +++ b/apps/sim/lib/mothership/tools/client/browser-tool-execution.test.ts @@ -1470,7 +1470,7 @@ describe('pre-dispatch drops still resolve the waiter', () => { expect(mockReportCompletion).toHaveBeenCalledWith( 'stale-call-1', 'error', - expect.stringContaining('not run this time'), + expect.stringContaining('too late to run safely'), expect.objectContaining({ staleEvent: true }) ) const [, , message] = mockReportCompletion.mock.calls[0] ?? [] diff --git a/apps/sim/lib/mothership/tools/client/browser-tool-execution.ts b/apps/sim/lib/mothership/tools/client/browser-tool-execution.ts index 95c908ec0af..cc066474fd6 100644 --- a/apps/sim/lib/mothership/tools/client/browser-tool-execution.ts +++ b/apps/sim/lib/mothership/tools/client/browser-tool-execution.ts @@ -110,7 +110,7 @@ const OUTCOME_UNKNOWN_MESSAGE = const REPLAY_OUTCOME_UNKNOWN_MESSAGE = 'This browser action was recorded before the Sim page reloaded, but its terminal result could not be recovered. It may already have taken effect. Do not retry it automatically; take a fresh browser snapshot before deciding what to do.' const STALE_OBSERVATION_NOT_RUN_MESSAGE = - 'Not run: this browser observation reached the Sim desktop app too late to run safely, so it was not run this time and has no result. An observation changes nothing in the browser. Do not retry it in this turn; tell the user to keep this chat open in the Sim desktop app, or to ask again later.' + 'Not run: this browser observation reached the Sim desktop app too late to run safely, so it has no result. An observation changes nothing in the browser. Do not retry it in this turn; tell the user to keep this chat open in the Sim desktop app, or to ask again later.' const STALE_STATEFUL_OUTCOME_UNKNOWN_MESSAGE = 'This browser action was delivered too late to recover its exact result. It may already have taken effect. Do not retry it automatically; take a fresh browser snapshot before deciding what to do.' const REPLAY_GUARD_CAPACITY_MESSAGE = diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index b5a3588ea24..418ef9d8bd6 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -44,6 +44,7 @@ const sessionStorageAdapter = { const initialState = { queues: {} as Record, editing: {} as Record, + cleared: {} as Record, } const omitKey = (record: Record, key: string): Record => { @@ -67,6 +68,7 @@ export const useMothershipQueueStore = create()( enqueue: (chatKey, message) => set((state) => ({ + cleared: omitKey(state.cleared, chatKey), queues: setQueueForChat(state.queues, chatKey, [ ...(state.queues[chatKey] ?? []), message, @@ -75,6 +77,8 @@ export const useMothershipQueueStore = create()( insertAt: (chatKey, index, message) => set((state) => { + /** A restore that lands after its chat was cleared (deleted) must not recreate it. */ + if (state.cleared[chatKey]) return state const current = state.queues[chatKey] ?? [] if (current.some((m) => m.id === message.id)) return state const next = [...current] @@ -93,6 +97,8 @@ export const useMothershipQueueStore = create()( queuedSendHandoff, resumeUserMessageId: _staleResume, retryRequired: _retry, + heldUntilOnline: _held, + heldSurface: _surface, ...rest } = next[index] next[index] = { @@ -142,7 +148,11 @@ export const useMothershipQueueStore = create()( // Merge defensively in case a stale bucket survived in // sessionStorage. FIFO: existing first, then the resolved stream. const existing = state.queues[toKey] ?? [] - queues[toKey] = [...existing, ...fromQueue] + /** A chat-bound key is stable, so its messages no longer need a surface to adopt them. */ + queues[toKey] = [ + ...existing, + ...fromQueue.map(({ heldSurface: _surface, ...message }) => message), + ] } const editing = omitKey(state.editing, fromKey) if (fromEditing !== undefined) { @@ -151,10 +161,47 @@ export const useMothershipQueueStore = create()( return { queues, editing } }), + releaseHeldUntilOnline: () => + set((state) => { + let released = false + const queues: Record = {} + for (const [chatKey, queue] of Object.entries(state.queues)) { + queues[chatKey] = queue.map((message) => { + if (!message.heldUntilOnline) return message + released = true + const { retryRequired: _retry, heldUntilOnline: _held, ...rest } = message + return rest + }) + } + return released ? { queues } : state + }), + + adoptHeldSends: (toKey, surface) => + set((state) => { + const adopted: QueuedMothershipMessage[] = [] + let queues = state.queues + for (const [chatKey, queue] of Object.entries(state.queues)) { + if (chatKey === toKey) continue + const held = queue.filter((message) => message.heldSurface === surface) + if (held.length === 0) continue + adopted.push(...held) + queues = setQueueForChat( + queues, + chatKey, + queue.filter((message) => message.heldSurface !== surface) + ) + } + if (adopted.length === 0) return state + return { + queues: setQueueForChat(queues, toKey, [...(queues[toKey] ?? []), ...adopted]), + } + }), + clearChat: (chatKey) => set((state) => ({ queues: omitKey(state.queues, chatKey), editing: omitKey(state.editing, chatKey), + cleared: { ...state.cleared, [chatKey]: true }, })), reset: () => set(initialState), diff --git a/apps/sim/stores/mothership-queue/types.ts b/apps/sim/stores/mothership-queue/types.ts index 74a3fe553a4..af41dea707c 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -13,6 +13,17 @@ export type QueuedMothershipMessage = QueuedMessage & { queuedSendHandoff?: QueuedSendHandoffSeed /** A failed dispatch remains queued until the user retries or edits it. */ retryRequired?: boolean + /** + * The failed dispatch got no response, so the server may or may not have + * admitted it; the browser being online releases it for dispatch under the + * same id, which the server deduplicates if it did. + */ + heldUntilOnline?: boolean + /** + * Set on a send held by a chatless surface, whose queue key dies with its + * mount: the next chatless surface for the same owner and workflow adopts it. + */ + heldSurface?: string /** * Message id of a prior attempt at this send that an unmount cleanup * withdrew. Reused when the entry is dispatched so the server deduplicates @@ -37,6 +48,11 @@ export type QueuedMessageEditPatch = Pick< 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. + */ + cleared: Record enqueue: (chatKey: string, message: QueuedMothershipMessage) => void insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void @@ -44,6 +60,10 @@ export interface MothershipQueueState { remove: (chatKey: string, id: string) => void setEditing: (chatKey: string, id: string | null) => void migrate: (fromKey: string, toKey: string) => void + /** Releases every send held for the network for dispatch. */ + releaseHeldUntilOnline: () => void + /** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */ + adoptHeldSends: (toKey: string, surface: string) => void clearChat: (chatKey: string) => void reset: () => void }