From 246c0ef9940c03837be0b133a4ddc5a4d214fe87 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 14:53:12 -0700 Subject: [PATCH 1/2] 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. --- .../async-runs/orphaned-runs.integration.ts | 93 ++++++++++++++++++- .../mothership/async-runs/orphaned-runs.ts | 55 ++++++++--- 2 files changed, 134 insertions(+), 14 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index dae25a4dfa6..067b1aa6d06 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -30,6 +30,7 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { import { db } from '@sim/db' import { + copilotAsyncToolCalls, copilotChats, copilotRequestStops, copilotRuns, @@ -37,6 +38,7 @@ import { user, workspace, } from '@sim/db/schema' +import { createDeferred } from '@sim/testing' import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import { randomInt } from '@sim/utils/random' @@ -48,7 +50,11 @@ import { settleStoppedRunWithoutController, sweepOrphanedRuns, } from '@/lib/mothership/async-runs/orphaned-runs' -import { requestRunStop, updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { + claimSimToolExecution, + requestRunStop, + updateRunStatus, +} from '@/lib/mothership/async-runs/repository' import { chatPubSub } from '@/lib/mothership/chat-status' import { abortRun } from '@/lib/mothership/request/application/controls' import { claimRunController } from '@/lib/mothership/request/lifecycle/controller-ownership' @@ -184,6 +190,37 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { await requestRunStop({ userId, workspaceId, streamId: run.streamId, chatId: run.chatId }) } + /** A Sim tool call the worker dispatched on the run, not yet admitted for execution. */ + async function dispatchedTool(runId: string) { + const toolCallId = generateId() + await db.insert(copilotAsyncToolCalls).values({ runId, toolCallId, toolName: 'run_workflow' }) + return { toolCallId, runId, userId, ownerToken: generateId() } + } + + /** The backend queued on a lock behind any of these, once one is. */ + async function lockWaiterBehind(...blockers: number[]) { + const pids = sql`ARRAY[${sql.join( + blockers.map((pid) => sql`${pid}::int`), + sql`, ` + )}]` + let waiter: number | undefined + await expect + .poll( + async () => { + const [row] = await db.execute<{ pid: number }>(sql` + SELECT pid FROM pg_stat_activity WHERE datname = current_database() + AND wait_event_type = 'Lock' AND pid <> ALL(${pids}) + AND pg_blocking_pids(pid) && ${pids} LIMIT 1 + `) + waiter = row?.pid + return waiter + }, + { interval: 5, timeout: 5000 } + ) + .toBeDefined() + return waiter! + } + async function stored(runId: string) { const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId)) const [chat] = await db @@ -207,6 +244,60 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect(run.marker).toBeNull() }) + it('never settles a run while one of its Sim tools holds a live execution lease', async () => { + /** A long tool call writes nothing to the run; only its execution heartbeat shows it is alive. */ + const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' }) + const tool = await dispatchedTool(orphan.runId) + expect(await claimSimToolExecution(tool)).toEqual({ outcome: 'claimed' }) + + expect((await sweepOrphanedRuns()).settledRunIds).not.toContain(orphan.runId) + const live = await stored(orphan.runId) + expect(live.status).toBe('paused_waiting_for_tool') + expect(live.toolAdmissionClosedAt).toBeNull() + expect(live.marker).toBe(orphan.streamId) + + /** Its owner died: the heartbeat stopped renewing the lease. */ + await db + .update(copilotAsyncToolCalls) + .set({ executionLeaseExpiresAt: sql`now() - interval '1 second'` }) + .where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId)) + + expect((await sweepOrphanedRuns()).settledRunIds).toContain(orphan.runId) + expect((await stored(orphan.runId)).status).toBe('error') + }) + + it('never settles a run whose Sim tool was admitted while the sweep waited to settle it', async () => { + const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' }) + const tool = await dispatchedTool(orphan.runId) + const locked = createDeferred() + const release = createDeferred() + /** Holds the run row so the tool's admission and then the sweep queue behind it, in that order. */ + const holding = db.transaction(async (tx) => { + await tx + .select({ id: copilotRuns.id }) + .from(copilotRuns) + .where(eq(copilotRuns.id, orphan.runId)) + .for('update') + const [backend] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`) + locked.resolve(backend.pid) + await release.promise + }) + const holder = await locked.promise + + const claim = claimSimToolExecution(tool) + const claimant = await lockWaiterBehind(holder) + const sweep = sweepOrphanedRuns() + await lockWaiterBehind(holder, claimant) + release.resolve() + await holding + + expect(await claim).toEqual({ outcome: 'claimed' }) + expect((await sweep).settledRunIds).not.toContain(orphan.runId) + const run = await stored(orphan.runId) + expect(run.status).toBe('paused_waiting_for_tool') + expect(run.toolAdmissionClosedAt).toBeNull() + }) + it('settles a run stopped while no controller owned it as cancelled', async () => { const orphan = await admittedRun({ idleMinutes: 90, stopped: true }) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index f656084f646..f21f17102e0 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -1,6 +1,7 @@ import { db } from '@sim/db' import { type CopilotRunStatus, + copilotAsyncToolCalls, copilotChats, copilotOrganizationRequestStops, copilotRequestStops, @@ -19,6 +20,7 @@ import { isNull, lt, lte, + not, notInArray, or, type SQL, @@ -51,11 +53,12 @@ const UNFINISHED_RUN_STATUSES: CopilotRunStatus[] = [ * How long a leased run must go without a status write before the sweep may settle it. * * This is a recovery window, not a liveness test, and is independent of any run - * deadline. Liveness comes only from the chat lock: a live controller renews it by - * heartbeat for as long as it runs, however long that is, and a run whose stream holds - * the lock is never settled. For a run with no lock holder, this window and the replay - * buffer's `:seq` key (whose TTL each write renews, `COPILOT_STREAM_TTL_SECONDS`, - * one hour by default) leave a reconnect time to resume it; the sweep waits for both. + * deadline. Liveness comes only from heartbeats, which run for as long as their work + * does, however long that is: a run whose stream holds the chat lock, which its live + * controller renews, or whose Sim tool holds an execution lease, is never settled. For a + * run with neither, this window and the replay buffer's `:seq` key (whose TTL each write + * renews, `COPILOT_STREAM_TTL_SECONDS`, one hour by default) leave a reconnect time to + * resume it; the sweep waits for both. * A TTL configured below this window shortens only that resume window, never safety. */ export const ORPHANED_RUN_GRACE_MS = 60 * 60 * 1000 @@ -92,7 +95,20 @@ function idleFor(ms: number): SQL { return sql`${copilotRuns.updatedAt} < now() - make_interval(secs => ${ms / 1000})` } -const leasedRunIdle = and(isNotNull(controllerToken), idleFor(ORPHANED_RUN_GRACE_MS)) +/** + * One of the run's Sim tools is still executing. A tool call writes nothing to its run + * while it runs, however long that is; its owner only renews this execution lease by + * heartbeat, so an unexpired lease is live Sim work the worker is still waiting on. + */ +const toolExecuting = sql`EXISTS (SELECT 1 FROM ${copilotAsyncToolCalls} t + WHERE t.run_id = ${copilotRuns.id} AND t.execution_settled_at IS NULL + AND t.execution_revoked_at IS NULL AND t.execution_lease_expires_at > clock_timestamp())` + +const leasedRunIdle = and( + isNotNull(controllerToken), + idleFor(ORPHANED_RUN_GRACE_MS), + not(toolExecuting) +) const legacyRunIdle = and( isNull(controllerToken), lt(copilotRuns.toolExecutionVersion, SIM_TOOL_EXECUTION_VERSION), @@ -143,9 +159,11 @@ function terminalValues(reason: 'orphaned' | 'legacy') { * which write the same row, wins or loses atomically against it. * * Chat rows are locked first, in id order, as a controller's claim does, so the two - * never wait on each other in opposite orders. A legacy run keeps its last write as its - * completion and retention time. The chat marker is released without - * touching the chat's ordering timestamp. + * never wait on each other in opposite orders. The run rows are locked next, before the + * guarded update takes its snapshot: a tool's admission locks its run row, so the update + * then sees any execution lease an admission committed, and a later admission sees the + * run settled. A legacy run keeps its last write as its completion and retention time. + * The chat marker is released without touching the chat's ordering timestamp. */ async function settleRuns( tx: DbTransaction, @@ -160,6 +178,17 @@ async function settleRuns( .where(inArray(copilotChats.id, chatIds)) .orderBy(asc(copilotChats.id)) .for('update') + await tx + .select({ id: copilotRuns.id }) + .from(copilotRuns) + .where( + inArray( + copilotRuns.id, + runs.map((run) => run.id) + ) + ) + .orderBy(asc(copilotRuns.id)) + .for('update') const settled: UnownedRun[] = [] const apply = async (batch: UnownedRun[], owner: SQL, values: object) => { @@ -336,10 +365,10 @@ async function settleBatch(candidates: UnownedRun[]): Promise { /** * Settles runs that no controller will ever finish: a leased run whose stream holds no - * chat lock and has no replay buffer left, idle past the recovery window, and a legacy - * run from before the current protocol. Each sweep resumes where the last one stopped - * and wraps to the first run, so no run is starved by the ones before it. A failed - * batch is logged and skipped. + * chat lock and has no replay buffer left, with no Sim tool still executing, idle past + * the recovery window, and a legacy run from before the current protocol. Each sweep + * resumes where the last one stopped and wraps to the first run, so no run is starved by + * the ones before it. A failed batch is logged and skipped. */ export async function sweepOrphanedRuns(): Promise<{ settledRunIds: string[] }> { const settledRunIds: string[] = [] From 53ece36113ac21e42382a64caf056e4b8623c098 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 30 Sep 2026 15:38:28 -0700 Subject: [PATCH 2/2] 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. --- .../async-runs/orphaned-runs.integration.ts | 95 +++++++++++++++---- .../mothership/async-runs/orphaned-runs.ts | 24 +++-- 2 files changed, 93 insertions(+), 26 deletions(-) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts index 067b1aa6d06..05ef2d58341 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts @@ -44,6 +44,7 @@ import { generateId } from '@sim/utils/id' import { randomInt } from '@sim/utils/random' import { eq, inArray, sql } from 'drizzle-orm' import { closeRedisConnection, getRedisClient } from '@/lib/core/config/redis' +import type { DbTransaction } from '@/lib/db/types' import { LEGACY_RUN_ERROR, ORPHANED_RUN_ERROR, @@ -221,6 +222,33 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { return waiter! } + /** + * Runs `lock` in a transaction held open until `run` settles, then commits it, so a + * failed step never leaves the rows locked behind the test. `run` returns the work + * queued behind the lock wrapped, never as a bare promise it would wait on. + */ + async function whileHolding( + lock: (tx: DbTransaction) => Promise, + run: (holder: number) => Promise + ): Promise { + const locked = createDeferred() + const release = createDeferred() + const holding = db.transaction(async (tx) => { + await lock(tx) + const [backend] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`) + locked.resolve(backend.pid) + await release.promise + }) + holding.catch(locked.reject) + const holder = await locked.promise + try { + return await run(holder) + } finally { + release.resolve() + await holding + } + } + async function stored(runId: string) { const [run] = await db.select().from(copilotRuns).where(eq(copilotRuns.id, runId)) const [chat] = await db @@ -269,27 +297,22 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { it('never settles a run whose Sim tool was admitted while the sweep waited to settle it', async () => { const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' }) const tool = await dispatchedTool(orphan.runId) - const locked = createDeferred() - const release = createDeferred() /** Holds the run row so the tool's admission and then the sweep queue behind it, in that order. */ - const holding = db.transaction(async (tx) => { - await tx - .select({ id: copilotRuns.id }) - .from(copilotRuns) - .where(eq(copilotRuns.id, orphan.runId)) - .for('update') - const [backend] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`) - locked.resolve(backend.pid) - await release.promise - }) - const holder = await locked.promise - - const claim = claimSimToolExecution(tool) - const claimant = await lockWaiterBehind(holder) - const sweep = sweepOrphanedRuns() - await lockWaiterBehind(holder, claimant) - release.resolve() - await holding + const { claim, sweep } = await whileHolding( + (tx) => + tx + .select({ id: copilotRuns.id }) + .from(copilotRuns) + .where(eq(copilotRuns.id, orphan.runId)) + .for('update'), + async (holder) => { + const claim = claimSimToolExecution(tool) + const claimant = await lockWaiterBehind(holder) + const sweep = sweepOrphanedRuns() + await lockWaiterBehind(holder, claimant) + return { claim, sweep } + } + ) expect(await claim).toEqual({ outcome: 'claimed' }) expect((await sweep).settledRunIds).not.toContain(orphan.runId) @@ -298,6 +321,38 @@ describe.runIf(Boolean(redisUrl))('Chat runs no controller owns', () => { expect(run.toolAdmissionClosedAt).toBeNull() }) + it('never settles a run whose Sim tool lease a heartbeat renewed as the sweep settled it', async () => { + const orphan = await admittedRun({ idleMinutes: 90, status: 'paused_waiting_for_tool' }) + const tool = await dispatchedTool(orphan.runId) + expect(await claimSimToolExecution(tool)).toEqual({ outcome: 'claimed' }) + await db + .update(copilotAsyncToolCalls) + .set({ executionLeaseExpiresAt: sql`clock_timestamp() - interval '1 second'` }) + .where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId)) + + /** + * A heartbeat that passed its expiry check just before the lease ran out, and has + * not committed yet: the sweep sees the old, expired lease until it does. + */ + const { sweep } = await whileHolding( + (tx) => + tx + .update(copilotAsyncToolCalls) + .set({ executionLeaseExpiresAt: sql`clock_timestamp() + interval '1 minute'` }) + .where(eq(copilotAsyncToolCalls.toolCallId, tool.toolCallId)), + async (holder) => { + const sweep = sweepOrphanedRuns() + await lockWaiterBehind(holder) + return { sweep } + } + ) + + expect((await sweep).settledRunIds).not.toContain(orphan.runId) + const run = await stored(orphan.runId) + expect(run.status).toBe('paused_waiting_for_tool') + expect(run.toolAdmissionClosedAt).toBeNull() + }) + it('settles a run stopped while no controller owned it as cancelled', async () => { const orphan = await admittedRun({ idleMinutes: 90, stopped: true }) diff --git a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts index f21f17102e0..cbee9b04d24 100644 --- a/apps/sim/lib/mothership/async-runs/orphaned-runs.ts +++ b/apps/sim/lib/mothership/async-runs/orphaned-runs.ts @@ -162,8 +162,11 @@ function terminalValues(reason: 'orphaned' | 'legacy') { * never wait on each other in opposite orders. The run rows are locked next, before the * guarded update takes its snapshot: a tool's admission locks its run row, so the update * then sees any execution lease an admission committed, and a later admission sees the - * run settled. A legacy run keeps its last write as its completion and retention time. - * The chat marker is released without touching the chat's ordering timestamp. + * run settled. Their unsettled tool executions are locked last: a lease heartbeat + * writes only the tool row, so one already past its expiry check commits before the + * update reads the lease, and a later one finds the lease expired. A legacy run keeps + * its last write as its completion and retention time. The chat marker is released + * without touching the chat's ordering timestamp. */ async function settleRuns( tx: DbTransaction, @@ -178,16 +181,25 @@ async function settleRuns( .where(inArray(copilotChats.id, chatIds)) .orderBy(asc(copilotChats.id)) .for('update') + const runIds = runs.map((run) => run.id) await tx .select({ id: copilotRuns.id }) .from(copilotRuns) + .where(inArray(copilotRuns.id, runIds)) + .orderBy(asc(copilotRuns.id)) + .for('update') + await tx + .select({ id: copilotAsyncToolCalls.id }) + .from(copilotAsyncToolCalls) .where( - inArray( - copilotRuns.id, - runs.map((run) => run.id) + and( + inArray(copilotAsyncToolCalls.runId, runIds), + isNotNull(copilotAsyncToolCalls.executionOwnerToken), + isNull(copilotAsyncToolCalls.executionSettledAt), + isNull(copilotAsyncToolCalls.executionRevokedAt) ) ) - .orderBy(asc(copilotRuns.id)) + .orderBy(asc(copilotAsyncToolCalls.id)) .for('update') const settled: UnownedRun[] = []