Skip to content

Commit 1f114dd

Browse files
authored
fix(mothership): re-read the transcript until the server has saved a finished turn (#8676)
* fix(mothership): re-read the transcript until the server has saved a finished turn A tab finalizes a turn on the stream's `complete` event and reads the saved transcript. The server sends that event before it persists the turn and clears the chat's stream marker, so the read can land in the gap and return the in-flight copy: the finished stream still listed as active and the answer under its live id. That copy matches the optimistic one, so the history effect saw nothing new, and the chat's `completed` event only marks a cached live stream stale, so the tab kept it until a reload. On the local QA stack this happened in roughly one turn in four: the finished answer had no Fork action, and the cached stream marker kept a queued message from draining on a later return. Finalize now re-reads with a short backoff while the server still lists the stream it saw end. The first pass joins finalize's own read, so a turn saved in time is still read once. * 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. * improvement(mothership): pace the saved-turn re-read with the shared jittered backoff
1 parent c8ef18d commit 1f114dd

2 files changed

Lines changed: 176 additions & 2 deletions

File tree

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

Lines changed: 119 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2138,6 +2138,125 @@ describe('useChat remount send recovery', () => {
21382138
?.messages.map((message) => message.id)
21392139
).toEqual(['saved-user', 'saved-assistant'])
21402140
})
2141+
2142+
/**
2143+
* The tab finalizes on the `complete` event, which reaches it before the server
2144+
* saves the turn. A transcript read in that gap is the server's in-flight copy;
2145+
* 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. That holds
2147+
* when the save is slow, and when a follow-up is queued but not yet sent.
2148+
*/
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+
})
2197+
})
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)
2214+
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+
}
2246+
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+
)
2259+
21412260
describe.each([
21422261
{
21432262
kind: 'browser action',

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

Lines changed: 57 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -299,6 +299,11 @@ const STREAM_BATCH_FETCH_TIMEOUT_MS = 10_000
299299
const STREAM_IDLE_TIMEOUT_MS = 45_000
300300
const STREAM_CHAT_ID_RESOLVE_TIMEOUT_MS = 10_000
301301
const CHAT_HISTORY_RECOVERY_TIMEOUT_MS = 10_000
302+
/** Backoff for re-reading a transcript the server has not yet saved a finished turn into. */
303+
const PERSISTED_TURN_REFETCH_BASE_MS = 250
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
302307
const STOP_REQUEST_TIMEOUT_MS = 15_000
303308
const DETACHED_CHAT_RETRY_BASE_MS = 1000
304309
const DETACHED_CHAT_RETRY_MAX_MS = 30_000
@@ -992,6 +997,8 @@ export function useChat(
992997
// the copy-request-ID button functional after refetch).
993998
const streamRequestIdRef = useRef<string | undefined>(undefined)
994999
const locallyTerminalStreamIdRef = useRef<string | undefined>(undefined)
1000+
/** The finished stream whose saved transcript is being waited for, if any. */
1001+
const persistedTurnWaitRef = useRef<string | null>(null)
9951002
const lastCursorRef = useRef('0')
9961003
const logResyncedStreamIdRef = useRef<string | null>(null)
9971004
const activeStreamReturnRecoveryRef = useRef<ActiveStreamRecovery | null>(null)
@@ -1915,6 +1922,47 @@ export function useChat(
19151922
resetHomeChatState()
19161923
}, [isHomePage, resetHomeChatState])
19171924

1925+
/**
1926+
* This tab finalizes on the stream's `complete` event, which the server sends
1927+
* before it saves the turn, so the transcript read right after can still be the
1928+
* in-flight copy: the stream listed as active and the answer under its live id.
1929+
* That copy matches the optimistic one, so nothing would read it again; re-read
1930+
* until the saved turn is there. The first pass joins finalize's own read. The
1931+
* wait ends as soon as this chat moves on: another send, or another chat.
1932+
*/
1933+
const awaitPersistedTurn = useCallback(
1934+
async (chatId: string, streamId: string) => {
1935+
if (persistedTurnWaitRef.current === streamId) return
1936+
persistedTurnWaitRef.current = streamId
1937+
const deadline = Date.now() + PERSISTED_TURN_WAIT_MS
1938+
try {
1939+
for (let attempt = 0; Date.now() < deadline; attempt++) {
1940+
if (attempt > 0) {
1941+
await sleep(
1942+
backoffWithJitter(attempt, null, {
1943+
baseMs: PERSISTED_TURN_REFETCH_BASE_MS,
1944+
maxMs: PERSISTED_TURN_REFETCH_MAX_DELAY_MS,
1945+
})
1946+
)
1947+
}
1948+
if (locallyTerminalStreamIdRef.current !== streamId || chatIdRef.current !== chatId)
1949+
return
1950+
await queryClient.refetchQueries(
1951+
{ queryKey: mothershipChatKeys.detail(chatId), exact: true },
1952+
{ cancelRefetch: false }
1953+
)
1954+
const history = queryClient.getQueryData<MothershipChatHistory>(
1955+
mothershipChatKeys.detail(chatId)
1956+
)
1957+
if (history?.activeStreamId !== streamId) return
1958+
}
1959+
} finally {
1960+
if (persistedTurnWaitRef.current === streamId) persistedTurnWaitRef.current = null
1961+
}
1962+
},
1963+
[queryClient]
1964+
)
1965+
19181966
useEffect(() => {
19191967
if (!chatHistory) return
19201968

@@ -3347,9 +3395,12 @@ export function useChat(
33473395
if (completedActivityTracker?.generation === streamGenRef.current) {
33483396
clearResourceActivity(completedActivityTracker, true)
33493397
}
3398+
const terminalStreamId =
3399+
options?.streamTerminal !== false
3400+
? (streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined)
3401+
: undefined
33503402
if (options?.streamTerminal !== false) {
3351-
locallyTerminalStreamIdRef.current =
3352-
streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined
3403+
locallyTerminalStreamIdRef.current = terminalStreamId
33533404
}
33543405
clearActiveTurn()
33553406
setTransportIdle()
@@ -3358,9 +3409,13 @@ export function useChat(
33583409
includeDetail: !hasQueuedFollowUp,
33593410
...(options?.targetChatId ? { targetChatId: options.targetChatId } : {}),
33603411
})
3412+
if (terminalStreamId && completedChatId) {
3413+
void awaitPersistedTurn(completedChatId, terminalStreamId)
3414+
}
33613415
notifyTurnEnded({ error: isError })
33623416
},
33633417
[
3418+
awaitPersistedTurn,
33643419
clearResourceActivity,
33653420
clearActiveTurn,
33663421
invalidateChatQueries,

0 commit comments

Comments
 (0)