Skip to content

Commit 39fce7c

Browse files
committed
fix(chat): mark a completed live-stream detail stale instead of skipping it
1 parent f3924ce commit 39fce7c

3 files changed

Lines changed: 105 additions & 16 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@ import {
8282
} from '@/app/workspace/[workspaceId]/home/hooks/send-handoff'
8383
import { useChat } from '@/app/workspace/[workspaceId]/home/hooks/use-chat'
8484
import { type MothershipChatHistory, mothershipChatKeys } from '@/hooks/queries/mothership-chats'
85+
import { handleMothershipChatStatusEvent } from '@/hooks/use-mothership-chat-events'
8586
import { useExecutionStore } from '@/stores/execution/store'
8687
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
8788

@@ -1821,4 +1822,81 @@ describe('useChat remount send recovery', () => {
18211822
expect(getResult().messageQueue.map((entry) => entry.id)).toEqual(['unsent-entry'])
18221823
expect(state.postBodies).toHaveLength(0)
18231824
})
1825+
1826+
it('loads the saved transcript once when its own stream completes', async () => {
1827+
const chatId = 'chat-own-completion'
1828+
const history: MothershipChatHistory = {
1829+
id: chatId,
1830+
mode: 'agent',
1831+
title: 'Own stream',
1832+
messages: [],
1833+
activeStreamId: null,
1834+
resources: [],
1835+
}
1836+
const saved = [
1837+
{ id: 'saved-user', role: 'user', content: 'Summarize the run' },
1838+
{ id: 'saved-assistant', role: 'assistant', content: 'Done.' },
1839+
]
1840+
const detailRequests: string[] = []
1841+
mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => {
1842+
if (contract.path !== '/api/mothership/chats/[chatId]') {
1843+
return Promise.resolve({ chats: [] })
1844+
}
1845+
detailRequests.push(chatId)
1846+
return Promise.resolve({ chat: { ...history, messages: saved } })
1847+
})
1848+
let stream: ReadableStreamDefaultController<Uint8Array> | undefined
1849+
let streamId: string | undefined
1850+
const emit = (event: Omit<MothershipStreamV1EventEnvelope, 'v' | 'ts' | 'stream'>) =>
1851+
stream?.enqueue(
1852+
new TextEncoder().encode(
1853+
`data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n`
1854+
)
1855+
)
1856+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1857+
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
1858+
return fetchStub(input, init)
1859+
}
1860+
streamId = JSON.parse(String(init.body)).userMessageId
1861+
return new Response(
1862+
new ReadableStream<Uint8Array>({
1863+
start(controller) {
1864+
stream = controller
1865+
},
1866+
}),
1867+
{ headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } }
1868+
)
1869+
})
1870+
const { getResult } = renderUseChatInChat(chatId, history)
1871+
1872+
await act(async () => {
1873+
void getResult().sendMessage('Summarize the run')
1874+
})
1875+
await waitFor(() => stream !== undefined)
1876+
emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } })
1877+
await waitFor(
1878+
() =>
1879+
queryClient.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
1880+
?.activeStreamId === streamId
1881+
)
1882+
/** The server publishes `completed` after persisting and before closing the stream. */
1883+
handleMothershipChatStatusEvent(queryClient, 'ws-1', {
1884+
chatId,
1885+
type: 'completed',
1886+
streamId,
1887+
})
1888+
emit({ seq: 2, type: 'complete', payload: { status: 'complete' } })
1889+
stream?.close()
1890+
1891+
await waitFor(() => !getResult().isSending && detailRequests.length > 0)
1892+
await act(async () => {
1893+
await sleep(50)
1894+
})
1895+
expect(detailRequests).toHaveLength(1)
1896+
expect(
1897+
queryClient
1898+
.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
1899+
?.messages.map((message) => message.id)
1900+
).toEqual(['saved-user', 'saved-assistant'])
1901+
})
18241902
})

