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
148 changes: 147 additions & 1 deletion apps/sim/lib/mothership/async-runs/orphaned-runs.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,25 +30,32 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => {

import { db } from '@sim/db'
import {
copilotAsyncToolCalls,
copilotChats,
copilotRequestStops,
copilotRuns,
permissions,
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'
Expand Down Expand Up @@ -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()
Comment thread
waleedlatif1 marked this conversation as resolved.
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<T>(
lock: (tx: DbTransaction) => Promise<unknown>,
run: (holder: number) => Promise<T>
): Promise<T> {
const locked = createDeferred<number>()
const release = createDeferred<void>()
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)
Comment thread
waleedlatif1 marked this conversation as resolved.
} 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
Expand All @@ -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 })

Expand Down
67 changes: 54 additions & 13 deletions apps/sim/lib/mothership/async-runs/orphaned-runs.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { db } from '@sim/db'
import {
type CopilotRunStatus,
copilotAsyncToolCalls,
copilotChats,
copilotOrganizationRequestStops,
copilotRequestStops,
Expand All @@ -19,6 +20,7 @@ import {
isNull,
lt,
lte,
not,
notInArray,
or,
type SQL,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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,
Expand All @@ -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 })
Comment thread
waleedlatif1 marked this conversation as resolved.
.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) => {
Expand Down Expand Up @@ -336,10 +377,10 @@ async function settleBatch(candidates: UnownedRun[]): Promise<UnownedRun[]> {

/**
* 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[] = []
Expand Down
Loading