From 8ea86217bdb0c882b9c893619d42a79537a7243c Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 05:09:11 -0700 Subject: [PATCH 1/4] refactor(mothership): a pure resend verdict, one discard helper, and Send-now honours it - resendVerdict(entry, history | null) returns send, wait or drop, and the queue drain acts on it, instead of a boolean check with side effects - discardQueuedSend replaces the three copies of clear handoff, clear claim, remove - Send-now checks the history too and drops a message the server already accepted instead of resending it (G5); a history it cannot read does not hold back a send the user asked for - the own-id conflict test reaches that branch through a restored Send-now, since a message the history shows accepted is now dropped first --- .../home/hooks/send-queue-policy.test.ts | 43 +++++++++ .../home/hooks/send-queue-policy.ts | 39 ++++++++ .../home/hooks/use-chat.dom.test.tsx | 77 ++++++++++++---- .../[workspaceId]/home/hooks/use-chat.ts | 89 +++++++++++-------- 4 files changed, 193 insertions(+), 55 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts index 90498050a3a..9a5ff2b54bd 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts @@ -1,9 +1,11 @@ import { describe, expect, it } from 'vitest' import { requeuedFields, + resendVerdict, sendPayload, withoutRequeueFields, } from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy' +import type { MothershipChatHistory } from '@/hooks/queries/mothership-chats' describe('requeuedFields', () => { it('holds an offline send for the network, on its chatless surface', () => { @@ -72,3 +74,44 @@ describe('sendPayload', () => { ).toEqual({ content: 'hello', requestMode: 'assistant' }) }) }) + +describe('resendVerdict', () => { + const history = (accepted: string[] = [], activeStreamId: string | null = null) => + ({ + id: 'chat-A', + mode: 'agent', + title: 'A', + messages: accepted.map((id) => ({ + id, + role: 'user' as const, + content: 'sent', + timestamp: new Date(0).toISOString(), + })), + activeStreamId, + resources: [], + }) satisfies MothershipChatHistory + const resumed = { + id: 'm1', + content: 'hello', + resumeUserMessageId: 'attempt-1', + admissionUnknown: true, + } + + it('sends a message the server cannot already hold, without reading history', () => { + expect(resendVerdict({ id: 'm1', content: 'hello' }, null)).toBe('send') + expect(resendVerdict({ ...resumed, admissionUnknown: false }, null)).toBe('send') + }) + + it('drops a message the history shows accepted, as a message or as the running turn', () => { + expect(resendVerdict(resumed, history(['attempt-1']))).toBe('drop') + expect(resendVerdict(resumed, history([], 'attempt-1'))).toBe('drop') + }) + + it('sends a message the history does not show', () => { + expect(resendVerdict(resumed, history(['other']))).toBe('send') + }) + + it('waits when the history could not be read', () => { + expect(resendVerdict(resumed, null)).toBe('wait') + }) +}) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts index 7695e0487c1..55d03410bc9 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts @@ -1,5 +1,7 @@ import { backoffWithJitter } from '@sim/utils/retry' import type { SendPayload } from '@/app/workspace/[workspaceId]/home/types' +import type { MothershipChatHistory } from '@/hooks/queries/mothership-chats' +import { reusedRequestId } from '@/stores/mothership-queue/store' import type { QueuedMothershipMessage, SendRetry } from '@/stores/mothership-queue/types' /** @@ -87,3 +89,40 @@ export function sendPayload(source: SendPayload): SendPayload { : {}), } } + +/** Ids of the sends a chat's history shows the server accepted: its user messages and running turn. */ +export function acceptedMessageIds(history: MothershipChatHistory): Set { + const ids = new Set( + history.messages.filter((message) => message.role === 'user').map((message) => message.id) + ) + if (history.activeStreamId) ids.add(history.activeStreamId) + return ids +} + +/** + * Whether a queued message must be checked against its chat's history before it + * goes out: it may already be a turn on the server, under the id it reuses. + */ +export function needsResendCheck(entry: QueuedMothershipMessage): boolean { + return entry.admissionUnknown === true && reusedRequestId(entry) !== undefined +} + +/** What to do with a queued message about to go out. */ +export type ResendVerdict = 'send' | 'wait' | 'drop' + +/** + * Whether a queued message may go out, given its chat's history read fresh + * (`null` when the read failed). The server deduplicates a resend only while the + * earlier attempt's claim lasts, which a long outage outlives, so a message the + * history shows accepted is dropped (`drop`). One whose history cannot be read + * waits (`wait`): the read failing says nothing about whether it ran. + */ +export function resendVerdict( + entry: QueuedMothershipMessage, + history: MothershipChatHistory | null +): ResendVerdict { + const requestId = reusedRequestId(entry) + if (!needsResendCheck(entry) || requestId === undefined) return 'send' + if (!history) return 'wait' + return acceptedMessageIds(history).has(requestId) ? 'drop' : 'send' +} 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 cd223be7034..f737c361524 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 @@ -2671,10 +2671,10 @@ describe('useChat remount send recovery', () => { }) /** - * A resumed message whose earlier attempt is the very turn now running, sent - * now: its Stop and its resend share one id, and the server answers the resend - * as a duplicate of that turn. That is not a refusal, so the message must - * never come back editable. + * A Send-now restored after a reload, whose earlier attempt is the very turn it + * was stopping: its Stop and its resend share one id, and the server answers + * the resend as a duplicate of that turn. That is not a refusal, so the message + * must never come back editable. */ it('never treats a conflict naming the resent id itself as a refusal', async () => { const history: MothershipChatHistory = { @@ -2682,9 +2682,10 @@ describe('useChat remount send recovery', () => { mode: 'agent', title: 'Own id', messages: [], - activeStreamId: 'earlier-attempt', + activeStreamId: null, resources: [], } + /** Neither the cache nor a fresh read shows the earlier attempt yet. */ mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { const url = String(input) @@ -2709,25 +2710,25 @@ describe('useChat remount send recovery', () => { } return fetchStub(input, init) }) - const { getResult } = renderUseChatInChat(history.id, history) - await waitFor(() => getResult().isSending) - useMothershipQueueStore.getState().enqueue(history.id, { + writeQueuedSendHandoffState({ id: 'resumed', - content: 'sent earlier with no answer', - resumeUserMessageId: 'earlier-attempt', + chatId: history.id, + workspaceId: 'ws-1', + supersededStreamId: 'earlier-attempt', + userMessageId: 'earlier-attempt', + message: 'sent earlier with no answer', + stopRequired: true, admissionUnknown: true, + requestedAt: Date.now(), }) - await act(async () => { - /** Reattaches to the running turn, which stays open. */ - void getResult() - .sendNow('resumed') - .catch(() => {}) - }) + renderUseChatInChat(history.id, history) + await waitFor(() => state.postBodies.length === 1) await act(async () => { await sleep(300) }) + expect(state.abortBodies.map((body) => body.streamId)).toEqual(['earlier-attempt']) expect(state.postBodies.map((body) => body.userMessageId)).toEqual(['earlier-attempt']) const requeued = useMothershipQueueStore .getState() @@ -2881,6 +2882,50 @@ describe('useChat remount send recovery', () => { expect(getResult().messageQueue.map((message) => message.content)).toEqual(['Follow-up']) }) + /** + * Send-now on a message that may already be a turn checks the chat's history + * first, as the queue drain does: if the server shows the id accepted, the + * message is already in the chat and sending it again could run a second turn. + */ + it('drops a Send-now whose id the server already accepted instead of resending it', async () => { + const cached: MothershipChatHistory = { + id: 'chat-send-now-accepted', + mode: 'agent', + title: 'Accepted', + messages: [], + activeStreamId: null, + resources: [], + } + const fresh: MothershipChatHistory = { + ...cached, + messages: [ + { + id: 'earlier-attempt', + role: 'user', + content: 'sent earlier with no answer', + timestamp: new Date().toISOString(), + }, + ], + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: fresh })) + useMothershipQueueStore.getState().enqueue(cached.id, { + id: 'resumed', + content: 'sent earlier with no answer', + resumeUserMessageId: 'earlier-attempt', + admissionUnknown: true, + hold: 'user', + }) + const { getResult } = renderUseChatInChat(cached.id, cached) + + await act(async () => { + void getResult().sendNow('resumed') + await sleep(200) + }) + + expect(state.postBodies).toHaveLength(0) + expect(useMothershipQueueStore.getState().queues[cached.id]).toBeUndefined() + }) + 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 e46753b164e..995a85f18f5 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -182,7 +182,11 @@ import { writeQueuedSendHandoffState, } from './send-handoff' import { + acceptedMessageIds, + needsResendCheck, + type ResendVerdict, requeuedFields, + resendVerdict, sendPayload, sendRetry, type WithdrawalReason, @@ -718,13 +722,11 @@ export function getWorkflowCopilotUseChatOptions( } } -/** Ids of the sends a chat's history shows the server accepted: its user messages and running turn. */ -function acceptedMessageIds(history: MothershipChatHistory): Set { - const ids = new Set( - history.messages.filter((message) => message.role === 'user').map((message) => message.id) - ) - if (history.activeStreamId) ids.add(history.activeStreamId) - return ids +/** Removes a queued send and the handoff state and claim kept for it. */ +function discardQueuedSend(chatKey: string, id: string): void { + clearQueuedSendHandoffState(id) + clearQueuedSendHandoffClaim(id) + useMothershipQueueStore.getState().remove(chatKey, id) } export function useChat( @@ -5051,36 +5053,34 @@ export function useChat( ) /** - * Whether the queue drain must not send a message that may already have been - * admitted, checked against its chat's history (read fresh). The server - * deduplicates a resend only while the earlier attempt's claim lasts, which a - * long outage, online or off, outlives. An entry the history shows is dropped. - * One whose history cannot be read waits on a growing delay instead: the read - * failing says nothing about whether the attempt ran. + * The resend verdict for a queued message, reading its chat's history fresh + * only when the message may already be a turn there (`needsResendCheck`). */ - const mustNotResend = useCallback( - async (chatKey: string, msg: QueuedMothershipMessage): Promise => { - const requestId = reusedRequestId(msg) - if (!msg.admissionUnknown || !requestId || chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) - return false + const checkResend = useCallback( + async (chatKey: string, msg: QueuedMothershipMessage): Promise => { + if (!needsResendCheck(msg) || chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) return 'send' const history = await queryClient .fetchQuery({ ...mothershipChatHistoryQueryOptions(chatKey), staleTime: 0 }) - .catch(() => undefined) - if (history && !acceptedMessageIds(history).has(requestId)) return false + .catch(() => null) + return resendVerdict(msg, history) + }, + [queryClient] + ) + + /** Holds back a message the verdict says must wait, or discards one already sent. */ + const applyHeldResend = useCallback( + (chatKey: string, msg: QueuedMothershipMessage, verdict: 'wait' | 'drop') => { /** Sent by hand meanwhile: that dispatch owns the entry now. */ - if (queuedMessageDispatchIds.has(msg.id)) return true - if (!history) { - useMothershipQueueStore - .getState() - .deferRetry(chatKey, msg.id, sendRetry((msg.retry?.attempt ?? 0) + 1)) - return true + if (queuedMessageDispatchIds.has(msg.id)) return + if (verdict === 'drop') { + discardQueuedSend(chatKey, msg.id) + return } - clearQueuedSendHandoffState(msg.id) - clearQueuedSendHandoffClaim(msg.id) - useMothershipQueueStore.getState().remove(chatKey, msg.id) - return true + useMothershipQueueStore + .getState() + .deferRetry(chatKey, msg.id, sendRetry((msg.retry?.attempt ?? 0) + 1)) }, - [queryClient] + [] ) const runQueueDispatchLoop = useCallback(async () => { @@ -5107,7 +5107,11 @@ export function useChat( // 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 - if (await mustNotResend(activeChatKey, msg)) continue + const verdict = await checkResend(activeChatKey, msg) + if (verdict !== 'send') { + applyHeldResend(activeChatKey, msg, verdict) + continue + } await dispatchQueuedMessage(msg, { epoch: action.epoch }) } @@ -5123,7 +5127,7 @@ export function useChat( void queueDispatchLoopRef.current() } }) - }, [dispatchQueuedMessage, hasPendingChatAdmission, mustNotResend]) + }, [dispatchQueuedMessage, hasPendingChatAdmission, checkResend, applyHeldResend]) queueDispatchLoopRef.current = runQueueDispatchLoop const enqueueQueueDispatch = useCallback((action: QueueDispatchActionInput) => { @@ -5139,9 +5143,7 @@ export function useChat( if (queuedMessageDispatchIds.has(id)) { userRemovedDuringDispatch.add(id) } - clearQueuedSendHandoffState(id) - clearQueuedSendHandoffClaim(id) - useMothershipQueueStore.getState().remove(chatKeyRef.current, id) + discardQueuedSend(chatKeyRef.current, id) }, []) const sendQueuedMessageImmediately = useCallback( @@ -5152,6 +5154,15 @@ export function useChat( const msg = id === undefined ? queue?.[0] : queue?.find((queued) => queued.id === id) if (!msg || queueState.editing[chatKey] === msg.id) return if (queuedMessageDispatchIds.has(msg.id)) return + /* Sent by hand, it still must not run a second turn: if the history shows it + accepted, it is already in the chat. A history that cannot be read does not + hold it back here; the user asked to send it, and the server deduplicates it + while the earlier attempt's claim lasts. */ + if ((await checkResend(chatKey, msg)) === 'drop') { + applyHeldResend(chatKey, msg, 'drop') + return + } + if (queuedMessageDispatchIds.has(msg.id)) return const admissionPending = hasPendingChatAdmission() // Explicit queue sends should supersede any older auto-drain work scheduled by finalize(). @@ -5203,6 +5214,8 @@ export function useChat( organizationId, scopeKey, hasPendingChatAdmission, + checkResend, + applyHeldResend, ] ) @@ -5264,9 +5277,7 @@ export function useChat( if (queuedMessageDispatchIds.has(queued.id)) continue const requestId = reusedRequestId(queued) if (!requestId || !accepted.has(requestId)) continue - clearQueuedSendHandoffState(queued.id) - clearQueuedSendHandoffClaim(queued.id) - useMothershipQueueStore.getState().remove(chatHistory.id, queued.id) + discardQueuedSend(chatHistory.id, queued.id) } }, [chatHistory, messageQueue]) From d9d35894f67cf1d446e7ea393a4bba12737c72af Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 05:20:58 -0700 Subject: [PATCH 2/4] test(mothership): lock the queue upgrade path and tighten held waits - rehydrate rows for the production shape (retryRequired: false, no hold) and for a queue already in the new shape - the offline-hold waits check hold === 'online' instead of any hold - the one-lookup comment names where the one-hop invariant is enforced --- .../[workspaceId]/home/hooks/use-chat.dom.test.tsx | 6 +++--- apps/sim/stores/mothership-queue/store.dom.test.ts | 9 +++++++++ apps/sim/stores/mothership-queue/store.ts | 5 ++++- 3 files changed, 16 insertions(+), 4 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index f737c361524..4045fd20592 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 @@ -3480,7 +3480,7 @@ describe('useChat remount send recovery', () => { await waitFor(() => state.postBodies.length === 1) await waitFor( - () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined + () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold === 'online' ) const queued = useMothershipQueueStore.getState().queues[history.id] ?? [] @@ -3575,7 +3575,7 @@ describe('useChat remount send recovery', () => { await act(async () => { await first.getResult().sendMessage('First message, sent offline') }) - await waitFor(() => allQueuedMessages().some((message) => message.hold !== undefined)) + await waitFor(() => allQueuedMessages().some((message) => message.hold === 'online')) first.unmount() const second = renderUseChat() @@ -4352,7 +4352,7 @@ describe('useChat remount send recovery', () => { await first.getResult().sendMessage('Held while I was elsewhere') }) await waitFor( - () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined + () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold === 'online' ) first.unmount() diff --git a/apps/sim/stores/mothership-queue/store.dom.test.ts b/apps/sim/stores/mothership-queue/store.dom.test.ts index 3417a18bc7e..e6035ea900a 100644 --- a/apps/sim/stores/mothership-queue/store.dom.test.ts +++ b/apps/sim/stores/mothership-queue/store.dom.test.ts @@ -21,6 +21,13 @@ describe('useMothershipQueueStore rehydration', () => { { id: 'for-network', content: 'b', retryRequired: true, heldUntilOnline: true }, { id: 'retrying', content: 'c', sendRetries: 2, notBefore: 1_000 }, { id: 'plain', content: 'd' }, + { id: 'main-shape', content: 'e', retryRequired: false }, + { + id: 'new-shape', + content: 'f', + hold: 'online', + retry: { attempt: 1, notBefore: 5 }, + }, ], }, }, @@ -35,6 +42,8 @@ describe('useMothershipQueueStore rehydration', () => { { id: 'for-network', content: 'b', hold: 'online' }, { id: 'retrying', content: 'c', retry: { attempt: 2, notBefore: 1_000 } }, { id: 'plain', content: 'd' }, + { id: 'main-shape', content: 'e' }, + { id: 'new-shape', content: 'f', hold: 'online', retry: { attempt: 1, notBefore: 5 } }, ]) }) diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index 090e2106e4f..ea2785b213d 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -147,7 +147,10 @@ export function liveQueuePosition( aheadIds: readonly string[] ): { chatKey: string; index: number } { const { migratedTo, queues } = useMothershipQueueStore.getState() - /** One lookup: only a new-chat key moves, and only to its chat's key, which never does. */ + /* One lookup: only a new-chat key moves, and only to its chat's key, which never + does. \`useChat\` migrates only from its pending sentinel key to the resolved chat + id (the chat-resolution effect and the detached chat resolution), never from a + chat key. */ const migration = migratedTo[chatKey] const key = migration?.key ?? chatKey const ahead = new Set([...aheadIds, ...(migration?.ahead ?? [])]) From 2c5608bac4dab76153cee0cce0dbd30824a0173f Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 05:43:32 -0700 Subject: [PATCH 3/4] fix(mothership): don't stop the running turn for a Send-now gone during its history read Send-now reads the chat's history before stopping the running turn. A message removed, edited or dispatched in that time, or a view that moved on, still stopped the turn. It now re-reads the queue and checks the chat and mount are current before going on. --- .../home/hooks/use-chat.dom.test.tsx | 50 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 13 ++++- 2 files changed, 62 insertions(+), 1 deletion(-) 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 4045fd20592..9c45f0d750f 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 @@ -2887,6 +2887,56 @@ describe('useChat remount send recovery', () => { * first, as the queue drain does: if the server shows the id accepted, the * message is already in the chat and sending it again could run a second turn. */ + /** + * Send-now reads the history before it stops the running turn. A message the + * user removes during that read is no longer theirs to send, so the running + * turn must not be stopped for it. + */ + it('does not stop the running turn for a Send-now removed while its history is read', async () => { + const { getResult } = renderUseChatInChat('chat-a') + await act(async () => { + void getResult().sendMessage('Original request') + }) + await waitFor(() => state.postBodies.length === 1 && getResult().isSending) + let answerHistory: (() => void) | undefined + mockRequestJson.mockImplementation( + () => + new Promise((resolve) => { + answerHistory = () => + resolve({ + chat: { + id: 'chat-a', + mode: 'agent', + title: 'A', + messages: [], + activeStreamId: null, + resources: [], + }, + }) + }) + ) + useMothershipQueueStore.getState().enqueue('chat-a', { + id: 'resumed', + content: 'sent earlier with no answer', + resumeUserMessageId: 'earlier-attempt', + admissionUnknown: true, + }) + + await act(async () => { + void getResult().sendNow('resumed') + }) + await waitFor(() => answerHistory !== undefined) + await act(async () => { + getResult().removeFromQueue('resumed') + answerHistory?.() + await sleep(100) + }) + + expect(state.abortBodies).toHaveLength(0) + expect(state.postBodies).toHaveLength(1) + expect(getResult().isSending).toBe(true) + }) + it('drops a Send-now whose id the server already accepted instead of resending it', async () => { const cached: MothershipChatHistory = { id: 'chat-send-now-accepted', 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 995a85f18f5..968e35b8365 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -5162,7 +5162,18 @@ export function useChat( applyHeldResend(chatKey, msg, 'drop') return } - if (queuedMessageDispatchIds.has(msg.id)) return + /* The read took time. Only what is still this view's queued, unedited and + undispatched message may stop the running turn and go out. */ + const afterRead = useMothershipQueueStore.getState() + if ( + chatKeyRef.current !== chatKey || + !surfaceMountedRef.current || + queuedMessageDispatchIds.has(msg.id) || + afterRead.editing[chatKey] === msg.id || + !afterRead.queues[chatKey]?.some((queued) => queued.id === msg.id) + ) { + return + } const admissionPending = hasPendingChatAdmission() // Explicit queue sends should supersede any older auto-drain work scheduled by finalize(). From cb91ded87e6a6b8bed27e73d08191969e9639bf2 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 06:21:18 -0700 Subject: [PATCH 4/4] test(mothership): cover Send-now's re-checks after its history read Send-now no longer stops a turn in a chat the user moved to during the read, or from a surface that unmounted during it, and leaves a message the drain dispatched meanwhile to that dispatch. applyHeldResend is renamed applyResendVerdict. --- .../home/hooks/use-chat.dom.test.tsx | 155 +++++++++++++++++- .../[workspaceId]/home/hooks/use-chat.ts | 14 +- 2 files changed, 158 insertions(+), 11 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 9c45f0d750f..a65a61778c1 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 @@ -2882,11 +2882,6 @@ describe('useChat remount send recovery', () => { expect(getResult().messageQueue.map((message) => message.content)).toEqual(['Follow-up']) }) - /** - * Send-now on a message that may already be a turn checks the chat's history - * first, as the queue drain does: if the server shows the id accepted, the - * message is already in the chat and sending it again could run a second turn. - */ /** * Send-now reads the history before it stops the running turn. A message the * user removes during that read is no longer theirs to send, so the running @@ -2937,6 +2932,156 @@ describe('useChat remount send recovery', () => { expect(getResult().isSending).toBe(true) }) + /** + * The drain may already be reading the history for the message the user then + * sends by hand. When the drain's read lands first it dispatches the message, + * and Send-now must leave that dispatch alone rather than stop the turn it + * just started. + */ + it('does not stop the turn the drain started for a Send-now waiting on the same history read', async () => { + const history: MothershipChatHistory = { + id: 'chat-a', + mode: 'agent', + title: 'A', + messages: [], + activeStreamId: null, + resources: [], + } + const answers: Array<() => void> = [] + mockRequestJson.mockImplementation( + () => + new Promise((resolve) => { + answers.push(() => resolve({ chat: history })) + }) + ) + useMothershipQueueStore.getState().enqueue('chat-a', { + id: 'resumed', + content: 'sent earlier with no answer', + resumeUserMessageId: 'earlier-attempt', + admissionUnknown: true, + }) + const { getResult } = renderUseChatInChat('chat-a', history) + await waitFor(() => answers.length > 0) + + await act(async () => { + void getResult().sendNow('resumed') + }) + await act(async () => { + for (const answer of answers) answer() + await sleep(100) + }) + + expect(state.abortBodies).toHaveLength(0) + expect(state.postBodies).toHaveLength(1) + expect(state.postBodies[0]).toMatchObject({ userMessageId: 'earlier-attempt' }) + expect(getResult().isSending).toBe(true) + }) + + /** + * A Send-now belongs to the chat it was pressed in. If the user moves to + * another chat during its history read (which the move may cancel), the turn + * running there is not the one it was meant to stop. + */ + it("does not stop another chat's turn for a Send-now whose chat was left during its history read", async () => { + const chat = (id: string): MothershipChatHistory => ({ + id, + mode: 'agent', + title: id, + messages: [], + activeStreamId: null, + resources: [], + }) + let answerHistory: (() => void) | undefined + mockRequestJson.mockImplementation( + () => + new Promise((resolve) => { + answerHistory = () => resolve({ chat: chat('chat-a') }) + }) + ) + useMothershipQueueStore.getState().enqueue('chat-a', { + id: 'resumed', + content: 'sent earlier with no answer', + resumeUserMessageId: 'earlier-attempt', + admissionUnknown: true, + hold: 'user', + }) + const { getResult, navigate } = renderUseChatInChat('chat-a', chat('chat-a')) + await act(async () => { + void getResult().sendNow('resumed') + }) + await waitFor(() => answerHistory !== undefined) + const answerChatA = answerHistory + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: chat('chat-b') })) + + navigate('chat-b', chat('chat-b')) + await act(async () => { + void getResult().sendMessage('Chat B request') + }) + await waitFor(() => state.postBodies.length === 1) + await act(async () => { + answerChatA?.() + await sleep(100) + }) + + expect(state.abortBodies).toHaveLength(0) + expect(state.postBodies).toHaveLength(1) + expect(getResult().isSending).toBe(true) + expect(useMothershipQueueStore.getState().queues['chat-a']?.map((queued) => queued.id)).toEqual( + ['resumed'] + ) + }) + + /** A surface that unmounts during the read leaves the running turn to whoever owns it next. */ + it('does not stop the running turn for a Send-now whose surface unmounted during its history read', async () => { + const { getResult, unmount } = renderUseChatInChat('chat-a') + await act(async () => { + void getResult().sendMessage('Original request') + }) + await waitFor(() => state.postBodies.length === 1 && getResult().isSending) + let answerHistory: (() => void) | undefined + mockRequestJson.mockImplementation( + () => + new Promise((resolve) => { + answerHistory = () => + resolve({ + chat: { + id: 'chat-a', + mode: 'agent', + title: 'A', + messages: [], + activeStreamId: null, + resources: [], + }, + }) + }) + ) + useMothershipQueueStore.getState().enqueue('chat-a', { + id: 'resumed', + content: 'sent earlier with no answer', + resumeUserMessageId: 'earlier-attempt', + admissionUnknown: true, + }) + await act(async () => { + void getResult().sendNow('resumed') + }) + await waitFor(() => answerHistory !== undefined) + const abortsBeforeUnmount = state.abortBodies.length + + unmount() + await act(async () => { + answerHistory?.() + await sleep(100) + }) + + expect(state.abortBodies).toHaveLength(abortsBeforeUnmount) + expect(state.postBodies).toHaveLength(1) + }) + + /** + * Send-now on a message that may already be a turn checks the chat's history + * first, as the queue drain does: if the server shows the id accepted, the + * message is already in the chat and sending it again could run a second turn. + */ it('drops a Send-now whose id the server already accepted instead of resending it', async () => { const cached: MothershipChatHistory = { id: 'chat-send-now-accepted', 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 968e35b8365..a22c53070d5 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -5068,7 +5068,7 @@ export function useChat( ) /** Holds back a message the verdict says must wait, or discards one already sent. */ - const applyHeldResend = useCallback( + const applyResendVerdict = useCallback( (chatKey: string, msg: QueuedMothershipMessage, verdict: 'wait' | 'drop') => { /** Sent by hand meanwhile: that dispatch owns the entry now. */ if (queuedMessageDispatchIds.has(msg.id)) return @@ -5109,7 +5109,7 @@ export function useChat( if (queueState.editing[activeChatKey] === msg.id) continue const verdict = await checkResend(activeChatKey, msg) if (verdict !== 'send') { - applyHeldResend(activeChatKey, msg, verdict) + applyResendVerdict(activeChatKey, msg, verdict) continue } @@ -5127,7 +5127,7 @@ export function useChat( void queueDispatchLoopRef.current() } }) - }, [dispatchQueuedMessage, hasPendingChatAdmission, checkResend, applyHeldResend]) + }, [dispatchQueuedMessage, hasPendingChatAdmission, checkResend, applyResendVerdict]) queueDispatchLoopRef.current = runQueueDispatchLoop const enqueueQueueDispatch = useCallback((action: QueueDispatchActionInput) => { @@ -5159,11 +5159,13 @@ export function useChat( hold it back here; the user asked to send it, and the server deduplicates it while the earlier attempt's claim lasts. */ if ((await checkResend(chatKey, msg)) === 'drop') { - applyHeldResend(chatKey, msg, 'drop') + applyResendVerdict(chatKey, msg, 'drop') return } /* The read took time. Only what is still this view's queued, unedited and - undispatched message may stop the running turn and go out. */ + undispatched message may stop the running turn and go out. An entry that + needed the read cannot be in the editor (editing refuses one possibly + sent); the editing check covers an entry that skipped it. */ const afterRead = useMothershipQueueStore.getState() if ( chatKeyRef.current !== chatKey || @@ -5226,7 +5228,7 @@ export function useChat( scopeKey, hasPendingChatAdmission, checkResend, - applyHeldResend, + applyResendVerdict, ] )