Skip to content

Commit 11e9bbd

Browse files
committed
fix(mothership): bound a silent replay, never skip past a trimmed cursor, and re-sync an expired ring
- The worker replay is bounded like a stream leg: no response headers, or no bytes including keepalives, for the idle timeout ends it so the reader re-attaches. - A ring read that starts after the reader's next event (the ring trimmed its head between the gap check and the read) is never delivered; the reader re-attaches and re-syncs from the log, in both the live tail and batch reads. - An empty ring serves only a reader starting from cursor 0; a live run whose buffer expired under a reader re-syncs from the log, and a finished one answers its terminal since its transcript is persisted. - A leg that delivers events after a failure starts a fresh 30 s reachable window; its three retries still refill only after five minutes of delivered events. - The two new integration suites close their worker server and restore env even when they are skipped.
1 parent 4e115aa commit 11e9bbd

10 files changed

Lines changed: 303 additions & 29 deletions

File tree

‎apps/sim/app/api/copilot/chat/stream/route.test.ts‎

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,4 +219,54 @@ describe('copilot chat stream replay route', () => {
219219
)
220220
trace.disable()
221221
})
222+
223+
it('never delivers a ring read that starts past the reader cursor, and ends without a terminal', async () => {
224+
getLatestRunForStream.mockResolvedValue({
225+
status: 'active',
226+
executionId: 'exec-1',
227+
id: 'run-1',
228+
})
229+
readEvents.mockResolvedValue([
230+
{
231+
stream: { streamId: 'stream-1', cursor: '5' },
232+
seq: 5,
233+
trace: { requestId: 'req-1' },
234+
type: MothershipStreamV1EventType.text,
235+
payload: { channel: 'assistant', text: 'the middle of the turn' },
236+
},
237+
])
238+
239+
const response = await GET(
240+
new NextRequest('http://localhost:3000/api/copilot/chat/stream?streamId=stream-1&after=0')
241+
)
242+
const text = (await readAllChunks(response)).join('')
243+
244+
expect(text).not.toContain('the middle of the turn')
245+
expect(text).not.toContain(`"type":"${MothershipStreamV1EventType.complete}"`)
246+
})
247+
248+
it('serves a batch read that starts past the reader cursor no events', async () => {
249+
getLatestRunForStream.mockResolvedValue({
250+
status: 'active',
251+
executionId: 'exec-1',
252+
id: 'run-1',
253+
})
254+
readEvents.mockResolvedValue([
255+
{
256+
stream: { streamId: 'stream-1', cursor: '5' },
257+
seq: 5,
258+
trace: { requestId: 'req-1' },
259+
type: MothershipStreamV1EventType.text,
260+
payload: { channel: 'assistant', text: 'the middle of the turn' },
261+
},
262+
])
263+
264+
const response = await GET(
265+
new NextRequest(
266+
'http://localhost:3000/api/copilot/chat/stream?streamId=stream-1&after=0&batch=true'
267+
)
268+
)
269+
270+
await expect(response.json()).resolves.toMatchObject({ success: true, events: [] })
271+
})
222272
})

‎apps/sim/app/api/copilot/chat/stream/route.ts‎

Lines changed: 37 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,15 @@ const RING_CHECK_EVERY_BUSY_POLLS = 8
5959
*/
6060
const MAX_STREAM_MS = 60 * 60 * 1000 - 60_000
6161

