Files
stack/packages/conversation/tests/claim.test.mjs
T
jason.woltjeandClaude Opus 5.5 243e153c8b feat(conversation): CHAT-03 I1, mediated control of a sealed headless Pi (#1507)
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]>
2026-10-04 15:47:53 -05:00

708 lines
34 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// CHAT-03 §2 writer claim and live-session guard (#1507): W1–W9, W11–W17,
// W20, G1–G3. Fixtures that need a dead or crashed owner run the controller
// in a child process (ctrl-child.mjs); the owner check sees this test's own
// pid as live.
import { test, after } from "node:test";
import assert from "node:assert/strict";
import { appendFileSync, copyFileSync, linkSync, mkdirSync, statSync, mkdtempSync, readFileSync, readdirSync, rmSync, symlinkSync, unlinkSync, writeFileSync, lstatSync } from "node:fs";
import { homedir, tmpdir } from "node:os";
import { join } from "node:path";
import { spawnSync } from "node:child_process";
import { createHash } from "node:crypto";
import { ALREADY_ACTIVE, ClaimStore, FOREIGN_HOST, UNSAFE_REPLACEMENT, held, machineId } from "../src/claim.mjs";
import { Controller } from "../src/controller.mjs";
import { ConversationClient } from "../src/client.mjs";
import { AUTHORITY, scopeAvailable } from "../src/cohort.mjs";
import { FixtureVerifier } from "../src/records.mjs";
import { LIVE_SESSION_REFUSED } from "../src/guard.mjs";
import { ENGINE_PIN_MISMATCH, PI_PACKAGE } from "../src/pi-pin.mjs";
import { bootId } from "../../discord/src/journal.mjs";
import { FakeLauncher } from "./fake-pi.mjs";
import {
REPO, FAST, fixture, controllerFor, started, spawnController, killChildren, cleanupAll, reap, claimRecords,
userEntry, assistantEntry, thinkingEntry, receiptState, noUnits, tick,
} from "./harness.mjs";
const reaped = [];
after(() => {
for (const fx of reaped) reap(fx);
killChildren();
cleanupAll();
});
const track = (fx) => (reaped.push(fx), fx);
const SCOPE = scopeAvailable();
const OTHER_BOOT = "00000000-0000-4000-8000-000000000000";
// A pid that existed and has exited: an owner proven gone.
function deadPid() {
return spawnSync(process.execPath, ["-e", "process.stdout.write(String(process.pid))"], { encoding: "utf8" }).stdout.trim() * 1;
}
const deadOwner = () => ({ pid: deadPid(), start: "1", boot: bootId(), incarnation: "dead".padEnd(32, "0") });
const pins = { engineVersion: "0.85.1", enginePin: "x", argvDigest: "y" };
const fields = (extra = {}) => ({ bindingId: "binding-1", harness: "pi", conversation: "pi-c1", branch: "main", leaf: "b2c3d4e5", pins, owner: deadOwner(), generation: 1, ...extra });
// The keys a controller for `fx` uses, without starting it.
function keysFor(fx, opts = {}) {
const { ctrl } = controllerFor(fx, opts);
return { store: ctrl.store, seatK: ctrl.seatK, sessionK: ctrl.sessionK, ctrl };
}
function revisionTexts(store, key) {
return store.revisions(key).map((r) => r.text);
}
function treeDigest(dir) {
const out = {};
const walk = (p) => {
const st = lstatSync(p);
if (st.isDirectory()) for (const n of readdirSync(p).sort()) walk(join(p, n));
else out[p] = createHash("sha256").update(readFileSync(p)).digest("hex");
};
walk(dir);
return out;
}
async function connectClient(socketPath) {
const c = new ConversationClient({ socketPath });
await c.connect();
return c;
}
// Force stop through the wire, with a confirmation.
async function forceStop(c) {
const r = await c.confirmed("force-stop");
assert.equal(r.outcome, "force-stop-fenced", JSON.stringify(r));
return r.stop;
}
async function waitBinding(c, states, ms = 8000) {
const want = new Set([states].flat());
const ok = await c.waitFor(() => want.has(c.binding?.state), ms);
assert.ok(ok, `binding stayed ${c.binding?.state}; wanted ${[...want]}`);
}
// ---- W1–W4: acquisition --------------------------------------------------
test("W1: two processes acquire the same pair at once; exactly one claim", async () => {
const fx = track(fixture());
const a = spawnController({ fx, launcher: "pgroup", units: "absent" });
const b = spawnController({ fx, launcher: "pgroup", units: "absent" });
const [ra, rb] = await Promise.all([a.next((m) => m.ready || m.error), b.next((m) => m.ready || m.error)]);
const results = [ra, rb];
assert.equal(results.filter((r) => r.ready).length, 1, JSON.stringify(results));
assert.equal(results.find((r) => r.error).error, ALREADY_ACTIVE);
const ids = new Set(claimRecords(fx).map((r) => r.record.claimId));
assert.equal(ids.size, 1);
for (const ch of [a, b]) ch.send("kill-engine");
await Promise.all([a.exited, b.exited]);
});
test("W1: two writers publish the same revision at once: one wins, the other gets null, the winner's record stays", async () => {
// Both hold at temp-synced, so both reach link() with the name free when
// they read the head. A rename would let the second replace the first.
const fx = track(fixture());
const { seatK } = keysFor(fx);
let arrived = 0, go;
const ready = new Promise((r) => (go = r));
const barrier = async (name) => {
if (name !== "temp-synced") return;
if (++arrived === 2) go();
await ready;
};
const stores = [0, 1].map(() => new ClaimStore({ root: fx.claimRoot, barrier }));
const results = await Promise.all(stores.map((s, i) => s.publish(seatK, 1, { claimId: `claim-${i}`, state: "reserved" })));
assert.equal(results.filter((r) => r === null).length, 1, JSON.stringify(results));
const winner = results.find((r) => r !== null);
const texts = revisionTexts(stores[0], seatK);
assert.equal(texts.length, 1);
assert.equal(JSON.parse(texts[0]).claimId, winner.claimId, "the published revision is the winner's, never replaced");
assert.deepEqual(readdirSync(stores[0].dir(seatK)).filter((n) => n.startsWith(".tmp-")), [], "no temp file left behind");
});
test("W1: a revision name appears only after its bytes are synced; before that, only a temp file exists", async () => {
const fx = track(fixture());
const { seatK } = keysFor(fx);
let release = null;
const barrier = (name) => (name === "temp-written" ? new Promise((r) => (release = r)) : undefined);
const store = new ClaimStore({ root: fx.claimRoot, barrier });
const p = store.publish(seatK, 1, { claimId: "claim-a", state: "reserved" });
for (let i = 0; !release && i < 200; i++) await tick(10);
assert.ok(release, "publish reached temp-written");
assert.deepEqual(revisionTexts(store, seatK), [], "no revision is visible while its bytes are unsynced");
assert.equal(readdirSync(store.dir(seatK)).filter((n) => n.startsWith(".tmp-")).length, 1);
release();
assert.equal((await p).claimId, "claim-a");
assert.equal(revisionTexts(store, seatK).length, 1);
});
test("W2: acquire while a claim is reserved or active refuses already-active", async () => {
const fx = fixture();
const { store, seatK, sessionK } = keysFor(fx);
const claim = await store.acquire(seatK, sessionK, fields());
await assert.rejects(store.acquire(seatK, sessionK, fields()), { code: ALREADY_ACTIVE });
await store.advance(claim, { state: "active" });
await assert.rejects(store.acquire(seatK, sessionK, fields()), { code: ALREADY_ACTIVE });
// A live controller holds it: a second controller refuses before it
// touches the first one's socket.
const fx2 = fixture();
const h = await started({ fx: fx2 });
try {
const sock = lstatSync(h.ctrl.socketPath);
const { ctrl } = controllerFor(fx2);
await assert.rejects(ctrl.start(), { code: ALREADY_ACTIVE });
assert.equal(lstatSync(h.ctrl.socketPath).ino, sock.ino);
assert.equal(h.launcher.engines.length, 1);
} finally {
await h.close();
}
});
test("W3: acquire while stopping, uncertain, or stopped without proof refuses unsafe-replacement", async () => {
for (const patch of [{ state: "stopping" }, { state: "uncertain" }, { state: "stopped", proof: null }]) {
const fx = fixture();
const { store, seatK, sessionK } = keysFor(fx);
const claim = await store.acquire(seatK, sessionK, fields());
await store.advance(claim, patch);
await assert.rejects(store.acquire(seatK, sessionK, fields()), { code: UNSAFE_REPLACEMENT }, JSON.stringify(patch));
}
});
test("W4: same session with another seat tuple, and the reverse, both refuse; a loser on the seat key closes it no-unit", async () => {
const fx = fixture();
const { store, seatK, sessionK } = keysFor(fx);
await store.acquire(seatK, sessionK, fields());
const otherSeat = { ...seatK, seat: "other-seat" };
const otherSession = { ...sessionK, session: { harness: "pi", nativeSession: "other" } };
await assert.rejects(store.acquire(otherSeat, sessionK, fields()), { code: ALREADY_ACTIVE });
await assert.rejects(store.acquire(seatK, otherSession, fields()), { code: ALREADY_ACTIVE });
assert.equal(store.head(otherSeat).n, 0, "a contender that read a held key publishes nothing");
// The race: the loser publishes on its seat key, then finds the session
// key taken.
const fx2 = fixture();
const k = keysFor(fx2);
let winner = null;
const loserStore = new ClaimStore({
root: fx2.claimRoot,
barrier: async (name) => {
if (name === "between-keys" && !winner) winner = await k.store.acquire({ ...k.seatK, seat: "winner-seat" }, k.sessionK, fields());
},
});
const winnerBefore = () => [revisionTexts(k.store, { ...k.seatK, seat: "winner-seat" }), revisionTexts(k.store, k.sessionK)];
await assert.rejects(loserStore.acquire(k.seatK, k.sessionK, fields()), { code: ALREADY_ACTIVE });
const head = k.store.head(k.seatK);
assert.equal(head.n, 2);
assert.equal(head.record.state, "stopped");
assert.deepEqual(head.record.proof, { kind: "no-unit", ref: null });
assert.equal(head.record.spawnMarker, false);
assert.equal(held(head), false);
const [ws, wss] = winnerBefore();
assert.equal(ws.length, 1);
assert.equal(wss.length, 1);
assert.equal(JSON.parse(wss[0]).claimId, winner.claimId);
});
// ---- W5: SIGKILL at every publication barrier ------------------------------
// The claim invariants after any crash: every visible revision parses, no two
// claim IDs hold the pair, and a pair with any non-stopped revision is not
// free.
function assertClaimInvariants(fx, label) {
const recs = claimRecords(fx);
for (const r of recs) assert.ok(r.record && r.record.kind === "writer-claim", `${label}: incomplete revision ${r.path}`);
const k = keysFor(fx);
const heads = [k.store.head(k.seatK), k.store.head(k.sessionK)];
const holders = new Set(heads.filter(held).map((h) => h.record.claimId));
assert.ok(holders.size <= 1, `${label}: two holders`);
return { k, heads };
}
async function lifecycleTrace(fx, launcher) {
const ch = spawnController({ fx, launcher, units: launcher === "scope" ? undefined : "absent", trace: true });
const ready = await ch.next((m) => m.ready || m.error, 15000);
assert.ok(ready.ready, JSON.stringify(ready));
return { ch, ready };
}
// W4 through the controller's own key mapping: the session key is the Pi
// header ID, so a second path to the same session collides with the first.
for (const [how, place] of [["hard link", linkSync], ["copy", copyFileSync]]) {
test(`W4: a ${how} of one session under another seat is the same session: the second controller refuses already-active and launches nothing`, async () => {
const fx = track(fixture());
const sessionsB = join(fx.proj, ".pi", "state", "seat-b", "sessions");
mkdirSync(sessionsB, { recursive: true });
const fileB = join(sessionsB, "s1.jsonl");
place(fx.sessionFile, fileB);
if (how === "hard link") assert.equal(statSync(fileB).ino, statSync(fx.sessionFile).ino);
const fxB = { ...fx, seat: "seat-b", sessionFile: fileB, socketDir: join(fx.base, "sock-b") };
const h = await started({ fx });
try {
const { ctrl: b, launcher } = controllerFor(fxB);
assert.notEqual(b.conversation, h.ctrl.conversation, "two conversation IDs");
assert.deepEqual(b.sessionK, h.ctrl.sessionK, "one session key");
await assert.rejects(b.start(), { code: ALREADY_ACTIVE });
assert.equal(launcher.launches.length, 0, "no second engine");
assert.equal(h.ctrl.binding.state, "active");
await b.close();
} finally {
await h.close();
}
});
}
test("W4: a session header ID that changes after construction refuses target; nothing is claimed or launched", async () => {
const fx = track(fixture());
const { ctrl, launcher } = controllerFor(fx);
writeFileSync(fx.sessionFile, readFileSync(fx.sessionFile, "utf8").replace(/"id":"[^"]+"/, '"id":"11111111-2222-4333-8444-555555555555"'));
await assert.rejects(ctrl.start(), { code: "target" });
assert.equal(launcher.launches.length, 0);
assert.equal(ctrl.store.head(ctrl.seatK).n, 0, "the seat key is untouched");
});
test("W5: SIGKILL between every publication barrier of acquire and transition; restart never finds two holders or a lost claim", { timeout: 240000 }, async () => {
const probe = track(fixture());
const { ch } = await lifecycleTrace(probe, "pgroup");
ch.send("kill-engine");
await ch.exited;
const points = ch.msgs.filter((m) => m.barrier).map((m) => [m.barrier, m.n]);
assert.ok(points.length >= 20, `barrier trace too short: ${JSON.stringify(points)}`);
for (const [name, n] of points) {
const fx = track(fixture());
const dying = spawnController({ fx, launcher: "pgroup", units: "absent", dieAt: { [name]: n } });
await dying.next((m) => m.dying || m.ready || m.error, 15000);
await dying.exited;
reap(fx);
const label = `${name}#${n}`;
const { heads } = assertClaimInvariants(fx, label);
const anyHeld = heads.some(held);
const restart = spawnController({ fx, launcher: "pgroup", units: "absent" });
const r = await restart.next((m) => m.ready || m.error, 15000);
assert.ok(r.ready, `${label}: restart refused ${JSON.stringify(r)}`);
// A held pair is completed or classified under its claim ID; never
// launched over.
if (anyHeld) assert.equal(r.started.launched, false, `${label}: launched over a held pair`);
const after = assertClaimInvariants(fx, `${label} after restart`);
if (anyHeld) {
const before = heads.find(held).record.claimId;
for (const h of after.heads) if (held(h)) assert.equal(h.record.claimId, before, `${label}: claim lost`);
}
restart.send("kill-engine");
await restart.exited;
reap(fx);
}
});
test("W5: SIGKILL between every publication barrier of release; restart finishes or holds the release", { skip: !SCOPE && "systemd user scopes unavailable", timeout: 300000 }, async () => {
// One traced force stop enumerates the release barriers: every barrier
// after the force stop is sent.
const runRelease = async (fx, dieAt = null) => {
const ch = spawnController({ fx, launcher: "scope", trace: !dieAt, dieAt: dieAt ?? undefined });
const ready = await ch.next((m) => m.ready || m.error, 15000);
assert.ok(ready.ready, JSON.stringify(ready));
const c = await connectClient(ready.socketPath);
await c.takeover();
await c.waitFor(() => false, 200); // the generation publication settles
const mark = ch.msgs.length;
const r = await c.confirmed("force-stop");
if (!dieAt) {
assert.equal(r.outcome, "force-stop-fenced", JSON.stringify(r));
await waitBinding(c, "stopped", 15000);
} else await ch.next((m) => m.dying, 20000);
c.close();
return { ch, mark };
};
const probe = track(fixture());
const { ch, mark } = await runRelease(probe);
ch.send("close");
await ch.exited;
const points = ch.msgs.slice(mark).filter((m) => m.barrier).map((m) => [m.barrier, m.n]);
assert.ok(points.some(([b]) => b === "phase-kill"));
for (const [name, n] of points) {
const fx = track(fixture());
const { ch: dying } = await runRelease(fx, { [name]: n });
await dying.exited;
const label = `release ${name}#${n}`;
const { heads } = assertClaimInvariants(fx, label);
const claimId = heads.find((h) => h.n > 0).record.claimId;
const anyHeld = heads.some(held);
const restart = spawnController({ fx, launcher: "scope" });
const r = await restart.next((m) => m.ready || m.error, 15000);
assert.ok(r.ready, `${label}: restart refused ${JSON.stringify(r)}`);
if (!anyHeld) {
// The release had published on both keys: the pair is free.
assert.equal(r.started.launched, true, label);
restart.send("kill-engine");
await restart.exited;
reap(fx);
continue;
}
assert.equal(r.started.launched, false, `${label}: launched over a release in progress`);
const after = assertClaimInvariants(fx, `${label} after restart`);
for (const h of after.heads) assert.equal(h.record.claimId, claimId, `${label}: claim lost`);
// Either both keys are stopped with one proof, or the pair is held.
const states = after.heads.map((h) => h.record.state);
if (states.includes("stopped") && after.heads.every((h) => !held(h))) {
assert.deepEqual(after.heads[0].record.proof, after.heads[1].record.proof, label);
} else assert.ok(after.heads.some(held), `${label}: ${states}`);
restart.send("close");
await restart.exited;
reap(fx);
}
});
// ---- W6, W12–W15, W20: crash and restart -----------------------------------
test("W6: controller killed mid-turn while the engine lives; restart is uncertain, no launch, prompts refuse", async () => {
const fx = track(fixture());
const log = join(fx.base, "fake.log");
const fakeEnv = { FAKE_PI_CONTROL: join(fx.base, "fake.sock"), FAKE_PI_LOG: log, FAKE_PI_SCRIPT: JSON.stringify([[{ pause: "mid" }, { text: "late" }]]) };
const a = spawnController({ fx, launcher: "pgroup", units: "absent", fakeEnv });
const ra = await a.next((m) => m.ready);
const c = await connectClient(ra.socketPath);
await c.takeover();
const r = await c.prompt("long turn");
await receiptState(c, r.receipt.id, "working");
a.proc.kill("SIGKILL");
await a.exited;
c.close();
const engineCount = () => readFileSync(log, "utf8").split("\n").filter((l) => l.includes('"t":"argv"')).length;
const prompts = () => readFileSync(log, "utf8").split("\n").filter((l) => l.includes('"type":"prompt"')).length;
const b = spawnController({ fx, launcher: "pgroup", units: "absent", fakeEnv });
const rb = await b.next((m) => m.ready || m.error);
assert.ok(rb.ready, JSON.stringify(rb));
assert.equal(rb.started.launched, false);
assert.equal(rb.started.classified.state, "uncertain");
assert.equal(engineCount(), 1, "no second engine");
const c2 = await connectClient(rb.socketPath);
assert.equal(c2.binding.state, "uncertain");
assert.equal((await c2.takeover()).refusal, "fenced");
const rc = await c2.confirmed("acquire-recovery-control");
assert.equal(rc.outcome, "recovery-control-acquired");
const p = await c2.prompt("after restart");
assert.ok(p.outcome.startsWith("refused:"), JSON.stringify(p));
assert.equal(prompts(), 1);
c2.close();
b.send("close");
await b.exited;
});
test("W12: a live owner paused with SIGSTOP; a second controller refuses already-active and changes nothing", async () => {
const fx = track(fixture());
const a = spawnController({ fx, launcher: "pgroup", units: "absent" });
await a.next((m) => m.ready);
const k = keysFor(fx);
const before = [revisionTexts(k.store, k.seatK), revisionTexts(k.store, k.sessionK)];
a.proc.kill("SIGSTOP");
try {
await assert.rejects(k.ctrl.start(), { code: ALREADY_ACTIVE });
} finally {
a.proc.kill("SIGCONT");
}
assert.deepEqual([revisionTexts(k.store, k.seatK), revisionTexts(k.store, k.sessionK)], before);
a.send("kill-engine");
await a.exited;
});
test("W13: crash after the engine spawns, before active; restart finds the live unit: uncertain, no second spawn, force stop only", { skip: !SCOPE && "systemd user scopes unavailable", timeout: 60000 }, async () => {
const fx = track(fixture());
const log = join(fx.base, "fake.log");
const fakeEnv = { FAKE_PI_CONTROL: join(fx.base, "fake.sock"), FAKE_PI_LOG: log };
const a = spawnController({ fx, launcher: "scope", dieAt: { spawned: 1 }, fakeEnv });
await a.next((m) => m.dying, 15000);
await a.exited;
const b = spawnController({ fx, launcher: "scope", fakeEnv });
const rb = await b.next((m) => m.ready || m.error, 15000);
assert.ok(rb.ready, JSON.stringify(rb));
assert.equal(rb.started.launched, false);
assert.equal(rb.started.classified.state, "uncertain");
assert.equal(rb.started.classified.unit, "alive");
const engines = readFileSync(log, "utf8").split("\n").filter((l) => l.includes('"t":"argv"')).length;
assert.equal(engines, 1, "no second spawn");
const c = await connectClient(rb.socketPath);
assert.equal((await c.prompt("x")).refusal, "controller");
assert.equal((await c.confirmed("acquire-recovery-control")).outcome, "recovery-control-acquired");
await forceStop(c);
await waitBinding(c, "stopped", 15000);
const k = keysFor(fx);
for (const key of [k.seatK, k.sessionK]) {
const h = k.store.head(key);
assert.equal(h.record.state, "stopped");
assert.equal(h.record.proof.kind, "cohortProof");
}
c.close();
b.send("close");
await b.exited;
});
test("W14: crash after reservation, before the spawn marker: stopped with a no-unit observation; the pair is free", async () => {
const fx = track(fixture());
const a = spawnController({ fx, launcher: "pgroup", units: "absent", dieAt: { reserved: 1 } });
await a.next((m) => m.dying);
await a.exited;
const b = spawnController({ fx, launcher: "pgroup", units: "absent" });
const rb = await b.next((m) => m.ready || m.error);
assert.deepEqual(rb.started, { launched: false, classified: { state: "stopped", proofKind: "no-unit" } });
b.send("close");
await b.exited;
const k = keysFor(fx);
for (const key of [k.seatK, k.sessionK]) {
const h = k.store.head(key);
assert.equal(h.record.state, "stopped");
assert.deepEqual(h.record.proof, { kind: "no-unit", ref: null });
assert.equal(held(h), false);
}
const c = spawnController({ fx, launcher: "pgroup", units: "absent" });
const rc = await c.next((m) => m.ready || m.error);
assert.equal(rc.started.launched, true);
c.send("kill-engine");
await c.exited;
});
test("W20: crash after the spawn marker, scope collected; uncertain in both runs, the marker is copied, no launch until a boot proof", async () => {
// Run 1: marker on both keys. Run 2: marker on the seat key only (the
// second between-keys is the spawn-marker advance's).
for (const dieAt of [{ "spawn-marker": 1 }, { "between-keys": 2 }]) {
const fx = track(fixture());
const a = spawnController({ fx, launcher: "pgroup", units: "absent", dieAt });
await a.next((m) => m.dying);
await a.exited;
const k = keysFor(fx);
const seatMarker = k.store.head(k.seatK).record.spawnMarker;
const sessionMarker = k.store.head(k.sessionK).record.spawnMarker;
assert.equal(seatMarker, true);
assert.equal(sessionMarker, !dieAt["between-keys"]);
const b = spawnController({ fx, launcher: "pgroup", units: "absent" });
const rb = await b.next((m) => m.ready || m.error);
assert.ok(rb.ready, JSON.stringify(rb));
assert.equal(rb.started.launched, false);
assert.equal(rb.started.classified.state, "uncertain");
b.send("close");
await b.exited;
for (const key of [k.seatK, k.sessionK]) {
assert.equal(k.store.head(key).record.spawnMarker, true, "the marker is copied, never dropped");
assert.equal(k.store.head(key).record.state, "uncertain");
}
// Only a boot proof moves it on: the same host under another boot.
const { ctrl } = controllerFor(fx, { host: { machineId, bootId: () => OTHER_BOOT } });
const res = await ctrl.start();
assert.deepEqual(res, { launched: false, classified: { state: "stopped", proofKind: "boot" } });
await ctrl.close();
}
});
test("W15: crash between the two keys during release; restart finishes it under the same claim ID", async () => {
const fx = fixture();
const k = keysFor(fx);
const claim = await k.store.acquire(k.seatK, k.sessionK, fields());
await k.store.advance(claim, { spawnMarker: true, state: "active" });
await k.store.advance(claim, { state: "stopping" });
const crashing = new ClaimStore({ root: fx.claimRoot, barrier: async (name) => {
if (name === "between-keys") throw Object.assign(new Error("crash"), { code: "crash" });
} });
await assert.rejects(crashing.finish({ ...claim, record: { ...claim.record } }, { state: "stopped", proof: { kind: "cohortProof", ref: "cohort-proof-1", effects: "effects-1" }, leafAtProof: "b2c3d4e5", branchAtProof: "main" }), { code: "crash" });
assert.equal(k.store.head(k.seatK).record.state, "stopped");
assert.equal(k.store.head(k.sessionK).record.state, "stopping");
assert.ok(held(k.store.head(k.sessionK)));
const res = await k.ctrl.start();
assert.equal(res.launched, false);
assert.equal(res.classified.state, "stopped");
for (const key of [k.seatK, k.sessionK]) {
const h = k.store.head(key);
assert.equal(h.record.claimId, claim.claimId);
assert.equal(h.record.state, "stopped");
assert.equal(h.record.proof.ref, "cohort-proof-1");
}
await k.ctrl.close();
});
// ---- W7–W9: boot proof and resume ------------------------------------------
test("W7: recorded boot ID differs on the same machine: stopped with a boot proof; open tool calls become uncertain", async () => {
const toolCall = { type: "message", id: "c3d4e5f6", parentId: "b2c3d4e5", timestamp: new Date().toISOString(), message: { role: "assistant", content: [{ type: "toolCall", id: "tool-open-1", name: "bash", arguments: {} }], stopReason: "toolUse" } };
const fx = track(fixture({ entries: [thinkingEntry(), userEntry("a1b2c3d4", "f0e1d2c3", "hello"), assistantEntry("b2c3d4e5", "a1b2c3d4", "hi"), toolCall] }));
const a = spawnController({ fx, launcher: "pgroup", units: "absent", host: { bootId: OTHER_BOOT } });
await a.next((m) => m.ready);
a.proc.kill("SIGKILL");
await a.exited;
reap(fx);
const posted = [];
const verifier = new FixtureVerifier({ authorities: [AUTHORITY] });
const post = verifier.post.bind(verifier);
verifier.post = (p) => (posted.push(p), post(p));
const { ctrl } = controllerFor(fx, { verifier });
const res = await ctrl.start();
assert.deepEqual(res, { launched: false, classified: { state: "stopped", proofKind: "boot" } });
const proof = posted.find((p) => p.kind === "cohortProof");
const effects = posted.find((p) => p.kind === "effectReport");
assert.ok(proof.id.startsWith("boot-proof-"));
assert.deepEqual(effects.invocations.map((i) => [i.id, i.disposition]), [["tool-open-1", "uncertain"]]);
const k = keysFor(fx);
for (const key of [k.seatK, k.sessionK]) assert.deepEqual(k.store.head(key).record.proof, { kind: "boot", ref: proof.id, effects: effects.id });
await ctrl.close();
});
// A proven stop in process: the fixture launcher's cohort stop.
async function provenStop(fx, opts = {}) {
const h = await started({ fx, ...opts });
await forceStop(h.client);
await waitBinding(h.client, "stopped");
const claim = h.ctrl.claim.record;
await h.close();
return claim;
}
test("W8: resume after a proven stop with the same pins: new claim ID, generation +1, same conversation, branch and leaf", async () => {
const fx = fixture();
const first = await provenStop(fx);
const h = await started({ fx, control: false });
try {
const next = h.ctrl.claim.record;
assert.notEqual(next.claimId, first.claimId);
assert.equal(next.generation, first.generation + 1);
assert.equal(next.conversation, first.conversation);
assert.equal(next.branch, first.branch);
assert.equal(next.leaf, first.leafAtProof);
assert.deepEqual(next.prior, { claimId: first.claimId, stop: first.stop.id });
} finally {
await h.close();
}
});
test("W9: resume with a changed binary, argv digest, branch or leaf is refused and the claim is unchanged", async () => {
const cases = {
binary: (fx) => {
const root = mkdtempSync(join(fx.base, "pin-"));
const lock = { packages: { [`node_modules/${PI_PACKAGE}`]: { version: "0.85.2", integrity: "sha512-other" } } };
mkdirSync(join(root, "node_modules"));
writeFileSync(join(root, "package-lock.json"), JSON.stringify(lock));
writeFileSync(join(root, "node_modules", ".package-lock.json"), JSON.stringify(lock));
return [{ pinRoot: root }, ENGINE_PIN_MISMATCH];
},
argv: () => [{ engine: { extraArgs: ["--model", "other"] } }, ENGINE_PIN_MISMATCH],
leaf: (fx) => {
appendFileSync(fx.sessionFile, JSON.stringify(userEntry("d4e5f6a7", "b2c3d4e5", "outside", 9)) + "\n");
return [{}, "target"];
},
branch: (fx) => {
appendFileSync(fx.sessionFile, JSON.stringify(userEntry("e5f6a7b8", "a1b2c3d4", "a fork", 9)) + "\n");
return [{}, "target"];
},
};
for (const [name, change] of Object.entries(cases)) {
const fx = fixture();
await provenStop(fx);
const k = keysFor(fx);
const before = [revisionTexts(k.store, k.seatK), revisionTexts(k.store, k.sessionK)];
const [opts, code] = change(fx);
const { ctrl, launcher } = controllerFor(fx, opts);
await assert.rejects(ctrl.start(), { code }, name);
assert.equal(launcher.engines.length, 0, name);
assert.deepEqual([revisionTexts(k.store, k.seatK), revisionTexts(k.store, k.sessionK)], before, name);
}
});
// ---- W11, W16, W17 ----------------------------------------------------------
test("W11: the controller writes no session file; only the fake engine's own appends appear", async () => {
const fx = fixture();
const original = readFileSync(fx.sessionFile, "utf8");
const listing = readdirSync(fx.sessions).sort();
const h = await started({ fx });
try {
const r = await h.client.prompt("one");
await receiptState(h.client, r.receipt.id, "finished");
await h.client.interrupt();
await forceStop(h.client);
await waitBinding(h.client, "stopped");
} finally {
await h.close();
}
assert.deepEqual(readdirSync(fx.sessions).sort(), listing);
const lines = readFileSync(fx.sessionFile, "utf8").slice(original.length).split("\n").filter(Boolean).map((l) => JSON.parse(l).id);
assert.ok(readFileSync(fx.sessionFile, "utf8").startsWith(original));
assert.deepEqual(lines, h.engine.appends, "every added line is the fake engine's");
});
test("W16: a highest revision that won't parse holds the pair uncertain; the older stopped revision is not reused", async () => {
const fx = fixture();
const k = keysFor(fx);
const claim = await k.store.acquire(k.seatK, k.sessionK, fields());
await k.store.finish(claim, { state: "stopped", proof: { kind: "no-unit", ref: null } });
const dir = k.store.dir(k.sessionK);
const top = readdirSync(dir).filter((n) => n.startsWith("r")).sort().at(-1);
const next = `r${String(Number(top.slice(1, 11)) + 1).padStart(10, "0")}.json`;
writeFileSync(join(dir, next), "{\"version\":1,\"kind\":\"writer-cl");
assert.equal(k.store.head(k.sessionK).damaged, true);
await assert.rejects(k.store.acquire(k.seatK, k.sessionK, fields()), { code: UNSAFE_REPLACEMENT });
await assert.rejects(k.ctrl.start(), { code: UNSAFE_REPLACEMENT });
assert.equal(readdirSync(dir).filter((n) => n.startsWith("r")).length, 3);
});
test("W17: a claim root copied from another host refuses foreign-host and promotes nothing", async () => {
const fx = fixture();
const k = keysFor(fx);
const foreign = new ClaimStore({ root: fx.claimRoot, host: { machineId: () => "f".repeat(32), bootId } });
await foreign.acquire(k.seatK, k.sessionK, fields());
const before = treeDigest(fx.claimRoot);
await assert.rejects(k.store.acquire(k.seatK, k.sessionK, fields()), { code: FOREIGN_HOST });
await assert.rejects(k.ctrl.start(), { code: FOREIGN_HOST });
assert.deepEqual(treeDigest(fx.claimRoot), before);
});
// ---- G1–G3: live-session guard ---------------------------------------------
test("G1: a session path or claim root under .pi/state, ~/.claude, the data root or a registration refuses at construction", () => {
const fx = fixture();
const base = { claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat, launcher: new FakeLauncher(), units: noUnits };
const repoSession = join(REPO, ".pi", "state", "x-seat", "sessions", "s.jsonl");
const claudeRoot = join(homedir(), ".claude", "mosaic-claims");
const dataRoot = join(fx.base, "data");
mkdirSync(dataRoot);
const regDir = join(fx.base, "registered");
const cases = [
{ fixtureRoot: REPO, sessionFile: repoSession, claimRoot: join(REPO, "tmp-claims"), socketDir: join(REPO, "tmp-sock") },
{ fixtureRoot: homedir(), claimRoot: claudeRoot, sessionFile: join(homedir(), "x", ".pi", "state", "s", "sessions", "a.jsonl"), socketDir: join(homedir(), "x-sock") },
{ fixtureRoot: fx.base, claimRoot: join(dataRoot, "claims"), guardOptions: { dataRoots: [dataRoot] } },
{ fixtureRoot: fx.base, guardOptions: { registrations: [{ sessionsDir: fx.sessions }] } },
{ fixtureRoot: fx.base, claimRoot: join(regDir, "claims"), guardOptions: { registrations: [{ workspace: regDir }] } },
];
for (const c of cases) {
assert.throws(() => new Controller({ ...base, ...c }), { code: LIVE_SESSION_REFUSED }, JSON.stringify(c));
}
assert.equal(readdirSync(fx.base).includes("claims"), false, "nothing created");
});
test("G2: a symlink inside the fixture root to a live session file is refused by the real-path check", () => {
const fx = fixture();
const live = mkdtempSync(join(tmpdir(), "chat03-live-"));
try {
const liveFile = join(live, "live.jsonl");
writeFileSync(liveFile, readFileSync(fx.sessionFile));
unlinkSync(fx.sessionFile);
symlinkSync(liveFile, fx.sessionFile);
// Construction already applies the real-path check; bind repeats it.
assert.throws(() => controllerFor(fx, { guardOptions: { registrations: [{ sessionsDir: live }] } }), { code: LIVE_SESSION_REFUSED });
} finally {
rmSync(live, { recursive: true, force: true });
}
});
test("G3: a fixture path swapped for a live path after construction is refused at bind", async () => {
const fx = fixture();
const live = mkdtempSync(join(tmpdir(), "chat03-live-"));
try {
const liveFile = join(live, "live.jsonl");
writeFileSync(liveFile, readFileSync(fx.sessionFile));
const { ctrl, launcher } = controllerFor(fx, { guardOptions: { registrations: [{ sessionsDir: live }] } });
unlinkSync(fx.sessionFile);
symlinkSync(liveFile, fx.sessionFile);
await assert.rejects(ctrl.start(), { code: LIVE_SESSION_REFUSED });
assert.equal(launcher.engines.length, 0);
assert.equal(claimRecords(fx).length, 0);
} finally {
rmSync(live, { recursive: true, force: true });
}
});