Files
stack/packages/queue/src/lock.mjs
T
jason.woltjeandClaude Opus 5.5 34a72af912 feat(queue): queue as data A1, journal, lock, CLI and verify (#1508)
packages/queue, scripts/queue-commit.sh, scripts/git-hooks and
scripts/test-queue.sh, plus docs/plans/BRIEF-TEMPLATE.md. There is no
queue.json yet, so verify skips until the genesis commit after A2.

Darkwing built it, and Filbert reviewed R0 (6933b885, changes requested)
and r1 (e464be6c, approved). The 20 files match manifest 85a8a453. The
nine suites passed on an index export, including the new queue suite.
test-queue.sh joins the suite list in AGENTS.md. Lead decisions 20, 23
and 26.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-09-26 19:07:48 -05:00

201 lines
8.8 KiB
JavaScript

// The queue lock (8.4): <root>/.git/mosaic-queue.lock, published by link()
// from an O_EXCL temp file, and the unlock gate <root>/.git/mosaic-queue.unlock
// made the same way. Nothing is ever removed because of its age. Only
// `unlock` removes another process's lock, and only a dead or mismatched one.
import { randomBytes } from "node:crypto";
import { hostname } from "node:os";
import { join } from "node:path";
import { bootId, pidAlive, processStart, validBoot, validStart } from "../../discord/src/journal.mjs";
import { QueueError } from "./errors.mjs";
import { errno, lstatOrNull, readOrNull, sleepMs, unlinkQuiet, writeAll } from "./io.mjs";
export const LOCK_NAME = "mosaic-queue.lock";
export const GATE_NAME = "mosaic-queue.unlock";
export const realProc = Object.freeze({
pid: process.pid,
host: () => hostname(),
processStart: (pid) => processStart(pid),
bootId: () => bootId(),
pidAlive: (pid) => pidAlive(pid),
});
const RECORD_KEYS = ["pid", "start", "boot", "host", "op", "verb", "at"];
function ownRecord(proc, op, verb) {
const start = proc.processStart(proc.pid);
const boot = proc.bootId();
if (start === null || boot === null) throw new QueueError("cannot read this process's start time or boot id from /proc; refusing to take the lock", 2);
return { pid: proc.pid, start, boot, host: proc.host(), op: op ?? null, verb, at: new Date().toISOString() };
}
export function parseRecord(bytes) {
if (bytes === null || bytes.length === 0) return null;
let rec;
try { rec = JSON.parse(bytes.toString("utf8")); } catch { return null; }
if (rec === null || typeof rec !== "object" || Array.isArray(rec)) return null;
if (Object.keys(rec).join() !== RECORD_KEYS.join()) return null;
if (!Number.isInteger(rec.pid) || rec.pid <= 0) return null;
if (!validStart(rec.start) || !validBoot(rec.boot)) return null;
if (typeof rec.host !== "string" || rec.host.length === 0) return null;
if (rec.op !== null && typeof rec.op !== "string") return null;
if (typeof rec.verb !== "string" || typeof rec.at !== "string") return null;
return rec;
}
// The 8.4 table, in order; the first test that applies decides.
export function classify(bytes, proc) {
const rec = parseRecord(bytes);
if (rec === null) return { state: "invalid", rec: null };
const selfStart = proc.processStart(proc.pid);
const selfBoot = proc.bootId();
if (selfStart === null || selfBoot === null) return { state: "unknown", rec, why: "this process cannot read its own /proc identity" };
if (rec.host !== proc.host()) return { state: "unknown", rec, why: `recorded on host ${rec.host}` };
if (rec.boot !== selfBoot) return { state: "mismatch", rec, why: "recorded in a previous boot" };
if (!proc.pidAlive(rec.pid)) return { state: "dead", rec };
const start = proc.processStart(rec.pid);
if (start === null) return { state: "unknown", rec, why: `pid ${rec.pid} is alive but its start time is unreadable` };
if (start !== rec.start) return { state: "mismatch", rec, why: `pid ${rec.pid} was reused` };
return { state: "live", rec };
}
function describe(c) {
if (c.rec === null) return c.state;
const r = c.rec;
return `${c.state}: pid ${r.pid}, ${r.verb}${r.op ? ` ${r.op}` : ""} since ${r.at}${c.why ? `; ${c.why}` : ""}`;
}
// Temp file with the whole record, fsynced, read back, then link()ed to
// `target`. Returns {linked, dev, ino, bytes}; linked is false on EEXIST.
function publish(io, target, bytes, hook, waitMs, stepMs) {
const tmp = `${target}.${process.pid}.${randomBytes(6).toString("hex")}.tmp`;
let fd;
let st;
try {
fd = io.openExcl(tmp, 0o600);
} catch (err) {
throw new QueueError(`cannot create ${tmp}: ${errno(err)}`, 1);
}
try {
writeAll(io, fd, bytes);
io.fsync(fd);
io.close(fd);
fd = null;
const back = io.readFile(tmp);
if (!back.equals(bytes)) throw Object.assign(new Error("read-back differs from the record"), { code: "EREADBACK" });
// The link gives the target this inode, so read it before linking:
// nothing that can fail runs between a successful link and the return.
st = io.stat(tmp);
} catch (err) {
if (fd !== null) { try { io.close(fd); } catch { /* already failing */ } }
unlinkQuiet(io, tmp);
throw new QueueError(`cannot write the lock record ${tmp}: ${errno(err)}; no lock taken`, 1);
}
hook("lock-temp-written");
const deadline = Date.now() + waitMs;
try {
for (;;) {
try {
io.link(tmp, target);
} catch (err) {
if (err.code !== "EEXIST") throw new QueueError(`cannot link ${target}: ${errno(err)}; no lock taken`, 1);
if (Date.now() >= deadline) return { linked: false };
sleepMs(stepMs);
continue;
}
return { linked: true, dev: st.dev, ino: st.ino, bytes };
}
} finally {
unlinkQuiet(io, tmp);
}
}
export function acquire({ gitDir, io, proc = realProc, op = null, verb, waitMs = 10000, stepMs = 100, hook = () => {} }) {
const path = join(gitDir, LOCK_NAME);
const rec = ownRecord(proc, op, verb);
const bytes = Buffer.from(JSON.stringify(rec) + "\n");
const got = publish(io, path, bytes, hook, waitMs, stepMs);
if (!got.linked) {
const c = classify(readOrNull(io, path), proc);
if (c.state === "live") throw new QueueError(`queue lock held by ${c.rec.verb}${c.rec.op ? ` ${c.rec.op}` : ""} since ${c.rec.at}; retry the same op later`, 2);
if (c.state === "dead" || c.state === "mismatch") throw new QueueError(`queue lock owner is ${describe(c)}; run \`scripts/mosaic queue unlock\` once nothing is running`, 2);
if (c.state === "unknown") throw new QueueError(`queue lock owner is ${describe(c)}; unlock refuses this too; diagnose ${path} by hand`, 2);
throw new QueueError(`queue lock record is invalid; inspect ${path} by hand`, 2);
}
const handle = { path, dev: got.dev, ino: got.ino, bytes: got.bytes };
hook("lock-linked");
const gate = join(gitDir, GATE_NAME);
let c = null;
try {
if (lstatOrNull(io, gate) !== null) c = classify(readOrNull(io, gate), proc);
} catch (err) {
const left = release(handle, io);
throw new QueueError(`cannot check the unlock gate ${gate}: ${errno(err)}; ${left ?? "lock released"}`, 1);
}
if (c !== null) {
release(handle, io);
throw new QueueError(`unlock gate ${gate} is present (${describe(c)}); check it with \`scripts/mosaic queue unlock --check-gate\``, 2);
}
return handle;
}
// Unlinks only the lock this process published: same inode, same record.
export function release(handle, io) {
const st = lstatOrNull(io, handle.path);
if (st === null) return `lock ${handle.path} was already gone`;
const bytes = readOrNull(io, handle.path);
if (st.dev !== handle.dev || st.ino !== handle.ino || bytes === null || !bytes.equals(handle.bytes)) {
return `lock ${handle.path} is not the one this process took; left in place`;
}
io.unlink(handle.path);
return null;
}
export function unlock({ gitDir, io, proc = realProc, hook = () => {} }) {
const lockPath = join(gitDir, LOCK_NAME);
const gatePath = join(gitDir, GATE_NAME);
const rec = ownRecord(proc, null, "unlock");
const bytes = Buffer.from(JSON.stringify(rec) + "\n");
const got = publish(io, gatePath, bytes, () => {}, 0, 0);
if (!got.linked) {
const c = classify(readOrNull(io, gatePath), proc);
throw new QueueError(`unlock gate ${gatePath} is held (${describe(c)}); another unlock is running or a stale gate needs \`--check-gate\``, 2);
}
const gate = { path: gatePath, dev: got.dev, ino: got.ino, bytes };
hook("gate-held");
let result;
let failure = null;
try {
const lockBytes = readOrNull(io, lockPath);
if (lockBytes === null) {
result = "no queue lock present; nothing removed";
} else {
const c = classify(lockBytes, proc);
if (c.state !== "dead" && c.state !== "mismatch") {
throw new QueueError(`queue lock owner is ${describe(c)}; unlock refuses`, 2);
}
io.unlink(lockPath);
result = `removed queue lock (${describe(c)}): ${lockBytes.toString("utf8").trim()}`;
}
} catch (err) {
failure = err;
}
let msg;
try { msg = release(gate, io); } catch (err) { msg = `cannot release the unlock gate (${errno(err)})`; }
// Like the lock, a swapped gate is reported on success and on refusal (8.4).
if (msg && failure instanceof Error) failure.message += `\nwarning: ${msg}`;
if (failure) throw failure;
return msg ? `${result}\nwarning: ${msg}` : result;
}
export function checkGate({ gitDir, io, proc = realProc }) {
const gatePath = join(gitDir, GATE_NAME);
const bytes = readOrNull(io, gatePath);
if (bytes === null) return { state: "absent", line: "no unlock gate present" };
const c = classify(bytes, proc);
const advice = c.state === "dead" || c.state === "mismatch"
? `; remove ${gatePath} by hand only once no queue command is running`
: "; leave it for diagnosis";
return { state: c.state, line: `unlock gate ${describe(c)}${advice}` };
}