diff --git a/doc/DEVELOPING.md b/doc/DEVELOPING.md index 2615e61611..c88890e98a 100644 --- a/doc/DEVELOPING.md +++ b/doc/DEVELOPING.md @@ -6,6 +6,9 @@ This project can run fully in local dev without setting up PostgreSQL manually. For mode definitions and intended CLI behavior, see `doc/DEPLOYMENT-MODES.md`. +For sandbox file synchronization, lock ownership, and the required upgrade +procedure from directory locks, see [Workspace restore locks](workspace-restore-locks.md). + Current implementation status: - canonical model: `local_trusted` and `authenticated` (with `private/public` exposure) diff --git a/doc/workspace-restore-locks.md b/doc/workspace-restore-locks.md new file mode 100644 index 0000000000..60a00fc0b7 --- /dev/null +++ b/doc/workspace-restore-locks.md @@ -0,0 +1,56 @@ +# Workspace restore locks + +Workspace restore and agent-file collection serialize writes to each canonical +target directory. Their lock files live in +`/locks/directory-merge`, outside the writable target. All writers +must use the same instance root and a filesystem with reliable SQLite file +locking. Network filesystems that do not provide that locking are not supported. + +Each target has a permanent `.lock.sqlite` file. An open SQLite +`BEGIN IMMEDIATE` transaction holds its reserved file lock for the entire write +operation. Contenders retry without blocking the Node event loop, up to the +existing 30-second limit. Closing the connection releases the lock; the operating +system also releases it when the process exits or crashes. Independent targets +use different files and can proceed concurrently. + +This uses the same built-in `node:sqlite` dependency as workspace manifests. +See [SQLite file locking](https://www.sqlite.org/lockingv3.html) for the reserved +lock contract. No application database or schema migration is involved. + +The adjacent `.lock.owner.json` records a PID and creation time only for +bounded timeout diagnostics. Missing, stale, or incorrect metadata cannot grant +or retain ownership. PID reuse, PID namespaces, and wall-clock changes do not +decide whether a writer owns the lock. + +**Never delete, replace, or move a `.lock.sqlite` file while an instance can +write to it.** Its stable inode is part of the locking contract. Replacing it +could let two writers lock different files for the same target. Files remain +after release, including for targets that no longer exist. Their diagnostic +owner sidecars are normally removed after release. + +## Upgrade from directory locks + +Older versions created `.lock/owner.json` and checked only whether the +recorded PID existed. A server restart could reuse that PID and leave an orphaned +lock permanently protected. Those records cannot establish a process lifetime or +PID namespace, so the new implementation never guesses that a legacy holder is +dead. An existing legacy directory continues to block admission. + +Do not run old and new lock protocols concurrently against the same instance +root. An old process does not participate in the SQLite lock protocol. + +1. Drain and stop **all** old server and worker processes that can write to the + instance root, including processes on other hosts or in other containers. +2. Preserve any unfinished-run evidence needed for recovery. After all writers + are stopped, move leftover legacy `.lock/` directories to an operator + scratch directory outside the lock root. Do not remove `.lock.sqlite` files. +3. Start all writers on the new version. Verify that a run completes both the + provider turn and file collection/restore. A successful model response alone + does not prove that its local file changes were saved. + +The same drain requirement applies to rollback. Existing permanent SQLite files +can remain on disk; older versions ignore them. Do not infer that a legacy lock +is safe to remove from its age, an absent PID, or a successful task response. + +This change prevents new orphaned ownership. It cannot recover file changes +that an earlier failed collection discarded. diff --git a/packages/adapter-utils/src/directory-merge-lock.test.ts b/packages/adapter-utils/src/directory-merge-lock.test.ts new file mode 100644 index 0000000000..c05305fb8d --- /dev/null +++ b/packages/adapter-utils/src/directory-merge-lock.test.ts @@ -0,0 +1,182 @@ +import { execFile, spawn, type ChildProcess } from "node:child_process"; +import { createHash } from "node:crypto"; +import { once } from "node:events"; +import { access, link, mkdir, mkdtemp, readFile, realpath, rm, stat, symlink, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { promisify } from "node:util"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { withDirectoryMergeLock, WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE } from "./workspace-restore-merge.js"; + +describe("directory merge lock process lifetime", () => { + const loader = fileURLToPath(new URL("../../../cli/node_modules/tsx/dist/cli.mjs", import.meta.url)); + const module = fileURLToPath(new URL("./workspace-restore-merge.ts", import.meta.url)); + const directories: string[] = []; + const children: ChildProcess[] = []; + afterEach(async () => { + vi.restoreAllMocks(); + for (const child of children.splice(0)) { + if (child.exitCode === null && child.signalCode === null) { + const exited = once(child, "exit"); + child.kill("SIGKILL"); + await exited; + } + } + await Promise.all(directories.splice(0).map((dir) => rm(dir, { recursive: true, force: true }))); + }); + + async function fixture() { + const root = await mkdtemp(path.join(os.tmpdir(), "paperclip-lock-process-")); + directories.push(root); + const target = path.join(root, "target"); + await mkdir(target); + const env = { ...process.env, PAPERCLIP_HOME: path.join(root, "home"), PAPERCLIP_INSTANCE_ID: "test" }; + const key = createHash("sha256").update(await realpath(target)).digest("hex"); + const lock = path.join(env.PAPERCLIP_HOME, "instances", "test", "locks", "directory-merge", `${key}.lock`); + return { target, env, lock }; + } + + async function holder(target: string, env: NodeJS.ProcessEnv) { + const child = spawn(process.execPath, [loader, "--eval", ` + import { withDirectoryMergeLock } from ${JSON.stringify(module)}; + withDirectoryMergeLock(${JSON.stringify(target)}, async () => { + process.send?.("locked"); + await new Promise(resolve => process.once("message", resolve)); + }).then(() => process.exit(0)).catch(error => { console.error(error); process.exit(1); }); + `], { env, stdio: ["ignore", "ignore", "pipe", "ipc"] }); + children.push(child); + let stderr = ""; + child.stderr?.on("data", (chunk) => { stderr += chunk; }); + await new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error(`Child did not acquire lock: ${stderr}`)), 10_000); + child.once("message", () => { clearTimeout(timer); resolve(); }); + child.once("error", (error) => { clearTimeout(timer); reject(error); }); + child.once("exit", (code) => { clearTimeout(timer); reject(new Error(`Child exited ${code}: ${stderr}`)); }); + }); + return child; + } + + // Both protocols carry diagnostic owner records. This also lets the crash + // regression demonstrate the old implementation's PID-reuse failure. + async function ownerPath(lock: string) { + return await access(`${lock}.owner.json`).then(() => `${lock}.owner.json`, () => path.join(lock, "owner.json")); + } + + function expireWait() { + const now = Date.now(); + return vi.spyOn(Date, "now").mockReturnValueOnce(now).mockReturnValue(now + 60_000); + } + + it("recovers a killed holder even when its recorded PID has been reused", async () => { + const { target, env, lock } = await fixture(); + const child = await holder(target, env); + const exited = once(child, "exit"); + child.kill("SIGKILL"); + await exited; + const recordPath = await ownerPath(lock); + const record = JSON.parse(await readFile(recordPath, "utf8")); + await writeFile(recordPath, JSON.stringify({ ...record, pid: process.pid })); + const clock = expireWait(); + try { + await expect(withDirectoryMergeLock(target, async () => "restored", env)).resolves.toBe("restored"); + } finally { clock.mockRestore(); } + }, 15_000); + + it("protects a live holder in another process regardless of diagnostic PID or age", async () => { + const { target, env, lock } = await fixture(); + const child = await holder(target, env); + const recordPath = await ownerPath(lock); + await writeFile(recordPath, JSON.stringify({ pid: 2_147_483_647, createdAt: "2000-01-01T00:00:00.000Z" })); + const contender = vi.fn(); + const clock = expireWait(); + try { + await expect(withDirectoryMergeLock(target, contender, env)).rejects.toMatchObject({ code: WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE }); + } finally { clock.mockRestore(); } + expect(contender).not.toHaveBeenCalled(); + const exited = once(child, "exit"); + child.send("release"); + await exited; + await expect(withDirectoryMergeLock(target, async () => "released", env)).resolves.toBe("released"); + }, 15_000); + + it("releases ownership when the protected operation throws", async () => { + const { target, env } = await fixture(); + await expect(withDirectoryMergeLock(target, async () => { throw new Error("operation failed"); }, env)).rejects.toThrow("operation failed"); + await expect(withDirectoryMergeLock(target, async () => "next", env)).resolves.toBe("next"); + }); + + it("keeps the OS lock when a same-process contender closes its connection", async () => { + const { target, env } = await fixture(); + await withDirectoryMergeLock(target, async () => { + const clock = expireWait(); + try { + await expect(withDirectoryMergeLock(target, async () => undefined, env)).rejects.toMatchObject({ code: WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE }); + } finally { clock.mockRestore(); } + // A same-process test alone cannot prove that the OS lock survived: on + // POSIX, closing an unmanaged descriptor can drop process-wide locks. + const result = await promisify(execFile)(process.execPath, [loader, "--eval", ` + import { withDirectoryMergeLock, WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE } from ${JSON.stringify(module)}; + const now = Date.now(); + let calls = 0; + Date.now = () => now + calls++ * 60_000; + withDirectoryMergeLock(${JSON.stringify(target)}, async () => "entered") + .then(() => { console.error("Entered a live holder's lock"); process.exit(1); }) + .catch(error => { + if (error.code !== WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE) { console.error(error); process.exit(1); } + console.log("blocked"); + }); + `], { env, timeout: 10_000 }); + expect(result.stdout.trim()).toBe("blocked"); + }, env); + }, 15_000); + + it("does not serialize unrelated target directories", async () => { + const { target, env } = await fixture(); + const other = `${target}-other`; + await mkdir(other); + await withDirectoryMergeLock(target, async () => { + await expect(withDirectoryMergeLock(other, async () => "independent", env)).resolves.toBe("independent"); + }, env); + }); + + it("retains the same database inode across holders", async () => { + const { target, env, lock } = await fixture(); + await withDirectoryMergeLock(target, async () => undefined, env); + const before = await stat(`${lock}.sqlite`); + await withDirectoryMergeLock(target, async () => undefined, env); + const after = await stat(`${lock}.sqlite`); + expect({ dev: after.dev, ino: after.ino }).toEqual({ dev: before.dev, ino: before.ino }); + }); + + it("closes the kernel lock when writing diagnostics fails", async () => { + const { target, env, lock } = await fixture(); + await mkdir(`${lock}.owner.json`, { recursive: true }); + await expect(withDirectoryMergeLock(target, async () => undefined, env)).rejects.toThrow(); + await rm(`${lock}.owner.json`, { recursive: true }); + await expect(withDirectoryMergeLock(target, async () => "next", env)).resolves.toBe("next"); + }); + + it.each(["missing", "malformed", "dead PID"])("does not guess ownership of a legacy lock with %s metadata", async (kind) => { + const { target, env, lock } = await fixture(); + await mkdir(lock, { recursive: true }); + if (kind !== "missing") await writeFile(path.join(lock, "owner.json"), kind === "malformed" ? "{broken" : JSON.stringify({ pid: 2_147_483_647, createdAt: "2000-01-01T00:00:00.000Z" })); + const contender = vi.fn(); + const clock = expireWait(); + try { + await expect(withDirectoryMergeLock(target, contender, env)).rejects.toMatchObject({ code: WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE }); + } finally { clock.mockRestore(); } + expect(contender).not.toHaveBeenCalled(); + expect((await stat(lock)).isDirectory()).toBe(true); + }); + + it.skipIf(process.platform === "win32").each(["symlink", "hard link"])("rejects a lock database that is a %s", async (kind) => { + const { target, env, lock } = await fixture(); + await mkdir(path.dirname(lock), { recursive: true }); + const decoy = `${target}-decoy`; + await writeFile(decoy, "keep"); + await (kind === "symlink" ? symlink(decoy, `${lock}.sqlite`) : link(decoy, `${lock}.sqlite`)); + await expect(withDirectoryMergeLock(target, async () => undefined, env)).rejects.toThrow("not a plain, unshared file"); + await expect(readFile(decoy, "utf8")).resolves.toBe("keep"); + }); +}); diff --git a/packages/adapter-utils/src/workspace-restore-merge.test.ts b/packages/adapter-utils/src/workspace-restore-merge.test.ts index aaacd5d7da..bfa57996a9 100644 --- a/packages/adapter-utils/src/workspace-restore-merge.test.ts +++ b/packages/adapter-utils/src/workspace-restore-merge.test.ts @@ -376,7 +376,7 @@ describe("workspace restore merge", () => { } }); - it("creates the lock root at mode 0o700 and removes the lock directory after release", async () => { + it("keeps the private lock database and removes diagnostic ownership after release", async () => { const rootDir = await mkdtemp(path.join(os.tmpdir(), "paperclip-restore-merge-")); cleanupDirs.push(rootDir); const paperclipHome = path.join(rootDir, "paperclip-home"); @@ -392,8 +392,11 @@ describe("workspace restore merge", () => { }); expect((await stat(lockRootDir)).mode & 0o777).toBe(0o700); - expect(entriesDuringLock).toHaveLength(1); - await expect(readdir(lockRootDir)).resolves.toHaveLength(0); + expect(entriesDuringLock.filter((name) => name.endsWith(".owner.json"))).toHaveLength(1); + const entriesAfterRelease = await readdir(lockRootDir); + expect(entriesAfterRelease).toHaveLength(1); + expect(entriesAfterRelease[0]).toMatch(/\.lock\.sqlite$/); + expect((await stat(path.join(lockRootDir, entriesAfterRelease[0]!))).mode & 0o777).toBe(0o600); }); it("classifies the real lock-timeout error by its stable code, never by the message text", async () => { @@ -469,9 +472,9 @@ describe("workspace restore merge", () => { }); } finally { clock.mockRestore(); } expect(contender).not.toHaveBeenCalled(); - expect(await readdir(lockRoot)).toHaveLength(1); + expect((await readdir(lockRoot)).filter((name) => name.endsWith(".owner.json"))).toHaveLength(1); }); - expect(await readdir(lockRoot)).toHaveLength(0); + expect(await readdir(lockRoot)).toEqual([expect.stringMatching(/\.lock\.sqlite$/)]); }); it("delivers the timeout when the diagnostic owner read stalls and ignores its late rejection", async () => { diff --git a/packages/adapter-utils/src/workspace-restore-merge.ts b/packages/adapter-utils/src/workspace-restore-merge.ts index 755bce9be4..60bdd8cc90 100644 --- a/packages/adapter-utils/src/workspace-restore-merge.ts +++ b/packages/adapter-utils/src/workspace-restore-merge.ts @@ -2,6 +2,7 @@ import { createHash, randomUUID } from "node:crypto"; import { createReadStream } from "node:fs"; import { constants as fsConstants, promises as fs } from "node:fs"; import path from "node:path"; +import { DatabaseSync } from "node:sqlite"; import { createWorkspaceManifest, WorkspaceManifestMap, workspacePathMatcher, type PathManifest, type WorkspacePaths, type WorkspaceManifestWriter } from "./workspace-manifest.js"; import { shouldExcludePath } from "./exclude-patterns.js"; import { resolvePaperclipInstanceRootForAdapter } from "./server-utils.js"; @@ -200,7 +201,7 @@ function entriesMatch(left: SnapshotEntry | null | undefined, right: SnapshotEnt return false; } -const LOCK_STALE_MS = 30_000; +const LOCK_WAIT_MS = 30_000; const LOCK_DIAGNOSTIC_READ_TIMEOUT_MS = 100; const activeDirectoryMergeLocks = new Set(); const MAX_LOCK_DIAGNOSTIC_AGE_MS = 7 * 24 * 60 * 60 * 1000; @@ -212,7 +213,7 @@ export type DirectoryMergeLockOperation = /** Evidence only: neither process age nor this module's holder set can prove * that a lock in another process or PID namespace is safe to reclaim. */ -async function directoryMergeLockDiagnostics(lockDir: string, waitMs: number): Promise> { +async function directoryMergeLockDiagnostics(lockDir: string, waitMs: number, ownerPath: string): Promise> { const diagnostics: Record = { ownerState: "unknown", knownLocalHolder: activeDirectoryMergeLocks.has(lockDir), @@ -224,7 +225,7 @@ async function directoryMergeLockDiagnostics(lockDir: string, waitMs: number): P // Abort is best effort: race the read as well so a stalled filesystem // cannot keep the original lock timeout from reaching its caller. const raw = await Promise.race([ - fs.readFile(path.join(lockDir, "owner.json"), { encoding: "utf8", signal: controller.signal }), + fs.readFile(ownerPath, { encoding: "utf8", signal: controller.signal }), new Promise((resolve) => { timeout = setTimeout(() => resolve(undefined), LOCK_DIAGNOSTIC_READ_TIMEOUT_MS); }), @@ -339,75 +340,86 @@ export function describeWorkspaceRestoreFailure(code: WorkspaceRestoreFailureCod } } -async function isLockStale(lockDir: string): Promise { - try { - const raw = await fs.readFile(path.join(lockDir, "owner.json"), "utf8"); - const owner = JSON.parse(raw) as { pid?: unknown }; - const pid = typeof owner.pid === "number" && Number.isFinite(owner.pid) && owner.pid > 0 ? owner.pid : null; - if (pid === null) { - // Owner record is unparseable / missing pid — treat as stale. - return true; - } - try { - process.kill(pid, 0); - return false; - } catch { - return true; - } - } catch { - // owner.json is missing or unreadable. A live holder also passes through - // this exact state, briefly, between its own `fs.mkdir(lockDir)` and its - // `fs.writeFile(owner.json)` below. Reading "missing" as "stale" here would - // let a concurrent acquirer delete a live holder's lock directory during - // that window. Mirror the materializePaperclipSkillCopy lock pattern: fall - // back to the lock directory's own mtime, and only call it stale once the - // directory itself has outlived the stale threshold. - const stat = await fs.stat(lockDir).catch(() => null); - return !stat || Date.now() - stat.mtimeMs > LOCK_STALE_MS; - } -} - async function acquireDirectoryMergeLock(lockDir: string, operation?: DirectoryMergeLockOperation): Promise<() => Promise> { const startedAt = performance.now(); - const deadline = Date.now() + LOCK_STALE_MS; - while (true) { - try { - await fs.mkdir(lockDir); - await fs.writeFile( - path.join(lockDir, "owner.json"), - `${JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() })}\n`, - "utf8", + const deadline = Date.now() + LOCK_WAIT_MS; + const databasePath = `${lockDir}.sqlite`; + const ownerPath = `${lockDir}.owner.json`; + async function waitForLock(diagnosticOwnerPath: string) { + if (Date.now() >= deadline) { + const timeoutError: NodeJS.ErrnoException & { workspaceRestoreLock?: Record } = new Error( + `Timed out waiting for workspace restore lock at ${lockDir}`, ); - activeDirectoryMergeLocks.add(lockDir); - return async () => { - try { - await fs.rm(lockDir, { recursive: true, force: true }).catch(() => undefined); - } finally { - activeDirectoryMergeLocks.delete(lockDir); - } - }; - } catch (error) { - const code = error && typeof error === "object" ? (error as { code?: unknown }).code : null; - if (code !== "EEXIST") throw error; - // Stale-lock detection: if the owner PID is dead (SIGKILL / OOM / crash), - // the lockDir would otherwise persist forever and stall restores. Mirror - // the materializePaperclipSkillCopy lock pattern — remove and retry. - if (await isLockStale(lockDir)) { - await fs.rm(lockDir, { recursive: true, force: true }).catch(() => undefined); - continue; - } - if (Date.now() >= deadline) { - const timeoutError: NodeJS.ErrnoException & { workspaceRestoreLock?: Record } = new Error( - `Timed out waiting for workspace restore lock at ${lockDir}`, - ); - timeoutError.code = WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE; - // Keep the original timeout if the diagnostic read itself fails. - timeoutError.workspaceRestoreLock = await directoryMergeLockDiagnostics(lockDir, performance.now() - startedAt).catch(() => undefined); - if (operation && timeoutError.workspaceRestoreLock) timeoutError.workspaceRestoreLock.operation = operation; - throw timeoutError; - } - await new Promise((resolve) => setTimeout(resolve, 50)); + timeoutError.code = WORKSPACE_RESTORE_LOCK_TIMEOUT_CODE; + // Keep the original timeout if the diagnostic read itself fails. + timeoutError.workspaceRestoreLock = await directoryMergeLockDiagnostics(lockDir, performance.now() - startedAt, diagnosticOwnerPath).catch(() => undefined); + if (operation && timeoutError.workspaceRestoreLock) timeoutError.workspaceRestoreLock.operation = operation; + throw timeoutError; } + await new Promise((resolve) => setTimeout(resolve, 50)); + } + + // SQLite's RESERVED file lock is the authority, including across processes + // and PID namespaces. It is released by the OS on a crash. The empty database + // is permanent: unlinking it would let contenders lock different inodes. + // node:sqlite is already required for workspace manifests; no native add-on + // or external flock command is needed on macOS, Linux, or Windows. + const databaseStat = await fs.lstat(databasePath).catch((error: NodeJS.ErrnoException) => { + if (error.code === "ENOENT") return null; + throw error; + }); + if (databaseStat && (!databaseStat.isFile() || databaseStat.isSymbolicLink() || databaseStat.nlink !== 1)) { + throw new Error("Directory merge lock database is not a plain, unshared file."); + } + // Let SQLite create and manage every descriptor for this inode. On POSIX, + // closing a raw fs.open descriptor could release another connection's locks. + // The parent is private (0700), including while a new file is chmodded. + const database = new DatabaseSync(databasePath, { allowExtension: false }); + try { + await fs.chmod(databasePath, 0o600); + // Never block the event loop while another async operation holds the lock. + database.exec("PRAGMA busy_timeout=0;"); + while (true) { + try { + database.exec("BEGIN IMMEDIATE;"); + break; + } catch (error) { + const code = (error as { errcode?: number }).errcode; + if (typeof code !== "number" || (code & 0xff) !== 5) throw error; // SQLITE_BUSY + await waitForLock(ownerPath); + } + } + + // Old processes do not participate in the SQLite protocol. Never infer + // that a legacy owner is dead from PID existence, age, or missing metadata. + // Drain old writers before upgrading. A leftover legacy directory requires + // explicit offline cleanup; a live legacy holder can still release normally. + while (await fs.lstat(lockDir).then(() => true, (error: NodeJS.ErrnoException) => { + if (error.code === "ENOENT") return false; + throw error; + })) await waitForLock(path.join(lockDir, "owner.json")); + + const owner = await fs.open(ownerPath, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_TRUNC | fsConstants.O_NOFOLLOW, 0o600); + try { + await owner.writeFile(`${JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() })}\n`, "utf8"); + } finally { await owner.close(); } + activeDirectoryMergeLocks.add(lockDir); + let released = false; + return async () => { + if (released) return; + released = true; + try { + // This sidecar is diagnostic only. A failed removal cannot retain + // ownership, and the next holder replaces it while holding the DB lock. + await fs.unlink(ownerPath).catch(() => undefined); + } finally { + activeDirectoryMergeLocks.delete(lockDir); + database.close(); + } + }; + } catch (error) { + database.close(); + throw error; } }