Skip to content

Commit f1d878a

Browse files
committed
fix(mothership): never sweep a current headless run
A headless turn has no chat lease and no heartbeat, so its age says nothing about whether it is still running once runs have no deadline. The sweep's lease-less rule now applies only to runs admitted before the current tool-execution protocol: every run the current code admits records the current version, so after a deploy no such row can be live. A current headless run is left to its own lifecycle, which always settles it. The protocol version moves beside the other async-run constants so the sweep can read it without importing the repository.
1 parent bc9369c commit f1d878a

4 files changed

Lines changed: 52 additions & 25 deletions

File tree

‎apps/sim/lib/mothership/async-runs/lifecycle.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,9 @@ import {
44
MothershipStreamV1ToolOutcome,
55
} from '@/lib/mothership/generated/mothership-stream-v1'
66

7+
/** Recorded on every run the current code admits; older values mark runs from earlier protocols. */
8+
export const SIM_TOOL_EXECUTION_VERSION = 2
9+
710
export const ASYNC_TOOL_STATUS = MothershipStreamV1AsyncToolRecordStatus
811

912
export const EXECUTABLE_TOOL_PERMISSION_DECISIONS = [

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

Lines changed: 22 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -43,10 +43,10 @@ import { randomInt } from '@sim/utils/random'
4343
import { eq, inArray, sql } from 'drizzle-orm'
4444
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
4545
import {
46+
LEGACY_RUN_ERROR,
4647
ORPHANED_RUN_ERROR,
4748
settleStoppedRunWithoutController,
4849
sweepOrphanedRuns,
49-
UNLEASED_RUN_ERROR,
5050
} from '@/lib/mothership/async-runs/orphaned-runs'
5151
import { updateRunStatus } from '@/lib/mothership/async-runs/repository'
5252
import { abortRun } from '@/lib/mothership/request/application/controls'
@@ -119,6 +119,8 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
119119
controllerToken?: string | null
120120
stopped?: boolean
121121
superseded?: boolean
122+
/** Admitted by code predating the current tool-execution protocol. */
123+
legacy?: boolean
122124
} = {}
123125
) {
124126
const chatId = generateId()
@@ -144,7 +146,7 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
144146
userId,
145147
workspaceId,
146148
streamId,
147-
toolExecutionVersion: 2,
149+
toolExecutionVersion: options.legacy ? 0 : 2,
148150
status: options.status ?? 'active',
149151
requestContext: controllerToken
150152
? { requestId: generateId(), controllerToken, recovery: { kind: 'interactive_stream' } }
@@ -201,9 +203,13 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
201203
expect(run.error).toBeTruthy()
202204
})
203205

204-
it('settles a run without a lease only after the unleased ceiling', async () => {
205-
const recent = await admittedRun({ idleMinutes: 90, controllerToken: null })
206-
const abandoned = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null })
206+
it('settles a legacy run without a lease only after the legacy ceiling', async () => {
207+
const recent = await admittedRun({ idleMinutes: 90, controllerToken: null, legacy: true })
208+
const abandoned = await admittedRun({
209+
idleMinutes: 25 * 60,
210+
controllerToken: null,
211+
legacy: true,
212+
})
207213
const [before] = await db
208214
.select({ updatedAt: copilotRuns.updatedAt })
209215
.from(copilotRuns)
@@ -217,11 +223,21 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
217223
/** Its retention clock keeps running from its last real write, and it reads as never finalized. */
218224
expect(settled.updatedAt).toEqual(before.updatedAt)
219225
expect(settled.completedAt).toEqual(before.updatedAt)
220-
expect(settled.error).toBe(UNLEASED_RUN_ERROR)
226+
expect(settled.error).toBe(LEGACY_RUN_ERROR)
221227
expect(settled.error).not.toBe(ORPHANED_RUN_ERROR)
222228
expect((await stored(recent.runId)).status).toBe('active')
223229
})
224230

