diff --git a/doc/agent-files.md b/doc/agent-files.md index 187323d2d7..2cfd67c668 100644 --- a/doc/agent-files.md +++ b/doc/agent-files.md @@ -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 diff --git a/packages/adapter-utils/src/workspace-restore-merge.ts b/packages/adapter-utils/src/workspace-restore-merge.ts index 60bdd8cc90..6b05b135b9 100644 --- a/packages/adapter-utils/src/workspace-restore-merge.ts +++ b/packages/adapter-utils/src/workspace-restore-merge.ts @@ -206,6 +206,7 @@ const LOCK_DIAGNOSTIC_READ_TIMEOUT_MS = 100; const activeDirectoryMergeLocks = new Set(); 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" diff --git a/server/src/__tests__/agent-directory-working-copies.test.ts b/server/src/__tests__/agent-directory-working-copies.test.ts index 2a768289b9..97ca234e07 100644 --- a/server/src/__tests__/agent-directory-working-copies.test.ts +++ b/server/src/__tests__/agent-directory-working-copies.test.ts @@ -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>; @@ -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)[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(resolve => { stopCommitted = resolve; }); + const proceed = new Promise(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(); } + }); + }); diff --git a/server/src/services/agent-directory-working-copies.ts b/server/src/services/agent-directory-working-copies.ts index 80efb17eb0..869a497e97 100644 --- a/server/src/services/agent-directory-working-copies.ts +++ b/server/src/services/agent-directory-working-copies.ts @@ -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 { 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 { + async function destroyedRemoteCopy(row: Copy): Promise { + 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 { + // 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[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"); + }, }; } diff --git a/server/src/services/agent-instruction-working-copies.ts b/server/src/services/agent-instruction-working-copies.ts index f31b2e3d32..c88c053772 100644 --- a/server/src/services/agent-instruction-working-copies.ts +++ b/server/src/services/agent-instruction-working-copies.ts @@ -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) { + // 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(), diff --git a/server/src/services/run-failure-diagnostics.ts b/server/src/services/run-failure-diagnostics.ts index 27ccc1f719..cf5e7ad390 100644 --- a/server/src/services/run-failure-diagnostics.ts +++ b/server/src/services/run-failure-diagnostics.ts @@ -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");