fix: make workspace registry lock process-bound
This commit is contained in:
@@ -1,7 +1,5 @@
|
|||||||
import { execFile } from "node:child_process";
|
import { execFile, spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
|
||||||
import {
|
import { lstatSync, mkdirSync } from "node:fs";
|
||||||
closeSync, constants, lstatSync, mkdirSync, openSync, readFileSync, unlinkSync, writeFileSync,
|
|
||||||
} from "node:fs";
|
|
||||||
import { mkdir } from "node:fs/promises";
|
import { mkdir } from "node:fs/promises";
|
||||||
import { basename, isAbsolute, join } from "node:path";
|
import { basename, isAbsolute, join } from "node:path";
|
||||||
import { promisify } from "node:util";
|
import { promisify } from "node:util";
|
||||||
@@ -228,78 +226,88 @@ export class WorkspaceRepositoryLock {
|
|||||||
this.queue = new Promise<void>((resolve) => { releaseQueue = resolve; });
|
this.queue = new Promise<void>((resolve) => { releaseQueue = resolve; });
|
||||||
await previous;
|
await previous;
|
||||||
|
|
||||||
let descriptor: number | undefined;
|
let holder: ChildProcessWithoutNullStreams | undefined;
|
||||||
const lockPath = join(this.locksPath, "repository.lock");
|
const lockPath = join(this.locksPath, "repository.lock");
|
||||||
try {
|
try {
|
||||||
mkdirSync(this.locksPath, { recursive: true, mode: 0o700 });
|
try {
|
||||||
assertDirectory(this.locksPath);
|
mkdirSync(this.locksPath, { recursive: true, mode: 0o700 });
|
||||||
} catch (error) {
|
assertDirectory(this.locksPath);
|
||||||
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable");
|
} catch (error) {
|
||||||
}
|
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable");
|
||||||
try {
|
}
|
||||||
descriptor = this.acquire(lockPath);
|
try {
|
||||||
} catch (error) {
|
holder = await this.acquire(lockPath);
|
||||||
throw this.lockError(error);
|
} catch (error) {
|
||||||
}
|
throw this.lockError(error);
|
||||||
try {
|
}
|
||||||
return await operation();
|
return await operation();
|
||||||
} finally {
|
} 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 {
|
try {
|
||||||
process.kill(pid, 0);
|
if (holder !== undefined) await this.release(holder);
|
||||||
return false;
|
} finally {
|
||||||
} catch (probeError) {
|
releaseQueue();
|
||||||
if (!(typeof probeError === "object" && probeError !== null && "code" in probeError && probeError.code === "ESRCH")) {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
unlinkSync(lockPath);
|
|
||||||
return true;
|
|
||||||
} catch {
|
|
||||||
return false;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async acquire(lockPath: string): Promise<ChildProcessWithoutNullStreams> {
|
||||||
|
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<void>((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<void> {
|
||||||
|
if (!holder.stdin.destroyed) holder.stdin.end();
|
||||||
|
await new Promise<void>((resolve) => holder.once("exit", () => resolve()));
|
||||||
|
}
|
||||||
|
|
||||||
private lockError(error: unknown): WorkspaceRegistryError {
|
private lockError(error: unknown): WorkspaceRegistryError {
|
||||||
|
if (error instanceof WorkspaceRegistryError) return error;
|
||||||
if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") {
|
if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") {
|
||||||
return new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy");
|
return new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy");
|
||||||
}
|
}
|
||||||
return new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable");
|
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");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import { tmpdir } from "node:os";
|
|||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import { promisify } from "node:util";
|
import { promisify } from "node:util";
|
||||||
import { afterEach, expect, test } from "vitest";
|
import { afterEach, expect, test } from "vitest";
|
||||||
|
import { WorkspaceRepositoryLock } from "../src/workspaces/git-repository.js";
|
||||||
import { WorkspaceRegistry } from "../src/workspaces/registry.js";
|
import { WorkspaceRegistry } from "../src/workspaces/registry.js";
|
||||||
import type { WorkspaceRegistryConfig } from "../src/workspaces/types.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 remote = await fixture();
|
||||||
const root = join(remote.root, "registry");
|
const root = join(remote.root, "registry");
|
||||||
const registry = new WorkspaceRegistry(config(root, remote.remote));
|
const registry = new WorkspaceRegistry(config(root, remote.remote));
|
||||||
mkdirSync(join(root, "locks"), { recursive: true });
|
const lock = new WorkspaceRepositoryLock(join(root, "locks"));
|
||||||
writeFileSync(join(root, "locks", "repository.lock"), "held");
|
let release!: () => void;
|
||||||
|
let started!: () => void;
|
||||||
|
const held = lock.run(async () => {
|
||||||
|
started();
|
||||||
|
await new Promise<void>((resolve) => { release = resolve; });
|
||||||
|
});
|
||||||
|
await new Promise<void>((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 () => {
|
test("rejects a symbolic-link registry root before creating a lock below it", async () => {
|
||||||
|
|||||||
@@ -137,3 +137,44 @@ test("maps lock filesystem failures to a stable redacted error", async () => {
|
|||||||
});
|
});
|
||||||
expect((error as Error).message).not.toContain(file);
|
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<string>((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);
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user