Skip to content

Commit 2e54a59

Browse files
committed
fix(mothership): keep waiting for a slow save, and wait with a follow-up queued
- The re-read gave up after six reads (about 8s), so a slow save left the finished turn on its in-flight copy. It now keeps reading on a capped backoff for up to two minutes, and still stops as soon as the chat moves on (another send or another chat). - It skipped the wait when a follow-up was queued, but a queued follow-up that does not go out at once (held for an edit) left the finished stream listed as running in the cache. It now runs for every finished turn; a follow-up that does go out ends it, since its send cancels the read and replaces the stream.
1 parent 603bf6f commit 2e54a59

2 files changed

Lines changed: 125 additions & 89 deletions

File tree

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

Lines changed: 109 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -2143,94 +2143,119 @@ describe('useChat remount send recovery', () => {
21432143
* The tab finalizes on the `complete` event, which reaches it before the server
21442144
* saves the turn. A transcript read in that gap is the server's in-flight copy;
21452145
* the tab must read again rather than keep it (live ids, the finished stream
2146-
* still listed as running) until something else happens to refetch.
2146+
* still listed as running) until something else happens to refetch. That holds
2147+
* when the save is slow, and when a follow-up is queued but not yet sent.
21472148
*/
2148-
it('re-reads a transcript fetched before the server saved the finished turn', async () => {
2149-
const chatId = 'chat-saved-after-complete'
2150-
const history: MothershipChatHistory = {
2151-
id: chatId,
2152-
mode: 'agent',
2153-
title: 'Saved late',
2154-
messages: [],
2155-
activeStreamId: null,
2156-
resources: [],
2157-
}
2158-
let streamId: string | undefined
2159-
let completed = false
2160-
const detailRequests: number[] = []
2161-
mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => {
2162-
if (contract.path !== '/api/mothership/chats/[chatId]') {
2163-
return Promise.resolve({ chats: [] })
2164-
}
2165-
if (!completed) return Promise.resolve({ chat: history })
2166-
detailRequests.push(Date.now())
2167-
if (detailRequests.length === 1 && streamId) {
2168-
return Promise.resolve({
2169-
chat: {
2170-
...history,
2171-
activeStreamId: streamId,
2172-
messages: [
2173-
{ id: streamId, role: 'user', content: 'Summarize the run' },
2174-
{ id: `live-assistant:${streamId}`, role: 'assistant', content: 'Done.' },
2175-
],
2176-
},
2149+
it.each([
2150+
{ label: 'right after the first read', unsavedReads: 1, heldFollowUp: false },
2151+
{ label: 'only after a slow save', unsavedReads: 9, heldFollowUp: false },
2152+
{ label: 'with a follow-up queued but held', unsavedReads: 1, heldFollowUp: true },
2153+
])(
2154+
're-reads a transcript fetched before the server saved the finished turn ($label)',
2155+
async ({ unsavedReads, heldFollowUp }) => {
2156+
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
2157+
try {
2158+
const chatId = `chat-saved-after-complete-${unsavedReads}-${heldFollowUp}`
2159+
const history: MothershipChatHistory = {
2160+
id: chatId,
2161+
mode: 'agent',
2162+
title: 'Saved late',
2163+
messages: [],
2164+
activeStreamId: null,
2165+
resources: [],
2166+
}
2167+
let streamId: string | undefined
2168+
let completed = false
2169+
let detailReads = 0
2170+
mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => {
2171+
if (contract.path !== '/api/mothership/chats/[chatId]') {
2172+
return Promise.resolve({ chats: [] })
2173+
}
2174+
if (!completed) return Promise.resolve({ chat: history })
2175+
detailReads++
2176+
if (detailReads <= unsavedReads && streamId) {
2177+
return Promise.resolve({
2178+
chat: {
2179+
...history,
2180+
activeStreamId: streamId,
2181+
messages: [
2182+
{ id: streamId, role: 'user', content: 'Summarize the run' },
2183+
{ id: `live-assistant:${streamId}`, role: 'assistant', content: 'Done.' },
2184+
],
2185+
},
2186+
})
2187+
}
2188+
return Promise.resolve({
2189+
chat: {
2190+
...history,
2191+
messages: [
2192+
{ id: streamId, role: 'user', content: 'Summarize the run' },
2193+
{ id: 'saved-assistant', role: 'assistant', content: 'Done.' },
2194+
],
2195+
},
2196+
})
21772197
})
2178-
}
2179-
return Promise.resolve({
2180-
chat: {
2181-
...history,
2182-
messages: [
2183-
{ id: streamId, role: 'user', content: 'Summarize the run' },
2184-
{ id: 'saved-assistant', role: 'assistant', content: 'Done.' },
2185-
],
2186-
},
2187-
})
2188-
})
2189-
let stream: ReadableStreamDefaultController<Uint8Array> | undefined
2190-
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2191-
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
2192-
return fetchStub(input, init)
2193-
}
2194-
streamId = JSON.parse(String(init.body)).userMessageId
2195-
return new Response(
2196-
new ReadableStream<Uint8Array>({
2197-
start(controller) {
2198-
stream = controller
2199-
},
2200-
}),
2201-
{ headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } }
2202-
)
2203-
})
2204-
const { getResult } = renderUseChatInChat(chatId, history)
2198+
let stream: ReadableStreamDefaultController<Uint8Array> | undefined
2199+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2200+
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
2201+
return fetchStub(input, init)
2202+
}
2203+
streamId = JSON.parse(String(init.body)).userMessageId
2204+
return new Response(
2205+
new ReadableStream<Uint8Array>({
2206+
start(controller) {
2207+
stream = controller
2208+
},
2209+
}),
2210+
{ headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } }
2211+
)
2212+
})
2213+
const { getResult } = renderUseChatInChat(chatId, history)
22052214

