Skip to content

Commit 55d877b

Browse files
authored
fix(mothership): never sweep a Chat run while one of its Sim tools is executing (#8482)
* fix(mothership): never sweep a Chat run while one of its Sim tools is executing The orphaned-run sweep settled a leased run once it had gone an hour without a status write and no controller held its chat lock. A Sim tool call writes nothing to its run while it executes; only its execution lease heartbeat shows it is alive. So a long tool call (a workflow run can take well over an hour) could have its run settled as interrupted while the worker still held the run, closing tool admission under it. The default one-hour run deadline masked this. The sweep now skips a leased run with an unsettled, unrevoked tool execution whose lease has not expired, both when it selects candidates and in the guarded update. It also locks the run rows before that update, so a tool admission (which locks its run row) either commits a lease the update then sees, or finds the run settled and is refused. * fix(mothership): lock a run's unsettled tool executions before the sweep settles it A lease heartbeat writes only the tool row, so one that passed its expiry check before the lease ran out could commit after the sweep read the old lease and settled the run. The sweep now locks those tool rows after the run rows, so the guarded update sees a committed renewal and a later heartbeat finds the lease expired. Lock-holder tests release on failure.
1 parent 517a5f3 commit 55d877b

2 files changed

Lines changed: 201 additions & 14 deletions

File tree

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

Lines changed: 147 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,25 +30,32 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {
3030

3131
import { db } from '@sim/db'
3232
import {
33+
copilotAsyncToolCalls,
3334
copilotChats,
3435
copilotRequestStops,
3536
copilotRuns,
3637
permissions,
3738
user,
3839
workspace,
3940
} from '@sim/db/schema'
41+
import { createDeferred } from '@sim/testing'
4042
import { sleep } from '@sim/utils/helpers'
4143
import { generateId } from '@sim/utils/id'
4244
import { randomInt } from '@sim/utils/random'
4345
import { eq, inArray, sql } from 'drizzle-orm'
4446
import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis'
47+
import type { DbTransaction } from '@/lib/db/types'
4548
import {
4649
LEGACY_RUN_ERROR,
4750
ORPHANED_RUN_ERROR,
4851
settleStoppedRunWithoutController,
4952
sweepOrphanedRuns,
5053
} from '@/lib/mothership/async-runs/orphaned-runs'
51-
import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository'
54+
import {
55+
claimSimToolExecution,
56+
requestRunStop,
57+
updateRunStatus,
58+
} from '@/lib/mothership/async-runs/repository'
5259
import { chatPubSub } from '@/lib/mothership/chat-status'
5360
import { abortRun } from '@/lib/mothership/request/application/controls'
5461
import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership'
@@ -184,6 +191,64 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
184191
await requestRunStop({ userId, workspaceId, streamId: run.streamId, chatId: run.chatId })
185192
}
186193

194+
/** A Sim tool call the worker dispatched on the run, not yet admitted for execution. */
195+
async function dispatchedTool(runId: string) {
196+
const toolCallId = generateId()
197+
await db.insert(copilotAsyncToolCalls).values({ runId, toolCallId, toolName: 'run_workflow' })
198+
return { toolCallId, runId, userId, ownerToken: generateId() }
199+
}
200+
201+
/** The backend queued on a lock behind any of these, once one is. */
202+
async function lockWaiterBehind(...blockers: number[]) {
203+
const pids = sql`ARRAY[${sql.join(
204+
blockers.map((pid) => sql`${pid}::int`),
205+
sql`, `
206+
)}]`
207+
let waiter: number | undefined
208+
await expect
209+
.poll(
210+
async () => {
211+
const [row] = await db.execute<{ pid: number }>(sql`
212+
SELECT pid FROM pg_stat_activity WHERE datname = current_database()
213+
AND wait_event_type = 'Lock' AND pid <> ALL(${pids})
214+
AND pg_blocking_pids(pid) && ${pids} LIMIT 1
215+
`)
216+
waiter = row?.pid
217+
return waiter
218+
},
219+
{ interval: 5, timeout: 5000 }
220+
)
221+
.toBeDefined()
222+
return waiter!
223+
}
224+
225+
/**
226+
* Runs `lock` in a transaction held open until `run` settles, then commits it, so a
227+
* failed step never leaves the rows locked behind the test. `run` returns the work
228+
* queued behind the lock wrapped, never as a bare promise it would wait on.
229+
*/
230+
async function whileHolding<T>(
231+
lock: (tx: DbTransaction) => Promise<unknown>,
232+
run: (holder: number) => Promise<T>
233+
): Promise<T> {
234+
const locked = createDeferred<number>()
235+
const release = createDeferred<void>()
236+
const holding = db.transaction(async (tx) => {
237+
await lock(tx)
238+
const [backend] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`)
239+
locked.resolve(backend.pid)
240+
await release.promise
241+
})
242+
holding.catch(locked.reject)
243+
const holder = await locked.promise
244+
try {
245+
return await run(holder)
246+
} finally {
247+
release.resolve()
248+
await holding
249+
}
250+
}
251+
187252
async function stored(runId: string) {
188253
const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId))
189254
const [chat] = await db
@@ -207,6 +272,87 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => {
207272
expect(run.marker).toBeNull()
208273
})
209274

275+
it('never settles a run while one of its Sim tools holds a live execution lease', async () => {
276+
/** A long tool call writes nothing to the run; only its execution heartbeat shows it is alive. */
277+
const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' })
278+
const tool = await dispatchedTool(orphan.runId)
279+
expect(await claimSimToolExecution(tool)).toEqual({ outcome: 'claimed' })
280+
281+
expect((await sweepOrphanedRuns()).settledRunIds).not.toContain(orphan.runId)
282+
const live = await stored(orphan.runId)
283+
expect(live.status).toBe('paused_waiting_for_tool')
284+
expect(live.toolAdmissionClosedAt).toBeNull()
285+
expect(live.marker).toBe(orphan.streamId)
286+
287+
/** Its owner died: the heartbeat stopped renewing the lease. */
288+
await db
289+
.update(copilotAsyncToolCalls)
290+
.set({ executionLeaseExpiresAt: sql`now() - interval '1 second'` })
291+
.where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId))
292+
293+
expect((await sweepOrphanedRuns()).settledRunIds).toContain(orphan.runId)
294+
expect((await stored(orphan.runId)).status).toBe('error')
295+
})
296+
297+
it('never settles a run whose Sim tool was admitted while the sweep waited to settle it', async () => {
298+
const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' })
299+
const tool = await dispatchedTool(orphan.runId)
300+
/** Holds the run row so the tool's admission and then the sweep queue behind it, in that order. */
301+
const { claim, sweep } = await whileHolding(
302+
(tx) =>
303+
tx
304+
.select({ id: copilotRuns.id })
305+
.from(copilotRuns)
306+
.where(eq(copilotRuns.id, orphan.runId))
307+
.for('update'),
308+
async (holder) => {
309+
const claim = claimSimToolExecution(tool)
310+
const claimant = await lockWaiterBehind(holder)
311+
const sweep = sweepOrphanedRuns()
312+
await lockWaiterBehind(holder, claimant)
313+
return { claim, sweep }
314+
}
315+
)
316+
317+
expect(await claim).toEqual({ outcome: 'claimed' })
318+
expect((await sweep).settledRunIds).not.toContain(orphan.runId)
319+
const run = await stored(orphan.runId)
320+
expect(run.status).toBe('paused_waiting_for_tool')
321+
expect(run.toolAdmissionClosedAt).toBeNull()
322+
})
323+
324+
it('never settles a run whose Sim tool lease a heartbeat renewed as the sweep settled it', async () => {
325+
const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' })
326+
const tool = await dispatchedTool(orphan.runId)
327+
expect(await claimSimToolExecution(tool)).toEqual({ outcome: 'claimed' })
328+
await db
329+
.update(copilotAsyncToolCalls)
330+
.set({ executionLeaseExpiresAt: sql`clock_timestamp() - interval '1 second'` })
331+
.where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId))
332+
333+
/**
334+
* A heartbeat that passed its expiry check just before the lease ran out, and has
335+
* not committed yet: the sweep sees the old, expired lease until it does.
336+
*/
337+
const { sweep } = await whileHolding(
338+
(tx) =>
339+
tx
340+
.update(copilotAsyncToolCalls)
341+
.set({ executionLeaseExpiresAt: sql`clock_timestamp() + interval '1 minute'` })
342+
.where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId)),
343+
async (holder) => {
344+
const sweep = sweepOrphanedRuns()
345+
await lockWaiterBehind(holder)
346+
return { sweep }
347+
}
348+
)
349+
350+
expect((await sweep).settledRunIds).not.toContain(orphan.runId)
351+
const run = await stored(orphan.runId)
352+
expect(run.status).toBe('paused_waiting_for_tool')
353+
expect(run.toolAdmissionClosedAt).toBeNull()
354+
})
355+
210356
it('settles a run stopped while no controller owned it as cancelled', async () => {
211357
const orphan = await admittedRun({ idleMinutes: 90, stopped: true })
212358

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

Lines changed: 54 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { db } from '@sim/db'
22
import {
33
type CopilotRunStatus,
4+
copilotAsyncToolCalls,
45
copilotChats,
56
copilotOrganizationRequestStops,
67
copilotRequestStops,
@@ -19,6 +20,7 @@ import {
1920
isNull,
2021
lt,
2122
lte,
23+
not,
2224
notInArray,
2325
or,
2426
type SQL,
@@ -51,11 +53,12 @@ const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [
5153
* How long a leased run must go without a status write before the sweep may settle it.
5254
*
5355
* This is a recovery window, not a liveness test, and is independent of any run
54-
* deadline. Liveness comes only from the chat lock: a live controller renews it by
55-
* heartbeat for as long as it runs, however long that is, and a run whose stream holds
56-
* the lock is never settled. For a run with no lock holder, this window and the replay
57-
* buffer's `:seq` key (whose TTL each write renews, `COPILOT_STREAM_TTL_SECONDS`,
58-
* one hour by default) leave a reconnect time to resume it; the sweep waits for both.
56+
* deadline. Liveness comes only from heartbeats, which run for as long as their work
57+
* does, however long that is: a run whose stream holds the chat lock, which its live
58+
* controller renews, or whose Sim tool holds an execution lease, is never settled. For a
59+
* run with neither, this window and the replay buffer's `:seq` key (whose TTL each write
60+
* renews, `COPILOT_STREAM_TTL_SECONDS`, one hour by default) leave a reconnect time to
61+
* resume it; the sweep waits for both.
5962
* A TTL configured below this window shortens only that resume window, never safety.
6063
*/
6164
export const ORPHANED_RUN_GRACE_MS = 60 * 60 * 1000
@@ -92,7 +95,20 @@ function idleFor(ms: number): SQL {
9295
return sql`${copilotRuns.updatedAt} < now() - make_interval(secs => ${ms / 1000})`
9396
}
9497

95-
const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS))
98+
/**
99+
* One of the run's Sim tools is still executing. A tool call writes nothing to its run
100+
* while it runs, however long that is; its owner only renews this execution lease by
101+
* heartbeat, so an unexpired lease is live Sim work the worker is still waiting on.
102+
*/
103+
const toolExecuting = sql`EXISTS (SELECT 1 FROM ${copilotAsyncToolCalls} t
104+
WHERE t.run_id = ${copilotRuns.id} AND t.execution_settled_at IS NULL
105+
AND t.execution_revoked_at IS NULL AND t.execution_lease_expires_at > clock_timestamp())`
106+
107+
const leasedRunIdle = and(
108+
isNotNull(controllerToken),
109+
idleFor(ORPHANED_RUN_GRACE_MS),
110+
not(toolExecuting)
111+
)
96112
const legacyRunIdle = and(
97113
isNull(controllerToken),
98114
lt(copilotRuns.toolExecutionVersion, SIM_TOOL_EXECUTION_VERSION),
@@ -143,9 +159,14 @@ function terminalValues(reason: 'orphaned' | 'legacy') {
143159
* which write the same row, wins or loses atomically against it.
144160
*
145161
* Chat rows are locked first, in id order, as a controller's claim does, so the two
146-
* never wait on each other in opposite orders. A legacy run keeps its last write as its
147-
* completion and retention time. The chat marker is released without
148-
* touching the chat's ordering timestamp.
162+
* never wait on each other in opposite orders. The run rows are locked next, before the
163+
* guarded update takes its snapshot: a tool's admission locks its run row, so the update
164+
* then sees any execution lease an admission committed, and a later admission sees the
165+
* run settled. Their unsettled tool executions are locked last: a lease heartbeat
166+
* writes only the tool row, so one already past its expiry check commits before the
167+
* update reads the lease, and a later one finds the lease expired. A legacy run keeps
168+
* its last write as its completion and retention time. The chat marker is released
169+
* without touching the chat's ordering timestamp.
149170
*/
150171
async function settleRuns(
151172
tx: DbTransaction,
@@ -160,6 +181,26 @@ async function settleRuns(
160181
.where(inArray(copilotChats.id, chatIds))
161182
.orderBy(asc(copilotChats.id))
162183
.for('update')
184+
const runIds = runs.map((run) => run.id)
185+
await tx
186+
.select({ id: copilotRuns.id })
187+
.from(copilotRuns)
188+
.where(inArray(copilotRuns.id, runIds))
189+
.orderBy(asc(copilotRuns.id))
190+
.for('update')
191+
await tx
192+
.select({ id: copilotAsyncToolCalls.id })
193+
.from(copilotAsyncToolCalls)
194+
.where(
195+
and(
196+
inArray(copilotAsyncToolCalls.runId, runIds),
197+
isNotNull(copilotAsyncToolCalls.executionOwnerToken),
198+
isNull(copilotAsyncToolCalls.executionSettledAt),
199+
isNull(copilotAsyncToolCalls.executionRevokedAt)
200+
)
201+
)
202+
.orderBy(asc(copilotAsyncToolCalls.id))
203+
.for('update')
163204

164205
const settled: UnownedRun[] = []
165206
const apply = async (batch: UnownedRun[], owner: SQL, values: object) => {
@@ -336,10 +377,10 @@ async function settleBatch(candidates: UnownedRun[]): Promise<UnownedRun[]> {
336377

337378
/**
338379
* Settles runs that no controller will ever finish: a leased run whose stream holds no
339-
* chat lock and has no replay buffer left, idle past the recovery window, and a legacy
340-
* run from before the current protocol. Each sweep resumes where the last one stopped
341-
* and wraps to the first run, so no run is starved by the ones before it. A failed
342-
* batch is logged and skipped.
380+
* chat lock and has no replay buffer left, with no Sim tool still executing, idle past
381+
* the recovery window, and a legacy run from before the current protocol. Each sweep
382+
* resumes where the last one stopped and wraps to the first run, so no run is starved by
383+
* the ones before it. A failed batch is logged and skipped.
343384
*/
344385
export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> {
345386
const settledRunIds: string[] = []

0 commit comments

Comments
 (0)