Files
stack/packages/discord/tests/engine.test.mjs
T
jason.woltjeandClaude Opus 5.5 3edb15eb96 fix(discord): engine tests stop in finally; a turn pi never started stops pi instead of guessing (#1509)
Every engine test stops its engine in finally, and commands() tolerates a
log that doesn't exist yet, so a failed assertion no longer leaks a fake pi
and hangs the six-package union. A turn that failed client-side stays at
the front of the queue and holds the next prompt. If pi has sent no
agent_start ABORT_GRACE_MS (30 s) after the failure, the engine marks
itself wedged, fails held prompts with engine-wedged, refuses new ones with
engine-down, and stops pi. The exit reaches onExit, the connector exits 1,
and the unit restarts it. Pi's events carry no prompt id, so R1's approach,
dropping the turn and sending on, let a late run answer the next prompt.
Rocko rejected R1 and approved R2.

Tests: engine 17/17 (R1 fails 4, HEAD fails 5). Union 408/408 and the eight
suites green on the committed index. Record:
agents/darkwing/work/discord-engine-busy/.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-09-26 15:50:37 -05:00

312 lines
17 KiB
JavaScript

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();
}
});