Skip to content

Commit e368316

Browse files
committed
fix(mothership): hand off a failed replay append again, bounded by the recovery budget
Ending the turn on any leased append error that was not a lost lease turned a short Redis blip (a failover READONLY or a connection reset that outlasts the append retries) into a failed Chat turn and a worker stop. Before, that turn was handed off and a successor completed it. - A failed append that is not a budget refusal hands off again, as on staging. The per-run recovery budget bounds a persistent failure: MAX_RECOVERY_ATTEMPTS takeovers with backoff, and then the run ends as stream_recovery_exhausted. - Drop StreamPersistenceFailedError and stream_persistence_failed. The writer only latches budget refusals again, and finalizeAsError surfaces a failed terminal append as it did before. - Keep the change that an unreadable lease is not a hand-off. Its test now fails a lease read alone: a corrupted lock key also fails the fenced append, which is now a hand-off. - Storm suite: a single failed append hands off and the next controller completes the turn. A persistent append failure exhausts the budget, the same as a lease lost with no successor.
1 parent b154dfd commit e368316

6 files changed

Lines changed: 169 additions & 210 deletions

File tree

‎apps/sim/lib/mothership/request/lifecycle/start.ts‎

Lines changed: 12 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -62,11 +62,7 @@ import {
6262
} from '@/lib/mothership/request/session/controller-lease'
6363
import { StreamReplayBudgetExhaustedError } from '@/lib/mothership/request/session/replay-budget'
6464
import { SSE_RESPONSE_HEADERS } from '@/lib/mothership/request/session/sse'
65-
import {
66-
StreamPersistenceFailedError,
67-
type StreamTurnFailure,
68-
turnFailure,
69-
} from '@/lib/mothership/request/session/turn-failure'
65+
import { type StreamTurnFailure, turnFailure } from '@/lib/mothership/request/session/turn-failure'
7066
import { TraceCollector } from '@/lib/mothership/request/trace'
7167
import type { OrchestratorResult } from '@/lib/mothership/request/types'
7268
import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url'
@@ -448,17 +444,18 @@ export function createSSEStream(params: StreamingOrchestrationParams): ReadableS
448444
await publisher.publish(event)
449445
} catch (error) {
450446
/*
451-
Only a lost lease hands the turn to a successor. A refusal, or a
452-
Redis error that outlasted its retries, is terminal: leaving it
453-
recoverable made each replacement re-receive and fail the same
454-
event, re-POSTing the run on every reconnect poll.
447+
A refused write is a terminal failure, not a handoff: leaving it
448+
recoverable made each replacement re-receive and re-refuse the same
449+
event. Any other failure to persist means this controller can no
450+
longer prove ownership of the replay, so a successor takes over; the
451+
run's recovery budget bounds how often that repeats.
455452
*/
456-
const failure =
457-
error instanceof StreamControllerSupersededError
458-
? error
459-
: (turnFailure(error) ?? new StreamPersistenceFailedError(error))
460-
if (!abortController.signal.aborted) abortController.abort(failure)
461-
throw failure
453+
if (!abortController.signal.aborted) {
454+
abortController.abort(
455+
turnFailure(error) ?? new StreamControllerSupersededError()
456+
)
457+
}
458+
throw error
462459
}
463460
},
464461
onAbortObserved: (reason) => {

‎apps/sim/lib/mothership/request/session/recovery-storm.integration.ts‎

Lines changed: 47 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,8 @@
22
* How often reconnecting to an orphaned Chat run re-POSTs it to the worker, against real
33
* Redis and PostgreSQL. The production reconnect route, stream recovery, chat lifecycle
44
* and finalization run unmodified; a local HTTP server stands in for the worker. Faults
5-
* are injected at one seam, the leased replay append: an append that fails for a reason
6-
* other than a lost lease, or a lease lost with no successor to take over.
5+
* are injected at one seam, the leased replay append: an append that fails once or every
6+
* time, or a lease lost with no successor to take over.
77
*/
88
import { authMock, authMockFns } from '@sim/testing/mocks/auth.mock'
99
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
@@ -15,7 +15,10 @@ const { redisUrl, inheritedEnv, worker, faults } = await vi.hoisted(async () =>
1515
/** Which leased appends fail, and how; `undefined` lets every append through. */
1616
append: undefined as
1717
| undefined
18-
| { frames: 'any_tool' | 'tool_result'; effect: 'throw' | 'lose_lease' },
18+
| {
19+
frames: 'any_tool' | 'tool_result'
20+
effect: 'throw' | 'throw_once' | 'lose_lease'
21+
},
1922
}
2023
const worker = {
2124
posts: [] as Array<{ path: string; at: number }>,
@@ -76,6 +79,11 @@ vi.mock('@/lib/mothership/request/session/buffer', async (importOriginal) => {
7679
(envelope.payload as { phase?: string }).phase === 'result')
7780
)
7881
if (hit && fault.effect === 'throw') throw new Error('simulated Redis write failure')
82+
/** A Redis blip that outlasts the append retries, after which Redis is healthy again. */
83+
if (hit && fault.effect === 'throw_once') {
84+
faults.append = undefined
85+
throw new Error('simulated transient Redis write failure')
86+
}
7987
/** The lock expires under a live controller, and nobody else holds it. */
8088
if (hit && fault.effect === 'lose_lease') await getRedisClient()!.del(lease.key)
8189
return actual.appendEvents(...args)
@@ -105,10 +113,7 @@ import { isTerminalStreamStatus } from '@/lib/mothership/request/session'
105113
import { appendEvents } from '@/lib/mothership/request/session/buffer'
106114
import { chatStreamLockKey } from '@/lib/mothership/request/session/controller-lease'
107115
import { createEvent } from '@/lib/mothership/request/session/event'
108-
import {
109-
StreamPersistenceFailedError,
110-
StreamRecoveryExhaustedError,
111-
} from '@/lib/mothership/request/session/turn-failure'
116+
import { StreamRecoveryExhaustedError } from '@/lib/mothership/request/session/turn-failure'
112117
import { GET as streamGET } from '@/app/api/copilot/chat/stream/route'
113118

114119
const userId = generateId()
@@ -331,56 +336,57 @@ describe.runIf(Boolean(redisUrl))('reconnecting to an orphaned Chat run', () =>
331336
30_000
332337
)
333338

339+
it('hands off a turn whose append fails once, and the next controller completes it', async () => {
340+
worker.mode = 'frames'
341+
const run = await orphanedRun()
342+
faults.append = { frames: 'any_tool', effect: 'throw_once' }
343+
try {
344+
const { turns, stops } = await reconnect(run, { windowMs: 15_000 })
345+
346+
expect(turns).toHaveLength(2)
347+
expect(stops).toHaveLength(0)
348+
expect(await storedRun(run.runId)).toMatchObject({
349+
status: 'complete',
350+
recoveryBackoff: { attempts: 2 },
351+
})
352+
} finally {
353+
faults.append = undefined
354+
}
355+
}, 60_000)
356+
334357
it.each([
335-
{ frames: 'any_tool', tails: 1 },
336-
{ frames: 'any_tool', tails: 3 },
337-
{ frames: 'tool_result', tails: 1 },
338-
{ frames: 'tool_result', tails: 3 },
358+
{ frames: 'any_tool', effect: 'throw', tails: 1 },
359+
{ frames: 'any_tool', effect: 'throw', tails: 3 },
360+
{ frames: 'tool_result', effect: 'throw', tails: 1 },
361+
{ frames: 'tool_result', effect: 'throw', tails: 3 },
362+
{ frames: 'any_tool', effect: 'lose_lease', tails: 3 },
339363
] as const)(
340-
'ends the run after one worker request when a $frames frame cannot be persisted ($tails tails)',
341-
async ({ frames, tails }) => {
364+
'gives up on a run whose recovered controllers keep failing ($effect on $frames, $tails tails)',
365+
async ({ frames, effect, tails }) => {
342366
worker.mode = 'frames'
343367
const run = await orphanedRun()
344-
faults.append = { frames, effect: 'throw' }
368+
faults.append = { frames, effect }
345369
try {
346-
const { turns, stops } = await reconnect(run, { tails, windowMs: 15_000 })
370+
const { turns, stops } = await reconnect(run, { tails, windowMs: 60_000 })
347371

348-
expect(turns).toHaveLength(1)
372+
expect(turns).toHaveLength(MAX_RECOVERY_ATTEMPTS)
373+
/** Each takeover waits at least the jittered floor of the one before it. */
374+
turns.slice(1).forEach((turn, i) => {
375+
expect(turn.at - turns[i].at).toBeGreaterThanOrEqual(0.7 * 1_000 * 2 ** i)
376+
})
349377
expect(stops).toHaveLength(1)
350378
expect(await storedRun(run.runId)).toMatchObject({
351379
status: 'error',
352-
error: new StreamPersistenceFailedError(undefined).userMessage,
380+
error: new StreamRecoveryExhaustedError().userMessage,
381+
recoveryBackoff: { attempts: MAX_RECOVERY_ATTEMPTS + 1 },
353382
})
354383
} finally {
355384
faults.append = undefined
356385
}
357386
},
358-
60_000
387+
120_000
359388
)
360389

361-
it('gives up on a run whose recovered controllers keep losing their lease', async () => {
362-
worker.mode = 'frames'
363-
const run = await orphanedRun()
364-
faults.append = { frames: 'any_tool', effect: 'lose_lease' }
365-
try {
366-
const { turns, stops } = await reconnect(run, { tails: 3, windowMs: 60_000 })
367-
368-
expect(turns).toHaveLength(MAX_RECOVERY_ATTEMPTS)
369-
/** Each takeover waits at least the jittered floor of the one before it. */
370-
turns.slice(1).forEach((turn, i) => {
371-
expect(turn.at - turns[i].at).toBeGreaterThanOrEqual(0.7 * 1_000 * 2 ** i)
372-
})
373-
expect(stops).toHaveLength(1)
374-
expect(await storedRun(run.runId)).toMatchObject({
375-
status: 'error',
376-
error: new StreamRecoveryExhaustedError().userMessage,
377-
recoveryBackoff: { attempts: MAX_RECOVERY_ATTEMPTS + 1 },
378-
})
379-
} finally {
380-
faults.append = undefined
381-
}
382-
}, 120_000)
383-
384390
it('waits out the backoff of a recent takeover, then takes over and completes the run', async () => {
385391
worker.mode = 'frames'
386392
const now = Date.now()

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

Lines changed: 34 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {
1717
onAbort: undefined as (() => Promise<void>) | undefined,
1818
/** The read-only replay's answer; a worker that does not know the run by default. */
1919
replay: { status: 404, frames: [] as unknown[] },
20+
/** How many coming reads of the controller's own lease fail, as a Redis error would. */
21+
leaseReadFailures: 0,
2022
}
2123
const replayRequests: Array<Record<string, unknown>> = []
2224
const server = createHttpServer(async (request, response) => {
@@ -72,6 +74,20 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {
7274
})
7375

7476
vi.mock('@/lib/auth', () => authMock)
77+
vi.mock('@/lib/mothership/request/session/controller-lease', async (importOriginal) => {
78+
const actual =
79+
await importOriginal<typeof import('@/lib/mothership/request/session/controller-lease')>()
80+
return {
81+
...actual,
82+
assertChatStreamLease: async (...args: Parameters<typeof actual.assertChatStreamLease>) => {
83+
if (worker.hooks.leaseReadFailures > 0) {
84+
worker.hooks.leaseReadFailures--
85+
throw new Error('simulated Redis read failure')
86+
}
87+
return actual.assertChatStreamLease(...args)
88+
},
89+
}
90+
})
7591
vi.mock('@/lib/mothership/request/lifecycle/run', () => ({
7692
/**
7793
* Stands in for the worker leg: forwards each scripted event to the controller's
@@ -147,7 +163,6 @@ import {
147163
StreamReplayBudgetExhaustedError,
148164
} from '@/lib/mothership/request/session/replay-budget'
149165
import { STREAM_STRING_PREVIEW_UNITS } from '@/lib/mothership/request/session/replay-compaction'
150-
import { STREAM_PERSISTENCE_FAILED_CODE } from '@/lib/mothership/request/session/turn-failure'
151166
import type { StreamEvent } from '@/lib/mothership/request/session/types'
152167
import { StreamWriter } from '@/lib/mothership/request/session/writer'
153168
import type { StreamingContext } from '@/lib/mothership/request/types'
@@ -659,7 +674,7 @@ describe.runIf(Boolean(redisUrl))('a turn whose stream exhausts its replay budge
659674
}
660675
)
661676

662-
it('ends a turn whose append fails while it still holds the lease, and stops the worker', async () => {
677+
it('leaves a run it handed off recoverable after a transient append failure', async () => {
663678
let eventsKey = ''
664679
const { streamId, runId, frames } = await runTurn(
665680
[
@@ -674,36 +689,27 @@ describe.runIf(Boolean(redisUrl))('a turn whose stream exhausts its replay budge
674689
}
675690
)
676691

677-
expect(frames.map((frame) => [frame.type, frame.payload.code ?? frame.payload.status])).toEqual(
678-
[
679-
['session', undefined],
680-
['error', STREAM_PERSISTENCE_FAILED_CODE],
681-
['complete', 'error'],
682-
]
683-
)
692+
expect(frames.map((frame) => frame.type)).toEqual(['session'])
693+
expect(await redis().ttl(`mothership_stream:${streamId}:events`)).toBeGreaterThan(300)
684694
const [stored] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId))
685-
expect(stored.status).toBe('error')
686-
expect(worker.abortRequests).toEqual([expect.objectContaining({ messageId: streamId })])
687-
expect(worker.runs[0].dispatched).toEqual([])
688-
expect(await redis().get(chatStreamLockKey(chatId))).toBeNull()
695+
expect(stored.status).toBe('active')
696+
expect(worker.abortRequests).toEqual([])
689697
})
690698

691-
it('settles a finished turn whose lease became unreadable instead of handing it off', async () => {
692-
const { runId, frames } = await runTurn([
693-
text('Done.'),
694-
async () => {
695-
// A lock key of the wrong type makes every read of the lease fail, as a Redis error would.
696-
await redis().del(chatStreamLockKey(chatId))
697-
await redis().hset(chatStreamLockKey(chatId), 'corrupt', '1')
698-
},
699-
])
700-
699+
it('settles a finished turn whose lease could not be read instead of handing it off', async () => {
701700
try {
701+
const { runId, frames } = await runTurn([
702+
text('Done.'),
703+
async () => {
704+
worker.hooks.leaseReadFailures = 1
705+
},
706+
])
707+
702708
expect(frames.at(-1)).toMatchObject({ type: 'complete', payload: { status: 'complete' } })
703709
const [stored] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId))
704710
expect(stored.status).toBe('complete')
705711
} finally {
706-
await redis().del(chatStreamLockKey(chatId))
712+
worker.hooks.leaseReadFailures = 0
707713
}
708714
})
709715

@@ -1065,7 +1071,7 @@ describe.runIf(Boolean(redisUrl))('a turn whose stream exhausts its replay budge
10651071
)
10661072
}
10671073

1068-
it('marks its run terminal even when the final events cannot be persisted', async () => {
1074+
it('marks its run terminal even when the final events cannot be published', async () => {
10691075
const controllerToken = `owner\n${generateId()}`
10701076
const { streamId, runId } = await pausedRun(controllerToken)
10711077
const lease = { key: chatStreamLockKey(generateId()), value: controllerToken }
@@ -1074,8 +1080,10 @@ describe.runIf(Boolean(redisUrl))('a turn whose stream exhausts its replay budge
10741080
await redis().set(`mothership_stream:${streamId}:events`, 'not a sorted set', 'EX', 60)
10751081
const publisher = new StreamWriter({ streamId, requestId: generateId(), lease })
10761082

1077-
await finalizeAsError(publisher, runId)
1083+
const failure = await finalizeAsError(publisher, runId).catch((error: unknown) => error)
10781084

1085+
expect(failure).toBeInstanceOf(Error)
1086+
expect(failure).not.toBeInstanceOf(StreamControllerSupersededError)
10791087
const [stored] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId))
10801088
expect(stored.status).toBe('error')
10811089
await redis().del(lease.key)

‎apps/sim/lib/mothership/request/session/turn-failure.ts‎

Lines changed: 3 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,7 @@
11
/**
2-
* A failure that ends the turn as an error instead of handing it to a successor.
3-
* Only a controller that provably lost its chat lease hands off
4-
* ({@link StreamControllerSupersededError}); a replacement would meet any other failure
5-
* again, re-POSTing the same run to the worker on every reconnect poll.
2+
* A failure that ends the turn as an error instead of handing it to a successor
3+
* ({@link StreamControllerSupersededError}): a replacement would meet it again,
4+
* re-POSTing the same run to the worker on every reconnect poll.
65
*/
76
export abstract class StreamTurnFailure extends Error {
87
abstract readonly code: string
@@ -15,29 +14,9 @@ export function turnFailure(value: unknown): StreamTurnFailure | undefined {
1514
return value instanceof StreamTurnFailure ? value : undefined
1615
}
1716

18-
/** Run-error code for a turn whose controller could not persist or verify its stream. */
19-
export const STREAM_PERSISTENCE_FAILED_CODE = 'stream_persistence_failed'
20-
2117
/** Run-error code for a run whose controllers kept dying before it could finish. */
2218
export const STREAM_RECOVERY_EXHAUSTED_CODE = 'stream_recovery_exhausted'
2319

24-
/**
25-
* The controller could not record an event to the replay buffer, or could not read its
26-
* own lease, for a reason other than losing it (a Redis error that outlasted the retries).
27-
*/
28-
export class StreamPersistenceFailedError extends StreamTurnFailure {
29-
readonly code = STREAM_PERSISTENCE_FAILED_CODE
30-
31-
constructor(cause: unknown) {
32-
super('Stream persistence failed', { cause })
33-
this.name = 'StreamPersistenceFailedError'
34-
}
35-
36-
get userMessage(): string {
37-
return 'This response was stopped because it could not be saved while streaming. The work it already completed has been saved — send a message to continue from there.'
38-
}
39-
}
40-
4120
/** Recovery took over the run too many times without a controller staying alive. */
4221
export class StreamRecoveryExhaustedError extends StreamTurnFailure {
4322
readonly code = STREAM_RECOVERY_EXHAUSTED_CODE

0 commit comments

Comments
 (0)