Files
stack/packages/queue/src/lock.mjs
T
jason.woltjeandClaude Opus 5.5 6ca116b7ba feat(queue): queue as data A2, migration, render and dispatch (#1508)
Filbert approved round 1 (f167b85e). Manifest 782bcb62, 21 files, plus
the QUEUE.md markers and the TOOLS.md section. Lead decision 35.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-09-26 20:14:09 -05:00

212 lines
9.3 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 = releaseOrWarn(handle, io);
throw new QueueError(`cannot check the unlock gate ${gate}: ${errno(err)}; ${left ?? "lock released"}`, 1);
}
if (c !== null) {
const left = releaseOrWarn(handle, io);
throw new QueueError(`unlock gate ${gate} is present (${describe(c)}); check it with \`scripts/mosaic queue unlock --check-gate\`${left ? `\nwarning: ${left}` : ""}`, 2);
}
return handle;
}
// release() on a path that is already reporting something: a throw becomes
// a message naming the lock that may be left behind (P1).
export function releaseOrWarn(handle, io) {
try {
return release(handle, io);
} catch (err) {
return `cannot release the queue lock (${errno(err)}); ${handle.path} may be left in place`;
}
}
// 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;
// Apart, so the caller never splits a lock record's text to find it (P3).
return { result, warning: msg };
}
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}` };
}