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 795b1deeb9e..6c836ec81af 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 @@ -2151,11 +2151,6 @@ describe('useChat remount send recovery', () => { expect(allQueuedMessages()).toHaveLength(0) }) - /** - * After a Stop of the first message, only the follow-up was the user's - * intent: the Stop's POST is left to the server, nothing withdraws it, and the - * next mount sends just the follow-up, once. - */ /** * A first message held at the queue head after a remount may already be a * turn on the server. Editing it would send different text under a new id, @@ -2204,6 +2199,199 @@ describe('useChat remount send recovery', () => { }) }) + /** A chat with a turn running, so anything sent to it waits in its queue. */ + function renderBusyChat(id: string) { + const history: MothershipChatHistory = { + id, + mode: 'agent', + title: 'Busy', + messages: [], + activeStreamId: 'turn-still-running', + resources: [], + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).includes('/api/mothership/chat/stream')) { + if (String(input).includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + return { history, ...renderUseChatInChat(id, history) } + } + + /** + * A send another surface withdrew arrives here under its original id, through + * the send event or the stored handoff. Queued behind a running turn, it may + * already be a turn on the server, so it can't be edited either. + */ + it('does not let a withdrawn send handed to a busy chat be edited', async () => { + const { history, getResult } = renderBusyChat('chat-busy-on-handoff') + await waitFor(() => getResult().isSending) + await act(async () => { + await getResult().sendMessage('handed over from another surface', undefined, undefined, { + resumeUserMessageId: 'withdrawn-attempt', + }) + }) + const queued = useMothershipQueueStore.getState().queues[history.id]?.[0] + expect(queued?.resumeUserMessageId).toBe('withdrawn-attempt') + + let edited: ReturnType['editQueuedMessage']> + await act(async () => { + edited = getResult().editQueuedMessage(queued?.id ?? '') + }) + + expect(edited).toBeUndefined() + expect(getResult().editingQueuedId).toBeNull() + }) + + /** A follow-up whose dispatch got no answer may have reached the server too. */ + it('does not let a queued follow-up be edited after its send got no answer', async () => { + const history: MothershipChatHistory = { + id: 'chat-follow-up-unanswered', + mode: 'agent', + title: 'Unanswered', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) === '/api/mothership/chat' && init?.method === 'POST') { + state.postBodies.push(JSON.parse(String(init.body))) + throw new TypeError('Failed to fetch') + } + return fetchStub(input, init) + }) + useMothershipQueueStore + .getState() + .enqueue(history.id, { id: 'follow-up', content: 'and the second invoice' }) + const { getResult } = renderUseChatInChat(history.id, history) + await waitFor( + () => + useMothershipQueueStore.getState().queues[history.id]?.[0]?.resumeUserMessageId !== + undefined + ) + + let edited: ReturnType['editQueuedMessage']> + await act(async () => { + edited = getResult().editQueuedMessage('follow-up') + }) + + expect(edited).toBeUndefined() + expect(useMothershipQueueStore.getState().queues[history.id]?.[0]).toMatchObject({ + content: 'and the second invoice', + resumeUserMessageId: state.postBodies[0].userMessageId, + }) + }) + + /** + * Send-now on a resumed message whose Stop of the running turn does not + * settle sends nothing. That says nothing about the earlier attempt the + * message resumes, so it must stay uneditable. + */ + 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') + await act(async () => { + void getResult().sendMessage('Original request') + }) + await waitFor(() => state.postBodies.length === 1 && getResult().isSending) + await act(async () => { + await getResult().sendMessage('handed over from another surface', undefined, undefined, { + resumeUserMessageId: 'withdrawn-attempt', + }) + }) + await waitFor(() => useMothershipQueueStore.getState().queues['chat-a']?.length === 1) + + await act(async () => { + await getResult() + .sendNow() + .catch(() => {}) + await sleep(200) + }) + const queued = useMothershipQueueStore.getState().queues['chat-a']?.[0] + let edited: ReturnType['editQueuedMessage']> + await act(async () => { + edited = getResult().editQueuedMessage(queued?.id ?? '') + }) + + expect(state.postBodies).toHaveLength(1) + expect(queued).toMatchObject({ + content: 'handed over from another surface', + resumeUserMessageId: 'withdrawn-attempt', + admissionUnknown: true, + }) + expect(edited).toBeUndefined() + }) + + /** + * A held message the server then refuses as busy is known not to be a turn + * there: the server answers a retry of an admitted id as a duplicate, never + * as busy. The user can edit it again. + */ + it('lets a held message be edited again once the server refuses it as busy', async () => { + const history: MothershipChatHistory = { + id: 'chat-held-then-refused', + mode: 'agent', + title: 'Refused', + messages: [], + activeStreamId: null, + resources: [], + } + mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history })) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url === '/api/mothership/chat' && init?.method === 'POST') { + state.postBodies.push(JSON.parse(String(init.body))) + return Response.json( + { + error: 'A response is already in progress for this chat.', + activeStreamId: 'turn-from-another-tab', + }, + { status: 409 } + ) + } + if (url.includes('/api/mothership/chat/stream')) { + if (url.includes('batch=true')) { + return Response.json({ success: true, events: [], status: 'streaming' }) + } + return new Response(new ReadableStream(), { + headers: { 'Content-Type': 'text/event-stream' }, + }) + } + return fetchStub(input, init) + }) + useMothershipQueueStore.getState().enqueue(history.id, { + id: 'held-first', + content: 'inspect the workspace', + resumeUserMessageId: 'first-attempt', + admissionUnknown: true, + }) + const { getResult } = renderUseChatInChat(history.id, history) + await waitFor(() => state.postBodies.length === 1) + await waitFor( + () => useMothershipQueueStore.getState().queues[history.id]?.[0]?.id === 'held-first' + ) + + let edited: ReturnType['editQueuedMessage']> + await act(async () => { + edited = getResult().editQueuedMessage('held-first') + }) + + expect(edited?.content).toBe('inspect the workspace') + expect(getResult().editingQueuedId).toBe('held-first') + }) + + /** + * After a Stop of the first message, only the follow-up was the user's + * intent: the Stop's POST is left to the server, nothing withdraws it, and the + * next mount sends just the follow-up, once. + */ it('sends only the follow-up after a failed Stop when the new-chat surface remounts', async () => { stubFirstPostPendingThenAdmitted() const first = renderHomeLikeSurface() 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 28249857f0c..c11bd188ddf 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -234,6 +234,18 @@ interface WithdrawnSendResult { busy?: boolean /** Not sent at all (its Stop handoff failed); kept queued for the user to send. */ held?: boolean + /** + * The server refused this id outright (busy, or a predecessor still shutting + * down). It answers a retry of an admitted id as a duplicate instead, so the + * server is known not to have it, and its queue entry can be edited. + */ + notAdmitted?: boolean + /** + * This attempt never reached the server (its Stop did not settle). That says + * nothing about an earlier attempt the message resumes, whose uncertainty it + * keeps. + */ + neverSent?: boolean } /** @@ -3887,7 +3899,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 } + return { userMessageId, held: true, neverSent: true } } } @@ -4004,7 +4016,7 @@ export function useChat( } if (viewOnSend) setError('Previous response is still shutting down; queued message was restored.') - return { userMessageId, held: true } + return { userMessageId, held: true, notAdmitted: true } } /** Withdraws this refused send so the queue retries it, under the same id, later. */ const releaseRefusedSend = () => { @@ -4040,7 +4052,7 @@ export function useChat( exact: true, refetchType: 'none', }) - return { userMessageId, busy: true } + return { userMessageId, busy: true, notAdmitted: true } } /* "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 @@ -4414,6 +4426,11 @@ export function useChat( : {}), ...(result.held ? { retryRequired: true } : {}), ...(result.busy ? busyRetry(1) : {}), + admissionUnknown: result.notAdmitted + ? false + : result.neverSent + ? options?.resumeUserMessageId !== undefined + : true, ...((result.unreachable || result.busy) && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX) ? { heldSurface: heldSendSurface } : {}), @@ -5043,6 +5060,17 @@ export function useChat( ? { heldSurface: heldSendSurface } : {}), ...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}), + /* A refusal of this id settles it; an attempt that never left keeps the + earlier uncertainty; any other withdrawal may have reached the server. */ + ...(withdrawn + ? { + admissionUnknown: withdrawn.notAdmitted + ? false + : withdrawn.neverSent + ? dispatched.admissionUnknown === true + : true, + } + : {}), }) } diff --git a/apps/sim/app/workspace/[workspaceId]/home/types.ts b/apps/sim/app/workspace/[workspaceId]/home/types.ts index 25d249ce301..7e0fbeed555 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/types.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/types.ts @@ -35,9 +35,10 @@ export interface QueuedMessage { assistantSearch?: WorkspaceSearchFilters assistantSearchLevel?: AssistantSearchLevel /** - * A first message withdrawn before the server answered. The server may - * already hold it as sent, so it goes out exactly as written, under its - * original id, and cannot be edited into a different message. + * 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 + * attempt's id, and cannot be edited into a different message. False once + * the server is known not to have it (it refused it, or it was never sent). */ admissionUnknown?: boolean } diff --git a/apps/sim/stores/mothership-queue/store.dom.test.ts b/apps/sim/stores/mothership-queue/store.dom.test.ts new file mode 100644 index 00000000000..06f840d1d96 --- /dev/null +++ b/apps/sim/stores/mothership-queue/store.dom.test.ts @@ -0,0 +1,42 @@ +/** + * @vitest-environment jsdom + */ +import { beforeEach, describe, expect, it } from 'vitest' +import { useMothershipQueueStore } from '@/stores/mothership-queue/store' + +describe('useMothershipQueueStore rehydration', () => { + beforeEach(() => { + useMothershipQueueStore.getState().reset() + sessionStorage.clear() + }) + + it('treats a resumed message saved before the edit guard as possibly sent', async () => { + sessionStorage.setItem( + 'mothership-queue', + JSON.stringify({ + state: { + queues: { + 'chat-A': [ + { id: 'saved-before', content: 'original', resumeUserMessageId: 'attempt-1' }, + { + id: 'refused', + content: 'original', + resumeUserMessageId: 'attempt-2', + admissionUnknown: false, + }, + { id: 'plain', content: 'never sent' }, + ], + }, + }, + version: 0, + }) + ) + + await useMothershipQueueStore.persist.rehydrate() + + const [savedBefore, refused, plain] = useMothershipQueueStore.getState().queues['chat-A'] ?? [] + expect(savedBefore?.admissionUnknown).toBe(true) + expect(refused?.admissionUnknown).toBe(false) + expect(plain?.admissionUnknown).toBeUndefined() + }) +}) diff --git a/apps/sim/stores/mothership-queue/store.test.ts b/apps/sim/stores/mothership-queue/store.test.ts index 3b01b3232ed..dee3c41d3f3 100644 --- a/apps/sim/stores/mothership-queue/store.test.ts +++ b/apps/sim/stores/mothership-queue/store.test.ts @@ -30,6 +30,24 @@ describe('useMothershipQueueStore', () => { }) describe('replaceAt', () => { + it('treats any message resuming an earlier attempt as possibly sent, unless told otherwise', () => { + useMothershipQueueStore + .getState() + .enqueue('chat-A', { id: 'resumed', content: 'original', resumeUserMessageId: 'attempt-1' }) + useMothershipQueueStore.getState().insertAt('chat-A', 0, { + id: 'refused', + content: 'original', + resumeUserMessageId: 'attempt-2', + admissionUnknown: false, + }) + useMothershipQueueStore.getState().replaceAt('chat-A', 'resumed', { content: 'edited' }) + useMothershipQueueStore.getState().replaceAt('chat-A', 'refused', { content: 'edited' }) + + const [refused, resumed] = useMothershipQueueStore.getState().queues['chat-A'] ?? [] + expect(resumed).toMatchObject({ content: 'original', resumeUserMessageId: 'attempt-1' }) + expect(refused?.content).toBe('edited') + }) + it('leaves a first message the server may already hold unchanged, with its id', () => { useMothershipQueueStore.getState().enqueue('chat-A', { id: 'm1', @@ -50,6 +68,8 @@ describe('useMothershipQueueStore', () => { content: 'original', retryRequired: true, resumeUserMessageId: 'prior-request', + /** Its Stop never settled, so it was never sent: the server cannot hold it. */ + admissionUnknown: false, queuedSendHandoff: { id: 'm1', chatId: 'chat-A', diff --git a/apps/sim/stores/mothership-queue/store.ts b/apps/sim/stores/mothership-queue/store.ts index 9cfa62533a6..1c2427abd8f 100644 --- a/apps/sim/stores/mothership-queue/store.ts +++ b/apps/sim/stores/mothership-queue/store.ts @@ -1,5 +1,6 @@ import { createLogger } from '@sim/logger' 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' @@ -50,6 +51,38 @@ const initialState = { cleared: {} as Record, } +/** + * A message resuming an earlier attempt (`resumeUserMessageId`) may already be + * a turn on the server, unless the writer knows it is not + * (`admissionUnknown: false`). Every queue write goes through this, so no path + * can queue such a message as editable by leaving the flag out. + */ +function withAdmissionGuard(message: QueuedMothershipMessage): QueuedMothershipMessage { + if (message.resumeUserMessageId === undefined || message.admissionUnknown !== undefined) { + return message + } + return { ...message, admissionUnknown: true } +} + +function isQueuedMessage(value: unknown): value is QueuedMothershipMessage { + const record = toRecordOrNull(value) + return record !== null && typeof record.id === 'string' && typeof record.content === 'string' +} + +/** + * Queues saved to this tab's session, guarded on the way back in: an entry + * saved before `admissionUnknown` existed would otherwise be editable. + */ +function restoredQueues(persisted: unknown): Record { + const queues: Record = {} + for (const [chatKey, queue] of Object.entries(toRecord(toRecord(persisted).queues))) { + if (!Array.isArray(queue)) continue + const messages = queue.filter(isQueuedMessage).map(withAdmissionGuard) + if (messages.length > 0) queues[chatKey] = messages + } + return queues +} + const omitKey = (record: Record, key: string): Record => { if (!(key in record)) return record const { [key]: _removed, ...rest } = record @@ -75,7 +108,7 @@ export const useMothershipQueueStore = create()( return { queues: setQueueForChat(state.queues, chatKey, [ ...(state.queues[chatKey] ?? []), - message, + withAdmissionGuard(message), ]), } }), @@ -87,7 +120,7 @@ export const useMothershipQueueStore = create()( const current = state.queues[chatKey] ?? [] if (current.some((m) => m.id === message.id)) return state const next = [...current] - next.splice(Math.max(0, Math.min(index, next.length)), 0, message) + next.splice(Math.max(0, Math.min(index, next.length)), 0, withAdmissionGuard(message)) return { queues: setQueueForChat(state.queues, chatKey, next) } }), @@ -250,6 +283,7 @@ export const useMothershipQueueStore = create()( // edit text is component-local and empty after reload, so a persisted // editing flag would render an in-edit row with nothing bound. partialize: (state) => ({ queues: state.queues }), + merge: (persisted, current) => ({ ...current, queues: restoredQueues(persisted) }), } ), { name: 'mothership-queue-store' }