231+
it('never settles a current headless run, however long it has run', async () => {
232+
/** A headless turn has no lease or heartbeat; only its own lifecycle can end it. */
233+
const headless = await admittedRun({ idleMinutes: 25 * 60, controllerToken: null })
234+
235+
const { settledRunIds } = await sweepOrphanedRuns()
236+
237+
expect(settledRunIds).not.toContain(headless.runId)
238+
expect((await stored(headless.runId)).status).toBe('active')
239+
})
240+
225241
it('never settles a run whose stream holds its chat lock, or that is still recoverable', async () => {
226242
const leased = await admittedRun({ idleMinutes: 90 })
227243
await redis().set(chatStreamLockKey(leased.chatId), leased.controllerToken!, 'EX', 60)

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

Lines changed: 26 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,13 @@ import {
1717
inArray,
1818
isNotNull,
1919
isNull,
20+
lt,
2021
notInArray,
2122
or,
2223
type SQL,
2324
sql,
2425
} from 'drizzle-orm'
26+
import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle'
2527
import { publishChatStatusChanged } from '@/lib/mothership/chat-status'
2628
import { findStreamsWithReplay } from '@/lib/mothership/request/session/buffer'
2729
import { findStreamsHoldingChatLock } from '@/lib/mothership/request/session/controller-lease'
@@ -50,14 +52,16 @@ const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [
5052
export const ORPHANED_RUN_GRACE_MS = 60 * 60 * 1000
5153

5254
/**
53-
* Runs admitted without a chat lease (headless turns and rows from before the lease
54-
* protocol) have no liveness signal, so only an age far past any process lifetime
55-
* proves them dead.
55+
* Runs admitted by code predating the current tool-execution protocol. Every run the
56+
* current code admits records the current version, so once a deploy has replaced the
57+
* processes that admitted these, none can be live; the age only leaves room for a
58+
* rollout. A current run without a lease (a headless turn) is never swept: it has no
59+
* liveness signal and its own lifecycle always settles it.
5660
*/
57-
export const UNLEASED_RUN_GRACE_MS = 24 * 60 * 60 * 1000
61+
export const LEGACY_RUN_GRACE_MS = 24 * 60 * 60 * 1000
5862

5963
export const ORPHANED_RUN_ERROR = 'This response was interrupted before it finished.'
60-
export const UNLEASED_RUN_ERROR = 'Run was never finalized (no controller lease).'
64+
export const LEGACY_RUN_ERROR = 'Run was never finalized (pre-lease run).'
6165

6266
const SWEEP_BATCH_SIZE = 500
6367
const SWEEP_MAX_ROWS_PER_RUN = 5_000
@@ -71,8 +75,12 @@ function idleFor(ms: number): SQL {
7175
}
7276

7377
const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS))
74-
const unleasedRunIdle = and(isNull(controllerToken), idleFor(UNLEASED_RUN_GRACE_MS))
75-
const orphanIdle = or(leasedRunIdle, unleasedRunIdle)
78+
const legacyRunIdle = and(
79+
isNull(controllerToken),
80+
lt(copilotRuns.toolExecutionVersion, SIM_TOOL_EXECUTION_VERSION),
81+
idleFor(LEGACY_RUN_GRACE_MS)
82+
)
83+
const orphanIdle = or(leasedRunIdle, legacyRunIdle)
7684

7785
/** The user pressed Stop on this stream; a newer turn also closes tool admission, without one. */
7886
const stopRequested = sql`(EXISTS (SELECT 1 FROM ${copilotRequestStops} s
@@ -105,11 +113,11 @@ const unownedRunColumns = {
105113
type Transaction = Parameters<Parameters<typeof db.transaction>[0]>[0]
106114

107115
/** A stopped run ends cancelled; the sweep ends it cancelled only if its user pressed Stop. */
108-
function terminalValues(reason: 'stopped' | 'orphaned' | 'unleased') {
116+
function terminalValues(reason: 'stopped' | 'orphaned' | 'legacy') {
109117
if (reason === 'stopped') {
110118
return { status: sql`'cancelled'::copilot_run_status`, error: sql`NULL::text` }
111119
}
112-
const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : UNLEASED_RUN_ERROR
120+
const error = reason === 'orphaned' ? ORPHANED_RUN_ERROR : LEGACY_RUN_ERROR
113121
return {
114122
status: sql`(CASE WHEN ${stopRequested} THEN 'cancelled' ELSE 'error' END)::copilot_run_status`,
115123
error: sql`CASE WHEN ${stopRequested} THEN NULL ELSE ${error}::text END`,
@@ -122,8 +130,8 @@ function terminalValues(reason: 'stopped' | 'orphaned' | 'unleased') {
122130
* which write the same row, wins or loses atomically against it.
123131
*
124132
* Chat rows are locked first, in id order, as a controller's claim does, so the two
125-
* never wait on each other in opposite orders. A run without a lease keeps its last
126-
* write as its completion and retention time. The chat marker is released without
133+
* never wait on each other in opposite orders. A legacy run keeps its last write as its
134+
* completion and retention time. The chat marker is released without
127135
* touching the chat's ordering timestamp.
128136
*/
129137
async function settleRuns(
@@ -169,7 +177,7 @@ async function settleRuns(
169177
await apply(
170178
runs.filter((run) => run.controllerToken === null),
171179
isNull(controllerToken),
172-
{ ...terminalValues('unleased'), completedAt: sql`${copilotRuns.updatedAt}` }
180+
{ ...terminalValues('legacy'), completedAt: sql`${copilotRuns.updatedAt}` }
173181
)
174182
for (const run of runs) {
175183
if (run.controllerToken === null) continue
@@ -212,28 +220,28 @@ function announceSettled(runs: UnownedRun[]): void {
212220
/** The candidates no controller owns; leased runs are skipped when ownership is unreadable. */
213221
async function withoutOwners(candidates: UnownedRun[]): Promise<UnownedRun[]> {
214222
const leased = candidates.filter((run) => run.controllerToken !== null)
215-
const unleased = candidates.filter((run) => run.controllerToken === null)
216-
if (leased.length === 0) return unleased
223+
const legacy = candidates.filter((run) => run.controllerToken === null)
224+
if (leased.length === 0) return legacy
217225
try {
218226
const [locked, replayable] = await Promise.all([
219227
findStreamsHoldingChatLock(leased),
220228
findStreamsWithReplay(leased.map((run) => run.streamId)),
221229
])
222-
return unleased.concat(
230+
return legacy.concat(
223231
leased.filter((run) => !locked.has(run.streamId) && !replayable.has(run.streamId))
224232
)
225233
} catch (error) {
226234
logger.warn('Chat stream ownership is unreadable; leaving leased runs for a later sweep', {
227235
error: getErrorMessage(error),
228236
})
229-
return unleased
237+
return legacy
230238
}
231239
}
232240

233241
/**
234242
* Settles runs that no controller will ever finish: a leased run whose stream holds no
235-
* chat lock and has no replay buffer left, idle past the recovery window, and a run
236-
* without a lease idle past any process lifetime. A failed batch is logged and skipped.
243+
* chat lock and has no replay buffer left, idle past the recovery window, and a legacy
244+
* run from before the current protocol. A failed batch is logged and skipped.
237245
*/
238246
export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> {
239247
const settledRunIds: string[] = []

‎apps/sim/lib/mothership/async-runs/repository.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ import {
4242
type AsyncTerminalStatus,
4343
DESKTOP_TOOL_CLAIM_OWNER,
4444
EXECUTABLE_TOOL_PERMISSION_DECISIONS,
45+
SIM_TOOL_EXECUTION_VERSION,
4546
} from '@/lib/mothership/async-runs/lifecycle'
4647
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
4748
import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1'
@@ -55,7 +56,6 @@ import { chatSandboxSessionKey } from '@/lib/mothership/tools/sandbox-session-ke
5556

5657
const logger = createLogger('CopilotAsyncRunsRepo')
5758
const WORKFLOW_EXECUTION_CLAIM_PREFIX = 'workflow:'
58-
const SIM_TOOL_EXECUTION_VERSION = 2
5959
const TERMINAL_RUN_STATUSES: CopilotRunStatus[] = ['complete', 'error', 'cancelled']
6060
// Resolve the tracer lazily per-call to avoid capturing the NoOp tracer
6161
// before NodeSDK installs the global TracerProvider (Next.js 16/Turbopack

0 commit comments

Comments
 (0)