Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,9 @@ import {
chatStreamLockKey,
} from '@/lib/mothership/request/session/controller-lease'

/** A recovering controller's first takeover of a run. */
const FIRST_RECOVERY = { attempts: 1, claimedAt: 0, notBefore: 0 }

function redis() {
const client = getRedisClient()
if (!client) throw new Error('The integration suite requires TEST_REDIS_URL')
Expand Down Expand Up @@ -466,6 +469,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
chatId: orphan.chatId,
previousToken: orphan.controllerToken!,
token: `${orphan.streamId}\n${generateId()}`,
recoveryBackoff: FIRST_RECOVERY,
}),
sweepOrphanedRuns(),
])
Expand Down Expand Up @@ -533,6 +537,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
chatId: orphan.chatId,
previousToken: orphan.controllerToken!,
token: `${orphan.streamId}\n${generateId()}`,
recoveryBackoff: FIRST_RECOVERY,
})
)
),
Expand Down Expand Up @@ -570,6 +575,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
chatId: orphan.chatId,
previousToken: orphan.controllerToken!,
token: `${orphan.streamId}\n${generateId()}`,
recoveryBackoff: FIRST_RECOVERY,
}),
settleStoppedRunWithoutController(orphan.runId),
])
Expand Down Expand Up @@ -604,6 +610,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
chatId: orphan.chatId,
previousToken: orphan.controllerToken!,
token: lease.value,
recoveryBackoff: FIRST_RECOVERY,
})
return { owned, claimed }
} finally {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,8 @@ vi.mock('@/lib/mothership/request/session/controller-lease', async (original) =>
...(await original<typeof import('@/lib/mothership/request/session/controller-lease')>()),
assertChatStreamLease: hoisted.assertLease,
}))
vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', () => ({
vi.mock('@/lib/mothership/request/lifecycle/controller-ownership', async (original) => ({
...(await original<typeof import('@/lib/mothership/request/lifecycle/controller-ownership')>()),
claimRunController: hoisted.claim,
}))
vi.mock('@/lib/mothership/request/session/buffer', () => ({
Expand Down Expand Up @@ -187,6 +188,7 @@ describe('authorized chat stream recovery', () => {
chatId: '22222222-2222-4222-8222-222222222222',
previousToken: 'old-controller',
token: 'stream\nnew-controller',
recoveryBackoff: expect.objectContaining({ attempts: 1 }),
})
expect(mocks.start).toHaveBeenCalledOnce()
const params = mocks.start.mock.calls[0][0]
Expand Down
45 changes: 32 additions & 13 deletions apps/sim/lib/mothership/request/application/recover-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,10 @@ import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/a
import { resolveOwnedChatContext } from '@/lib/mothership/chat/application/context'
import { buildOnComplete, buildOnError } from '@/lib/mothership/chat/completion'
import { restoreBillingAdmission } from '@/lib/mothership/request/lifecycle/admission'
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
import {
claimRunController,
planRecovery,
} from '@/lib/mothership/request/lifecycle/controller-ownership'
import { StreamRecoveryConfigSchema } from '@/lib/mothership/request/lifecycle/recovery-config'
import { createSSEStream } from '@/lib/mothership/request/lifecycle/start'
import { isTerminalStreamStatus } from '@/lib/mothership/request/session'
Expand All @@ -28,6 +31,7 @@ import { getLatestSeq, readEvents } from '@/lib/mothership/request/session/buffe
import { assertChatStreamLease } from '@/lib/mothership/request/session/controller-lease'
import { eventToStreamEvent } from '@/lib/mothership/request/session/event'
import { startsAtReplayHead } from '@/lib/mothership/request/session/recovery'
import { StreamRecoveryExhaustedError } from '@/lib/mothership/request/session/turn-failure'
import { getUserEntityPermissions } from '@/lib/workspaces/permissions/utils'

const logger = createLogger('MothershipStreamRecovery')
Expand Down Expand Up @@ -81,21 +85,11 @@ export const readChatStream = defineAuthorizedChatUseCase({
) {
throw new OrchestrationError('validation', 'Saved stream identity does not match its chat')
}
const plan = planRecovery(saved.recoveryBackoff, Date.now())
if (plan.kind === 'wait') return run
if (!(await acquirePendingChatStream(chatId, run.streamId, 0))) return run
const lease = getLocalChatStreamLease(chatId, run.streamId)!
try {
await assertChatStreamLease(lease)
if (
!(await claimRunController({
runId: run.id,
chatId,
previousToken: saved.controllerToken,
token: lease.value,
}))
) {
await releasePendingChatStream(chatId, run.streamId, lease)
return (await getLatestRunForStream(run.streamId, userId)) ?? run
}
if (isHosted && !config.data.billingAdmission)
throw new OrchestrationError(
'forbidden',
Expand Down Expand Up @@ -132,6 +126,30 @@ export const readChatStream = defineAuthorizedChatUseCase({
const recoveredEvents = ringIntact ? events : []
const lastEvent = recoveredEvents.at(-1)
const resumeSeq = lastEvent ? lastEvent.seq : ((await getLatestSeq(run.streamId)) ?? 0)
/**
* Claim last: everything before it can fail without touching the run, so a takeover
* that cannot start neither spends the recovery budget nor refreshes the run, and
* an exhausted claim always reaches the terminal path below.
*/
await assertChatStreamLease(lease)
if (
!(await claimRunController({
runId: run.id,
chatId,
previousToken: saved.controllerToken,
token: lease.value,
recoveryBackoff: plan.backoff,
}))
) {
await releasePendingChatStream(chatId, run.streamId, lease)
return (await getLatestRunForStream(run.streamId, userId)) ?? run
}
logger.info('Claimed stream run for recovery', {
runId: run.id,
streamId: run.streamId,
attempt: plan.backoff.attempts,
exhausted: plan.kind === 'exhausted',
})
const requestId = typeof saved?.requestId === 'string' ? saved.requestId : generateId()
const completion = {
chatId,
Expand Down Expand Up @@ -163,6 +181,7 @@ export const readChatStream = defineAuthorizedChatUseCase({
message: '',
titleModel: '',
resumeSeq,
...(plan.kind === 'exhausted' ? { failure: new StreamRecoveryExhaustedError() } : {}),
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
waleedlatif1 marked this conversation as resolved.
orchestrateOptions: {
userId,
workspaceId,
Expand Down
4 changes: 2 additions & 2 deletions apps/sim/lib/mothership/request/go/parser.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { createLogger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
import { readSSELines } from '@/lib/core/utils/sse'
import { StreamControllerSupersededError } from '@/lib/mothership/request/session/controller-lease'
import { StreamReplayBudgetExhaustedError } from '@/lib/mothership/request/session/replay-budget'
import { StreamTurnFailure } from '@/lib/mothership/request/session/turn-failure'

const logger = createLogger('CopilotSseParser')

Expand Down Expand Up @@ -55,7 +55,7 @@ export async function processSSEStream(
if (
error instanceof FatalSseEventError ||
error instanceof StreamControllerSupersededError ||
error instanceof StreamReplayBudgetExhaustedError
error instanceof StreamTurnFailure
)
throw error
logger.warn('Failed to handle SSE event', {
Expand Down
60 changes: 58 additions & 2 deletions apps/sim/lib/mothership/request/lifecycle/controller-ownership.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,66 @@
import { db } from '@sim/db'
import { copilotChats, copilotRuns } from '@sim/db/schema'
import { backoffWithJitter } from '@sim/utils/retry'
import { and, eq, notInArray, sql } from 'drizzle-orm'
import { z } from 'zod'

/** Serializes takeover with assistant persistence; Redis alone cannot fence a delayed DB write. */
/**
* Takeovers of one run allowed in a row before it ends as an error. A controller keeps
* its chat lock while it streams and while it waits on parked, permission-gated or
* client-executed tools, so only a lost controller (a pod deploy or crash, or a lost
* lease) needs one. Five in a row is a crash loop, not a deploy.
*/
export const MAX_RECOVERY_ATTEMPTS = 5

/** A takeover this long after the previous one starts a fresh budget: that controller lived. */
export const RECOVERY_BUDGET_RESET_MS = 5 * 60_000

const RECOVERY_BACKOFF = { baseMs: 1_000, maxMs: 60_000 } as const

/** Stored on the run beside its controller token, in epoch milliseconds. */
const RecoveryBackoffSchema = z.object({
attempts: z.number().int().positive(),
claimedAt: z.number(),
notBefore: z.number(),
})

export type RecoveryBackoff = z.infer<typeof RecoveryBackoffSchema>

export type RecoveryPlan =
| { kind: 'wait' }
| { kind: 'claim' | 'exhausted'; backoff: RecoveryBackoff }

/**
* Exponential backoff with a restart limit for taking over a run, as a supervisor
* restarts a crashing child. The first takeover is immediate; each later one waits
* until the previous one's `notBefore`. Seq progress does not reset the budget: a
* recovered leg re-persists the frames the worker replays since its last checkpoint,
* so a crash loop advances the replay ring on every attempt.
*/
export function planRecovery(saved: unknown, now: number): RecoveryPlan {
const previous = RecoveryBackoffSchema.safeParse(saved)
const fresh = !previous.success || now - previous.data.claimedAt >= RECOVERY_BUDGET_RESET_MS
if (previous.success && !fresh && now < previous.data.notBefore) return { kind: 'wait' }
const attempts = fresh ? 1 : previous.data.attempts + 1
const backoff = {
attempts,
claimedAt: now,
notBefore: now + backoffWithJitter(attempts, null, RECOVERY_BACKOFF),
}
return { kind: attempts > MAX_RECOVERY_ATTEMPTS ? 'exhausted' : 'claim', backoff }
}

/**
* Serializes takeover with assistant persistence; Redis alone cannot fence a delayed DB write.
* Every claim replaces the controller token it compares, so the recovery budget written
* beside it is as atomic as the claim, across pods, tabs and callers.
*/
export async function claimRunController(input: {
runId: string
chatId: string
previousToken: string
token: string
recoveryBackoff: RecoveryBackoff
}): Promise<boolean> {
return db.transaction(async (tx) => {
await tx
Expand All @@ -18,7 +71,10 @@ export async function claimRunController(input: {
const [run] = await tx
.update(copilotRuns)
.set({
requestContext: sql`jsonb_set(${copilotRuns.requestContext}, '{controllerToken}', ${JSON.stringify(input.token)}::jsonb)`,
requestContext: sql`${copilotRuns.requestContext} || ${JSON.stringify({
controllerToken: input.token,
recoveryBackoff: input.recoveryBackoff,
})}::jsonb`,
updatedAt: new Date(),
})
.where(
Expand Down
24 changes: 12 additions & 12 deletions apps/sim/lib/mothership/request/lifecycle/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ import { StreamRetryWindow } from '@/lib/mothership/request/lifecycle/stream-ret
import { recordDegraded } from '@/lib/mothership/request/metrics'
import { AbortReason } from '@/lib/mothership/request/session/abort-reason'
import { StreamControllerSupersededError } from '@/lib/mothership/request/session/controller-lease'
import { replayRefusal } from '@/lib/mothership/request/session/replay-budget'
import { turnFailure } from '@/lib/mothership/request/session/turn-failure'
import {
getToolCallTerminalData,
requireToolCallStateResult,
Expand Down Expand Up @@ -562,27 +562,27 @@ export async function runCopilotLifecycle(
// the work the user watched succeed.
const backendFinishedTurn =
context.completionStatus === MothershipStreamV1CompletionStatus.complete
// A refused replay write aborts the turn to stop it, but the turn failed; it
// was not stopped by the user.
const refusal = replayRefusal(lifecycleOptions.abortSignal?.reason)
// A turn failure (such as a refused replay write) aborts the turn to stop it, but
// the turn failed; it was not stopped by the user.
const failure = turnFailure(lifecycleOptions.abortSignal?.reason)
// Consult the lifecycle signal as well as the flag. `context.wasAborted` is
// only reached from a fanout leg through the (deliberately asymmetric) merge
// in `mergeResumeLegOutputs`, so a Stop landing mid-fanout could otherwise
// classify the turn as a success. Mirrors the check already used below on
// the throw path.
const turnWasAborted =
!refusal &&
!failure &&
(context.completionStatus === MothershipStreamV1CompletionStatus.cancelled ||
context.wasAborted ||
(lifecycleOptions.abortSignal?.aborted ?? false))
const succeeded =
!refusal &&
!failure &&
!turnWasAborted &&
(backendFinishedTurn || (!context.completionStatus && context.errors.length === 0))
// The worker sends an error terminal with no `error` event only when it replays a run
// that already ended (for example at its deadline) to a resume or reattach, because
// that replay does not carry the run's stored reason. Say so rather than leave the turn
// to a generic failure; a reported reason or a replay refusal always wins.
// to a generic failure; a reported reason or a turn failure always wins.
const endedWithoutReason =
!turnWasAborted &&
context.completionStatus === MothershipStreamV1CompletionStatus.error &&
Expand All @@ -606,7 +606,7 @@ export async function runCopilotLifecycle(
chatId: context.chatId,
requestId: context.requestId,
...(endedWithoutReason ? { error: ENDED_RUN_MESSAGE } : {}),
...(refusal ? { error: refusal.userMessage, errorCode: refusal.code } : {}),
...(failure ? { error: failure.userMessage, errorCode: failure.code } : {}),
errors: !succeeded && context.errors.length ? context.errors : undefined,
usage: context.usage,
cost: context.cost,
Expand Down Expand Up @@ -645,8 +645,8 @@ export async function runCopilotLifecycle(
// partial content can be appended.
// Return `cancelled: true` so upstream classification stays
// consistent with the success-path cancel result.
const refusal = replayRefusal(lifecycleOptions.abortSignal?.reason)
const wasCancelled = !refusal && (lifecycleOptions.abortSignal?.aborted ?? false)
const failure = turnFailure(lifecycleOptions.abortSignal?.reason)
const wasCancelled = !failure && (lifecycleOptions.abortSignal?.aborted ?? false)
// Preserve whatever streamed before the throw for both terminals. A thrown
// backend error (as opposed to an `error` SSE event that lets the loop finish
// normally) must still carry the partial assistant turn so onError can
Expand All @@ -661,8 +661,8 @@ export async function runCopilotLifecycle(
toolCalls: buildToolCallSummaries(context),
chatId: context.chatId,
requestId: context.requestId,
error: refusal?.userMessage ?? err.message,
...(refusal ? { errorCode: refusal.code } : {}),
error: failure?.userMessage ?? err.message,
...(failure ? { errorCode: failure.code } : {}),
errors: context.errors.length ? context.errors : undefined,
usage: context.usage,
cost: context.cost,
Expand Down
Loading
Loading