2206-
await act(async () => {
2207-
void getResult().sendMessage('Summarize the run')
2208-
})
2209-
await waitFor(() => stream !== undefined)
2210-
const emit = (event: Omit<MothershipStreamV1EventEnvelope, 'v' | 'ts' | 'stream'>) =>
2211-
stream?.enqueue(
2212-
new TextEncoder().encode(
2213-
`data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n`
2214-
)
2215-
)
2216-
emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } })
2217-
completed = true
2218-
emit({ seq: 2, type: 'complete', payload: { status: 'complete' } })
2219-
stream?.close()
2215+
await act(async () => {
2216+
void getResult().sendMessage('Summarize the run')
2217+
await vi.advanceTimersByTimeAsync(50)
2218+
})
2219+
expect(stream).toBeDefined()
2220+
if (heldFollowUp) {
2221+
useMothershipQueueStore
2222+
.getState()
2223+
.enqueue(chatId, { id: 'held-follow-up', content: 'And the next one' })
2224+
useMothershipQueueStore.getState().setEditing(chatId, 'held-follow-up')
2225+
}
2226+
const emit = (event: Omit<MothershipStreamV1EventEnvelope, 'v' | 'ts' | 'stream'>) =>
2227+
stream?.enqueue(
2228+
new TextEncoder().encode(
2229+
`data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n`
2230+
)
2231+
)
2232+
const saved = () =>
2233+
queryClient
2234+
.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
2235+
?.messages.some((message) => message.id === 'saved-assistant') === true
2236+
await act(async () => {
2237+
emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } })
2238+
completed = true
2239+
emit({ seq: 2, type: 'complete', payload: { status: 'complete' } })
2240+
stream?.close()
2241+
await vi.advanceTimersByTimeAsync(50)
2242+
})
2243+
for (let second = 0; second < 90 && !saved(); second++) {
2244+
await act(async () => vi.advanceTimersByTimeAsync(1_000))
2245+
}
22202246

