Skip to content

Commit 662bd70

Browse files
authored
fix(mothership): re-attach at once when the user returns to a recovered stream (#8675)
* fix(mothership): re-attach at once when the user returns to a recovered stream Once the user left a running chat and came back, the return recovery owned the stream for the rest of the turn, and every later online/visible/pageshow event just joined it. So a network drop after that waited out whatever the recovery was doing: a tail that went silent held the stream until the 45s idle timeout, and a tail that failed slept out a reconnect backoff of up to 30s. The same drop on a stream the send still owned re-attached immediately, because the return signal supersedes the send's reader. On the local QA stack the stream resumed 48s after the network came back. A return signal now supersedes an in-flight recovery the same way, and the new recovery re-attaches from the cursor. * fix(mothership): finish a turn that ended while its reader was silent A return recovery used the chat history's `activeStreamId` to decide what to re-attach to. When the turn had ended while this surface was not listening (its reader stalled, or a later return superseded the recovery that held it), the history listed no running turn, so recovery returned without touching the stream this surface still showed as running, and the chat stayed on Stop until an idle timeout or backoff happened to run into the terminal state. When the history lists no running turn but this surface is still sending, recovery now resolves that stream: its terminal status replays the remaining events and finalizes. * fix(mothership): leave a send waiting for admission alone on return A send shows as running before its POST is admitted, but the chat cannot list it yet. Resolving that stream on a return event read it as ended and aborted the POST. Recovery now resolves only a locally running stream whose POST was admitted. * fix(mothership): finish an admitted turn whose POST never answered on return Recovery left a send alone while its POST had not answered, so a turn the server admitted and finished, whose answer never reached the client, kept the chat on Stop with nothing to clear it. The send now counts as admitted once the loaded chat holds its message, and recovery resolves its stream as any other.
1 parent 1f114dd commit 662bd70

2 files changed

Lines changed: 270 additions & 5 deletions

File tree

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

Lines changed: 242 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1136,6 +1136,248 @@ describe('useChat remount send recovery', () => {
11361136
}
11371137
})
11381138

