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]>
This commit is contained in:
2026-09-26 15:50:37 -05:00
co-authored by Claude Opus 5.5
parent f55b94e866
commit 3edb15eb96
11 changed files with 1843 additions and 102 deletions
+49
View File
@@ -2997,3 +2997,52 @@ captured the proxy's real headers: Host 127.0.0.1:<port>, no Origin. One low
finding is open: the notes' wording on raw Host syntax is broader than the
code (see DEFERRED Done). The running board keeps the old code until Sage
restarts it.
## 2026-09-26: Discord engine test cleanup and the busy gap (#1509, Darkwing, reviewed by Rocko)
Before, the engine tests could leak a fake pi. Seven tests stopped the engine
outside `finally`, and one read the fake's `commands.jsonl` 20 ms after a
prompt, which got ENOENT under load. A failed assertion then left the fake
running and the test file never exited. On clean HEAD the six-package union
(control-board, webui, seat, mosaic, ledger, discord) hung at the 240 s cap at
199 ok in four of four runs. In `engine-pi.mjs`, `busy` ignored a turn that
timed out before its `agent_start` was read, so the next prompt went to a pi
that was still running and pi refused it as streaming.
After, every engine test stops its engine in `finally`, and `commands()`
tolerates a log that doesn't exist yet. A failed turn stays at the front of
the queue and holds the next prompt, so late events land on it. If pi has
sent no `agent_start` 30 s after the turn failed (`ABORT_GRACE_MS`), the
engine marks itself wedged. It fails held prompts with `engine-wedged`,
refuses anything new with `engine-down`, and stops pi (SIGTERM, then SIGKILL
after 5 s). The exit reaches `onExit`, `cli.mjs` exits 1, and the unit's
`Restart=on-failure` starts a new connector. A run pi did start still ends
on its `agent_end` or a settle, as before.
Correction: R1 handled the grace by dropping the failed turn and sending the
next prompt. Rocko rejected it (F1, High). Pi's events carry no prompt id,
so if the old run came late, its answer and tool records went to the new
prompt. With a real child, R1 answered "after stall" with "echo: stalled".
R1's README claimed late events would be dropped. That was wrong.
Evidence: `packages/discord/tests/engine.test.mjs` 17/17. R1's engine fails 4
of these tests and HEAD's fails 5. Three mutations of R2 (no busy check, no
wedge, held prompts not failed) each fail. On `git archive` of 401cc850 plus
the three files, the union passed 406/406 three times, 23 to 24 s each, with
no fake pi left running. The eight suites passed at 24/90/43/17/14/15/63/18
on that snapshot. On the committed index (base f55b94e8, which adds two board
tests), the union passed 408/408 and the eight suites passed again. Rocko approved
R2 (engine 77077b7f, tests 47a99817, fake a8e54cc3) in
`agents/rocko/work/discord-engine-busy-r2-review-2026-09-26.md`. Record:
`agents/darkwing/work/discord-engine-busy/`.
Known limit: the unit's start limit (5 starts in 600 s) does not bound
repeated wedges. The live binding has turnTimeoutSeconds 600, so one wedge
cycle takes at least 645 s, which is longer than the window. A pi that wedges
on every prompt restarts once per prompt and never hits the limit. It happens
only when someone sends a message, and they get an error reply. Sage accepted
this and is recording two follow-ups in DEFERRED: a cumulative wedge stop, and
a startup journal line with HEAD, the scoped dirty state and a runtime file
digest. Also unchanged from HEAD: a run pi started that never ends holds
prompts until each one times out. No connector restart and no push were done
here.
@@ -0,0 +1,162 @@
# Discord engine: guaranteed test cleanup and the timeout gap in `busy` (#1509), R2 candidate
Sage assigned this on 2026-09-26 as 6b, source only. Rocko reviews. Base is
HEAD 401cc850. Not committed. The live connector runs from this checkout, so
Sage is holding its restart until this is approved and committed. Nobody
should restart it from a working copy.
## Defects (DEFERRED Open, "Discord engine: leaked fake pi…")
(a) `engine.test.mjs` read the fake's `commands.jsonl` 20 ms after a prompt and
got ENOENT under load. Seven tests stopped the engine outside `finally`, so a
failed assertion left the fake pi running and the test file never exited.
(b) `busy` was `state.busy || pending.some((t) => !t.done)`. If a turn timed out
before its `agent_start` was read, it was done while `state.busy` was still
false. The next prompt then went straight to pi, which refused it as
streaming.
## R1 and Rocko's finding
R1 held later prompts behind a failed turn. If pi had sent no `agent_start` for
it within a grace period, R1 dropped that turn from the queue and sent the next
prompt. Rocko rejected it (F1, High), in
`agents/rocko/work/discord-engine-busy-r1-review-2026-09-26.md`, sha256
047dbd8f.
Pi's events carry no prompt id. The engine attributes them to the front of its
queue. Silence until the grace ends does not prove the old run will never come.
If pi then runs it, its events land on the new prompt, which R1 had just put at
the front. Rocko's reproducer got the old run's answer and its `old.md` tool
record back as the new prompt's result. My R1 README said such events "find no
live head and are dropped". That was wrong.
The R1 files stay here as `r1-manifest.sha256` and `r1.patch`.
## Change (R2)
`packages/discord/src/engine-pi.mjs`:
- `engineBusy()` is `state.busy || state.pending.length > 0`. A failed turn
still in the queue holds the next prompt back, and stays at the front, so any
late events for it land on it. `prompt()`, `sendHeld()` and the `busy` getter
use it. This part is unchanged from R1.
- The bound is now a stop, not a drop. When a turn fails while it is still in
the queue, `failTurn` starts a timer, `abortGraceMs` (default
`ABORT_GRACE_MS`, 30 s, an engine option, not binding config). When it
fires:
- If pi has sent `agent_start` (`state.busy`), nothing happens. That run
ends on its `agent_end` or a settle, as on HEAD.
- Otherwise `wedge()` sets `state.wedged`, fails every held prompt with
code `engine-wedged`, and stops pi: stdin closed, SIGTERM, then SIGKILL
after 5 s. The failed turn stays at the front until the exit.
- While wedged, nothing is written to that child. `write()`, `sendHeld()` and
`prompt()` refuse, and a new prompt fails at once with `engine-down`. The exit
runs the usual `failAll` and `onExit`.
- `stop()`'s body moved into `stopChild()`, which both `stop()` and `wedge()`
call.
- `release()` clears the timer wherever a turn leaves the queue: `agent_end`,
settle, a refused send, and process exit. As in R1, the settle handler removes
turns before failing them.
What recovery looks like live: `cli.mjs` handles `onExit` with `shutdown(1)`.
The unit's `Restart=on-failure` starts a new connector and a new pi 15 s later,
within its limit of five tries in ten minutes. This change doesn't touch the
unit or the restart policy. A wedge now costs one connector restart. R1 would
have kept the same pi and risked a wrong answer.
`packages/discord/tests/fake-pi.mjs`:
- `mute`: accepted and never run.
- `stall <ms>`: accepted, then the fake reads nothing for `<ms>`, runs the
stalled prompt, and only then reads what came in meanwhile. This is the
order in Rocko's case.
`packages/discord/tests/engine.test.mjs`:
- `withEngine()` stops the engine in `finally`. Every test that starts an
engine uses it, or has its own `try/finally` in the exit test.
- `commands()` returns `[]` until the fake creates its log. The held-prompt
test waits for the first prompt with `until()` instead of a 20 ms sleep, and
its first prompt is `slow 300`.
- The manual-timer test fires the turn timer before any pi event is read. The
next prompt must wait for the settle and get its own answer.
- New or changed for R2:
- `mute` with `abortGraceMs: 150`. The held prompt fails with
`engine-wedged` after the grace, a later prompt fails with
`engine-down`, `onExit` fires, and pi saw only `mute` and `abort`.
- `stall 400` with the same grace, which is Rocko's case with a real
child. The held prompt fails with `engine-wedged`, pi exits, and "after
stall" never reaches pi.
- Rocko's reproducer as an in-memory test, run twice. The old prompt's
response comes either before its timeout or only with the late events.
After the grace, the old run's start, tool pair, answer, end and settle
arrive while pi is still exiting. The held prompt stays failed with
`engine-wedged` and a later prompt fails with `engine-down`. Pi saw only
`old` and `abort`, then SIGTERM, then SIGKILL at 5 s. Only the exit
reaches `onExit`. The test reads recorded outcomes after a tick instead
of awaiting, so a regression fails instead of hanging.
- `late 400` with the same grace. Pi started that run, so the grace does
not stop pi, and the next prompt gets its own answer when the run ends.
## Evidence
- `engine.test.mjs`: 17/17.
- R1's engine (d5bf24b5, from `r1.patch`) against these tests fails 4: `mute`,
`stall`, and both in-memory runs. In that run a probe shows R1 answering
"after stall" with "echo: stalled". An earlier draft of the in-memory test
awaited the held prompt and hung on R1 until the 120 s cap. It now fails in
milliseconds.
- HEAD's engine against these tests fails 5: the manual-timer test and the
same four.
- Mutations of R2:
- Without the `state.busy` check, the `late 400` test fails.
- Without the `wedge()` call, 4 fail.
- Without failing held prompts in `wedge()`, 4 fail.
- Test union (control-board, webui, seat, mosaic, ledger, discord) at default
concurrency on `git archive` of 401cc850 plus the three files: 406/406
three times, 23 to 24 s each. No fake pi was left running.
- Eight suites green on that snapshot: config 24, task 90, foundation 43,
conductor 17, release 14, auth 15, discord 63, extension-package 18.
Logs: `/tmp/dw-6b-r2-conc-{1,2,3}.txt`. R1's evidence runs:
`/tmp/dw-6b-conc-{1,2,3}.txt`, `/tmp/dw-6b-serial.txt`. HEAD's hang control:
`/tmp/dw-1509-headctl-{1,2,3}.txt`, `/tmp/dw-1509-ef00-1.txt`. There, HEAD hit
the 240 s cap at 199 ok under the union's load.
## Not covered
- A run pi started and never ends, even after abort, still holds prompts
until pi settles or exits. Each held prompt fails at its own timeout ("while
waiting for the engine"). HEAD behaves the same way through `state.busy`, and
Rocko did not block on it. Only a pi restart clears it.
- A wedge ends the connector process, and the recovery is systemd's restart.
Nothing here changes the unit, and the restart limit still applies.
- No live restart, and no change to the binding schema.
## Frozen files
`r2-manifest.sha256` holds the three R2 hashes. `r2.patch` is `git diff
packages/discord` at freeze time.
## Review
Rocko, R1, 2026-09-26: request changes, F1 High, as described above. Report:
`agents/rocko/work/discord-engine-busy-r1-review-2026-09-26.md`, sha256
047dbd8f.
Rocko, R2, 2026-09-26: approved the three pinned files. Report:
`agents/rocko/work/discord-engine-busy-r2-review-2026-09-26.md`, sha256
ed5510a0. He checked the manifests before and after, ran 17/17 himself, and
read the CLI shutdown path, `connector.stop` and the unit template. Sage asked
him three operational questions:
- A wedge exits 1, never 3. Exit 3 remains the supervised startup refusal.
- The unit's start limit (5 starts in 600 s) is a rate limit. It does not
bound repeated wedges. With the default 180 s turn timeout, the 30 s grace
and the 15 s restart delay, a cycle takes at least 225 s. That stays under
the limit, so a pi that wedges every time could restart indefinitely.
Stopping for good after repeated wedges would need a separate policy. This
change does not add one.
- He recommends, as a nonblocking follow-up, that the connector journal
record at startup: HEAD, dirty state scoped to runtime source, and a digest
of the runtime files. A wedge restart loads whatever the checkout holds.
This section was added after approval, so the README hash no longer matches
the one Rocko pinned (69350f29). The three source files are unchanged.
@@ -0,0 +1,3 @@
d5bf24b59c07c85067f4087c03b54ca8b4df923c1d591441dedd2e8a7ff2ae39 packages/discord/src/engine-pi.mjs
f0abee9c243d46d66dd2271abc3fd89089c350ac6a66ab49131bce80adfcdc33 packages/discord/tests/engine.test.mjs
fa1bf44e3f33eb970a714ada1c686abbc1932baaf418679276edf8813abbe6de packages/discord/tests/fake-pi.mjs
@@ -0,0 +1,414 @@
diff --git a/packages/discord/src/engine-pi.mjs b/packages/discord/src/engine-pi.mjs
index 5c8fd0a9..9dcaeb42 100644
--- a/packages/discord/src/engine-pi.mjs
+++ b/packages/discord/src/engine-pi.mjs
@@ -16,9 +16,12 @@
// from `tool_execution_start`/`tool_execution_end` into the result so the
// turn record shows what was read. An `agent_end` with `willRetry` is not
// the end of the run. A timeout sends `abort` and fails that turn; the
-// process stays. A malformed JSONL line from pi fails the current turn (its
-// outcome is now unknowable) and the process stays. Process exit fails
-// every pending turn and is reported through `onExit`.
+// process stays. The failed turn holds later prompts back until its
+// agent_end or a settle. If pi has not started it within ABORT_GRACE_MS, it
+// is dropped and the next prompt goes out; a run pi did start holds them
+// until it ends, as any run does. A malformed JSONL line from pi fails the
+// current turn (its outcome is now unknowable) and the process stays.
+// Process exit fails every pending turn and is reported through `onExit`.
//
// Framing follows pi's RPC doc: split on "\n" only, strip a trailing "\r".
// Node readline is not used because it also splits on U+2028/U+2029.
@@ -64,12 +67,19 @@ export function assistantText(message) {
.trim();
}
+// How long a turn that failed here (timeout, protocol error) may wait for
+// pi's agent_start before it stops holding the next prompt back. Without a
+// bound, a prompt pi accepted but never ran would queue every later prompt
+// until restart.
+export const ABORT_GRACE_MS = 30000;
+
export function createEngine({
command, args, cwd, env = {},
spawn = nodeSpawn,
setTimeoutImpl = globalThis.setTimeout, clearTimeoutImpl = globalThis.clearTimeout,
log = () => {},
onExit = () => {},
+ abortGraceMs = ABORT_GRACE_MS,
} = {}) {
if (typeof command !== "string" || command.length === 0) throw new DiscordError("engine: command required", 1);
if (!Array.isArray(args)) throw new DiscordError("engine: args required", 1);
@@ -80,15 +90,35 @@ export function createEngine({
// A turn that fails on the client side (timeout, protocol error) stays in
// the pending queue, marked done, until pi's own turn_end for it arrives.
- // Otherwise that turn_end would be attributed to the next prompt.
+ // Otherwise that turn_end would be attributed to the next prompt. It holds
+ // later prompts back; if pi has not started it within abortGraceMs, it goes.
function failTurn(turn, code, message) {
if (turn.done) return;
turn.done = true;
if (turn.timer !== null) clearTimeoutImpl(turn.timer);
turn.timer = null;
+ if (state.pending.includes(turn)) {
+ turn.grace = setTimeoutImpl(() => {
+ turn.grace = null;
+ // No agent_start by now: pi never started this run and will send no
+ // agent_end for it, so it leaves the queue and cannot take the next
+ // prompt's. A run pi did start keeps its place until it ends.
+ if (!state.busy) {
+ const i = state.pending.indexOf(turn);
+ if (i !== -1) state.pending.splice(i, 1);
+ }
+ sendHeld();
+ }, abortGraceMs);
+ }
turn.reject(new DiscordError(message, 1, { code }));
}
+ // Call when a turn leaves the pending queue.
+ function release(turn) {
+ if (turn.grace !== null) clearTimeoutImpl(turn.grace);
+ turn.grace = null;
+ }
+
function settleTurn(turn, value) {
if (turn.done) return;
turn.done = true;
@@ -99,7 +129,10 @@ export function createEngine({
function failAll(code, message) {
const pending = state.pending.splice(0);
- for (const t of pending) failTurn(t, code, message);
+ for (const t of pending) {
+ release(t);
+ failTurn(t, code, message);
+ }
for (const h of state.held.splice(0)) failTurn(h.turn, code, message);
for (const [, r] of state.responses) r.reject(new DiscordError(message, 1, { code }));
state.responses.clear();
@@ -174,6 +207,7 @@ export function createEngine({
// Attribute the run to the head even if it failed client-side, so the
// next prompt's agent_end is not taken for this one.
const run = state.pending.shift();
+ if (run) release(run);
if (!run || run.done) return;
const messages = Array.isArray(event.messages) ? event.messages.filter((m) => m && m.role === "assistant") : [];
const message = messages.length > 0 ? messages[messages.length - 1] : run.last;
@@ -193,13 +227,14 @@ export function createEngine({
// this settle and still has no agent_end will never get one: fail it now
// instead of waiting for its timeout. Turns whose prompt response has
// not arrived yet belong to a later run and stay.
+ const dropped = [];
const keep = [];
- for (const t of state.pending) {
- if (t.done) continue;
- if (t.accepted) failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
- else keep.push(t);
- }
+ for (const t of state.pending) (t.done || t.accepted ? dropped : keep).push(t);
state.pending = keep;
+ for (const t of dropped) {
+ release(t);
+ failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
+ }
sendHeld();
}
}
@@ -214,15 +249,23 @@ export function createEngine({
// Never accepted: pi will not emit a turn_end for it, so remove it.
const i = state.pending.indexOf(turn);
if (i !== -1) state.pending.splice(i, 1);
+ release(turn);
failTurn(turn, (err.details && err.details.code) || "engine-refused", err.message);
sendHeld();
});
}
+ // Pi is busy from our side while any sent prompt is still queued, even one
+ // that already failed here: a turn that timed out before its agent_start
+ // was read leaves state.busy false while pi runs it, and sending then would
+ // be refused as streaming. It leaves the queue on its agent_end, on a
+ // settle, on a refused send, or when its grace ends before pi started it.
+ const engineBusy = () => state.busy || state.pending.length > 0;
+
// After a settle (or a refused send) the oldest held prompt goes out.
function sendHeld() {
if (state.exited !== null) return;
- if (state.busy || state.pending.some((t) => !t.done)) return;
+ if (engineBusy()) return;
const next = state.held.shift();
if (next) send(next.turn, next.command);
}
@@ -281,7 +324,7 @@ export function createEngine({
// with DiscordError carrying details.code for the turn record.
prompt(text, { timeoutMs = 180000 } = {}) {
if (typeof text !== "string" || text.length === 0) throw new DiscordError("prompt text required", 1);
- const turn = { resolve: null, reject: null, timer: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
+ const turn = { resolve: null, reject: null, timer: null, grace: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
const done = new Promise((resolve, reject) => {
turn.resolve = resolve;
turn.reject = reject;
@@ -310,13 +353,13 @@ export function createEngine({
failTurn(turn, "engine-down", "engine is not running");
return done;
}
- if (state.busy || state.pending.some((t) => !t.done) || state.held.length > 0) state.held.push({ turn, command });
+ if (engineBusy() || state.held.length > 0) state.held.push({ turn, command });
else send(turn, command);
return done;
},
get busy() {
- return state.busy || state.pending.some((t) => !t.done) || state.held.length > 0;
+ return engineBusy() || state.held.length > 0;
},
get pendingCount() {
return state.pending.filter((t) => !t.done).length + state.held.length;
diff --git a/packages/discord/tests/engine.test.mjs b/packages/discord/tests/engine.test.mjs
index 59674d1e..62dd6017 100644
--- a/packages/discord/tests/engine.test.mjs
+++ b/packages/discord/tests/engine.test.mjs
@@ -28,7 +28,20 @@ function start(root, extra = {}) {
log: (m) => logs.push(m), ...extra,
});
engine.start();
- return { engine, logs, commands: () => readFileSync(logPath, "utf8").trim().split("\n").filter(Boolean).map((l) => JSON.parse(l)) };
+ // 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", () => {
@@ -56,8 +69,7 @@ test("engine: with tools, buildPiArgs turns pi's own tools off, loads the extens
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", async () => {
- const { engine } = start(makeRoot());
+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);
@@ -71,45 +83,38 @@ test("engine: a run with tool turns settles once, on the answer, with every tool
assert.equal(plain.turns, 1);
await idle(engine);
assert.equal(engine.busy, false);
- await engine.stop();
-});
+}));
-test("engine: a run that ends on a tool-only turn fails the prompt as empty; a retried run settles on the real end", async () => {
- const { engine } = start(makeRoot());
+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");
- await engine.stop();
-});
+}));
-test("engine: one prompt, one turn, text and usage come back", async () => {
- const { engine } = start(makeRoot());
- try {
- 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);
- } finally {
- await engine.stop();
- }
-});
+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", async () => {
- const { engine, commands } = start(makeRoot());
- const first = engine.prompt("slow 150");
- await new Promise((r) => setTimeout(r, 20));
+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);
- await new Promise((r) => setTimeout(r, 20));
- assert.equal(commands().filter((c) => c.type === "prompt").length, 1, "the second prompt is not sent while pi is busy");
+ 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 = commands().filter((c) => c.type === "prompt");
+ 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).
@@ -117,11 +122,9 @@ test("engine: a prompt while streaming is held until pi settles, then sent as it
assert.equal(prompts[1].streamingBehavior, undefined);
await idle(engine);
assert.equal(engine.busy, false);
- await engine.stop();
-});
+}));
-test("engine: a held prompt that times out before pi settles fails on its own and is never sent", async () => {
- const { engine, commands } = start(makeRoot());
+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));
@@ -130,50 +133,85 @@ test("engine: a held prompt that times out before pi settles fails on its own an
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");
- await engine.stop();
-});
+}));
-test("engine: timeout sends abort and fails only that turn; the process stays", async () => {
- const { engine, commands, logs } = start(makeRoot());
+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");
- await engine.stop();
+}));
+
+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: tool events from a run that outlived its timeout never land in the next prompt's record", async () => {
- const { engine } = start(makeRoot());
- try {
- await assert.rejects(engine.prompt("late 200", { timeoutMs: 40 }), (err) => err.details.code === "timeout");
- const r = await engine.prompt("after late");
+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, "the dead run's turns are not counted here");
- } finally {
- await engine.stop();
- }
+ 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);
+ });
});
-test("engine: a malformed JSONL line fails the turn, not the process", async () => {
- const { engine, logs } = start(makeRoot());
+// "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: a timed-out turn pi never started holds the next prompt only for the abort grace, then leaves the queue", () => withEngine({ abortGraceMs: 150 }, async ({ engine, commands }) => {
+ 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();
+ const r = await engine.prompt("after mute", { timeoutMs: 5000 });
+ assert.equal(r.text, "echo: after mute");
+ assert.ok(Date.now() - started >= 100, "held for the grace, not sent at once");
+ assert.deepEqual(commands().map((c) => (c.type === "prompt" ? c.message : c.type)), ["mute", "abort", "after mute"]);
+ await idle(engine);
+ assert.equal(engine.busy, false);
+}));
+
+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");
- await engine.stop();
-});
+}));
test("engine: a turn that ends in error rejects with the error code; process exit fails pending turns", async () => {
- const root = makeRoot();
let exited = null;
- const { engine } = start(root, { onExit: (e) => (exited = e) });
- 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/);
+ 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();
+ }
});
diff --git a/packages/discord/tests/fake-pi.mjs b/packages/discord/tests/fake-pi.mjs
index 94919306..ecc04e3b 100644
--- a/packages/discord/tests/fake-pi.mjs
+++ b/packages/discord/tests/fake-pi.mjs
@@ -8,6 +8,7 @@
// then a second turn that answers "read <n> file(s)"
// "toolonly" a run whose only turn calls a tool and never answers
// "retry" an agent_end with willRetry, then the real answer
+// "mute" accept the prompt and emit nothing, staying idle
// "late <ms>" ignore abort; after <ms> emit a tool pair and a tool turn,
// then answer "late reply", like a run that outlives its
// client-side timeout
@@ -30,6 +31,7 @@ function assistant(text, stopReason = "stop") {
}
function run(text) {
+ if (text === "mute") return;
busy = true;
out({ type: "agent_start" });
out({ type: "turn_start" });
@@ -0,0 +1,3 @@
77077b7fbd5a933ffd352094eb073227c299ba47b7aea52d4e60fdc55cc7101e packages/discord/src/engine-pi.mjs
47a998179c6eb46827f43ab2c6b0f6b062ef94fb47da402cfb9af8c4f180f38e packages/discord/tests/engine.test.mjs
a8e54cc3f4b670eef2c06755b63e9e6bfeb44b1b1efde3aca9bfaa91583c3ef3 packages/discord/tests/fake-pi.mjs
@@ -0,0 +1,651 @@
diff --git a/packages/discord/src/engine-pi.mjs b/packages/discord/src/engine-pi.mjs
index 5c8fd0a9..46ef1f88 100644
--- a/packages/discord/src/engine-pi.mjs
+++ b/packages/discord/src/engine-pi.mjs
@@ -16,9 +16,14 @@
// from `tool_execution_start`/`tool_execution_end` into the result so the
// turn record shows what was read. An `agent_end` with `willRetry` is not
// the end of the run. A timeout sends `abort` and fails that turn; the
-// process stays. A malformed JSONL line from pi fails the current turn (its
-// outcome is now unknowable) and the process stays. Process exit fails
-// every pending turn and is reported through `onExit`.
+// process stays. The failed turn holds later prompts back until its
+// agent_end or a settle. If pi has not started it within ABORT_GRACE_MS, the
+// engine stops pi instead of sending again: pi's events carry no prompt id,
+// so a late run of the failed prompt would be taken for the next one's. A
+// run pi did start holds later prompts until it ends, as any run does. A
+// malformed JSONL line from pi fails the current turn (its outcome is now
+// unknowable) and the process stays. Process exit fails every pending turn
+// and is reported through `onExit`.
//
// Framing follows pi's RPC doc: split on "\n" only, strip a trailing "\r".
// Node readline is not used because it also splits on U+2028/U+2029.
@@ -64,31 +69,67 @@ export function assistantText(message) {
.trim();
}
+// How long a turn that failed here (timeout, protocol error) may wait for
+// pi's agent_start before the engine stops pi. Without a bound, a prompt pi
+// accepted but never ran would hold every later prompt until restart.
+export const ABORT_GRACE_MS = 30000;
+
export function createEngine({
command, args, cwd, env = {},
spawn = nodeSpawn,
setTimeoutImpl = globalThis.setTimeout, clearTimeoutImpl = globalThis.clearTimeout,
log = () => {},
onExit = () => {},
+ abortGraceMs = ABORT_GRACE_MS,
} = {}) {
if (typeof command !== "string" || command.length === 0) throw new DiscordError("engine: command required", 1);
if (!Array.isArray(args)) throw new DiscordError("engine: args required", 1);
// pending: prompts sent to pi, oldest first. held: prompts waiting for pi
// to settle before they are sent, oldest first.
- const state = { child: null, buffer: "", pending: [], held: [], responses: new Map(), nextId: 1, busy: false, exited: null };
+ // wedged: set when the engine gave up on pi and is stopping it. Nothing
+ // is sent to that child again.
+ const state = { child: null, buffer: "", pending: [], held: [], responses: new Map(), nextId: 1, busy: false, exited: null, wedged: false };
// A turn that fails on the client side (timeout, protocol error) stays in
// the pending queue, marked done, until pi's own turn_end for it arrives.
- // Otherwise that turn_end would be attributed to the next prompt.
+ // Otherwise that turn_end would be attributed to the next prompt. It holds
+ // later prompts back; if pi has not started it within abortGraceMs, the
+ // engine stops pi.
function failTurn(turn, code, message) {
if (turn.done) return;
turn.done = true;
if (turn.timer !== null) clearTimeoutImpl(turn.timer);
turn.timer = null;
+ if (state.pending.includes(turn)) {
+ turn.grace = setTimeoutImpl(() => {
+ turn.grace = null;
+ // A run pi started keeps its place until its agent_end or a settle.
+ if (state.busy) return;
+ // No agent_start yet. Pi may never run this prompt, or its events
+ // may still be on the way; with no prompt id in them, nothing sent
+ // now could be told apart from it. Stop pi: held prompts fail, and
+ // the exit fails the rest and reaches onExit.
+ log(`engine: no agent_start ${abortGraceMs} ms after a failed turn; stopping pi`);
+ wedge();
+ }, abortGraceMs);
+ }
turn.reject(new DiscordError(message, 1, { code }));
}
+ function wedge() {
+ if (state.wedged || state.exited !== null) return;
+ state.wedged = true;
+ for (const h of state.held.splice(0)) failTurn(h.turn, "engine-wedged", "engine stopped: pi did not start an aborted turn");
+ stopChild();
+ }
+
+ // Call when a turn leaves the pending queue.
+ function release(turn) {
+ if (turn.grace !== null) clearTimeoutImpl(turn.grace);
+ turn.grace = null;
+ }
+
function settleTurn(turn, value) {
if (turn.done) return;
turn.done = true;
@@ -99,7 +140,10 @@ export function createEngine({
function failAll(code, message) {
const pending = state.pending.splice(0);
- for (const t of pending) failTurn(t, code, message);
+ for (const t of pending) {
+ release(t);
+ failTurn(t, code, message);
+ }
for (const h of state.held.splice(0)) failTurn(h.turn, code, message);
for (const [, r] of state.responses) r.reject(new DiscordError(message, 1, { code }));
state.responses.clear();
@@ -174,6 +218,7 @@ export function createEngine({
// Attribute the run to the head even if it failed client-side, so the
// next prompt's agent_end is not taken for this one.
const run = state.pending.shift();
+ if (run) release(run);
if (!run || run.done) return;
const messages = Array.isArray(event.messages) ? event.messages.filter((m) => m && m.role === "assistant") : [];
const message = messages.length > 0 ? messages[messages.length - 1] : run.last;
@@ -193,13 +238,14 @@ export function createEngine({
// this settle and still has no agent_end will never get one: fail it now
// instead of waiting for its timeout. Turns whose prompt response has
// not arrived yet belong to a later run and stay.
+ const dropped = [];
const keep = [];
- for (const t of state.pending) {
- if (t.done) continue;
- if (t.accepted) failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
- else keep.push(t);
- }
+ for (const t of state.pending) (t.done || t.accepted ? dropped : keep).push(t);
state.pending = keep;
+ for (const t of dropped) {
+ release(t);
+ failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
+ }
sendHeld();
}
}
@@ -214,21 +260,29 @@ export function createEngine({
// Never accepted: pi will not emit a turn_end for it, so remove it.
const i = state.pending.indexOf(turn);
if (i !== -1) state.pending.splice(i, 1);
+ release(turn);
failTurn(turn, (err.details && err.details.code) || "engine-refused", err.message);
sendHeld();
});
}
+ // Pi is busy from our side while any sent prompt is still queued, even one
+ // that already failed here: a turn that timed out before its agent_start
+ // was read leaves state.busy false while pi runs it, and sending then would
+ // be refused as streaming. It leaves the queue on its agent_end, on a
+ // settle, on a refused send, or at process exit.
+ const engineBusy = () => state.busy || state.pending.length > 0;
+
// After a settle (or a refused send) the oldest held prompt goes out.
function sendHeld() {
- if (state.exited !== null) return;
- if (state.busy || state.pending.some((t) => !t.done)) return;
+ if (state.exited !== null || state.wedged) return;
+ if (engineBusy()) return;
const next = state.held.shift();
if (next) send(next.turn, next.command);
}
function write(command) {
- if (!state.child || state.exited !== null) throw new DiscordError("engine is not running", 1, { code: "engine-down" });
+ if (!state.child || state.exited !== null || state.wedged) throw new DiscordError("engine is not running", 1, { code: "engine-down" });
state.child.stdin.write(JSON.stringify(command) + "\n");
}
@@ -245,6 +299,30 @@ export function createEngine({
});
}
+ function stopChild({ graceMs = 5000 } = {}) {
+ const child = state.child;
+ if (!child || state.exited !== null) return Promise.resolve(state.exited);
+ return new Promise((resolve) => {
+ const timer = setTimeoutImpl(() => {
+ try {
+ child.kill("SIGKILL");
+ } catch {
+ // already gone
+ }
+ }, graceMs);
+ child.once("exit", () => {
+ clearTimeoutImpl(timer);
+ resolve(state.exited);
+ });
+ try {
+ child.stdin.end();
+ child.kill("SIGTERM");
+ } catch {
+ // already gone
+ }
+ });
+ }
+
return {
start() {
if (state.child) throw new DiscordError("engine already started", 1);
@@ -281,7 +359,7 @@ export function createEngine({
// with DiscordError carrying details.code for the turn record.
prompt(text, { timeoutMs = 180000 } = {}) {
if (typeof text !== "string" || text.length === 0) throw new DiscordError("prompt text required", 1);
- const turn = { resolve: null, reject: null, timer: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
+ const turn = { resolve: null, reject: null, timer: null, grace: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
const done = new Promise((resolve, reject) => {
turn.resolve = resolve;
turn.reject = reject;
@@ -306,44 +384,24 @@ export function createEngine({
}
failTurn(turn, "timeout", `turn timed out after ${timeoutMs} ms`);
}, timeoutMs);
- if (state.exited !== null) {
+ if (state.exited !== null || state.wedged) {
failTurn(turn, "engine-down", "engine is not running");
return done;
}
- if (state.busy || state.pending.some((t) => !t.done) || state.held.length > 0) state.held.push({ turn, command });
+ if (engineBusy() || state.held.length > 0) state.held.push({ turn, command });
else send(turn, command);
return done;
},
get busy() {
- return state.busy || state.pending.some((t) => !t.done) || state.held.length > 0;
+ return engineBusy() || state.held.length > 0;
},
get pendingCount() {
return state.pending.filter((t) => !t.done).length + state.held.length;
},
- stop({ graceMs = 5000 } = {}) {
- const child = state.child;
- if (!child || state.exited !== null) return Promise.resolve(state.exited);
- return new Promise((resolve) => {
- const timer = setTimeoutImpl(() => {
- try {
- child.kill("SIGKILL");
- } catch {
- // already gone
- }
- }, graceMs);
- child.once("exit", () => {
- clearTimeoutImpl(timer);
- resolve(state.exited);
- });
- try {
- child.stdin.end();
- child.kill("SIGTERM");
- } catch {
- // already gone
- }
- });
+ stop(options) {
+ return stopChild(options);
},
};
}
diff --git a/packages/discord/tests/engine.test.mjs b/packages/discord/tests/engine.test.mjs
index 59674d1e..f3bd7525 100644
--- a/packages/discord/tests/engine.test.mjs
+++ b/packages/discord/tests/engine.test.mjs
@@ -4,6 +4,8 @@ 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");
@@ -28,7 +30,20 @@ function start(root, extra = {}) {
log: (m) => logs.push(m), ...extra,
});
engine.start();
- return { engine, logs, commands: () => readFileSync(logPath, "utf8").trim().split("\n").filter(Boolean).map((l) => JSON.parse(l)) };
+ // 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", () => {
@@ -56,8 +71,7 @@ test("engine: with tools, buildPiArgs turns pi's own tools off, loads the extens
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", async () => {
- const { engine } = start(makeRoot());
+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);
@@ -71,45 +85,38 @@ test("engine: a run with tool turns settles once, on the answer, with every tool
assert.equal(plain.turns, 1);
await idle(engine);
assert.equal(engine.busy, false);
- await engine.stop();
-});
+}));
-test("engine: a run that ends on a tool-only turn fails the prompt as empty; a retried run settles on the real end", async () => {
- const { engine } = start(makeRoot());
+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");
- await engine.stop();
-});
+}));
-test("engine: one prompt, one turn, text and usage come back", async () => {
- const { engine } = start(makeRoot());
- try {
- 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);
- } finally {
- await engine.stop();
- }
-});
+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", async () => {
- const { engine, commands } = start(makeRoot());
- const first = engine.prompt("slow 150");
- await new Promise((r) => setTimeout(r, 20));
+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);
- await new Promise((r) => setTimeout(r, 20));
- assert.equal(commands().filter((c) => c.type === "prompt").length, 1, "the second prompt is not sent while pi is busy");
+ 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 = commands().filter((c) => c.type === "prompt");
+ 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).
@@ -117,11 +124,9 @@ test("engine: a prompt while streaming is held until pi settles, then sent as it
assert.equal(prompts[1].streamingBehavior, undefined);
await idle(engine);
assert.equal(engine.busy, false);
- await engine.stop();
-});
+}));
-test("engine: a held prompt that times out before pi settles fails on its own and is never sent", async () => {
- const { engine, commands } = start(makeRoot());
+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));
@@ -130,50 +135,177 @@ test("engine: a held prompt that times out before pi settles fails on its own an
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");
- await engine.stop();
-});
+}));
-test("engine: timeout sends abort and fails only that turn; the process stays", async () => {
- const { engine, commands, logs } = start(makeRoot());
+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");
- await engine.stop();
+}));
+
+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: tool events from a run that outlived its timeout never land in the next prompt's record", async () => {
- const { engine } = start(makeRoot());
- try {
- await assert.rejects(engine.prompt("late 200", { timeoutMs: 40 }), (err) => err.details.code === "timeout");
- const r = await engine.prompt("after late");
+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, "the dead run's turns are not counted here");
- } finally {
- await engine.stop();
- }
+ 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");
+ });
});
-test("engine: a malformed JSONL line fails the turn, not the process", async () => {
- const { engine, logs } = start(makeRoot());
+// 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");
- await engine.stop();
-});
+}));
test("engine: a turn that ends in error rejects with the error code; process exit fails pending turns", async () => {
- const root = makeRoot();
let exited = null;
- const { engine } = start(root, { onExit: (e) => (exited = e) });
- 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/);
+ 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();
+ }
});
diff --git a/packages/discord/tests/fake-pi.mjs b/packages/discord/tests/fake-pi.mjs
index 94919306..bf94c013 100644
--- a/packages/discord/tests/fake-pi.mjs
+++ b/packages/discord/tests/fake-pi.mjs
@@ -8,9 +8,13 @@
// then a second turn that answers "read <n> file(s)"
// "toolonly" a run whose only turn calls a tool and never answers
// "retry" an agent_end with willRetry, then the real answer
+// "mute" accept the prompt and emit nothing, staying idle
// "late <ms>" ignore abort; after <ms> emit a tool pair and a tool turn,
// then answer "late reply", like a run that outlives its
// client-side timeout
+// "stall <ms>" accept the prompt, then read nothing for <ms> (pi stuck
+// before agent_start); then run it, answering "echo:
+// stalled", and only then read what came in meanwhile
// anything else answer "echo: <text>" immediately
// A prompt received while busy without streamingBehavior is refused, as pi
// does. A prompt with streamingBehavior followUp is folded into the running
@@ -30,6 +34,7 @@ function assistant(text, stopReason = "stop") {
}
function run(text) {
+ if (text === "mute") return;
busy = true;
out({ type: "agent_start" });
out({ type: "turn_start" });
@@ -109,11 +114,16 @@ function run(text) {
let current = null;
let buffer = "";
+let stalled = false;
process.stdin.setEncoding("utf8");
process.stdin.on("data", (chunk) => {
buffer += chunk;
+ drain();
+});
+
+function drain() {
let idx;
- while ((idx = buffer.indexOf("\n")) !== -1) {
+ while (!stalled && (idx = buffer.indexOf("\n")) !== -1) {
const line = buffer.slice(0, idx);
buffer = buffer.slice(idx + 1);
if (!line) continue;
@@ -125,7 +135,15 @@ process.stdin.on("data", (chunk) => {
continue;
}
out({ id: cmd.id, type: "response", command: "prompt", success: true });
- if (busy) queue.push(cmd.message);
+ const sm = /^stall (\d+)$/.exec(cmd.message);
+ if (sm) {
+ stalled = true;
+ setTimeout(() => {
+ run("stalled");
+ stalled = false;
+ drain();
+ }, Number(sm[1]));
+ } else if (busy) queue.push(cmd.message);
else run(cmd.message);
} else if (cmd.type === "abort") {
out({ id: cmd.id, type: "response", command: "abort", success: true });
@@ -139,5 +157,5 @@ process.stdin.on("data", (chunk) => {
out({ id: cmd.id, type: "response", command: cmd.type, success: true, data: {} });
}
}
-});
+}
process.stdin.on("end", () => process.exit(0));
@@ -0,0 +1,166 @@
# #1509 6b engine fixes — Rocko R1 review
VERDICT: REQUEST CHANGES
Base supplied: `ef0020ad`. Source-only review, 2026-09-26, for Darkwing and
Sage. The three candidate hashes passed `sha256sum -c` before inspection and
again after the reproducer:
| File | SHA-256 |
|---|---|
| packages/discord/src/engine-pi.mjs | d5bf24b59c07c85067f4087c03b54ca8b4df923c1d591441dedd2e8a7ff2ae39 |
| packages/discord/tests/engine.test.mjs | f0abee9c243d46d66dd2271abc3fd89089c350ac6a66ab49131bce80adfcdc33 |
| packages/discord/tests/fake-pi.mjs | fa1bf44e3f33eb970a714ada1c686abbc1932baaf418679276edf8813abbe6de |
## F1 — High, blocking: grace expiry destroys the attribution boundary
**Location:** engine-pi.mjs lines 103–110, combined with lines 167–168 and
205–221. The latter attribute events to the current pending head, without
an event-level request identifier.
The new timer infers that no observed agent_start by the deadline means pi
never started the old prompt and will never emit its end. That inference
does not follow from silence. A delayed child or delayed stdout delivery
can emit the old run's events after the grace expires.
**Deterministically reproduced against the pinned engine, without processes,
files, providers or live connector effects:**
1. Old prompt is accepted; parent has not observed agent_start.
2. Fire its client timeout. Its failed queue entry remains temporarily.
3. Enqueue a new prompt and fire the old turn's grace timer.
4. Grace removes the failed entry and sends the new prompt, making it the
live queue head.
5. Deliver the old run's delayed start, tool pair, answer/end and settled,
in order. Then deliver the new prompt's success response and own run.
6. The new promise returns **OLD RUN ANSWER**, carries **old.md** as its tool
evidence, and its actual **NEW RUN ANSWER** is dropped.
No streaming refusal is needed. This schedule preserves FIFO ordering of
the old run's events followed by the new response/run. The README assertion
that late events find no live head is false: the grace callback just put
the next prompt there. A later refusal also cannot undo an already-resolved
promise, but the reproduced accepted-response schedule is sufficient.
This is a cross-turn answer/evidence integrity failure, not merely reduced
availability. In this connector, attributed tool records may include write,
git and approval-related metadata. No actual live incident or cross-user
data exposure is claimed by this isolated reproduction.
**Required correction:** preserve the failed turn's attribution boundary
until an authoritative termination/quiescence signal, or invalidate this
child generation before any new prompt is sent. A bounded grace can fail
held work and mark the engine unavailable; it cannot establish quiescence
by removing the only correlation guard and continuing on the same stream.
If using child replacement, require confirmed old-child termination and
isolate old events from the replacement. Do not turn this review into a
live restart or automatic restart authority change.
**Required regression:** timeout before observed start, grace expiry, then
late old start/tool/end/settled with a new prompt waiting. The new promise
must either fail closed or receive only its own reply/evidence from a safely
established generation. Test both late acceptance and refusal timing.
Retain the never-started/no-events test, but make its expected recovery
consistent with the chosen fail-closed behavior.
## Answer to Sage's bounded-wait question
- **Never observed started:** R1 adds a bound, and its `mute` test passes,
but the bound is unsafe for the indistinguishable late-events case (F1).
- **Observed started, never ends:** R1 deliberately leaves engine busy and
retains its pending tombstone. Each held prompt's own arrival-time timeout
still rejects it. This bounds each caller's wait, not recovery of engine
availability. That limitation is disclosed and inherited through
state.busy; it is not the new blocking finding.
- **Exit/settle/refused-send cleanup:** release clears the new grace timer
at the enumerated removal points. Keeping failed pending entries in
engineBusy is sound while their correlation boundary is retained.
## Verification and acceptable parts
`timeout 30s node --test packages/discord/tests/engine.test.mjs` completed
with **13/13 passing**, 0 failures, about 1.47 seconds. No force-exit was
used. The suite's existing assertions do not exercise F1.
The withEngine finally cleanup, empty command-log fallback, and polling
instead of a fixed startup sleep address the reported test cleanup/race
defects. Those changes do not require reversal. I did not repeat Darkwing's
union/eight-suite runs after finding the blocking integrity failure; broader
green results cannot cover the missing event schedule.
Only this report was added to the repository. The targeted suite created
its usual isolated temporary fixtures; the additional reproducer below is
entirely in memory. No source edits, commits, credentials, provider calls,
connector restart, or live-state changes were made. R1 is **not approved**
for the pending commit/restart gate.
## Standalone in-memory reproducer
Run from the canonical repository against the pinned engine with
`node --input-type=module` (stdin). It asserts the unsafe observed result;
after a fix this assertion must no longer hold.
```js
import { EventEmitter } from 'node:events';
import { PassThrough } from 'node:stream';
import assert from 'node:assert/strict';
import { createEngine } from './packages/discord/src/engine-pi.mjs';
const timers = [], commands = [];
const child = new EventEmitter();
child.stdout = new PassThrough(); child.stderr = new PassThrough();
child.stdin = {
write(s) { commands.push(JSON.parse(s)); return true; }, end() {}
};
child.kill = () => { child.emit('exit', 0, null); return true; };
const engine = createEngine({
command: 'memory-only', args: [], spawn: () => child, abortGraceMs: 150,
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(t => t.ms === ms && t.active);
assert.ok(t); t.active = false; t.fn();
};
const message = text => ({
role: 'assistant', content: [{ type: 'text', text }], stopReason: 'stop'
});
const finish = text => {
emit({ type: 'turn_end', message: message(text) });
emit({ type: 'agent_end', messages: [message(text)] });
emit({ type: 'agent_settled' });
};
engine.start();
try {
const first = engine.prompt('old', { timeoutMs: 50 })
.catch(e => e.details.code);
emit({ type: 'response', id: commands[0].id,
command: 'prompt', success: true });
await Promise.resolve();
fire(50); assert.equal(await first, 'timeout');
const next = engine.prompt('new', { timeoutMs: 2000 });
fire(150);
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 }
} });
finish('OLD RUN ANSWER');
const second = commands.filter(c => c.type === 'prompt')[1];
emit({ type: 'response', id: second.id,
command: 'prompt', success: true });
emit({ type: 'agent_start' }); finish('NEW RUN ANSWER');
const r = await next;
assert.equal(r.text, 'OLD RUN ANSWER');
assert.equal(r.tools[0].path, 'old.md');
console.log({ reproduced: true, returned: r.text,
tools: r.tools.map(t => t.path), busy: engine.busy });
} finally { await engine.stop(); }
```
Observed output: `reproduced: true`, `returned: 'OLD RUN ANSWER'`,
`tools: ['old.md']`, `busy: false`.
@@ -0,0 +1,85 @@
# Discord 6b R2 review
Verdict: **approve** the pinned source change. Rocko, 2026-09-26.
No blocking source findings. The restart-rate qualification below matters
operationally; this approval does not claim an unconditional restart cap.
Base: `401cc850`. Verified before review and again after tests:
- engine-pi.mjs: `77077b7fbd5a933ffd352094eb073227c299ba47b7aea52d4e60fdc55cc7101e`
- engine.test.mjs: `47a998179c6eb46827f43ab2c6b0f6b062ef94fb47da402cfb9af8c4f180f38e`
- fake-pi.mjs: `a8e54cc3f4b670eef2c06755b63e9e6bfeb44b1b1efde3aca9bfaa91583c3ef3`
- r2.patch: `4307a879f469d2b032ea6f90d1803a6085c5e949c8f88c34329c45e6501fef39`
- README.md: `69350f299f62a1cbd9283e994ead998ca3e8669eae75d15c6b12fa07036bb23e`
## R1 F1 is closed
Grace expiry marks the child wedged before failing held turns or beginning
termination. prompt, write and sendHeld all refuse that child thereafter.
The old tombstone survives grace expiry. A late agent_end may remove it
before exit, but the irreversible wedged flag still prevents any new live
head or command. Thus the old answer and tool evidence have no new prompt
to settle. Late accepted responses likewise cannot reopen sending.
The tests exercise both prompt-response timings with deterministic timers,
late start/tools/end/settle while the child is still alive, exact write
history (old plus abort), held engine-wedged, new engine-down, TERM then
KILL, and onExit. The real-child stall and mute cases also pass. The
started-run case proves the grace does not kill a run already observed
running. Cleanup releases grace timers when pending turns leave.
The no-start bound is 30 seconds after client-side failure to quarantine,
then a TERM attempt and KILL after five seconds if exit has not arrived.
Held prompts may fail earlier at their own deadlines. Actual process exit
still depends on the OS. Started runs that never settle retain the disclosed
HEAD behavior: later prompts time out individually, without a new engine
watchdog. This is an honest bounded change, not a claim of total liveness.
## Operational findings requested by Sage
1. **Wedge selects exit 1, not 3.** Engine exit invokes cli.mjs onExit,
which calls shutdown(1). Shutdown waits for connector.stop and then
process.exit(exitCode); the connector waits for failed turn handlers and
delivery work and stops the already-exited engine. Exit 3 is the separate
supervised startup refusal, not this wedge path. If an operator shutdown
or other fatal shutdown began first, the existing first-shutdown-wins
guard retains that earlier code; it still does not turn a wedge into 3.
A later supervised startup can independently refuse with 3 if STOP or a
held binding is present, which is the intended operator brake.
2. **Medium operational qualification: the start limit is rate-based.**
The template sets Restart=on-failure, RestartSec=15,
StartLimitIntervalSec=600, StartLimitBurst=5 and
RestartPreventExitStatus=3. Rapid failures that exhaust that start budget
stop automatic restarting. Repeated wedges do NOT necessarily exhaust
it. The default turn timeout is 180 seconds; adding 30 seconds grace and
15 seconds restart delay gives at least 225 seconds per cycle, even
before request arrival, shutdown work or KILL delay. Such cycles remain
below five starts per 600 seconds and may continue indefinitely.
Therefore “every repeatedly wedging pi eventually stays down” is not
established by this unit. This is an existing unit policy limitation,
not a defect in R2's event isolation. If a cumulative wedge circuit
breaker is required, specify it separately; do not treat this rate limit
as one. No live unit was started or fault-injected for this review.
3. **Recommended follow-up: journal startup source identity.** Yes: record
repository HEAD, dirty status scoped to runtime source, and preferably a
digest of the actual runtime files along with startup time and context
digest. A commit id alone conceals uncommitted source; a dirty boolean
alone cannot identify which candidate ran. Do not log source contents or
secret values. The live-checkout process rule remains necessary, because
startup fingerprints do not prevent an edit during module loading or a
later dynamic read. Restart-on-failure was already configured for other
failures; R2 adds a new way to reach it. This follow-up is nonblocking.
## Independent verification
`timeout 30s node --test packages/discord/tests/engine.test.mjs`:
17 passed, zero failed, exit 0. Manifest checks passed before and after.
Reviewed the actual CLI shutdown and connector.stop paths and the versioned
unit template. Author-reported 406-test repetitions and eight suites were
not rerun or represented as independent evidence.
Only this report was written. No source edits, commit, restart, live message
or service-policy change. Approval is for the pinned R2 source; integration
and any authorized restart remain with the lead.
+98 -40
View File
@@ -16,9 +16,14 @@
// from `tool_execution_start`/`tool_execution_end` into the result so the
// turn record shows what was read. An `agent_end` with `willRetry` is not
// the end of the run. A timeout sends `abort` and fails that turn; the
// process stays. A malformed JSONL line from pi fails the current turn (its
// outcome is now unknowable) and the process stays. Process exit fails
// every pending turn and is reported through `onExit`.
// process stays. The failed turn holds later prompts back until its
// agent_end or a settle. If pi has not started it within ABORT_GRACE_MS, the
// engine stops pi instead of sending again: pi's events carry no prompt id,
// so a late run of the failed prompt would be taken for the next one's. A
// run pi did start holds later prompts until it ends, as any run does. A
// malformed JSONL line from pi fails the current turn (its outcome is now
// unknowable) and the process stays. Process exit fails every pending turn
// and is reported through `onExit`.
//
// Framing follows pi's RPC doc: split on "\n" only, strip a trailing "\r".
// Node readline is not used because it also splits on U+2028/U+2029.
@@ -64,31 +69,67 @@ export function assistantText(message) {
.trim();
}
// How long a turn that failed here (timeout, protocol error) may wait for
// pi's agent_start before the engine stops pi. Without a bound, a prompt pi
// accepted but never ran would hold every later prompt until restart.
export const ABORT_GRACE_MS = 30000;
export function createEngine({
command, args, cwd, env = {},
spawn = nodeSpawn,
setTimeoutImpl = globalThis.setTimeout, clearTimeoutImpl = globalThis.clearTimeout,
log = () => {},
onExit = () => {},
abortGraceMs = ABORT_GRACE_MS,
} = {}) {
if (typeof command !== "string" || command.length === 0) throw new DiscordError("engine: command required", 1);
if (!Array.isArray(args)) throw new DiscordError("engine: args required", 1);
// pending: prompts sent to pi, oldest first. held: prompts waiting for pi
// to settle before they are sent, oldest first.
const state = { child: null, buffer: "", pending: [], held: [], responses: new Map(), nextId: 1, busy: false, exited: null };
// wedged: set when the engine gave up on pi and is stopping it. Nothing
// is sent to that child again.
const state = { child: null, buffer: "", pending: [], held: [], responses: new Map(), nextId: 1, busy: false, exited: null, wedged: false };
// A turn that fails on the client side (timeout, protocol error) stays in
// the pending queue, marked done, until pi's own turn_end for it arrives.
// Otherwise that turn_end would be attributed to the next prompt.
// Otherwise that turn_end would be attributed to the next prompt. It holds
// later prompts back; if pi has not started it within abortGraceMs, the
// engine stops pi.
function failTurn(turn, code, message) {
if (turn.done) return;
turn.done = true;
if (turn.timer !== null) clearTimeoutImpl(turn.timer);
turn.timer = null;
if (state.pending.includes(turn)) {
turn.grace = setTimeoutImpl(() => {
turn.grace = null;
// A run pi started keeps its place until its agent_end or a settle.
if (state.busy) return;
// No agent_start yet. Pi may never run this prompt, or its events
// may still be on the way; with no prompt id in them, nothing sent
// now could be told apart from it. Stop pi: held prompts fail, and
// the exit fails the rest and reaches onExit.
log(`engine: no agent_start ${abortGraceMs} ms after a failed turn; stopping pi`);
wedge();
}, abortGraceMs);
}
turn.reject(new DiscordError(message, 1, { code }));
}
function wedge() {
if (state.wedged || state.exited !== null) return;
state.wedged = true;
for (const h of state.held.splice(0)) failTurn(h.turn, "engine-wedged", "engine stopped: pi did not start an aborted turn");
stopChild();
}
// Call when a turn leaves the pending queue.
function release(turn) {
if (turn.grace !== null) clearTimeoutImpl(turn.grace);
turn.grace = null;
}
function settleTurn(turn, value) {
if (turn.done) return;
turn.done = true;
@@ -99,7 +140,10 @@ export function createEngine({
function failAll(code, message) {
const pending = state.pending.splice(0);
for (const t of pending) failTurn(t, code, message);
for (const t of pending) {
release(t);
failTurn(t, code, message);
}
for (const h of state.held.splice(0)) failTurn(h.turn, code, message);
for (const [, r] of state.responses) r.reject(new DiscordError(message, 1, { code }));
state.responses.clear();
@@ -174,6 +218,7 @@ export function createEngine({
// Attribute the run to the head even if it failed client-side, so the
// next prompt's agent_end is not taken for this one.
const run = state.pending.shift();
if (run) release(run);
if (!run || run.done) return;
const messages = Array.isArray(event.messages) ? event.messages.filter((m) => m && m.role === "assistant") : [];
const message = messages.length > 0 ? messages[messages.length - 1] : run.last;
@@ -193,13 +238,14 @@ export function createEngine({
// this settle and still has no agent_end will never get one: fail it now
// instead of waiting for its timeout. Turns whose prompt response has
// not arrived yet belong to a later run and stay.
const dropped = [];
const keep = [];
for (const t of state.pending) {
if (t.done) continue;
if (t.accepted) failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
else keep.push(t);
}
for (const t of state.pending) (t.done || t.accepted ? dropped : keep).push(t);
state.pending = keep;
for (const t of dropped) {
release(t);
failTurn(t, "engine-settled-without-turn", "engine settled without answering this prompt");
}
sendHeld();
}
}
@@ -214,21 +260,29 @@ export function createEngine({
// Never accepted: pi will not emit a turn_end for it, so remove it.
const i = state.pending.indexOf(turn);
if (i !== -1) state.pending.splice(i, 1);
release(turn);
failTurn(turn, (err.details && err.details.code) || "engine-refused", err.message);
sendHeld();
});
}
// Pi is busy from our side while any sent prompt is still queued, even one
// that already failed here: a turn that timed out before its agent_start
// was read leaves state.busy false while pi runs it, and sending then would
// be refused as streaming. It leaves the queue on its agent_end, on a
// settle, on a refused send, or at process exit.
const engineBusy = () => state.busy || state.pending.length > 0;
// After a settle (or a refused send) the oldest held prompt goes out.
function sendHeld() {
if (state.exited !== null) return;
if (state.busy || state.pending.some((t) => !t.done)) return;
if (state.exited !== null || state.wedged) return;
if (engineBusy()) return;
const next = state.held.shift();
if (next) send(next.turn, next.command);
}
function write(command) {
if (!state.child || state.exited !== null) throw new DiscordError("engine is not running", 1, { code: "engine-down" });
if (!state.child || state.exited !== null || state.wedged) throw new DiscordError("engine is not running", 1, { code: "engine-down" });
state.child.stdin.write(JSON.stringify(command) + "\n");
}
@@ -245,6 +299,30 @@ export function createEngine({
});
}
function stopChild({ graceMs = 5000 } = {}) {
const child = state.child;
if (!child || state.exited !== null) return Promise.resolve(state.exited);
return new Promise((resolve) => {
const timer = setTimeoutImpl(() => {
try {
child.kill("SIGKILL");
} catch {
// already gone
}
}, graceMs);
child.once("exit", () => {
clearTimeoutImpl(timer);
resolve(state.exited);
});
try {
child.stdin.end();
child.kill("SIGTERM");
} catch {
// already gone
}
});
}
return {
start() {
if (state.child) throw new DiscordError("engine already started", 1);
@@ -281,7 +359,7 @@ export function createEngine({
// with DiscordError carrying details.code for the turn record.
prompt(text, { timeoutMs = 180000 } = {}) {
if (typeof text !== "string" || text.length === 0) throw new DiscordError("prompt text required", 1);
const turn = { resolve: null, reject: null, timer: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
const turn = { resolve: null, reject: null, timer: null, grace: null, done: false, accepted: false, tools: new Map(), turns: 0, last: null };
const done = new Promise((resolve, reject) => {
turn.resolve = resolve;
turn.reject = reject;
@@ -306,44 +384,24 @@ export function createEngine({
}
failTurn(turn, "timeout", `turn timed out after ${timeoutMs} ms`);
}, timeoutMs);
if (state.exited !== null) {
if (state.exited !== null || state.wedged) {
failTurn(turn, "engine-down", "engine is not running");
return done;
}
if (state.busy || state.pending.some((t) => !t.done) || state.held.length > 0) state.held.push({ turn, command });
if (engineBusy() || state.held.length > 0) state.held.push({ turn, command });
else send(turn, command);
return done;
},
get busy() {
return state.busy || state.pending.some((t) => !t.done) || state.held.length > 0;
return engineBusy() || state.held.length > 0;
},
get pendingCount() {
return state.pending.filter((t) => !t.done).length + state.held.length;
},
stop({ graceMs = 5000 } = {}) {
const child = state.child;
if (!child || state.exited !== null) return Promise.resolve(state.exited);
return new Promise((resolve) => {
const timer = setTimeoutImpl(() => {
try {
child.kill("SIGKILL");
} catch {
// already gone
}
}, graceMs);
child.once("exit", () => {
clearTimeoutImpl(timer);
resolve(state.exited);
});
try {
child.stdin.end();
child.kill("SIGTERM");
} catch {
// already gone
}
});
stop(options) {
return stopChild(options);
},
};
}
+191 -59
View File
@@ -4,6 +4,8 @@ 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");
@@ -28,7 +30,20 @@ function start(root, extra = {}) {
log: (m) => logs.push(m), ...extra,
});
engine.start();
return { engine, logs, commands: () => readFileSync(logPath, "utf8").trim().split("\n").filter(Boolean).map((l) => JSON.parse(l)) };
// 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", () => {
@@ -56,8 +71,7 @@ test("engine: with tools, buildPiArgs turns pi's own tools off, loads the extens
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", async () => {
const { engine } = start(makeRoot());
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);
@@ -71,45 +85,38 @@ test("engine: a run with tool turns settles once, on the answer, with every tool
assert.equal(plain.turns, 1);
await idle(engine);
assert.equal(engine.busy, false);
await engine.stop();
});
}));
test("engine: a run that ends on a tool-only turn fails the prompt as empty; a retried run settles on the real end", async () => {
const { engine } = start(makeRoot());
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");
await engine.stop();
});
}));
test("engine: one prompt, one turn, text and usage come back", async () => {
const { engine } = start(makeRoot());
try {
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);
} finally {
await engine.stop();
}
});
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", async () => {
const { engine, commands } = start(makeRoot());
const first = engine.prompt("slow 150");
await new Promise((r) => setTimeout(r, 20));
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);
await new Promise((r) => setTimeout(r, 20));
assert.equal(commands().filter((c) => c.type === "prompt").length, 1, "the second prompt is not sent while pi is busy");
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 = commands().filter((c) => c.type === "prompt");
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).
@@ -117,11 +124,9 @@ test("engine: a prompt while streaming is held until pi settles, then sent as it
assert.equal(prompts[1].streamingBehavior, undefined);
await idle(engine);
assert.equal(engine.busy, false);
await engine.stop();
});
}));
test("engine: a held prompt that times out before pi settles fails on its own and is never sent", async () => {
const { engine, commands } = start(makeRoot());
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));
@@ -130,50 +135,177 @@ test("engine: a held prompt that times out before pi settles fails on its own an
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");
await engine.stop();
});
}));
test("engine: timeout sends abort and fails only that turn; the process stays", async () => {
const { engine, commands, logs } = start(makeRoot());
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");
await engine.stop();
}));
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: tool events from a run that outlived its timeout never land in the next prompt's record", async () => {
const { engine } = start(makeRoot());
try {
await assert.rejects(engine.prompt("late 200", { timeoutMs: 40 }), (err) => err.details.code === "timeout");
const r = await engine.prompt("after late");
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, "the dead run's turns are not counted here");
} finally {
await engine.stop();
}
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);
});
});
test("engine: a malformed JSONL line fails the turn, not the process", async () => {
const { engine, logs } = start(makeRoot());
// "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");
await engine.stop();
});
}));
test("engine: a turn that ends in error rejects with the error code; process exit fails pending turns", async () => {
const root = makeRoot();
let exited = null;
const { engine } = start(root, { onExit: (e) => (exited = e) });
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/);
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();
}
});
+21 -3
View File
@@ -8,9 +8,13 @@
// then a second turn that answers "read <n> file(s)"
// "toolonly" a run whose only turn calls a tool and never answers
// "retry" an agent_end with willRetry, then the real answer
// "mute" accept the prompt and emit nothing, staying idle
// "late <ms>" ignore abort; after <ms> emit a tool pair and a tool turn,
// then answer "late reply", like a run that outlives its
// client-side timeout
// "stall <ms>" accept the prompt, then read nothing for <ms> (pi stuck
// before agent_start); then run it, answering "echo:
// stalled", and only then read what came in meanwhile
// anything else answer "echo: <text>" immediately
// A prompt received while busy without streamingBehavior is refused, as pi
// does. A prompt with streamingBehavior followUp is folded into the running
@@ -30,6 +34,7 @@ function assistant(text, stopReason = "stop") {
}
function run(text) {
if (text === "mute") return;
busy = true;
out({ type: "agent_start" });
out({ type: "turn_start" });
@@ -109,11 +114,16 @@ function run(text) {
let current = null;
let buffer = "";
let stalled = false;
process.stdin.setEncoding("utf8");
process.stdin.on("data", (chunk) => {
buffer += chunk;
drain();
});
function drain() {
let idx;
while ((idx = buffer.indexOf("\n")) !== -1) {
while (!stalled && (idx = buffer.indexOf("\n")) !== -1) {
const line = buffer.slice(0, idx);
buffer = buffer.slice(idx + 1);
if (!line) continue;
@@ -125,7 +135,15 @@ process.stdin.on("data", (chunk) => {
continue;
}
out({ id: cmd.id, type: "response", command: "prompt", success: true });
if (busy) queue.push(cmd.message);
const sm = /^stall (\d+)$/.exec(cmd.message);
if (sm) {
stalled = true;
setTimeout(() => {
run("stalled");
stalled = false;
drain();
}, Number(sm[1]));
} else if (busy) queue.push(cmd.message);
else run(cmd.message);
} else if (cmd.type === "abort") {
out({ id: cmd.id, type: "response", command: "abort", success: true });
@@ -139,5 +157,5 @@ process.stdin.on("data", (chunk) => {
out({ id: cmd.id, type: "response", command: cmd.type, success: true, data: {} });
}
}
});
}
process.stdin.on("end", () => process.exit(0));