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 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 " ignore abort; after 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" });