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]>
708 lines
34 KiB
JavaScript
708 lines
34 KiB
JavaScript
// 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 });
|
||
}
|
||
});
|