Controller, claim store, live-session guard, engine link and seal, turn tracker, cohort force stop and recovery, client library, transcript and mediated terminal, with the fake engine and tests. Fixtures only; no live cutover. Dewey built it. Darkwing (comment 26690) and Filbert (comment 26694) approved round 2. Manifest I1-r2-manifest.sha256 (2b48e333, 27 files). Suites on an export: conversation 152/152, control-board 124, webui 14, seat 19, chat-00/01/01c checks, and all nine scripts/test-*.sh green. Follow-ups for I3 are in DEFERRED. Gate E stays with Jason. Co-Authored-By: Claude Opus 5.5 <[email protected]>
355 lines
16 KiB
JavaScript
355 lines
16 KiB
JavaScript
// The writer-claim record (#1507, CHAT-03 §2, D1).
|
|
//
|
|
// A claim holds two keys, the seat tuple (seat, project, workspace) and the
|
|
// native session identity. Each key is a directory under the claim root and
|
|
// each revision a numbered file in it (r0000000001.json, ...). A revision is
|
|
// published by writing a temporary file in the same directory, fsyncing it,
|
|
// link()ing it to the next revision name and fsyncing the directory. link()
|
|
// fails when the name exists, so publication is exclusive and a revision name
|
|
// only ever points at complete contents. Revisions are never rewritten.
|
|
//
|
|
// Keys are taken seat first, then session. Every step of the lifecycle
|
|
// publishes on the seat key and then on the session key. The pair's state is
|
|
// the more conservative of the two: uncertain > stopping > active > reserved
|
|
// > stopped. A key is held unless its highest revision is `stopped` with a
|
|
// proof reference. A highest revision that does not parse holds the key as
|
|
// `uncertain`; an older revision is never reused.
|
|
//
|
|
// The store is internal. It never crosses the wire, and it stores only the
|
|
// CHAT-01 binding fields it needs.
|
|
|
|
import { closeSync, fsyncSync, linkSync, mkdirSync, openSync, readdirSync, readFileSync, unlinkSync, writeSync } from "node:fs";
|
|
import { join } from "node:path";
|
|
import { createHash, randomBytes } from "node:crypto";
|
|
import { bootId as readBootId, identityOf, ownerState } from "../../discord/src/journal.mjs";
|
|
import { ControlRefusal } from "./safe-fs.mjs";
|
|
import { ID } from "./parts.mjs";
|
|
|
|
export const CLAIM_VERSION = 1;
|
|
export const STATES = Object.freeze(["uncertain", "stopping", "active", "reserved", "stopped"]);
|
|
export const PROOF_KINDS = Object.freeze(["cohortProof", "boot", "no-unit"]);
|
|
export const ALREADY_ACTIVE = "already-active";
|
|
export const UNSAFE_REPLACEMENT = "unsafe-replacement";
|
|
export const FOREIGN_HOST = "foreign-host";
|
|
|
|
const REVISION = /^r(\d{10})\.json$/;
|
|
const revName = (n) => `r${String(n).padStart(10, "0")}.json`;
|
|
const rank = (state) => STATES.indexOf(state);
|
|
|
|
export function machineId() {
|
|
try {
|
|
const id = readFileSync("/etc/machine-id", "utf8").trim();
|
|
return /^[0-9a-f]{32}$/.test(id) ? id : null;
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
export const defaultHost = Object.freeze({ machineId, bootId: readBootId });
|
|
|
|
export function newClaimId() {
|
|
return "c" + randomBytes(12).toString("hex");
|
|
}
|
|
|
|
export function unitNameFor(claimId) {
|
|
return `mosaic-chat-${claimId}`;
|
|
}
|
|
|
|
// The more conservative of two states.
|
|
export function conservative(a, b) {
|
|
return rank(a) <= rank(b) ? a : b;
|
|
}
|
|
|
|
const keyHash = (key) => createHash("sha256").update(JSON.stringify(key)).digest("hex").slice(0, 32);
|
|
|
|
export function seatKey({ seat, project, workspace }) {
|
|
return { kind: "seat", seat, project, workspace };
|
|
}
|
|
|
|
export function sessionKey(session) {
|
|
return { kind: "session", session };
|
|
}
|
|
|
|
function valid(rec, key, n) {
|
|
return rec && typeof rec === "object" && rec.version === CLAIM_VERSION && rec.kind === "writer-claim" &&
|
|
rec.revision === n && JSON.stringify(rec.key) === JSON.stringify(key) && typeof rec.claimId === "string" &&
|
|
ID.test(rec.claimId) && STATES.includes(rec.state) && rec.host && typeof rec.host === "object" &&
|
|
(rec.proof === null || (rec.proof && PROOF_KINDS.includes(rec.proof.kind)));
|
|
}
|
|
|
|
// A key's highest revision is free only when it is `stopped` with a proof.
|
|
export function held(head) {
|
|
if (head.n === 0) return false;
|
|
if (head.damaged) return true;
|
|
return !(head.record.state === "stopped" && head.record.proof);
|
|
}
|
|
|
|
export function headState(head) {
|
|
if (head.n === 0) return "stopped";
|
|
if (head.damaged) return "uncertain";
|
|
if (head.record.state === "stopped" && !head.record.proof) return "uncertain";
|
|
return head.record.state;
|
|
}
|
|
|
|
export class ClaimStore {
|
|
constructor({ root, host = defaultHost, identity = identityOf, barrier = null, now = () => new Date() }) {
|
|
if (typeof root !== "string" || !root) throw new ControlRefusal("configuration", "a claim root is required; CHAT-03 has no default");
|
|
this.root = root;
|
|
this.host = host;
|
|
this.identity = identity;
|
|
this.barrier = barrier;
|
|
this.now = now;
|
|
}
|
|
|
|
dir(key) {
|
|
return join(this.root, key.kind, keyHash(key));
|
|
}
|
|
|
|
async pause(name, detail) {
|
|
if (this.barrier) await this.barrier(name, detail);
|
|
}
|
|
|
|
head(key) {
|
|
const dir = this.dir(key);
|
|
let names;
|
|
try {
|
|
names = readdirSync(dir);
|
|
} catch (err) {
|
|
if (err.code === "ENOENT") return { key, n: 0, record: null, damaged: false };
|
|
throw err;
|
|
}
|
|
let n = 0;
|
|
for (const name of names) {
|
|
const m = REVISION.exec(name);
|
|
if (m) n = Math.max(n, Number(m[1]));
|
|
}
|
|
if (n === 0) return { key, n: 0, record: null, damaged: false };
|
|
let record = null;
|
|
try {
|
|
record = JSON.parse(readFileSync(join(dir, revName(n)), "utf8"));
|
|
} catch {
|
|
record = null;
|
|
}
|
|
if (!valid(record, key, n)) return { key, n, record: null, damaged: true };
|
|
return { key, n, record, damaged: false };
|
|
}
|
|
|
|
revisions(key) {
|
|
const dir = this.dir(key);
|
|
let names;
|
|
try {
|
|
names = readdirSync(dir);
|
|
} catch {
|
|
return [];
|
|
}
|
|
return names.filter((x) => REVISION.test(x)).sort().map((name) => ({ name, text: readFileSync(join(dir, name), "utf8") }));
|
|
}
|
|
|
|
// Publishes revision `n` on `key`. Returns null when another writer
|
|
// published that revision first.
|
|
async publish(key, n, body) {
|
|
const dir = this.dir(key);
|
|
mkdirSync(dir, { recursive: true, mode: 0o700 });
|
|
const record = { ...body, key, revision: n, writtenAt: this.now().toISOString() };
|
|
const tmp = join(dir, `.tmp-${randomBytes(8).toString("hex")}`);
|
|
const fd = openSync(tmp, "wx", 0o400);
|
|
try {
|
|
writeSync(fd, JSON.stringify(record) + "\n");
|
|
await this.pause("temp-written", { key, n });
|
|
fsyncSync(fd);
|
|
} finally {
|
|
closeSync(fd);
|
|
}
|
|
await this.pause("temp-synced", { key, n });
|
|
let won = true;
|
|
try {
|
|
linkSync(tmp, join(dir, revName(n)));
|
|
} catch (err) {
|
|
if (err.code !== "EEXIST") throw err;
|
|
won = false;
|
|
}
|
|
unlinkSync(tmp);
|
|
if (!won) return null;
|
|
await this.pause("linked", { key, n });
|
|
const dfd = openSync(dir, "r");
|
|
try {
|
|
fsyncSync(dfd);
|
|
} finally {
|
|
closeSync(dfd);
|
|
}
|
|
await this.pause("dir-synced", { key, n });
|
|
return record;
|
|
}
|
|
|
|
owner(incarnation) {
|
|
const id = this.identity(process.pid);
|
|
return { pid: process.pid, start: id.start, boot: id.boot, incarnation };
|
|
}
|
|
|
|
// Why a held pair refuses a new claim.
|
|
refusal(seat, session) {
|
|
const heads = [seat, session].filter(held);
|
|
const mine = this.host.machineId();
|
|
if (heads.some((h) => !h.damaged && h.record.host.machineId !== mine)) return FOREIGN_HOST;
|
|
if (heads.some((h) => h.damaged)) return UNSAFE_REPLACEMENT;
|
|
if (heads.length === 2 && heads[0].record.claimId !== heads[1].record.claimId) return UNSAFE_REPLACEMENT;
|
|
const state = heads.map(headState).reduce(conservative, "stopped");
|
|
return state === "reserved" || state === "active" ? ALREADY_ACTIVE : UNSAFE_REPLACEMENT;
|
|
}
|
|
|
|
refuse(seat, session) {
|
|
const code = this.refusal(seat, session);
|
|
return new ControlRefusal(code, `the writer claim for this seat or session is held (${code})`);
|
|
}
|
|
|
|
// Reserves both keys under a new claim ID. `fields` holds the binding
|
|
// fields the record stores. Refuses when either key is held.
|
|
async acquire(seatK, sessionK, fields) {
|
|
for (let attempt = 0; attempt < 4; attempt++) {
|
|
const seat = this.head(seatK), session = this.head(sessionK);
|
|
if (held(seat) || held(session)) throw this.refuse(seat, session);
|
|
const claimId = fields.claimId ?? newClaimId();
|
|
const body = {
|
|
version: CLAIM_VERSION, kind: "writer-claim", claimId,
|
|
bindingId: fields.bindingId, harness: fields.harness, conversation: fields.conversation,
|
|
branch: fields.branch, leaf: fields.leaf, pins: fields.pins,
|
|
host: { machineId: this.host.machineId(), bootId: this.host.bootId() },
|
|
owner: fields.owner, unitName: unitNameFor(claimId), spawnMarker: false,
|
|
invocationId: null, engine: null, shim: null, generation: fields.generation,
|
|
state: "reserved", proof: null, stop: null, prior: fields.prior ?? null,
|
|
execution: fields.execution ?? null, cohortRef: fields.cohortRef ?? null,
|
|
};
|
|
const first = await this.publish(seatK, seat.n + 1, body);
|
|
if (!first) continue; // lost the seat key: re-read and re-decide
|
|
await this.pause("between-keys", { claimId });
|
|
const now = this.head(sessionK);
|
|
const second = held(now) ? null : await this.publish(sessionK, now.n + 1, body);
|
|
if (!second) {
|
|
// The session key is held by someone else: close our seat revision
|
|
// with a no-unit proof. Nothing was spawned.
|
|
await this.publish(seatK, seat.n + 2, { ...body, state: "stopped", proof: { kind: "no-unit", ref: null } });
|
|
throw this.refuse(this.head(seatK), this.head(sessionK));
|
|
}
|
|
return { claimId, seatKey: seatK, sessionKey: sessionK, n: { seat: first.revision, session: second.revision }, record: body };
|
|
}
|
|
throw new ControlRefusal(UNSAFE_REPLACEMENT, "the writer claim changed during every attempt");
|
|
}
|
|
|
|
// Publishes `patch` on both keys, seat first. Both heads must still be this
|
|
// claim's own last revisions.
|
|
async advance(claim, patch) {
|
|
const next = { ...claim.record, ...patch };
|
|
const n = { ...claim.n };
|
|
for (const kind of ["seat", "session"]) {
|
|
const key = kind === "seat" ? claim.seatKey : claim.sessionKey;
|
|
const head = this.head(key);
|
|
if (head.damaged || head.n !== n[kind] || head.record.claimId !== claim.claimId) {
|
|
throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key changed under claim ${claim.claimId}`);
|
|
}
|
|
const rec = await this.publish(key, head.n + 1, next);
|
|
if (!rec) throw new ControlRefusal(UNSAFE_REPLACEMENT, `another writer published on the ${kind} key of claim ${claim.claimId}`);
|
|
n[kind] = rec.revision;
|
|
if (kind === "seat") await this.pause("between-keys", { claimId: claim.claimId, patch });
|
|
}
|
|
claim.record = next;
|
|
claim.n = n;
|
|
return claim;
|
|
}
|
|
|
|
// Restart classification for a pair whose recorded owner may be gone. Acts
|
|
// only when the owner is proven gone: a different boot, or no process with
|
|
// the recorded pid and start time. `units.lookup(record)` reports the
|
|
// intended unit as { state: "absent" | "alive" | "unknown" }. `bootProof`
|
|
// issues the boot proof: it returns { proof, effects } once the verifier
|
|
// accepted both, or null, which leaves the pair `uncertain`. `adopt` is the
|
|
// owner record a recovery controller publishes when it takes over an
|
|
// orphan; without it nothing is adopted. `stoppedPatch` adds fields (the
|
|
// leaf at proof) to a `stopped` revision this classification publishes.
|
|
async classify(seatK, sessionK, { units, bootProof = null, adopt = null, stoppedPatch = {} } = {}) {
|
|
const seat = this.head(seatK), session = this.head(sessionK);
|
|
if (!held(seat) && !held(session)) return { state: "free" };
|
|
const code = this.refusal(seat, session);
|
|
if (code === FOREIGN_HOST) return { state: "uncertain", refusal: FOREIGN_HOST };
|
|
if (seat.damaged || session.damaged) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, damaged: true };
|
|
const heads = [seat, session].filter(held);
|
|
if (heads.length === 2 && heads[0].record.claimId !== heads[1].record.claimId) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT };
|
|
const claimId = heads[0].record.claimId;
|
|
const ours = [seat, session].filter((h) => h.n > 0 && !h.damaged && h.record.claimId === claimId);
|
|
const latest = ours.reduce((a, b) => (a.record.writtenAt >= b.record.writtenAt ? a : b)).record;
|
|
const owner = latest.owner;
|
|
if (adopt && owner?.incarnation === adopt.incarnation && owner?.pid === adopt.pid) return { state: "owned", claimId, record: latest };
|
|
const status = ownerState(owner ? { pid: owner.pid, start: owner.start ?? null, boot: owner.boot ?? null } : { invalid: true }, { identity: this.identity });
|
|
if (status === "live" || status === "unknown" || status === "invalid") return { state: headState(heads[0]), refusal: ALREADY_ACTIVE, claimId, record: latest };
|
|
|
|
const pair = heads.map(headState).reduce(conservative, "stopped");
|
|
const marker = ours.some((h) => h.record.spawnMarker);
|
|
const claim = this.claimFrom(seatK, sessionK, claimId, seat, session, latest);
|
|
|
|
// A different boot on the same host: every process of that boot is gone.
|
|
if (latest.host.bootId !== this.host.bootId()) {
|
|
if (!bootProof) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, claimId, record: latest, needs: "boot-proof" };
|
|
const issued = await bootProof(latest);
|
|
if (!issued) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, claimId, record: latest, needs: "boot-proof" };
|
|
const { proof, effects = null } = issued;
|
|
await this.finish(claim, { ...stoppedPatch, state: "stopped", proof: { kind: "boot", ref: proof.id, effects: effects?.id ?? null }, spawnMarker: marker });
|
|
return { state: "stopped", proofKind: "boot", proof, effects, claimId, record: claim.record };
|
|
}
|
|
|
|
// Keys disagree because a release stopped between them: finish it under
|
|
// the same claim ID and proof.
|
|
const stoppedHead = ours.find((h) => h.record.state === "stopped" && h.record.proof);
|
|
if (stoppedHead) {
|
|
await this.finish(claim, { state: "stopped", proof: stoppedHead.record.proof, spawnMarker: marker || latest.spawnMarker, leafAtProof: stoppedHead.record.leafAtProof ?? null, branchAtProof: stoppedHead.record.branchAtProof ?? null });
|
|
return { state: "stopped", proofKind: stoppedHead.record.proof.kind, claimId, record: claim.record, completed: true };
|
|
}
|
|
|
|
const unit = await units.lookup(latest);
|
|
if (pair === "reserved" && !marker && unit.state === "absent") {
|
|
await this.finish(claim, { ...stoppedPatch, state: "stopped", proof: { kind: "no-unit", ref: null } });
|
|
return { state: "stopped", proofKind: "no-unit", claimId, record: claim.record };
|
|
}
|
|
|
|
// A marker with no unit, a unit that exists, or an engine that may run:
|
|
// the pair is uncertain. A collected scope is an absent observation, not
|
|
// proof that every process ended. The marker is copied, never dropped.
|
|
const patch = { state: latest.state === "stopping" && latest.stop ? "stopping" : "uncertain", spawnMarker: marker };
|
|
if (adopt) patch.owner = adopt;
|
|
await this.finish(claim, patch);
|
|
return { state: patch.state, claimId, record: claim.record, unit, adopted: Boolean(adopt), claim, resumeStop: latest.stop && latest.state === "stopping" ? latest.stop : null };
|
|
}
|
|
|
|
claimFrom(seatK, sessionK, claimId, seat, session, latest) {
|
|
return {
|
|
claimId, seatKey: seatK, sessionKey: sessionK,
|
|
n: { seat: seat.n, session: session.n },
|
|
heads: { seat, session },
|
|
record: { ...latest },
|
|
};
|
|
}
|
|
|
|
// Brings both keys of a claim to `patch`, seat first. A key whose head
|
|
// belongs to an older claim (a half-done acquire) gains the revision too.
|
|
async finish(claim, patch) {
|
|
const next = { ...claim.record, ...patch };
|
|
for (const kind of ["seat", "session"]) {
|
|
const key = kind === "seat" ? claim.seatKey : claim.sessionKey;
|
|
const head = this.head(key);
|
|
if (head.damaged) throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key is unreadable`);
|
|
if (head.n > 0 && head.record.claimId !== claim.claimId && held(head)) throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key belongs to another claim`);
|
|
const same = head.n > 0 && head.record.claimId === claim.claimId && head.record.state === next.state &&
|
|
JSON.stringify(head.record.proof) === JSON.stringify(next.proof) && head.record.spawnMarker === next.spawnMarker &&
|
|
JSON.stringify(head.record.owner) === JSON.stringify(next.owner);
|
|
if (same) {
|
|
claim.n[kind] = head.n;
|
|
continue;
|
|
}
|
|
const rec = await this.publish(key, head.n + 1, next);
|
|
if (!rec) throw new ControlRefusal(UNSAFE_REPLACEMENT, `another writer published on the ${kind} key of claim ${claim.claimId}`);
|
|
claim.n[kind] = rec.revision;
|
|
if (kind === "seat") await this.pause("between-keys", { claimId: claim.claimId, patch });
|
|
}
|
|
claim.record = next;
|
|
return claim;
|
|
}
|
|
}
|