From e2d1117614f6247da831c9db1706eaf7e3cde03b Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 3 Aug 2026 22:42:57 +0200 Subject: [PATCH] fix: make workspace registry lock process-bound --- backend/src/workspaces/git-repository.ts | 134 ++++++++++-------- backend/test/workspace-registry.test.ts | 20 ++- .../test/workspaces-git-repository.test.ts | 41 ++++++ 3 files changed, 128 insertions(+), 67 deletions(-) diff --git a/backend/src/workspaces/git-repository.ts b/backend/src/workspaces/git-repository.ts index fef70512..f53a6e02 100644 --- a/backend/src/workspaces/git-repository.ts +++ b/backend/src/workspaces/git-repository.ts @@ -1,7 +1,5 @@ -import { execFile } from "node:child_process"; -import { - closeSync, constants, lstatSync, mkdirSync, openSync, readFileSync, unlinkSync, writeFileSync, -} from "node:fs"; +import { execFile, spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import { lstatSync, mkdirSync } from "node:fs"; import { mkdir } from "node:fs/promises"; import { basename, isAbsolute, join } from "node:path"; import { promisify } from "node:util"; @@ -228,78 +226,88 @@ export class WorkspaceRepositoryLock { this.queue = new Promise((resolve) => { releaseQueue = resolve; }); await previous; - let descriptor: number | undefined; + let holder: ChildProcessWithoutNullStreams | undefined; const lockPath = join(this.locksPath, "repository.lock"); try { - mkdirSync(this.locksPath, { recursive: true, mode: 0o700 }); - assertDirectory(this.locksPath); - } catch (error) { - throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"); - } - try { - descriptor = this.acquire(lockPath); - } catch (error) { - throw this.lockError(error); - } - try { + try { + mkdirSync(this.locksPath, { recursive: true, mode: 0o700 }); + assertDirectory(this.locksPath); + } catch (error) { + throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"); + } + try { + holder = await this.acquire(lockPath); + } catch (error) { + throw this.lockError(error); + } return await operation(); } finally { - if (descriptor !== undefined) closeSync(descriptor); - if (descriptor !== undefined) { - try { unlinkSync(lockPath); } catch { /* stale lock cleanup is retried by the operator */ } - } - releaseQueue(); - } - } - - private acquire(lockPath: string): number { - try { - return this.createProcessLock(lockPath); - } catch (error) { - if (!this.recoverDeadProcessLock(lockPath, error)) throw error; - return this.createProcessLock(lockPath); - } - } - - private createProcessLock(lockPath: string): number { - const descriptor = openSync(lockPath, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600); - try { - writeFileSync(descriptor, JSON.stringify({ pid: process.pid }), "utf8"); - return descriptor; - } catch (error) { - closeSync(descriptor); - try { unlinkSync(lockPath); } catch { /* incomplete lock is never treated as recoverable */ } - throw error; - } - } - - private recoverDeadProcessLock(lockPath: string, error: unknown): boolean { - if (!(typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST")) { - return false; - } - try { - const record = JSON.parse(readFileSync(lockPath, "utf8")) as { pid?: unknown }; - const pid = record.pid; - if (typeof pid !== "number" || !Number.isSafeInteger(pid) || pid <= 0 || pid === process.pid) return false; try { - process.kill(pid, 0); - return false; - } catch (probeError) { - if (!(typeof probeError === "object" && probeError !== null && "code" in probeError && probeError.code === "ESRCH")) { - return false; - } + if (holder !== undefined) await this.release(holder); + } finally { + releaseQueue(); } - unlinkSync(lockPath); - return true; - } catch { - return false; } } + private async acquire(lockPath: string): Promise { + try { + const entry = lstatSync(lockPath); + if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("invalid lock path"); + } catch (error) { + if (!(typeof error === "object" && error !== null && "code" in error && error.code === "ENOENT")) { + throw error; + } + } + const holder = spawn("python3", ["-c", WorkspaceRepositoryLock.HOLDER_PROGRAM, lockPath], { + stdio: ["pipe", "pipe", "pipe"], + }); + await new Promise((resolve, reject) => { + let output = ""; + const fail = (error: WorkspaceRegistryError) => { + holder.stdout.removeAllListeners("data"); + reject(error); + }; + holder.once("error", () => fail(new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"))); + holder.once("exit", (code) => { + fail(new WorkspaceRegistryError( + code === 73 ? "workspace_stale" : "git_unavailable", + code === 73 ? "Workspace registry is busy" : "Workspace registry lock is unavailable", + )); + }); + holder.stdout.on("data", (chunk: Buffer) => { + output += chunk.toString("utf8"); + if (output === "locked\n") { + holder.stdout.removeAllListeners("data"); + resolve(); + } + }); + }); + return holder; + } + + private async release(holder: ChildProcessWithoutNullStreams): Promise { + if (!holder.stdin.destroyed) holder.stdin.end(); + await new Promise((resolve) => holder.once("exit", () => resolve())); + } + private lockError(error: unknown): WorkspaceRegistryError { + if (error instanceof WorkspaceRegistryError) return error; if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") { return new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy"); } return new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"); } + + private static readonly HOLDER_PROGRAM = [ + "import fcntl, os, sys", + "fd = os.open(sys.argv[1], os.O_RDWR | os.O_CREAT | getattr(os, 'O_NOFOLLOW', 0), 0o600)", + "try:", + " fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)", + "except BlockingIOError:", + " sys.exit(73)", + "sys.stdout.write('locked\\n')", + "sys.stdout.flush()", + "sys.stdin.buffer.read()", + ].join("\n"); } diff --git a/backend/test/workspace-registry.test.ts b/backend/test/workspace-registry.test.ts index 86a0dd4e..18a42428 100644 --- a/backend/test/workspace-registry.test.ts +++ b/backend/test/workspace-registry.test.ts @@ -6,6 +6,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; import { afterEach, expect, test } from "vitest"; +import { WorkspaceRepositoryLock } from "../src/workspaces/git-repository.js"; import { WorkspaceRegistry } from "../src/workspaces/registry.js"; import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js"; @@ -113,14 +114,25 @@ test("keeps the last valid snapshot when a pulled commit has invalid YAML", asyn }); }); -test("does not bypass an existing advisory repository lock", async () => { +test("does not bypass an existing live advisory repository lock", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); const registry = new WorkspaceRegistry(config(root, remote.remote)); - mkdirSync(join(root, "locks"), { recursive: true }); - writeFileSync(join(root, "locks", "repository.lock"), "held"); + const lock = new WorkspaceRepositoryLock(join(root, "locks")); + let release!: () => void; + let started!: () => void; + const held = lock.run(async () => { + started(); + await new Promise((resolve) => { release = resolve; }); + }); + await new Promise((resolve) => { started = resolve; }); - await expect(registry.bootstrap()).rejects.toMatchObject({ code: "workspace_stale" }); + try { + await expect(registry.bootstrap()).rejects.toMatchObject({ code: "workspace_stale" }); + } finally { + release(); + await held; + } }); test("rejects a symbolic-link registry root before creating a lock below it", async () => { diff --git a/backend/test/workspaces-git-repository.test.ts b/backend/test/workspaces-git-repository.test.ts index cdf95f5a..c656fc99 100644 --- a/backend/test/workspaces-git-repository.test.ts +++ b/backend/test/workspaces-git-repository.test.ts @@ -137,3 +137,44 @@ test("maps lock filesystem failures to a stable redacted error", async () => { }); expect((error as Error).message).not.toContain(file); }); + +test("releases its queue after a malformed lock failure so a later attempt can acquire", async () => { + const root = mkdtempSync(join(tmpdir(), "thoth-workspace-git-queue-")); + temporaryRoots.push(root); + const locks = join(root, "locks"); + mkdirSync(join(locks, "repository.lock"), { recursive: true }); + const lock = new WorkspaceRepositoryLock(locks); + + await expect(lock.run(async () => "unreachable")).rejects.toMatchObject({ code: "git_unavailable" }); + rmSync(join(locks, "repository.lock"), { recursive: true, force: true }); + + const result = await Promise.race([ + lock.run(async () => "recovered"), + new Promise((resolve) => setTimeout(() => resolve("timed out"), 250)), + ]); + expect(result).toBe("recovered"); +}); + +test("parallel contenders recover a stale lock file without overlapping critical sections", async () => { + const root = mkdtempSync(join(tmpdir(), "thoth-workspace-git-contenders-")); + temporaryRoots.push(root); + const locks = join(root, "locks"); + mkdirSync(locks); + writeFileSync(join(locks, "repository.lock"), JSON.stringify({ pid: 999_999_999 })); + const first = new WorkspaceRepositoryLock(locks); + const second = new WorkspaceRepositoryLock(locks); + let active = 0; + let maximum = 0; + const critical = async () => { + active += 1; + maximum = Math.max(maximum, active); + await new Promise((resolve) => setTimeout(resolve, 25)); + active -= 1; + }; + + const results = await Promise.allSettled([first.run(critical), second.run(critical)]); + + expect(results.filter((result) => result.status === "fulfilled")).toHaveLength(1); + expect(results.filter((result) => result.status === "rejected")).toHaveLength(1); + expect(maximum).toBe(1); +});