2221-
await waitFor(
2222-
() =>
2223-
queryClient
2224-
.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
2225-
?.messages.some((message) => message.id === 'saved-assistant') === true
2226-
)
2227-
expect(
2228-
queryClient.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
2229-
?.activeStreamId
2230-
).toBeNull()
2231-
expect(detailRequests).toHaveLength(2)
2232-
expect(getResult().isSending).toBe(false)
2233-
})
2247+
expect(saved()).toBe(true)
2248+
expect(
2249+
queryClient.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
2250+
?.activeStreamId
2251+
).toBeNull()
2252+
expect(detailReads).toBe(unsavedReads + 1)
2253+
expect(getResult().isSending).toBe(false)
2254+
} finally {
2255+
vi.useRealTimers()
2256+
}
2257+
}
2258+
)
22342259

22352260
describe.each([
22362261
{

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts‎

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -301,7 +301,9 @@ const STREAM_CHAT_ID_RESOLVE_TIMEOUT_MS = 10_000
301301
const CHAT_HISTORY_RECOVERY_TIMEOUT_MS = 10_000
302302
/** Backoff for re-reading a transcript the server has not yet saved a finished turn into. */
303303
const PERSISTED_TURN_REFETCH_BASE_MS = 250
304-
const PERSISTED_TURN_REFETCH_ATTEMPTS = 6
304+
const PERSISTED_TURN_REFETCH_MAX_DELAY_MS = 5_000
305+
/** How long a finished turn's save is waited for; a slow save still lands well inside it. */
306+
const PERSISTED_TURN_WAIT_MS = 120_000
305307
const STOP_REQUEST_TIMEOUT_MS = 15_000
306308
const DETACHED_CHAT_RETRY_BASE_MS = 1000
307309
const DETACHED_CHAT_RETRY_MAX_MS = 30_000
@@ -1920,15 +1922,24 @@ export function useChat(
19201922
* before it saves the turn, so the transcript read right after can still be the
19211923
* in-flight copy: the stream listed as active and the answer under its live id.
19221924
* That copy matches the optimistic one, so nothing would read it again; re-read
1923-
* until the saved turn is there. The first pass joins finalize's own read.
1925+
* until the saved turn is there. The first pass joins finalize's own read. The
1926+
* wait ends as soon as this chat moves on: another send, or another chat.
19241927
*/
19251928
const awaitPersistedTurn = useCallback(
19261929
async (chatId: string, streamId: string) => {
19271930
if (persistedTurnWaitRef.current === streamId) return
19281931
persistedTurnWaitRef.current = streamId
1932+
const deadline = Date.now() + PERSISTED_TURN_WAIT_MS
19291933
try {
1930-
for (let attempt = 0; attempt < PERSISTED_TURN_REFETCH_ATTEMPTS; attempt++) {
1931-
if (attempt > 0) await sleep(PERSISTED_TURN_REFETCH_BASE_MS * 2 ** (attempt - 1))
1934+
for (let attempt = 0; Date.now() < deadline; attempt++) {
1935+
if (attempt > 0) {
1936+
await sleep(
1937+
Math.min(
1938+
PERSISTED_TURN_REFETCH_BASE_MS * 2 ** (attempt - 1),
1939+
PERSISTED_TURN_REFETCH_MAX_DELAY_MS
1940+
)
1941+
)
1942+
}
19321943
if (locallyTerminalStreamIdRef.current !== streamId || chatIdRef.current !== chatId)
19331944
return
19341945
await queryClient.refetchQueries(
@@ -3392,7 +3403,7 @@ export function useChat(
33923403
includeDetail: !hasQueuedFollowUp,
33933404
...(options?.targetChatId ? { targetChatId: options.targetChatId } : {}),
33943405
})
3395-
if (terminalStreamId && completedChatId && !hasQueuedFollowUp) {
3406+
if (terminalStreamId && completedChatId) {
33963407
void awaitPersistedTurn(completedChatId, terminalStreamId)
33973408
}
33983409
notifyTurnEnded({ error: isError })

0 commit comments

Comments
 (0)