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]>
873 lines
38 KiB
JavaScript
873 lines
38 KiB
JavaScript
// CHAT-03 §5 control races (#1507): H1–H4 and H9–H23. H5–H8 are approval
|
||
// races on Claude permission requests and run in I4. Each race lands at an
|
||
// exact step through the controller's barriers or the fake engine's pause
|
||
// points, and asserts the engine bytes, the receipts and the events.
|
||
|
||
import { test, after } from "node:test";
|
||
import assert from "node:assert/strict";
|
||
import { existsSync, readFileSync } from "node:fs";
|
||
import { join } from "node:path";
|
||
import { PassThrough, Writable } from "node:stream";
|
||
import { ALREADY_ACTIVE } from "../src/claim.mjs";
|
||
import { ConversationClient, OUTCOME_UNKNOWN } from "../src/client.mjs";
|
||
import { BUSY, NO_TURN, STALE_INCARNATION, TRANSPORT_UNKNOWN } from "../src/controller.mjs";
|
||
import { EngineLink } from "../src/engine.mjs";
|
||
import { newId } from "../src/records.mjs";
|
||
import { ControlClient, FakeLauncher } from "./fake-pi.mjs";
|
||
import { scopeAvailable } from "../src/cohort.mjs";
|
||
import { controllerFor, started, receiptState, spawnController, killChildren, cleanupAll, reap, fixture, 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();
|
||
|
||
async function until(pred, ms = 4000, what = "condition") {
|
||
const end = Date.now() + ms;
|
||
while (!pred()) {
|
||
if (Date.now() > end) throw new Error(`timed out waiting for ${what}`);
|
||
await tick(5);
|
||
}
|
||
}
|
||
|
||
// Holds the controller at named barriers. `hold(name)` holds the next
|
||
// occurrence; `release(name)` lets it go.
|
||
function gate() {
|
||
const want = new Set(), held = new Map();
|
||
return {
|
||
barrier: async (name) => {
|
||
if (!want.has(name)) return;
|
||
want.delete(name);
|
||
await new Promise((r) => held.set(name, r));
|
||
},
|
||
hold: (name) => want.add(name),
|
||
waitHeld: (name) => until(() => held.has(name), 4000, `barrier ${name}`),
|
||
release: (name) => {
|
||
const r = held.get(name);
|
||
held.delete(name);
|
||
r?.();
|
||
},
|
||
};
|
||
}
|
||
|
||
const cmds = (engine, type) => engine.commands.filter((c) => c.type === type);
|
||
const draft = () => ({ draft: newId("draft"), draftRevision: 1 });
|
||
const promptEnvelope = (c) => c.envelope("prompt", draft());
|
||
|
||
async function confirm(c, operation) {
|
||
const issued = await c.request("issue-confirmation", { operationToConfirm: operation });
|
||
assert.equal(issued.outcome, "confirmation-issued", JSON.stringify(issued));
|
||
const id = issued.data.confirmation.id;
|
||
const answered = await c.request("answer-confirmation", { confirmation: id, answer: "confirm" });
|
||
assert.equal(answered.outcome, "confirmation-confirmed", JSON.stringify(answered));
|
||
return id;
|
||
}
|
||
|
||
async function heldRun(h, c = h.client) {
|
||
h.engine.script([{ pause: "p1" }, { text: "never", stop: "stop" }]);
|
||
h.engine.arm("p1");
|
||
const r = await c.prompt("long turn");
|
||
assert.equal(r.outcome, "admitted");
|
||
await receiptState(c, r.receipt.id, "working");
|
||
await h.engine.waitPaused("p1");
|
||
return r.receipt.id;
|
||
}
|
||
|
||
const waitState = (h, state, ms = 6000) => until(() => h.ctrl.binding.state === state, ms, `binding ${state}`);
|
||
const synced = (h, ...cs) => until(() => cs.every((c) => c.binding?.controllerGeneration === h.ctrl.binding.controllerGeneration), 2000, "binding pushes");
|
||
|
||
test("H1: two takeovers with the same expected generation: one wins, +1; the other refuses generation", async () => {
|
||
const h = await started({ clients: 3 });
|
||
try {
|
||
const [, c2, c3] = h.clients;
|
||
const gen = h.ctrl.binding.controllerGeneration;
|
||
await synced(h, c2, c3);
|
||
const before = h.engine.bytes().length;
|
||
const [r2, r3] = await Promise.all([c2.send(c2.envelope("takeover")), c3.send(c3.envelope("takeover"))]);
|
||
const outcomes = [r2.outcome, r3.outcome].sort();
|
||
assert.deepEqual(outcomes, ["refused:generation", "transferred"]);
|
||
assert.equal(h.ctrl.binding.controllerGeneration, gen + 1);
|
||
const winner = r2.outcome === "transferred" ? c2 : c3;
|
||
assert.equal(h.ctrl.binding.controllerConnection, winner.connection.id);
|
||
assert.equal(h.engine.bytes().length, before);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H2: the old controller's prompt after a takeover commits is refused with zero engine bytes", async () => {
|
||
const h = await started({ clients: 2 });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
await synced(h, c2);
|
||
const stale = promptEnvelope(c1);
|
||
assert.equal((await c2.takeover()).outcome, "transferred");
|
||
const before = h.engine.bytes().length;
|
||
const r = await c1.send(stale, "late");
|
||
assert.equal(r.refusal, "generation");
|
||
await synced(h, c1);
|
||
const fresh = await c1.prompt("late, current generation");
|
||
assert.equal(fresh.refusal, "controller");
|
||
await tick(50);
|
||
assert.equal(h.engine.bytes().length, before);
|
||
assert.equal(h.ctrl.receipts.size, 0);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H3: a takeover while a prompt holds the dispatch lock: written under the old actor, or refused; never both", async () => {
|
||
// The prompt holds the lock at its recheck: the takeover waits, the write
|
||
// completes first, and the receipt keeps the old controller.
|
||
{
|
||
const g = gate();
|
||
const h = await started({ clients: 2, barrier: g.barrier });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
await synced(h, c2);
|
||
g.hold("before-recheck");
|
||
const p = await c1.prompt("held at the lock");
|
||
await g.waitHeld("before-recheck");
|
||
let took = null;
|
||
const t = c2.takeover().then((r) => (took = r));
|
||
await tick(50);
|
||
assert.equal(took, null, "the takeover waits for the dispatch lock");
|
||
g.release("before-recheck");
|
||
await t;
|
||
assert.equal(took.outcome, "transferred");
|
||
const r = await receiptState(c1, p.receipt.id, "finished");
|
||
assert.equal(r.state, "finished");
|
||
const req = h.ctrl.requests.get(p.receipt.request);
|
||
assert.equal(req.connection, c1.connection.id, "the receipt keeps the old actor's request");
|
||
assert.equal(cmds(h.engine, "prompt").length, 1);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// The takeover commits before the prompt reaches the lock: the recheck
|
||
// refuses it and nothing is written.
|
||
{
|
||
const g = gate();
|
||
const h = await started({ clients: 2, barrier: g.barrier });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
await synced(h, c2);
|
||
const before = h.engine.bytes().length;
|
||
g.hold("admitted");
|
||
const p = await c1.prompt("refused at recheck");
|
||
await g.waitHeld("admitted");
|
||
assert.equal((await c2.takeover()).outcome, "transferred");
|
||
g.release("admitted");
|
||
const r = await receiptState(c1, p.receipt.id, "dispatch-refused");
|
||
assert.ok(["controller", "generation"].includes(r.reasonCode), r.reasonCode);
|
||
assert.equal(h.engine.bytes().length, before, "never dispatched under either controller");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// Control leaves and comes back (A, B, A) while the prompt waits: the
|
||
// controller matches again, so only the generation check refuses it.
|
||
{
|
||
const g = gate();
|
||
const h = await started({ clients: 2, barrier: g.barrier });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
await synced(h, c2);
|
||
const before = h.engine.bytes().length;
|
||
g.hold("admitted");
|
||
const p = await c1.prompt("refused after A, B, A");
|
||
await g.waitHeld("admitted");
|
||
assert.equal((await c2.takeover()).outcome, "transferred");
|
||
await synced(h, c1);
|
||
assert.equal((await c1.takeover()).outcome, "transferred");
|
||
assert.equal(h.ctrl.binding.controllerConnection, c1.connection.id);
|
||
g.release("admitted");
|
||
const r = await receiptState(c1, p.receipt.id, "dispatch-refused");
|
||
assert.equal(r.reasonCode, "generation");
|
||
assert.equal(h.engine.bytes().length, before);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
test("H4: self-takeover is refused", async () => {
|
||
const h = await started();
|
||
try {
|
||
const gen = h.ctrl.binding.controllerGeneration;
|
||
const r = await h.client.takeover();
|
||
assert.equal(r.refusal, "already-controller");
|
||
assert.equal(h.ctrl.binding.controllerGeneration, gen);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H9: Interrupt racing a prompt's dispatch: before the write, dispatch-refused and no-turn; after, §3 rules", async () => {
|
||
{
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
const before = h.engine.bytes().length;
|
||
g.hold("admitted");
|
||
const p = await h.client.prompt("never written");
|
||
await g.waitHeld("admitted");
|
||
const i = await h.client.interrupt();
|
||
assert.equal(i.refusal, NO_TURN);
|
||
assert.deepEqual(i.effect, { dispatchRefused: p.receipt.id });
|
||
g.release("admitted");
|
||
const r = await receiptState(h.client, p.receipt.id, "dispatch-refused");
|
||
assert.equal(r.reasonCode, "fenced");
|
||
await tick(50);
|
||
assert.equal(h.engine.bytes().length, before, "no engine bytes");
|
||
assert.equal(h.ctrl.stops.size, 0, "no stop record");
|
||
assert.equal(h.ctrl.binding.admission, "open");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
{
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
h.engine.script([{ pause: "p1" }, { text: "never" }]);
|
||
h.engine.arm("p1");
|
||
g.hold("written");
|
||
const p = await h.client.prompt("written first");
|
||
await g.waitHeld("written");
|
||
const i = await h.client.interrupt();
|
||
assert.equal(i.outcome, "interrupt-fenced");
|
||
g.release("written");
|
||
await until(() => h.ctrl.events.some((e) => e.type === "reconciled" && e.stop === i.stop.id), 6000, "reconciled");
|
||
const r = await receiptState(h.client, p.receipt.id, "failed");
|
||
assert.equal(r.reasonCode, "interrupted");
|
||
assert.equal(cmds(h.engine, "prompt").length, 1);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
test("H10: Interrupt and force stop together: one stop chain, force stop supersedes", async () => {
|
||
// Force stop lands while the Interrupt's exchange is under way.
|
||
{
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
await heldRun(h);
|
||
g.hold("before-abort");
|
||
const i = await h.client.interrupt();
|
||
assert.equal(i.outcome, "interrupt-fenced");
|
||
await g.waitHeld("before-abort");
|
||
await synced(h, h.client);
|
||
const fs = await h.client.confirmed("force-stop");
|
||
assert.equal(fs.outcome, "force-stop-fenced", JSON.stringify(fs));
|
||
g.release("before-abort");
|
||
await waitState(h, "stopped");
|
||
const istop = h.ctrl.stops.get(i.stop.id), fstop = h.ctrl.stops.get(fs.stop.id);
|
||
assert.equal(istop.state, "superseded");
|
||
assert.equal(fstop.supersedes, istop.id);
|
||
assert.equal(fstop.state, "stopped");
|
||
assert.equal(h.ctrl.binding.stop, fstop.id);
|
||
assert.equal(cmds(h.engine, "abort").length, 0, "the superseded Interrupt writes nothing more");
|
||
assert.ok(!h.ctrl.events.some((e) => e.type === "reconciled"));
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// Both sent at once on a confirmation issued before either. The Interrupt
|
||
// yields once between its fence and the dispatch lock, so either may
|
||
// land first; there is one live stop either way. A stop the force stop
|
||
// didn't see changes its confirmation's context (H17).
|
||
for (const order of ["interrupt first", "force stop first"]) {
|
||
const h = await started();
|
||
try {
|
||
await heldRun(h);
|
||
const x = await confirm(h.client, "force-stop");
|
||
const ei = h.client.envelope("interrupt"), ef = h.client.envelope("force-stop", { confirmation: x });
|
||
const [ri, rf] = order === "interrupt first"
|
||
? await Promise.all([h.client.send(ei), h.client.send(ef)])
|
||
: await Promise.all([h.client.send(ef), h.client.send(ei)]).then(([f, i]) => [i, f]);
|
||
const won = [ri.outcome === "interrupt-fenced", rf.outcome === "force-stop-fenced"];
|
||
assert.equal(won.filter(Boolean).length, 1, `exactly one lands: ${ri.outcome} / ${rf.outcome}`);
|
||
if (won[0]) assert.equal(rf.refusal, "confirmation");
|
||
else assert.equal(ri.refusal, "fenced");
|
||
const live = [...h.ctrl.stops.values()].filter((s) => s.state !== "superseded");
|
||
assert.equal(live.length, 1, "one stop chain");
|
||
if (won[1]) await waitState(h, "stopped");
|
||
else await until(() => !["fenced", "cancelling", "stopping"].includes(h.ctrl.stops.get(ri.stop.id).state), 6000, "interrupt settles");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
// n2: an overlap read while the Interrupt waits before its abort means no
|
||
// abort, which would run whatever was queued. Mutant r2-n2 (no recheck
|
||
// after the pause) fails here.
|
||
test("H10: an overlap during the pause before the abort: no abort, the stop ends uncertain", async () => {
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
await heldRun(h);
|
||
g.hold("before-abort");
|
||
const i = await h.client.interrupt();
|
||
assert.equal(i.outcome, "interrupt-fenced");
|
||
await g.waitHeld("before-abort");
|
||
h.engine.queue("followUp", "queued in the pause");
|
||
await until(() => h.ctrl.evidence.overlaps.some((o) => o.signal === "O5"), 4000, "the overlap");
|
||
g.release("before-abort");
|
||
await until(() => h.ctrl.stops.get(i.stop.id)?.state === "uncertain", 6000, "stop uncertain");
|
||
assert.equal(cmds(h.engine, "abort").length, 0, "no abort after the overlap");
|
||
assert.equal(h.engine.runs.length, 1, "the queued item never ran");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H10: a no-turn Interrupt lifts only its own fence; admission stays closed under force stop, overlap or revocation", async () => {
|
||
const cases = {
|
||
"force stop": async (h) => {
|
||
const fs = await h.client.confirmed("force-stop");
|
||
assert.equal(fs.outcome, "force-stop-fenced", JSON.stringify(fs));
|
||
return "force-stop";
|
||
},
|
||
"overlap signal": async (h) => {
|
||
h.engine.emit({ type: "queue_update", steering: ["from an extension"], followUp: [] });
|
||
await until(() => h.ctrl.closers.has("overlap"), 2000, "overlap closer");
|
||
return "overlap";
|
||
},
|
||
revocation: async (h) => {
|
||
h.ctrl.revokeConnection(h.client.connection.id);
|
||
return "revocation";
|
||
},
|
||
};
|
||
for (const [name, close] of Object.entries(cases)) {
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
g.hold("interrupt-fenced");
|
||
const pending = h.client.interrupt();
|
||
await g.waitHeld("interrupt-fenced");
|
||
const reason = await close(h);
|
||
g.release("interrupt-fenced");
|
||
const r = await pending;
|
||
assert.equal(r.refusal, NO_TURN, name);
|
||
assert.equal(h.ctrl.binding.admission, "closed", `${name}: admission stays closed`);
|
||
assert.ok(h.ctrl.closers.has(reason), `${name}: the surviving reason holds`);
|
||
assert.ok(![...h.ctrl.closers].some((k) => k.startsWith("interrupt:")), `${name}: the Interrupt's own fence is gone`);
|
||
assert.ok([...h.ctrl.stops.values()].every((s) => s.mode !== "interrupt"), `${name}: no interrupt stop`);
|
||
if (name === "force stop") await waitState(h, "stopped");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
test("H11: the controller disconnects mid-turn: work continues, the claim is unchanged, control stays put", async () => {
|
||
const h = await started({ clients: 2 });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
const id = await heldRun(h, c1);
|
||
const claimBefore = structuredClone(h.ctrl.claim.record);
|
||
const controller = h.ctrl.binding.controllerConnection;
|
||
c1.close();
|
||
await until(() => h.ctrl.connections.get(controller).rec.state === "disconnected", 2000, "disconnected");
|
||
h.engine.resume("p1");
|
||
await until(() => h.ctrl.receipts.get(id).state === "finished", 4000, "the turn finishes");
|
||
assert.deepEqual(h.ctrl.claim.record, claimBefore, "the claim is unchanged");
|
||
assert.equal(h.ctrl.binding.controllerConnection, controller, "control stays with the disconnected connection");
|
||
await synced(h, c2);
|
||
assert.equal((await c2.prompt("observer")).refusal, "controller");
|
||
await tick(100);
|
||
assert.equal(h.ctrl.binding.controllerConnection, controller, "nothing transfers automatically");
|
||
assert.equal((await c2.takeover()).outcome, "transferred", "an explicit takeover");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H12: an exact retry after reconnecting to the same incarnation returns the same receipt; one dispatch", async () => {
|
||
const h = await started();
|
||
try {
|
||
const c = h.client;
|
||
const e = promptEnvelope(c);
|
||
const first = await c.send(e, "once");
|
||
assert.equal(first.outcome, "admitted");
|
||
await receiptState(c, first.receipt.id, "finished");
|
||
const incarnation = c.incarnation;
|
||
c.close();
|
||
await c.connect();
|
||
assert.equal(c.incarnation, incarnation);
|
||
const again = await c.send({ ...e, connection: c.connection.id }, "once");
|
||
assert.equal(again.outcome, "existing:finished");
|
||
assert.equal(again.receipt.id, first.receipt.id);
|
||
assert.equal(cmds(h.engine, "prompt").length, 1);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H13: a retry with the same request ID and different text is refused", async () => {
|
||
const h = await started();
|
||
try {
|
||
const e = promptEnvelope(h.client);
|
||
const first = await h.client.send(e, "original");
|
||
await receiptState(h.client, first.receipt.id, "finished");
|
||
const r = await h.client.send(e, "changed");
|
||
assert.equal(r.refusal, "conflicting-request");
|
||
assert.equal(cmds(h.engine, "prompt").length, 1);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
// The old engine's stdout stays open after its stop, as a pipe with bytes
|
||
// still in it would.
|
||
class LateLauncher extends FakeLauncher {
|
||
async launch(a) {
|
||
const r = await super.launch(a);
|
||
const e = this.last;
|
||
const kill = e.kill;
|
||
e.kill = () => {
|
||
const end = e.end;
|
||
e.end = () => {};
|
||
kill();
|
||
e.end = end;
|
||
};
|
||
return r;
|
||
}
|
||
}
|
||
|
||
test("H14: late stdout from the old engine after a replacement is dropped by incarnation, counted, never rendered", async () => {
|
||
const launcher = new LateLauncher();
|
||
const h = await started({ launcher });
|
||
try {
|
||
const c = h.client;
|
||
const old = launcher.last;
|
||
const fs = await c.confirmed("force-stop");
|
||
assert.equal(fs.outcome, "force-stop-fenced");
|
||
await waitState(h, "stopped");
|
||
await synced(h, c);
|
||
const rec = await c.confirmed("recover", { stop: fs.stop.id });
|
||
assert.equal(rec.outcome, "recovery-eligible", JSON.stringify(rec));
|
||
await h.ctrl.launch(rec.data.eligibility);
|
||
await waitState(h, "active");
|
||
assert.notEqual(launcher.last, old);
|
||
const dropped = h.ctrl.evidence.dropped.lines;
|
||
const events = c.events.length;
|
||
old.emit({ type: "agent_start" });
|
||
old.emit({ type: "message_start", message: { role: "assistant", content: [{ type: "text", text: "LATE FROM THE OLD ENGINE" }] } });
|
||
await tick(100);
|
||
assert.equal(h.ctrl.evidence.dropped.lines, dropped + 2, "counted in the evidence");
|
||
assert.equal(c.events.length, events, "never rendered");
|
||
assert.ok(!JSON.stringify(c.pushes).includes("LATE FROM THE OLD ENGINE"));
|
||
assert.deepEqual(h.ctrl.evidence.overlaps, [], "no signal on the new execution");
|
||
assert.equal(h.ctrl.binding.admission, "open");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H15: a revoked connection's command is refused; the revocation fence holds", async () => {
|
||
const h = await started({ clients: 2 });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
const gen = h.ctrl.binding.controllerGeneration;
|
||
const before = h.engine.bytes().length;
|
||
h.ctrl.revokeConnection(c1.connection.id);
|
||
assert.equal(h.ctrl.binding.controllerConnection, null);
|
||
assert.equal(h.ctrl.binding.controllerGeneration, gen + 1);
|
||
assert.equal(h.ctrl.binding.admission, "closed");
|
||
assert.ok([...h.ctrl.stops.values()].some((s) => s.mode === "revocation" && s.revokedConnection === c1.connection.id));
|
||
assert.equal((await c1.prompt("after revocation")).refusal, "channel");
|
||
await synced(h, c2);
|
||
assert.equal((await c2.takeover()).refusal, "fenced", "CHAT-03 never lifts the revocation fence");
|
||
await tick(50);
|
||
assert.equal(h.engine.bytes().length, before);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H16: a second controller for the same session refuses already-active; the first is untouched", async () => {
|
||
const h = await started();
|
||
try {
|
||
const claim = structuredClone(h.ctrl.claim.record);
|
||
const { ctrl: second, launcher } = controllerFor(h.fx);
|
||
await assert.rejects(() => second.start(), (e) => e.code === ALREADY_ACTIVE);
|
||
assert.equal(launcher.launches.length, 0);
|
||
assert.deepEqual(h.ctrl.claim.record, claim);
|
||
assert.ok(existsSync(h.ctrl.socketPath), "the first controller's socket is still there");
|
||
const p = await h.client.prompt("still mine");
|
||
await receiptState(h.client, p.receipt.id, "finished");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H10: a second force stop while the first escalation runs refuses fenced; one escalation, and the claim records only the first stop's phases", async () => {
|
||
{
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
const x1 = await confirm(h.client, "force-stop");
|
||
g.hold("phase-term");
|
||
assert.equal((await h.client.request("force-stop", { confirmation: x1 })).outcome, "force-stop-fenced");
|
||
await g.waitHeld("phase-term");
|
||
const first = h.ctrl.binding.stop;
|
||
const bytes = h.engine.bytes().length;
|
||
await synced(h, h.client);
|
||
// Issued after the first stop began, so the stop it binds is current.
|
||
const x2 = await confirm(h.client, "force-stop");
|
||
const second = await h.client.request("force-stop", { confirmation: x2 });
|
||
assert.equal(second.refusal, "fenced");
|
||
assert.equal(h.ctrl.binding.stop, first, "no second stop record");
|
||
assert.notEqual(h.ctrl.confirmations.get(x2).state, "consumed", "the refusal comes before the confirmation is used");
|
||
assert.deepEqual(h.ctrl.store.head(h.ctrl.sessionK).record.stop.id, first);
|
||
assert.equal(h.ctrl.store.head(h.ctrl.sessionK).record.stop.phaseStarted, "term");
|
||
g.release("phase-term");
|
||
await waitState(h, "stopped");
|
||
assert.equal(h.launcher.stops, 1, "one escalation");
|
||
for (const key of [h.ctrl.seatK, h.ctrl.sessionK]) {
|
||
const recs = h.ctrl.store.revisions(key).map((r) => JSON.parse(r.text));
|
||
assert.ok(recs.every((r) => !r.stop || r.stop.id === first), "every stop on the claim is the first");
|
||
assert.deepEqual(recs.at(-1).stop.phaseStarted, "kill");
|
||
assert.equal(recs.at(-1).state, "stopped");
|
||
}
|
||
assert.equal(h.ctrl.events.filter((e) => e.type === "stopping").length, 1);
|
||
assert.equal(h.engine.bytes().length, bytes, "no engine bytes after the first fence");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// Once the first ends uncertain, a force stop can be retried.
|
||
{
|
||
const h = await started({ launcher: new FakeLauncher({ stopOutcome: "unavailable" }) });
|
||
try {
|
||
const x1 = await confirm(h.client, "force-stop");
|
||
assert.equal((await h.client.request("force-stop", { confirmation: x1 })).outcome, "force-stop-fenced");
|
||
await waitState(h, "uncertain");
|
||
await until(() => h.ctrl.escalating === null, 2000, "the first escalation ends");
|
||
await synced(h, h.client);
|
||
assert.equal((await h.client.request("force-stop", { confirmation: x1 })).refusal, "confirmation", "H17: a used confirmation is refused");
|
||
assert.equal((await h.client.request("force-stop", { confirmation: await confirm(h.client, "force-stop") })).outcome, "force-stop-fenced");
|
||
await until(() => h.launcher.stops === 2, 2000, "a second escalation");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
test("H17: a confirmation reused, answered from another connection, or used after the stop changed is refused", async () => {
|
||
// Reused while its force stop is still running.
|
||
{
|
||
const g = gate();
|
||
const h = await started({ barrier: g.barrier });
|
||
try {
|
||
const x = await confirm(h.client, "force-stop");
|
||
g.hold("force-stop-recorded");
|
||
const first = await h.client.request("force-stop", { confirmation: x });
|
||
assert.equal(first.outcome, "force-stop-fenced");
|
||
await g.waitHeld("force-stop-recorded");
|
||
await synced(h, h.client);
|
||
// One escalation at a time: the reuse refuses `fenced` before its
|
||
// confirmation is looked at. The reuse after it ends is below.
|
||
const reuse = await h.client.request("force-stop", { confirmation: x });
|
||
assert.equal(reuse.refusal, "fenced");
|
||
g.release("force-stop-recorded");
|
||
await waitState(h, "stopped");
|
||
assert.equal(h.launcher.stops, 1);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// Answered from another connection.
|
||
{
|
||
const h = await started({ clients: 2 });
|
||
try {
|
||
const [c1, c2] = h.clients;
|
||
await synced(h, c2);
|
||
const issued = await c1.request("issue-confirmation", { operationToConfirm: "force-stop" });
|
||
const id = issued.data.confirmation.id;
|
||
const r = await c2.request("answer-confirmation", { confirmation: id, answer: "confirm" });
|
||
assert.equal(r.refusal, "confirmation");
|
||
assert.equal(h.ctrl.confirmations.get(id).state, "pending");
|
||
assert.equal((await c2.request("force-stop", { confirmation: id })).refusal, "controller");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
// Used after an Interrupt changed the stop context.
|
||
{
|
||
const h = await started();
|
||
try {
|
||
const x = await confirm(h.client, "force-stop");
|
||
await heldRun(h);
|
||
const i = await h.client.interrupt();
|
||
await until(() => h.ctrl.events.some((e) => e.type === "reconciled" && e.stop === i.stop.id), 6000, "reconciled");
|
||
await synced(h, h.client);
|
||
const r = await h.client.request("force-stop", { confirmation: x });
|
||
assert.equal(r.refusal, "confirmation");
|
||
assert.equal(h.ctrl.binding.state, "active");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
}
|
||
});
|
||
|
||
test("H18: two prompts before any native output: the second refuses busy; one engine write", async () => {
|
||
const h = await started();
|
||
try {
|
||
h.engine.arm("843");
|
||
const [a, b] = await Promise.all([h.client.send(promptEnvelope(h.client), "first"), h.client.send(promptEnvelope(h.client), "second")]);
|
||
assert.equal(a.outcome, "admitted");
|
||
assert.equal(b.refusal, BUSY);
|
||
await h.engine.waitPaused("843");
|
||
h.engine.resume("843");
|
||
await receiptState(h.client, a.receipt.id, "finished");
|
||
assert.equal(cmds(h.engine, "prompt").length, 1);
|
||
assert.ok(!h.engine.bytes().toString().includes("second"));
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H19: the pipe fails mid-line under a large prompt: delivery-unknown transport-unknown, poisoned, no later write", async () => {
|
||
const h = await started();
|
||
try {
|
||
const before = h.engine.bytes().length;
|
||
h.engine.stdinGate.left = 40000;
|
||
const text = "x".repeat(250000);
|
||
const p = await h.client.prompt(text);
|
||
const r = await receiptState(h.client, p.receipt.id, "delivery-unknown");
|
||
assert.equal(r.reasonCode, TRANSPORT_UNKNOWN);
|
||
assert.equal(h.engine.bytes().length - before, 40000, "the engine got a partial line");
|
||
assert.equal(cmds(h.engine, "prompt").length, 0);
|
||
assert.ok(h.ctrl.exec.link.poisoned);
|
||
await waitState(h, "uncertain");
|
||
const after = h.engine.bytes().length;
|
||
h.engine.stdinGate.left = null;
|
||
assert.equal((await h.client.prompt("again")).refusal, "fenced");
|
||
await tick(100);
|
||
assert.equal(h.engine.bytes().length, after, "no later write and no retry");
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
});
|
||
|
||
test("H19: the link itself never writes again after an unknown outcome, whoever calls it", async () => {
|
||
// The controller checks the poison before every write; this pins the
|
||
// link's own refusal, so a caller that forgets the check still can't
|
||
// write after a partial line.
|
||
const calls = [];
|
||
const stdin = new Writable({
|
||
write(chunk, _enc, cb) {
|
||
calls.push(chunk.length);
|
||
cb(calls.length === 1 ? Object.assign(new Error("broken pipe"), { code: "EPIPE" }) : null);
|
||
},
|
||
});
|
||
stdin.on("error", () => {});
|
||
const link = new EngineLink({ execution: "exec-x", stdin, stdout: new PassThrough(), onLine() {}, onGap() {} });
|
||
assert.deepEqual(await link.write({ type: "prompt", id: "p1", message: "x" }), { outcome: "unknown", reason: "EPIPE" });
|
||
assert.equal(link.poisoned, "EPIPE");
|
||
assert.deepEqual(await link.write({ type: "get_state", id: "s1" }), { outcome: "refused", reason: "poisoned" });
|
||
assert.equal(calls.length, 1, "one write reached the pipe");
|
||
assert.equal(link.bytesWritten, 0);
|
||
});
|
||
|
||
// ---- the controller in a child process ------------------------------------
|
||
|
||
const logOf = (path) => (existsSync(path) ? readFileSync(path, "utf8").split("\n").filter(Boolean).map((l) => JSON.parse(l)) : []);
|
||
const promptCmds = (path) => logOf(path).filter((v) => v.t === "cmd" && v.cmd.type === "prompt").length;
|
||
const engineCount = (path) => logOf(path).filter((v) => v.t === "argv").length;
|
||
const inText = (path) => Buffer.concat(logOf(path).filter((v) => v.t === "in").map((v) => Buffer.from(v.b64, "base64"))).toString();
|
||
|
||
function childFixture() {
|
||
const fx = track(fixture());
|
||
const log = join(fx.base, "fake.log");
|
||
return { fx, log, fakeEnv: { FAKE_PI_CONTROL: join(fx.base, "fake.sock"), FAKE_PI_LOG: log } };
|
||
}
|
||
|
||
async function connectTo(socketPath, client = null) {
|
||
const c = client ?? new ConversationClient({ socketPath });
|
||
c.socketPath = socketPath;
|
||
await c.connect();
|
||
return c;
|
||
}
|
||
|
||
test("H19: the controller dies mid-write of a large line: after restart the outcome is unknown and nothing is resent", async () => {
|
||
const { fx, log, fakeEnv } = childFixture();
|
||
const a = spawnController({ fx, launcher: "pgroup", units: "absent", fakeEnv });
|
||
const ra = await a.next((m) => m.ready || m.error);
|
||
assert.ok(ra.ready, JSON.stringify(ra));
|
||
const fake = new ControlClient(fakeEnv.FAKE_PI_CONTROL);
|
||
await fake.connect();
|
||
const c = await connectTo(ra.socketPath);
|
||
await c.takeover();
|
||
await fake.call("stall");
|
||
const p = await c.prompt("y".repeat(250000));
|
||
assert.equal(p.outcome, "admitted");
|
||
await tick(300);
|
||
assert.notEqual(c.receipt(p.receipt.id).state, "dispatched", "the write never completed");
|
||
a.proc.kill("SIGKILL");
|
||
await a.exited;
|
||
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");
|
||
await connectTo(rb.socketPath, c);
|
||
assert.ok(c.unknown.some((u) => u.receipt === p.receipt.id && u.display === OUTCOME_UNKNOWN), "shown as outcome unknown");
|
||
const st = (await fake.call("state")).result;
|
||
assert.ok(!st.commands.some((x) => x.type === "prompt"), "the engine never read a complete prompt");
|
||
assert.equal(engineCount(log), 1);
|
||
c.close();
|
||
fake.close();
|
||
b.send("close");
|
||
await b.exited;
|
||
reap(fx);
|
||
});
|
||
|
||
test("H20: the line is written but the ack is lost when the controller dies: orphan, outcome unknown, nothing resent", async () => {
|
||
const { fx, log, fakeEnv } = childFixture();
|
||
const a = spawnController({ fx, launcher: "pgroup", units: "absent", fakeEnv, dieAt: { written: 1 } });
|
||
const ra = await a.next((m) => m.ready || m.error);
|
||
assert.ok(ra.ready, JSON.stringify(ra));
|
||
const c = await connectTo(ra.socketPath);
|
||
await c.takeover();
|
||
const p = await c.prompt("written, then the controller dies");
|
||
await a.next((m) => m.dying === "written");
|
||
await a.exited;
|
||
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");
|
||
await connectTo(rb.socketPath, c);
|
||
assert.ok(c.unknown.some((u) => u.receipt === p.receipt.id && u.display === OUTCOME_UNKNOWN));
|
||
assert.equal(c.binding.state, "uncertain");
|
||
await tick(200);
|
||
assert.equal(promptCmds(log), 1, "nothing resent");
|
||
assert.equal(engineCount(log), 1);
|
||
c.close();
|
||
b.send("close");
|
||
await b.exited;
|
||
reap(fx);
|
||
});
|
||
|
||
// A controller that dies after its prompt line is written, and its restart.
|
||
async function crashAfterWrite(launcher) {
|
||
const { fx, log, fakeEnv } = childFixture();
|
||
const units = launcher === "scope" ? undefined : "absent";
|
||
const a = spawnController({ fx, launcher, units, fakeEnv, dieAt: { written: 1 } });
|
||
const ra = await a.next((m) => m.ready || m.error);
|
||
assert.ok(ra.ready, JSON.stringify(ra));
|
||
const c = await connectTo(ra.socketPath);
|
||
await c.takeover();
|
||
const oldToken = c.incarnation;
|
||
const e = promptEnvelope(c);
|
||
const text = "dispatched before the crash";
|
||
const sent = c.send(e, text);
|
||
await a.next((m) => m.dying === "written");
|
||
await a.exited;
|
||
await sent;
|
||
const b = spawnController({ fx, launcher, units, fakeEnv });
|
||
const rb = await b.next((m) => m.ready || m.error);
|
||
assert.ok(rb.ready, JSON.stringify(rb));
|
||
assert.equal(rb.started.launched, false);
|
||
await connectTo(rb.socketPath, c);
|
||
assert.notEqual(c.incarnation, oldToken);
|
||
return { fx, log, c, b, e, text, oldToken };
|
||
}
|
||
|
||
test("H21: a retry of the exact request with the old token after a crash is stale-incarnation; no second write", async () => {
|
||
const { fx, log, c, b, e, text, oldToken } = await crashAfterWrite("pgroup");
|
||
const retry = await c.send({ ...e, connection: c.connection.id }, text, { incarnation: oldToken });
|
||
assert.equal(retry.refusal, STALE_INCARNATION);
|
||
assert.equal(retry.display, OUTCOME_UNKNOWN);
|
||
await tick(100);
|
||
assert.equal(promptCmds(log), 1, "no second engine write");
|
||
// The process-group fallback can't prove the orphan's cohort, so its
|
||
// force stop ends uncertain; recovery needs a scope (H22).
|
||
assert.equal((await c.confirmed("acquire-recovery-control")).outcome, "recovery-control-acquired");
|
||
assert.equal((await c.confirmed("force-stop")).outcome, "force-stop-fenced");
|
||
assert.ok(await c.waitFor(() => c.binding?.state === "uncertain" && b.msgs.length > 0, 10000));
|
||
c.close();
|
||
b.send("close");
|
||
await b.exited;
|
||
reap(fx);
|
||
});
|
||
|
||
test("H22: after H21 and a valid recovery, a new request with the new token is admitted", { skip: !SCOPE && "systemd user scopes unavailable", timeout: 60000 }, async () => {
|
||
const { fx, log, c, b, e, text, oldToken } = await crashAfterWrite("scope");
|
||
assert.equal((await c.send({ ...e, connection: c.connection.id }, text, { incarnation: oldToken })).refusal, STALE_INCARNATION);
|
||
assert.equal((await c.confirmed("acquire-recovery-control")).outcome, "recovery-control-acquired");
|
||
const fs = await c.confirmed("force-stop");
|
||
assert.equal(fs.outcome, "force-stop-fenced", JSON.stringify(fs));
|
||
assert.ok(await c.waitFor(() => c.binding?.state === "stopped", 15000), `binding ${c.binding?.state}`);
|
||
const rec = await c.confirmed("recover", { stop: fs.stop.id });
|
||
assert.equal(rec.outcome, "recovery-eligible", JSON.stringify(rec));
|
||
b.send(`launch ${rec.data.eligibility}`);
|
||
const launched = await b.next((m) => m.launched !== undefined, 15000);
|
||
assert.ok(launched.launched, JSON.stringify(launched));
|
||
assert.equal(launched.incarnation, c.incarnation, "the token the client holds is the new one");
|
||
if (!(await c.waitFor(() => c.binding?.state === "active" && c.binding.admission === "open", 15000))) {
|
||
b.send("evidence");
|
||
const ev = await b.next((m) => m.evidence !== undefined);
|
||
assert.fail(`binding ${c.binding?.state}: ${JSON.stringify({ uncertain: ev.evidence.uncertain, stops: ev.evidence.stops, internal: ev.evidence.internal, gaps: ev.evidence.gaps, overlaps: ev.evidence.overlaps, stderr: ev.evidence.stderrTail.slice(-800) })}`);
|
||
}
|
||
assert.equal((await c.takeover()).outcome, "transferred");
|
||
const p = await c.prompt("a new request");
|
||
assert.equal(p.outcome, "admitted", JSON.stringify(p));
|
||
await receiptState(c, p.receipt.id, "finished", 8000);
|
||
assert.equal(promptCmds(log), 2, "one per engine; the old request was never resent");
|
||
assert.equal(inText(log).split(text).length - 1, 1);
|
||
c.close();
|
||
b.send("kill-engine");
|
||
await b.exited;
|
||
reap(fx);
|
||
});
|
||
|
||
test("H23: requests pending at a restart are not resent; each shows outcome unknown", async () => {
|
||
const { fx, log, fakeEnv } = childFixture();
|
||
const a = spawnController({ fx, launcher: "pgroup", units: "absent", fakeEnv, holdAt: ["before-recheck"] });
|
||
const ra = await a.next((m) => m.ready || m.error);
|
||
assert.ok(ra.ready, JSON.stringify(ra));
|
||
const c = await connectTo(ra.socketPath);
|
||
await c.takeover();
|
||
const p1 = await c.prompt("held at the recheck");
|
||
assert.equal(p1.outcome, "admitted");
|
||
await a.next((m) => m.held === "before-recheck");
|
||
// Both wait behind the held dispatch lock.
|
||
const p2 = c.prompt("pending one");
|
||
const ob = c.observe();
|
||
await tick(100);
|
||
assert.equal(c.pending.size, 2);
|
||
a.proc.kill("SIGKILL");
|
||
await a.exited;
|
||
const [r2, rob] = await Promise.all([p2, ob]);
|
||
assert.equal(r2.outcome, "outcome-unknown");
|
||
assert.equal(rob.outcome, "outcome-unknown");
|
||
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));
|
||
await connectTo(rb.socketPath, c);
|
||
await tick(200);
|
||
assert.equal(c.pending.size, 0, "the library resends none");
|
||
assert.equal(promptCmds(log), 0, "zero engine bytes for them");
|
||
assert.ok(!inText(log).includes("pending one") && !inText(log).includes("held at the recheck"));
|
||
const shown = c.unknown.filter((u) => u.display === OUTCOME_UNKNOWN);
|
||
assert.equal(shown.filter((u) => u.operation === "prompt" && u.request).length, 1, "the pending prompt");
|
||
assert.equal(shown.filter((u) => u.operation === "observe").length, 1, "the pending observe");
|
||
assert.ok(shown.some((u) => u.receipt === p1.receipt.id), "the admitted prompt");
|
||
c.close();
|
||
b.send("close");
|
||
await b.exited;
|
||
reap(fx);
|
||
});
|