Skip to content

Commit 6c7bf90

Browse files
committed
fix(mothership): claim a recovering run only once its takeover can start
Recovery claimed the run before reading its replay, resolving billing and checking access. A takeover that then failed released the lock but had already spent an attempt and refreshed the run, so an exhausted run whose takeover kept failing there stayed active, and the refreshed row kept the orphan sweep away. Claim last: a takeover that fails before it starts leaves the run untouched, and an exhausted claim always reaches the terminal path.
1 parent d8bf6cc commit 6c7bf90

2 files changed

Lines changed: 68 additions & 21 deletions

File tree

‎apps/sim/lib/mothership/request/application/recover-stream.ts‎

Lines changed: 24 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -90,25 +90,6 @@ export const readChatStream = defineAuthorizedChatUseCase({
9090
if (!(await acquirePendingChatStream(chatId, run.streamId, 0))) return run
9191
const lease = getLocalChatStreamLease(chatId, run.streamId)!
9292
try {
93-
await assertChatStreamLease(lease)
94-
if (
95-
!(await claimRunController({
96-
runId: run.id,
97-
chatId,
98-
previousToken: saved.controllerToken,
99-
token: lease.value,
100-
recoveryBackoff: plan.backoff,
101-
}))
102-
) {
103-
await releasePendingChatStream(chatId, run.streamId, lease)
104-
return (await getLatestRunForStream(run.streamId, userId)) ?? run
105-
}
106-
logger.info('Claimed stream run for recovery', {
107-
runId: run.id,
108-
streamId: run.streamId,
109-
attempt: plan.backoff.attempts,
110-
exhausted: plan.kind === 'exhausted',
111-
})
11293
if (isHosted && !config.data.billingAdmission)
11394
throw new OrchestrationError(
11495
'forbidden',
@@ -145,6 +126,30 @@ export const readChatStream = defineAuthorizedChatUseCase({
145126
const recoveredEvents = ringIntact ? events : []
146127
const lastEvent = recoveredEvents.at(-1)
147128
const resumeSeq = lastEvent ? lastEvent.seq : ((await getLatestSeq(run.streamId)) ?? 0)
129+
/**
130+
* Claim last: everything before it can fail without touching the run, so a takeover
131+
* that cannot start neither spends the recovery budget nor refreshes the run, and
132+
* an exhausted claim always reaches the terminal path below.
133+
*/
134+
await assertChatStreamLease(lease)
135+
if (
136+
!(await claimRunController({
137+
runId: run.id,
138+
chatId,
139+
previousToken: saved.controllerToken,
140+
token: lease.value,
141+
recoveryBackoff: plan.backoff,
142+
}))
143+
) {
144+
await releasePendingChatStream(chatId, run.streamId, lease)
145+
return (await getLatestRunForStream(run.streamId, userId)) ?? run
146+
}
147+
logger.info('Claimed stream run for recovery', {
148+
runId: run.id,
149+
streamId: run.streamId,
150+
attempt: plan.backoff.attempts,
151+
exhausted: plan.kind === 'exhausted',
152+
})
148153
const requestId = typeof saved?.requestId === 'string' ? saved.requestId : generateId()
149154
const completion = {
150155
chatId,

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

Lines changed: 44 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,9 @@
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 once or every
6-
* time, or a lease lost with no successor to take over.
5+
* are injected at two seams: the leased replay append (an append that fails once or every
6+
* time, or a lease lost with no successor to take over) and a recovering controller's read
7+
* of the replay ring.
78
*/
89
import { authMock, authMockFns } from '@sim/testing/mocks/auth.mock'
910
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
@@ -19,6 +20,8 @@ const { redisUrl, inheritedEnv, worker, faults } = await vi.hoisted(async () =>
1920
frames: 'any_tool' | 'tool_result'
2021
effect: 'throw' | 'throw_once' | 'lose_lease'
2122
},
23+
/** Fails a recovering controller's read of the whole replay ring. */
24+
recoveryRead: false,
2225
}
2326
const worker = {
2427
posts: [] as Array<{ path: string; at: number }>,
@@ -88,6 +91,12 @@ vi.mock('@/lib/mothership/request/session/buffer', async (importOriginal) => {
8891
if (hit && fault.effect === 'lose_lease') await getRedisClient()!.del(lease.key)
8992
return actual.appendEvents(...args)
9093
},
94+
readEvents: async (...args: Parameters<typeof actual.readEvents>) => {
95+
const [, afterCursor] = args
96+
if (faults.recoveryRead && afterCursor === '0')
97+
throw new Error('simulated Redis read failure')
98+
return actual.readEvents(...args)
99+
},
91100
}
92101
})
93102

@@ -420,6 +429,39 @@ describe.runIf(Boolean(redisUrl))('reconnecting to an orphaned Chat run', () =>
420429
expect(stored.recoveryBackoff).toMatchObject({ attempts: 1 })
421430
}, 30_000)
422431

432+
it('leaves the run untouched when a takeover fails before it starts, then ends it once exhausted', async () => {
433+
worker.mode = 'frames'
434+
const claimedAt = Date.now() - 1_000
435+
const recoveryBackoff = { attempts: MAX_RECOVERY_ATTEMPTS, claimedAt, notBefore: claimedAt }
436+
const run = await orphanedRun({ recoveryBackoff })
437+
const { controllerToken } = (await storedRun(run.runId)).requestContext as {
438+
controllerToken: string
439+
}
440+
faults.recoveryRead = true
441+
try {
442+
const { turns } = await reconnect(run, { windowMs: 2_000 })
443+
444+
expect(turns).toHaveLength(0)
445+
expect(await storedRun(run.runId)).toMatchObject({
446+
status: 'active',
447+
requestContext: { controllerToken },
448+
recoveryBackoff,
449+
})
450+
} finally {
451+
faults.recoveryRead = false
452+
}
453+
454+
const { turns, stops } = await reconnect(run, { windowMs: 10_000 })
455+
456+
expect(turns).toHaveLength(0)
457+
expect(stops).toHaveLength(1)
458+
expect(await storedRun(run.runId)).toMatchObject({
459+
status: 'error',
460+
error: new StreamRecoveryExhaustedError().userMessage,
461+
recoveryBackoff: { attempts: MAX_RECOVERY_ATTEMPTS + 1 },
462+
})
463+
}, 30_000)
464+
423465
it('never takes over a parked run whose controller still holds the chat lock', async () => {
424466
worker.mode = 'frames'
425467
const run = await orphanedRun({ status: 'paused_waiting_for_tool' })

0 commit comments

Comments
 (0)