1139+
/**
1140+
* After the user leaves and returns once, the return recovery owns the stream for
1141+
* the rest of the turn. When the network then drops, its tail either goes silent
1142+
* (the socket stalls) or fails into the reconnect backoff, which grows to 30s.
1143+
* Coming back online must re-attach at once, as it does while the send still owns
1144+
* the stream, instead of waiting out the idle timeout or the backoff.
1145+
*/
1146+
it.each(['stalled', 'failed'] as const)(
1147+
're-attaches at once when the network returns to a return recovery whose tail %s',
1148+
async (tailOutcome) => {
1149+
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
1150+
try {
1151+
let online = true
1152+
let backOnline = false
1153+
let tailOpenedAfterReturn = false
1154+
let failedReconnects = 0
1155+
const openTails: ReadableStreamDefaultController<Uint8Array>[] = []
1156+
const history: MothershipChatHistory = {
1157+
id: `chat-recovery-${tailOutcome}`,
1158+
mode: 'agent',
1159+
title: 'Recovery',
1160+
messages: [],
1161+
activeStreamId: null,
1162+
resources: [],
1163+
}
1164+
mockRequestJson.mockImplementation(() =>
1165+
Promise.resolve({
1166+
chat: { ...history, activeStreamId: state.postBodies[0]?.userMessageId ?? null },
1167+
})
1168+
)
1169+
state.postBehavior = 'accept'
1170+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1171+
const url = String(input)
1172+
if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init)
1173+
if (!online) {
1174+
failedReconnects++
1175+
throw new TypeError('Failed to fetch')
1176+
}
1177+
if (url.includes('batch=true')) {
1178+
return Response.json({ success: true, events: [], status: 'streaming' })
1179+
}
1180+
if (backOnline) tailOpenedAfterReturn = true
1181+
return new Response(
1182+
new ReadableStream<Uint8Array>({
1183+
start: (controller) => void openTails.push(controller),
1184+
}),
1185+
{ headers: { 'Content-Type': 'text/event-stream' } }
1186+
)
1187+
})
1188+
const { getResult } = renderUseChatInChat(history.id, history)
1189+
await act(async () => {
1190+
void getResult().sendMessage('Keep going while I am away')
1191+
})
1192+
await act(async () => vi.advanceTimersByTimeAsync(100))
1193+
await act(async () => {
1194+
window.dispatchEvent(new Event('pageshow'))
1195+
await vi.advanceTimersByTimeAsync(100)
1196+
})
1197+
expect(openTails.length).toBeGreaterThan(0)
1198+
1199+
online = false
1200+
if (tailOutcome === 'failed') {
1201+
await act(async () => {
1202+
for (const tail of openTails.splice(0)) tail.error(new TypeError('network error'))
1203+
await vi.advanceTimersByTimeAsync(0)
1204+
})
1205+
for (let second = 0; second < 120 && failedReconnects < 6; second++) {
1206+
await act(async () => vi.advanceTimersByTimeAsync(1_000))
1207+
}
1208+
expect(failedReconnects).toBeGreaterThanOrEqual(6)
1209+
} else {
1210+
await act(async () => vi.advanceTimersByTimeAsync(20_000))
1211+
}
1212+
1213+
online = true
1214+
backOnline = true
1215+
await act(async () => {
1216+
window.dispatchEvent(new Event('online'))
1217+
await vi.advanceTimersByTimeAsync(500)
1218+
})
1219+
1220+
expect(tailOpenedAfterReturn).toBe(true)
1221+
expect(getResult().isSending).toBe(true)
1222+
expect(state.postBodies).toHaveLength(1)
1223+
} finally {
1224+
vi.useRealTimers()
1225+
}
1226+
}
1227+
)
1228+
1229+
/**
1230+
* The turn ends on the server while this surface's reader is silent (a stalled
1231+
* socket, or a recovery a later return superseded), so it never sees `complete`.
1232+
* The next return reads a chat with no running turn; it must resolve the stream
1233+
* it still shows as running instead of leaving the chat stuck on Stop.
1234+
*/
1235+
it('finishes a turn that ended while its reader was silent when the user returns', async () => {
1236+
let turnRunning = true
1237+
const history: MothershipChatHistory = {
1238+
id: 'chat-ended-while-silent',
1239+
mode: 'agent',
1240+
title: 'Ended while silent',
1241+
messages: [],
1242+
activeStreamId: null,
1243+
resources: [],
1244+
}
1245+
mockRequestJson.mockImplementation(() =>
1246+
Promise.resolve({
1247+
chat: {
1248+
...history,
1249+
activeStreamId: turnRunning ? (state.postBodies[0]?.userMessageId ?? null) : null,
1250+
},
1251+
})
1252+
)
1253+
state.postBehavior = 'accept'
1254+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1255+
const url = String(input)
1256+
if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init)
1257+
if (url.includes('batch=true')) {
1258+
return Response.json({
1259+
success: true,
1260+
events: [],
1261+
status: turnRunning ? 'streaming' : 'complete',
1262+
})
1263+
}
1264+
return new Response(new ReadableStream<Uint8Array>(), {
1265+
headers: { 'Content-Type': 'text/event-stream' },
1266+
})
1267+
})
1268+
const { getResult } = renderUseChatInChat(history.id, history)
1269+
await act(async () => {
1270+
void getResult().sendMessage('Finish while I am away')
1271+
})
1272+
await act(async () => {
1273+
window.dispatchEvent(new Event('pageshow'))
1274+
await sleep(100)
1275+
})
1276+
expect(getResult().isSending).toBe(true)
1277+
1278+
turnRunning = false
1279+
await act(async () => {
1280+
window.dispatchEvent(new Event('online'))
1281+
})
1282+
await waitFor(() => !getResult().isSending)
1283+
1284+
expect(state.postBodies).toHaveLength(1)
1285+
})
1286+
1287+
/**
1288+
* Before its POST is admitted a send shows as running, but the chat cannot list
1289+
* it yet. A return event in that window must leave the POST alone.
1290+
*/
1291+
it('does not abort a send still waiting for admission when the user returns', async () => {
1292+
const history: MothershipChatHistory = {
1293+
id: 'chat-pending-admission',
1294+
mode: 'agent',
1295+
title: 'Pending admission',
1296+
messages: [],
1297+
activeStreamId: null,
1298+
resources: [],
1299+
}
1300+
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
1301+
let postSignal: AbortSignal | undefined
1302+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1303+
if (String(input) === '/api/mothership/chat' && init?.method === 'POST') {
1304+
state.postBodies.push(JSON.parse(String(init.body)))
1305+
postSignal = init.signal ?? undefined
1306+
return new Promise<Response>(() => {})
1307+
}
1308+
return fetchStub(input, init)
1309+
})
1310+
const { getResult } = renderUseChatInChat(history.id, history)
1311+
await act(async () => {
1312+
void getResult().sendMessage('Still being admitted')
1313+
})
1314+
await waitFor(() => postSignal !== undefined)
1315+
1316+
await act(async () => {
1317+
window.dispatchEvent(new Event('online'))
1318+
await sleep(200)
1319+
})
1320+
1321+
expect(postSignal?.aborted).toBe(false)
1322+
expect(getResult().isSending).toBe(true)
1323+
})
1324+
1325+
/**
1326+
* The server admitted the send and finished its turn, but the POST's answer
1327+
* never arrived. Once the chat holds the message, a return resolves the turn
1328+
* rather than leaving the chat on Stop behind a POST that will not answer.
1329+
*/
1330+
it('finishes an admitted turn whose POST never answered when the user returns', async () => {
1331+
let admitted = false
1332+
const history: MothershipChatHistory = {
1333+
id: 'chat-admitted-unanswered',
1334+
mode: 'agent',
1335+
title: 'Admitted, unanswered',
1336+
messages: [],
1337+
activeStreamId: null,
1338+
resources: [],
1339+
}
1340+
mockRequestJson.mockImplementation(() => {
1341+
const userMessageId = state.postBodies[0]?.userMessageId
1342+
return Promise.resolve({
1343+
chat: {
1344+
...history,
1345+
messages:
1346+
admitted && userMessageId
1347+
? [
1348+
{ id: userMessageId, role: 'user', content: 'Answer lost' },
1349+
{ id: 'saved-answer', role: 'assistant', content: 'Done.' },
1350+
]
1351+
: [],
1352+
},
1353+
})
1354+
})
1355+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1356+
const url = String(input)
1357+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
1358+
state.postBodies.push(JSON.parse(String(init.body)))
1359+
return new Promise<Response>(() => {})
1360+
}
1361+
if (url.includes('/api/mothership/chat/stream')) {
1362+
return Response.json({ success: true, events: [], status: 'complete' })
1363+
}
1364+
return fetchStub(input, init)
1365+
})
1366+
const { getResult } = renderUseChatInChat(history.id, history)
1367+
await act(async () => {
1368+
void getResult().sendMessage('Answer lost')
1369+
})
1370+
await waitFor(() => state.postBodies.length === 1)
1371+
1372+
admitted = true
1373+
await act(async () => {
1374+
window.dispatchEvent(new Event('online'))
1375+
})
1376+
await waitFor(() => !getResult().isSending)
1377+
1378+
expect(state.postBodies).toHaveLength(1)
1379+
})
1380+
11391381
it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => {
11401382
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
11411383
try {

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

Lines changed: 28 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3036,12 +3036,17 @@ export function useChat(
30363036
if (!chatId) return
30373037

30383038
const subjectKey = buildRecoverySubjectKey(startingChatId, startingSelectedChatId)
3039+
/* A return signal supersedes a recovery already in flight, as it supersedes the
3040+
send's own reader: that recovery's tail may have gone silent while the tab was
3041+
away or offline, or it may be sleeping out a reconnect backoff, and either
3042+
would hold the stream for up to the idle timeout or the backoff. */
30393043
const existingRecovery = activeStreamReturnRecoveryRef.current
3040-
if (existingRecovery?.subjectKey === subjectKey) {
3041-
return existingRecovery.promise
3042-
}
30433044
if (existingRecovery) {
3044-
existingRecovery.controller.abort('replaced_by_new_recovery_subject')
3045+
existingRecovery.controller.abort(
3046+
existingRecovery.subjectKey === subjectKey
3047+
? 'superseded_by_return'
3048+
: 'replaced_by_new_recovery_subject'
3049+
)
30453050
activeStreamReturnRecoveryRef.current = null
30463051
}
30473052

@@ -3059,8 +3064,26 @@ export function useChat(
30593064
const fallbackStreamId =
30603065
streamIdRef.current ?? activeTurnRef.current?.userMessageId ?? cached?.activeStreamId
30613066
const loadedStream = await getActiveStreamIdForChat(chatId, recoveryController.signal)
3067+
/* The chat no longer lists a running turn, but this surface is still showing
3068+
one: it ended while nothing here was listening (its reader went silent, or a
3069+
recovery it superseded was attached). Resolve that stream instead of
3070+
leaving it running: its terminal state replays the rest and finalizes. A
3071+
send whose POST has not answered yet is such a stream only once the loaded
3072+
chat holds its message (the server admitted it, and the answer is lost);
3073+
until then it may still be on its way, and recovering it would abort it. */
3074+
const pendingAdmission = pendingChatAdmissionRef.current
3075+
const admitted =
3076+
!pendingAdmission ||
3077+
(loadedStream.loaded &&
3078+
queryClient
3079+
.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
3080+
?.messages.some((message) => message.id === pendingAdmission.userMessageId) === true)
3081+
const locallyRunningStreamId =
3082+
sendingRef.current && admitted
3083+
? (streamIdRef.current ?? activeTurnRef.current?.userMessageId)
3084+
: undefined
30623085
const streamId = loadedStream.loaded
3063-
? (loadedStream.streamId ?? undefined)
3086+
? (loadedStream.streamId ?? locallyRunningStreamId)
30643087
: fallbackStreamId
30653088
if (
30663089
!isSameRecoverySubject() ||

0 commit comments

Comments
 (0)