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 }