// 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, killShims, shimsGone } from "./harness.mjs"; const reaped = []; after(async () => { for (const fx of reaped) reap(fx); killShims(); const left = await shimsGone(); killChildren(); cleanupAll(); assert.deepEqual(left, [], "no shim outlives this file (#1533)"); }); 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("a force stop whose fence throws leaves no escalation flag behind, so the next force stop runs (Darkwing F2, #1507)", async () => { const h = await started(); try { const link = h.ctrl.exec.link; const poison = link.poison; link.poison = () => { throw new Error("poison failed"); }; const failed = await h.client.request("force-stop", { confirmation: await confirm(h.client, "force-stop") }); assert.equal(failed.outcome, "error", JSON.stringify(failed)); assert.equal(h.ctrl.escalating, null, "the flag is cleared on the throw"); assert.equal(h.launcher.stops ?? 0, 0, "no escalation ran"); link.poison = poison; await synced(h, h.client); const next = await h.client.request("force-stop", { confirmation: await confirm(h.client, "force-stop") }); assert.equal(next.outcome, "force-stop-fenced", JSON.stringify(next)); await waitState(h, "stopped"); assert.equal(h.launcher.stops, 1); assert.equal(h.ctrl.escalating, null); } 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); }); // Last in the file. A shim ignores SIGTERM and outlives its controller, so a // test that leaves one running relies on the after hook's kill (#1533). test("no shim from this file's tests is left running for the after hook", async () => { assert.deepEqual(await shimsGone(), [], "a test left its shim running"); });