Skip to content

Commit be26984

Browse files
committed
fix(mothership): renew the Chat reconnect budget after a tail delivers events
A reconnect attempt whose re-attached tail delivered new events now restarts the retry budget at the base delay, so separate network drops hours apart in a long turn no longer add up to the ten-attempt exhaustion.
1 parent f376c31 commit be26984

2 files changed

Lines changed: 78 additions & 1 deletion

File tree

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

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1089,6 +1089,74 @@ describe('useChat remount send recovery', () => {
10891089
}
10901090
})
10911091

1092+
it('keeps re-attaching a long turn whose tails deliver events between separate network failures', async () => {
1093+
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
1094+
try {
1095+
let tails = 0
1096+
const history: MothershipChatHistory = {
1097+
id: 'chat-long-turn',
1098+
mode: 'agent',
1099+
title: 'Long turn',
1100+
messages: [],
1101+
activeStreamId: null,
1102+
resources: [],
1103+
}
1104+
mockRequestJson.mockImplementation(() =>
1105+
Promise.resolve({
1106+
chat: { ...history, activeStreamId: state.postBodies[0]?.userMessageId ?? null },
1107+
})
1108+
)
1109+
state.postBehavior = 'accept'
1110+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1111+
const url = String(input)
1112+
if (!url.includes('/api/mothership/chat/stream')) return fetchStub(input, init)
1113+
if (url.includes('batch=true')) {
1114+
return Response.json({ success: true, events: [], status: 'streaming' })
1115+
}
1116+
tails++
1117+
const streamId = state.postBodies[0]?.userMessageId ?? ''
1118+
const event: MothershipStreamV1EventEnvelope = {
1119+
v: 1,
1120+
seq: tails,
1121+
ts: new Date().toISOString(),
1122+
type: 'text',
1123+
stream: { streamId, cursor: String(tails) },
1124+
payload: { channel: 'assistant', text: `part ${tails} ` },
1125+
}
1126+
return new Response(
1127+
new ReadableStream<Uint8Array>({
1128+
start(controller) {
1129+
controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(event)}\n\n`))
1130+
},
1131+
pull(controller) {
1132+
controller.error(new TypeError('network error'))
1133+
},
1134+
}),
1135+
{ headers: { 'Content-Type': 'text/event-stream' } }
1136+
)
1137+
})
1138+
const { getResult } = renderUseChatInChat(history.id, history)
1139+
await act(async () => {
1140+
void getResult().sendMessage('Keep going for hours')
1141+
})
1142+
const errors = new Set<string>()
1143+
let seconds = 0
1144+
for (; seconds < 600 && tails < 15; seconds++) {
1145+
await act(async () => vi.advanceTimersByTimeAsync(1_000))
1146+
const error = getResult().error
1147+
if (error) errors.add(error)
1148+
}
1149+
1150+
expect(tails).toBeGreaterThanOrEqual(15)
1151+
/* Each failure after a tail that delivered events retries at the base delay. */
1152+
expect(seconds).toBeLessThan(60)
1153+
expect([...errors]).toEqual([])
1154+
expect(getResult().isSending).toBe(true)
1155+
} finally {
1156+
vi.useRealTimers()
1157+
}
1158+
})
1159+
10921160
it('sends a queued correction after stopping with more than 10 MiB of tool input', async () => {
10931161
state.postBehavior = 'tool'
10941162
state.toolInputPadding = 'x'.repeat(11 * 1024 * 1024)

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

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2765,8 +2765,16 @@ export function useChat(
27652765
abortControllerRef.current?.signal.aborted === true ||
27662766
shouldContinue?.() === false
27672767

2768-
for (let attempt = 0; attempt <= MAX_RECONNECT_ATTEMPTS; attempt++) {
2768+
/**
2769+
* An attempt whose tail delivered new events re-attached successfully, so
2770+
* the failure after it starts a fresh budget at the base delay. Only
2771+
* failures without progress count toward exhaustion, which keeps separate
2772+
* network drops hours apart in a long turn from adding up.
2773+
*/
2774+
let attempt = 0
2775+
while (attempt <= MAX_RECONNECT_ATTEMPTS) {
27692776
if (isStaleReconnect()) return true
2777+
const cursorBeforeAttempt = lastCursorRef.current
27702778

27712779
if (attempt > 0) {
27722780
const delayMs = Math.min(
@@ -2868,6 +2876,7 @@ export function useChat(
28682876
error: toError(err).message,
28692877
})
28702878
}
2879+
attempt = lastCursorRef.current !== cursorBeforeAttempt ? 1 : attempt + 1
28712880
}
28722881

28732882
logger.error('All reconnect attempts exhausted', {

0 commit comments

Comments
 (0)