Skip to content

Commit 6585682

Browse files
committed
fix(chat): preserve live turns and recover interrupted streams
1 parent 6aac637 commit 6585682

2 files changed

Lines changed: 114 additions & 23 deletions

File tree

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

Lines changed: 70 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1153,9 +1153,13 @@ describe('useChat remount send recovery', () => {
11531153
}
11541154
)
11551155

1156-
it.each([false, true])(
1157-
'keeps the owning optimistic user through lagging resource history (model output: %s)',
1158-
async (hasOutput) => {
1156+
it.each(
1157+
[false, true].flatMap((hasOutput) =>
1158+
(['assistant', 'user', 'neither'] as const).map((retained) => ({ hasOutput, retained }))
1159+
)
1160+
)(
1161+
'keeps the live turn through lagging resource history (output: $hasOutput, retained: $retained)',
1162+
async ({ hasOutput, retained }) => {
11591163
state.postBehavior = 'task'
11601164
const history: MothershipChatHistory = {
11611165
id: 'chat-lagging-user',
@@ -1208,19 +1212,22 @@ describe('useChat remount send recovery', () => {
12081212
.messages.find((message) => message.role === 'assistant')!
12091213
await act(async () => getResult().addResource(search))
12101214
await waitFor(() => historyReads > 0)
1215+
const beforeHistoryRefresh = getResult()
12111216
await act(async () =>
12121217
staleHistory.resolve({
12131218
chat: {
12141219
...history,
12151220
title: 'Hydrated during run',
12161221
activeStreamId: sentUser.id,
1217-
messages: [liveAssistant],
1222+
messages:
1223+
retained === 'assistant' ? [liveAssistant] : retained === 'user' ? [sentUser] : [],
12181224
resources: [search],
12191225
},
12201226
})
12211227
)
12221228
await waitFor(
12231229
() =>
1230+
getResult() !== beforeHistoryRefresh &&
12241231
queryClient.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(history.id))
12251232
?.title === 'Hydrated during run'
12261233
)
@@ -1229,6 +1236,9 @@ describe('useChat remount send recovery', () => {
12291236
liveAssistant.id,
12301237
])
12311238
expect(getResult().messages[0].content).toBe('Orion acceptance policy')
1239+
expect(getResult().isSending).toBe(true)
1240+
expect(getResult().messages[1].content).toBe(liveAssistant.content)
1241+
if (hasOutput) expect(getResult().messages[1].contentBlocks).toHaveLength(1)
12321242
await act(async () =>
12331243
queryClient.setQueryData(mothershipChatKeys.detail(history.id), {
12341244
...history,
@@ -1411,6 +1421,62 @@ describe('useChat remount send recovery', () => {
14111421
expect(state.postBodies[0]).toHaveProperty('effort')
14121422
})
14131423

1424+
it('recovers a running turn after reconnect exhaustion without reloading or resending', async () => {
1425+
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
1426+
try {
1427+
let online = false
1428+
let recoveredTail = false
1429+
let failedReconnects = 0
1430+
const history: MothershipChatHistory = {
1431+
id: 'chat-reconnect-exhausted',
1432+
mode: 'agent',
1433+
title: 'Reconnect',
1434+
messages: [],
1435+
activeStreamId: null,
1436+
resources: [],
1437+
}
1438+
mockRequestJson.mockImplementation(() =>
1439+
Promise.resolve({
1440+
chat: {
1441+
...history,
1442+
activeStreamId: state.postBodies[0]?.userMessageId ?? null,
1443+
},
1444+
})
1445+
)
1446+
state.postBehavior = 'accept'
1447+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1448+
const url = String(input)
1449+
if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init)
1450+
if (!online) {
1451+
failedReconnects++
1452+
throw new TypeError('Failed to fetch')
1453+
}
1454+
if (url.includes('batch=true')) {
1455+
return Response.json({ success: true, events: [], status: 'streaming' })
1456+
}
1457+
recoveredTail = true
1458+
return new Response(new ReadableStream<Uint8Array>(), {
1459+
headers: { 'Content-Type': 'text/event-stream' },
1460+
})
1461+
})
1462+
const { getResult } = renderUseChatInChat(history.id, history)
1463+
await act(async () => {
1464+
void getResult().sendMessage('Continue working')
1465+
})
1466+
for (let second = 0; second < 240 && failedReconnects < 12; second++) {
1467+
await act(async () => vi.advanceTimersByTimeAsync(1_000))
1468+
}
1469+
expect(failedReconnects).toBeGreaterThanOrEqual(12)
1470+
online = true
1471+
await act(async () => vi.advanceTimersByTimeAsync(30_000))
1472+
expect(recoveredTail).toBe(true)
1473+
expect(getResult().isSending).toBe(true)
1474+
expect(state.postBodies).toHaveLength(1)
1475+
} finally {
1476+
vi.useRealTimers()
1477+
}
1478+
})
1479+
14141480
it('preserves a visible workflow watch when Stop persists the partial response', async () => {
14151481
state.postBehavior = 'task'
14161482
const { getResult } = renderUseChatInChat('chat-a')

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

Lines changed: 44 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -187,6 +187,13 @@ export interface SendMessageOptions {
187187
assistantSearchLevel?: AssistantSearchLevel
188188
}
189189

190+
interface FinalizeOptions {
191+
error?: boolean
192+
targetChatId?: string
193+
/** A lost transport must remain recoverable while the server is still running. */
194+
streamTerminal?: boolean
195+
}
196+
190197
/**
191198
* `true` when the send owns the transcript (rendered, or handed to reconnect),
192199
* `false` when the caller should restore the queue entry, and the object form
@@ -883,9 +890,7 @@ export function useChat(
883890
const resolveDetachedChatForStreamRef = useRef<
884891
(streamId: string, signal?: AbortSignal) => Promise<DetachedChatResolution>
885892
>(async () => ({ terminal: false }))
886-
const finalizeRef = useRef<(options?: { error?: boolean; targetChatId?: string }) => void>(
887-
() => {}
888-
)
893+
const finalizeRef = useRef<(options?: FinalizeOptions) => void>(() => {})
889894
const recoveringQueuedSendHandoffRef = useRef<ActiveQueuedSendHandoffRecovery | null>(null)
890895
const recoverActiveStreamRef = useRef<
891896
(reason: 'pageshow' | 'visible' | 'online' | 'exhausted_recheck') => Promise<void>
@@ -1238,21 +1243,34 @@ export function useChat(
12381243
options?.requestMode ?? chatHistory?.mode ?? (organizationId ? 'assistant' : 'agent')
12391244
const pendingTurn =
12401245
chatHistory?.id === (initialChatId ?? chatIdRef.current) ? activeTurnRef.current : null
1246+
const liveContent = streamingContentRef.current
1247+
const liveBlocks = streamingBlocksRef.current
12411248
const messages = useMemo(() => {
12421249
const source = chatHistory?.messages.map(toDisplayMessage) ?? [...pendingMessages]
1243-
/** A resource-history read can lag admission; keep this chat's own optimistic user visible. */
1244-
if (pendingTurn && !source.some((message) => message.id === pendingTurn.userMessageId)) {
1245-
const assistantIndex = source.findIndex(
1246-
(message) => message.id === pendingTurn.assistantMessageId
1247-
)
1248-
source.splice(
1249-
assistantIndex < 0 ? source.length : assistantIndex,
1250-
0,
1251-
pendingTurn.optimisticUserMessage
1252-
)
1250+
/** History reads can lag the live turn; retain both its owner and current response. */
1251+
if (pendingTurn) {
1252+
let ownerIndex = source.findIndex((message) => message.id === pendingTurn.userMessageId)
1253+
if (ownerIndex < 0) {
1254+
const assistantIndex = source.findIndex(
1255+
(message) => message.id === pendingTurn.assistantMessageId
1256+
)
1257+
ownerIndex = assistantIndex < 0 ? source.length : assistantIndex
1258+
source.splice(ownerIndex, 0, pendingTurn.optimisticUserMessage)
1259+
}
1260+
const assistant = source[ownerIndex + 1]
1261+
const liveAssistant = {
1262+
...pendingTurn.optimisticAssistantMessage,
1263+
content: liveContent,
1264+
contentBlocks: liveBlocks,
1265+
}
1266+
if (assistant?.id === pendingTurn.assistantMessageId) {
1267+
source[ownerIndex + 1] = { ...assistant, ...liveAssistant }
1268+
} else if (assistant?.role !== 'assistant') {
1269+
source.splice(ownerIndex + 1, 0, liveAssistant)
1270+
}
12531271
}
12541272
return source.map((m) => restoreRevealedSimKeysForMessage(m, revealedSimKeysRef.current))
1255-
}, [chatHistory, pendingMessages, pendingTurn])
1273+
}, [chatHistory, pendingMessages, pendingTurn, liveContent, liveBlocks])
12561274
const addResource = useCallback(
12571275
(resourceUpdate: MothershipResourceUpdate): boolean => {
12581276
// The single fan-in for tab creation, so the invariant lives here.
@@ -2006,7 +2024,11 @@ export function useChat(
20062024
}
20072025
if (!succeeded && streamGenRef.current === gen) {
20082026
try {
2009-
finalizeRef.current({ error: true, targetChatId: chatHistory.id })
2027+
finalizeRef.current({
2028+
error: true,
2029+
targetChatId: chatHistory.id,
2030+
streamTerminal: false,
2031+
})
20102032
} catch {
20112033
setTransportIdle()
20122034
abortControllerRef.current = null
@@ -2941,7 +2963,7 @@ export function useChat(
29412963
shouldContinue: isSameRecoverySubject,
29422964
})
29432965
if (!succeeded && streamGenRef.current === recoveryGen && isSameRecoverySubject()) {
2944-
finalizeRef.current({ error: true, targetChatId: chatId })
2966+
finalizeRef.current({ error: true, targetChatId: chatId, streamTerminal: false })
29452967
}
29462968
}
29472969
})()
@@ -3170,7 +3192,7 @@ export function useChat(
31703192
)
31713193

31723194
const finalize = useCallback(
3173-
(options?: { error?: boolean; targetChatId?: string }) => {
3195+
(options?: FinalizeOptions) => {
31743196
const isError = !!options?.error
31753197
if (isError) {
31763198
const blocks = streamingBlocksRef.current
@@ -3219,8 +3241,10 @@ export function useChat(
32193241
if (completedActivityTracker?.generation === streamGenRef.current) {
32203242
clearResourceActivity(completedActivityTracker, true)
32213243
}
3222-
locallyTerminalStreamIdRef.current =
3223-
streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined
3244+
if (options?.streamTerminal !== false) {
3245+
locallyTerminalStreamIdRef.current =
3246+
streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? undefined
3247+
}
32243248
clearActiveTurn()
32253249
setTransportIdle()
32263250
abortControllerRef.current = null
@@ -3804,6 +3828,7 @@ export function useChat(
38043828
if (gen !== undefined && streamGenRef.current === gen) {
38053829
finalize({
38063830
error: true,
3831+
streamTerminal: false,
38073832
...(streamTargetChatId ? { targetChatId: streamTargetChatId } : {}),
38083833
})
38093834
}

0 commit comments

Comments
 (0)