Skip to content

Commit b154dfd

Browse files
committed
fix(mothership): bound stream recovery and hand off only on a lost lease
A recovered Chat controller that failed to persist an event for any reason other than a budget refusal was treated as superseded, so it left its run recoverable without finalizing. Every reconnect poll then claimed the run again and re-POSTed it to the worker, for as long as a tab stayed open. - Only a lease another controller holds is a hand-off. A failed leased append latches StreamPersistenceFailedError in the writer, like a budget refusal: the turn ends as an error through the existing terminal path, which settles the run and stops the worker run. An unreadable lease proves nothing, so the fenced writes that follow decide ownership. - Refusals, persistence failures and recovery exhaustion share one StreamTurnFailure base that the lifecycle classifies as an error, never a cancellation or a hand-off. - Recovery takeovers carry a per-run budget in the run's request context, written in the same token-compared claim UPDATE, so it covers every pod, tab and caller. The first takeover is immediate; later ones wait an exponential backoff (1 s base, 60 s cap). After 5 takeovers without a controller staying alive for 5 minutes, the next claim ends the run as an error and stops the worker run. Each claim logs one line with its attempt number.
1 parent 90e3b91 commit b154dfd

13 files changed

Lines changed: 810 additions & 153 deletions

‎apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,9 @@ import {
6969
chatStreamLockKey,
7070
} from '@/lib/mothership/request/session/controller-lease'
7171

