From 0b920afe5d878fc1fc0b50f62d0c35c05a22cc17 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 06:09:38 -0700 Subject: [PATCH] refactor(mothership): keep a queued send's reused id in one field resumeUserMessageId is the only id a queue entry goes out under; the Stop handoff seed no longer carries a copy and reusedRequestId is gone. The stored handoff record keeps its userMessageId and is converted where it enters and leaves the queue. Queues saved with the id on the seed move it to resumeUserMessageId on rehydrate. --- .../home/hooks/send-queue-policy.ts | 5 +- .../home/hooks/use-chat.dom.test.tsx | 8 ++-- .../[workspaceId]/home/hooks/use-chat.ts | 16 ++++--- .../stores/mothership-queue/store.dom.test.ts | 46 +++++++++++++++++++ .../sim/stores/mothership-queue/store.test.ts | 10 ++-- apps/sim/stores/mothership-queue/store.ts | 23 ++++------ apps/sim/stores/mothership-queue/types.ts | 7 ++- 7 files changed, 80 insertions(+), 35 deletions(-) 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 55d03410bc9..10570935554 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,7 +1,6 @@ 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' /** @@ -104,7 +103,7 @@ export function acceptedMessageIds(history: MothershipChatHistory): Set * 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 + return entry.admissionUnknown === true && entry.resumeUserMessageId !== undefined } /** What to do with a queued message about to go out. */ @@ -121,7 +120,7 @@ export function resendVerdict( entry: QueuedMothershipMessage, history: MothershipChatHistory | null ): ResendVerdict { - const requestId = reusedRequestId(entry) + const requestId = entry.resumeUserMessageId 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 a65a61778c1..2035ec6d383 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 @@ -1811,7 +1811,7 @@ describe('useChat remount send recovery', () => { await waitFor(() => state.postBodies.length === 2) expect(state.postBodies[1]).toMatchObject({ message: queued.content, - userMessageId: failed.queuedSendHandoff?.userMessageId, + userMessageId: failed.resumeUserMessageId, }) expect(state.abortBodies).toHaveLength(4) expect(state.abortBodies[3]).toEqual(state.abortBodies[0]) @@ -1967,8 +1967,8 @@ describe('useChat remount send recovery', () => { expect.objectContaining({ id: 'queued-correction', hold: 'user', + resumeUserMessageId: 'prepared-correction-request', queuedSendHandoff: expect.objectContaining({ - userMessageId: 'prepared-correction-request', supersededStreamId: 'previous-response', stopRequired: true, }), @@ -1999,11 +1999,11 @@ describe('useChat remount send recovery', () => { id: 'earlier-correction', content: 'inspect the second invoice instead', hold: 'user', + resumeUserMessageId: 'prepared-correction', queuedSendHandoff: { id: 'earlier-correction', chatId: 'chat-a', supersededStreamId: 'earlier-response', - userMessageId: 'prepared-correction', stopRequired: true, }, }) @@ -2039,9 +2039,9 @@ describe('useChat remount send recovery', () => { expect(state.postBodies).toHaveLength(1) expect(allQueuedMessages()[0]).toMatchObject({ hold: 'user', + resumeUserMessageId: 'prepared-correction', queuedSendHandoff: { supersededStreamId: newerStreamId, - userMessageId: 'prepared-correction', stopRequired: 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 a22c53070d5..2e12636dccf 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -137,7 +137,6 @@ import { useMothershipEffortStore } from '@/stores/mothership-effort/store' import { liveQueueKey, liveQueuePosition, - reusedRequestId, useMothershipQueueStore, } from '@/stores/mothership-queue/store' import type { @@ -3551,7 +3550,7 @@ export function useChat( /* A retry of a withdrawn send reuses its id so the server deduplicates the two attempts; anything else mints a fresh one. */ - const reusedId = queuedSendHandoff?.userMessageId ?? options?.resumeUserMessageId + const reusedId = options?.resumeUserMessageId const userMessageId = reusedId ?? generateId() /* Whether the server may already hold `userMessageId`: the one fact that keeps a queued message from being edited into a second turn. A reused id may have been @@ -4583,9 +4582,10 @@ export function useChat( id: handoff.id, chatId: handoff.chatId, supersededStreamId: handoff.supersededStreamId, - userMessageId: handoff.userMessageId, ...(handoff.stopRequired ? { stopRequired: true } : {}), }, + /** The stored record's id is the one this entry goes out under. */ + resumeUserMessageId: handoff.userMessageId, ...(handoff.admissionUnknown !== undefined ? { admissionUnknown: handoff.admissionUnknown } : {}), @@ -4935,7 +4935,6 @@ export function useChat( handoff?: QueuedSendHandoffSeed, withdrawn?: WithdrawnSendResult ) => { - const withdrawnUserMessageId = withdrawn?.userMessageId /* 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( @@ -4944,13 +4943,16 @@ export function useChat( ) const chatless = restoreKey.startsWith(PENDING_CHAT_KEY_PREFIX) const savedHandoff = readQueuedSendHandoffState() + /** The id it went out under: the withdrawal's, else its stored handoff record's. */ + const restoredRequestId = + withdrawn?.userMessageId ?? + (savedHandoff?.id === msg.id ? savedHandoff.userMessageId : undefined) const retainedHandoff = savedHandoff?.id === msg.id ? { id: savedHandoff.id, chatId: savedHandoff.chatId, supersededStreamId: savedHandoff.supersededStreamId, - userMessageId: savedHandoff.userMessageId, stopRequired: savedHandoff.stopRequired, } : handoff @@ -4999,7 +5001,7 @@ export function useChat( dispatched.retry?.attempt ?? 0, chatless ? heldSendSurface : undefined ), - ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), + ...(restoredRequestId ? { resumeUserMessageId: restoredRequestId } : {}), ...(withdrawn ? { admissionUnknown: withdrawn.admissionUnknown } : {}), }) } @@ -5288,7 +5290,7 @@ export function useChat( const accepted = acceptedMessageIds(chatHistory) for (const queued of messageQueue) { if (queuedMessageDispatchIds.has(queued.id)) continue - const requestId = reusedRequestId(queued) + const requestId = queued.resumeUserMessageId if (!requestId || !accepted.has(requestId)) continue discardQueuedSend(chatHistory.id, queued.id) } diff --git a/apps/sim/stores/mothership-queue/store.dom.test.ts b/apps/sim/stores/mothership-queue/store.dom.test.ts index e6035ea900a..90e5d7b9c85 100644 --- a/apps/sim/stores/mothership-queue/store.dom.test.ts +++ b/apps/sim/stores/mothership-queue/store.dom.test.ts @@ -76,4 +76,50 @@ describe('useMothershipQueueStore rehydration', () => { expect(refused?.admissionUnknown).toBe(false) expect(plain?.admissionUnknown).toBeUndefined() }) + + it('moves the reused id a saved Stop handoff carried onto the entry', async () => { + const seed = { chatId: 'chat-A', supersededStreamId: 'previous-response', stopRequired: true } + sessionStorage.setItem( + 'mothership-queue', + JSON.stringify({ + state: { + queues: { + 'chat-A': [ + { + id: 'seed-only', + content: 'a', + queuedSendHandoff: { id: 'seed-only', ...seed, userMessageId: 'attempt-1' }, + }, + { + id: 'never-sent', + content: 'b', + admissionUnknown: false, + queuedSendHandoff: { id: 'never-sent', ...seed, userMessageId: 'attempt-2' }, + }, + ], + }, + }, + version: 0, + }) + ) + + await useMothershipQueueStore.persist.rehydrate() + + expect(useMothershipQueueStore.getState().queues['chat-A']).toEqual([ + { + id: 'seed-only', + content: 'a', + resumeUserMessageId: 'attempt-1', + admissionUnknown: true, + queuedSendHandoff: { id: 'seed-only', ...seed }, + }, + { + id: 'never-sent', + content: 'b', + resumeUserMessageId: 'attempt-2', + admissionUnknown: false, + queuedSendHandoff: { id: 'never-sent', ...seed }, + }, + ]) + }) }) diff --git a/apps/sim/stores/mothership-queue/store.test.ts b/apps/sim/stores/mothership-queue/store.test.ts index d3607ab9cda..53ad30f1d6a 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -39,11 +39,11 @@ describe('useMothershipQueueStore', () => { useMothershipQueueStore.getState().enqueue('chat-A', { id: 'legacy', content: 'original', + resumeUserMessageId: 'earlier-attempt', queuedSendHandoff: { id: 'legacy', chatId: 'chat-A', supersededStreamId: 'previous-response', - userMessageId: 'earlier-attempt', stopRequired: true, }, }) @@ -59,21 +59,21 @@ describe('useMothershipQueueStore', () => { useMothershipQueueStore.getState().enqueue('chat-A', { id: 'sent', content: 'original', + resumeUserMessageId: 'send-now-request', queuedSendHandoff: { id: 'sent', chatId: 'chat-A', supersededStreamId: 'previous-response', - userMessageId: 'send-now-request', }, }) useMothershipQueueStore.getState().enqueue('chat-A', { id: 'waiting', content: 'original', + resumeUserMessageId: 'not-sent-yet', queuedSendHandoff: { id: 'waiting', chatId: 'chat-A', supersededStreamId: 'previous-response', - userMessageId: 'not-sent-yet', stopRequired: true, }, /** A fresh id still waiting on its Stop, as the hook records it. */ @@ -85,7 +85,7 @@ describe('useMothershipQueueStore', () => { const [sent, waiting] = useMothershipQueueStore.getState().queues['chat-A'] ?? [] expect(sent).toMatchObject({ content: 'original', - queuedSendHandoff: { userMessageId: 'send-now-request' }, + resumeUserMessageId: 'send-now-request', }) expect(waiting?.content).toBe('edited') }) @@ -134,7 +134,6 @@ describe('useMothershipQueueStore', () => { id: 'm1', chatId: 'chat-A', supersededStreamId: 'previous-response', - userMessageId: 'prior-request', stopRequired: true, }, }) @@ -150,7 +149,6 @@ describe('useMothershipQueueStore', () => { stopRequired: true, }, }) - expect(edited?.queuedSendHandoff?.userMessageId).toBeUndefined() expect(edited?.resumeUserMessageId).toBeUndefined() expect(edited?.hold).toBeUndefined() }) diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index ea2785b213d..40d56c61dd1 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -56,15 +56,6 @@ const initialState = { migratedTo: {} as Record, } -/** - * The earlier attempt's id a queued message goes out under, if it reuses one: - * its Stop handoff's, else the withdrawn send's. `startSendMessage` picks the id - * in the same order. - */ -export function reusedRequestId(message: QueuedMothershipMessage): string | undefined { - return message.queuedSendHandoff?.userMessageId ?? message.resumeUserMessageId -} - /** * `admissionUnknown` is decided where a send chooses its id (`startSendMessage`) * and carried on the entry. A writer that has no say (a session saved before the @@ -73,7 +64,7 @@ export function reusedRequestId(message: QueuedMothershipMessage): string | unde * write goes through this, so no path can queue such a message as editable. */ function withAdmissionGuard(message: QueuedMothershipMessage): QueuedMothershipMessage { - if (reusedRequestId(message) === undefined || message.admissionUnknown !== undefined) { + if (message.resumeUserMessageId === undefined || message.admissionUnknown !== undefined) { return message } return { ...message, admissionUnknown: true } @@ -93,6 +84,14 @@ function withCurrentWaitFields(value: unknown): unknown { const record = toRecordOrNull(value) if (!record) return value const { retryRequired, heldUntilOnline, sendRetries, notBefore, ...rest } = record + /* An entry saved when its Stop handoff carried the reused id: that id is the + entry's `resumeUserMessageId` now. */ + const seed = toRecordOrNull(rest.queuedSendHandoff) + if (seed && typeof seed.userMessageId === 'string') { + const { userMessageId: seedRequestId, ...seedRest } = seed + rest.queuedSendHandoff = seedRest + if (rest.resumeUserMessageId === undefined) rest.resumeUserMessageId = seedRequestId + } return { ...rest, ...(retryRequired === true && rest.hold === undefined @@ -214,9 +213,7 @@ export const useMothershipQueueStore = create()( } = next[index] next[index] = { ...rest, - ...(queuedSendHandoff?.stopRequired - ? { queuedSendHandoff: { ...queuedSendHandoff, userMessageId: undefined } } - : {}), + ...(queuedSendHandoff?.stopRequired ? { queuedSendHandoff } : {}), content: patch.content, fileAttachments: patch.fileAttachments, contexts: patch.contexts, diff --git a/apps/sim/stores/mothership-queue/types.ts b/apps/sim/stores/mothership-queue/types.ts index b060de5911e..b0ad092668a 100644 --- a/apps/sim/stores/mothership-queue/types.ts +++ b/apps/sim/stores/mothership-queue/types.ts @@ -1,11 +1,14 @@ import type { QueuedMessage } from '@/app/workspace/[workspaceId]/home/types' -/** Durable predecessor and request identity for an outgoing message. */ +/** + * Durable predecessor of an outgoing message (the stream its Send-now stops). + * The id the message goes out under is the entry's `resumeUserMessageId`; only + * the stored handoff record (`QueuedSendHandoffState`) carries its own copy. + */ export interface QueuedSendHandoffSeed { id: string chatId?: string supersededStreamId: string | null - userMessageId?: string stopRequired?: boolean }