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 affe8a1c15c..679dcdf97f9 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 @@ -2138,6 +2138,125 @@ describe('useChat remount send recovery', () => { ?.messages.map((message) => message.id) ).toEqual(['saved-user', 'saved-assistant']) }) + + /** + * The tab finalizes on the `complete` event, which reaches it before the server + * saves the turn. A transcript read in that gap is the server's in-flight copy; + * the tab must read again rather than keep it (live ids, the finished stream + * still listed as running) until something else happens to refetch. That holds + * when the save is slow, and when a follow-up is queued but not yet sent. + */ + it.each([ + { label: 'right after the first read', unsavedReads: 1, heldFollowUp: false }, + { label: 'only after a slow save', unsavedReads: 9, heldFollowUp: false }, + { label: 'with a follow-up queued but held', unsavedReads: 1, heldFollowUp: true }, + ])( + 're-reads a transcript fetched before the server saved the finished turn ($label)', + async ({ unsavedReads, heldFollowUp }) => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) + try { + const chatId = `chat-saved-after-complete-${unsavedReads}-${heldFollowUp}` + const history: MothershipChatHistory = { + id: chatId, + mode: 'agent', + title: 'Saved late', + messages: [], + activeStreamId: null, + resources: [], + } + let streamId: string | undefined + let completed = false + let detailReads = 0 + mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => { + if (contract.path !== '/api/mothership/chats/[chatId]') { + return Promise.resolve({ chats: [] }) + } + if (!completed) return Promise.resolve({ chat: history }) + detailReads++ + if (detailReads <= unsavedReads && streamId) { + return Promise.resolve({ + chat: { + ...history, + activeStreamId: streamId, + messages: [ + { id: streamId, role: 'user', content: 'Summarize the run' }, + { id: `live-assistant:${streamId}`, role: 'assistant', content: 'Done.' }, + ], + }, + }) + } + return Promise.resolve({ + chat: { + ...history, + messages: [ + { id: streamId, role: 'user', content: 'Summarize the run' }, + { id: 'saved-assistant', role: 'assistant', content: 'Done.' }, + ], + }, + }) + }) + let stream: ReadableStreamDefaultController | undefined + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') { + return fetchStub(input, init) + } + streamId = JSON.parse(String(init.body)).userMessageId + return new Response( + new ReadableStream({ + start(controller) { + stream = controller + }, + }), + { headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } } + ) + }) + const { getResult } = renderUseChatInChat(chatId, history) + + await act(async () => { + void getResult().sendMessage('Summarize the run') + await vi.advanceTimersByTimeAsync(50) + }) + expect(stream).toBeDefined() + if (heldFollowUp) { + useMothershipQueueStore + .getState() + .enqueue(chatId, { id: 'held-follow-up', content: 'And the next one' }) + useMothershipQueueStore.getState().setEditing(chatId, 'held-follow-up') + } + const emit = (event: Omit) => + stream?.enqueue( + new TextEncoder().encode( + `data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n` + ) + ) + const saved = () => + queryClient + .getQueryData(mothershipChatKeys.detail(chatId)) + ?.messages.some((message) => message.id === 'saved-assistant') === true + await act(async () => { + emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } }) + completed = true + emit({ seq: 2, type: 'complete', payload: { status: 'complete' } }) + stream?.close() + await vi.advanceTimersByTimeAsync(50) + }) + for (let second = 0; second < 90 && !saved(); second++) { + await act(async () => vi.advanceTimersByTimeAsync(1_000)) + } + + expect(saved()).toBe(true) + expect( + queryClient.getQueryData(mothershipChatKeys.detail(chatId)) + ?.activeStreamId + ).toBeNull() + expect(detailReads).toBe(unsavedReads + 1) + expect(getResult().isSending).toBe(false) + } finally { + vi.useRealTimers() + } + } + ) + describe.each([ { kind: 'browser action', 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 feb4fde0e4a..1a6ff1ba3a4 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -299,6 +299,11 @@ const STREAM_BATCH_FETCH_TIMEOUT_MS = 10_000 const STREAM_IDLE_TIMEOUT_MS = 45_000 const STREAM_CHAT_ID_RESOLVE_TIMEOUT_MS = 10_000 const CHAT_HISTORY_RECOVERY_TIMEOUT_MS = 10_000 +/** Backoff for re-reading a transcript the server has not yet saved a finished turn into. */ +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 const STOP_REQUEST_TIMEOUT_MS = 15_000 const DETACHED_CHAT_RETRY_BASE_MS = 1000 const DETACHED_CHAT_RETRY_MAX_MS = 30_000 @@ -992,6 +997,8 @@ export function useChat( // the copy-request-ID button functional after refetch). const streamRequestIdRef = useRef(undefined) const locallyTerminalStreamIdRef = useRef(undefined) + /** The finished stream whose saved transcript is being waited for, if any. */ + const persistedTurnWaitRef = useRef(null) const lastCursorRef = useRef('0') const logResyncedStreamIdRef = useRef(null) const activeStreamReturnRecoveryRef = useRef(null) @@ -1910,6 +1917,47 @@ export function useChat( resetHomeChatState() }, [isHomePage, resetHomeChatState]) + /** + * This tab finalizes on the stream's `complete` event, which the server sends + * before it saves the turn, so the transcript read right after can still be the + * in-flight copy: the stream listed as active and the answer under its live id. + * That copy matches the optimistic one, so nothing would read it again; re-read + * until the saved turn is there. The first pass joins finalize's own read. The + * wait ends as soon as this chat moves on: another send, or another chat. + */ + const awaitPersistedTurn = useCallback( + async (chatId: string, streamId: string) => { + if (persistedTurnWaitRef.current === streamId) return + persistedTurnWaitRef.current = streamId + const deadline = Date.now() + PERSISTED_TURN_WAIT_MS + try { + for (let attempt = 0; Date.now() < deadline; attempt++) { + if (attempt > 0) { + await sleep( + backoffWithJitter(attempt, null, { + baseMs: PERSISTED_TURN_REFETCH_BASE_MS, + maxMs: PERSISTED_TURN_REFETCH_MAX_DELAY_MS, + }) + ) + } + if (locallyTerminalStreamIdRef.current !== streamId || chatIdRef.current !== chatId) + return + await queryClient.refetchQueries( + { queryKey: mothershipChatKeys.detail(chatId), exact: true }, + { cancelRefetch: false } + ) + const history = queryClient.getQueryData( + mothershipChatKeys.detail(chatId) + ) + if (history?.activeStreamId !== streamId) return + } + } finally { + if (persistedTurnWaitRef.current === streamId) persistedTurnWaitRef.current = null + } + }, + [queryClient] + ) + useEffect(() => { if (!chatHistory) return @@ -3341,9 +3389,12 @@ export function useChat( if (completedActivityTracker?.generation === streamGenRef.current) { clearResourceActivity(completedActivityTracker, true) } + const terminalStreamId = + options?.streamTerminal !== false + ? (streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined) + : undefined if (options?.streamTerminal !== false) { - locallyTerminalStreamIdRef.current = - streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined + locallyTerminalStreamIdRef.current = terminalStreamId } clearActiveTurn() setTransportIdle() @@ -3352,9 +3403,13 @@ export function useChat( includeDetail: !hasQueuedFollowUp, ...(options?.targetChatId ? { targetChatId: options.targetChatId } : {}), }) + if (terminalStreamId && completedChatId) { + void awaitPersistedTurn(completedChatId, terminalStreamId) + } notifyTurnEnded({ error: isError }) }, [ + awaitPersistedTurn, clearResourceActivity, clearActiveTurn, invalidateChatQueries,