From 4fc822c9e4899309529d3fd5016de0299fb24119 Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Fri, 9 Oct 2026 23:49:00 -0400 Subject: [PATCH] fix: preserve filesystem lock ownership during recovery --- CHANGELOG.md | 5 + docs/0.5.0-api.md | 8 +- src/directory-swap.ts | 54 +- src/file-changes.ts | 72 +-- src/filesystem-lock.ts | 147 ++++++ tests/cli-smoke.test.ts | 13 +- tests/filesystem-locks.test.ts | 490 ++++++++++++++++++ .../fixtures/filesystem-lock-crash-worker.ts | 19 + tests/public-contract.json | 1 + 9 files changed, 684 insertions(+), 125 deletions(-) create mode 100644 src/filesystem-lock.ts create mode 100644 tests/filesystem-locks.test.ts create mode 100644 tests/fixtures/filesystem-lock-crash-worker.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 7017618..654c10f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Prepare file and directory locks with unique owner records before publishing + them. Recover only the observed dead owner so delayed stale-lock recovery + cannot delete a successor's lock and allow overlapping updates. Preserve + ambiguous lock contents and junctions; report failed private-stage cleanup. + - Publish generated OpenAPI clients through the same recoverable directory swap as other commands. Keep complete published output when backup cleanup fails, clean partial stages, and report retained stages after cleanup failure. diff --git a/docs/0.5.0-api.md b/docs/0.5.0-api.md index b36cd45..3ff9f8c 100644 --- a/docs/0.5.0-api.md +++ b/docs/0.5.0-api.md @@ -48,6 +48,12 @@ The OpenAPI generator now uses that same publication routine and preserves its d Create reports success after publication completes. Child-process tests inject package-manager nonzero exit, missing-command and timeout results, then assert an error exit with no published destination or abandoned project stage. A successful install followed by failed publication also reports no success. These are subprocess result injections, not qualification of a real installer timeout or cancellation bound. +## Filesystem lock failure qualification + +File transactions and directory publication now share one lock protocol: prepare a unique owner record in a sibling directory, publish the nonempty lock with a rename, then unlink only the observed dead owner or this operation's own owner and remove the directory only if empty. A delayed stale reaper cannot unlink a successor's different owner record. This replaces the two duplicated recursive-lock-removal implementations; the ten-second wait bound is unchanged. Existing dead `owner.json` records and aged empty/malformed legacy locks can be recovered. Ambiguous contents, symbolic links and owner-record junctions stop recovery and preserve unrelated paths. + +The executed matrix covers two simultaneous stale-owner observations for file and directory operations, exact stale-write rejection, unrelated lock contents, junctions, fresh malformed-record timeout, aged legacy recovery, partial owner writes, failed preparation cleanup and an exclusive preparation collision. It also verifies successor ownership during release, denied process probes and filesystem errors during acquisition. The original six minimal cases failed before the fix. Separate child-process cases terminate after preparing an owner but before publishing the lock: data is unchanged, no incomplete live lock is exposed, retry succeeds and the abandoned private preparation remains available for manual cleanup. Existing file/directory termination and recovery cases also run with the shared protocol. These results qualify the tested process checkpoints and cooperating current CLI commands, not power-loss durability or concurrent use of older lock protocols. + ## Release boundaries -The normal packed consumer checks both retained declaration names, nine rejected old imports and private paths under strict TypeScript 6 and 7, then runs the installed CLI help/version commands. No production dependency, package version or peer range changes in this cut. Full candidate graph, generated website API and maintainer review still gate release. CLI hardening issue #159 stays open until its remaining filesystem-failure qualification is complete. +The normal packed consumer checks both retained declaration names, nine rejected old imports and private paths under strict TypeScript 6 and 7, then runs the installed CLI help/version commands. No production dependency, package version or peer range changes in this cut. Full candidate graph, generated website API and maintainer review still gate release. CLI hardening issue #159 stays open until its remaining qualification is complete. diff --git a/src/directory-swap.ts b/src/directory-swap.ts index eb175ca..2fe28e6 100644 --- a/src/directory-swap.ts +++ b/src/directory-swap.ts @@ -2,6 +2,7 @@ import { randomUUID } from "node:crypto"; import { constants, type BigIntStats } from "node:fs"; import fs from "node:fs/promises"; import path from "node:path"; +import { acquireFilesystemLock } from "./filesystem-lock"; async function directoryStat(target: string): Promise { let stat: BigIntStats; @@ -134,70 +135,21 @@ async function recoverPublication(target: string): Promise { } } -const LOCK_RETRY_MS = 10; -const LOCK_TIMEOUT_MS = 10_000; -const ORPHANED_LOCK_AGE_MS = 30_000; - function isNodeError(error: unknown, code: string): error is NodeJS.ErrnoException { return error instanceof Error && "code" in error && error.code === code; } -async function removeOrphanedLock(lock: string): Promise { - try { - const owner = JSON.parse(await fs.readFile(path.join(lock, "owner.json"), "utf8")) as { - pid?: unknown; - }; - if (Number.isInteger(owner.pid) && (owner.pid as number) > 0) { - try { - process.kill(owner.pid as number, 0); - return false; - } catch (error) { - if (!isNodeError(error, "ESRCH") && !isNodeError(error, "EINVAL")) return false; - } - } else { - return false; - } - } catch { - const stat = await fs.stat(lock).catch(() => null); - if (!stat || Date.now() - stat.mtimeMs < ORPHANED_LOCK_AGE_MS) return false; - } - await fs.rm(lock, { recursive: true, force: true }); - return true; -} - export async function withDirectoryTargetLock( target: string, operation: () => Promise, ): Promise { const lock = `${path.resolve(target)}.askr-lock`; - const deadline = Date.now() + LOCK_TIMEOUT_MS; - while (true) { - try { - await fs.mkdir(lock); - await fs.writeFile( - path.join(lock, "owner.json"), - `${JSON.stringify({ pid: process.pid })}\n`, - { - flag: "wx", - }, - ); - break; - } catch (error) { - if (!isNodeError(error, "EEXIST")) { - await fs.rm(lock, { recursive: true, force: true }).catch(() => undefined); - throw error; - } - if (await removeOrphanedLock(lock)) continue; - if (Date.now() >= deadline) - throw new Error(`Timed out waiting for directory lock: ${target}`); - await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_MS)); - } - } + const release = await acquireFilesystemLock(lock, `directory lock for ${JSON.stringify(target)}`); try { await recoverPublication(path.resolve(target)); return await operation(); } finally { - await fs.rm(lock, { recursive: true, force: true }); + await release(); } } diff --git a/src/file-changes.ts b/src/file-changes.ts index 5256b99..2e1c3f9 100644 --- a/src/file-changes.ts +++ b/src/file-changes.ts @@ -1,6 +1,7 @@ import fs from "node:fs/promises"; import path from "node:path"; import { randomUUID } from "node:crypto"; +import { acquireFilesystemLock } from "./filesystem-lock"; export interface FileChange { readonly filePath: string; @@ -19,80 +20,25 @@ export interface FileChangeWriterOptions { } interface FileLock { - readonly lockPath: string; + readonly release: () => Promise; } -const LOCK_RETRY_MS = 10; -const LOCK_TIMEOUT_MS = 10_000; -const ORPHANED_LOCK_AGE_MS = 30_000; - function isNodeError(error: unknown, code: string): error is NodeJS.ErrnoException { return error instanceof Error && "code" in error && error.code === code; } -function delay(milliseconds: number): Promise { - return new Promise((resolve) => setTimeout(resolve, milliseconds)); -} - -async function ownerIsAlive(lockPath: string): Promise { - try { - const owner = JSON.parse(await fs.readFile(path.join(lockPath, "owner.json"), "utf8")) as { - pid?: unknown; - }; - if (!Number.isInteger(owner.pid) || (owner.pid as number) <= 0) return undefined; - try { - process.kill(owner.pid as number, 0); - return true; - } catch (error) { - if (isNodeError(error, "ESRCH")) return false; - return true; - } - } catch { - return undefined; - } -} - -async function removeOrphanedLock(lockPath: string): Promise { - const ownerAlive = await ownerIsAlive(lockPath); - if (ownerAlive === true) return false; - if (ownerAlive === undefined) { - const stat = await fs.stat(lockPath).catch(() => null); - if (!stat || Date.now() - stat.mtimeMs < ORPHANED_LOCK_AGE_MS) return false; - } - await fs.rm(lockPath, { recursive: true, force: true }); - return true; -} - async function acquireFileLock(filePath: string): Promise { const lockPath = path.join(path.dirname(filePath), `.${path.basename(filePath)}.askr-lock`); - const deadline = Date.now() + LOCK_TIMEOUT_MS; - while (true) { - try { - await fs.mkdir(lockPath); - await fs.writeFile( - path.join(lockPath, "owner.json"), - `${JSON.stringify({ pid: process.pid })}\n`, - { flag: "wx" }, - ); - return { lockPath }; - } catch (error) { - if (!isNodeError(error, "EEXIST")) { - await fs.rm(lockPath, { recursive: true, force: true }).catch(() => undefined); - throw error; - } - if (await removeOrphanedLock(lockPath)) continue; - if (Date.now() >= deadline) { - throw new Error(`Timed out waiting for file transaction lock: ${filePath}`); - } - await delay(LOCK_RETRY_MS); - } - } + return { + release: await acquireFilesystemLock( + lockPath, + `file transaction lock for ${JSON.stringify(filePath)}`, + ), + }; } async function releaseFileLocks(locks: readonly FileLock[]): Promise { - await Promise.all( - [...locks].reverse().map((lock) => fs.rm(lock.lockPath, { recursive: true, force: true })), - ); + await Promise.all([...locks].reverse().map((lock) => lock.release())); } async function readCurrentContent(filePath: string): Promise { diff --git a/src/filesystem-lock.ts b/src/filesystem-lock.ts new file mode 100644 index 0000000..5b1f730 --- /dev/null +++ b/src/filesystem-lock.ts @@ -0,0 +1,147 @@ +import fs from "node:fs/promises"; +import path from "node:path"; +import { randomUUID } from "node:crypto"; + +const LOCK_RETRY_MS = 10; +const LOCK_TIMEOUT_MS = 10_000; +const ORPHANED_LOCK_AGE_MS = 30_000; +const OWNER_FILE = /^owner-[\da-f]{8}-[\da-f]{4}-[\da-f]{4}-[\da-f]{4}-[\da-f]{12}\.json$/i; +const RENAME_CONTENTION_CODES = new Set(["EEXIST", "ENOTEMPTY", "EPERM", "EBUSY", "EACCES"]); + +function isNodeError(error: unknown, code: string): error is NodeJS.ErrnoException { + return error instanceof Error && "code" in error && error.code === code; +} + +async function statIfPresent(target: string) { + try { + return await fs.lstat(target); + } catch (error) { + if (isNodeError(error, "ENOENT")) return null; + throw error; + } +} + +function ambiguousLock(lock: string): Error { + return new Error( + `Cannot recover filesystem lock ${JSON.stringify(lock)}: its contents do not identify one regular owner record. Inspect the lock and preserve unrelated files before retrying.`, + ); +} + +async function removeEmptyLock(lock: string): Promise { + try { + await fs.rmdir(lock); + return true; + } catch (error) { + if (isNodeError(error, "ENOENT")) return true; + if (isNodeError(error, "ENOTEMPTY") || isNodeError(error, "EEXIST")) return false; + throw error; + } +} + +async function removeOrphanedLock(lock: string): Promise { + const stat = await statIfPresent(lock); + if (!stat) return false; + if (!stat.isDirectory() || stat.isSymbolicLink()) throw ambiguousLock(lock); + let entries: string[]; + try { + entries = await fs.readdir(lock); + } catch (error) { + if (isNodeError(error, "ENOENT")) return false; + throw error; + } + if (entries.length === 0) { + if (Date.now() - stat.mtimeMs < ORPHANED_LOCK_AGE_MS) return false; + return removeEmptyLock(lock); + } + if (entries.length !== 1 || (entries[0] !== "owner.json" && !OWNER_FILE.test(entries[0]))) { + throw ambiguousLock(lock); + } + const ownerPath = path.join(lock, entries[0]); + const ownerStat = await statIfPresent(ownerPath); + if (!ownerStat) return false; + if (!ownerStat.isFile() || ownerStat.isSymbolicLink()) throw ambiguousLock(lock); + let pid: unknown; + try { + const owner = JSON.parse(await fs.readFile(ownerPath, "utf8")) as { pid?: unknown } | null; + pid = owner?.pid; + } catch (error) { + if (isNodeError(error, "ENOENT")) return false; + if (!(error instanceof SyntaxError)) throw error; + } + if (Number.isInteger(pid) && (pid as number) > 0 && (pid as number) <= 2_147_483_647) { + try { + process.kill(pid as number, 0); + return false; + } catch (error) { + if (!isNodeError(error, "ESRCH") && !isNodeError(error, "EINVAL")) return false; + } + } else if (Date.now() - stat.mtimeMs < ORPHANED_LOCK_AGE_MS) { + return false; + } + // Only one reaper can unlink this observed token. A new acquisition publishes + // a different token atomically, so a delayed reaper cannot remove its owner. + try { + await fs.unlink(ownerPath); + } catch (error) { + if (isNodeError(error, "ENOENT")) return false; + throw error; + } + return removeEmptyLock(lock); +} + +/** Acquires a cooperative lock and returns a release operation for its unique owner. */ +export async function acquireFilesystemLock( + lock: string, + description: string, +): Promise<() => Promise> { + const ownerName = `owner-${randomUUID()}.json`; + const stage = `${lock}.stage-${randomUUID()}`; + // Build a nonempty lock privately before publishing it. Never expose a newly + // acquired empty directory that another reaper could mistake for an orphan. + await fs.mkdir(stage); + try { + await fs.writeFile(path.join(stage, ownerName), `${JSON.stringify({ pid: process.pid })}\n`, { + flag: "wx", + }); + const deadline = Date.now() + LOCK_TIMEOUT_MS; + let lastError: unknown; + while (true) { + if (Date.now() >= deadline) { + throw new Error( + `Timed out waiting for ${description} at ${JSON.stringify(lock)}. Inspect its owner before retrying.`, + { cause: lastError }, + ); + } + const lockStat = await statIfPresent(lock); + if (lockStat && (!lockStat.isDirectory() || lockStat.isSymbolicLink())) { + throw ambiguousLock(lock); + } + try { + await fs.rename(stage, lock); + return async () => { + await fs.unlink(path.join(lock, ownerName)); + // Another contender can already have replaced the empty old lock with + // its nonempty lock. Never recursively remove a successor's contents. + await removeEmptyLock(lock); + }; + } catch (error) { + const code = error instanceof Error && "code" in error ? String(error.code) : ""; + if (!RENAME_CONTENTION_CODES.has(code)) throw error; + lastError = error; + if (await removeOrphanedLock(lock)) continue; + await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_MS)); + } + } + } catch (error) { + try { + await fs.rm(stage, { recursive: true, force: true }); + } catch (cleanupFailure) { + throw new AggregateError( + [error, cleanupFailure], + `Filesystem lock preparation failed; retained private lock stage ${JSON.stringify(stage)}. Remove this stage after resolving the filesystem error and retry.`, + { cause: error }, + ); + } + throw error; + } +} diff --git a/tests/cli-smoke.test.ts b/tests/cli-smoke.test.ts index aac0881..63baed3 100644 --- a/tests/cli-smoke.test.ts +++ b/tests/cli-smoke.test.ts @@ -2714,7 +2714,6 @@ test("should ensure skill sync never swaps the live tree while another sync is c }); const originalCp = fs.cp.bind(fs); const originalRename = fs.rename.bind(fs); - const originalMkdir = fs.mkdir.bind(fs); const cp = vi.spyOn(fs, "cp").mockImplementation((async (...args: Parameters) => { if (path.resolve(String(args[0])) !== skillsRoot) return originalCp(...args); copiesInFlight += 1; @@ -2738,19 +2737,14 @@ test("should ensure skill sync never swaps the live tree while another sync is c code: "EPERM", }); } - return originalRename(...args); - }) as typeof fs.rename); - const mkdir = vi.spyOn(fs, "mkdir").mockImplementation((async ( - ...args: Parameters - ) => { try { - return await originalMkdir(...args); + return await originalRename(...args); } catch (error) { // A second sync waiting on the lock lets the first one finish its copy. - if (path.resolve(String(args[0])) === lockPath) releaseGate(); + if (path.resolve(String(args[1])) === lockPath) releaseGate(); throw error; } - }) as typeof fs.mkdir); + }) as typeof fs.rename); try { const results = await Promise.all([ @@ -2766,7 +2760,6 @@ test("should ensure skill sync never swaps the live tree while another sync is c } finally { cp.mockRestore(); rename.mockRestore(); - mkdir.mockRestore(); await fs.rm(tempRoot, { recursive: true, force: true }); } }); diff --git a/tests/filesystem-locks.test.ts b/tests/filesystem-locks.test.ts new file mode 100644 index 0000000..3d7e6e7 --- /dev/null +++ b/tests/filesystem-locks.test.ts @@ -0,0 +1,490 @@ +import fs from "node:fs/promises"; +import path from "node:path"; +import os from "node:os"; +import { AsyncLocalStorage } from "node:async_hooks"; +import { fork } from "node:child_process"; +import { once } from "node:events"; +import { fileURLToPath } from "node:url"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { withDirectoryTargetLock } from "../src/directory-swap"; +import { writeFileChanges } from "../src/file-changes"; + +const roots: string[] = []; +const context = new AsyncLocalStorage(); +function deferred() { + let resolve!: () => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} +async function fixture(kind: string) { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "askr-lock-ownership-")); + roots.push(root); + const target = path.join(root, kind === "directory" ? "output" : "manifest.json"); + if (kind === "file") await fs.writeFile(target, "old"); + const lock = + kind === "directory" ? `${target}.askr-lock` : path.join(root, ".manifest.json.askr-lock"); + return { root, target, lock }; +} +async function operate(kind: string, target: string, operation: () => Promise) { + if (kind === "directory") return withDirectoryTargetLock(target, operation); + return writeFileChanges( + [{ filePath: target, content: context.getStore() ?? "new", expectedContent: "old" }], + { + replace: async (temporary, file) => { + await operation(); + await fs.rename(temporary, file); + }, + }, + ); +} +afterEach(async () => { + vi.restoreAllMocks(); + await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); +}); + +describe("filesystem lock ownership", () => { + it.each(["directory", "file"])( + "does not remove a successor's %s owner during release", + async (kind) => { + const { target, lock } = await fixture(kind); + const rmdir = fs.rmdir.bind(fs); + const successor = path.join(lock, "owner-aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa.json"); + let handedOff = false; + vi.spyOn(fs, "rmdir").mockImplementation(async (...args) => { + if (String(args[0]) === lock && !handedOff) { + handedOff = true; + await rmdir(...args); + await fs.mkdir(lock); + await fs.writeFile(successor, JSON.stringify({ pid: process.pid })); + // The successor published after our owner was unlinked but before the + // releasing process's empty-directory removal reached the filesystem. + } + return rmdir(...args); + }); + await expect(operate(kind, target, async () => {})).resolves.toBeUndefined(); + expect(handedOff).toBe(true); + expect(JSON.parse(await fs.readFile(successor, "utf8"))).toEqual({ pid: process.pid }); + if (kind === "file") expect(await fs.readFile(target, "utf8")).toBe("new"); + }, + ); + + it("preserves an owner whose process cannot be probed instead of treating permission denial as death", async () => { + const { root, target, lock } = await fixture("file"); + await fs.mkdir(lock); + const ownerPath = path.join(lock, "owner.json"); + await fs.writeFile(ownerPath, '{"pid":2147483647}'); + let now = Date.now(); + vi.spyOn(Date, "now").mockImplementation(() => now); + vi.spyOn(process, "kill").mockImplementation(() => { + now += 10_001; + throw Object.assign(new Error("injected process probe denial"), { code: "EPERM" }); + }); + await expect(operate("file", target, async () => {})).rejects.toThrow("Timed out waiting"); + expect(await fs.readFile(ownerPath, "utf8")).toBe('{"pid":2147483647}'); + expect(await fs.readFile(target, "utf8")).toBe("old"); + expect((await fs.readdir(root)).sort()).toEqual([path.basename(lock), "manifest.json"].sort()); + }); + + it.each(["stat", "readdir", "readFile", "unlink", "rmdir", "rename"] as const)( + "preserves a filesystem %s failure while acquiring a lock", + async (method) => { + const { root, target, lock } = await fixture("file"); + await fs.mkdir(lock); + const ownerPath = path.join(lock, "owner.json"); + await fs.writeFile(ownerPath, '{"pid":2147483647}'); + const failure = Object.assign(new Error(`injected ${method} failure`), { + code: method === "rename" ? "EIO" : "EACCES", + }); + if (method === "stat") { + const lstat = fs.lstat.bind(fs); + vi.spyOn(fs, "lstat").mockImplementation((async (...args: Parameters) => { + if (String(args[0]) === lock) throw failure; + return lstat(...args); + }) as typeof fs.lstat); + } else if (method === "readdir") { + const readdir = fs.readdir.bind(fs); + vi.spyOn(fs, "readdir").mockImplementation((async ( + ...args: Parameters + ) => { + if (String(args[0]) === lock) throw failure; + return readdir(...args); + }) as typeof fs.readdir); + } else if (method === "readFile") { + const readFile = fs.readFile.bind(fs); + vi.spyOn(fs, "readFile").mockImplementation((async ( + ...args: Parameters + ) => { + if (String(args[0]) === ownerPath) throw failure; + return readFile(...args); + }) as typeof fs.readFile); + } else if (method === "unlink") { + const unlink = fs.unlink.bind(fs); + vi.spyOn(fs, "unlink").mockImplementation(async (...args) => { + if (String(args[0]) === ownerPath) throw failure; + return unlink(...args); + }); + } else if (method === "rmdir") { + const rmdir = fs.rmdir.bind(fs); + vi.spyOn(fs, "rmdir").mockImplementation(async (...args) => { + if (String(args[0]) === lock) throw failure; + return rmdir(...args); + }); + } else { + const rename = fs.rename.bind(fs); + vi.spyOn(fs, "rename").mockImplementation(async (from, to) => { + if (String(to) === lock) throw failure; + return rename(from, to); + }); + } + await expect(operate("file", target, async () => {})).rejects.toBe(failure); + vi.restoreAllMocks(); + expect(await fs.readFile(target, "utf8")).toBe("old"); + expect( + (await fs.readdir(root)).some((entry) => entry.startsWith(`${path.basename(lock)}.stage-`)), + ).toBe(false); + if (method === "rmdir") expect(await fs.readdir(lock)).toEqual([]); + else expect(await fs.readFile(ownerPath, "utf8")).toBe('{"pid":2147483647}'); + }, + ); + + it.each(["directory", "file"])( + "can retry after a process is killed before publishing its prepared %s lock", + async (kind) => { + const { target, lock } = await fixture(kind); + const child = fork( + fileURLToPath(new URL("./fixtures/filesystem-lock-crash-worker.ts", import.meta.url)), + [kind, target], + { silent: true, execArgv: ["--import", "tsx"] }, + ); + const exited = once(child, "exit"); + let stderr = ""; + child.stderr?.on("data", (chunk: Buffer) => { + stderr += chunk.toString(); + }); + let stage = ""; + try { + const [checkpoint] = await Promise.race([ + once(child, "message", { signal: AbortSignal.timeout(5_000) }), + exited.then(() => { + throw new Error(`Lock worker exited before preparation: ${stderr}`); + }), + ]); + stage = checkpoint.stage; + expect(child.kill("SIGKILL")).toBe(true); + await exited; + } finally { + if (child.exitCode === null && child.signalCode === null) child.kill("SIGKILL"); + await exited; + } + await expect(fs.access(lock)).rejects.toMatchObject({ code: "ENOENT" }); + if (kind === "file") expect(await fs.readFile(target, "utf8")).toBe("old"); + const entries = await fs.readdir(stage); + expect(entries).toHaveLength(1); + expect(JSON.parse(await fs.readFile(path.join(stage, entries[0]), "utf8"))).toEqual({ + pid: child.pid, + }); + let entered = false; + await operate(kind, target, async () => { + entered = true; + }); + expect(entered).toBe(true); + if (kind === "file") expect(await fs.readFile(target, "utf8")).toBe("new"); + // This abandoned preparation belongs to the killed process, not the retry. + expect(await fs.readdir(stage)).toEqual(entries); + await fs.rm(stage, { recursive: true }); + }, + 10_000, + ); + + it.each(["directory", "file"])( + "recovers an old empty %s lock without leaving private preparation directories", + async (kind) => { + const { root, target, lock } = await fixture(kind); + await fs.mkdir(lock); + const old = new Date(Date.now() - 60_000); + await fs.utimes(lock, old, old); + let entered = false; + await operate(kind, target, async () => { + entered = true; + }); + expect(entered).toBe(true); + expect(await fs.readdir(root)).toEqual(kind === "directory" ? [] : ["manifest.json"]); + }, + ); + + it.each(["directory", "file"])( + "recovers an aged malformed legacy owner for a %s lock", + async (kind) => { + const { root, target, lock } = await fixture(kind); + await fs.mkdir(lock); + await fs.writeFile(path.join(lock, "owner.json"), "{"); + const old = new Date(Date.now() - 60_000); + await fs.utimes(lock, old, old); + await expect(operate(kind, target, async () => {})).resolves.toBeUndefined(); + expect(await fs.readdir(root)).toEqual(kind === "directory" ? [] : ["manifest.json"]); + }, + ); + + it.each(["directory", "file"])( + "times out without deleting a fresh malformed %s owner record", + async (kind) => { + const { root, target, lock } = await fixture(kind); + await fs.mkdir(lock); + const ownerPath = path.join(lock, "owner.json"); + await fs.writeFile(ownerPath, "{"); + let now = Date.now(); + vi.spyOn(Date, "now").mockImplementation(() => now); + const readFile = fs.readFile.bind(fs); + vi.spyOn(fs, "readFile").mockImplementation((async ( + ...args: Parameters + ) => { + const contents = await readFile(...args); + if (String(args[0]) === ownerPath) now += 10_001; + return contents; + }) as typeof fs.readFile); + await expect(operate(kind, target, async () => {})).rejects.toThrow("Timed out waiting"); + expect(await fs.readFile(ownerPath, "utf8")).toBe("{"); + expect((await fs.readdir(root)).sort()).toEqual( + kind === "directory" + ? [path.basename(lock)] + : [path.basename(lock), "manifest.json"].sort(), + ); + }, + ); + + it.each(["directory", "file"])( + "cleans its partial owner record if preparing a %s lock fails", + async (kind) => { + const { root, target, lock } = await fixture(kind); + const writeFile = fs.writeFile.bind(fs); + const failure = Object.assign(new Error("injected partial owner write"), { code: "ENOSPC" }); + let injected = false; + vi.spyOn(fs, "writeFile").mockImplementation(async (...args) => { + if (path.dirname(String(args[0])).startsWith(`${lock}.stage-`)) { + injected = true; + await writeFile(args[0], "{"); + throw failure; + } + return writeFile(...args); + }); + await expect( + operate(kind, target, async () => { + throw new Error("Operation must not run"); + }), + ).rejects.toBe(failure); + expect(injected).toBe(true); + expect(await fs.readdir(root)).toEqual(kind === "directory" ? [] : ["manifest.json"]); + if (kind === "file") expect(await fs.readFile(target, "utf8")).toBe("old"); + }, + ); + + it("reports the private preparation directory when partial owner cleanup also fails", async () => { + const { target, lock } = await fixture("directory"); + const writeFile = fs.writeFile.bind(fs); + const rm = fs.rm.bind(fs); + let stage = ""; + const writeFailure = Object.assign(new Error("injected partial owner write"), { + code: "ENOSPC", + }); + const cleanupFailure = Object.assign(new Error("injected preparation cleanup failure"), { + code: "EACCES", + }); + vi.spyOn(fs, "writeFile").mockImplementation(async (...args) => { + if (path.dirname(String(args[0])).startsWith(`${lock}.stage-`)) { + stage = path.dirname(String(args[0])); + await writeFile(args[0], "{"); + throw writeFailure; + } + return writeFile(...args); + }); + vi.spyOn(fs, "rm").mockImplementation(async (...args) => { + if (String(args[0]) === stage) throw cleanupFailure; + return rm(...args); + }); + let error: unknown; + try { + await withDirectoryTargetLock(target, async () => {}); + } catch (caught) { + error = caught; + } + expect(error).toBeInstanceOf(AggregateError); + expect((error as AggregateError).errors).toEqual([writeFailure, cleanupFailure]); + expect((error as Error).message).toContain(JSON.stringify(stage)); + const entries = await fs.readdir(stage); + expect(entries).toHaveLength(1); + expect(await fs.readFile(path.join(stage, entries[0]), "utf8")).toBe("{"); + await expect(fs.access(lock)).rejects.toMatchObject({ code: "ENOENT" }); + }); + + it("preserves an unrelated directory that collides with private lock preparation", async () => { + const { target, lock } = await fixture("directory"); + const mkdir = fs.mkdir.bind(fs); + let collision = ""; + vi.spyOn(fs, "mkdir").mockImplementation((async (...args: Parameters) => { + if (String(args[0]).startsWith(`${lock}.stage-`)) { + collision = String(args[0]); + await mkdir(...args); + await fs.writeFile(path.join(collision, "unrelated.txt"), "unrelated"); + } + return mkdir(...args); + }) as typeof fs.mkdir); + await expect(withDirectoryTargetLock(target, async () => {})).rejects.toMatchObject({ + code: "EEXIST", + }); + expect(await fs.readFile(path.join(collision, "unrelated.txt"), "utf8")).toBe("unrelated"); + await expect(fs.access(lock)).rejects.toMatchObject({ code: "ENOENT" }); + }); + + it.each(["directory", "file"])( + "preserves a junction used as a %s lock owner record", + async (kind) => { + const { root, target, lock } = await fixture(kind); + await fs.mkdir(lock); + const outside = path.join(root, "unrelated"); + await fs.mkdir(outside); + await fs.writeFile(path.join(outside, "untouched.txt"), "unrelated"); + const ownerPath = path.join(lock, "owner.json"); + await fs.symlink(outside, ownerPath, "junction"); + const old = new Date(Date.now() - 60_000); + await fs.utimes(lock, old, old); + await expect(operate(kind, target, async () => {})).rejects.toThrow(JSON.stringify(lock)); + expect((await fs.lstat(ownerPath)).isSymbolicLink()).toBe(true); + expect(await fs.readFile(path.join(outside, "untouched.txt"), "utf8")).toBe("unrelated"); + }, + ); + + it.each(["directory", "file"])( + "keeps %s operations exclusive when two stale-owner observations race", + async (kind) => { + const { target, lock } = await fixture(kind); + await fs.mkdir(lock); + await fs.writeFile(path.join(lock, "owner.json"), '{"pid":2147483647}'); + const readFile = fs.readFile.bind(fs); + const rm = fs.rm.bind(fs); + const unlink = fs.unlink.bind(fs); + const snapshots = deferred(); + const firstEntered = deferred(); + const releaseFirst = deferred(); + let reads = 0; + let delayedRemoval = false; + let active = 0; + let maximumActive = 0; + vi.spyOn(fs, "readFile").mockImplementation((async ( + ...args: Parameters + ) => { + const value = await readFile(...args); + if (path.dirname(String(args[0])) === lock) { + const owner = JSON.parse(value.toString()); + if (owner.pid === 2147483647 && reads < 2) { + if (++reads === 2) snapshots.resolve(); + await snapshots.promise; + } else if (owner.pid === process.pid && context.getStore() === "second" && active === 1) { + // Correct recovery observes the first live owner before it can enter. + releaseFirst.resolve(); + } + } + return value; + }) as typeof fs.readFile); + async function pauseStaleRemoval(name: string) { + if ( + context.getStore() === "second" && + !delayedRemoval && + (name === lock || path.dirname(name) === lock) + ) { + delayedRemoval = true; + await firstEntered.promise; + } + } + vi.spyOn(fs, "rm").mockImplementation(async (...args) => { + await pauseStaleRemoval(String(args[0])); + return rm(...args); + }); + vi.spyOn(fs, "unlink").mockImplementation(async (...args) => { + await pauseStaleRemoval(String(args[0])); + return unlink(...args); + }); + const critical = async () => { + active += 1; + maximumActive = Math.max(maximumActive, active); + if (context.getStore() === "first") { + firstEntered.resolve(); + await releaseFirst.promise; + } else { + // The pre-fix stale reaper deletes the new lock and enters concurrently. + releaseFirst.resolve(); + } + active -= 1; + }; + const runs = ["first", "second"].map((label) => + context.run(label, () => operate(kind, target, critical)), + ); + let timer: ReturnType | undefined; + try { + const outcomes = await Promise.race([ + Promise.allSettled(runs), + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error("Stale-owner race checkpoint timed out")), + 3_000, + ); + }), + ]); + expect(reads).toBe(2); + expect(delayedRemoval).toBe(true); + expect(maximumActive).toBe(1); + if (kind === "directory") + expect(outcomes.map((item) => item.status)).toEqual(["fulfilled", "fulfilled"]); + else { + expect(outcomes[0].status).toBe("fulfilled"); + expect(outcomes[1]).toMatchObject({ + status: "rejected", + reason: { message: `File changed before writing: ${target}` }, + }); + expect(await fs.readFile(target, "utf8")).toBe("first"); + } + } finally { + if (timer) clearTimeout(timer); + firstEntered.resolve(); + snapshots.resolve(); + releaseFirst.resolve(); + await Promise.allSettled(runs); + } + await expect(fs.access(lock)).rejects.toMatchObject({ code: "ENOENT" }); + }, + 15_000, + ); + + it.each(["directory", "file"])( + "preserves unrelated contents of an ambiguous %s lock directory", + async (kind) => { + const { target, lock } = await fixture(kind); + await fs.mkdir(lock); + await fs.writeFile(path.join(lock, "unrelated.txt"), "unrelated"); + const old = new Date(Date.now() - 60_000); + await fs.utimes(lock, old, old); + await expect(operate(kind, target, async () => {})).rejects.toThrow(JSON.stringify(lock)); + expect(await fs.readFile(path.join(lock, "unrelated.txt"), "utf8")).toBe("unrelated"); + if (kind === "file") expect(await fs.readFile(target, "utf8")).toBe("old"); + }, + ); + + it.each(["directory", "file"])( + "rejects a %s lock junction without deleting it or following its owner record", + async (kind) => { + const { root, target, lock } = await fixture(kind); + const outside = path.join(root, "unrelated"); + await fs.mkdir(outside); + await fs.writeFile(path.join(outside, "owner.json"), '{"pid":2147483647}'); + await fs.writeFile(path.join(outside, "untouched.txt"), "unrelated"); + await fs.symlink(outside, lock, "junction"); + await expect(operate(kind, target, async () => {})).rejects.toThrow(JSON.stringify(lock)); + expect((await fs.lstat(lock)).isSymbolicLink()).toBe(true); + expect(await fs.readFile(path.join(outside, "untouched.txt"), "utf8")).toBe("unrelated"); + expect(await fs.readFile(path.join(outside, "owner.json"), "utf8")).toBe( + '{"pid":2147483647}', + ); + }, + ); +}); diff --git a/tests/fixtures/filesystem-lock-crash-worker.ts b/tests/fixtures/filesystem-lock-crash-worker.ts new file mode 100644 index 0000000..204c933 --- /dev/null +++ b/tests/fixtures/filesystem-lock-crash-worker.ts @@ -0,0 +1,19 @@ +import fs from "node:fs/promises"; +import path from "node:path"; +import { withDirectoryTargetLock } from "../../src/directory-swap"; +import { writeFileChanges } from "../../src/file-changes"; + +const [kind, target] = process.argv.slice(2); +const writeFile = fs.writeFile.bind(fs); +fs.writeFile = async (name, ...args) => { + await writeFile(name, ...args); + if ( + path.basename(String(name)).startsWith("owner-") && + path.dirname(String(name)).includes(".askr-lock.stage-") + ) { + process.send!({ stage: path.dirname(String(name)) }); + await new Promise((resolve) => process.once("message", () => resolve())); + } +}; +if (kind === "directory") await withDirectoryTargetLock(target, async () => {}); +else await writeFileChanges([{ filePath: target, content: "new", expectedContent: "old" }]); diff --git a/tests/public-contract.json b/tests/public-contract.json index 3c882c6..763b27a 100644 --- a/tests/public-contract.json +++ b/tests/public-contract.json @@ -28,6 +28,7 @@ "generate/generator", "directory-swap", "file-changes", + "filesystem-lock", "ssg/sitemap", "ssg/output-report", "update",