Skip to content

Commit ddcf0c8

Browse files
committed
fix(mothership): bound stream reconnection attempts
1 parent a2a38b0 commit ddcf0c8

3 files changed

Lines changed: 85 additions & 15 deletions

File tree

‎apps/sim/lib/mothership/request/lifecycle/run.test.ts‎

Lines changed: 45 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1992,12 +1992,12 @@ describe('runCopilotLifecycle', () => {
19921992
}
19931993
})
19941994

1995-
it('keeps the same request and partial response across a longer worker outage', async () => {
1995+
it('keeps the same request and partial response through the allowed reconnects', async () => {
19961996
vi.useFakeTimers()
19971997
try {
19981998
const bodies: Record<string, unknown>[] = []
19991999
const headers: Headers[] = []
2000-
for (let attempt = 0; attempt < 4; attempt++) {
2000+
for (let attempt = 0; attempt < 3; attempt++) {
20012001
mockRunStreamLoop.mockImplementationOnce(
20022002
async (_url, request, context: StreamingContext) => {
20032003
bodies.push(JSON.parse(String(request.body)))
@@ -2041,19 +2041,58 @@ describe('runCopilotLifecycle', () => {
20412041
content: 'Saved partial answer recovered',
20422042
})
20432043
expect(result.errors).toBeUndefined()
2044-
expect(bodies).toHaveLength(5)
2045-
expect(bodies.map((body) => body.messageId)).toEqual(Array(5).fill('outage-stream'))
2044+
expect(bodies).toHaveLength(4)
2045+
expect(bodies.map((body) => body.messageId)).toEqual(Array(4).fill('outage-stream'))
20462046
expect(headers.map((header) => header.get('X-Sim-Request-ID'))).toEqual(
2047-
Array(5).fill('outage-request')
2047+
Array(4).fill('outage-request')
20482048
)
20492049
expect(bodies.slice(1).map((body) => body.receivedTextChars)).toEqual(
2050-
Array(4).fill('Saved partial answer'.length)
2050+
Array(3).fill('Saved partial answer'.length)
20512051
)
20522052
} finally {
20532053
vi.useRealTimers()
20542054
}
20552055
})
20562056

2057+
it('ends an empty-stream outage after the initial attempt and three reconnects', async () => {
2058+
vi.useFakeTimers()
2059+
try {
2060+
for (let attempt = 0; attempt < 4; attempt++) {
2061+
mockRunStreamLoop.mockImplementationOnce(
2062+
async (_url, _request, context: StreamingContext) => {
2063+
context.errors.push(STREAM_ENDED_WITHOUT_TERMINAL_MESSAGE)
2064+
throw new StreamEndedWithoutTerminalError('/api/mothership')
2065+
}
2066+
)
2067+
}
2068+
const pending = runCopilotLifecycle(
2069+
{ message: 'hello', messageId: 'bounded-outage' },
2070+
{
2071+
userId: 'user-1',
2072+
workspaceId: 'ws-1',
2073+
chatId: 'chat-1',
2074+
executionId: 'exec-1',
2075+
runId: 'run-1',
2076+
executionContext: {
2077+
userId: 'user-1',
2078+
workflowId: '',
2079+
workspaceId: 'ws-1',
2080+
chatId: 'chat-1',
2081+
},
2082+
}
2083+
)
2084+
await vi.advanceTimersByTimeAsync(30_000)
2085+
expect(await pending).toMatchObject({
2086+
success: false,
2087+
cancelled: false,
2088+
errors: [STREAM_ENDED_WITHOUT_TERMINAL_MESSAGE],
2089+
})
2090+
expect(mockRunStreamLoop).toHaveBeenCalledTimes(4)
2091+
} finally {
2092+
vi.useRealTimers()
2093+
}
2094+
})
2095+
20572096
it('honors Stop while waiting to reconnect to the worker', async () => {
20582097
vi.useFakeTimers()
20592098
try {

‎apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts‎

Lines changed: 32 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,45 @@
11
import { afterEach, describe, expect, it, vi } from 'vitest'
2-
import { CopilotBackendError } from '@/lib/mothership/request/go/stream'
2+
import {
3+
CopilotBackendError,
4+
StreamEndedWithoutTerminalError,
5+
} from '@/lib/mothership/request/go/stream'
36
import { StreamRetryWindow } from '@/lib/mothership/request/lifecycle/stream-retry'
47

58
afterEach(() => vi.useRealTimers())
69

710
describe('stream recovery budget', () => {
8-
it('survives multiple unavailable connections without resetting its deadline', () => {
11+
it.each([
12+
new TypeError('fetch failed'),
13+
new StreamEndedWithoutTerminalError('/api/mothership'),
14+
new CopilotBackendError('Unavailable', { status: 503 }),
15+
])('stops after three retries despite a long task budget: %s', (error) => {
916
vi.useFakeTimers()
10-
const retry = new StreamRetryWindow(120_000)
11-
for (let index = 0; index < 6; index++) {
12-
const delay = retry.nextDelay(new TypeError('fetch failed'))
17+
const retry = new StreamRetryWindow()
18+
for (let index = 0; index < 3; index++) {
19+
const delay = retry.nextDelay(error)
1320
expect(delay).not.toBeNull()
1421
vi.advanceTimersByTime(delay ?? 0)
1522
}
16-
expect(retry.attempt).toBe(6)
17-
expect(retry.remainingMs()).toBeLessThan(120_000)
23+
expect(retry.nextDelay(error)).toBeNull()
24+
expect(retry.attempt).toBe(3)
25+
expect(retry.remainingMs()).toBeGreaterThan(3_500_000)
26+
})
27+
28+
it('bounds the recovery period from the first failure without shortening healthy work', () => {
29+
vi.useFakeTimers()
30+
const retry = new StreamRetryWindow()
31+
vi.advanceTimersByTime(600_000)
32+
expect(retry.nextDelay(new TypeError('fetch failed'))).not.toBeNull()
33+
vi.advanceTimersByTime(30_000)
34+
expect(retry.nextDelay(new TypeError('fetch failed'))).toBeNull()
35+
expect(retry.remainingMs()).toBe(2_970_000)
36+
})
37+
38+
it('never extends the original execution deadline', () => {
39+
vi.useFakeTimers()
40+
const retry = new StreamRetryWindow(120_000)
41+
vi.advanceTimersByTime(119_999)
42+
expect(retry.nextDelay(new TypeError('fetch failed'))).toBeNull()
1843
vi.advanceTimersByTime(120_000)
1944
expect(retry.nextDelay(new TypeError('fetch failed'))).toBeNull()
2045
expect(() => retry.remainingMs()).toThrow('could not be restored')

‎apps/sim/lib/mothership/request/lifecycle/stream-retry.ts‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,9 +6,13 @@ import {
66
StreamEndedWithoutTerminalError,
77
} from '@/lib/mothership/request/go/stream'
88

9-
/** A connection failure leaves the durable run unresolved; its leg keeps the existing time budget. */
9+
const MAX_STREAM_RETRIES = 3
10+
const STREAM_RECOVERY_WINDOW_MS = 30_000
11+
12+
/** Recovery is bounded independently of the healthy run's execution budget. */
1013
export class StreamRetryWindow {
1114
private readonly deadline: number
15+
private recoveryDeadline?: number
1216
attempt = 0
1317

1418
constructor(timeoutMs = ORCHESTRATION_TIMEOUT_MS) {
@@ -24,8 +28,10 @@ export class StreamRetryWindow {
2428

2529
nextDelay(error: unknown, signal?: AbortSignal): number | null {
2630
if (signal?.aborted || !isRetryableStreamError(error)) return null
31+
this.recoveryDeadline ??= Date.now() + STREAM_RECOVERY_WINDOW_MS
32+
if (this.attempt >= MAX_STREAM_RETRIES) return null
2733
const delay = backoffWithJitter(this.attempt + 1, null, { baseMs: 250, maxMs: 5_000 })
28-
if (Date.now() + delay >= this.deadline) return null
34+
if (Date.now() + delay >= Math.min(this.deadline, this.recoveryDeadline)) return null
2935
this.attempt++
3036
return delay
3137
}

0 commit comments

Comments
 (0)