mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
Defer agent directory cleanup until stop proof is available (#14866)
## Thinking Path > - Paperclip manages agents and their persistent files. > - Each run owns a temporary agent directory and a save receipt. > - Cleanup needs independent proof that the owning process stopped. > - A cleanup call without that proof currently waits for the directory lock anyway. > - A second lock failure can prevent environment release after the run already reported a failed save. > - This change skips cleanup that has no authority and retries unavailable remote copies after exact destruction proof. > - The save failure stays visible. Existing lock owners remain protected. ## Linked Issues or Issue Description Related work: Refs #14787 (lock diagnostics), #14695 (warm instruction ownership), #9667 (stale lock proposal), and #9872 (control-plane ownership proposal). I checked open PRs and issues. This change leaves the shared filesystem lock protocol in place and does not duplicate the warm-retention work in #14695. **What happened?** Heartbeat cleanup records an explicit unavailable instruction-save warning, then calls directory release before releasing the environment lease. Release can wait for a lock even though the copy has no process-stop proof and cannot be removed. That secondary timeout prevents the following lease-release step. If destruction proof arrives later, the unavailable copy is excluded from both recovery queries. **Expected behavior** Skip a release that cannot remove anything. Preserve the failed-save receipt and candidate fields. Once exact remote destruction is recorded, recover remote cleanup without running a provider command. Unavailable local copies retain their potentially uncollected edits even if local stop proof arrives later. A blocked cleanup must not prevent cleanup for other agents. **Steps to reproduce** 1. Prepare an agent directory, report its save unavailable, and leave process-stop proof absent. 2. Hold the shared directory lock and call release. Before this change, release waits and fails although removal is not authorized. 3. Record destruction of the copy's exact remote lease. Before this change, neither recovery sweep selects the unavailable copy. **Paperclip version or commit** Reproduced against `efc2e6810e9bc0dc8cb412b0e7647c0db9821caa`. **Deployment mode** Local and remote execution with persistent agent directories. Tests use an isolated embedded PostgreSQL database and fixture transports. ## What Changed - Re-read receipts and skip release before lock acquisition when stop proof is absent, the copy is superseded, or cleanup is complete. Keep the same checks inside the lock. - Recover unavailable remote copies only after exact destruction proof. Preserve their unavailable state, errors, candidate hash, and candidate bytes. Keep unavailable local copies and their uncollected edits unchanged. - Store destruction-only cleanup authority with the stop proof. Later cleanup honors it after a lost database response or restart, including when a transport remains cached. - Defer failed or unproven cleanup with bounded batches and a retry delay. Keep failed cleanup visible in logs and its receipt. - Serialize preparation of an existing run with cleanup. Fresh run preparation keeps its existing admission path. - Cover held locks, receipt scope, delayed proof, batch fairness, lost update responses, cached transports, and concurrent same-run preparation with database regressions. ## Verification - Focused directory, legacy instruction-copy, shared lock, and bounded diagnostic suites: 169 tests passed across four files. - `pnpm -r typecheck`: passed on the final source. - `pnpm build`: passed on the final source. - Completed all selected local `pnpm test:run` groups: 733 general server suites, 149 serialized suites, and 14 workspace projects. There are 13 known macOS `EACCES` failures in the unchanged runtime skill cache tests. Their exact signatures match earlier clean-base results, and the cache source and test blobs match both that base and this PR base (existing fix: #14290). One CLI import test timed out under concurrent load; its full file passed separately (17 tests). Broad coverage began before the review corrections; the final source has the focused 169-test run, typecheck, and build. This is a local verification limit, not a passing full local suite. - `git diff --check` and local Gitleaks plus private-identifier/PII diff scans passed. - Independent review of the final source found no remaining actionable issue. Its 17 targeted tests cover crash recovery, cached transports, same-run preparation, real local edit preservation, proof scope, and batch fairness. The main focused run also covers contained scheduling failures. - Final commit `35a24085f7`: Greptile 5/5 with no recommendations and zero unresolved review threads. - Final commit `35a24085f7`: all 53 checks passed, including Canary Dry Run and the security scan; two visual checks were intentionally skipped. The workspace shard passed on retry after GitHub reported that its first runner lost communication. An earlier Canary runner shut down after the release dry run passed. Neither interruption recorded an application assertion failure; the exact final-head checks are now green. ## Risks - This repairs cleanup ordering and recovery eligibility. It does not repair an ambiguous legacy lock owner or restore unsaved files. Actual collection still fails visibly when its lock cannot be acquired. - An unavailable remote copy is recovered only after exact destruction proof. A stopped but retained environment stays protected; recovery does not execute a command that could restart it. - Unavailable local copies with later stop proof still retain potentially uncollected edits. A general local recollection or reclamation policy remains outside this change. - Existing-run preparation now waits for the same lock as cleanup. The fresh-run path is unchanged. - The cleanup mode is stored in the existing private receipt JSON. No schema migration or public API change is required. - No deployment, task replay, or runtime lock deletion was performed. ## Model Used OpenAI GPT-6 (Codex), with reasoning, repository tools, and test execution. The runtime does not expose a more specific model suffix or context-window size. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub references) - [x] My branch name describes the change and contains no internal Paperclip ticket id - [ ] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
0f9e9be408
commit
f2e0f19630
6 files changed
+335
-17
No files matched your search
@@ -172,6 +172,16 @@ New current-file bytes are not database revision rows. Old instruction-only
|
||||
candidates are retained solely for upgrade compatibility.
|
||||
|
||||
Crash recovery can collect a stopped working copy without starting a model.
|
||||
Cleanup does not wait for a directory lock before process-stop proof exists, or
|
||||
after the copy is superseded or cleanup is complete. An unavailable copy keeps
|
||||
its failed-save receipt. Recovery can later clean an unavailable remote copy
|
||||
after destruction of its exact lease and executes no remote command. An
|
||||
unavailable local copy can still contain uncollected edits; this cleanup path
|
||||
preserves those bytes even if local stop proof arrives later.
|
||||
Deferred cleanup retries after a
|
||||
delay so one blocked copy does not prevent other copies from being cleaned.
|
||||
Re-preparing an existing run uses the same lock as cleanup and rechecks its
|
||||
receipt under that lock. Preparing a new run keeps its separate admission path.
|
||||
Missing stop proof or lost remote bytes produce a visible diagnostic, never a
|
||||
save receipt. An interrupted apply can replay its changed files with the same
|
||||
last-sync-wins rule. Cleanup resumes for terminal runs; no copy is retained as
|
||||
|
||||
@@ -206,6 +206,7 @@ const LOCK_DIAGNOSTIC_READ_TIMEOUT_MS = 100;
|
||||
const activeDirectoryMergeLocks = new Set<string>();
|
||||
const MAX_LOCK_DIAGNOSTIC_AGE_MS = 7 * 24 * 60 * 60 * 1000;
|
||||
export type DirectoryMergeLockOperation =
|
||||
| "agent_directory_prepare"
|
||||
| "agent_directory_release"
|
||||
| "agent_directory_collect"
|
||||
| "agent_directory_checkpoint"
|
||||
|
||||
@@ -19,6 +19,8 @@ import { buildNativeRuntimeContext } from "../services/native-runtime/runtime-co
|
||||
import type { EnvironmentRuntimeService } from "../services/environment-runtime.js";
|
||||
import { remoteTerminationReceipt } from "../services/remote-execution-termination.js";
|
||||
import { AgentDirectoryReuseInvalidatedError, agentDirectoryWorkingCopyService } from "../services/agent-directory-working-copies.js";
|
||||
import { withDirectoryMergeLock } from "@paperclipai/adapter-utils/workspace-restore-merge";
|
||||
import { heartbeatRunEvents } from "@paperclipai/db";
|
||||
|
||||
describe("persistent agent directories", () => {
|
||||
let database: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>>;
|
||||
@@ -729,4 +731,240 @@ describe("persistent agent directories", () => {
|
||||
expect(await fs.readFile(path.join(root, "note.txt"), "utf8")).toBe("one write");
|
||||
});
|
||||
|
||||
it.each(["prepared", "unavailable", "warm_saved", "superseded", "saved", "unchanged", "resolved"])(
|
||||
"does not acquire a held directory lock for a no-op %s release", async state => {
|
||||
const copy = await run();
|
||||
const stopped = ["superseded", "saved", "unchanged", "resolved"].includes(state);
|
||||
await db.update(agentInstructionWorkingCopies).set({ state, processStoppedAt: stopped ? new Date() : null,
|
||||
receipt: { ...copy.receipt, ...(stopped ? { cleanupPending: false } : {}) } })
|
||||
.where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
const before = await copies.get(companyId, copy.runId);
|
||||
await withDirectoryMergeLock(path.resolve(copy.localRoot, "../../.."), async () => {
|
||||
// Any attempted nested acquisition fails promptly instead of hanging the test.
|
||||
let now = Date.now();
|
||||
const clock = vi.spyOn(Date, "now").mockImplementation(() => now += 31_000);
|
||||
try { await copies.release(companyId, copy.runId); } finally { clock.mockRestore(); }
|
||||
expect(await copies.get(companyId, copy.runId)).toEqual(before);
|
||||
expect(await fs.readFile(path.join(copy.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
});
|
||||
});
|
||||
|
||||
async function unavailableRemote(options: { state?: "stopped" | "destroyed"; receiptMismatch?: string } = {}) {
|
||||
const copy = await run();
|
||||
const environmentId = randomUUID(), leaseId = randomUUID(), remoteCwd = "/fixture/cleanup-task";
|
||||
const lease = { id: leaseId, companyId, environmentId, heartbeatRunId: copy.runId, provider: "daytona", providerLeaseId: "cleanup-fixture" };
|
||||
await db.insert(environments).values({ id: environmentId, name: environmentId, driver: "sandbox" });
|
||||
const receipt = options.state ? remoteTerminationReceipt(lease, { providerLeaseId: lease.providerLeaseId, state: options.state })! : undefined;
|
||||
if (receipt && options.receiptMismatch) (receipt as Record<string, unknown>)[options.receiptMismatch] = randomUUID();
|
||||
await db.insert(environmentLeases).values({ ...lease, ...(receipt ? { status: "released", releasedAt: new Date(), cleanupStatus: "success",
|
||||
metadata: { remoteExecutionTermination: receipt } } : {}) });
|
||||
await db.update(heartbeatRuns).set({ status: "succeeded" }).where(eq(heartbeatRuns.id, copy.runId));
|
||||
await db.update(agentInstructionWorkingCopies).set({ state: "unavailable", errorCode: "INSTRUCTION_COLLECTION_UNAVAILABLE",
|
||||
errorMessage: "Fixture could not collect files", candidateHash: "preserved-candidate", candidateBase64: Buffer.from("preserved bytes").toString("base64"), attempts: 3,
|
||||
location: `remote:${environmentId}`, executionRoot: path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", agentId, copy.runId),
|
||||
receipt: { ...copy.receipt, cleanup: { leaseId, remoteCwd } } }).where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
return { copy: (await copies.get(companyId, copy.runId))!, lease, environmentId };
|
||||
}
|
||||
|
||||
it.each([undefined, "stopped", "companyId", "runId", "leaseId", "provider", "providerLeaseId"])(
|
||||
"defers unavailable cleanup without exact destruction proof (%s)", async proof => {
|
||||
const { copy } = await unavailableRemote({ state: proof === undefined ? undefined : proof === "stopped" ? "stopped" : "destroyed",
|
||||
receiptMismatch: proof && proof !== "stopped" ? proof : undefined });
|
||||
const execute = vi.fn();
|
||||
copies = agentInstructionWorkingCopyService(db, { environmentRuntime: { execute } as unknown as EnvironmentRuntimeService });
|
||||
const before = Date.now();
|
||||
await withDirectoryMergeLock(path.resolve(copy.localRoot, "../../.."), () => copies.recoverStopped());
|
||||
const deferred = (await copies.get(companyId, copy.runId))!;
|
||||
expect(deferred).toMatchObject({ state: "unavailable", processStoppedAt: null, attempts: 3,
|
||||
candidateHash: copy.candidateHash, candidateBase64: copy.candidateBase64, errorCode: copy.errorCode, errorMessage: copy.errorMessage, receipt: copy.receipt });
|
||||
expect(deferred.nextAttemptAt!.getTime()).toBeGreaterThanOrEqual(before + 30_000);
|
||||
expect(await copies.recoverStopped()).toBe(0);
|
||||
expect(await fs.readFile(path.join(copy.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
expect(execute).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("cleans an unavailable copy after exact destruction without claiming a save or executing remotely", async () => {
|
||||
const { copy, environmentId } = await unavailableRemote({ state: "destroyed" });
|
||||
await fs.writeFile(path.join(copy.localRoot, "unsaved.txt"), "unsaved edit");
|
||||
await db.delete(environments).where(eq(environments.id, environmentId));
|
||||
const execute = vi.fn();
|
||||
const shell = vi.spyOn(executionTargetTools, "runAdapterExecutionTargetShellCommand").mockRejectedValue(new Error("must not execute"));
|
||||
copies = agentInstructionWorkingCopyService(db, { environmentRuntime: { execute } as unknown as EnvironmentRuntimeService });
|
||||
try {
|
||||
await copies.recoverStopped();
|
||||
const recovered = (await copies.get(companyId, copy.runId))!;
|
||||
expect(recovered).toMatchObject({ state: "unavailable", candidateHash: copy.candidateHash, candidateBase64: copy.candidateBase64, attempts: 3,
|
||||
errorCode: copy.errorCode, errorMessage: copy.errorMessage, receipt: { cleanupPending: false } });
|
||||
expect(recovered.processStoppedAt).toBeInstanceOf(Date);
|
||||
expect(recovered.nextAttemptAt).toBeNull();
|
||||
await expect(fs.stat(copy.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
await expect(fs.stat(path.join(root, "unsaved.txt"))).rejects.toMatchObject({ code: "ENOENT" });
|
||||
expect(await fs.readFile(path.join(root, entryFile), "utf8")).toBe(initial);
|
||||
expect(execute).not.toHaveBeenCalled(); expect(shell).not.toHaveBeenCalled();
|
||||
} finally { shell.mockRestore(); }
|
||||
});
|
||||
|
||||
it.each([false, true])("preserves uncollected local edits and their unavailable receipt even after local stop proof (%s)", async stopped => {
|
||||
const copy = await run();
|
||||
await fs.writeFile(path.join(copy.localRoot, "unsaved-local.txt"), "recoverable local edit");
|
||||
await db.update(heartbeatRuns).set({ status: "failed", runtimeMode: "native" }).where(eq(heartbeatRuns.id, copy.runId));
|
||||
await copies.reportUnavailable(companyId, copy.runId);
|
||||
if (stopped) await db.insert(heartbeatRunEvents).values({ companyId, runId: copy.runId, agentId, seq: 1, eventType: "native.local_process_stopped", stream: "system", level: "info", message: "fixture stop" });
|
||||
const before = (await copies.get(companyId, copy.runId))!;
|
||||
expect(before).toMatchObject({ state: "unavailable", attempts: 0, processStoppedAt: null });
|
||||
await copies.recoverStopped();
|
||||
const patch = vi.fn(async () => { throw new Error("local copy must not be changed"); });
|
||||
await agentDirectoryWorkingCopyService(db, copies.get, patch).recoverUnavailable(before);
|
||||
expect(patch).not.toHaveBeenCalled();
|
||||
expect(await copies.get(companyId, copy.runId)).toEqual(before);
|
||||
expect(await fs.readFile(path.join(copy.localRoot, "unsaved-local.txt"), "utf8")).toBe("recoverable local edit");
|
||||
await expect(fs.stat(path.join(root, "unsaved-local.txt"))).rejects.toMatchObject({ code: "ENOENT" });
|
||||
});
|
||||
|
||||
it("serializes same-run preparation after recovery commits stop proof and before it removes bytes", async () => {
|
||||
const { copy, lease } = await unavailableRemote({ state: "destroyed" });
|
||||
const remoteCwd = "/fixture/cleanup-task";
|
||||
const executionTarget = { kind: "remote" as const, transport: "ssh" as const, environmentId: lease.environmentId, leaseId: lease.id, remoteCwd,
|
||||
spec: { host: "unused.invalid", port: 22, username: "test", remoteCwd } };
|
||||
const shell = vi.spyOn(executionTargetTools, "runAdapterExecutionTargetShellCommand").mockResolvedValue({ exitCode: 0, signal: null, timedOut: false, stdout: "", stderr: "" });
|
||||
const stage = vi.spyOn(ssh, "syncDirectoryToSsh").mockResolvedValue(undefined);
|
||||
let stopCommitted!: () => void, finishCleanup!: () => void;
|
||||
const committed = new Promise<void>(resolve => { stopCommitted = resolve; });
|
||||
const proceed = new Promise<void>(resolve => { finishCleanup = resolve; });
|
||||
const directories = agentDirectoryWorkingCopyService(db, copies.get, async (row, values) => {
|
||||
const [updated] = await db.update(agentInstructionWorkingCopies).set(values).where(eq(agentInstructionWorkingCopies.runId, row.runId)).returning();
|
||||
if (values.processStoppedAt) { stopCommitted(); await proceed; }
|
||||
return updated!;
|
||||
});
|
||||
const recovering = directories.recoverUnavailable(copy);
|
||||
await committed;
|
||||
let prepareSettled = false;
|
||||
const preparing = copies.prepare({ ...target(), runId: copy.runId, cwd: home, target: executionTarget }).then(row => { prepareSettled = true; return row; });
|
||||
try {
|
||||
try {
|
||||
await new Promise(resolve => setTimeout(resolve, 100));
|
||||
expect(prepareSettled).toBe(false);
|
||||
expect((await copies.get(companyId, copy.runId))?.state).toBe("unavailable");
|
||||
} finally { finishCleanup(); await recovering; }
|
||||
const prepared = await preparing;
|
||||
expect(prepared).toMatchObject({ state: "prepared", processStoppedAt: null });
|
||||
expect(await fs.readFile(path.join(copy.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
// A stale cleanup callback now observes the live lifecycle and does nothing.
|
||||
await directories.release(copy);
|
||||
expect(await fs.readFile(path.join(copy.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
} finally { shell.mockRestore(); stage.mockRestore(); }
|
||||
});
|
||||
|
||||
it("keeps stopped pending cleanup behind a held lock and lets recovery proceed to other agents", async () => {
|
||||
const blocked = await run();
|
||||
await db.update(agentInstructionWorkingCopies).set({ state: "unavailable", processStoppedAt: new Date() }).where(eq(agentInstructionWorkingCopies.runId, blocked.runId));
|
||||
// A distinct canonical agent has its own lifecycle lock.
|
||||
const firstAgent = agentId, firstRoot = root;
|
||||
agentId = randomUUID();
|
||||
root = resolveManagedInstructionsRoot({ ...target(), id: agentId, name: "Other", adapterConfig: {} });
|
||||
await db.insert(agents).values({ id: agentId, companyId, name: "Other", adapterConfig: { instructionsBundleMode: "managed", instructionsRootPath: root, instructionsEntryFile: entryFile } });
|
||||
await db.insert(companyMemberships).values({ companyId, principalType: "agent", principalId: agentId, membershipRole: "member" });
|
||||
await db.update(principalPermissionGrants).set({ scope: { agentIds: [firstAgent, agentId] } }).where(eq(principalPermissionGrants.companyId, companyId));
|
||||
await fs.mkdir(path.dirname(path.join(root, entryFile)), { recursive: true }); await fs.writeFile(path.join(root, entryFile), initial);
|
||||
const other = await run();
|
||||
await db.update(agentInstructionWorkingCopies).set({ state: "unavailable", processStoppedAt: new Date() }).where(eq(agentInstructionWorkingCopies.runId, other.runId));
|
||||
agentId = firstAgent; root = firstRoot;
|
||||
await withDirectoryMergeLock(path.resolve(blocked.localRoot, "../../.."), async () => {
|
||||
let now = Date.now();
|
||||
const clock = vi.spyOn(Date, "now").mockImplementation(() => now += 31_000);
|
||||
try {
|
||||
await expect(copies.release(companyId, blocked.runId)).rejects.toMatchObject({ code: "ERR_WORKSPACE_RESTORE_LOCK_TIMEOUT" });
|
||||
await copies.recoverCaptured();
|
||||
} finally { clock.mockRestore(); }
|
||||
expect((await copies.get(companyId, blocked.runId))?.nextAttemptAt).toBeInstanceOf(Date);
|
||||
expect(await fs.readFile(path.join(blocked.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
await expect(fs.stat(other.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
});
|
||||
});
|
||||
|
||||
it.each(["environment", "path", "lease-run"])("rejects destruction outside the registered copy binding (%s)", async mismatch => {
|
||||
const { copy, lease } = await unavailableRemote({ state: "destroyed" });
|
||||
if (mismatch === "environment") await db.update(agentInstructionWorkingCopies).set({ location: `remote:${randomUUID()}` }).where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
if (mismatch === "path") await db.update(agentInstructionWorkingCopies).set({ executionRoot: "/fixture/foreign" }).where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
if (mismatch === "lease-run") {
|
||||
const other = await run();
|
||||
const moved = { ...lease, heartbeatRunId: other.runId };
|
||||
await db.update(environmentLeases).set({ heartbeatRunId: other.runId, metadata: { remoteExecutionTermination: remoteTerminationReceipt(moved, { providerLeaseId: moved.providerLeaseId, state: "destroyed" }) } }).where(eq(environmentLeases.id, lease.id));
|
||||
}
|
||||
await copies.recoverStopped();
|
||||
expect((await copies.get(companyId, copy.runId))?.processStoppedAt).toBeNull();
|
||||
expect(await fs.readFile(path.join(copy.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
});
|
||||
|
||||
it("backs off an unproven batch so the next unavailable copy can be cleaned", async () => {
|
||||
const deferred: string[] = [];
|
||||
for (let i = 0; i < 20; i++) {
|
||||
const { copy } = await unavailableRemote();
|
||||
deferred.push(copy.runId);
|
||||
await db.update(agentInstructionWorkingCopies).set({ updatedAt: new Date(i) }).where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
}
|
||||
const { copy: ready } = await unavailableRemote({ state: "destroyed" });
|
||||
await db.update(agentInstructionWorkingCopies).set({ updatedAt: new Date(20) }).where(eq(agentInstructionWorkingCopies.runId, ready.runId));
|
||||
await copies.recoverStopped();
|
||||
expect((await copies.get(companyId, ready.runId))?.processStoppedAt).toBeNull();
|
||||
await copies.recoverStopped();
|
||||
expect((await copies.get(companyId, ready.runId))?.receipt?.cleanupPending).toBe(false);
|
||||
for (const runId of deferred) expect(await copies.get(companyId, runId)).toMatchObject({ state: "unavailable", processStoppedAt: null, attempts: 3 });
|
||||
});
|
||||
|
||||
it("continues cleanup when both retry-reservation writes fail for the first copy", async () => {
|
||||
const first = await run(), second = await run();
|
||||
for (const [index, copy] of [first, second].entries()) {
|
||||
await db.update(agentInstructionWorkingCopies).set({ state: "unavailable", processStoppedAt: new Date(), updatedAt: new Date(index),
|
||||
errorCode: "INSTRUCTION_COLLECTION_UNAVAILABLE", candidateHash: "preserved", candidateBase64: Buffer.from("candidate").toString("base64") })
|
||||
.where(eq(agentInstructionWorkingCopies.runId, copy.runId));
|
||||
}
|
||||
const before = await copies.get(companyId, first.runId);
|
||||
const failReservation = () => { throw new Error("fixture retry-reservation write failure"); };
|
||||
const update = vi.spyOn(db, "update").mockImplementationOnce(failReservation).mockImplementationOnce(failReservation);
|
||||
try { await copies.recoverCaptured(); } finally { update.mockRestore(); }
|
||||
expect(await copies.get(companyId, first.runId)).toEqual(before);
|
||||
expect(await fs.readFile(path.join(first.localRoot, entryFile), "utf8")).toBe(initial);
|
||||
expect(await copies.get(companyId, second.runId)).toMatchObject({ state: "unavailable", errorCode: "INSTRUCTION_COLLECTION_UNAVAILABLE", candidateHash: "preserved", receipt: { cleanupPending: false } });
|
||||
await expect(fs.stat(second.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
});
|
||||
|
||||
it.each([false, true])("does not execute a cached transport after destruction, including an ambiguous stop-proof commit (%s)", async interrupted => {
|
||||
const runId = randomUUID(), environmentId = randomUUID(), leaseId = randomUUID(), remoteCwd = "/fixture/cached-task";
|
||||
const lease = { id: leaseId, companyId, environmentId, heartbeatRunId: runId, provider: "daytona", providerLeaseId: "cached-fixture" };
|
||||
await db.insert(heartbeatRuns).values({ id: runId, companyId, agentId, invocationSource: "on_demand", responsibleUserId: userId });
|
||||
await db.insert(environments).values({ id: environmentId, name: environmentId, driver: "sandbox" });
|
||||
await db.insert(environmentLeases).values(lease);
|
||||
const execute = vi.fn(), restoreWorkspace = vi.fn(), cleanupWorkspaceSnapshot = vi.fn().mockResolvedValue(undefined);
|
||||
const executionTarget = { kind: "remote" as const, transport: "sandbox" as const, environmentId, leaseId, remoteCwd, runner: { execute } };
|
||||
const shell = vi.spyOn(executionTargetTools, "runAdapterExecutionTargetShellCommand").mockResolvedValue({ exitCode: 0, signal: null, timedOut: false, stdout: "", stderr: "" });
|
||||
const stage = vi.spyOn(executionTargetTools, "prepareAdapterExecutionTargetRuntime").mockImplementation(async () => ({ target: executionTarget,
|
||||
workspaceRemoteDir: path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", agentId, runId), runtimeRootDir: null,
|
||||
assetDirs: {}, additionalSourceDirs: {}, additionalSourceFailures: [], workspaceSyncSnapshot: null, restoreWorkspace, cleanupWorkspaceSnapshot }));
|
||||
try {
|
||||
const copy = (await copies.prepare({ ...target(), runId, cwd: home, target: executionTarget }))!;
|
||||
await copies.reportUnavailable(companyId, runId);
|
||||
await copies.release(companyId, runId);
|
||||
expect(cleanupWorkspaceSnapshot).not.toHaveBeenCalled();
|
||||
shell.mockClear(); stage.mockClear();
|
||||
await db.update(heartbeatRuns).set({ status: "failed" }).where(eq(heartbeatRuns.id, runId));
|
||||
await db.update(environmentLeases).set({ status: "released", releasedAt: new Date(), cleanupStatus: "success",
|
||||
metadata: { remoteExecutionTermination: remoteTerminationReceipt(lease, { providerLeaseId: lease.providerLeaseId, state: "destroyed" }) } }).where(eq(environmentLeases.id, leaseId));
|
||||
if (interrupted) {
|
||||
const directories = agentDirectoryWorkingCopyService(db, copies.get, async (row, values) => {
|
||||
const [updated] = await db.update(agentInstructionWorkingCopies).set(values).where(eq(agentInstructionWorkingCopies.runId, row.runId)).returning();
|
||||
if (values.processStoppedAt) throw new Error("fixture lost update response after commit");
|
||||
return updated!;
|
||||
});
|
||||
await expect(directories.recoverUnavailable((await copies.get(companyId, runId))!)).rejects.toThrow("lost update response");
|
||||
expect((await copies.get(companyId, runId))?.receipt?.cleanupDestroyedOnly).toBe(true);
|
||||
await copies.recoverCaptured();
|
||||
} else await copies.recoverStopped();
|
||||
expect(await copies.get(companyId, runId)).toMatchObject({ state: "unavailable", errorCode: "INSTRUCTION_COLLECTION_UNAVAILABLE", receipt: { cleanupPending: false } });
|
||||
await expect(fs.stat(copy.localRoot)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
expect(cleanupWorkspaceSnapshot).toHaveBeenCalledOnce();
|
||||
expect(execute).not.toHaveBeenCalled(); expect(restoreWorkspace).not.toHaveBeenCalled(); expect(shell).not.toHaveBeenCalled(); expect(stage).not.toHaveBeenCalled();
|
||||
} finally { shell.mockRestore(); stage.mockRestore(); }
|
||||
});
|
||||
|
||||
});
|
||||
@@ -66,7 +66,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
workspaceBaseline: baseline(row), workspaceGitSnapshot: null, workspaceFileMode: "all",
|
||||
workspaceExclude: [".paperclip-runtime", ".paperclip-runtime/**"] });
|
||||
}
|
||||
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; warm?: boolean; reuseRunId?: string; onWarmHandoff?: (copy: Copy) => void }) {
|
||||
async function prepareCopy(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; warm?: boolean; reuseRunId?: string; onWarmHandoff?: (copy: Copy) => void }, serialHeld = false): Promise<Copy | null> {
|
||||
const [agent] = await db.select().from(agents).where(and(eq(agents.id, input.agentId), eq(agents.companyId, input.companyId)));
|
||||
if (!agent) throw notFound("Agent not found");
|
||||
if (agentInstructionsBundleMode(agent) !== "managed") return null;
|
||||
@@ -120,6 +120,13 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
const location = input.target?.kind === "remote" ? `remote:${input.target.environmentId ?? ""}` : "local";
|
||||
let row = await get(input.companyId, input.runId);
|
||||
if (row && (row.localRoot !== localRoot || row.executionRoot !== executionRoot || row.agentId !== input.agentId || row.location !== location)) throw conflict("Agent directory belongs to a different execution environment");
|
||||
if (row && !serialHeld) {
|
||||
// Restart preparation can replace a completed copy under the same run ID.
|
||||
// Join cleanup's lock before resetting its receipt or creating new bytes.
|
||||
// Interrupted admission may have registered the row before making this root.
|
||||
await fs.mkdir(path.resolve(localRoot, "../../.."), { recursive: true, mode: 0o700 });
|
||||
return serial(row, () => prepareCopy(input, true), "agent_directory_prepare");
|
||||
}
|
||||
if (row && !completed.has(row.state) && row.state !== "preparing") {
|
||||
if (input.target?.kind === "remote" && !transports.has(key(row))) transports.set(key(row), await transport(row, input.target, true));
|
||||
return row;
|
||||
@@ -289,17 +296,19 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
await release(row);
|
||||
return (await get(row.companyId, row.runId))!;
|
||||
}
|
||||
function releaseIsNoop(row: Copy) {
|
||||
return row.state === "superseded" || !row.processStoppedAt ||
|
||||
completed.has(row.state) && row.receipt?.cleanupPending === false;
|
||||
}
|
||||
async function release(row: Copy, target?: AdapterExecutionTarget | null) {
|
||||
if (row.state === "superseded" || !await owns(row) || row.receipt?.warm === true && !row.processStoppedAt) return;
|
||||
// Collection and environment teardown can both release the same copy.
|
||||
// A compact successful cleanup receipt is final, even after restart.
|
||||
if (completed.has(row.state) && row.processStoppedAt && row.receipt?.cleanupPending === false) return;
|
||||
if (releaseIsNoop(row) || !await owns(row)) return;
|
||||
const destroyedOnly = row.receipt?.cleanupDestroyedOnly === true;
|
||||
const runtime = transports.get(key(row));
|
||||
transports.delete(key(row));
|
||||
let cleanupPending = false;
|
||||
await runtime?.cleanupWorkspaceSnapshot?.().catch(() => { cleanupPending = true; });
|
||||
const cleanupTarget = target ?? runtime?.target;
|
||||
if (row.processStoppedAt && completed.has(row.state) && cleanupTarget?.kind === "remote") {
|
||||
if (!destroyedOnly && row.processStoppedAt && completed.has(row.state) && cleanupTarget?.kind === "remote") {
|
||||
const expected = path.posix.join(cleanupTarget.remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, String(row.receipt?.directoryRunId ?? row.runId));
|
||||
if (row.executionRoot !== expected) throw new Error("Agent directory cleanup path changed");
|
||||
const quoted = `'${expected.replaceAll("'", `'"'"'`)}'`;
|
||||
@@ -309,7 +318,7 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
} else if (row.processStoppedAt && completed.has(row.state) && row.location !== "local") {
|
||||
// Rebind only the original host-owned lease. Never acquire a replacement
|
||||
// sandbox or persist transport credentials in a working-copy receipt.
|
||||
cleanupPending ||= !await cleanupRemoteAfterRestart(row).catch(() => false);
|
||||
cleanupPending ||= !await cleanupRemoteAfterRestart(row, destroyedOnly).catch(() => false);
|
||||
}
|
||||
if (completed.has(row.state) && row.processStoppedAt) {
|
||||
try { await fs.rm(path.dirname(row.localRoot), { recursive: true, force: true }); }
|
||||
@@ -320,11 +329,44 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
}
|
||||
await patch(row, { receipt: { schema: AGENT_FILES_CONTRACT, state: row.state, appliedCandidateHash: row.candidateHash, storageWarning: row.receipt?.storageWarning ?? null, cleanupPending,
|
||||
...(row.receipt?.directoryRunId ? { directoryRunId: row.receipt.directoryRunId } : {}),
|
||||
...(cleanupPending ? { cleanup: row.receipt?.cleanup } : {}) },
|
||||
...(cleanupPending ? { cleanup: row.receipt?.cleanup, ...(destroyedOnly ? { cleanupDestroyedOnly: true } : {}) } : {}) },
|
||||
nextAttemptAt: cleanupPending ? new Date(Date.now() + 30_000) : null });
|
||||
}
|
||||
}
|
||||
async function cleanupRemoteAfterRestart(row: Copy): Promise<boolean> {
|
||||
async function destroyedRemoteCopy(row: Copy): Promise<boolean> {
|
||||
const cleanup = row.receipt?.cleanup as { leaseId?: string; remoteCwd?: string } | undefined;
|
||||
const leases = await db.select().from(environmentLeases).where(and(
|
||||
eq(environmentLeases.companyId, row.companyId), eq(environmentLeases.heartbeatRunId, row.runId),
|
||||
cleanup?.leaseId ? eq(environmentLeases.id, cleanup.leaseId)
|
||||
: eq(environmentLeases.environmentId, row.location.slice("remote:".length)),
|
||||
));
|
||||
if (leases.length !== 1) return false;
|
||||
const lease = leases[0]!;
|
||||
if (lease.environmentId !== null && row.location !== `remote:${lease.environmentId}`) return false;
|
||||
const remoteCwd = cleanup?.remoteCwd ?? lease.metadata?.remoteCwd;
|
||||
return typeof remoteCwd === "string" &&
|
||||
row.executionRoot === path.posix.join(remoteCwd, ".paperclip-runtime", "agent-files", row.agentId, String(row.receipt?.directoryRunId ?? row.runId)) &&
|
||||
hasRemoteTerminationReceipt(lease) && (lease.metadata?.remoteExecutionTermination as { state?: string }).state === "destroyed";
|
||||
}
|
||||
async function recoverUnavailable(row: Copy) {
|
||||
const eligible = (current: Copy) => isAgentDirectoryCopy(current) && current.location.startsWith("remote:") &&
|
||||
current.state === "unavailable" && !current.processStoppedAt;
|
||||
// Do not wait for a directory lock when no cleanup can be authorized.
|
||||
if (!eligible(row) || !await destroyedRemoteCopy(row)) return;
|
||||
await serial(row, async current => {
|
||||
if (!eligible(current) || !await destroyedRemoteCopy(current)) return;
|
||||
// The no-execute authority must survive a crash after this update. Every
|
||||
// later release, including one with a cached transport, honors it.
|
||||
const confirmed = await patch(current, { processStoppedAt: new Date(),
|
||||
receipt: { ...current.receipt, cleanupDestroyedOnly: true } });
|
||||
// A concurrent preparation can win the receipt CAS. Its bytes stay live.
|
||||
if (confirmed.state === "unavailable" && confirmed.processStoppedAt) await release(confirmed);
|
||||
}, "agent_directory_release");
|
||||
}
|
||||
async function cleanupRemoteAfterRestart(row: Copy, destroyedOnly = false): Promise<boolean> {
|
||||
// Newly recovered loss receipts permit no command that could wake a stopped
|
||||
// provider. Recheck the exact destruction binding even after lock acquisition.
|
||||
if (destroyedOnly) return destroyedRemoteCopy(row);
|
||||
const cleanup = row.receipt?.cleanup as { leaseId?: string; remoteCwd?: string } | undefined;
|
||||
// Older preview receipts did not record the lease ID. Accept only a unique
|
||||
// lease still bound to this run and environment in that compatibility case.
|
||||
@@ -355,9 +397,15 @@ export function agentDirectoryWorkingCopyService(db: Db, get: (companyId: string
|
||||
return fn(current);
|
||||
}, process.env, operation);
|
||||
}
|
||||
return { prepare, hasChanges, canReuse,
|
||||
return { prepare: (input: Parameters<typeof prepareCopy>[0]) => prepareCopy(input), hasChanges, canReuse, recoverUnavailable,
|
||||
checkpointWarm: (row: Copy, target?: AdapterExecutionTarget | null) => serial(row, current => current.receipt?.warm === true ? checkpoint(current, target) : Promise.resolve(current), "agent_directory_checkpoint"),
|
||||
collectStopped: (row: Copy, target?: AdapterExecutionTarget | null) => serial(row, current => collectStopped(current, target), "agent_directory_collect"),
|
||||
release: (row: Copy) => serial(row, release, "agent_directory_release"),
|
||||
release: async (row: Copy) => {
|
||||
const current = await get(row.companyId, row.runId);
|
||||
// This check changes no receipt or bytes. A concurrent stop merely leaves
|
||||
// cleanup to its owner or the recovery sweep; it grants no new authority.
|
||||
if (!current || releaseIsNoop(current)) return;
|
||||
await serial(current, release, "agent_directory_release");
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -20,6 +20,7 @@ import { hasNativeLocalProcessStop } from "./native-local-process-stop.js";
|
||||
import { remoteExecutionHasStopped } from "./remote-execution-termination.js";
|
||||
import type { AuthorizationActor } from "./authorization.js";
|
||||
import type { EnvironmentRuntimeService } from "./environment-runtime.js";
|
||||
import { logger } from "../middleware/logger.js";
|
||||
|
||||
const execFile = promisify(execFileCallback);
|
||||
type Copy = typeof copies.$inferSelect;
|
||||
@@ -71,6 +72,21 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
|
||||
return { type: "agent", companyId: row.companyId, agentId: row.agentId, runId: row.runId, onBehalfOfUserId: row.responsibleUserId };
|
||||
}
|
||||
const directories = agentDirectoryWorkingCopyService(db, get, patch, options.environmentRuntime);
|
||||
async function recoverDirectoryCleanup(row: Copy, cleanup: () => Promise<unknown>) {
|
||||
// Reserve a later retry before I/O, including lock waits. A blocked owner
|
||||
// must not occupy every sweep or prevent other copies from being reclaimed.
|
||||
try {
|
||||
await patch(row, { nextAttemptAt: new Date(Date.now() + 30_000) });
|
||||
await cleanup();
|
||||
}
|
||||
catch (error) {
|
||||
try { await patch(row, { nextAttemptAt: new Date(Date.now() + 30_000) }); }
|
||||
catch (retryError) {
|
||||
logger.warn({ err: retryError, runId: row.runId }, "Agent file cleanup retry could not be scheduled");
|
||||
}
|
||||
logger.warn({ err: error, runId: row.runId }, "Agent file cleanup deferred; save receipt unchanged");
|
||||
}
|
||||
}
|
||||
async function prepare(input: { companyId: string; agentId: string; runId: string; target?: AdapterExecutionTarget | null; cwd: string; legacy?: boolean; warm?: boolean; reuseRunId?: string; onWarmHandoff?: (copy: Copy) => void }) {
|
||||
const existing = await get(input.companyId, input.runId);
|
||||
if (isAgentDirectoryCopy(existing) || (!existing && !input.legacy)) return directories.prepare(input);
|
||||
@@ -288,7 +304,7 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
|
||||
or(sql`${copies.receipt} ? 'baseline'`, sql`${copies.receipt}->>'cleanupPending' = 'true'`),
|
||||
or(isNull(copies.nextAttemptAt), lte(copies.nextAttemptAt, new Date())),
|
||||
)).orderBy(asc(copies.updatedAt)).limit(20);
|
||||
for (const row of cleanup) await directories.release(row);
|
||||
for (const row of cleanup) await recoverDirectoryCleanup(row, () => directories.release(row));
|
||||
return pending.length;
|
||||
}
|
||||
|
||||
@@ -298,13 +314,18 @@ export function agentInstructionWorkingCopyService(db: Db, options: { environmen
|
||||
async function recoverStopped() {
|
||||
const pending = await db.select({ copy: copies, runtimeMode: heartbeatRuns.runtimeMode }).from(copies)
|
||||
.innerJoin(heartbeatRuns, and(eq(heartbeatRuns.companyId, copies.companyId), eq(heartbeatRuns.id, copies.runId)))
|
||||
.where(and(or(inArray(copies.state, ["prepared", "pending_collection", "warm_saved"]),
|
||||
.where(and(or(and(or(inArray(copies.state, ["prepared", "pending_collection", "warm_saved"]),
|
||||
and(eq(copies.state, "preparing"), sql`${copies.receipt}->>'schema' = 'paperclip.agent-files.v1'`)),
|
||||
lte(copies.attempts, MAX_COLLECTION_ATTEMPTS - 1)),
|
||||
and(eq(copies.state, "unavailable"), isNull(copies.processStoppedAt),
|
||||
sql`${copies.location} like 'remote:%'`,
|
||||
sql`${copies.receipt}->>'schema' = 'paperclip.agent-files.v1'`)),
|
||||
inArray(heartbeatRuns.status, ["succeeded", "failed", "cancelled", "timed_out", "interrupted"]),
|
||||
or(isNull(copies.nextAttemptAt), lte(copies.nextAttemptAt, new Date())),
|
||||
lte(copies.attempts, MAX_COLLECTION_ATTEMPTS - 1))).orderBy(asc(copies.updatedAt)).limit(20);
|
||||
or(isNull(copies.nextAttemptAt), lte(copies.nextAttemptAt, new Date())))).orderBy(asc(copies.updatedAt)).limit(20);
|
||||
for (const { copy: row, runtimeMode } of pending) {
|
||||
if (row.state === "preparing" && isAgentDirectoryCopy(row)) {
|
||||
if (row.state === "unavailable" && isAgentDirectoryCopy(row)) {
|
||||
await recoverDirectoryCleanup(row, () => directories.recoverUnavailable(row));
|
||||
} else if (row.state === "preparing" && isAgentDirectoryCopy(row)) {
|
||||
// A provider cannot launch until preparation records "prepared". With
|
||||
// the owning run terminal, this is an interrupted staging copy only.
|
||||
await directories.release(await patch(row, { state: "unavailable", processStoppedAt: new Date(),
|
||||
|
||||
@@ -186,7 +186,7 @@ export function collectRunFailureDiagnostics(run: Run, options: RunFailureReport
|
||||
if (execution.restoreLockOwnerState === undefined && read(error, "code") === "ERR_WORKSPACE_RESTORE_LOCK_TIMEOUT") {
|
||||
const lock = read(error, "workspaceRestoreLock");
|
||||
const operation = read(lock, "operation");
|
||||
if (["agent_directory_release", "agent_directory_collect", "agent_directory_checkpoint", "agent_directory_handoff"].some(value => value === operation)) {
|
||||
if (["agent_directory_prepare", "agent_directory_release", "agent_directory_collect", "agent_directory_checkpoint", "agent_directory_handoff"].some(value => value === operation)) {
|
||||
execution.restoreLockOperation = operation as string;
|
||||
}
|
||||
const ownerState = read(lock, "ownerState");
|
||||
|
||||
Reference in new issue
Block a user