From 197a9ce45bbcb60228a8d08e9ac0e643939526bd Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 02:32:16 -0700 Subject: [PATCH 1/2] refactor(mothership): one re-queue policy for a withdrawn send A withdrawn send's outcome is one reason (withdrawn, offline, unreachable, busy, stop-failed) instead of five exclusive booleans, and a pure requeuedFields(reason, previousAttempts, chatlessSurface) is the policy both re-queue sites apply. The restore strips the hold, retry and surface fields an earlier outcome left before applying it, so a stale heldUntilOnline can no longer let the browser coming online send a message waiting for the user. --- .../home/hooks/send-queue-policy.test.ts | 73 ++++++++++++++ .../home/hooks/send-queue-policy.ts | 98 +++++++++++++++++++ .../home/hooks/use-chat.dom.test.tsx | 38 +++++++ .../[workspaceId]/home/hooks/use-chat.ts | 95 ++++++------------ .../app/workspace/[workspaceId]/home/types.ts | 8 +- 5 files changed, 247 insertions(+), 65 deletions(-) create mode 100644 apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts create mode 100644 apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts 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 new file mode 100644 index 00000000000..50a9c2ace5b --- /dev/null +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts @@ -0,0 +1,73 @@ +import { describe, expect, it } from 'vitest' +import { + requeuedFields, + sendPayload, + withoutRequeueFields, +} from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy' + +describe('requeuedFields', () => { + it('holds an offline send for the network, on its chatless surface', () => { + expect(requeuedFields('offline', 0, 'ws-1:home')).toEqual({ + retryRequired: true, + heldUntilOnline: true, + heldSurface: 'ws-1:home', + }) + }) + + it.each(['unreachable', 'busy'] as const)('retries a %s send on a growing delay', (reason) => { + const before = Date.now() + const fields = requeuedFields(reason, 2, undefined) + + expect(fields.sendRetries).toBe(3) + expect(fields.notBefore).toBeGreaterThan(before) + expect(fields.retryRequired).toBeUndefined() + expect(fields.heldSurface).toBeUndefined() + }) + + it.each(['stop-failed', 'failed'] as const)( + 'leaves a %s send for the user, on any surface', + (reason) => { + expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({ retryRequired: true }) + } + ) + + it('sends a withdrawn message again as soon as the queue drains', () => { + expect(requeuedFields('withdrawn', 1, undefined)).toEqual({}) + }) +}) + +describe('withoutRequeueFields', () => { + it('drops every hold, retry and surface field and keeps the message itself', () => { + expect( + withoutRequeueFields({ + id: 'm1', + content: 'hello', + resumeUserMessageId: 'attempt-1', + admissionUnknown: true, + retryRequired: true, + heldUntilOnline: true, + sendRetries: 2, + notBefore: 123, + heldSurface: 'ws-1:home', + }) + ).toEqual({ + id: 'm1', + content: 'hello', + resumeUserMessageId: 'attempt-1', + admissionUnknown: true, + }) + }) +}) + +describe('sendPayload', () => { + it('keeps only the fields a send sets', () => { + expect( + sendPayload({ + content: 'hello', + fileAttachments: undefined, + requestMode: 'assistant', + assistantSearchLevel: undefined, + }) + ).toEqual({ content: 'hello', requestMode: 'assistant' }) + }) +}) 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 new file mode 100644 index 00000000000..455b5965ab0 --- /dev/null +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts @@ -0,0 +1,98 @@ +import { backoffWithJitter } from '@sim/utils/retry' +import type { SendPayload } from '@/app/workspace/[workspaceId]/home/types' +import type { QueuedMothershipMessage, ScheduledRetry } from '@/stores/mothership-queue/types' + +/** + * Why a send came back to its caller instead of going out: + * - `withdrawn`: an unmount cleanup withdrew it before the server answered; + * - `offline`: its POST got no answer and the browser is offline; + * - `unreachable`: its POST failed at the network while the browser reports + * itself online (a dropped connection), so no `online` event will come; + * - `busy`: the server refused it because another turn held the chat; + * - `stop-failed`: the Stop it waited on did not settle, so it was not sent. + */ +export type WithdrawalReason = 'withdrawn' | 'offline' | 'unreachable' | 'busy' | 'stop-failed' + +/** A withdrawal, or `failed`: the send failed outright and the user decides what next. */ +export type RequeueReason = WithdrawalReason | 'failed' + +/** The queue fields that say when, and on which surface, a re-queued message goes out. */ +type RequeueFields = Pick< + QueuedMothershipMessage, + 'retryRequired' | 'heldUntilOnline' | 'sendRetries' | 'notBefore' | 'heldSurface' +> + +const SEND_RETRY_BASE_MS = 1_000 +const SEND_RETRY_MAX_MS = 30_000 + +/** Queue fields for the `attempt`th automatic retry of a message: when it may be sent again. */ +export function sendRetry(attempt: number): ScheduledRetry { + return { + sendRetries: attempt, + notBefore: + Date.now() + + backoffWithJitter(attempt, null, { baseMs: SEND_RETRY_BASE_MS, maxMs: SEND_RETRY_MAX_MS }), + } +} + +/** + * The one re-queue policy, for a message going back to its queue: + * - `offline` waits for the browser to come back online (or the user); + * - `unreachable` and `busy` retry on a growing delay after `previousAttempts`; + * - `stop-failed` and `failed` wait for the user; + * - `withdrawn` goes out again as soon as the queue drains. + * + * A message held by a chatless surface carries that surface (`chatlessSurface`), + * whose queue key dies with its mount, so the next mount of it adopts the + * message. Only sends that wait on the network or the server are held that way. + */ +export function requeuedFields( + reason: RequeueReason, + previousAttempts: number, + chatlessSurface: string | undefined +): RequeueFields { + const surface = chatlessSurface ? { heldSurface: chatlessSurface } : {} + switch (reason) { + case 'offline': + return { retryRequired: true, heldUntilOnline: true, ...surface } + case 'unreachable': + case 'busy': + return { ...sendRetry(previousAttempts + 1), ...surface } + case 'stop-failed': + case 'failed': + return { retryRequired: true } + case 'withdrawn': + return {} + } +} + +/** + * A queue entry without the fields an earlier outcome set, so a re-queue applies + * only the policy for the outcome it is handling. A stale `heldUntilOnline`, for + * one, would let the browser coming online send a message waiting for the user. + */ +export function withoutRequeueFields(entry: QueuedMothershipMessage): QueuedMothershipMessage { + const { + retryRequired: _retryRequired, + heldUntilOnline: _heldUntilOnline, + sendRetries: _sendRetries, + notBefore: _notBefore, + heldSurface: _heldSurface, + ...rest + } = entry + return rest +} + +/** The payload of a send, without fields it does not set. */ +export function sendPayload(source: SendPayload): SendPayload { + return { + content: source.content, + ...(source.fileAttachments ? { fileAttachments: source.fileAttachments } : {}), + ...(source.contexts ? { contexts: source.contexts } : {}), + ...(source.requestMode ? { requestMode: source.requestMode } : {}), + ...(source.assistantSearch ? { assistantSearch: source.assistantSearch } : {}), + ...(source.assistantSearchLevel !== undefined + ? { assistantSearchLevel: source.assistantSearchLevel } + : {}), + } +} 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 1c2bcd4d5b2..8e35da8d46b 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 @@ -2341,6 +2341,44 @@ describe('useChat remount send recovery', () => { * settle sends nothing. That says nothing about the earlier attempt the * message resumes, so it must stay uneditable. */ + /** + * A message held for the network that the user then sends by hand, over a turn + * whose Stop does not settle, goes back waiting for the user. The hold it had + * before must not outlive that: the browser coming online must not send it. + */ + it('keeps a Send-now whose Stop failed waiting for the user when the browser comes online', async () => { + state.abortSettlements = [false, false, false, false] + const { getResult } = renderUseChatInChat('chat-a') + await act(async () => { + void getResult().sendMessage('Original request') + }) + await waitFor(() => state.postBodies.length === 1 && getResult().isSending) + useMothershipQueueStore.getState().enqueue('chat-a', { + id: 'held-offline', + content: 'written while offline', + resumeUserMessageId: 'offline-attempt', + admissionUnknown: true, + retryRequired: true, + heldUntilOnline: true, + }) + + await act(async () => { + await getResult() + .sendNow('held-offline') + .catch(() => {}) + await sleep(200) + }) + await act(async () => { + window.dispatchEvent(new Event('online')) + await sleep(100) + }) + + expect(state.postBodies).toHaveLength(1) + const queued = useMothershipQueueStore.getState().queues['chat-a']?.[0] + expect(queued).toMatchObject({ id: 'held-offline', retryRequired: true }) + expect(queued?.heldUntilOnline).toBeUndefined() + }) + it('keeps a resumed message uneditable when its Send-now Stop does not settle', async () => { state.abortSettlements = [false, false, false, false] const { getResult } = renderUseChatInChat('chat-a') 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 e659f6a55f9..6c95b1b42f9 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -176,6 +176,12 @@ import { writeQueuedSendHandoffClaim, writeQueuedSendHandoffState, } from './send-handoff' +import { + requeuedFields, + sendRetry, + type WithdrawalReason, + withoutRequeueFields, +} from './send-queue-policy' import { buildReplayStream, createStreamSchemaValidationError, @@ -222,21 +228,12 @@ type ActiveStreamRecoveryReason = /** * 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 failed before any response: dispatching it again at once would fail the - * same way. Its POST failing at the network layer while the browser still - * reports itself online (a dropped connection, `ERR_NETWORK_CHANGED`) is retried - * on a growing delay, since no `online` event will come. Otherwise it is held - * (`heldUntilOnline`) until the browser is back online or the user sends it. + * a retry reuses so the server deduplicates the two attempts; `reason` decides + * how it goes back to the queue (`requeuedFields`). */ interface WithdrawnSendResult { userMessageId: string - unreachable?: boolean - heldUntilOnline?: 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 + reason: WithdrawalReason /** Whether the server may hold `userMessageId`; see `admissionUnknown` in `startSendMessage`. */ admissionUnknown: boolean } @@ -345,8 +342,6 @@ 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 refused because the chat was busy, or that could not reach Sim. */ -const SEND_RETRY_BASE_MS = 1_000 -const SEND_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 @@ -732,16 +727,6 @@ function acceptedMessageIds(history: MothershipChatHistory): Set { return ids } -/** Queue fields for the `attempt`th automatic retry of a message: when it may be sent again. */ -function sendRetry(attempt: number): ScheduledRetry { - return { - sendRetries: attempt, - notBefore: - Date.now() + - backoffWithJitter(attempt, null, { baseMs: SEND_RETRY_BASE_MS, maxMs: SEND_RETRY_MAX_MS }), - } -} - export function useChat( owner: string | { organizationId: string }, initialChatId?: string, @@ -3910,7 +3895,7 @@ export function useChat( setError(getErrorMessage(err, 'Failed to stop the previous response')) /* 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, admissionUnknown } + return { userMessageId, reason: 'stop-failed', admissionUnknown } } } @@ -4037,7 +4022,7 @@ export function useChat( } if (viewOnSend) setError('Previous response is still shutting down; queued message was restored.') - return { userMessageId, held: true, admissionUnknown: false } + return { userMessageId, reason: 'stop-failed', admissionUnknown: false } } /** Withdraws this refused send so the queue retries it, under the same id, later. */ const releaseRefusedSend = () => { @@ -4076,7 +4061,7 @@ export function useChat( /* Only the chat lock refuses without naming this id. Admission's "superseded" conflict, where another attempt took this id's claim and may admit it, is answered as a duplicate naming this id instead. */ - return { userMessageId, busy: true, admissionUnknown: false } + return { userMessageId, reason: 'busy', admissionUnknown: false } } /* "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 @@ -4098,7 +4083,7 @@ export function useChat( ) if (!dedupedStreamExists) { releaseRefusedSend() - return { userMessageId, busy: true, admissionUnknown } + return { userMessageId, reason: 'busy', admissionUnknown } } /** The user may have moved on (another chat, another send) during the check. */ if (streamGenRef.current !== gen) return consumedByTranscript @@ -4209,7 +4194,7 @@ export function useChat( server deduplicates it against that turn instead of billing another one. */ rollbackOptimisticSend() - return { userMessageId, admissionUnknown } + return { userMessageId, reason: 'withdrawn', admissionUnknown } } return consumedByTranscript } @@ -4251,9 +4236,8 @@ export function useChat( ) return { userMessageId, + reason: retryLater ? 'unreachable' : 'offline', admissionUnknown, - unreachable: true, - ...(retryLater ? {} : { heldUntilOnline: true }), } } @@ -4434,12 +4418,8 @@ export function useChat( ? { assistantSearchLevel: options?.assistantSearchLevel } : {}), } - if ( - !result.unreachable && - !result.held && - !result.busy && - activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) - ) { + const chatless = activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + if (result.reason === 'withdrawn' && chatless) { handOffWithdrawnSend(withdrawn) return } @@ -4455,13 +4435,8 @@ export function useChat( options?.assistantSearch, options?.assistantSearchLevel ), - ...(result.heldUntilOnline ? { retryRequired: true, heldUntilOnline: true } : {}), - ...(result.held ? { retryRequired: true } : {}), - ...((result.unreachable && !result.heldUntilOnline) || result.busy ? sendRetry(1) : {}), + ...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined), admissionUnknown: result.admissionUnknown, - ...((result.unreachable || result.busy) && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) - ? { heldSurface: heldSendSurface } - : {}), }) }, [ @@ -5022,7 +4997,7 @@ export function useChat( withdrawn?: WithdrawnSendResult ) => { const withdrawnUserMessageId = withdrawn?.userMessageId - const retriesOnItsOwn = withdrawn !== undefined && !withdrawn.unreachable && !withdrawn.held + const chatless = dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) const savedHandoff = readQueuedSendHandoffState() const retainedHandoff = savedHandoff?.id === msg.id @@ -5043,7 +5018,10 @@ export function useChat( 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))) { + if ( + epochMoved && + (!withdrawn || (withdrawn.reason === 'stop-failed' && !surfaceMountedRef.current)) + ) { return } // If the user explicitly removed this message during dispatch, honor @@ -5056,12 +5034,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 ( - withdrawn && - retriesOnItsOwn && - !withdrawn.busy && - dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) - ) { + if (withdrawn?.reason === 'withdrawn' && chatless) { clearQueuedSendHandoffState(msg.id) handOffWithdrawnSend({ content: dispatched.content, @@ -5079,19 +5052,15 @@ export function useChat( /** Once restored, the queue owns recovery; a second handoff reader must not resend it. */ clearQueuedSendHandoffState(msg.id) useMothershipQueueStore.getState().insertAt(dispatchChatKey, originalIndex, { - ...dispatched, + /* Only this outcome's policy applies: what an earlier one set (a hold, a + retry delay, a surface) must not outlive it. */ + ...withoutRequeueFields(dispatched), ...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}), - retryRequired: withdrawn?.unreachable - ? withdrawn.heldUntilOnline === true - : !retriesOnItsOwn, - ...((withdrawn?.unreachable && !withdrawn.heldUntilOnline) || withdrawn?.busy - ? sendRetry((dispatched.sendRetries ?? 0) + 1) - : {}), - ...(withdrawn?.heldUntilOnline ? { heldUntilOnline: true } : {}), - ...((withdrawn?.unreachable || withdrawn?.busy) && - dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) - ? { heldSurface: heldSendSurface } - : {}), + ...requeuedFields( + withdrawn?.reason ?? 'failed', + dispatched.sendRetries ?? 0, + chatless ? heldSendSurface : undefined + ), ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), ...(withdrawn ? { admissionUnknown: withdrawn.admissionUnknown } : {}), }) diff --git a/apps/sim/app/workspace/[workspaceId]/home/types.ts b/apps/sim/app/workspace/[workspaceId]/home/types.ts index 7e0fbeed555..8f7c102a7e9 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/types.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/types.ts @@ -26,14 +26,18 @@ export interface FileAttachmentForApi { /** Assistant searches as the signed-in person and uses their connected accounts. */ export type ChatRequestMode = MothershipChat['mode'] -export interface QueuedMessage { - id: string +/** What a chat send carries, whether it goes out now, waits in the queue or is handed over. */ +export interface SendPayload { content: string fileAttachments?: FileAttachmentForApi[] contexts?: ChatContext[] requestMode?: ChatRequestMode assistantSearch?: WorkspaceSearchFilters assistantSearchLevel?: AssistantSearchLevel +} + +export interface QueuedMessage extends SendPayload { + id: string /** * An earlier attempt at this message got no answer, so the server may * already hold it as sent. It goes out exactly as written, under that From 2525c4b1b8869d3871a7fedb875989a76d936ea1 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 02:37:32 -0700 Subject: [PATCH 2/2] refactor(mothership): pass a chat send as one SendPayload A send's content, attachments, contexts, mode and search settings travel as one SendPayload instead of positional parameters (createQueuedMessage, sendMothershipMessage) and about ten hand-spread copies; sendPayload() is the one place that drops the fields a send does not set. --- .../app/workspace/[workspaceId]/home/home.tsx | 25 +-- .../[workspaceId]/home/hooks/use-chat.ts | 196 +++++------------- .../components/log-details/log-details.tsx | 2 +- .../components/terminal/terminal.tsx | 2 +- .../workflow-block/workflow-block.tsx | 9 +- .../components/search-modal/search-modal.tsx | 3 +- apps/sim/lib/mothership/events.ts | 22 +- apps/sim/stores/terminal/console/store.ts | 2 +- 8 files changed, 81 insertions(+), 180 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/home.tsx b/apps/sim/app/workspace/[workspaceId]/home/home.tsx index 165312717ba..0971951fd9a 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/home.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/home.tsx @@ -25,6 +25,7 @@ import { ChatResourcePanel } from '@/app/workspace/[workspaceId]/home/components import { RESOURCE_HEADER_CLASSES } from '@/app/workspace/[workspaceId]/home/components/mothership-view/components/resource-tabs/resource-tab-controls' import { SuggestedActions } from '@/app/workspace/[workspaceId]/home/components/suggested-actions' import { HomeFallback } from '@/app/workspace/[workspaceId]/home/home-fallback' +import { sendPayload } from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy' import { useChatResourcePanel, useResourcePanelController, @@ -242,13 +243,13 @@ function HomeContent({ chatId, userName, userId }: HomeProps) { if (!detail?.message) return e.preventDefault() prepareResourceViewForAgentTurn() - sendMessage(detail.message, detail.fileAttachments, detail.contexts, { + const { content, fileAttachments, contexts, ...sendOptions } = sendPayload({ + ...detail, + content: detail.message, + }) + sendMessage(content, fileAttachments, contexts, { + ...sendOptions, ...(detail.resumeUserMessageId ? { resumeUserMessageId: detail.resumeUserMessageId } : {}), - ...(detail.requestMode ? { requestMode: detail.requestMode } : {}), - ...(detail.assistantSearch ? { assistantSearch: detail.assistantSearch } : {}), - ...(detail.assistantSearchLevel !== undefined - ? { assistantSearchLevel: detail.assistantSearchLevel } - : {}), }) } window.addEventListener(MOTHERSHIP_SEND_MESSAGE_EVENT, handler) @@ -279,15 +280,15 @@ function HomeContent({ chatId, userName, userId }: HomeProps) { if (!handoff) return if (handoff.message) { prepareResourceViewForAgentTurn() - sendMessage(handoff.message, handoff.fileAttachments, handoff.contexts, { + const { content, fileAttachments, contexts, ...sendOptions } = sendPayload({ + ...handoff, + content: handoff.message, + }) + sendMessage(content, fileAttachments, contexts, { + ...sendOptions, ...(handoff.resumeUserMessageId ? { resumeUserMessageId: handoff.resumeUserMessageId } : {}), - ...(handoff.requestMode ? { requestMode: handoff.requestMode } : {}), - ...(handoff.assistantSearch ? { assistantSearch: handoff.assistantSearch } : {}), - ...(handoff.assistantSearchLevel !== undefined - ? { assistantSearchLevel: handoff.assistantSearchLevel } - : {}), }) return } 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 6c95b1b42f9..76635ab2bf0 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -138,7 +138,6 @@ import { reusedRequestId, useMothershipQueueStore } from '@/stores/mothership-qu import type { QueuedMothershipMessage, QueuedSendHandoffSeed, - ScheduledRetry, } from '@/stores/mothership-queue/types' import type { ChatContext } from '@/stores/panel' import { useTableViewPinStore } from '@/stores/table/view-pin/store' @@ -153,6 +152,7 @@ import type { MothershipResource, MothershipResourceType, QueuedMessage, + SendPayload, ToolCallInfo, } from '../types' import { @@ -178,6 +178,7 @@ import { } from './send-handoff' import { requeuedFields, + sendPayload, sendRetry, type WithdrawalReason, withoutRequeueFields, @@ -277,14 +278,8 @@ interface PendingChatAdmission { } /** A send an unmount cleanup withdrew, as handed to the next chat surface. */ -interface WithdrawnSend { - content: string - fileAttachments?: FileAttachmentForApi[] - contexts?: ChatContext[] +interface WithdrawnSend extends SendPayload { userMessageId: string - requestMode?: ChatRequestMode - assistantSearch?: WorkspaceSearchFilters - assistantSearchLevel?: AssistantSearchLevel } export interface UseChatReturn { @@ -3391,15 +3386,7 @@ export function useChat( ) const createQueuedMessage = useCallback( - ( - message: string, - fileAttachments?: FileAttachmentForApi[], - contexts?: ChatContext[], - resumeUserMessageId?: string, - requestMode?: ChatRequestMode, - assistantSearch?: WorkspaceSearchFilters, - assistantSearchLevel?: AssistantSearchLevel - ): QueuedMothershipMessage => { + (payload: SendPayload, resumeUserMessageId?: string): QueuedMothershipMessage => { const id = generateId() const handoffChatId = selectedChatIdRef.current ?? chatIdRef.current const cachedActiveStreamId = handoffChatId @@ -3415,13 +3402,8 @@ export function useChat( return { id, - content: message, - fileAttachments, - contexts, + ...sendPayload(payload), ...(resumeUserMessageId ? { resumeUserMessageId } : {}), - ...(requestMode ? { requestMode } : {}), - ...(assistantSearch ? { assistantSearch } : {}), - ...(assistantSearchLevel !== undefined ? { assistantSearchLevel } : {}), ...(supersededStreamId || handoffChatId ? { queuedSendHandoff: { @@ -3610,9 +3592,18 @@ export function useChat( if (!admittedThisSend || latestChoice !== effortChoice) saveMothershipChatEffort(queryClient, chatId, latestChoice) } + const payload = sendPayload({ + content: message, + fileAttachments, + contexts, + requestMode: options?.requestMode, + assistantSearch: options?.assistantSearch, + assistantSearchLevel: options?.assistantSearchLevel, + }) const writeQueuedSendHandoff = (chatId?: string) => { if (!queuedSendHandoff) return if (!chatId && !queuedSendHandoff.supersededStreamId) return + const { content, ...payloadFields } = payload writeQueuedSendHandoffState({ id: queuedSendHandoff.id, ...(chatId ? { chatId } : {}), @@ -3622,14 +3613,8 @@ export function useChat( ...(queuedSendHandoff.stopRequired ? { stopRequired: true } : {}), admissionUnknown, userMessageId, - message, - ...(fileAttachments ? { fileAttachments } : {}), - ...(contexts ? { contexts } : {}), - ...(options?.requestMode ? { requestMode: options.requestMode } : {}), - ...(options?.assistantSearch ? { assistantSearch: options.assistantSearch } : {}), - ...(options?.assistantSearchLevel !== undefined - ? { assistantSearchLevel: options?.assistantSearchLevel } - : {}), + message: content, + ...payloadFields, requestedAt: Date.now(), }) } @@ -3802,17 +3787,7 @@ 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 } - : {}), - }, + send: { ...payload, userMessageId }, } pendingChatAdmissionRef.current = admission } @@ -4297,31 +4272,11 @@ export function useChat( (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, - send.contexts, - send.fileAttachments, - send.userMessageId, - send.requestMode, - send.assistantSearch, - send.assistantSearchLevel - ) - ) { - return - } + const payload = sendPayload(send) + if (sendMothershipMessage(payload, send.userMessageId)) return + const { content, ...payloadFields } = payload MothershipHandoffStorage.store( - { - message: send.content, - ...(send.contexts?.length ? { contexts: send.contexts } : {}), - ...(send.fileAttachments?.length ? { fileAttachments: send.fileAttachments } : {}), - resumeUserMessageId: send.userMessageId, - ...(send.requestMode ? { requestMode: send.requestMode } : {}), - ...(send.assistantSearch ? { assistantSearch: send.assistantSearch } : {}), - ...(send.assistantSearchLevel !== undefined - ? { assistantSearchLevel: send.assistantSearchLevel } - : {}), - }, + { message: content, ...payloadFields, resumeUserMessageId: send.userMessageId }, organizationId ? { organizationId } : workspaceId! ) }, @@ -4365,6 +4320,14 @@ export function useChat( } options = { ...options, requestMode: options?.requestMode ?? requestModeRef.current } + const payload = sendPayload({ + content: message, + fileAttachments, + contexts, + requestMode: options.requestMode, + assistantSearch: options.assistantSearch, + assistantSearchLevel: options.assistantSearchLevel, + }) // An in-flight send drains the queue from `finalize`; a pending stop kicks // the dispatcher itself, since nothing else will once the stop settles. @@ -4380,18 +4343,7 @@ export function useChat( queuedAheadCount ) ) { - queueStore.enqueue( - activeChatKey, - createQueuedMessage( - message, - fileAttachments, - contexts, - options?.resumeUserMessageId, - options?.requestMode, - options?.assistantSearch, - options?.assistantSearchLevel - ) - ) + queueStore.enqueue(activeChatKey, createQueuedMessage(payload, options.resumeUserMessageId)) if (pendingStopPromiseRef.current || (queuedAheadCount > 0 && !sendingRef.current)) { void enqueueQueueDispatchRef.current({ type: 'send_head' }) } @@ -4407,34 +4359,17 @@ export function useChat( 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, - contexts, - userMessageId: result.userMessageId, - ...(options?.requestMode ? { requestMode: options.requestMode } : {}), - ...(options?.assistantSearch ? { assistantSearch: options.assistantSearch } : {}), - ...(options?.assistantSearchLevel !== undefined - ? { assistantSearchLevel: options?.assistantSearchLevel } - : {}), - } const chatless = activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) if (result.reason === 'withdrawn' && chatless) { - handOffWithdrawnSend(withdrawn) + handOffWithdrawnSend({ ...payload, userMessageId: result.userMessageId }) return } /* Back at the head: a direct send only goes out with nothing queued ahead of - it, so anything queued while its POST was out was written after it. */ + it, so anything queued while its POST was out was written after it. The one + exception is a held send adopted from a dead mount of this surface in that + window, which can be older; it lands behind this one. */ useMothershipQueueStore.getState().insertAt(activeChatKey, 0, { - ...createQueuedMessage( - message, - fileAttachments, - contexts, - result.userMessageId, - options?.requestMode, - options?.assistantSearch, - options?.assistantSearchLevel - ), + ...createQueuedMessage(payload, result.userMessageId), ...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined), admissionUnknown: result.admissionUnknown, }) @@ -4634,14 +4569,7 @@ export function useChat( /** Recovered sends join the queue so dispatch, failure and retry have one owner. */ useMothershipQueueStore.getState().insertAt(chatHistory.id, 0, { id: handoff.id, - content: handoff.message, - fileAttachments: handoff.fileAttachments, - contexts: handoff.contexts, - ...(handoff.requestMode ? { requestMode: handoff.requestMode } : {}), - ...(handoff.assistantSearch ? { assistantSearch: handoff.assistantSearch } : {}), - ...(handoff.assistantSearchLevel !== undefined - ? { assistantSearchLevel: handoff.assistantSearchLevel } - : {}), + ...sendPayload({ ...handoff, content: handoff.message }), queuedSendHandoff: { id: handoff.id, chatId: handoff.chatId, @@ -5037,14 +4965,7 @@ export function useChat( if (withdrawn?.reason === 'withdrawn' && chatless) { clearQueuedSendHandoffState(msg.id) handOffWithdrawnSend({ - content: dispatched.content, - fileAttachments: dispatched.fileAttachments, - contexts: dispatched.contexts, - ...(dispatched.requestMode ? { requestMode: dispatched.requestMode } : {}), - ...(dispatched.assistantSearch ? { assistantSearch: dispatched.assistantSearch } : {}), - ...(dispatched.assistantSearchLevel !== undefined - ? { assistantSearchLevel: dispatched.assistantSearchLevel } - : {}), + ...sendPayload(dispatched), userMessageId: withdrawn.userMessageId, }) return @@ -5083,27 +5004,19 @@ export function useChat( dispatched = liveMsg activeQueuedSendHandoff = options.queuedSendHandoff ?? liveMsg.queuedSendHandoff - const sendResult = await startSendMessage( - liveMsg.content, - liveMsg.fileAttachments, - liveMsg.contexts, - { - pendingStop: options.pendingStop, - onOptimisticSendApplied: removeQueuedMessage, - queuedSendHandoff: activeQueuedSendHandoff, - ...(liveMsg.resumeUserMessageId - ? { resumeUserMessageId: liveMsg.resumeUserMessageId } - : {}), - ...(liveMsg.admissionUnknown !== undefined - ? { admissionUnknown: liveMsg.admissionUnknown } - : {}), - ...(liveMsg.requestMode ? { requestMode: liveMsg.requestMode } : {}), - ...(liveMsg.assistantSearch ? { assistantSearch: liveMsg.assistantSearch } : {}), - ...(liveMsg.assistantSearchLevel !== undefined - ? { assistantSearchLevel: liveMsg.assistantSearchLevel } - : {}), - } - ) + const { content, fileAttachments, contexts, ...sendOptions } = sendPayload(liveMsg) + const sendResult = await startSendMessage(content, fileAttachments, contexts, { + ...sendOptions, + pendingStop: options.pendingStop, + onOptimisticSendApplied: removeQueuedMessage, + queuedSendHandoff: activeQueuedSendHandoff, + ...(liveMsg.resumeUserMessageId + ? { resumeUserMessageId: liveMsg.resumeUserMessageId } + : {}), + ...(liveMsg.admissionUnknown !== undefined + ? { admissionUnknown: liveMsg.admissionUnknown } + : {}), + }) if (sendResult !== true) { restoreQueuedMessage( @@ -5400,16 +5313,9 @@ export function useChat( const { send } = withdrawing queueStore.insertAt(deadKey, 0, { id: generateId(), - content: send.content, + ...sendPayload(send), resumeUserMessageId: send.userMessageId, admissionUnknown: true, - ...(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) diff --git a/apps/sim/app/workspace/[workspaceId]/logs/components/log-details/log-details.tsx b/apps/sim/app/workspace/[workspaceId]/logs/components/log-details/log-details.tsx index b3a28c53f71..2bc06ab2b9a 100644 --- a/apps/sim/app/workspace/[workspaceId]/logs/components/log-details/log-details.tsx +++ b/apps/sim/app/workspace/[workspaceId]/logs/components/log-details/log-details.tsx @@ -475,7 +475,7 @@ export function LogDetailsContent({ log, onActiveTabChange }: LogDetailsContentP const message = workflowName ? `The "${workflowName}" workflow run failed. Investigate the error in this run and help me fix it.` : 'This workflow run failed. Investigate the error in this run and help me fix it.' - if (sendMothershipMessage(message, [context])) return + if (sendMothershipMessage({ content: message, contexts: [context] })) return if (MothershipHandoffStorage.store({ message, contexts: [context] }, workspaceId)) { router.push(`/workspace/${workspaceId}/home`) } diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/terminal/terminal.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/terminal/terminal.tsx index 4a20bc7bc9b..076c48f474a 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/terminal/terminal.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/terminal/terminal.tsx @@ -905,7 +905,7 @@ export const Terminal = memo(function Terminal() { const errorMessage = entry.error ? String(entry.error) : 'Unknown error' const blockName = entry.blockName || 'Unknown Block' const message = `${errorMessage}\n\nError in ${blockName}.\n\nPlease fix this.` - sendMothershipMessage(message) + sendMothershipMessage({ content: message }) closeLogRowMenu() }, [closeLogRowMenu] diff --git a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/workflow-block.tsx b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/workflow-block.tsx index 1d7eea3beca..59411449d2d 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/workflow-block.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/[workflowId]/components/workflow-block/workflow-block.tsx @@ -834,9 +834,12 @@ export const WorkflowBlock = memo(function WorkflowBlock({ workflow_id: currentWorkflowId, kind: sunset.kind, }) - sendMothershipMessage(sunset.prompt, [ - { kind: 'workflow_block', workflowId: currentWorkflowId, blockId: id, label: name }, - ]) + sendMothershipMessage({ + content: sunset.prompt, + contexts: [ + { kind: 'workflow_block', workflowId: currentWorkflowId, blockId: id, label: name }, + ], + }) } const canonicalIndex = useMemo( diff --git a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/components/search-modal/search-modal.tsx b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/components/search-modal/search-modal.tsx index 85a47af1cf1..4ed71a6506f 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/components/search-modal/search-modal.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/components/search-modal/search-modal.tsx @@ -969,7 +969,8 @@ function SearchModalContent({ if (!query) return const homeHref = `/workspace/${workspaceId}/home` - const sentToMountedHome = window.location.pathname === homeHref && sendMothershipMessage(query) + const sentToMountedHome = + window.location.pathname === homeHref && sendMothershipMessage({ content: query }) if (!sentToMountedHome) { /* One-shot auto-send handoff: Home's mount consumer sends it on arrival, diff --git a/apps/sim/lib/mothership/events.ts b/apps/sim/lib/mothership/events.ts index 304545c6b94..08805684e4d 100644 --- a/apps/sim/lib/mothership/events.ts +++ b/apps/sim/lib/mothership/events.ts @@ -4,6 +4,7 @@ import type { AssistantSearchLevel } from '@/lib/mothership/generated/assistant' import type { ChatRequestMode, FileAttachmentForApi, + SendPayload, } from '@/app/workspace/[workspaceId]/home/types' import type { ChatContext } from '@/stores/panel' @@ -54,28 +55,17 @@ export interface MothershipSendMessageDetail { * was listening — callers that can fall back (e.g. cross-route navigation) use * this to decide whether to persist a handoff instead. */ -export function sendMothershipMessage( - message: string, - contexts?: ChatContext[], - fileAttachments?: FileAttachmentForApi[], - resumeUserMessageId?: string, - requestMode?: ChatRequestMode, - assistantSearch?: WorkspaceSearchFilters, - assistantSearchLevel?: AssistantSearchLevel -): boolean { - const trimmed = message.trim() - if (!trimmed && !fileAttachments?.length) { +export function sendMothershipMessage(payload: SendPayload, resumeUserMessageId?: string): boolean { + const { content, ...payloadFields } = payload + const trimmed = content.trim() + if (!trimmed && !payloadFields.fileAttachments?.length) { logger.warn('sendMothershipMessage called with empty message') return false } const consumed = dispatchClaimable(MOTHERSHIP_SEND_MESSAGE_EVENT, { message: trimmed, - contexts, - fileAttachments, + ...payloadFields, ...(resumeUserMessageId ? { resumeUserMessageId } : {}), - ...(requestMode ? { requestMode } : {}), - ...(assistantSearch ? { assistantSearch } : {}), - ...(assistantSearchLevel !== undefined ? { assistantSearchLevel } : {}), }) logger.info('Dispatched mothership message event', { messageLength: trimmed.length, consumed }) return consumed diff --git a/apps/sim/stores/terminal/console/store.ts b/apps/sim/stores/terminal/console/store.ts index 448f0a747ab..dfc6e7dee9c 100644 --- a/apps/sim/stores/terminal/console/store.ts +++ b/apps/sim/stores/terminal/console/store.ts @@ -319,7 +319,7 @@ const notifyBlockError = ({ action: getDeploymentShape().chatEnabled ? { label: 'Fix in Chat', - onClick: () => sendMothershipMessage(copilotMessage), + onClick: () => sendMothershipMessage({ content: copilotMessage }), } : undefined, })