Dewey's round 3 candidate, manifest
agents/dewey/work/queue-40/candidate-manifest-r3.sha256 (d0aa0ded,
27 files, checked OK in the canonical tree).
- WebUI inbox, tasks, agents and trail views, read-only over /api/bus.
The README says the bus proof ends at the Console process.
- CHAT-03 seal: the engine command is fixed, the engine environment is
explicit, SEAL_FLAGS has --no-approve, escalating is cleared on throw.
- Terminal input typed after Ctrl-T or Ctrl-O is held. Only the run whose
own parse set held drains it (T1), and #run catches errors per action.
- DEFERRED keeps N2 and moves F2 to done, citing T1.
Reviews: Filbert approve (comment 27011, rev 260), Darkwing approve
(27013, rev 264). Landing gate on 8cad7722 plus the candidate: webui 22,
conversation 161, control-board 124, every scripts/test-*.sh green,
test-task 98/0. Mutant Mr survives; its flows test is the first
follow-up row.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
1044 lines
47 KiB
JavaScript
1044 lines
47 KiB
JavaScript
// 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, TEST_ENGINE, REPO_ROOT, 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 { ENGINE_ENV, PI_BIN, 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 every seal flag 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();
|
||
}
|
||
});
|
||
|
||
test("N24b: the seal covers the engine command and environment: config can't name either, the env is built from names, and a mutated command is refused at bind", async () => {
|
||
const fx = fixture();
|
||
const make = (extra) => new Controller({ fixtureRoot: fx.base, claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat, launcher: new FakeLauncher(), units: noUnits, timeouts: FAST, ...extra });
|
||
// The plain engine option takes extraArgs, cwd and envKeys; the command,
|
||
// pre-arguments and environment are the controller's (I3, #1507).
|
||
for (const engine of [{ command: "/bin/sh" }, { preArgs: [join(REPO_ROOT, PI_BIN)] }, { command: process.execPath, preArgs: [join(REPO_ROOT, PI_BIN), "--extension", "x"] }, { env: {} }, { env: process.env }, null, [], "pi"]) {
|
||
assert.throws(() => make({ engine }), (e) => e.code === UNSEALED_ENGINE, JSON.stringify(engine));
|
||
}
|
||
// A configuration is JSON, which can't carry the symbol-keyed test engine.
|
||
const parsed = JSON.parse(JSON.stringify({ engine: { cwd: fx.proj }, [TEST_ENGINE]: { command: "/bin/sh", preArgs: [], env: {} } }));
|
||
assert.equal(Object.getOwnPropertySymbols(parsed).length, 0);
|
||
// envKeys may name provider credentials only.
|
||
for (const envKeys of [["NODE_OPTIONS"], ["LD_PRELOAD"], ["PI_PACKAGE_DIR"], ["BASH_ENV"], ["ZAI_API_KEY", "PATH"], ["zai_api_key"], ["_API_KEY"], [3], "ZAI_API_KEY"]) {
|
||
assert.throws(() => make({ engine: { envKeys } }), (e) => e.code === UNSEALED_ENGINE, JSON.stringify(envKeys));
|
||
}
|
||
// The test engine is still sealed on its arguments.
|
||
assert.throws(() => make({ [TEST_ENGINE]: { command: process.execPath, preArgs: ["fake.mjs", "-e", "x"], env: {} } }), (e) => e.code === UNSEALED_ENGINE);
|
||
assert.throws(() => make({ [TEST_ENGINE]: { command: "", preArgs: [], env: {} } }), (e) => e.code === UNSEALED_ENGINE);
|
||
|
||
// The launch: this Node, the pinned bin, and an environment of named keys
|
||
// only, whatever the controller's own environment holds.
|
||
const planted = { NODE_OPTIONS: "--require /tmp/x.cjs", LD_PRELOAD: "/tmp/x.so", PI_PACKAGE_DIR: "/tmp/pkg", ZAI_API_KEY: "zai-test-value", OTHER_API_KEY: "other-test-value" };
|
||
const saved = Object.fromEntries(Object.keys(planted).map((k) => [k, process.env[k]]));
|
||
Object.assign(process.env, planted);
|
||
let h;
|
||
try {
|
||
h = await started({ fx: fixture(), engine: { envKeys: ["ZAI_API_KEY"] } });
|
||
} finally {
|
||
for (const [k, v] of Object.entries(saved)) if (v === undefined) delete process.env[k]; else process.env[k] = v;
|
||
}
|
||
try {
|
||
const l = h.launcher.launches[0];
|
||
assert.equal(l.command, process.execPath);
|
||
assert.equal(l.args[0], join(REPO_ROOT, PI_BIN));
|
||
assert.equal(l.env.ZAI_API_KEY, "zai-test-value");
|
||
for (const k of Object.keys(l.env)) assert.ok(ENGINE_ENV.includes(k) || k === "ZAI_API_KEY", k);
|
||
for (const k of ["NODE_OPTIONS", "LD_PRELOAD", "PI_PACKAGE_DIR", "OTHER_API_KEY"]) assert.equal(k in l.env, false, k);
|
||
if (process.env.PATH) assert.equal(l.env.PATH, process.env.PATH);
|
||
} finally {
|
||
await h.close();
|
||
}
|
||
|
||
// A command changed after construction is refused at bind; nothing launches.
|
||
const fx2 = fixture();
|
||
const launcher = new FakeLauncher();
|
||
const ctrl = new Controller({ fixtureRoot: fx2.base, claimRoot: fx2.claimRoot, socketDir: fx2.socketDir, sessionFile: fx2.sessionFile, seat: fx2.seat, launcher, units: noUnits, timeouts: FAST });
|
||
// Closed in finally: a start that wrongly succeeds holds the socket open.
|
||
try {
|
||
ctrl.engine.command = "/bin/sh";
|
||
await assert.rejects(ctrl.start(), (e) => e.code === UNSEALED_ENGINE);
|
||
assert.equal(launcher.launches.length, 0);
|
||
ctrl.engine.command = process.execPath;
|
||
ctrl.engine.preArgs = [join(fx2.base, "cli.js")];
|
||
await assert.rejects(ctrl.start(), (e) => e.code === UNSEALED_ENGINE);
|
||
assert.equal(launcher.launches.length, 0);
|
||
} finally {
|
||
await ctrl.close().catch(() => {});
|
||
}
|
||
});
|