diff --git a/doc/sandbox-work-folders.md b/doc/sandbox-work-folders.md index c34565d34b..4f96f4c85e 100644 --- a/doc/sandbox-work-folders.md +++ b/doc/sandbox-work-folders.md @@ -100,8 +100,18 @@ An incomplete resume keeps a provisional lease marker until provider verificatio succeeds. Failed-run cleanup retains that exact lease without stopping or deleting its sandbox, so a retry cannot silently create a replacement. Daytona refreshes the live sandbox state on explicit resume, including externally stopped resources -whose cached handles still say running. Failure to stop a reusable sandbox is -reported and retained for retry; it never falls back to deletion or orphan cleanup. +whose cached handles still say running. Terminal sandbox release claims ownership +in Postgres before provider calls; duplicate completion paths and an old run's +stale lease snapshot cannot stop a newer owner. Native runs persist their selected +resource disposition independently of workspace copy-back, so recovery retains a +successful warm sandbox consistently. + +An uncertain provider stop leaves a durable release claim and reports +`sandbox_release_recovery_required`. Startup cannot resume that resource or create +a replacement while its stop may still be in flight. Recovery requires verifying +that the original provider operation settled before resolving the matching claim; +time passing or an application restart never clears it automatically. The working +copy remains retained, and deletion or orphan cleanup is not a fallback. This applies to both runner generations and leaves distinct task/user bindings isolated. If Daytona rejects a command because its cached shell session no longer exists, diff --git a/server/src/__tests__/cli-invocation-safety.test.ts b/server/src/__tests__/cli-invocation-safety.test.ts index 3186fa48d2..8cb452b197 100644 --- a/server/src/__tests__/cli-invocation-safety.test.ts +++ b/server/src/__tests__/cli-invocation-safety.test.ts @@ -332,6 +332,8 @@ const SKIP_DIRS = new Set([ ".next", "coverage", ".paperclip", + // Ignored local databases and downloaded acceptance evidence are not guidance. + ".paperclip-runtime", "tmp", ]); diff --git a/server/src/__tests__/environment-runtime.test.ts b/server/src/__tests__/environment-runtime.test.ts index e74b44abd7..58f247b7e5 100644 --- a/server/src/__tests__/environment-runtime.test.ts +++ b/server/src/__tests__/environment-runtime.test.ts @@ -31,7 +31,7 @@ import { startEmbeddedPostgresTestDatabase, } from "./helpers/embedded-postgres.js"; import { bindLegacySandboxIdentity, taskUsesLegacySandboxWorkspace } from "../services/legacy-sandbox-workspace.js"; -import { workFolderSandboxKey } from "../services/work-folder-retention.js"; +import { retainUnsavedWorkFolderLease, workFolderSandboxKey } from "../services/work-folder-retention.js"; import { resolveEnvironmentDriverConfigForRuntime } from "../services/environment-config.ts"; import { SANDBOX_CAPABILITY_KEYS, @@ -470,6 +470,162 @@ describeEmbeddedPostgres("environmentRuntimeService", () => { return { pluginId, companyId, agentId, environment, runId, executionWorkspaceId, reusableLease }; } + + it.each(["failed", "pending_cleanup", "expired"] as const)("protects unsaved current copies in %s from cleanup", async (status) => { + const f = await seedReusablePluginSandboxLease(); + await db.update(environmentLeases).set({ status }).where(eq(environmentLeases.id, f.reusableLease.id)); + await db.insert(workFolderRuns).values({ runId: f.runId, companyId: f.companyId, state: "failed", + manifest: { version: 1, companyId: f.companyId, runId: f.runId, agentId: f.agentId, taskId: null, projectId: null, + responsibleUserId: null, leaseId: f.reusableLease.id, sandboxKey: workFolderSandboxKey(f.reusableLease), home: "/home/sandbox", + folders: { task: null, agent: null, user: null, project: null }, repositories: [] } }); + expect(await retainUnsavedWorkFolderLease(db, f.reusableLease)).toBe(true); + expect((await environmentService(db).getLeaseById(f.reusableLease.id))?.status).toBe("retained"); + }); + + it("retention compare-and-set loses to an ownership handoff between its read and write", async () => { + const f = await seedReusablePluginSandboxLease(); + const service = environmentService(db); + await db.insert(workFolderRuns).values({ runId: f.runId, companyId: f.companyId, state: "failed", + manifest: { version: 1, companyId: f.companyId, runId: f.runId, agentId: f.agentId, taskId: null, projectId: null, + responsibleUserId: null, leaseId: f.reusableLease.id, sandboxKey: workFolderSandboxKey(f.reusableLease), home: "/home/sandbox", + folders: { task: null, agent: null, user: null, project: null }, repositories: [] } }); + let reached!: () => void; let resume!: () => void; + const ready = new Promise(resolve => { reached = resolve; }); + const wait = new Promise(resolve => { resume = resolve; }); + const update = db.update.bind(db); + const delayedDb = new Proxy(db, { get(target, key) { + if (key !== "update") { const value = Reflect.get(target, key); return typeof value === "function" ? value.bind(target) : value; } + return (...args: Parameters) => { + const builder = update(...args); + return new Proxy(builder, { get(builderTarget, builderKey) { + if (builderKey !== "set") return Reflect.get(builderTarget, builderKey); + return (...values: Parameters) => { + const query = builder.set(...values); + return new Proxy(query, { get(queryTarget, queryKey, receiver) { + if (queryKey === "where") return (...where: Parameters) => { query.where(...where); return receiver; }; + if (queryKey !== "returning") { const value = Reflect.get(queryTarget, queryKey); return typeof value === "function" ? value.bind(queryTarget) : value; } + return async (...selection: Parameters) => { reached(); await wait; return query.returning(...selection); }; + } }); + }; + } }); + }; + } }); + const retention = retainUnsavedWorkFolderLease(delayedDb, f.reusableLease); + // The real update predicate executes after the old row has been retired. + await ready; + await service.releaseLease(f.reusableLease.id, "retained"); + const nextRunId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: nextRunId, companyId: f.companyId, agentId: f.agentId, status: "running" }); + await service.acquireLease({ companyId: f.companyId, environmentId: f.environment.id, executionWorkspaceId: f.executionWorkspaceId, + heartbeatRunId: nextRunId, leasePolicy: "reuse_by_environment", provider: "fake-plugin", providerLeaseId: f.reusableLease.providerLeaseId, + metadata: f.reusableLease.metadata, replacesReusableLeaseId: f.reusableLease.id }); + resume(); + expect(await retention).toBe(true); + expect((await service.getLeaseById(f.reusableLease.id))?.status).toBe("expired"); + }); + + it("does not resurrect a retired old lease when the new owner has unsaved files", async () => { + const f = await seedReusablePluginSandboxLease(); + const service = environmentService(db); + await service.releaseLease(f.reusableLease.id, "retained"); + const nextRunId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: nextRunId, companyId: f.companyId, agentId: f.agentId, status: "running" }); + const next = await service.acquireLease({ companyId: f.companyId, environmentId: f.environment.id, executionWorkspaceId: f.executionWorkspaceId, + heartbeatRunId: nextRunId, leasePolicy: "reuse_by_environment", provider: "fake-plugin", providerLeaseId: f.reusableLease.providerLeaseId, + metadata: f.reusableLease.metadata, replacesReusableLeaseId: f.reusableLease.id }); + await db.insert(workFolderRuns).values({ runId: nextRunId, companyId: f.companyId, state: "starting", + manifest: { version: 1, companyId: f.companyId, runId: nextRunId, agentId: f.agentId, taskId: null, projectId: null, + responsibleUserId: null, leaseId: next.id, sandboxKey: workFolderSandboxKey(next), home: "/home/sandbox", + folders: { task: null, agent: null, user: null, project: null }, repositories: [] } }); + expect(await retainUnsavedWorkFolderLease(db, f.reusableLease)).toBe(true); + expect((await service.getLeaseById(f.reusableLease.id))?.status).toBe("expired"); + expect((await service.getLeaseById(next.id))?.status).toBe("active"); + expect(await retainUnsavedWorkFolderLease(db, next)).toBe(true); + expect((await service.getLeaseById(next.id))?.metadata?.workFolderRecoveryRequired).toBe(true); + }); + + it("atomically fences duplicate release snapshots and handoff until the matching completion", async () => { + const f = await seedReusablePluginSandboxLease(); + const service = environmentService(db); + const input = { id: f.reusableLease.id, companyId: f.companyId, heartbeatRunId: f.runId, providerLeaseId: f.reusableLease.providerLeaseId }; + const claims = await Promise.all([service.claimRunLeaseRelease(input), service.claimRunLeaseRelease(input)]); + expect(claims.filter(Boolean)).toHaveLength(1); + const claim = claims.find(Boolean)!; + await service.updateLeaseMetadata(f.reusableLease.id, f.reusableLease.metadata); + expect((await service.getLeaseById(f.reusableLease.id))?.metadata?.sandboxReleasePending).toMatchObject({ token: claim.token }); + await service.releaseLease(f.reusableLease.id, "retained", { cleanupStatus: "success" }); + const nextRunId = randomUUID(); + await db.insert(heartbeatRuns).values({ id: nextRunId, companyId: f.companyId, agentId: f.agentId, invocationSource: "manual", status: "running" }); + const handoff = { companyId: f.companyId, environmentId: f.environment.id, executionWorkspaceId: f.executionWorkspaceId, + heartbeatRunId: nextRunId, leasePolicy: "reuse_by_environment" as const, provider: "fake-plugin", providerLeaseId: f.reusableLease.providerLeaseId, + metadata: f.reusableLease.metadata, replacesReusableLeaseId: f.reusableLease.id }; + await expect(service.acquireLease(handoff)).rejects.toThrow(/ownership changed/); + expect(await service.finishRunLeaseRelease({ id: input.id, companyId: f.companyId, token: randomUUID(), confirmed: true })).toBeNull(); + await expect(service.acquireLease(handoff)).rejects.toThrow(/ownership changed/); + await service.finishRunLeaseRelease({ id: input.id, companyId: f.companyId, token: claim.token, confirmed: true }); + const next = await service.acquireLease(handoff); + expect(next.heartbeatRunId).toBe(nextRunId); + await service.updateLeaseMetadata(f.reusableLease.id, f.reusableLease.metadata); + expect((await service.getLeaseById(f.reusableLease.id))?.metadata?.reusableLeaseReplacedByRunId).toBe(nextRunId); + // Generic old cleanup must not re-enable ownership after retirement. + await service.releaseLease(f.reusableLease.id, "retained"); + await expect(service.acquireLease(handoff)).rejects.toThrow(/ownership changed/); + await db.update(environmentLeases).set({ status: "active" }).where(eq(environmentLeases.id, f.reusableLease.id)); + // The permanent handoff marker fences even this stale status overwrite. + expect(await service.claimRunLeaseRelease(input)).toBeNull(); + expect((await service.getLeaseById(next.id))?.status).toBe("active"); + }); + + it("serializes real provider release calls without holding a database transaction open", async () => { + const f = await seedReusablePluginSandboxLease(); + let releaseProvider!: () => void; + let enteredProvider!: () => void; + const entered = new Promise(resolve => { enteredProvider = resolve; }); + const pending = new Promise(resolve => { releaseProvider = resolve; }); + const call = vi.fn(async (_id: string, method: string) => { + if (method !== "environmentReleaseLease") throw new Error("Unexpected provider operation"); + enteredProvider(); await pending; + }); + const manager = { isRunning: () => true, call, getWorker: () => ({ supportedMethods: ["environmentReleaseLease"] }) } as unknown as PluginWorkerManager; + const service = environmentRuntimeService(db, { pluginWorkerManager: manager }); + const first = service.releaseRunLeases(f.runId, "released", undefined, "stop_and_retain"); + await entered; + await service.releaseRunLeases(f.runId, "released", undefined, "keep_running"); + expect(call).toHaveBeenCalledTimes(1); + expect((await environmentService(db).getLeaseById(f.reusableLease.id))?.metadata?.sandboxReleasePending).toBeTruthy(); + releaseProvider(); await first; + const after = await environmentService(db).getLeaseById(f.reusableLease.id); + expect(after?.metadata?.sandboxReleasePending).toBeUndefined(); + expect(after?.status).toBe("released"); + }); + + it("keeps concurrent successful warm completion paths running without provider release", async () => { + const f = await seedReusablePluginSandboxLease(); + const call = vi.fn(); + const manager = { isRunning: () => true, call, getWorker: () => ({ supportedMethods: ["environmentReleaseLease"] }) } as unknown as PluginWorkerManager; + const service = environmentRuntimeService(db, { pluginWorkerManager: manager }); + await Promise.all([service.releaseRunLeases(f.runId, "released", undefined, "keep_running"), service.releaseRunLeases(f.runId, "released", undefined, "keep_running")]); + expect(call).not.toHaveBeenCalled(); + const after = await environmentService(db).getLeaseById(f.reusableLease.id); + expect(after?.status).toBe("retained"); + expect(after?.metadata?.sandboxReleasePending).toBeUndefined(); + }); + + it("retains uncertain provider release ownership and blocks same-run resume without fallback", async () => { + const f = await seedReusablePluginSandboxLease(); + const call = vi.fn(async () => { throw new Error("provider RPC timed out after dispatch"); }); + const manager = { isRunning: () => true, call, getWorker: () => ({ supportedMethods: ["environmentReleaseLease", "environmentResumeLease", "environmentDestroyLease"] }) } as unknown as PluginWorkerManager; + const service = environmentRuntimeService(db, { pluginWorkerManager: manager }); + await service.releaseRunLeases(f.runId, "released", undefined, "stop_and_retain"); + const after = await environmentService(db).getLeaseById(f.reusableLease.id); + expect(after?.status).toBe("retained"); + expect(after?.failureReason).toBe("sandbox_release_recovery_required"); + expect(after?.metadata?.sandboxReleasePending).toBeTruthy(); + await expect(service.acquireRunLease({ companyId: f.companyId, environment: f.environment, issueId: null, agentId: f.agentId, + heartbeatRunId: f.runId, persistedExecutionWorkspace: { id: f.executionWorkspaceId, mode: "shared_workspace" } })).rejects.toThrow("sandbox_release_recovery_required"); + expect(call).toHaveBeenCalledTimes(1); + }); + it.each(["codex_local", "paperclip_runner"])("retains a failed %s resume through run cleanup and retries the original sandbox", async (adapterType) => { const seeded = await seedReusablePluginSandboxLease(adapterType); let failResume = false; diff --git a/server/src/__tests__/native-sandbox-lifecycle.test.ts b/server/src/__tests__/native-sandbox-lifecycle.test.ts index be7a68486c..3806f8d37b 100644 --- a/server/src/__tests__/native-sandbox-lifecycle.test.ts +++ b/server/src/__tests__/native-sandbox-lifecycle.test.ts @@ -1,6 +1,8 @@ import { describe, expect, it } from "vitest"; import { providerResourceDispositionForTerminalRun, + recoveredNativeProviderResourceDisposition, + readNativeProviderResourceDisposition, resolveNativeSandboxLifecycle, resolveReusableSandboxLifecycle, } from "../services/heartbeat.js"; @@ -119,3 +121,26 @@ describe("paperclip_runner sandbox lifecycle", () => { ); }); }); + + +describe("recovered native scoped-folder disposition", () => { + it("keeps the selected warm lease without an obsolete workspace copy-back reference", () => { + expect(recoveredNativeProviderResourceDisposition({ nativeProviderResourceDisposition: "keep_running" }, "succeeded", true)).toBe("keep_running"); + }); + it.each(["failed", "cancelled", "timed_out", "running", null])("does not retain a warm process for %s", (status) => { + expect(recoveredNativeProviderResourceDisposition({ nativeProviderResourceDisposition: "keep_running" }, status, true)).toBe("stop_and_retain"); + }); + it("uses an older persisted workspace-sync disposition when no independent stamp exists", () => { + const reference = { schema: "paperclip.native-workspace-sync/v1", state: "prepared", + descriptorSha256: "a".repeat(64), baselineSha256: "a".repeat(64), finalHostSha256: null, + workspaceId: "workspace", leaseId: "lease", providerLeaseId: "sandbox", remoteCwd: "/workspace", resourceDisposition: "keep_running" }; + expect(recoveredNativeProviderResourceDisposition({ nativeWorkspaceSync: reference }, "succeeded", true)).toBe("keep_running"); + }); + it("preserves per-turn policies and conservative failure/unknown behavior", () => { + expect(recoveredNativeProviderResourceDisposition({ nativeProviderResourceDisposition: "destroy" }, "succeeded", true)).toBe("destroy"); + expect(recoveredNativeProviderResourceDisposition({ nativeProviderResourceDisposition: "stop_and_retain" }, "succeeded", true)).toBe("stop_and_retain"); + expect(recoveredNativeProviderResourceDisposition({ nativeProviderResourceDisposition: "destroy" }, "succeeded", false)).toBe("stop_and_retain"); + expect(recoveredNativeProviderResourceDisposition({}, "succeeded", true)).toBe("stop_and_retain"); + expect(readNativeProviderResourceDisposition("invalid")).toBeUndefined(); + }); +}); diff --git a/server/src/services/environment-runtime.ts b/server/src/services/environment-runtime.ts index 6452046793..6cbfc082a6 100644 --- a/server/src/services/environment-runtime.ts +++ b/server/src/services/environment-runtime.ts @@ -1,5 +1,5 @@ import { createHash, randomUUID } from "node:crypto"; -import { and, eq, inArray, sql } from "drizzle-orm"; +import { or, and, eq, inArray, sql } from "drizzle-orm"; import type { Db } from "@paperclipai/db"; import { companySecrets, companySecretVersions, environmentLeases, heartbeatRuns } from "@paperclipai/db"; import type { @@ -479,6 +479,7 @@ export interface EnvironmentDriverAcquireInput { } export interface EnvironmentDriverReleaseInput { + onProviderOperationStart?: () => void; environment: Environment; lease: EnvironmentLease; status: Extract; @@ -503,6 +504,7 @@ function resolvePluginSandboxRpcTimeoutMs(config: Record): numb } export interface EnvironmentDriverLeaseInput { + onProviderOperationStart?: () => void; environment: Environment; lease: EnvironmentLease; failureReason?: string; @@ -1322,11 +1324,6 @@ function createSandboxEnvironmentDriver( ): Promise { if (!input.heartbeatRunId) return lease; const metadata = { ...lease.metadata, sandboxResumePending: true }; - if (lease.heartbeatRunId === input.heartbeatRunId) { - const claimed = await environmentsSvc.updateLeaseMetadata(lease.id, metadata); - if (!claimed) throw new Error("Reusable sandbox claim disappeared before resume"); - return claimed; - } // Transfer ownership atomically before touching the provider. Two server // processes can observe the same released lease; only the winner of this // conditional update/insert may resume its sandbox. A failed resume leaves @@ -1345,7 +1342,9 @@ function createSandboxEnvironmentDriver( // will attest its current expiry; until then use only this run's deadline. expiresAt: input.requestedExpiresAt ?? null, metadata, - replacesReusableLeaseId: lease.id, + ...(lease.heartbeatRunId === input.heartbeatRunId + ? { reusesReusableLeaseId: lease.id } + : { replacesReusableLeaseId: lease.id }), }); } @@ -1834,12 +1833,14 @@ function createSandboxEnvironmentDriver( // finished. A follow-up must wait for that handoff: filtering the active // lease out below would otherwise create an empty competing workspace. // Different tasks and responsible users still receive isolated sandboxes. - if (parsed.config.reuseLease && input.issueId && input.heartbeatRunId && input.agentId && input.executionWorkspaceId) { + if (input.heartbeatRunId && input.agentId && input.executionWorkspaceId) { const handoffDeadline = Date.now() + 30_000; while (true) { const [holder] = await db.select({ providerLeaseId: environmentLeases.providerLeaseId, runStatus: heartbeatRuns.status, + releaseFailureReason: environmentLeases.failureReason, + releasePending: sql`coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'sandboxReleasePending'`, }).from(environmentLeases).innerJoin(heartbeatRuns, and( eq(heartbeatRuns.id, environmentLeases.heartbeatRunId), eq(heartbeatRuns.companyId, environmentLeases.companyId), @@ -1847,14 +1848,18 @@ function createSandboxEnvironmentDriver( eq(environmentLeases.companyId, input.companyId), eq(environmentLeases.environmentId, input.environment.id), eq(environmentLeases.executionWorkspaceId, input.executionWorkspaceId), - eq(environmentLeases.issueId, input.issueId), - eq(environmentLeases.status, "active"), - eq(environmentLeases.leasePolicy, "reuse_by_environment"), + sql`${environmentLeases.issueId} is not distinct from ${input.issueId}`, + or( + sql`coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'sandboxReleasePending'`, + and(sql`${parsed.config.reuseLease === true && input.issueId !== null}`, eq(environmentLeases.status, "active"), + eq(environmentLeases.leasePolicy, "reuse_by_environment"), + sql`${heartbeatRuns.id} <> ${input.heartbeatRunId}`), + ), eq(heartbeatRuns.agentId, input.agentId), sql`${heartbeatRuns.responsibleUserId} is not distinct from ${responsibleUserId}`, - sql`${heartbeatRuns.id} <> ${input.heartbeatRunId}`, )).limit(1); if (!holder) break; + if (holder.releasePending && (holder.releaseFailureReason === "sandbox_release_recovery_required" || Date.now() >= handoffDeadline)) throw new Error("sandbox_release_recovery_required"); if (!["succeeded", "interrupted", "failed", "cancelled", "timed_out"].includes(holder.runStatus) || Date.now() >= handoffDeadline) { throw new ReusableSandboxResumeError({ provider: parsed.config.provider, @@ -2539,6 +2544,7 @@ function createSandboxEnvironmentDriver( environment: input.environment, lease: input.lease, failureReason: "lease_expired", + onProviderOperationStart: input.onProviderOperationStart, }); } @@ -2570,6 +2576,7 @@ function createSandboxEnvironmentDriver( let cleanupStatus: "success" | "failed" = "success"; try { + input.onProviderOperationStart?.(); await releaseSandboxProviderLease({ config: parsed.config, providerLeaseId: input.lease.providerLeaseId, @@ -3027,6 +3034,7 @@ function createSandboxEnvironmentDriver( environment: input.environment, lease: input.lease, failureReason: input.failureReason ?? "lease_destroyed", + onProviderOperationStart: input.onProviderOperationStart, }); }, }; @@ -3168,6 +3176,7 @@ function createSandboxEnvironmentDriver( lease: input.lease, provider: providerKey, }); + input.onProviderOperationStart?.(); await runLeaseReleaseWithRunParent(input.lease.id, () => pluginWorkerManager.call(pluginId, "environmentReleaseLease", { driverKey: providerKey, @@ -3210,6 +3219,7 @@ function createSandboxEnvironmentDriver( } async function destroyReusableSandboxLease(input: { + onProviderOperationStart?: () => void; environment: Environment; lease: EnvironmentLease; failureReason: string; @@ -3235,6 +3245,7 @@ function createSandboxEnvironmentDriver( lease: input.lease, provider: providerKey, }); + input.onProviderOperationStart?.(); await runLeaseReleaseWithRunParent(input.lease.id, () => pluginWorkerManager.call(pluginId, "environmentDestroyLease", { driverKey: providerKey, @@ -3259,6 +3270,7 @@ function createSandboxEnvironmentDriver( if (parsed.driver !== "sandbox") { cleanupStatus = "failed"; } else { + input.onProviderOperationStart?.(); await destroySandboxProviderLease({ config: parsed.config, providerLeaseId: input.lease.providerLeaseId, @@ -3860,14 +3872,27 @@ export function environmentRuntimeService( // error through `onLeaseReleaseError` for its log path. Keep the order // serial. const released: EnvironmentRuntimeLeaseRecord[] = []; - for (const leaseRow of leaseRows) { + for (const candidateRow of leaseRows) { + const environment = candidateRow.environmentId + ? await environmentsSvc.getById(candidateRow.environmentId).catch((error) => { onLeaseReleaseError?.(candidateRow.id, error); return null; }) + : null; + if (!environment) continue; + const sandbox = getLeaseDriverKey(toEnvironmentLeaseSnapshot(candidateRow), environment) === "sandbox"; + const claim = sandbox ? await environmentsSvc.claimRunLeaseRelease({ + id: candidateRow.id, companyId: candidateRow.companyId, heartbeatRunId, + providerLeaseId: candidateRow.providerLeaseId, + }).catch((error) => { onLeaseReleaseError?.(candidateRow.id, error); return null; }) + : { token: null, lease: toEnvironmentLeaseSnapshot(candidateRow) }; + if (!claim) continue; + const leaseRow = claim.lease; + let providerOperationStarted = false; + let releaseConfirmed = false; try { - const environment = leaseRow.environmentId - ? await environmentsSvc.getById(leaseRow.environmentId) - : null; - if (!environment) continue; - const leaseSnapshot = toEnvironmentLeaseSnapshot(leaseRow); + if (claim.token) { + const { sandboxReleasePending: _hostClaim, ...providerMetadata } = leaseSnapshot.metadata ?? {}; + leaseSnapshot.metadata = providerMetadata; + } const pendingResume = await retainIncompleteSandboxResume(db, leaseSnapshot); if (pendingResume) { released.push({ @@ -3879,9 +3904,10 @@ export function environmentRuntimeService( (pendingResume.metadata?.executionWorkspaceMode as ExecutionWorkspace["mode"] | null | undefined) ?? null, }, }); + releaseConfirmed = true; continue; } - if (await retainUnsavedWorkFolderLease(db, leaseSnapshot)) continue; + if (await retainUnsavedWorkFolderLease(db, leaseSnapshot)) { releaseConfirmed = true; continue; } if ( providerResourceDisposition === "keep_running" && leaseSnapshot.leasePolicy === "reuse_by_environment" @@ -3902,6 +3928,7 @@ export function environmentRuntimeService( }, }); } + releaseConfirmed = Boolean(lease); continue; } const driver = getDriver(getLeaseDriverKey(leaseSnapshot, environment)); @@ -3934,11 +3961,13 @@ export function environmentRuntimeService( environment, lease: leaseSnapshot, failureReason: "paperclip_runner_destroy_after_turn", + onProviderOperationStart: () => { providerOperationStarted = true; }, }) : driver ? await driver.releaseRunLease({ environment, lease: leaseSnapshot, + onProviderOperationStart: () => { providerOperationStarted = true; }, // A stopped reusable provider resource must remain eligible // for exact-lease resume independently of turn outcome. status: @@ -3951,6 +3980,8 @@ export function environmentRuntimeService( leaseRow.id, providerResourceDisposition === "destroy" ? "expired" : status, ); + releaseConfirmed = Boolean(lease && (lease.cleanupStatus !== "failed" + || lease.failureReason === "work_folder_save_required" || lease.failureReason === "sandbox_resume_incomplete")); if (!lease) continue; released.push({ @@ -3964,6 +3995,11 @@ export function environmentRuntimeService( }); } catch (error) { onLeaseReleaseError?.(leaseRow.id, error); + } finally { + if (claim.token) await environmentsSvc.finishRunLeaseRelease({ + id: leaseRow.id, companyId: leaseRow.companyId, token: claim.token, + confirmed: releaseConfirmed || !providerOperationStarted, + }).catch((error) => { onLeaseReleaseError?.(leaseRow.id, error); }); } } diff --git a/server/src/services/environments.ts b/server/src/services/environments.ts index ae46d86f1b..0176b37ed8 100644 --- a/server/src/services/environments.ts +++ b/server/src/services/environments.ts @@ -1453,6 +1453,7 @@ export function environmentService(db: Db) { const retired = await tx .update(environmentLeases) .set({ + metadata: sql`coalesce(${environmentLeases.metadata}, '{}'::jsonb) || ${JSON.stringify({ reusableLeaseReplacedByRunId: input.heartbeatRunId })}::jsonb`, status: "expired", releasedAt: now, lastUsedAt: now, @@ -1462,6 +1463,8 @@ export function environmentService(db: Db) { .where( and( eq(environmentLeases.id, input.replacesReusableLeaseId), + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'sandboxReleasePending')`, + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'reusableLeaseReplacedByRunId')`, eq(environmentLeases.companyId, input.companyId), eq(environmentLeases.environmentId, input.environmentId), eq( @@ -1503,6 +1506,8 @@ export function environmentService(db: Db) { .where( and( eq(environmentLeases.id, input.reusesReusableLeaseId), + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'sandboxReleasePending')`, + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'reusableLeaseReplacedByRunId')`, eq(environmentLeases.companyId, input.companyId), eq(environmentLeases.environmentId, input.environmentId), eq( @@ -1552,6 +1557,38 @@ export function environmentService(db: Db) { return toEnvironmentLease(row); }, + /** One terminal release owns provider access; no transaction spans its RPC. */ + claimRunLeaseRelease: async (input: { id: string; companyId: string; heartbeatRunId: string; providerLeaseId: string | null }) => { + const token = randomUUID(); + const claim = { token, runId: input.heartbeatRunId, startedAt: new Date().toISOString() }; + const [row] = await db.update(environmentLeases).set({ + metadata: sql`coalesce(${environmentLeases.metadata}, '{}'::jsonb) || ${JSON.stringify({ sandboxReleasePending: claim })}::jsonb`, + updatedAt: new Date(), + }).where(and( + eq(environmentLeases.id, input.id), eq(environmentLeases.companyId, input.companyId), + eq(environmentLeases.heartbeatRunId, input.heartbeatRunId), eq(environmentLeases.status, "active"), + sql`${environmentLeases.providerLeaseId} is not distinct from ${input.providerLeaseId}`, + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'sandboxReleasePending')`, + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'reusableLeaseReplacedByRunId')`, + )).returning(); + return row ? { token, lease: toEnvironmentLease(row) } : null; + }, + + finishRunLeaseRelease: async (input: { id: string; companyId: string; token: string; confirmed: boolean }) => { + const [row] = await db.update(environmentLeases).set(input.confirmed ? { + metadata: sql`${environmentLeases.metadata} - 'sandboxReleasePending'`, updatedAt: new Date(), + } : { + // A timeout or crashed owner is not permission to stop/resume again. + // Keep the claim until explicit recovery proves the provider settled. + status: "retained", expiresAt: null, releasedAt: null, cleanupStatus: "failed", + failureReason: "sandbox_release_recovery_required", updatedAt: new Date(), + }).where(and( + eq(environmentLeases.id, input.id), eq(environmentLeases.companyId, input.companyId), + sql`${environmentLeases.metadata}->'sandboxReleasePending'->>'token' = ${input.token}`, + )).returning(); + return row ? toEnvironmentLease(row) : null; + }, + releaseLease: async ( id: string, status: Extract = "released", @@ -1660,7 +1697,11 @@ export function environmentService(db: Db) { const row = await db .update(environmentLeases) .set({ - metadata, + // Late provider metadata snapshots cannot erase host lifecycle fences. + metadata: sql`case when ${metadata === null} and not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ?| array['sandboxReleasePending','reusableLeaseReplacedByRunId']) then null else + (${JSON.stringify(metadata ?? {})}::jsonb - 'sandboxReleasePending' - 'reusableLeaseReplacedByRunId') + || case when ${environmentLeases.metadata} ? 'sandboxReleasePending' then jsonb_build_object('sandboxReleasePending', ${environmentLeases.metadata}->'sandboxReleasePending') else '{}'::jsonb end + || case when ${environmentLeases.metadata} ? 'reusableLeaseReplacedByRunId' then jsonb_build_object('reusableLeaseReplacedByRunId', ${environmentLeases.metadata}->'reusableLeaseReplacedByRunId') else '{}'::jsonb end end`, lastUsedAt: new Date(), updatedAt: new Date(), }) diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 5f003028c4..16a6049bc6 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -1776,6 +1776,21 @@ export function providerResourceDispositionForTerminalRun( return status === "succeeded" ? desired : "stop_and_retain"; } +/** The selected lease policy is durable even when scoped work folders replace copy-back. */ +export function readNativeProviderResourceDisposition(value: unknown): ProviderResourceDisposition | undefined { + return value === "keep_running" || value === "stop_and_retain" || value === "destroy" ? value : undefined; +} + +export function recoveredNativeProviderResourceDisposition( + profile: Record, status: string | null | undefined, workspaceSucceeded: boolean, +): ProviderResourceDisposition { + if (!workspaceSucceeded) return "stop_and_retain"; + const desired = readNativeProviderResourceDisposition(profile.nativeProviderResourceDisposition) + ?? readNativeWorkspaceSyncReference(profile.nativeWorkspaceSync)?.resourceDisposition + ?? "stop_and_retain"; + return providerResourceDispositionForTerminalRun(desired, status) ?? "stop_and_retain"; +} + export interface NativeSandboxLifecycle { runnerProcess: "per_turn" | "warm"; sandboxResource: "keep_running" | "stop_and_reuse" | "destroy_after_turn"; @@ -17026,18 +17041,16 @@ export function heartbeatService( succeeded: boolean; }) { const settledRun = await getRun(input.runId); - const workspaceSyncReference = readNativeWorkspaceSyncReference( - parseObject(settledRun?.runnerProfileJson).nativeWorkspaceSync, - ); + await releaseEnvironmentLeasesForRun({ runId: input.runId, companyId: input.companyId, agentId: input.agentId, status: settledRun?.status, failureReason: settledRun?.error ?? undefined, - providerResourceDisposition: input.succeeded - ? (workspaceSyncReference?.resourceDisposition ?? "stop_and_retain") - : "stop_and_retain", + providerResourceDisposition: recoveredNativeProviderResourceDisposition( + parseObject(settledRun?.runnerProfileJson), settledRun?.status, input.succeeded, + ), }); await releaseRuntimeServicesForRun(input.runId).catch(() => undefined); await finalizeAgentStatus( @@ -20980,6 +20993,9 @@ export function heartbeatService( throw new Error("native_runtime_mode_conflict"); } const lockedProfile = parseObject(lockedRun.runnerProfileJson); + providerResourceDispositionForRun = readNativeProviderResourceDisposition(lockedProfile.nativeProviderResourceDisposition) + ?? providerResourceDispositionForRun; + const persistedResourceDisposition = providerResourceDispositionForRun; await measureSandboxOperation("heartbeat.tx.update.set.where", { operationIndex: 124 }, async () => (tx .update(heartbeatRuns) .set({ @@ -20994,6 +21010,7 @@ export function heartbeatService( runnerProfileJson: { ...nativeRuntimeResolution.profile, ...lockedProfile, + ...(persistedResourceDisposition ? { nativeProviderResourceDisposition: persistedResourceDisposition } : {}), ...(providerTraceRequested ? { providerTrace: { diff --git a/server/src/services/work-folder-retention.ts b/server/src/services/work-folder-retention.ts index 9c0bb0e965..a80330bd58 100644 --- a/server/src/services/work-folder-retention.ts +++ b/server/src/services/work-folder-retention.ts @@ -1,15 +1,22 @@ import { createHash } from "node:crypto"; -import { and, desc, eq, sql } from "drizzle-orm"; +import { and, desc, eq, inArray, sql } from "drizzle-orm"; import { environmentLeases, workFolderRuns, type Db } from "@paperclipai/db"; export function workFolderSandboxKey(lease: { id: string; companyId: string; environmentId: string | null; provider: string | null; providerLeaseId: string | null }) { return lease.providerLeaseId ? createHash("sha256").update(JSON.stringify([lease.companyId, lease.environmentId, lease.provider, lease.providerLeaseId])).digest("hex") : lease.id; } -/** A periodic checkpoint is not permission to discard edits made after it. */ +/** + * True means this physical resource must not be torn down by the caller. + * A transferred old lease is protected without mutating or resurrecting its row. + * A periodic checkpoint is not permission to discard edits made after it. + */ export async function retainUnsavedWorkFolderLease(db: Db, lease: { id: string; companyId: string }) { const [row] = await db.select().from(environmentLeases).where(and(eq(environmentLeases.id, lease.id), eq(environmentLeases.companyId, lease.companyId))); if (!row) return false; + // This old row no longer owns the physical resource. Protect it from stale + // teardown callers without resurrecting the retired lease or changing its owner. + if (row.metadata?.reusableLeaseReplacedByRunId) return true; const sandboxKey = workFolderSandboxKey(row); const [run] = await db.select({ manifest: workFolderRuns.manifest, state: workFolderRuns.state }) .from(workFolderRuns).where(and(eq(workFolderRuns.companyId, lease.companyId), @@ -19,6 +26,12 @@ export async function retainUnsavedWorkFolderLease(db: Db, lease: { id: string; await db.update(environmentLeases).set({ status: "retained", expiresAt: null, failureReason: "work_folder_save_required", cleanupStatus: "failed", metadata: sql`coalesce(${environmentLeases.metadata}, '{}'::jsonb) || '{"workFolderRecoveryRequired":true}'::jsonb`, - }).where(and(eq(environmentLeases.id, lease.id), eq(environmentLeases.companyId, lease.companyId))); + }).where(and(eq(environmentLeases.id, lease.id), eq(environmentLeases.companyId, lease.companyId), + inArray(environmentLeases.status, ["active", "released", "retained", "failed", "pending_cleanup", "expired"]), + sql`not (coalesce(${environmentLeases.metadata}, '{}'::jsonb) ? 'reusableLeaseReplacedByRunId')`, + sql`${environmentLeases.heartbeatRunId} is not distinct from ${row.heartbeatRunId}`, + )).returning({ id: environmentLeases.id }); + // Losing the conditional update can mean a new owner won the handoff. + // The unsaved resource is still protected; never authorize stale teardown. return true; }