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]>
212 lines
9.3 KiB
JavaScript
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}` };
|
|
}
|