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..05ef2d58341 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,18 +38,24 @@ 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' 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, 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 +191,64 @@ 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! + } + + /** + * 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 @@ -207,6 +272,87 @@ 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) + /** Holds the run row so the tool's admission and then the sweep queue behind it, in that order. */ + 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) + const run = await stored(orphan.runId) + expect(run.status).toBe('paused_waiting_for_tool') + 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 f656084f646..cbee9b04d24 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,14 @@ 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. 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, @@ -160,6 +181,26 @@ 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( + and( + inArray(copilotAsyncToolCalls.runId, runIds), + isNotNull(copilotAsyncToolCalls.executionOwnerToken), + isNull(copilotAsyncToolCalls.executionSettledAt), + isNull(copilotAsyncToolCalls.executionRevokedAt) + ) + ) + .orderBy(asc(copilotAsyncToolCalls.id)) + .for('update') const settled: UnownedRun[] = [] const apply = async (batch: UnownedRun[], owner: SQL, values: object) => { @@ -336,10 +377,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[] = []