72+
/** A recovering controller's first takeover of a run. */
73+
const FIRST_RECOVERY = { attempts: 1, claimedAt: 0, notBefore: 0 }
74+
7275
function redis() {
7376
const client = getRedisClient()
7477
if (!client) throw new Error('The integration suite requires TEST_REDIS_URL')
@@ -466,6 +469,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
466469
chatId: orphan.chatId,
467470
previousToken: orphan.controllerToken!,
468471
token: `${orphan.streamId}\n${generateId()}`,
472+
recoveryBackoff: FIRST_RECOVERY,
469473
}),
470474
sweepOrphanedRuns(),
471475
])
@@ -533,6 +537,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
533537
chatId: orphan.chatId,
534538
previousToken: orphan.controllerToken!,
535539
token: `${orphan.streamId}\n${generateId()}`,
540+
recoveryBackoff: FIRST_RECOVERY,
536541
})
537542
)
538543
),
@@ -570,6 +575,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
570575
chatId: orphan.chatId,
571576
previousToken: orphan.controllerToken!,
572577
token: `${orphan.streamId}\n${generateId()}`,
578+
recoveryBackoff: FIRST_RECOVERY,
573579
}),
574580
settleStoppedRunWithoutController(orphan.runId),
575581
])
@@ -604,6 +610,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
604610
chatId: orphan.chatId,
605611
previousToken: orphan.controllerToken!,
606612
token: lease.value,
613+
recoveryBackoff: FIRST_RECOVERY,
607614
})
608615
return { owned, claimed }
609616
} finally {

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,8 @@ vi.mock('@/lib/mothership/request/session/controller-lease', async (original) =>
4343
...(await original<typeof import('@/lib/mothership/request/session/controller-lease')>()),
4444
assertChatStreamLease: hoisted.assertLease,
4545
}))
46-
vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', () => ({
46+
vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', async (original) => ({
47+
...(await original<typeof import('@/lib/mothership/request/lifecycle/controller-ownership')>()),
4748
claimRunController: hoisted.claim,
4849
}))
4950
vi.mock('@/lib/mothership/request/session/buffer', () => ({
@@ -187,6 +188,7 @@ describe('authorized chat stream recovery', () => {
187188
chatId: '22222222-2222-4222-8222-222222222222',
188189
previousToken: 'old-controller',
189190
token: 'stream\nnew-controller',
191+
recoveryBackoff: expect.objectContaining({ attempts: 1 }),
190192
})
191193
expect(mocks.start).toHaveBeenCalledOnce()
192194
const params = mocks.start.mock.calls[0][0]

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

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,10 @@ import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/a
1515
import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context'
1616
import { buildOnComplete, buildOnError } from '@/lib/mothership/chat/completion'
1717
import { restoreBillingAdmission } from '@/lib/mothership/request/lifecycle/admission'
18-
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
18+
import {
19+
claimRunController,
20+
planRecovery,
21+
} from '@/lib/mothership/request/lifecycle/controller-ownership'
1922
import { StreamRecoveryConfigSchema } from '@/lib/mothership/request/lifecycle/recovery-config'
2023
import { createSSEStream } from '@/lib/mothership/request/lifecycle/start'
2124
import { isTerminalStreamStatus } from '@/lib/mothership/request/session'
@@ -28,6 +31,7 @@ import { getLatestSeq, readEvents } from '@/lib/mothership/request/session/buffe
2831
import { assertChatStreamLease } from '@/lib/mothership/request/session/controller-lease'
2932
import { eventToStreamEvent } from '@/lib/mothership/request/session/event'
3033
import { startsAtReplayHead } from '@/lib/mothership/request/session/recovery'
34+
import { StreamRecoveryExhaustedError } from '@/lib/mothership/request/session/turn-failure'
3135
import { getUserEntityPermissions } from '@/lib/workspaces/permissions/utils'
3236

3337
const logger = createLogger('MothershipStreamRecovery')
@@ -81,6 +85,8 @@ export const readChatStream = defineAuthorizedChatUseCase({
8185
) {
8286
throw new OrchestrationError('validation', 'Saved stream identity does not match its chat')
8387
}
88+
const plan = planRecovery(saved.recoveryBackoff, Date.now())
89+
if (plan.kind === 'wait') return run
8490
if (!(await acquirePendingChatStream(chatId, run.streamId, 0))) return run
8591
const lease = getLocalChatStreamLease(chatId, run.streamId)!
8692
try {
@@ -91,11 +97,18 @@ export const readChatStream = defineAuthorizedChatUseCase({
9197
chatId,
9298
previousToken: saved.controllerToken,
9399
token: lease.value,
100+
recoveryBackoff: plan.backoff,
94101
}))
95102
) {
96103
await releasePendingChatStream(chatId, run.streamId, lease)
97104
return (await getLatestRunForStream(run.streamId, userId)) ?? run
98105
}
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+
})
99112
if (isHosted && !config.data.billingAdmission)
100113
throw new OrchestrationError(
101114
'forbidden',
@@ -163,6 +176,7 @@ export const readChatStream = defineAuthorizedChatUseCase({
163176
message: '',
164177
titleModel: '',
165178
resumeSeq,
179+
...(plan.kind === 'exhausted' ? { failure: new StreamRecoveryExhaustedError() } : {}),
166180
orchestrateOptions: {
167181
userId,
168182
workspaceId,

‎apps/sim/lib/mothership/request/go/parser.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { createLogger } from '@sim/logger'
22
import { toError } from '@sim/utils/errors'
33
import { readSSELines } from '@/lib/core/utils/sse'
44
import { StreamControllerSupersededError } from '@/lib/mothership/request/session/controller-lease'
5-
import { StreamReplayBudgetExhaustedError } from '@/lib/mothership/request/session/replay-budget'
5+
import { StreamTurnFailure } from '@/lib/mothership/request/session/turn-failure'
66

77
const logger = createLogger('CopilotSseParser')
88

@@ -55,7 +55,7 @@ export async function processSSEStream(
5555
if (
5656
error instanceof FatalSseEventError ||
5757
error instanceof StreamControllerSupersededError ||
58-
error instanceof StreamReplayBudgetExhaustedError
58+
error instanceof StreamTurnFailure
5959
)
6060
throw error
6161
logger.warn('Failed to handle SSE event', {

‎apps/sim/lib/mothership/request/lifecycle/controller-ownership.ts‎

Lines changed: 58 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,66 @@
11
import { db } from '@sim/db'
22
import { copilotChats, copilotRuns } from '@sim/db/schema'
3+
import { backoffWithJitter } from '@sim/utils/retry'
34
import { and, eq, notInArray, sql } from 'drizzle-orm'
5+
import { z } from 'zod'
46

5-
/** Serializes takeover with assistant persistence; Redis alone cannot fence a delayed DB write. */
7+
/**
8+
* Takeovers of one run allowed in a row before it ends as an error. A controller keeps
9+
* its chat lock while it streams and while it waits on parked, permission-gated or
10+
* client-executed tools, so only a lost controller (a pod deploy or crash, or a lost
11+
* lease) needs one. Five in a row is a crash loop, not a deploy.
12+
*/
13+
export const MAX_RECOVERY_ATTEMPTS = 5
14+
15+
/** A takeover this long after the previous one starts a fresh budget: that controller lived. */
16+
export const RECOVERY_BUDGET_RESET_MS = 5 * 60_000
17+
18+
const RECOVERY_BACKOFF = { baseMs: 1_000, maxMs: 60_000 } as const
19+
20+
/** Stored on the run beside its controller token, in epoch milliseconds. */
21+
const RecoveryBackoffSchema = z.object({
22+
attempts: z.number().int().positive(),
23+
claimedAt: z.number(),
24+
notBefore: z.number(),
25+
})
26+
27+
export type RecoveryBackoff = z.infer<typeof RecoveryBackoffSchema>
28+
29+
export type RecoveryPlan =
30+
| { kind: 'wait' }
31+
| { kind: 'claim' | 'exhausted'; backoff: RecoveryBackoff }
32+
33+
/**
34+
* Exponential backoff with a restart limit for taking over a run, as a supervisor
35+
* restarts a crashing child. The first takeover is immediate; each later one waits
36+
* until the previous one's `notBefore`. Seq progress does not reset the budget: a
37+
* recovered leg re-persists the frames the worker replays since its last checkpoint,
38+
* so a crash loop advances the replay ring on every attempt.
39+
*/
40+
export function planRecovery(saved: unknown, now: number): RecoveryPlan {
41+
const previous = RecoveryBackoffSchema.safeParse(saved)
42+
const fresh = !previous.success || now - previous.data.claimedAt >= RECOVERY_BUDGET_RESET_MS
43+
if (previous.success && !fresh && now < previous.data.notBefore) return { kind: 'wait' }
44+
const attempts = fresh ? 1 : previous.data.attempts + 1
45+
const backoff = {
46+
attempts,
47+
claimedAt: now,
48+
notBefore: now + backoffWithJitter(attempts, null, RECOVERY_BACKOFF),
49+
}
50+
return { kind: attempts > MAX_RECOVERY_ATTEMPTS ? 'exhausted' : 'claim', backoff }
51+
}
52+
53+
/**
54+
* Serializes takeover with assistant persistence; Redis alone cannot fence a delayed DB write.
55+
* Every claim replaces the controller token it compares, so the recovery budget written
56+
* beside it is as atomic as the claim, across pods, tabs and callers.
57+
*/
658
export async function claimRunController(input: {
759
runId: string
860
chatId: string
961
previousToken: string
1062
token: string
63+
recoveryBackoff: RecoveryBackoff
1164
}): Promise<boolean> {
1265
return db.transaction(async (tx) => {
1366
await tx
@@ -18,7 +71,10 @@ export async function claimRunController(input: {
1871
const [run] = await tx
1972
.update(copilotRuns)
2073
.set({
21-
requestContext: sql`jsonb_set(${copilotRuns.requestContext}, '{controllerToken}', ${JSON.stringify(input.token)}::jsonb)`,
74+
requestContext: sql`${copilotRuns.requestContext} || ${JSON.stringify({
75+
controllerToken: input.token,
76+
recoveryBackoff: input.recoveryBackoff,
77+
})}::jsonb`,
2278
updatedAt: new Date(),
2379
})
2480
.where(

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

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ import { StreamRetryWindow } from '@/lib/mothership/request/lifecycle/stream-ret
5353
import { recordDegraded } from '@/lib/mothership/request/metrics'
5454
import { AbortReason } from '@/lib/mothership/request/session/abort-reason'
5555
import { StreamControllerSupersededError } from '@/lib/mothership/request/session/controller-lease'
56-
import { replayRefusal } from '@/lib/mothership/request/session/replay-budget'
56+
import { turnFailure } from '@/lib/mothership/request/session/turn-failure'
5757
import {
5858
getToolCallTerminalData,
5959
requireToolCallStateResult,
@@ -562,27 +562,27 @@ export async function runCopilotLifecycle(
562562
// the work the user watched succeed.
563563
const backendFinishedTurn =
564564
context.completionStatus === MothershipStreamV1CompletionStatus.complete
565-
// A refused replay write aborts the turn to stop it, but the turn failed; it
566-
// was not stopped by the user.
567-
const refusal = replayRefusal(lifecycleOptions.abortSignal?.reason)
565+
// A turn failure (such as a refused replay write) aborts the turn to stop it, but
566+
// the turn failed; it was not stopped by the user.
567+
const failure = turnFailure(lifecycleOptions.abortSignal?.reason)
568568
// Consult the lifecycle signal as well as the flag. `context.wasAborted` is
569569
// only reached from a fanout leg through the (deliberately asymmetric) merge
570570
// in `mergeResumeLegOutputs`, so a Stop landing mid-fanout could otherwise
571571
// classify the turn as a success. Mirrors the check already used below on
572572
// the throw path.
573573
const turnWasAborted =
574-
!refusal &&
574+
!failure &&
575575
(context.completionStatus === MothershipStreamV1CompletionStatus.cancelled ||
576576
context.wasAborted ||
577577
(lifecycleOptions.abortSignal?.aborted ?? false))
578578
const succeeded =
579-
!refusal &&
579+
!failure &&
580580
!turnWasAborted &&
581581
(backendFinishedTurn || (!context.completionStatus && context.errors.length === 0))
582582
// The worker sends an error terminal with no `error` event only when it replays a run
583583
// that already ended (for example at its deadline) to a resume or reattach, because
584584
// that replay does not carry the run's stored reason. Say so rather than leave the turn
585-
// to a generic failure; a reported reason or a replay refusal always wins.
585+
// to a generic failure; a reported reason or a turn failure always wins.
586586
const endedWithoutReason =
587587
!turnWasAborted &&
588588
context.completionStatus === MothershipStreamV1CompletionStatus.error &&
@@ -606,7 +606,7 @@ export async function runCopilotLifecycle(
606606
chatId: context.chatId,
607607
requestId: context.requestId,
608608
...(endedWithoutReason ? { error: ENDED_RUN_MESSAGE } : {}),
609-
...(refusal ? { error: refusal.userMessage, errorCode: refusal.code } : {}),
609+
...(failure ? { error: failure.userMessage, errorCode: failure.code } : {}),
610610
errors: !succeeded && context.errors.length ? context.errors : undefined,
611611
usage: context.usage,
612612
cost: context.cost,
@@ -645,8 +645,8 @@ export async function runCopilotLifecycle(
645645
// partial content can be appended.
646646
// Return `cancelled: true` so upstream classification stays
647647
// consistent with the success-path cancel result.
648-
const refusal = replayRefusal(lifecycleOptions.abortSignal?.reason)
649-
const wasCancelled = !refusal && (lifecycleOptions.abortSignal?.aborted ?? false)
648+
const failure = turnFailure(lifecycleOptions.abortSignal?.reason)
649+
const wasCancelled = !failure && (lifecycleOptions.abortSignal?.aborted ?? false)
650650
// Preserve whatever streamed before the throw for both terminals. A thrown
651651
// backend error (as opposed to an `error` SSE event that lets the loop finish
652652
// normally) must still carry the partial assistant turn so onError can
@@ -661,8 +661,8 @@ export async function runCopilotLifecycle(
661661
toolCalls: buildToolCallSummaries(context),
662662
chatId: context.chatId,
663663
requestId: context.requestId,
664-
error: refusal?.userMessage ?? err.message,
665-
...(refusal ? { errorCode: refusal.code } : {}),
664+
error: failure?.userMessage ?? err.message,
665+
...(failure ? { errorCode: failure.code } : {}),
666666
errors: context.errors.length ? context.errors : undefined,
667667
usage: context.usage,
668668
cost: context.cost,

0 commit comments

Comments
 (0)