62+
/**
63+
* Whether ring events read after `cursor` start right after it. The ring can trim its
64+
* head between a gap check and the read, and a read that starts later would silently
65+
* skip part of the turn.
66+
*/
67+
function startsAfterCursor(events: readonly { seq: number }[], cursor: string): boolean {
68+
return events.length === 0 || events[0].seq <= Number(cursor || '0') + 1
69+
}
70+
6271
function extractCanonicalRequestId(value: unknown): string {
6372
return typeof value === 'string' && value.length > 0 ? value : ''
6473
}
@@ -253,8 +262,10 @@ async function handleResumeRequestBody({
253262
return []
254263
}),
255264
])
256-
// A reader the ring cannot serve is re-synced from the worker log by the live tail.
257-
const batchEvents = fromLog || gap ? [] : events.map(toStreamBatchEvent)
265+
// A reader the ring cannot serve, or whose next event it trimmed after the gap check,
266+
// is re-synced from the worker log by the live tail.
267+
const batchEvents =
268+
fromLog || gap || !startsAfterCursor(events, afterSeq) ? [] : events.map(toStreamBatchEvent)
258269
logger.info('[Resume] Batch response', {
259270
streamId,
260271
afterCursor: afterSeq,
@@ -287,9 +298,12 @@ async function handleResumeRequestBody({
287298
* no shared position to join on. The header tells the client to rebuild the turn
288299
* from an empty response, since the replay's cursors restart at 1.
289300
*/
290-
const gap = fromLog
301+
const ringGap = fromLog
291302
? null
292303
: await findReplayGap(streamId, afterCursor || '0', extractRunRequestId(run))
304+
// A finished run whose buffer expired answers its terminal; its transcript is persisted.
305+
const gap =
306+
ringGap && !(ringGap.latestSeq <= 0 && isTerminalStreamStatus(run.status)) ? ringGap : null
293307
const resyncFromLog = fromLog || gap !== null
294308
let replayBody: ReadableStream<Uint8Array> | null = null
295309
/** Releases the worker's replay once this response ends; the request signal may never fire. */
@@ -380,8 +394,13 @@ async function handleResumeRequestBody({
380394
}
381395
request.signal.addEventListener('abort', abortListener, { once: true })
382396

383-
const flushEvents = async (): Promise<number> => {
397+
/** Delivers the ring's events after the cursor, or returns null if it trimmed the next one. */
398+
const flushEvents = async (): Promise<number | null> => {
384399
const events = await readEvents(streamId, cursor)
400+
if (!startsAfterCursor(events, cursor)) {
401+
logger.warn('Replay ring trimmed past a reader cursor', { streamId, cursor })
402+
return null
403+
}
385404
if (events.length > 0) {
386405
logger.debug('[Resume] Flushing events', {
387406
streamId,
@@ -482,6 +501,7 @@ async function handleResumeRequestBody({
482501
}
483502

484503
let lastFlushed = await flushEvents()
504+
if (lastFlushed === null) return
485505

486506
let pollDelayMs = POLL_INTERVAL_MS
487507
while (!controllerClosed && Date.now() - startTime < MAX_STREAM_MS) {
@@ -493,14 +513,6 @@ async function handleResumeRequestBody({
493513
})
494514
return null
495515
})
496-
// The ring lost its head or restarted under this tail; the re-attach re-syncs. Only a
497-
// quiet ring can have restarted or have a dead controller that recovery re-sends
498-
// under, so a busy tail checks every few polls rather than on each one.
499-
const checkRing = lastFlushed === 0 || pollIterations % RING_CHECK_EVERY_BUSY_POLLS === 0
500-
if (checkRing && !ringCanServe(await readRingPosition(streamId, cursor))) {
501-
logger.warn('Replay ring can no longer serve a live tail', { streamId, cursor })
502-
break
503-
}
504516
if (!currentRun) {
505517
emitTerminalIfMissing(MothershipStreamV1CompletionStatus.error, {
506518
message: 'The stream could not be recovered because its run metadata is unavailable.',
@@ -509,10 +521,23 @@ async function handleResumeRequestBody({
509521
})
510522
break
511523
}
524+
// The ring lost its head, restarted or expired under this live tail; the re-attach
525+
// re-syncs, and a finished run answers its terminal instead. Only a quiet ring can
526+
// restart or be re-sent into by a recovery, so a busy tail checks every few polls.
527+
const checkRing = lastFlushed === 0 || pollIterations % RING_CHECK_EVERY_BUSY_POLLS === 0
528+
if (
529+
checkRing &&
530+
!isTerminalStreamStatus(currentRun.status) &&
531+
!ringCanServe(await readRingPosition(streamId, cursor))
532+
) {
533+
logger.warn('Replay ring can no longer serve a live tail', { streamId, cursor })
534+
break
535+
}
512536

513537
currentRequestId = extractRunRequestId(currentRun) || currentRequestId
514538

515539
const flushed = await flushEvents()
540+
if (flushed === null) break
516541
lastFlushed = flushed
517542
/* Adaptive tail: 4 Hz only while events are actually flowing; a quiet stream
518543
decays toward the cap so an attached client doesn't hammer Postgres + Redis

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

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -216,4 +216,18 @@ describe('stream recovery budget', () => {
216216
vi.advanceTimersByTime(30 * 60_000)
217217
expect(retry.nextDelay(error)).toBeNull()
218218
})
219+
220+
it('retries a leg that fails again minutes after it re-attached and delivered events', () => {
221+
vi.useFakeTimers()
222+
const error = new WorkerStreamInterruptedError(new Error('socket closed'))
223+
const retry = new StreamRetryWindow()
224+
const first = retry.nextDelay(error)
225+
expect(first).not.toBeNull()
226+
vi.advanceTimersByTime(first ?? 0)
227+
for (let second = 0; second < 120; second += 10) {
228+
retry.recovered()
229+
vi.advanceTimersByTime(10_000)
230+
}
231+
expect(retry.nextDelay(error)).not.toBeNull()
232+
})
219233
})

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

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,8 @@ const HEALTHY_STREAM_REPLENISH_MS = 5 * 60_000
3232
* Recovery is bounded independently of the healthy leg's lifetime, by two
3333
* budgets that never share state: an unreachable worker gets a two-minute
3434
* window from the moment it stopped answering, and any failure of a worker that
35-
* did answer gets three retries within 30 s, replenished only after
35+
* did answer gets three retries, each burst of them within 30 s of its first
36+
* failure, replenished only after
3637
* {@link HEALTHY_STREAM_REPLENISH_MS} of healthy streaming. A leg has no deadline
3738
* unless the caller sets one.
3839
*/
@@ -63,9 +64,13 @@ export class StreamRetryWindow {
6364
return remaining
6465
}
6566

66-
/** The worker delivered an event, so a later loss of it starts a fresh unreachable window. */
67+
/**
68+
* The worker delivered an event: a later loss starts a fresh unreachable window, and
69+
* a fresh 30 s reachable window. Only the three reachable retries carry over.
70+
*/
6771
recovered(): void {
6872
this.resetUnreachable()
73+
this.firstFailureAt = undefined
6974
this.lastEventAt = Date.now()
7075
this.streamingSince ??= this.lastEventAt
7176
}

‎apps/sim/lib/mothership/request/session/recovery.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,11 @@ vi.mock('./buffer', () => ({
1212
readEvents,
1313
}))
1414

15-
import { findReplayGap, replayGapTerminal } from '@/lib/mothership/request/session/recovery'
15+
import {
16+
findReplayGap,
17+
replayGapTerminal,
18+
ringCanServe,
19+
} from '@/lib/mothership/request/session/recovery'
1620

1721
describe('replay gap', () => {
1822
it('uses the latest buffered request id when run metadata is missing it', async () => {
@@ -33,4 +37,9 @@ describe('replay gap', () => {
3337
expect(result?.envelopes[0].trace.requestId).toBe('req-live-123')
3438
expect(result?.envelopes[1].trace.requestId).toBe('req-live-123')
3539
})
40+
41+
it('cannot serve a reader that holds a cursor from an empty ring', () => {
42+
expect(ringCanServe({ requestedAfterSeq: 12, oldestSeq: 0, latestSeq: 0 })).toBe(false)
43+
expect(ringCanServe({ requestedAfterSeq: 0, oldestSeq: 0, latestSeq: 0 })).toBe(true)
44+
})
3645
})

‎apps/sim/lib/mothership/request/session/recovery.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,10 +52,11 @@ export async function readRingPosition(
5252
* head: the events before its oldest are gone, and a cursor that was served from the
5353
* worker's log instead is not a position in the ring, so no cursor is trusted. Nor can
5454
* it serve a cursor ahead of its latest event, which only a buffer whose numbering
55-
* restarted after it expired produces.
55+
* restarted after it expired produces, nor any cursor from a buffer that expired.
5656
*/
5757
export function ringCanServe({ requestedAfterSeq, oldestSeq, latestSeq }: RingPosition): boolean {
58-
return latestSeq <= 0 || (startsAtReplayHead(oldestSeq) && requestedAfterSeq <= latestSeq)
58+
if (latestSeq <= 0) return requestedAfterSeq <= 0
59+
return startsAtReplayHead(oldestSeq) && requestedAfterSeq <= latestSeq
5960
}
6061

6162
/** The ring's position when it cannot serve `afterCursor` (see {@link ringCanServe}). */

‎apps/sim/lib/mothership/request/session/replay-gap.integration.ts‎

Lines changed: 50 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,15 @@ function fullResponse(streamId: string) {
134134
]
135135
}
136136

137+
/** Runs whether or not the suite does, so a skipped suite never leaks the worker or env. */
138+
afterAll(async () => {
139+
await new Promise<void>((resolve) => worker.server.close(() => resolve()))
140+
for (const [key, value] of Object.entries(inheritedEnv)) {
141+
if (value === undefined) delete process.env[key]
142+
else process.env[key] = value
143+
}
144+
})
145+
137146
describe.runIf(Boolean(redisUrl))('reconnects past the replay ring', () => {
138147
beforeAll(async () => {
139148
const now = new Date()
@@ -177,11 +186,6 @@ describe.runIf(Boolean(redisUrl))('reconnects past the replay ring', () => {
177186
await db.delete(workspace).where(eq(workspace.id, workspaceId))
178187
await db.delete(user).where(eq(user.id, userId))
179188
await closeRedisConnection()
180-
await new Promise<void>((resolve) => worker.server.close(() => resolve()))
181-
for (const [key, value] of Object.entries(inheritedEnv)) {
182-
if (value === undefined) delete process.env[key]
183-
else process.env[key] = value
184-
}
185189
})
186190

187191
it.each([
@@ -375,6 +379,47 @@ describe.runIf(Boolean(redisUrl))('reconnects past the replay ring', () => {
375379
expect(response.headers.get(MOTHERSHIP_STREAM_REPLAY_HEADER)).toBeNull()
376380
})
377381

382+
it('re-syncs a live run whose buffer expired under a reader cursor', async () => {
383+
const streamId = generateId()
384+
await db.insert(copilotRuns).values({
385+
id: generateId(),
386+
executionId: generateId(),
387+
chatId,
388+
userId,
389+
workspaceId,
390+
streamId,
391+
})
392+
worker.reply.frames = fullResponse(streamId)
393+
394+
const response = await reconnect(streamId, '6')
395+
const frames = dataFrames(await response.text())
396+
397+
expect(response.headers.get(MOTHERSHIP_STREAM_REPLAY_HEADER)).toBe('log')
398+
expect(frames.map((frame) => frame.type)).toEqual(['session', 'text', 'complete'])
399+
})
400+
401+
it('answers a finished run whose buffer expired with its terminal, not a replay', async () => {
402+
const streamId = generateId()
403+
await db.insert(copilotRuns).values({
404+
id: generateId(),
405+
executionId: generateId(),
406+
chatId,
407+
userId,
408+
workspaceId,
409+
streamId,
410+
status: 'complete',
411+
})
412+
413+
const response = await reconnect(streamId, '6')
414+
const frames = dataFrames(await response.text())
415+
416+
expect(response.headers.get(MOTHERSHIP_STREAM_REPLAY_HEADER)).toBeNull()
417+
expect(frames.map((frame) => [frame.type, frame.payload.status])).toEqual([
418+
['complete', 'complete'],
419+
])
420+
expect(worker.requests).toEqual([])
421+
})
422+
378423
it('serves no ring events to a batch read the ring can no longer serve', async () => {
379424
const { streamId } = await liveRunWithTrimmedRing()
380425

0 commit comments

Comments
 (0)