Fence sandbox cleanup across warm run ownership transfers
Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
parent
3497dec88b
commit
f98fb8059e
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
]);
|
||||
|
||||
|
|
|
|||
|
|
@ -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<void>(resolve => { reached = resolve; });
|
||||
const wait = new Promise<void>(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<typeof update>) => {
|
||||
const builder = update(...args);
|
||||
return new Proxy(builder, { get(builderTarget, builderKey) {
|
||||
if (builderKey !== "set") return Reflect.get(builderTarget, builderKey);
|
||||
return (...values: Parameters<typeof builder.set>) => {
|
||||
const query = builder.set(...values);
|
||||
return new Proxy(query, { get(queryTarget, queryKey, receiver) {
|
||||
if (queryKey === "where") return (...where: Parameters<typeof query.where>) => { 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<typeof query.returning>) => { 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<void>(resolve => { enteredProvider = resolve; });
|
||||
const pending = new Promise<void>(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;
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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<EnvironmentLeaseStatus, "released" | "expired" | "failed">;
|
||||
|
|
@ -503,6 +504,7 @@ function resolvePluginSandboxRpcTimeoutMs(config: Record<string, unknown>): numb
|
|||
}
|
||||
|
||||
export interface EnvironmentDriverLeaseInput {
|
||||
onProviderOperationStart?: () => void;
|
||||
environment: Environment;
|
||||
lease: EnvironmentLease;
|
||||
failureReason?: string;
|
||||
|
|
@ -1322,11 +1324,6 @@ function createSandboxEnvironmentDriver(
|
|||
): Promise<EnvironmentLease> {
|
||||
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<boolean>`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); });
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<EnvironmentLeaseStatus, "released" | "expired" | "failed" | "retained" | "pending_cleanup"> = "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(),
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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<string, unknown>, 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: {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue