// 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; } }