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 50a9c2ace5b..5e12e319850 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 @@ -25,9 +25,13 @@ describe('requeuedFields', () => { }) it.each(['stop-failed', 'failed'] as const)( - 'leaves a %s send for the user, on any surface', + 'leaves a %s send for the user, adoptable by its chatless surface', (reason) => { - expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({ retryRequired: true }) + expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({ + retryRequired: true, + heldSurface: 'ws-1:home', + }) + expect(requeuedFields(reason, 4, undefined)).toEqual({ retryRequired: true }) } ) 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 455b5965ab0..b7974ffd637 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 @@ -42,9 +42,10 @@ export function sendRetry(attempt: number): ScheduledRetry { * - `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. + * A message re-queued on a chatless surface carries that surface + * (`chatlessSurface`), whose queue key dies with its mount, so the next mount of + * it adopts the message. That holds for every reason: a re-queue can land after + * the surface unmounted, when nothing else marks the dead key's queue. */ export function requeuedFields( reason: RequeueReason, @@ -60,7 +61,7 @@ export function requeuedFields( return { ...sendRetry(previousAttempts + 1), ...surface } case 'stop-failed': case 'failed': - return { retryRequired: true } + return { retryRequired: true, ...surface } case 'withdrawn': return {} } 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 8e35da8d46b..74041435f82 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 @@ -2829,6 +2829,57 @@ describe('useChat remount send recovery', () => { expect(allQueuedMessages()).toHaveLength(0) }) + /** + * Send-now on the new-chat surface stops the first message, which the Stop sees + * admitted into a chat: the surface moves to that chat and its queue moves with + * it. A busy refusal of the follow-up arriving after that must go back to the + * chat's queue, where the surface shows and retries it, not to the dead + * new-chat key it was dispatched from. + */ + it('re-queues a Send-now refused as busy in the chat the new-chat surface moved to', async () => { + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if ( + url === '/api/mothership/chat' && + init?.method === 'POST' && + state.postBodies.length > 0 + ) { + state.postBodies.push(JSON.parse(String(init.body))) + return Response.json( + { error: 'A response is already in progress for this chat.' }, + { status: 409 } + ) + } + return fetchStub(input, init) + }) + const { getResult } = renderUseChat() + await act(async () => { + void getResult().sendMessage('inspect the workspace') + }) + await waitFor(() => state.postBodies.length === 1) + await act(async () => { + void getResult().sendMessage('Follow-up') + }) + await waitFor(() => allQueuedMessages().length === 1) + + await act(async () => { + void getResult() + .sendNow() + .catch(() => {}) + }) + await waitFor(() => state.postBodies.length === 2) + await act(async () => { + await sleep(200) + }) + + const queues = useMothershipQueueStore.getState().queues + expect(queues[DEDUPED_CHAT_ID]?.map((message) => message.content)).toEqual(['Follow-up']) + expect( + Object.entries(queues).filter(([key, queue]) => key.startsWith('pending::') && queue.length) + ).toEqual([]) + expect(getResult().messageQueue.map((message) => message.content)).toEqual(['Follow-up']) + }) + 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 76635ab2bf0..6f85df3a3e6 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -134,7 +134,12 @@ import { workflowKeys } from '@/hooks/queries/workflows' import { snapAllSmoothText } from '@/hooks/use-smooth-text' import { useChatPanelStore } from '@/stores/chat-panel/store' import { useMothershipEffortStore } from '@/stores/mothership-effort/store' -import { reusedRequestId, useMothershipQueueStore } from '@/stores/mothership-queue/store' +import { + liveQueueKey, + liveQueuePosition, + reusedRequestId, + useMothershipQueueStore, +} from '@/stores/mothership-queue/store' import type { QueuedMothershipMessage, QueuedSendHandoffSeed, @@ -336,7 +341,7 @@ const PERSISTED_TURN_REFETCH_BASE_MS = 250 const PERSISTED_TURN_REFETCH_MAX_DELAY_MS = 5_000 /** How long a finished turn's save is waited for; a slow save still lands well inside it. */ const PERSISTED_TURN_WAIT_MS = 120_000 -/** Pacing for re-sending a message refused because the chat was busy, or that could not reach Sim. */ +/** How long a Stop's abort request may take before the Stop counts as failed. */ const STOP_REQUEST_TIMEOUT_MS = 15_000 const DETACHED_CHAT_RETRY_BASE_MS = 1000 const DETACHED_CHAT_RETRY_MAX_MS = 30_000 @@ -4359,7 +4364,9 @@ 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 chatless = activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + /** The new-chat queue may have moved to its chat while the POST was out. */ + const { chatKey: requeueKey, index: requeueIndex } = liveQueuePosition(activeChatKey, []) + const chatless = requeueKey.startsWith(PENDING_CHAT_KEY_PREFIX) if (result.reason === 'withdrawn' && chatless) { handOffWithdrawnSend({ ...payload, userMessageId: result.userMessageId }) return @@ -4368,7 +4375,7 @@ export function useChat( 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, { + useMothershipQueueStore.getState().insertAt(requeueKey, requeueIndex, { ...createQueuedMessage(payload, result.userMessageId), ...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined), admissionUnknown: result.admissionUnknown, @@ -4899,8 +4906,10 @@ export function useChat( const dispatchChatKey = chatKeyRef.current const queueAtStart = useMothershipQueueStore.getState().queues[dispatchChatKey] ?? EMPTY_MESSAGE_QUEUE - let originalIndex = queueAtStart.findIndex((queued) => queued.id === msg.id) - if (originalIndex === -1) { + const startIndex = queueAtStart.findIndex((queued) => queued.id === msg.id) + /** What was queued ahead of it, which it goes back behind if it is restored. */ + let aheadIds = queueAtStart.slice(0, Math.max(0, startIndex)).map((queued) => queued.id) + if (startIndex === -1) { queuedMessageDispatchIds.delete(msg.id) return } @@ -4913,7 +4922,7 @@ export function useChat( return } removedFromQueue = true - useMothershipQueueStore.getState().remove(dispatchChatKey, msg.id) + useMothershipQueueStore.getState().remove(liveQueueKey(dispatchChatKey), msg.id) } /* What actually went out. `msg` is the snapshot from when the dispatch was @@ -4925,7 +4934,13 @@ export function useChat( withdrawn?: WithdrawnSendResult ) => { const withdrawnUserMessageId = withdrawn?.userMessageId - const chatless = dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) + /* The send may have waited on a Stop that saw the new chat's first message + admitted, which moved this queue to that chat. */ + const { chatKey: restoreKey, index: restoreIndex } = liveQueuePosition( + dispatchChatKey, + aheadIds + ) + const chatless = restoreKey.startsWith(PENDING_CHAT_KEY_PREFIX) const savedHandoff = readQueuedSendHandoffState() const retainedHandoff = savedHandoff?.id === msg.id @@ -4972,7 +4987,7 @@ 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, { + useMothershipQueueStore.getState().insertAt(restoreKey, restoreIndex, { /* Only this outcome's policy applies: what an earlier one set (a hold, a retry delay, a surface) must not outlive it. */ ...withoutRequeueFields(dispatched), @@ -4996,7 +5011,7 @@ export function useChat( if (currentIndex === -1) { return } - originalIndex = currentIndex + aheadIds = queueAtSend.slice(0, currentIndex).map((queued) => queued.id) // Re-read live: the user may have applied an in-place edit (`replaceAt`) // between dispatch scheduling and this send. @@ -5057,12 +5072,12 @@ export function useChat( if (!history) { useMothershipQueueStore .getState() - .deferRetry(chatKey, msg.id, sendRetry((msg.sendRetries ?? 0) + 1)) + .deferRetry(liveQueueKey(chatKey), msg.id, sendRetry((msg.sendRetries ?? 0) + 1)) return true } clearQueuedSendHandoffState(msg.id) clearQueuedSendHandoffClaim(msg.id) - useMothershipQueueStore.getState().remove(chatKey, msg.id) + useMothershipQueueStore.getState().remove(liveQueueKey(chatKey), msg.id) return true }, [queryClient] diff --git a/apps/sim/lib/mothership/events.ts b/apps/sim/lib/mothership/events.ts index 08805684e4d..008b9e7f79e 100644 --- a/apps/sim/lib/mothership/events.ts +++ b/apps/sim/lib/mothership/events.ts @@ -1,6 +1,7 @@ import { createLogger } from '@sim/logger' import type { WorkspaceSearchFilters } from '@/lib/api/contracts/knowledge/search' import type { AssistantSearchLevel } from '@/lib/mothership/generated/assistant' +import { sendPayload } from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy' import type { ChatRequestMode, FileAttachmentForApi, @@ -56,7 +57,7 @@ export interface MothershipSendMessageDetail { * this to decide whether to persist a handoff instead. */ export function sendMothershipMessage(payload: SendPayload, resumeUserMessageId?: string): boolean { - const { content, ...payloadFields } = payload + const { content, ...payloadFields } = sendPayload(payload) const trimmed = content.trim() if (!trimmed && !payloadFields.fileAttachments?.length) { logger.warn('sendMothershipMessage called with empty message') diff --git a/apps/sim/stores/mothership-queue/store.test.ts b/apps/sim/stores/mothership-queue/store.test.ts index 94e16fc6d50..a63c138e71b 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -1,5 +1,9 @@ import { beforeEach, describe, expect, it } from 'vitest' -import { useMothershipQueueStore } from '@/stores/mothership-queue/store' +import { + liveQueueKey, + liveQueuePosition, + useMothershipQueueStore, +} from '@/stores/mothership-queue/store' import type { QueuedMothershipMessage } from '@/stores/mothership-queue/types' const message = (id: string, content = `content-${id}`): QueuedMothershipMessage => ({ @@ -181,6 +185,50 @@ describe('useMothershipQueueStore', () => { }) describe('migrate', () => { + it('points a late write at the chat a new-chat queue moved to, even an empty one', () => { + useMothershipQueueStore.getState().migrate('pending::empty', 'chat-X') + useMothershipQueueStore.getState().migrate('chat-X', 'chat-X') + + expect(liveQueueKey('pending::empty')).toBe('chat-X') + expect(liveQueueKey('pending::never-moved')).toBe('pending::never-moved') + expect(liveQueueKey('chat-X')).toBe('chat-X') + }) + + it('keeps a late write behind the messages the chat queue already held', () => { + useMothershipQueueStore.getState().enqueue('chat-Y', message('older-1')) + useMothershipQueueStore.getState().enqueue('chat-Y', message('older-2')) + useMothershipQueueStore.getState().enqueue('pending::moved', message('moved')) + useMothershipQueueStore.getState().migrate('pending::moved', 'chat-Y') + + const position = liveQueuePosition('pending::moved', []) + useMothershipQueueStore.getState().insertAt(position.chatKey, position.index, message('late')) + + expect(position).toEqual({ chatKey: 'chat-Y', index: 2 }) + expect(useMothershipQueueStore.getState().queues['chat-Y']?.map((m) => m.id)).toEqual([ + 'older-1', + 'older-2', + 'late', + 'moved', + ]) + }) + + it('keeps a late write in order when messages ahead of it were removed meanwhile', () => { + useMothershipQueueStore.getState().enqueue('chat-Z', message('older')) + useMothershipQueueStore.getState().enqueue('pending::sent', message('later')) + useMothershipQueueStore.getState().migrate('pending::sent', 'chat-Z') + useMothershipQueueStore.getState().remove('chat-Z', 'older') + + const position = liveQueuePosition('pending::sent', []) + useMothershipQueueStore + .getState() + .insertAt(position.chatKey, position.index, message('follow-up')) + + expect(useMothershipQueueStore.getState().queues['chat-Z']?.map((m) => m.id)).toEqual([ + 'follow-up', + 'later', + ]) + }) + it('merges into an existing destination bucket instead of overwriting', () => { useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1')) useMothershipQueueStore.getState().enqueue('chat-X', message('existing-2')) diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index 88c036ba526..4886d9d8141 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -3,7 +3,11 @@ import { toError } from '@sim/utils/errors' import { toRecord, toRecordOrNull } from '@sim/utils/object' import { create } from 'zustand' import { createJSONStorage, devtools, persist } from 'zustand/middleware' -import type { MothershipQueueState, QueuedMothershipMessage } from '@/stores/mothership-queue/types' +import type { + MothershipQueueState, + QueuedMothershipMessage, + QueueMigration, +} from '@/stores/mothership-queue/types' const logger = createLogger('MothershipQueueStore') @@ -49,6 +53,7 @@ const initialState = { queues: {} as Record, editing: {} as Record, cleared: {} as Record, + migratedTo: {} as Record, } /** @@ -106,6 +111,45 @@ const setQueueForChat = ( ): Record => next.length === 0 ? omitKey(queues, chatKey) : { ...queues, [chatKey]: next } +/** + * The queue key a write captured before an `await` should use now: the key a + * new-chat queue migrated to once its chat became known, if it did. + */ +export function liveQueueKey(chatKey: string): string { + const { migratedTo } = useMothershipQueueStore.getState() + let key = chatKey + for (let hops = 0; hops < 8 && migratedTo[key] !== undefined; hops++) key = migratedTo[key].key + return key +} + +/** + * Where a message goes back into its queue after a write captured before an + * `await`: in the queue's live key, right after the last message still there + * that was ahead of it (`aheadIds`, plus whatever a chat's queue already held + * when a new-chat queue moved into it), else at the head. Anchoring on ids, not + * an index, keeps it in order however the queue changed meanwhile. + */ +export function liveQueuePosition( + chatKey: string, + aheadIds: readonly string[] +): { chatKey: string; index: number } { + const { migratedTo, queues } = useMothershipQueueStore.getState() + const ahead = new Set(aheadIds) + let key = chatKey + for (let hops = 0; hops < 8; hops++) { + const migration = migratedTo[key] + if (!migration) break + for (const id of migration.ahead) ahead.add(id) + key = migration.key + } + const queue = queues[key] ?? [] + let index = 0 + queue.forEach((message, position) => { + if (ahead.has(message.id)) index = position + 1 + }) + return { chatKey: key, index } +} + export const useMothershipQueueStore = create()( devtools( persist( @@ -191,9 +235,16 @@ export const useMothershipQueueStore = create()( migrate: (fromKey, toKey) => set((state) => { if (fromKey === toKey) return state + const migratedTo = { + ...state.migratedTo, + [fromKey]: { + key: toKey, + ahead: (state.queues[toKey] ?? []).map((message) => message.id), + }, + } const fromQueue = state.queues[fromKey] const fromEditing = state.editing[fromKey] - if (!fromQueue && fromEditing === undefined) return state + if (!fromQueue && fromEditing === undefined) return { migratedTo } const queues = omitKey(state.queues, fromKey) /** A chat deleted meanwhile takes nothing: its queue is gone with it. */ @@ -211,7 +262,7 @@ export const useMothershipQueueStore = create()( if (fromEditing !== undefined) { editing[toKey] = fromEditing } - return { queues, editing } + return { queues, editing, migratedTo } }), releaseHeldUntilOnline: () => diff --git a/apps/sim/stores/mothership-queue/types.ts b/apps/sim/stores/mothership-queue/types.ts index 04d1e827be5..06bc775f162 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -56,6 +56,13 @@ export type QueuedMessageEditPatch = Pick< | 'assistantSearchLevel' > +/** A new-chat queue's move to its chat's key. */ +export interface QueueMigration { + key: string + /** Ids of the messages the chat's queue already held, which stay ahead of the moved ones. */ + ahead: string[] +} + export interface MothershipQueueState { queues: Record editing: Record @@ -65,6 +72,13 @@ export interface MothershipQueueState { * handed back); restoring the chat lifts it. */ cleared: Record + /** + * Where each new-chat key's queue moved when its chat became known + * (`migrate`). A write that captured the old key before an `await` follows + * this (`liveQueueKey`, `liveQueuePosition`), so it lands in the chat's queue, + * not a dead key. + */ + migratedTo: Record enqueue: (chatKey: string, message: QueuedMothershipMessage) => void insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void