diff --git a/server/src/__tests__/heartbeat-orphaned-active-lease-sweep.test.ts b/server/src/__tests__/heartbeat-orphaned-active-lease-sweep.test.ts new file mode 100644 index 0000000000..126fba2583 --- /dev/null +++ b/server/src/__tests__/heartbeat-orphaned-active-lease-sweep.test.ts @@ -0,0 +1,499 @@ +import { randomUUID } from "node:crypto"; +import { eq } from "drizzle-orm"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { + agents, + companies, + createDb, + environmentLeases, + environments, + heartbeatRuns, +} from "@paperclipai/db"; +import { + getEmbeddedPostgresTestSupport, + startEmbeddedPostgresTestDatabase, +} from "./helpers/embedded-postgres.js"; + +const mockTelemetryClient = vi.hoisted(() => ({ track: vi.fn() })); +vi.mock("../telemetry.ts", () => ({ getTelemetryClient: () => mockTelemetryClient })); + +vi.mock("../middleware/logger.js", () => ({ + logger: { + child: vi.fn(function child() { + return this; + }), + trace: vi.fn(), + debug: vi.fn(), + info: vi.fn(), + warn: vi.fn(), + error: vi.fn(), + fatal: vi.fn(), + }, + httpLogger: vi.fn(), +})); + +import { logger } from "../middleware/logger.ts"; +import { heartbeatService, type HeartbeatEnvironmentRuntime } from "../services/heartbeat.ts"; + +const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport(); +const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe : describe.skip; + +if (!embeddedPostgresSupport.supported) { + console.warn( + `Skipping embedded Postgres orphaned active lease sweep tests on this host: ${ + embeddedPostgresSupport.reason ?? "unsupported environment" + }`, + ); +} + +describeEmbeddedPostgres("heartbeat sweepOrphanedActiveLeases", () => { + let db!: ReturnType; + let tempDb: Awaited> | null = null; + + beforeAll(async () => { + tempDb = await startEmbeddedPostgresTestDatabase("paperclip-orphaned-active-lease-sweep-"); + db = createDb(tempDb.connectionString); + }, 20_000); + + beforeEach(() => { + vi.mocked(logger.warn).mockClear(); + vi.mocked(logger.error).mockClear(); + }); + + afterEach(async () => { + await db.delete(environmentLeases); + await db.delete(heartbeatRuns); + await db.delete(environments); + await db.delete(agents); + await db.delete(companies); + }); + + afterAll(async () => { + await tempDb?.cleanup(); + }); + + async function seedCompanyAgentAndEnvironment() { + const companyId = randomUUID(); + const agentId = randomUUID(); + const environmentId = randomUUID(); + await db.insert(companies).values({ + id: companyId, + name: "Paperclip", + issuePrefix: `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`, + requireBoardApprovalForNewAgents: false, + }); + await db.insert(agents).values({ + id: agentId, + companyId, + name: "CodexCoder", + role: "engineer", + status: "idle", + adapterType: "codex_local", + adapterConfig: {}, + runtimeConfig: {}, + permissions: {}, + }); + await db.insert(environments).values({ + id: environmentId, + companyId, + name: "Fake Sandbox", + driver: "sandbox", + status: "active", + config: { provider: "fake", image: "ubuntu:24.04" }, + }); + return { companyId, agentId, environmentId }; + } + + async function insertHeartbeatRun(input: { + companyId: string; + agentId: string; + status: string; + }): Promise { + const id = randomUUID(); + await db.insert(heartbeatRuns).values({ + id, + companyId: input.companyId, + agentId: input.agentId, + invocationSource: "on_demand", + status: input.status, + nextEventSeq: 1, + }); + return id; + } + + async function insertActiveLease(input: { + companyId: string; + environmentId: string | null; + heartbeatRunId: string | null; + updatedAt: Date; + provider?: string; + providerLeaseId?: string; + status?: string; + }): Promise { + const id = randomUUID(); + await db.insert(environmentLeases).values({ + id, + companyId: input.companyId, + environmentId: input.environmentId, + heartbeatRunId: input.heartbeatRunId, + status: input.status ?? "active", + leasePolicy: "reuse_by_environment", + provider: input.provider ?? "fake", + providerLeaseId: input.providerLeaseId ?? `sandbox://fake/${id}`, + acquiredAt: input.updatedAt, + lastUsedAt: input.updatedAt, + createdAt: input.updatedAt, + updatedAt: input.updatedAt, + }); + return id; + } + + async function leaseRow(leaseId: string) { + return db + .select() + .from(environmentLeases) + .where(eq(environmentLeases.id, leaseId)) + .then((rows) => rows[0] ?? null); + } + + const oldEnough = () => new Date(Date.now() - 60 * 60 * 1000); + + it("test_flips_an_active_lease_when_its_run_is_failed", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 1 }); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("pending_cleanup"); + }); + + it("test_keeps_an_active_lease_when_its_run_is_running", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "running" }); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 0 }); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("active"); + }); + + it("test_flips_an_active_lease_when_its_heartbeat_run_id_is_null", async () => { + const { companyId, environmentId } = await seedCompanyAgentAndEnvironment(); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: oldEnough(), + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 1 }); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("pending_cleanup"); + }); + + it("test_keeps_an_active_lease_when_its_updated_at_is_inside_the_backoff_window", async () => { + const { companyId, environmentId } = await seedCompanyAgentAndEnvironment(); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: new Date(), + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 0 }); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("active"); + }); + + it("test_keeps_a_retained_lease_unchanged", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + status: "retained", + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 0 }); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("retained"); + }); + + it("test_skips_a_lease_when_another_live_lease_holds_the_same_provider_lease_id", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const sharedProviderLeaseId = "sandbox://fake/shared-resource"; + const orphanedLeaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + providerLeaseId: sharedProviderLeaseId, + }); + // A second lease still owns the same physical sandbox resource. + await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: new Date(), + providerLeaseId: sharedProviderLeaseId, + status: "retained", + }); + + const heartbeat = heartbeatService(db); + const result = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + expect(result).toEqual({ recovered: 0 }); + const row = await leaseRow(orphanedLeaseId); + expect(row?.status).toBe("active"); + }); + + it("test_defers_a_guarded_lease_instead_of_leaving_it_at_the_front_of_the_page", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const sharedProviderLeaseId = "sandbox://fake/shared-resource"; + const orphanedLeaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + providerLeaseId: sharedProviderLeaseId, + }); + await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: new Date(), + providerLeaseId: sharedProviderLeaseId, + status: "retained", + }); + + const heartbeat = heartbeatService(db); + await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + const row = await leaseRow(orphanedLeaseId); + expect(row?.status).toBe("active"); + // The guard moved the row's `updatedAt` forward, so a fresh backoff + // cutoff no longer selects it. The exact value is not the point; only + // that it moved out of the stale range the sweep reads by. + expect(row!.updatedAt.getTime()).toBeGreaterThan(oldEnough().getTime()); + }); + + it("test_a_full_page_of_guarded_leases_does_not_starve_a_later_eligible_orphan", async () => { + const { companyId, environmentId } = await seedCompanyAgentAndEnvironment(); + const veryOld = new Date(Date.now() - 3 * 60 * 60 * 1000); + const staleButNewer = new Date(Date.now() - 60 * 60 * 1000); + + // Fill one full sweep page (ORPHANED_ACTIVE_LEASE_SWEEP_PAGE_SIZE = 20) + // with orphan candidates that are each guarded by a live second owner of + // the same provider resource. Every guarded row is older than the + // eligible orphan below, so an unbounded query would return them first + // on every tick. + for (let i = 0; i < 20; i += 1) { + const providerLeaseId = `sandbox://fake/guarded-${i}`; + await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: veryOld, + providerLeaseId, + }); + await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: new Date(), + providerLeaseId, + status: "retained", + }); + } + + // This orphan has no other owner, so it is eligible for cleanup. It + // sorts behind the 20 guarded rows because it is newer than them, so + // the first sweep page does not reach it. + const eligibleLeaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: staleButNewer, + }); + + const heartbeat = heartbeatService(db); + + const firstTick = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + expect(firstTick).toEqual({ recovered: 0 }); + const eligibleAfterFirstTick = await leaseRow(eligibleLeaseId); + expect(eligibleAfterFirstTick?.status).toBe("active"); + + // The guarded rows moved out of the stale page on the first tick, so + // the second tick reaches the eligible orphan behind them. + const secondTick = await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + expect(secondTick).toEqual({ recovered: 1 }); + const eligibleAfterSecondTick = await leaseRow(eligibleLeaseId); + expect(eligibleAfterSecondTick?.status).toBe("pending_cleanup"); + }); + + it("test_writes_a_failure_reason_on_each_flipped_lease", async () => { + const { companyId, environmentId } = await seedCompanyAgentAndEnvironment(); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: null, + updatedAt: oldEnough(), + }); + + const heartbeat = heartbeatService(db); + await heartbeat.sweepOrphanedActiveLeases({ backoffMs: 5 * 60 * 1000 }); + + const row = await leaseRow(leaseId); + expect(row?.failureReason).toBeTruthy(); + }); + + it("test_stops_the_recovered_sandbox_on_the_same_reaper_tick", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + }); + + const destroyRunLease = vi.fn(async ({ lease }: { lease: { id: string } }) => { + const now = new Date(); + const row = await db + .update(environmentLeases) + .set({ status: "expired", cleanupStatus: "success", updatedAt: now }) + .where(eq(environmentLeases.id, lease.id)) + .returning() + .then((rows) => rows[0] ?? null); + return row ? { ...row, status: "expired" as const } : null; + }); + const heartbeat = heartbeatService(db, { + environmentRuntime: { + destroyRunLease, + } as unknown as HeartbeatEnvironmentRuntime, + }); + + // The periodic timer call uses a five-minute staleness threshold. The + // lease is older than that window, so the recovery flip and the + // pending_cleanup teardown both run in this one call. + await heartbeat.reapOrphanedRuns({ staleThresholdMs: 5 * 60 * 1000 }); + + expect(destroyRunLease).toHaveBeenCalledTimes(1); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("expired"); + }); + + it("test_stops_the_recovered_sandbox_on_the_startup_call", async () => { + const { companyId, agentId, environmentId } = await seedCompanyAgentAndEnvironment(); + const runId = await insertHeartbeatRun({ companyId, agentId, status: "failed" }); + const leaseId = await insertActiveLease({ + companyId, + environmentId, + heartbeatRunId: runId, + updatedAt: oldEnough(), + }); + + const destroyRunLease = vi.fn(async ({ lease }: { lease: { id: string } }) => { + const now = new Date(); + const row = await db + .update(environmentLeases) + .set({ status: "expired", cleanupStatus: "success", updatedAt: now }) + .where(eq(environmentLeases.id, lease.id)) + .returning() + .then((rows) => rows[0] ?? null); + return row ? { ...row, status: "expired" as const } : null; + }); + const heartbeat = heartbeatService(db, { + environmentRuntime: { + destroyRunLease, + } as unknown as HeartbeatEnvironmentRuntime, + }); + + // The startup call uses a zero staleness threshold, so the new sweep and + // the existing pending_cleanup sweep both act without a backoff delay, + // and the recovered lease reaches the provider teardown in this one call. + await heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 }); + + expect(destroyRunLease).toHaveBeenCalledTimes(1); + const row = await leaseRow(leaseId); + expect(row?.status).toBe("expired"); + }); + + it("test_logs_a_distinct_error_kind_and_no_exception_field_when_the_sweep_fails", async () => { + const sentinel = "Bearer sk-SENTINEL-a1b2c3"; + const realSelect = db.select.bind(db); + const selectSpy = vi + .spyOn(db, "select") + .mockImplementation((...args: Parameters) => { + const builder = realSelect(...args); + const realFrom = builder.from.bind(builder); + (builder as { from: unknown }).from = (table: unknown) => { + if (table === environmentLeases) { + const failure = new Error(`sweep query failed: ${sentinel}`); + failure.name = `SweepQueryError ${sentinel}`; + (failure as { code?: string }).code = `ESWEEP ${sentinel}`; + throw failure; + } + return realFrom(table as Parameters[0]); + }; + return builder; + }); + + try { + const heartbeat = heartbeatService(db, { + environmentRuntime: { + destroyRunLease: vi.fn(async () => null), + } as unknown as HeartbeatEnvironmentRuntime, + }); + + // The reaper isolates the sweep, so the reaper itself still resolves. + await expect(heartbeat.reapOrphanedRuns({ staleThresholdMs: 0 })).resolves.toBeDefined(); + + const sweepCall = vi + .mocked(logger.error) + .mock.calls.find((call) => call[1] === "orphaned active environment lease sweep failed"); + expect(sweepCall).toBeDefined(); + const record = sweepCall?.[0] as Record; + + expect(JSON.stringify(record)).not.toContain(sentinel); + expect(record).not.toHaveProperty("err"); + expect(record).not.toHaveProperty("errorName"); + expect(record).not.toHaveProperty("errorCode"); + expect(record).not.toHaveProperty("message"); + expect(record).not.toHaveProperty("stack"); + expect(record).toMatchObject({ errorKind: "orphaned_active_lease_sweep_failed" }); + } finally { + selectSpy.mockRestore(); + } + }); +}); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 9321867d35..3efdbbf5cd 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -662,6 +662,8 @@ const PENDING_CLEANUP_SWEEP_ATTEMPT_CAP = 5; // The reaper stores its retry state under these keys in the lease metadata. const PENDING_CLEANUP_ATTEMPTS_METADATA_KEY = "pendingCleanupRetryAttempts"; const PENDING_CLEANUP_CAP_WARNED_METADATA_KEY = "pendingCleanupRetryCapWarned"; +// The reaper sweeps at most this many orphaned active leases per tick. +const ORPHANED_ACTIVE_LEASE_SWEEP_PAGE_SIZE = 20; // A provider or plugin destroy rejection can carry a bearer credential, a // signed URL, or provider response detail in its name, code, message, cause, or @@ -670,6 +672,7 @@ const PENDING_CLEANUP_CAP_WARNED_METADATA_KEY = "pendingCleanupRetryCapWarned"; // catch site logs a constant, locally generated `errorKind` instead. const PENDING_CLEANUP_RETRY_ERROR_KIND = "destroy_failed"; const PENDING_CLEANUP_SWEEP_ERROR_KIND = "sweep_failed"; +const ORPHANED_ACTIVE_LEASE_SWEEP_ERROR_KIND = "orphaned_active_lease_sweep_failed"; // Read the stored retry attempt count as a safe value, directly in SQL. A // provider can write a malformed value under the attempts key. The type guard @@ -18166,6 +18169,119 @@ export function heartbeatService( ); } + // Move a guarded orphan candidate to the back of the sweep order. The + // shared-resource guard below can skip a row on every tick as long as the + // other owner stays live, retained, or pending cleanup. The select orders + // by `updatedAt` ascending and takes only the oldest page, so an untouched + // skipped row keeps refilling that same page and blocks every row behind + // it. The bump pushes the row past the fixed-size page, so the next tick + // reaches the rows behind it. It costs the row one extra backoff wait, + // which is safe because the guard means a physical sandbox still exists. + async function deferOrphanedActiveLease(leaseId: string): Promise { + const now = new Date(); + await db + .update(environmentLeases) + .set({ updatedAt: now }) + .where( + and( + eq(environmentLeases.id, leaseId), + eq(environmentLeases.status, "active"), + ), + ); + } + + // An active lease is reachable only while its heartbeat run keeps the + // running status. The reaper writes the run status and the lease release as + // two separate statements, so a restart between them can leave a terminal + // run and an active lease. `releaseRunLeases` also skips a lease whose + // environment row is gone, so that lease stays active too. No later query + // finds either lease, because every production query selects an active + // lease by its environment or by its heartbeat run id, never by age. This + // sweep finds both stranded classes and moves each lease to pending_cleanup, + // so the existing pending_cleanup sweep tears the sandbox down from the + // data already on the lease row. + async function sweepOrphanedActiveLeases(opts: { + backoffMs: number; + }): Promise<{ recovered: number }> { + const cutoff = new Date(Date.now() - opts.backoffMs); + + const rows = await db + .select({ lease: environmentLeases }) + .from(environmentLeases) + .leftJoin( + heartbeatRuns, + eq(environmentLeases.heartbeatRunId, heartbeatRuns.id), + ) + .where( + and( + eq(environmentLeases.status, "active"), + or( + isNull(environmentLeases.heartbeatRunId), + inArray(heartbeatRuns.status, [ + ...HEARTBEAT_RUN_TERMINAL_STATUSES, + ]), + ), + lte(environmentLeases.updatedAt, cutoff), + ), + ) + .orderBy(asc(environmentLeases.updatedAt)) + .limit(ORPHANED_ACTIVE_LEASE_SWEEP_PAGE_SIZE); + + let recovered = 0; + for (const { lease } of rows) { + // A provider resource id names one physical sandbox. A different lease + // row can still hold that same resource in a live status, so this sweep + // must not tear down a sandbox that a different lease still owns. + if (lease.provider && lease.providerLeaseId) { + const [otherOwner] = await db + .select({ id: environmentLeases.id }) + .from(environmentLeases) + .where( + and( + ne(environmentLeases.id, lease.id), + eq(environmentLeases.provider, lease.provider), + eq(environmentLeases.providerLeaseId, lease.providerLeaseId), + inArray(environmentLeases.status, [ + "active", + "retained", + "pending_cleanup", + ]), + ), + ) + .limit(1); + if (otherOwner) { + // Defer this row so the fixed-size page reaches the rows behind + // it next tick, instead of re-selecting the same guarded rows + // forever. + await deferOrphanedActiveLease(lease.id); + continue; + } + } + + // Keep the row's existing updatedAt value. The select above already + // proved the row is older than the backoff cutoff, so the + // pending_cleanup sweep can accept the same row in this same tick. A + // fresh timestamp here would push the row inside that sweep's own + // backoff window and delay the teardown by one full tick. + const flipped = await db + .update(environmentLeases) + .set({ + status: "pending_cleanup", + failureReason: "orphaned_active_lease_recovered", + }) + .where( + and( + eq(environmentLeases.id, lease.id), + eq(environmentLeases.status, "active"), + ), + ) + .returning({ id: environmentLeases.id }); + if (flipped.length > 0) recovered += 1; + } + + return { recovered }; + } + // Retry the leases stranded in "pending_cleanup". A failed destroy leaves a // lease in that state forever without this sweep. The reaper tick runs the // sweep. The backoff equals the reaper staleness threshold, so a lease waits @@ -18987,6 +19103,28 @@ export function heartbeatService( ); } + // Recover an active lease whose run already ended before the same-tick + // pending_cleanup sweep, so this tick can stop the recovered sandbox. + // Isolate the sweep so its failure never hides the reaper result. + try { + const orphanedActiveLeaseSweep = await sweepOrphanedActiveLeases({ + backoffMs: staleThresholdMs, + }); + if (orphanedActiveLeaseSweep.recovered > 0) { + logger.warn( + { recovered: orphanedActiveLeaseSweep.recovered }, + "recovered orphaned active environment leases", + ); + } + } catch { + // Log a constant errorKind only. The exception can carry a credential in + // its name, code, message, cause, or stack, so the sweep never reads it. + logger.error( + { errorKind: ORPHANED_ACTIVE_LEASE_SWEEP_ERROR_KIND }, + "orphaned active environment lease sweep failed", + ); + } + // Retry stranded pending_cleanup leases on the same tick. Isolate the sweep // so its failure never hides the reaper result. The backoff equals the // reaper staleness threshold. @@ -29033,6 +29171,7 @@ export function heartbeatService( reconcileHotRestartAdoption, recoverNativeRunsAfterRestart, reapOrphanedRuns, + sweepOrphanedActiveLeases, sweepPendingCleanupLeases, // Override-aware scheduling-suppression check (honors the worktree // run-execution experimental setting). Callers outside the service that