Files
stack/packages/conversation/tests/turns.test.mjs
T
jason.woltjeandClaude Opus 5.5 243e153c8b feat(conversation): CHAT-03 I1, mediated control of a sealed headless Pi (#1507)
Controller, claim store, live-session guard, engine link and seal,
turn tracker, cohort force stop and recovery, client library,
transcript and mediated terminal, with the fake engine and tests.
Fixtures only; no live cutover.

Dewey built it. Darkwing (comment 26690) and Filbert (comment 26694)
approved round 2. Manifest I1-r2-manifest.sha256 (2b48e333, 27 files).
Suites on an export: conversation 152/152, control-board 124, webui 14,
seat 19, chat-00/01/01c checks, and all nine scripts/test-*.sh green.
Follow-ups for I3 are in DEFERRED. Gate E stays with Jason.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-10-04 15:47:53 -05:00

984 lines
43 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// CHAT-03 §3 turns, dispatch and Interrupt (#1507): N1–N25 on the in-process
// fake engine. Unsealed-extension effects are simulated by the fake; no real
// extension loads (the binding would refuse it, N24).
import { test, after } from "node:test";
import assert from "node:assert/strict";
import { PassThrough } from "node:stream";
import { join } from "node:path";
import { Controller, HANDLED_WITHOUT_RUN, ACK_WITHOUT_START, NO_TURN, RUN_OVERLAP, TRANSPORT_UNKNOWN } from "../src/controller.mjs";
import { AUTHORITY } from "../src/cohort.mjs";
import { FixtureVerifier, sha256 } from "../src/records.mjs";
import { SEAL_FLAGS, UNSEALED_ENGINE, checkSeal } from "../src/pi-pin.mjs";
import { LineSplitter, encodeLine, parseLine } from "../src/framing.mjs";
import { BUSY_ERROR, FakeLauncher, FakePi } from "./fake-pi.mjs";
import { fixture, started, receiptState, outcomeUnknownPush, sessionText, cleanupAll, tick, noUnits, FAST } from "./harness.mjs";
after(() => cleanupAll());
// Records every proof the controller posts, so a test can see that no
// turnProof other than `turnState: interrupted` was ever produced (mutant 32).
class SpyVerifier extends FixtureVerifier {
constructor() {
super({ authorities: [AUTHORITY] });
this.posted = [];
}
post(p) {
this.posted.push(structuredClone(p));
return super.post(p);
}
}
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);
}
}
const cmds = (engine, type) => engine.commands.filter((c) => c.type === type);
const stopEvidence = (ctrl, id) => ctrl.evidence.stops.filter((e) => e.stop === id);
const lastStopEvidence = (ctrl, id) => stopEvidence(ctrl, id).at(-1) ?? null;
const signals = (ctrl) => ctrl.evidence.overlaps.map((o) => o.signal);
const reconciled = (ctrl, id) => ctrl.events.some((e) => e.type === "reconciled" && e.stop === id);
const turnProofs = (v) => v.posted.filter((p) => p.kind === "turnProof");
async function setup(opts = {}) {
const verifier = new SpyVerifier();
const h = await started({ verifier, ...opts });
return { ...h, verifier };
}
// Admits a prompt and returns its receipt id.
async function prompt(h, text = "do the thing") {
const r = await h.client.prompt(text);
assert.equal(r.outcome, "admitted", JSON.stringify(r));
return r.receipt.id;
}
// A Mosaic run held at step `p1` with its receipt `working`.
async function heldRun(h, rest = [{ text: "never", stop: "stop" }]) {
h.engine.script([{ pause: "p1" }, ...rest]);
h.engine.arm("p1");
const id = await prompt(h);
await receiptState(h.client, id, "working");
await h.engine.waitPaused("p1");
return id;
}
async function interrupt(h) {
const r = await h.client.interrupt();
assert.equal(r.outcome, "interrupt-fenced", JSON.stringify(r));
return r.stop.id;
}
async function stopSettled(h, id) {
await until(() => {
const s = h.ctrl.stops.get(id);
return s && (s.state === "uncertain" || s.state === "turn-interrupted" || s.state === "superseded");
}, 8000, `stop ${id} to settle`);
return h.ctrl.stops.get(id);
}
// The stop is uncertain: no reconciled, prompts refuse, force stop works.
async function assertUncertainStop(h, stopId, outcome) {
const s = await stopSettled(h, stopId);
assert.equal(s.state, "uncertain");
assert.equal(lastStopEvidence(h.ctrl, stopId)?.outcome, outcome);
assert.equal(reconciled(h.ctrl, stopId), false);
assert.equal(h.ctrl.binding.admission, "closed");
assert.ok(turnProofs(h.verifier).every((p) => p.turnState === "interrupted"), "no other turnState is ever written");
assert.ok(!turnProofs(h.verifier).some((p) => p.stop === stopId), "no turnProof for an uncertain stop");
const before = h.engine.bytes().length;
const p = await h.client.prompt("again");
assert.equal(p.refusal, "fenced");
assert.equal(h.engine.bytes().length, before, "a refused prompt writes no engine bytes");
const fs = await h.client.confirmed("force-stop");
assert.equal(fs.outcome, "force-stop-fenced", JSON.stringify(fs));
await until(() => h.ctrl.binding.state === "stopped", 4000, "force stop");
}
test("N25: ordinary Interrupt reconciles; a non-empty queue_update in the window is O5", async () => {
const h = await setup();
try {
const id = await heldRun(h);
const stopId = await interrupt(h);
await until(() => reconciled(h.ctrl, stopId), 6000, "reconciled");
assert.deepEqual(signals(h.ctrl), [], "Pi's empty queue_update before each clear response is no signal");
const s = h.ctrl.stops.get(stopId);
assert.equal(s.state, "turn-interrupted");
assert.equal(s.nativeQueue, "cleared");
assert.equal((await receiptState(h.client, id, "failed")).reasonCode, "interrupted");
// Each clear_queue answered after Pi's empty queue_update.
const clears = cmds(h.engine, "clear_queue");
assert.equal(clears.length, 2, "the interrupt clear and the post-settle clear");
const types = h.engine.sent.map((l) => JSON.parse(l)).map((v) => (v.type === "response" ? `response:${v.command}` : v.type));
for (let i = 0; i < types.length; i++) if (types[i] === "response:clear_queue") assert.equal(types[i - 1], "queue_update");
assert.equal(h.ctrl.binding.admission, "open");
const next = await h.client.prompt("next");
assert.equal(next.outcome, "admitted");
await receiptState(h.client, next.receipt.id, "finished");
} finally {
await h.close();
}
// The same window with a non-empty queue_update.
const h2 = await setup();
try {
await heldRun(h2);
h2.engine.arm("clear");
const stopId = await interrupt(h2);
await h2.engine.waitPaused("clear");
h2.engine.queue("steer", "slipped in");
h2.engine.resume("clear");
const s = await stopSettled(h2, stopId);
assert.equal(s.state, "uncertain");
assert.ok(signals(h2.ctrl).includes("O5"));
assert.equal(reconciled(h2.ctrl, stopId), false);
} finally {
await h2.close();
}
});
test("N1: an extension's follow-up queued after the fence is cleared before any abort; O5, Unknown", async () => {
const h = await setup();
try {
const id = await heldRun(h);
h.engine.arm("clear");
const stopId = await interrupt(h);
// The clear_queue is in Pi's hands but not yet run; the extension queues.
await h.engine.waitPaused("clear");
h.engine.queue("followUp", "external item");
h.engine.resume("clear");
const s = await stopSettled(h, stopId);
const order = h.engine.commands.map((c) => c.type).filter((t) => t === "clear_queue" || t === "abort");
assert.equal(order[0], "clear_queue");
assert.equal(cmds(h.engine, "abort").length, 0, "a non-empty clear stops before abort");
assert.equal(h.engine.runs.length, 1);
assert.ok(!h.engine.sent.some((l) => l.includes("external item") && l.includes('"message_start"')), "the follow-up never ran");
assert.ok(signals(h.ctrl).filter((x) => x === "O5").length >= 2, "the queue_update and the non-empty clear");
assert.equal(s.state, "uncertain");
assert.equal(s.nativeQueue, "unknown", "a non-empty clear leaves the native queue unknown");
assert.equal(lastStopEvidence(h.ctrl, stopId).outcome, "unknown");
const digest = sha256("external item");
assert.ok(h.ctrl.evidence.overlaps.some((o) => o.removed?.some((r) => r.digest === digest && r.bytes === 13)), "evidence holds the item's digest and size");
assert.equal(h.client.receipt(id).state, "working");
await until(() => outcomeUnknownPush(h.client, id), 2000, "the receipt shown as outcome unknown");
assert.equal(outcomeUnknownPush(h.client, id).reason, RUN_OVERLAP);
assert.ok(!JSON.stringify(h.ctrl.receipts.get(id)).includes("external item"), "the removed item is never attributed to the request");
assert.equal(cmds(h.engine, "prompt").length, 1, "nothing is resent");
await assertUncertainStop(h, stopId, "unknown");
} finally {
await h.close();
}
});
test("N1: a follow-up queued before the fence is O5 at once; the Interrupt refuses fenced", async () => {
const h = await setup();
try {
await heldRun(h);
h.engine.queue("followUp", "early");
await until(() => h.ctrl.binding.state === "uncertain", 2000, "uncertain");
const r = await h.client.interrupt();
assert.equal(r.refusal, "fenced");
assert.equal(h.ctrl.stops.size, 0);
} finally {
await h.close();
}
});
test("N2: with abort first, the fake runs the external item (the ordering guard has teeth)", async () => {
// The mutant's order, driven straight at the fake: abort before clear.
const drive = async (order) => {
const input = new PassThrough(), output = new PassThrough();
const fake = new FakePi({ input, output });
const lines = [];
output.on("data", (c) => lines.push(...c.toString().split("\n").filter(Boolean).map((l) => JSON.parse(l))));
fake.script([{ pause: "p1" }, { text: "x" }]);
fake.arm("p1");
input.write(encodeLine({ id: "p", type: "prompt", message: "mosaic" }));
await fake.waitPaused("p1");
fake.queue("followUp", "external item");
for (const [i, type] of order.entries()) {
input.write(encodeLine({ id: `c${i}`, type }));
await until(() => lines.some((v) => v.id === `c${i}`));
}
await fake.waitIdle();
return lines.some((v) => v.type === "message_start" && v.message?.role === "user" && v.message.content?.[0]?.text === "external item");
};
assert.equal(await drive(["abort", "clear_queue"]), true, "abort first continues the queued item in the same run");
assert.equal(await drive(["clear_queue", "abort"]), false, "clear first removes it");
});
test("N3: the fence lands in preflight, preflight errors, no run: failed, No run, uncertain", async () => {
const h = await setup();
try {
h.engine.arm("915");
h.engine.plan({ error: "auth check failed", at: "915" });
const id = await prompt(h);
await h.engine.waitPaused("915");
await receiptState(h.client, id, "dispatched");
const stopId = await interrupt(h);
await until(() => cmds(h.engine, "abort").length === 1, 3000, "abort");
h.engine.resume("915");
const r = await receiptState(h.client, id, "failed");
assert.equal(r.reasonCode, null);
assert.ok(h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && p.nativeError === "auth check failed"), "native error kept");
assert.equal(h.engine.runs.length, 0);
await assertUncertainStop(h, stopId, "no-run");
} finally {
await h.close();
}
});
test("N4: the ack arrives after the first abort and a run starts: clear and abort again; Interrupted", async () => {
const h = await setup();
try {
h.engine.arm("915");
h.engine.script([{ pause: "p1" }, { text: "never" }]);
h.engine.arm("p1");
const id = await prompt(h);
await h.engine.waitPaused("915");
const stopId = await interrupt(h);
await until(() => cmds(h.engine, "abort").length === 1, 3000, "first abort");
h.engine.resume("915");
await until(() => reconciled(h.ctrl, stopId), 8000, "reconciled");
assert.equal(cmds(h.engine, "abort").length, 2, "a second abort for the late run");
assert.equal(cmds(h.engine, "clear_queue").length, 3, "two interrupt clears and the post-settle clear");
const r = await receiptState(h.client, id, "failed");
assert.equal(r.reasonCode, "interrupted");
assert.equal(h.ctrl.stops.get(stopId).state, "turn-interrupted");
assert.equal(turnProofs(h.verifier).at(-1).turnState, "interrupted");
assert.equal(h.ctrl.binding.admission, "open");
} finally {
await h.close();
}
});
test("N5: an input handler takes the prompt: ack, no run, delivery-unknown handled-without-run", async () => {
const h = await setup();
try {
h.engine.plan({ handled: true });
const id = await prompt(h);
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, HANDLED_WITHOUT_RUN);
assert.ok(!h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && ["working", "finished"].includes(p.receipt.state)));
await tick(100);
assert.equal(cmds(h.engine, "prompt").length, 1, "nothing is resent");
assert.equal(h.engine.runs.length, 0);
assert.equal(h.ctrl.binding.admission, "open");
} finally {
await h.close();
}
});
test("N6: an extension queues between clear_queue and abort: O5 and O6, Unknown", async () => {
const h = await setup();
try {
const id = await heldRun(h);
h.engine.arm("abort");
const stopId = await interrupt(h);
await h.engine.waitPaused("abort");
h.engine.queue("followUp", "between");
h.engine.resume("abort");
const s = await stopSettled(h, stopId);
assert.ok(signals(h.ctrl).includes("O5"));
assert.ok(signals(h.ctrl).includes("O6"), "abort continued the queued item in the same run");
assert.equal(h.engine.runs.length, 1);
assert.equal(s.state, "uncertain");
assert.equal(h.client.receipt(id).state, "working");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
assert.equal(outcomeUnknownPush(h.client, id).reason, RUN_OVERLAP);
await assertUncertainStop(h, stopId, "unknown");
} finally {
await h.close();
}
});
test("N7: clear_queue times out: no abort, nativeQueue unknown, force stop still ends it", async () => {
const h = await setup();
try {
const id = await heldRun(h);
h.engine.dropResponse("clear_queue");
const stopId = await interrupt(h);
const s = await stopSettled(h, stopId);
assert.equal(cmds(h.engine, "abort").length, 0, "abort would run whatever is queued");
assert.equal(s.nativeQueue, "unknown");
assert.equal(s.state, "uncertain");
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.client.receipt(id).state, "working");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
assert.equal(outcomeUnknownPush(h.client, id).reason, TRANSPORT_UNKNOWN);
assert.equal(lastStopEvidence(h.ctrl, stopId).outcome, "unknown");
const fs = await h.client.confirmed("force-stop");
assert.equal(fs.outcome, "force-stop-fenced");
await until(() => h.ctrl.binding.state === "stopped", 4000, "force stop");
} finally {
await h.close();
}
});
test("N7: clear_queue answers an error: no abort, nativeQueue unknown, the link not poisoned", async () => {
// An error response is the one clear failure that leaves the link
// unpoisoned, so nothing but rule 2 keeps the abort off the pipe.
const h = await setup();
try {
const id = await heldRun(h);
h.engine.failResponse("clear_queue");
const stopId = await interrupt(h);
const s = await stopSettled(h, stopId);
assert.equal(h.ctrl.exec.link.poisoned, null);
assert.equal(cmds(h.engine, "abort").length, 0, "abort would run whatever is queued");
assert.equal(s.nativeQueue, "unknown");
assert.equal(s.state, "uncertain");
assert.equal(h.client.receipt(id).state, "working");
} finally {
await h.close();
}
});
test("N8: an extension prompt starts a run during Mosaic preflight; the losing settle is O3", async () => {
const h = await setup();
try {
h.engine.arm("915");
const id = await prompt(h);
await h.engine.waitPaused("915");
// The order: Mosaic ack, the other run's agent_start and user
// message_start, then the losing Mosaic prompt's agent_settled.
h.engine.arm("run-start");
h.engine.arm("state");
h.engine.resume("915");
await h.engine.waitPaused("run-start");
await h.engine.waitPaused("state");
h.engine.arm("ext-hold");
void h.engine.extensionPrompt({ text: "from an extension", steps: [{ pause: "ext-hold" }, { text: "ext done" }] });
await h.engine.waitPaused("ext-hold");
await receiptState(h.client, id, "working");
h.engine.resume("run-start");
await until(() => signals(h.ctrl).includes("O3"), 2000, "O3");
h.engine.resume("state");
h.engine.resume("ext-hold");
await h.engine.waitIdle();
await tick(50);
const wire = h.engine.sent.map((l) => JSON.parse(l)).filter((v) => v.type !== "message_update");
const ack = wire.findIndex((v) => v.type === "response" && v.command === "prompt");
const start = wire.findIndex((v, i) => i > ack && v.type === "agent_start");
const user = wire.findIndex((v, i) => i > start && v.type === "message_start" && v.message.role === "user");
const settle = wire.findIndex((v, i) => i > user && v.type === "agent_settled");
assert.ok(ack >= 0 && ack < start && start < user && user < settle, "the pinned wire order");
assert.equal(h.client.receipt(id).state, "working", "never finished or failed");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
assert.equal(outcomeUnknownPush(h.client, id).reason, RUN_OVERLAP);
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.ctrl.binding.admission, "closed");
assert.equal(cmds(h.engine, "prompt").length, 1, "nothing resent");
} finally {
await h.close();
}
});
test("N9: a run that started before the fence and ends aborted: failed interrupted, Interrupted", async () => {
const h = await setup();
try {
const id = await heldRun(h);
const stopId = await interrupt(h);
await until(() => reconciled(h.ctrl, stopId), 6000, "reconciled");
const r = await receiptState(h.client, id, "failed");
assert.equal(r.reasonCode, "interrupted");
assert.ok(h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && p.stop === stopId), "evidence links the stop");
assert.equal(lastStopEvidence(h.ctrl, stopId).outcome, "interrupted");
assert.equal(turnProofs(h.verifier).at(-1).turnState, "interrupted");
} finally {
await h.close();
}
});
// Lead decision 34: only the controller's abort produces `aborted`. Mutant
// r2-D34 (turns.mjs drops the aborted-without-stop overlap) fails here.
test("N9: decision 34: a run that ends aborted with no stop in progress: aborted-without-stop, uncertain, outcome unknown", async () => {
const h = await setup();
try {
h.engine.script([{ text: "", stop: "aborted", errorMessage: "Request was aborted" }]);
const id = await prompt(h);
await until(() => signals(h.ctrl).includes("aborted-without-stop"), 4000, "the overlap");
assert.equal(h.ctrl.stops.size, 0);
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.ctrl.binding.admission, "closed");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
assert.equal(outcomeUnknownPush(h.client, id).reason, RUN_OVERLAP);
assert.notEqual(h.client.receipt(id).state, "failed");
assert.equal(turnProofs(h.verifier).length, 0);
} finally {
await h.close();
}
});
// n1: the fence alone doesn't link an `aborted`; an abort must have been
// written in the stop's chain. Mutant r2-n1 (link to any stop in progress)
// fails here.
test("N9: an aborted that lands after the fence but before any abort is written: aborted-without-stop, Unknown", async () => {
const h = await setup();
try {
const id = await heldRun(h, [{ text: "", stop: "aborted", errorMessage: "Request was aborted" }]);
h.engine.arm("clear");
const stopId = await interrupt(h);
await h.engine.waitPaused("clear");
h.engine.resume("p1");
await until(() => signals(h.ctrl).includes("aborted-without-stop"), 4000, "the overlap");
h.engine.resume("clear");
assert.equal(cmds(h.engine, "abort").length, 0, "no abort was written");
assert.notEqual(h.client.receipt(id).state, "failed");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
await assertUncertainStop(h, stopId, "unknown");
} finally {
await h.close();
}
});
// N10: the fake's own behaviour, against pinned Pi's (agent-session.js
// 821–949, rpc-mode.js 298–335). Mutant 28 (a fake that queues a Mosaic
// prompt while streaming) fails here.
function rawFake(opts = {}) {
const input = new PassThrough(), output = new PassThrough();
const fake = new FakePi({ input, output, ...opts });
const lines = [];
output.pipe(new PassThrough()).on("data", () => {});
const splitter = new LineSplitter((l) => lines.push(parseLine(l).value));
output.on("data", (c) => splitter.push(c));
let n = 0;
const send = (type, fields = {}) => {
const id = `r${++n}`;
input.write(encodeLine({ id, type, ...fields }));
return id;
};
const response = (id) => until(() => lines.some((v) => v.type === "response" && v.id === id)).then(() => lines.find((v) => v.type === "response" && v.id === id));
return { fake, lines, send, response };
}
test("N10: fake conformance", async () => {
// Acks before running; agent_settled from a finally; several agent_start …
// agent_end pairs in one run.
{
const f = rawFake();
f.fake.script([{ continue: true }, { text: "ok", stop: "stop" }]);
const id = f.send("prompt", { message: "a" });
await until(() => f.lines.some((v) => v.type === "agent_settled"));
const types = f.lines.map((v) => v.type);
const ack = f.lines.findIndex((v) => v.type === "response" && v.id === id);
assert.ok(ack >= 0 && ack < types.indexOf("agent_start"), "ack before the run");
assert.equal(types.filter((t) => t === "agent_start").length, 2);
assert.equal(types.filter((t) => t === "agent_settled").length, 1);
assert.equal(types.at(-1), "agent_settled");
}
// A failing run still settles (finally).
{
const f = rawFake();
f.fake.script([{ failBeforeUser: "boom" }]);
f.send("prompt", { message: "a" });
await until(() => f.lines.some((v) => v.type === "agent_settled"));
assert.ok(!f.lines.some((v) => v.type === "message_start" && v.message.role === "user"));
}
// A prompt while streaming throws (no streamingBehavior); it is never
// queued.
{
const f = rawFake();
f.fake.script([{ pause: "p1" }, { text: "x" }]);
f.fake.arm("p1");
f.send("prompt", { message: "first" });
await f.fake.waitPaused("p1");
const before = f.lines.length;
const r = await f.response(f.send("prompt", { message: "second" }));
assert.equal(r.success, false);
assert.equal(r.error, BUSY_ERROR);
assert.ok(!f.lines.slice(before).some((v) => v.type === "queue_update"), "no queue_update: the prompt was not queued");
assert.equal(f.fake.steering.length + f.fake.followUp.length + f.fake.agentQueue.length + f.fake.nextTurn.length, 0);
f.fake.resume("p1");
await f.fake.waitIdle();
assert.equal(f.fake.runs.length, 1);
assert.ok(!f.lines.some((v) => v.type === "message_start" && v.message.content?.[0]?.text === "second"), "the second prompt never ran");
}
// Pauses at 843, 895 and 915, between the line-860 check and the run start:
// no ack until all three pass.
{
const f = rawFake();
for (const p of ["843", "895", "915"]) f.fake.arm(p);
const id = f.send("prompt", { message: "a" });
for (const p of ["843", "895", "915"]) {
await f.fake.waitPaused(p);
assert.ok(!f.lines.some((v) => v.id === id), `no ack while paused at ${p}`);
f.fake.resume(p);
}
await f.response(id);
await f.fake.waitIdle();
}
// A colliding prompt (past line 860 before the other run started) is
// acked, its throw swallowed, it settles with no agent_start, and leaves
// isStreaming false while the other run goes on.
{
const f = rawFake();
f.fake.script([{ pause: "p1" }, { text: "x" }]);
f.fake.arm("p1");
f.fake.arm("run-start");
const a = f.send("prompt", { message: "A" });
await f.fake.waitPaused("run-start");
await f.response(a);
f.fake.arm("915");
const b = f.send("prompt", { message: "B" });
await f.fake.waitPaused("915");
f.fake.resume("run-start");
await f.fake.waitPaused("p1");
const mark = f.lines.length;
f.fake.resume("915");
const rb = await f.response(b);
assert.equal(rb.success, true, "the colliding prompt is acked");
await until(() => f.lines.slice(mark).some((v) => v.type === "agent_settled"));
assert.ok(!f.lines.slice(mark).some((v) => v.type === "agent_start"), "no agent_start for the loser");
const st = await f.response(f.send("get_state"));
assert.equal(st.data.isStreaming, false, "isStreaming false while the other run goes on");
assert.ok(f.fake.run, "the other run is still going");
f.fake.resume("p1");
await f.fake.waitIdle();
}
// clear_queue: an empty queue_update before its response; returns steer and
// follow-up only; drops agent-level custom messages; keeps nextTurn.
{
const f = rawFake();
f.fake.queue("steer", "s");
f.fake.queue("followUp", "f");
f.fake.queue("agent", "custom");
f.fake.queue("nextTurn", "nt");
const mark = f.lines.length;
const id = f.send("clear_queue");
const r = await f.response(id);
const window = f.lines.slice(mark);
const at = window.findIndex((v) => v.id === id);
assert.deepEqual(window[at - 1], { type: "queue_update", steering: [], followUp: [] });
assert.deepEqual(r.data, { steering: ["s"], followUp: ["f"] });
assert.equal(f.fake.agentQueue.length, 0, "agent-level custom messages are dropped, not returned");
assert.equal(f.fake.nextTurn.length, 1, "nextTurn survives the clear");
}
});
test("N11: the run fails before any user message_start: delivery-unknown ack-without-start, never failed", async () => {
const h = await setup();
try {
const before = sessionText(h.fx);
h.engine.script([{ failBeforeUser: "provider refused" }]);
const id = await prompt(h);
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, ACK_WITHOUT_START);
assert.ok(!h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && p.receipt.state === "failed"));
const added = sessionText(h.fx).slice(before.length).split("\n").filter(Boolean).map((l) => JSON.parse(l));
assert.ok(!added.some((e) => e.message?.role === "user"), "the session gains no user entry");
assert.equal(cmds(h.engine, "prompt").length, 1);
} finally {
await h.close();
}
// The same schedule with the failure line unparseable.
const h2 = await setup();
try {
h2.engine.script([{ failBeforeUser: "x", raw: "{not json\n" }]);
const id = await prompt(h2);
const r = await receiptState(h2.client, id, "delivery-unknown");
assert.equal(r.reasonCode, TRANSPORT_UNKNOWN);
assert.equal(h2.ctrl.binding.state, "uncertain");
} finally {
await h2.close();
}
});
test("N12: input that starts a run after the final empty clear is O1 and not part of the stop's proof", async () => {
const h = await setup();
try {
await heldRun(h);
const stopId = await interrupt(h);
await until(() => reconciled(h.ctrl, stopId), 6000, "reconciled");
const ev = lastStopEvidence(h.ctrl, stopId);
const lastClear = ev.clears.at(-1);
const proof = turnProofs(h.verifier).find((p) => p.stop === stopId);
assert.equal(proof.observedAt, lastClear.at, "the stop records the final clear's observedAt");
void h.engine.extensionPrompt({ text: "late" });
await until(() => signals(h.ctrl).includes("O1"), 2000, "O1");
const o1 = h.ctrl.evidence.overlaps.find((o) => o.signal === "O1");
assert.equal(o1.why, "no slot held");
assert.ok(o1.idx > lastClear.idx, "the run came after the proof's observation");
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.ctrl.binding.admission, "closed");
assert.equal(h.ctrl.stops.get(stopId).state, "turn-interrupted", "the stopped turn's proof is unchanged");
} finally {
await h.close();
}
});
test("N13: agent_start with no slot held is O1; a later prompt refuses with zero engine bytes", async () => {
const h = await setup();
try {
h.engine.emit({ type: "agent_start" });
await until(() => signals(h.ctrl).includes("O1"), 2000, "O1");
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.ctrl.binding.admission, "closed");
assert.ok(h.ctrl.evidence.uncertain.some((u) => u.reason === RUN_OVERLAP));
const before = h.engine.bytes().length;
const r = await h.client.prompt("after");
assert.equal(r.refusal, "fenced");
await tick(50);
assert.equal(h.engine.bytes().length, before);
} finally {
await h.close();
}
});
test("N14: the run completes while clear_queue is in flight: finished, Completed first, uncertain", async () => {
const h = await setup();
try {
const id = await heldRun(h, [{ text: "done", stop: "stop" }]);
h.engine.arm("clear");
const stopId = await interrupt(h);
await h.engine.waitPaused("clear");
h.engine.resume("p1");
await receiptState(h.client, id, "finished");
h.engine.resume("clear");
await stopSettled(h, stopId);
assert.equal(cmds(h.engine, "abort").length, 1, "the abort reached an idle engine");
assert.ok(lastStopEvidence(h.ctrl, stopId).clears.every((c) => c.empty));
assert.equal(h.client.receipt(id).state, "finished", "never relabelled interrupted");
assert.equal(h.client.receipt(id).reasonCode, null);
await assertUncertainStop(h, stopId, "completed-first");
} finally {
await h.close();
}
});
test("N14: the run completes after the abort is written, before Pi applies it: finished, never relabelled", async () => {
const h = await setup();
try {
const id = await heldRun(h, [{ text: "done", stop: "stop" }]);
h.engine.arm("abort");
const stopId = await interrupt(h);
await h.engine.waitPaused("abort");
h.engine.resume("p1");
await receiptState(h.client, id, "finished");
h.engine.resume("abort");
await stopSettled(h, stopId);
assert.equal(cmds(h.engine, "abort").length, 1);
assert.equal(h.client.receipt(id).state, "finished", "a run that ended stop is never relabelled interrupted");
assert.equal(h.client.receipt(id).reasonCode, null);
await assertUncertainStop(h, stopId, "completed-first");
} finally {
await h.close();
}
});
test("N15: the fence lands in preflight, then an input handler takes it: handled-without-run, No run", async () => {
const h = await setup();
try {
h.engine.arm("843");
h.engine.plan({ handled: true });
const id = await prompt(h);
await h.engine.waitPaused("843");
const stopId = await interrupt(h);
await until(() => cmds(h.engine, "abort").length === 1, 3000, "abort");
h.engine.resume("843");
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, HANDLED_WITHOUT_RUN);
await stopSettled(h, stopId);
assert.equal(cmds(h.engine, "abort").length, 1, "no second abort");
assert.equal(cmds(h.engine, "clear_queue").length, 1, "no second clear");
await assertUncertainStop(h, stopId, "no-run");
} finally {
await h.close();
}
});
test("N16: Interrupt with no slot and no run refuses no-turn: no stop, no bytes, admission open", async () => {
const h = await setup();
try {
const before = h.engine.bytes().length;
const r = await h.client.interrupt();
assert.equal(r.refusal, NO_TURN);
assert.deepEqual(r.effect, { dispatchRefused: null });
assert.equal(h.ctrl.stops.size, 0);
assert.equal(h.ctrl.binding.stop, null);
assert.equal(h.engine.bytes().length, before);
assert.equal(h.ctrl.binding.admission, "open");
const p = await h.client.prompt("next");
assert.equal(p.outcome, "admitted");
await receiptState(h.client, p.receipt.id, "finished");
} finally {
await h.close();
}
});
test("N17: the run fails on its own during the exchange: failed, Failed on its own", async () => {
const h = await setup();
try {
const id = await heldRun(h, [{ text: "", stop: "error", errorMessage: "model error" }]);
h.engine.arm("clear");
const stopId = await interrupt(h);
await h.engine.waitPaused("clear");
h.engine.resume("p1");
const r = await receiptState(h.client, id, "failed");
assert.equal(r.reasonCode, null);
h.engine.resume("clear");
await assertUncertainStop(h, stopId, "failed-on-its-own");
} finally {
await h.close();
}
});
test("N18: no final assistant message_end, or a lost line: working stays working; before working, transport-unknown", async () => {
// No final assistant message during the exchange.
{
const h = await setup();
try {
const id = await heldRun(h, [{ text: "partial", final: false }]);
h.engine.arm("clear");
const stopId = await interrupt(h);
await h.engine.waitPaused("clear");
h.engine.resume("p1");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
h.engine.resume("clear");
const s = await stopSettled(h, stopId);
assert.equal(s.state, "uncertain");
assert.equal(h.client.receipt(id).state, "working");
assert.equal(lastStopEvidence(h.ctrl, stopId).outcome, "unknown");
assert.equal(h.ctrl.binding.state, "uncertain");
} finally {
await h.close();
}
}
// A lost (unparseable) line after working never moves it back.
{
const h = await setup();
try {
const id = await heldRun(h, [{ raw: "{lost\n" }, { text: "x" }]);
h.engine.resume("p1");
await until(() => outcomeUnknownPush(h.client, id), 2000, "outcome unknown");
await tick(50);
assert.equal(h.client.receipt(id).state, "working");
assert.equal(outcomeUnknownPush(h.client, id).reason, TRANSPORT_UNKNOWN);
assert.ok(!h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && p.receipt.state === "delivery-unknown"));
} finally {
await h.close();
}
}
// A lost line before working.
{
const h = await setup();
try {
h.engine.arm("run-start");
h.engine.arm("state");
const id = await prompt(h);
await h.engine.waitPaused("run-start");
await receiptState(h.client, id, "acknowledged");
h.engine.raw("{lost\n");
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, TRANSPORT_UNKNOWN);
h.engine.resume("state");
h.engine.resume("run-start");
} finally {
await h.close();
}
}
});
test("N19: a losing extension prompt settles inside the Mosaic run before its user message: O3, run-overlap", async () => {
const h = await setup();
try {
h.engine.arm("ext:915");
void h.engine.extensionPrompt({ text: "loser" });
await h.engine.waitPaused("ext:915");
h.engine.arm("after-start");
const id = await prompt(h);
await h.engine.waitPaused("after-start");
await receiptState(h.client, id, "acknowledged");
h.engine.resume("ext:915");
await until(() => signals(h.ctrl).includes("O3"), 2000, "O3");
h.engine.resume("after-start");
await h.engine.waitIdle();
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, RUN_OVERLAP, "never failed ack-without-start");
await tick(50);
assert.equal(h.client.receipt(id).state, "delivery-unknown");
assert.equal(h.ctrl.binding.state, "uncertain");
} finally {
await h.close();
}
// Before the Mosaic agent_start: O4.
const h2 = await setup();
try {
h2.engine.arm("run-start");
h2.engine.arm("state");
const id = await prompt(h2);
await h2.engine.waitPaused("run-start");
await h2.engine.waitPaused("state");
h2.engine.emit({ type: "agent_settled" });
await until(() => signals(h2.ctrl).includes("O4"), 2000, "O4");
const r = await receiptState(h2.client, id, "delivery-unknown");
assert.equal(r.reasonCode, RUN_OVERLAP);
h2.engine.resume("state");
h2.engine.resume("run-start");
await h2.engine.waitIdle();
await tick(50);
assert.equal(h2.client.receipt(id).state, "delivery-unknown");
assert.equal(h2.ctrl.binding.state, "uncertain");
} finally {
await h2.close();
}
});
test("N20: an extension triggerTurn during Mosaic preflight starts first; while streaming it queues with no signal", async () => {
const h = await setup();
try {
h.engine.arm("915");
const id = await prompt(h);
await h.engine.waitPaused("915");
h.engine.arm("ext-hold");
void h.engine.extensionPrompt({ text: "timer", preflight: false, steps: [{ pause: "ext-hold" }, { text: "t" }] });
await h.engine.waitPaused("ext-hold");
await until(() => signals(h.ctrl).includes("O1"), 2000, "O1");
h.engine.resume("915");
h.engine.resume("ext-hold");
await h.engine.waitIdle();
const r = await receiptState(h.client, id, "delivery-unknown");
assert.equal(r.reasonCode, RUN_OVERLAP);
await tick(50);
assert.ok(!h.client.pushes.some((p) => p.kind === "receipt" && p.receipt.id === id && ["finished", "failed"].includes(p.receipt.state)), "never finished or failed");
assert.equal(h.ctrl.binding.state, "uncertain");
} finally {
await h.close();
}
// While the Mosaic run streams: queued straight into the agent, no signal.
const h2 = await setup();
try {
const id = await heldRun(h2, [{ text: "done", stop: "stop" }]);
const q = await h2.engine.extensionPrompt({ text: "timer", preflight: false });
assert.deepEqual(q, { queued: true });
h2.engine.resume("p1");
await receiptState(h2.client, id, "finished");
assert.deepEqual(signals(h2.ctrl), [], "Limit 7: no signal fires");
assert.ok(h2.engine.sent.some((l) => l.includes('"customType":"fake"') && l.includes("timer")), "the custom message ran inside the Mosaic run");
} finally {
await h2.close();
}
});
test("N21: a losing settle after the receipt settled finished is O2; the receipt stays finished", async () => {
const h = await setup();
try {
const id = await prompt(h);
await receiptState(h.client, id, "finished");
h.engine.emit({ type: "agent_settled" });
await until(() => signals(h.ctrl).includes("O2"), 2000, "O2");
assert.equal(h.client.receipt(id).state, "finished");
const o2 = h.ctrl.evidence.overlaps.find((o) => o.signal === "O2");
assert.equal(o2.request, h.ctrl.receipts.get(id).request, "evidence records the overlap against it");
assert.equal(h.ctrl.binding.state, "uncertain");
assert.equal(h.ctrl.binding.admission, "closed");
assert.ok(h.ctrl.evidence.uncertain.some((u) => u.reason === RUN_OVERLAP));
} finally {
await h.close();
}
});
test("N22: an agent-level custom message is dropped by the clear with no signal; evidence names the seal", async () => {
const h = await setup();
try {
await heldRun(h);
h.engine.queue("agent", "custom from an extension");
const stopId = await interrupt(h);
await until(() => reconciled(h.ctrl, stopId), 6000, "reconciled");
assert.deepEqual(signals(h.ctrl), [], "Limit 7: no signal fires");
assert.equal(h.engine.agentQueue.length, 0);
assert.ok(!h.engine.sent.some((l) => l.includes("custom from an extension")), "the item is gone, never run");
const ev = lastStopEvidence(h.ctrl, stopId);
assert.equal(ev.agentLevelQueues, "unobservable");
assert.match(ev.queueBasis, /^seal:/);
assert.ok(!JSON.stringify(ev).includes("nothing removed"));
} finally {
await h.close();
}
});
test("N23: a nextTurn message survives clear and abort and attaches to the next prompt, with no signal", async () => {
const h = await setup();
try {
await heldRun(h);
h.engine.queue("nextTurn", "rides along");
const stopId = await interrupt(h);
await until(() => reconciled(h.ctrl, stopId), 6000, "reconciled");
assert.equal(h.engine.nextTurn.length, 1, "it survived");
const p = await h.client.prompt("next");
await receiptState(h.client, p.receipt.id, "finished");
assert.equal(h.engine.nextTurn.length, 0);
assert.ok(h.engine.sent.some((l) => l.includes("rides along") && l.includes('"message_start"')), "it attached to the next prompt");
assert.deepEqual(signals(h.ctrl), [], "Limit 7: no signal fires");
} finally {
await h.close();
}
});
test("N24: the seal is an allow-list: --extension, a missing --no-* flag, a second --mode or --session, a session or output flag, or a stray word refuses unsealed-engine; no engine starts", async () => {
const fx = fixture();
const outside = join(fx.base, "outside", "proj", ".pi", "state", "other", "sessions", "s1.jsonl");
for (const extraArgs of [
["--extension", "./ext.ts"], ["--extension=npm:some-ext"], ["-e", "git:github.com/x/y"], ["--extension", "npm:@scope/pkg"],
// Pi keeps the last --session and the last --mode (cli/args.js).
["--session", outside], ["--mode", "json"], ["--mode", "text"], ["--mode", "rpc"],
["--no-session"], ["--session-dir", fx.base], ["--session-id", "x"], ["--fork", fx.sessionFile], ["--continue"], ["-c"], ["--resume"],
["--print"], ["-p"], ["--export", join(fx.base, "out.html")], ["--prompt-template", "x"], ["--approve"], ["--no-extensions"],
// A bare word is a prompt, an @ word a file.
["hello"], ["@notes.md"],
// Allowed options: once each, one plain value.
["--model", "a", "--model", "b"], ["--model"], ["--model", "--mode"], ["--thinking", "-x"], ["--provider", "@p"], ["--model=other"],
]) {
const launcher = new FakeLauncher();
assert.throws(
() => new Controller({ fixtureRoot: fx.base, claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat, launcher, units: noUnits, timeouts: FAST, engine: { extraArgs } }),
(e) => e.code === UNSEALED_ENGINE,
JSON.stringify(extraArgs),
);
assert.equal(launcher.launches.length, 0);
}
{
const launcher = new FakeLauncher();
assert.throws(
() => new Controller({ fixtureRoot: fx.base, claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat, launcher, units: noUnits, engine: { preArgs: ["cli.js", "--extension", "x"] } }),
(e) => e.code === UNSEALED_ENGINE,
);
}
for (const engine of [{ extraArgs: "--model x" }, { extraArgs: 3 }, { preArgs: "cli.js" }]) {
assert.throws(
() => new Controller({ fixtureRoot: fx.base, claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat, launcher: new FakeLauncher(), units: noUnits, engine }),
(e) => e.code === UNSEALED_ENGINE,
JSON.stringify(engine),
);
}
// The controller always builds the prefix itself, so a missing flag, a
// reordered prefix or a relative session path is reachable only through
// checkSeal, which bind runs on the argv.
for (const flag of SEAL_FLAGS) {
const args = ["--mode", "rpc", ...SEAL_FLAGS.filter((f) => f !== flag), "--session", fx.sessionFile];
assert.throws(() => checkSeal(args), (e) => e.code === UNSEALED_ENGINE, flag);
}
for (const args of [
["--mode", "rpc", "--session", fx.sessionFile, ...SEAL_FLAGS],
["--mode", "json", ...SEAL_FLAGS, "--session", fx.sessionFile],
["--mode", "rpc", ...SEAL_FLAGS, "--session", "s1.jsonl"],
["--mode", "rpc", ...SEAL_FLAGS, "--session"],
]) {
assert.throws(() => checkSeal(args), (e) => e.code === UNSEALED_ENGINE, JSON.stringify(args));
}
assert.equal(checkSeal(["--mode", "rpc", ...SEAL_FLAGS, "--session", fx.sessionFile, "--model", "m", "--provider", "p", "--thinking", "off"]), true);
// The argv a real bind launches carries all three and no --extension.
const h = await started({ fx: fixture() });
try {
const args = h.launcher.launches[0].args;
for (const flag of SEAL_FLAGS) assert.ok(args.includes(flag), flag);
assert.ok(!args.some((a) => a === "-e" || a.startsWith("--extension")));
assert.deepEqual(h.engine.argv.slice(0, 2), ["--mode", "rpc"]);
} finally {
await h.close();
}
});