import { test } from "node:test"; import assert from "node:assert/strict"; import { readFileSync } from "node:fs"; import { join } from "node:path"; import { createEngine, buildPiArgs, PI_FIXED_ARGS, TOOLS_EXTENSION, READONLY_TOOLS_EXTENSION, assistantText } from "../src/engine-pi.mjs"; import { existsSync } from "node:fs"; import { EventEmitter } from "node:events"; import { PassThrough } from "node:stream"; import { makeRoot } from "./helpers.mjs"; const fakePi = join(import.meta.dirname, "fake-pi.mjs"); // pi writes agent_settled after agent_end, at times in the next stdout // chunk, and the fake mirrors a command to its log only once it has read // it. Both are a few milliseconds; wait for them instead of racing them. async function until(check, ms = 1000) { for (let i = 0; i < ms / 10; i += 1) { if (check()) return true; await new Promise((res) => setTimeout(res, 10)); } return check(); } const idle = (engine) => until(() => !engine.busy); function start(root, extra = {}) { const logPath = join(root, "commands.jsonl"); const logs = []; const engine = createEngine({ command: process.execPath, args: [fakePi], cwd: root, env: { ...process.env, FAKE_PI_LOG: logPath }, log: (m) => logs.push(m), ...extra, }); engine.start(); // The fake creates its log on the first command; until then there are none. const commands = () => (existsSync(logPath) ? readFileSync(logPath, "utf8").trim().split("\n").filter(Boolean).map((l) => JSON.parse(l)) : []); return { engine, logs, commands }; } // Every test stops its engine in finally: a fake pi left running after a // failed assertion keeps the test file from exiting. async function withEngine(extra, body) { const started = start(makeRoot(), extra); try { await body(started); } finally { await started.engine.stop(); } } test("engine: buildPiArgs carries the fixed flags, engine settings, session dir and prompt file", () => { const args = buildPiArgs({ provider: "zai", model: "glm-5.3", thinking: "high", sessionDir: "/s", appendSystemPromptFile: "/p.md", continueSession: true }); for (const f of PI_FIXED_ARGS) assert.ok(args.includes(f), f); assert.ok(args.includes("--no-tools") && args.includes("--offline")); assert.deepEqual(args.slice(-9), ["--provider", "zai", "--model", "glm-5.3", "--thinking", "high", "--session-dir", "/s", "--append-system-prompt", "/p.md", "--continue"].slice(-9)); assert.ok(!buildPiArgs({ provider: "p", model: "m", thinking: "off", sessionDir: "/s", appendSystemPromptFile: "/p", continueSession: false }).includes("--continue")); assert.equal(assistantText({ content: [{ type: "thinking", thinking: "x" }, { type: "text", text: " a " }, { type: "text", text: "b" }] }), "a b".replace(" ", " ")); assert.ok(!args.includes("--no-builtin-tools") && !args.includes("--extension"), "no extension without tools"); }); test("engine: with tools, buildPiArgs turns pi's own tools off, loads the extension explicitly and allowlists exactly our three", () => { const tools = { roots: [{ name: "docs", path: "/r" }], maxFileBytes: 4096, maxCallsPerTurn: 8 }; const args = buildPiArgs({ provider: "p", model: "m", thinking: "off", sessionDir: "/s", appendSystemPromptFile: "/p", continueSession: false, tools }); assert.ok(!args.includes("--no-tools"), "--no-tools would hide the extension's tools too"); assert.ok(args.includes("--no-extensions"), "discovery stays off; only the explicit path loads"); assert.ok(args.includes("--no-builtin-tools")); assert.equal(args[args.indexOf("--extension") + 1], TOOLS_EXTENSION); assert.equal(READONLY_TOOLS_EXTENSION, TOOLS_EXTENSION); assert.equal(args[args.indexOf("--tools") + 1], "list_dir,read_file,search"); assert.ok(existsSync(TOOLS_EXTENSION), TOOLS_EXTENSION); assert.ok(TOOLS_EXTENSION.endsWith("/packages/discord/extension/tools.mjs")); const rw = buildPiArgs({ provider: "p", model: "m", thinking: "off", sessionDir: "/s", appendSystemPromptFile: "/p", continueSession: false, tools: { ...tools, roots: [{ name: "docs", path: "/r" }, { name: "vault", path: "/v", write: true }] } }); assert.equal(rw[rw.indexOf("--tools") + 1], "list_dir,read_file,search,write_file,edit_file", "a writable root adds exactly the two write tools"); }); test("engine: a run with tool turns settles once, on the answer, with every tool call in the result", () => withEngine({}, async ({ engine }) => { const r = await engine.prompt("tools 3"); assert.equal(r.text, "read 3 file(s)"); assert.equal(r.turns, 2); assert.equal(r.tools.length, 3); assert.deepEqual(r.tools[0], { name: "read_file", root: "docs", path: "f1.md", ok: true, reason: null, bytes: 9, ms: 2 }); assert.equal(r.tools[2].ok, false); assert.match(r.tools[2].reason, /budget/); const plain = await engine.prompt("hello"); assert.equal(plain.text, "echo: hello"); assert.deepEqual(plain.tools, []); assert.equal(plain.turns, 1); await idle(engine); assert.equal(engine.busy, false); })); test("engine: a run that ends on a tool-only turn fails the prompt as empty; a retried run settles on the real end", () => withEngine({}, async ({ engine }) => { const r = await engine.prompt("toolonly"); assert.equal(r.text, "", "no text: the connector turns this into engine-empty"); assert.equal(r.tools.length, 1); const again = await engine.prompt("retry"); assert.equal(again.text, "after retry"); })); test("engine: one prompt, one turn, text and usage come back", () => withEngine({}, async ({ engine }) => { const r = await engine.prompt("hello"); assert.equal(r.text, "echo: hello"); assert.deepEqual(r.usage, { input: 3, output: 2 }); await idle(engine); assert.equal(engine.busy, false); })); test("engine: a prompt while streaming is held until pi settles, then sent as its own run, and answered in order", () => withEngine({}, async ({ engine, commands }) => { const first = engine.prompt("slow 300"); assert.equal(engine.busy, true); const second = engine.prompt("second"); assert.equal(engine.pendingCount, 2); const prompted = () => commands().filter((c) => c.type === "prompt"); assert.ok(await until(() => prompted().length > 0), "the first prompt reached pi"); assert.equal(prompted().length, 1, "the second prompt is not sent while pi is busy"); // The fake refuses a prompt without streamingBehavior while it runs one, so // an answered second prompt also proves it was not sent early. const [r1, r2] = await Promise.all([first, second]); assert.equal(r1.text, "slow reply"); assert.equal(r2.text, "echo: second"); const prompts = prompted(); assert.equal(prompts.length, 2); // Never a pi follow-up: pi would fold it into the first run and close both // answers with one agent_end (the live loss of 2026-09-17). assert.equal(prompts[0].streamingBehavior, undefined); assert.equal(prompts[1].streamingBehavior, undefined); await idle(engine); assert.equal(engine.busy, false); })); test("engine: a held prompt that times out before pi settles fails on its own and is never sent", () => withEngine({}, async ({ engine, commands }) => { const first = engine.prompt("slow 200"); await new Promise((r) => setTimeout(r, 20)); await assert.rejects(engine.prompt("late one", { timeoutMs: 50 }), (e) => e.details.code === "timeout" && /waiting for the engine/.test(e.message)); const r1 = await first; assert.equal(r1.text, "slow reply"); await idle(engine); assert.deepEqual(commands().filter((c) => c.type === "prompt").map((c) => c.message), ["slow 200"]); assert.deepEqual(commands().filter((c) => c.type === "abort"), [], "a held turn is not aborted; pi never had it"); })); test("engine: timeout sends abort and fails only that turn; the process stays", () => withEngine({}, async ({ engine, commands, logs }) => { await assert.rejects(engine.prompt("slow 5000", { timeoutMs: 100 }), (err) => err.details.code === "timeout"); assert.ok(await until(() => commands().some((c) => c.type === "abort")), "abort reached pi"); assert.ok(logs.some((l) => /timed out/.test(l))); const r = await engine.prompt("again"); assert.equal(r.text, "echo: again"); })); test("engine: tool events from a run that outlived its timeout never land in the next prompt's record", () => withEngine({}, async ({ engine }) => { await assert.rejects(engine.prompt("late 200", { timeoutMs: 40 }), (err) => err.details.code === "timeout"); const r = await engine.prompt("after late"); assert.equal(r.text, "echo: after late"); assert.deepEqual(r.tools, [], "the dead run's read is not this prompt's evidence"); assert.equal(r.turns, 1, "the dead run's turns are not counted here"); })); // The turn timer is fired by hand, before the engine has read any event from // pi, so the timed-out run is still pi's and state.busy is still false when // the next prompt arrives. Under load a real timer does the same. const TURN_MS = 60000; const manualTurnTimer = (fire) => ({ setTimeoutImpl: (fn, ms) => (ms === TURN_MS ? fire.push(fn) : setTimeout(fn, ms)), clearTimeoutImpl: (id) => { if (typeof id !== "number") clearTimeout(id); }, }); test("engine: a prompt after a turn that timed out before its agent_start waits for pi to settle instead of being refused", () => { const fire = []; return withEngine(manualTurnTimer(fire), async ({ engine, commands }) => { const late = engine.prompt("late 100", { timeoutMs: TURN_MS }); fire.shift()(); assert.equal(engine.busy, true, "pi is still running the prompt that timed out"); const next = engine.prompt("after late", { timeoutMs: 5000 }); assert.equal(engine.pendingCount, 1, "only the new prompt is live"); await assert.rejects(late, (err) => err.details.code === "timeout"); const r = await next; assert.equal(r.text, "echo: after late"); assert.deepEqual(r.tools, [], "the dead run's read is not this prompt's evidence"); assert.equal(r.turns, 1); assert.deepEqual(commands().map((c) => (c.type === "prompt" ? c.message : c.type)), ["late 100", "abort", "after late"]); await idle(engine); assert.equal(engine.busy, false); }); }); // "mute" is accepted and never run, so no agent_start, agent_end or settle // ever comes for it. Unbounded, it would hold every later prompt. test("engine: when pi has not started a timed-out turn by the end of the abort grace, the engine stops pi and fails held prompts", async () => { let exited = null; await withEngine({ abortGraceMs: 150, onExit: (e) => (exited = e) }, async ({ engine, commands, logs }) => { await assert.rejects(engine.prompt("mute", { timeoutMs: 50 }), (err) => err.details.code === "timeout"); assert.equal(engine.busy, true, "pi might still be running it"); const started = Date.now(); await assert.rejects(engine.prompt("after mute", { timeoutMs: 5000 }), (err) => err.details.code === "engine-wedged"); assert.ok(Date.now() - started >= 100, "held for the grace, not failed at once"); await assert.rejects(engine.prompt("later"), (err) => err.details.code === "engine-down"); assert.ok(await until(() => exited !== null), "pi exits and onExit hears of it"); assert.ok(logs.some((l) => /stopping pi/.test(l))); assert.deepEqual(commands().map((c) => (c.type === "prompt" ? c.message : c.type)), ["mute", "abort"]); }); }); // Rocko's 6b R1 case: pi is stuck before agent_start, then runs the old // prompt and only afterwards reads the next one. The events carry no prompt // id, so a prompt sent after the grace would get the old run's answer. test("engine: a timed-out turn pi starts only after the grace never answers a later prompt", async () => { let exited = null; await withEngine({ abortGraceMs: 150, onExit: (e) => (exited = e) }, async ({ engine, commands }) => { await assert.rejects(engine.prompt("stall 400", { timeoutMs: 50 }), (err) => err.details.code === "timeout"); await assert.rejects(engine.prompt("after stall", { timeoutMs: 5000 }), (err) => err.details.code === "engine-wedged"); assert.ok(await until(() => exited !== null), "pi exits and onExit hears of it"); await new Promise((res) => setTimeout(res, 400)); assert.ok(!commands().some((c) => c.message === "after stall"), "nothing was sent after the grace"); }); }); // The same case in memory, after Rocko's reproducer: the old run's events // arrive after the grace while pi is still exiting. They land on the failed // turn, nothing more is written to pi, and only the exit ends the engine. // Pi's response to the old prompt comes either before its timeout or only // with the late events. for (const lateResponse of [false, true]) test(`engine: late events of a run past its grace, before pi exits, answer nothing and nothing more is sent (${lateResponse ? "late" : "early"} prompt response)`, async () => { const timers = []; const written = []; const kills = []; const child = new EventEmitter(); child.stdout = new PassThrough(); child.stderr = new PassThrough(); child.stdin = { write: (s) => { written.push(JSON.parse(s)); return true; }, end: () => {} }; child.kill = (signal) => { kills.push(signal); return true; }; let exited = null; const engine = createEngine({ command: "memory-only", args: [], spawn: () => child, abortGraceMs: 150, onExit: (e) => (exited = e), setTimeoutImpl: (fn, ms) => { const t = { fn, ms, active: true }; timers.push(t); return t; }, clearTimeoutImpl: (t) => { t.active = false; }, }); const emit = (x) => child.stdout.write(JSON.stringify(x) + "\n"); const fire = (ms) => { const t = timers.find((x) => x.ms === ms && x.active); assert.ok(t, `timer ${ms}`); t.active = false; t.fn(); }; const message = (text) => ({ role: "assistant", content: [{ type: "text", text }], stopReason: "stop" }); const tick = () => new Promise((res) => setImmediate(res)); // Checked after a tick instead of awaited, so a regression fails here // rather than hanging on a promise nothing will settle. const outcome = (p) => { const o = { state: "pending", code: null, text: null }; p.then((v) => Object.assign(o, { state: "resolved", text: v.text }), (e) => Object.assign(o, { state: "rejected", code: e.details && e.details.code })); return o; }; engine.start(); const first = engine.prompt("old", { timeoutMs: 50 }); const accept = () => emit({ type: "response", id: written[0].id, command: "prompt", success: true }); if (!lateResponse) accept(); await tick(); fire(50); await assert.rejects(first, (err) => err.details.code === "timeout"); const next = outcome(engine.prompt("new", { timeoutMs: 2000 })); fire(150); await tick(); assert.deepEqual(next, { state: "rejected", code: "engine-wedged", text: null }); assert.deepEqual(kills, ["SIGTERM"]); if (lateResponse) accept(); emit({ type: "agent_start" }); emit({ type: "tool_execution_start", toolCallId: "old-call", toolName: "read_file", args: { root: "docs", path: "old.md" } }); emit({ type: "tool_execution_end", toolCallId: "old-call", toolName: "read_file", result: { details: { root: "docs", path: "old.md", ok: true } } }); emit({ type: "turn_end", message: message("OLD RUN ANSWER") }); emit({ type: "agent_end", messages: [message("OLD RUN ANSWER")] }); emit({ type: "agent_settled" }); await tick(); const after = outcome(engine.prompt("after settle", { timeoutMs: 2000 })); await tick(); assert.deepEqual(after, { state: "rejected", code: "engine-down", text: null }); assert.deepEqual(next, { state: "rejected", code: "engine-wedged", text: null }, "the old answer did not reach the new prompt"); assert.deepEqual(written.map((c) => (c.type === "prompt" ? c.message : c.type)), ["old", "abort"], "no prompt reached pi after the grace"); assert.equal(exited, null); fire(5000); assert.deepEqual(kills, ["SIGTERM", "SIGKILL"]); child.emit("exit", null, "SIGKILL"); assert.deepEqual(exited, { code: null, signal: "SIGKILL" }); }); test("engine: a timed-out run pi did start outlives the grace; the next prompt goes out when it ends", async () => { let exited = null; await withEngine({ abortGraceMs: 150, onExit: (e) => (exited = e) }, async ({ engine, commands }) => { await assert.rejects(engine.prompt("late 400", { timeoutMs: 50 }), (err) => err.details.code === "timeout"); const r = await engine.prompt("after late", { timeoutMs: 5000 }); assert.equal(r.text, "echo: after late"); assert.deepEqual(r.tools, []); assert.equal(exited, null, "pi was not stopped"); assert.deepEqual(commands().map((c) => (c.type === "prompt" ? c.message : c.type)), ["late 400", "abort", "after late"]); }); }); test("engine: a malformed JSONL line fails the turn, not the process", () => withEngine({}, async ({ engine, logs }) => { await assert.rejects(engine.prompt("garbage"), (err) => err.details.code === "engine-protocol"); assert.ok(logs.some((l) => /malformed/.test(l))); const r = await engine.prompt("still here"); assert.equal(r.text, "echo: still here"); })); test("engine: a turn that ends in error rejects with the error code; process exit fails pending turns", async () => { let exited = null; const { engine } = start(makeRoot(), { onExit: (e) => (exited = e) }); try { await assert.rejects(engine.prompt("error"), (err) => err.details.code === "engine-error" && /fake provider error/.test(err.message)); const pending = engine.prompt("slow 5000"); await new Promise((r) => setTimeout(r, 20)); await engine.stop(); await assert.rejects(pending, (err) => err.details.code === "engine-down"); assert.ok(exited); await assert.rejects(engine.prompt("x"), /not running/); } finally { await engine.stop(); } });