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 2cffbbc4628..1ab239790e9 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 @@ -19,7 +19,7 @@ * request it aborted was accepted. */ -import { act, type ReactNode, StrictMode, useEffect, useState } from 'react' +import { type ReactNode, act as reactAct, StrictMode, useEffect, useState } from 'react' import { authClientMock, authClientMockFns } from '@sim/testing/mocks/auth-client.mock' import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock' import { nextNavigationMock, nextNavigationMockFns } from '@sim/testing/mocks/next-navigation.mock' @@ -103,6 +103,36 @@ import { useExecutionStore } from '@/stores/execution/store' import { useMothershipEffortStore } from '@/stores/mothership-effort/store' import { useMothershipQueueStore } from '@/stores/mothership-queue/store' +/** Captured before any test fakes timers, so the act budget below runs in real time. */ +const realSetTimeout = globalThis.setTimeout +const realClearTimeout = globalThis.clearTimeout +/** Well under the 10s test timeout: the scope must end before the runner abandons the test. */ +const ACT_BUDGET_MS = 6_000 + +/** + * React's `act`, with async callbacks bounded. A callback that never settles (a + * regressed send stuck reconnecting) used to run into the test timeout while + * still inside React's act scope, and the renders of every later test queued + * behind it ("Hook result is not ready"). Failing the act after a budget ends + * the scope, so one regression is one red test. Sync callbacks stay synchronous. + */ +function act(callback: () => unknown): Promise { + return reactAct((): undefined | Promise => { + const result = callback() + if (!(result instanceof Promise)) return undefined + let budget: ReturnType | undefined + const budgetSpent = new Promise((_, reject) => { + budget = realSetTimeout( + () => reject(new Error(`act callback still pending after ${ACT_BUDGET_MS}ms`)), + ACT_BUDGET_MS + ) + }) + return Promise.race([result.then(() => undefined), budgetSpent]).finally(() => + realClearTimeout(budget) + ) + }) +} + authClientMockFns.mockUseSession.mockImplementation(() => ({ data: { user: { id: 'test-viewer' } }, })) @@ -2022,6 +2052,119 @@ describe('useChat remount send recovery', () => { expect(state.abortBodies[0]).not.toHaveProperty('chatId') }) + /** + * The first POST on the new-chat surface never answers; later POSTs open a + * turn in the chat the first message created. The abort endpoint and the + * stream lookup fail, as they would for a Stop that cannot reach the server. + */ + function stubFirstPostPendingThenAdmitted() { + 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) => { + init.signal?.addEventListener('abort', () => reject(init.signal?.reason), { + once: true, + }) + }) + } + return new Response( + new ReadableStream({ + start(controller) { + controller.close() + }, + }), + { + status: 200, + headers: { + 'Content-Type': 'text/event-stream', + 'x-mothership-chat-id': DEDUPED_CHAT_ID, + }, + } + ) + } + if (url.includes('/api/copilot/chat/abort')) { + state.abortBodies.push(JSON.parse(String(init?.body))) + return Response.json({ error: 'Internal error' }, { status: 500 }) + } + if (url.includes('/api/mothership/chat/stream') && state.postBodies.length === 1) { + return Response.json({ error: 'Internal error' }, { status: 500 }) + } + return fetchStub(input, init) + }) + } + + /** + * A follow-up typed on the new-chat surface while the first message waits for + * the server is queued behind it. If the surface remounts then, the first + * message is withdrawn and both must reach the next mount in the order they + * were written: the first message (under its own id), then the follow-up. + */ + it('sends a withdrawn first message before its follow-up when the new-chat surface remounts', async () => { + stubFirstPostPendingThenAdmitted() + const first = renderHomeLikeSurface() + await act(async () => { + void first.getResult().sendMessage('inspect the workspace') + }) + await waitFor(() => state.postBodies.length === 1) + await act(async () => { + void first.getResult().sendMessage('follow-up while admission pending') + }) + await waitFor(() => allQueuedMessages().length === 1) + first.unmount() + + const second = renderHomeLikeSurface() + await waitFor(() => state.postBodies.length >= 3, 4_000) + await act(async () => { + await sleep(300) + }) + + const afterRemount = state.postBodies.slice(1) + expect(afterRemount.map((body) => body.message)).toEqual([ + 'inspect the workspace', + 'follow-up while admission pending', + ]) + expect(afterRemount[0].userMessageId).toBe(state.postBodies[0].userMessageId) + expect(afterRemount[1].chatId).toBe(DEDUPED_CHAT_ID) + expect(second.claimedByOwnListener()).toBe(0) + expect(allQueuedMessages()).toHaveLength(0) + }) + + /** + * After a Stop of the first message, only the follow-up was the user's + * intent: the Stop's POST is left to the server, nothing withdraws it, and the + * next mount sends just the follow-up, once. + */ + it('sends only the follow-up after a failed Stop when the new-chat surface remounts', async () => { + stubFirstPostPendingThenAdmitted() + const first = renderHomeLikeSurface() + await act(async () => { + void first.getResult().sendMessage('inspect the workspace') + }) + await waitFor(() => state.postBodies.length === 1) + await act(async () => { + void first + .getResult() + .stopGeneration() + .catch(() => {}) + void first.getResult().sendMessage('Sent while the Stop was failing') + await sleep(1_000) + }) + first.unmount() + + renderHomeLikeSurface() + await waitFor(() => state.postBodies.length >= 2, 4_000) + await act(async () => { + await sleep(500) + }) + + expect(state.postBodies.slice(1).map((body) => body.message)).toEqual([ + 'Sent while the Stop was failing', + ]) + expect(allQueuedMessages()).toHaveLength(0) + }) + it('stopping a chat preserves an unrelated manual workflow execution', async () => { const executionStore = useExecutionStore.getState() executionStore.setIsExecuting('manual-workflow', true) 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 ff5db89ef70..bb405a336a9 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -267,6 +267,8 @@ interface PendingChatAdmission { chatKey: string controller: AbortController settled: Promise + /** What an unmount would withdraw, if it ran before the server answered. */ + send: WithdrawnSend } /** A send an unmount cleanup withdrew, as handed to the next chat surface. */ @@ -770,6 +772,10 @@ export function useChat( const onlineEventsRef = useRef(0) /** Identifies this chatless surface across mounts, for the sends it holds. */ const heldSendSurface = `${scopeKey}:${options?.workflowId ?? 'home'}` + const heldSendSurfaceRef = useRef(heldSendSurface) + heldSendSurfaceRef.current = heldSendSurface + /** Withdrawn first messages the unmount queued itself, so their handoff is skipped. */ + const withdrawnHeldAtUnmountRef = useRef | null>(null) const onToolResultRef = useRef(options?.onToolResult) onToolResultRef.current = options?.onToolResult const onTitleUpdateRef = useRef(options?.onTitleUpdate) @@ -3782,6 +3788,17 @@ export function useChat( settled: new Promise((resolve) => { resolveAdmission = resolve }), + send: { + content: message, + userMessageId, + ...(fileAttachments ? { fileAttachments } : {}), + ...(contexts ? { contexts } : {}), + ...(options?.requestMode ? { requestMode: options.requestMode } : {}), + ...(options?.assistantSearch ? { assistantSearch: options.assistantSearch } : {}), + ...(options?.assistantSearchLevel !== undefined + ? { assistantSearchLevel: options.assistantSearchLevel } + : {}), + }, } pendingChatAdmissionRef.current = admission } @@ -4025,9 +4042,10 @@ export function useChat( 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. */ + another reason proves nothing either way, so it is retried on purpose; that + includes the lookup's own timeout abort, which leaves this send's signal + untouched. The server deduplicates the retry by id. Only an abort of this + send itself (Stop, or the user moving on) is not retried. */ const dedupedStreamExists = await fetchStreamBatch( conflictStreamId, '0', @@ -4240,6 +4258,8 @@ export function useChat( */ const handOffWithdrawnSend = useCallback( (send: WithdrawnSend) => { + /** The unmount already queued it ahead of its follow-ups; see the unmount cleanup. */ + if (withdrawnHeldAtUnmountRef.current?.delete(send.userMessageId)) return if ( sendMothershipMessage( send.content, @@ -5302,6 +5322,40 @@ export function useChat( useEffect(() => { return () => { + /* A chatless mount's queue key dies with it, so messages still queued there + go to the next mount of this surface, as held sends do. A first message + this unmount withdraws (its POST not yet answered, and not stopped) goes + at their head: the follow-ups were written after it, and the next mount + would otherwise send them before its handoff arrives. Alone, it keeps + the usual cross-surface handoff. */ + const deadKey = chatKeyRef.current + if (deadKey.startsWith(PENDING_CHAT_KEY_PREFIX)) { + const queueStore = useMothershipQueueStore.getState() + const withdrawing = pendingChatAdmissionRef.current + if ( + withdrawing && + withdrawing.chatKey === deadKey && + abortControllerRef.current === withdrawing.controller && + (queueStore.queues[deadKey]?.length ?? 0) > 0 + ) { + const { send } = withdrawing + queueStore.insertAt(deadKey, 0, { + id: generateId(), + content: send.content, + resumeUserMessageId: send.userMessageId, + ...(send.fileAttachments ? { fileAttachments: send.fileAttachments } : {}), + ...(send.contexts ? { contexts: send.contexts } : {}), + ...(send.requestMode ? { requestMode: send.requestMode } : {}), + ...(send.assistantSearch ? { assistantSearch: send.assistantSearch } : {}), + ...(send.assistantSearchLevel !== undefined + ? { assistantSearchLevel: send.assistantSearchLevel } + : {}), + }) + withdrawnHeldAtUnmountRef.current ??= new Set() + withdrawnHeldAtUnmountRef.current.add(send.userMessageId) + } + queueStore.holdForSurface(deadKey, heldSendSurfaceRef.current) + } cancelActiveStreamRecovery() clearQueueDispatchState() streamGenRef.current++ diff --git a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts index e14c3394303..18bf0beccdc 100644 --- a/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts +++ b/apps/sim/hooks/queries/mothership-chat-history-restore.test.ts @@ -90,7 +90,7 @@ describe('lifting a chat delete', () => { useMothershipQueueStore.getState().clearChat(history.id) const answer = deferredAnswer({ chat: history }) const read = fetchMothershipChatHistory(history.id) - useMothershipQueueStore.getState().reopenChat(history.id) + useMothershipQueueStore.getState().reopenRestoredChat(history.id) useMothershipQueueStore.getState().clearChat(history.id) answer() await read @@ -110,7 +110,7 @@ describe('lifting a chat delete', () => { useMothershipQueueStore.getState().clearChat(history.id) await restore(history.id, () => { - useMothershipQueueStore.getState().reopenChat(history.id) + useMothershipQueueStore.getState().reopenRestoredChat(history.id) useMothershipQueueStore.getState().clearChat(history.id) }) diff --git a/apps/sim/hooks/queries/mothership-chats.ts b/apps/sim/hooks/queries/mothership-chats.ts index 5bcdc5966f4..b10e4fc5f96 100644 --- a/apps/sim/hooks/queries/mothership-chats.ts +++ b/apps/sim/hooks/queries/mothership-chats.ts @@ -321,7 +321,7 @@ export async function fetchMothershipChatHistory( ): Promise { const deleteSeen = useMothershipQueueStore.getState().cleared[chatId] const history = await readMothershipChatHistory(chatId, signal) - if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen) + if (deleteSeen !== undefined) useMothershipQueueStore.getState().liftDelete(chatId, deleteSeen) return history } @@ -385,7 +385,7 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) { onMutate: (chatId) => ({ deleteSeen: useMothershipQueueStore.getState().cleared[chatId] }), onSuccess: (_data, chatId, context) => { if (context?.deleteSeen === undefined) return - useMothershipQueueStore.getState().reopenChat(chatId, context.deleteSeen) + useMothershipQueueStore.getState().liftDelete(chatId, context.deleteSeen) }, onSettled: () => { queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) }) diff --git a/apps/sim/hooks/use-mothership-chat-events.ts b/apps/sim/hooks/use-mothership-chat-events.ts index 887eae2d3d5..71093a05341 100644 --- a/apps/sim/hooks/use-mothership-chat-events.ts +++ b/apps/sim/hooks/use-mothership-chat-events.ts @@ -125,7 +125,8 @@ export function handleMothershipChatStatusEvent( 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 === 'created') + useMothershipQueueStore.getState().reopenRestoredChat(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 145be699d68..7db72bf3f00 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -75,6 +75,21 @@ describe('useMothershipQueueStore', () => { }) }) + describe('holdForSurface', () => { + it("hands a dead chatless mount's queue to the next mount of its surface only", () => { + useMothershipQueueStore.getState().enqueue('pending::dead', message('m1')) + useMothershipQueueStore.getState().holdForSurface('pending::dead', 'ws-1:home') + + useMothershipQueueStore.getState().adoptHeldSends('pending::other', 'ws-1:workflow-1') + expect(useMothershipQueueStore.getState().queues['pending::other']).toBeUndefined() + + useMothershipQueueStore.getState().adoptHeldSends('pending::next', 'ws-1:home') + const state = useMothershipQueueStore.getState() + expect(state.queues['pending::next']?.map((m) => m.id)).toEqual(['m1']) + expect(state.queues['pending::dead']).toBeUndefined() + }) + }) + describe('migrate', () => { it('merges into an existing destination bucket instead of overwriting', () => { useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1')) @@ -92,15 +107,15 @@ describe('useMothershipQueueStore', () => { 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().reopenRestoredChat('chat-X') useMothershipQueueStore.getState().clearChat('chat-X') - useMothershipQueueStore.getState().reopenChat('chat-X', seen) + useMothershipQueueStore.getState().liftDelete('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().liftDelete('chat-X', latest) useMothershipQueueStore.getState().enqueue('chat-X', message('after-restore')) expect(useMothershipQueueStore.getState().queues['chat-X']?.map((m) => m.id)).toEqual([ 'after-restore', diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index c8003340501..50f477e16be 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -205,6 +205,19 @@ export const useMothershipQueueStore = create()( } }), + holdForSurface: (chatKey, surface) => + set((state) => { + const queue = state.queues[chatKey] + if (!queue?.some((message) => message.heldSurface !== surface)) return state + return { + queues: setQueueForChat( + state.queues, + chatKey, + queue.map((message) => ({ ...message, heldSurface: surface })) + ), + } + }), + clearChat: (chatKey) => set((state) => ({ queues: omitKey(state.queues, chatKey), @@ -212,13 +225,19 @@ export const useMothershipQueueStore = create()( 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) } - }), + liftDelete: (chatKey, deleteToken) => + set((state) => + state.cleared[chatKey] === deleteToken + ? { cleared: omitKey(state.cleared, chatKey) } + : state + ), + + reopenRestoredChat: (chatKey) => + set((state) => + state.cleared[chatKey] === undefined + ? state + : { 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 c5102343fed..33eb613476e 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -69,11 +69,23 @@ export interface MothershipQueueState { releaseHeldUntilOnline: () => void /** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */ adoptHeldSends: (toKey: string, surface: string) => void + /** + * Marks everything queued under a chatless mount's key as held for its + * `surface`, as that mount goes away: the key dies with it, and the next + * mount of the surface adopts the messages instead of losing them. + */ + holdForSurface: (chatKey: 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. + * Lifts the delete an operation saw when it began (`cleared[chatKey]` read + * then), for a read or restore that found the chat. A delete that landed + * after it has a newer token and stays. + */ + liftDelete: (chatKey: string, deleteToken: number) => void + /** + * Lifts any delete of a chat the server announced as restored (its `created` + * event, published after every delete before it). */ - reopenChat: (chatKey: string, deleteToken?: number) => void + reopenRestoredChat: (chatKey: string) => void reset: () => void }