‎apps/sim/hooks/use-mothership-chat-events.test.ts‎

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -223,17 +223,22 @@ describe('chat detail refetches driven by status events', () => {
223223
unsubscribe()
224224
})
225225

226-
it('reloads the transcript when a stream this viewer is not rendering completes', async () => {
227-
const { queryClient, fetchTranscript, unsubscribe } = mountDetail({
228-
...liveStream,
229-
messages: [{ id: 'stream-1' }] as MothershipChatHistory['messages'],
230-
})
226+
it('reloads the saved transcript when a cached mid-stream detail is opened after completion', async () => {
227+
const queryClient = new QueryClient()
228+
const fetchTranscript = vi.fn(async () => ({ ...liveStream, activeStreamId: null }))
229+
queryClient.setQueryData(mothershipChatKeys.detail('chat-1'), liveStream)
231230

232231
handleMothershipChatStatusEvent(queryClient, 'ws-1', {
233232
chatId: 'chat-1',
234233
type: 'completed',
235234
streamId: 'stream-1',
236235
})
236+
const unsubscribe = new QueryObserver(queryClient, {
237+
queryKey: mothershipChatKeys.detail('chat-1'),
238+
queryFn: fetchTranscript,
239+
staleTime: Number.POSITIVE_INFINITY,
240+
}).subscribe(() => {})
241+
237242
await vi.waitFor(() => expect(fetchTranscript).toHaveBeenCalledTimes(1))
238243
unsubscribe()
239244
})

‎apps/sim/hooks/use-mothership-chat-events.ts‎

Lines changed: 17 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -59,22 +59,14 @@ function hasNewerKnownActiveStream(current: MothershipChatHistory | undefined, s
5959
return activeIndex > eventStreamIndex
6060
}
6161

62-
/**
63-
* Returns true when refetching the chat detail for a stream event would only
64-
* reload the transcript this client already holds or is about to reload. A
65-
* completion of the viewer's own live stream is skipped because the server
66-
* persists the turn before closing that stream, and the client's own
67-
* finalization refetches the detail once it sees the close.
68-
*/
6962
function shouldSkipDetailInvalidationForStreamEvent(
7063
current: MothershipChatHistory | undefined,
7164
payload: ChatStatusEventPayload
7265
) {
7366
if (!current?.activeStreamId) return false
7467
if (!payload.streamId) return isLocalOptimisticActiveStream(current)
75-
if (current.activeStreamId === payload.streamId) {
76-
return payload.type === 'started' || isLocalOptimisticActiveStream(current)
77-
}
68+
if (payload.type === 'started' && current.activeStreamId === payload.streamId) return true
69+
if (current.activeStreamId === payload.streamId) return false
7870
if (hasNewerKnownActiveStream(current, payload.streamId)) return true
7971
return (
8072
payload.type === 'completed' &&
@@ -145,7 +137,21 @@ export function handleMothershipChatStatusEvent(
145137
mothershipChatKeys.detail(payload.chatId)
146138
)
147139
if (shouldSkipDetailInvalidationForStreamEvent(current, payload)) return
148-
queryClient.invalidateQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) })
140+
/**
141+
* A completion of the cached live stream only marks the detail stale. The
142+
* server persists the turn before closing the stream, so a surface rendering
143+
* it refetches through its own finalization; the live message alone cannot
144+
* tell this tab's stream from a server-loaded mid-stream snapshot, so any
145+
* other cached copy reloads the saved transcript on its next mount.
146+
*/
147+
const completesCachedLiveStream =
148+
payload.type === 'completed' &&
149+
current?.activeStreamId === payload.streamId &&
150+
isLocalOptimisticActiveStream(current)
151+
queryClient.invalidateQueries({
152+
queryKey: mothershipChatKeys.detail(payload.chatId),
153+
...(completesCachedLiveStream ? { refetchType: 'none' as const } : {}),
154+
})
149155
}
150156

151157
/**

0 commit comments

Comments
 (0)