// The engine: one `pi --mode rpc` child per binding, one conversation, one // turn at a time from the connector's point of view. A prompt sent while pi // is busy is held here and sent when pi settles, so a second Discord // message during a turn is neither lost nor run concurrently. It is not // sent as a pi follow-up (streamingBehavior followUp): pi folds a follow-up // into the running agent loop and closes both answers with one agent_end, // which lost the first answer live on 2026-09-17 (the second message's // reply was posted as the first message's, the second failed as settled // without a turn). One prompt per run keeps the events unambiguous. // // Each prompt resolves on the `agent_end` event that closes its run (one // run per prompt, in the order prompts were sent). A run holds one or more // pi turns: with tools, an assistant message that only calls tools ends a // turn and the next turn carries the answer. The reply is the last // assistant message of the run; every tool call in between is collected // 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. 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. // // This module can be replaced by the CHAT-03 conversation controller later // without the connector noticing: the contract is start(), prompt(), stop(). import { spawn as nodeSpawn } from "node:child_process"; import { fileURLToPath } from "node:url"; import { DiscordError } from "./errors.mjs"; import { enabledToolNames } from "./tools.mjs"; export const PI_FIXED_ARGS = Object.freeze([ "--mode", "rpc", "--no-extensions", "--no-context-files", "--no-skills", "--no-prompt-templates", "--no-themes", "--offline", ]); // Without tools: pi's own tools off, no extension. With tools: pi's // built-in tools off, our extension loaded explicitly, and an allowlist of // exactly its tool names. --no-extensions stays in both cases; it disables // discovery, not an explicit --extension. export const PI_NO_TOOLS_ARGS = Object.freeze(["--no-tools"]); export const TOOLS_EXTENSION = fileURLToPath(new URL("../extension/tools.mjs", import.meta.url)); // Older name, kept for the tests and docs that still use it. export const READONLY_TOOLS_EXTENSION = TOOLS_EXTENSION; export function buildPiArgs({ provider, model, thinking, sessionDir, appendSystemPromptFile, continueSession, tools = null }) { const args = [...PI_FIXED_ARGS]; if (tools) args.push("--no-builtin-tools", "--extension", TOOLS_EXTENSION, "--tools", enabledToolNames(tools).join(",")); else args.push(...PI_NO_TOOLS_ARGS); args.push("--provider", provider, "--model", model); if (thinking) args.push("--thinking", thinking); args.push("--session-dir", sessionDir, "--append-system-prompt", appendSystemPromptFile); if (continueSession) args.push("--continue"); return args; } export function assistantText(message) { if (!message || !Array.isArray(message.content)) return ""; return message.content .filter((c) => c && c.type === "text" && typeof c.text === "string") .map((c) => c.text) .join("") .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. // 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. 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; if (turn.timer !== null) clearTimeoutImpl(turn.timer); turn.timer = null; turn.resolve(value); } function failAll(code, message) { const pending = state.pending.splice(0); 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(); } function handleLine(line) { let event; try { event = JSON.parse(line); } catch { log("engine: malformed JSONL line from pi"); const head = state.pending.find((t) => !t.done); if (head) failTurn(head, "engine-protocol", "engine emitted a malformed line during the turn"); return; } if (!event || typeof event !== "object") return; if (event.type === "response") { const waiter = event.id !== undefined ? state.responses.get(event.id) : undefined; if (waiter) { state.responses.delete(event.id); if (event.success === false) waiter.reject(new DiscordError(`engine refused ${event.command}: ${event.error || "unknown error"}`, 1, { code: "engine-refused" })); else waiter.resolve(event.data); } return; } if (event.type === "agent_start") state.busy = true; // Tool and turn events belong to the run pi is executing, which is the // oldest queued prompt: agent_end shifts it off, and a prompt that failed // client-side (a timeout) stays at the front, marked done, until then. // Events while that front prompt is done belong to the dead run and are // dropped, so they never become another prompt's evidence. const front = state.pending[0] || null; const head = front && !front.done ? front : null; if (event.type === "tool_execution_start" && head) { head.tools.set(event.toolCallId, { name: event.toolName, startedAt: Date.now(), args: event.args || {} }); return; } if (event.type === "tool_execution_end" && head) { const open = head.tools.get(event.toolCallId) || { name: event.toolName, startedAt: Date.now(), args: {} }; const d = (event.result && event.result.details) || {}; head.tools.set(event.toolCallId, { ...open, done: true, record: { name: event.toolName, root: d.root ?? (typeof open.args.root === "string" ? open.args.root : null), path: d.path ?? (typeof open.args.path === "string" ? open.args.path : null), ...(d.url !== undefined || typeof open.args.url === "string" ? { url: d.url ?? open.args.url } : {}), ...(d.query !== undefined || typeof open.args.query === "string" ? { query: d.query ?? open.args.query } : {}), ...(d.status !== undefined ? { status: d.status } : {}), ...(d.hash !== undefined ? { hash: d.hash } : {}), ...(d.pushed !== undefined ? { pushed: d.pushed } : {}), ...(d.paths !== undefined ? { paths: d.paths } : {}), ...(d.requester !== undefined ? { requester: d.requester } : {}), ...(d.id !== undefined ? { id: d.id } : {}), ...(d.request !== undefined ? { request: d.request } : {}), ...(d.verb !== undefined ? { verb: d.verb } : {}), ...(d.key !== undefined ? { key: d.key } : {}), ...(d.revision !== undefined ? { revision: d.revision } : {}), ...(d.code !== undefined ? { code: d.code } : {}), ok: event.isError ? false : d.ok !== false, reason: d.reason ?? (event.isError ? "tool error" : null), bytes: d.bytes ?? null, ms: d.ms ?? Date.now() - open.startedAt, }, }); return; } if (event.type === "turn_end" && head) { head.turns += 1; head.last = event.message || head.last; return; } if (event.type === "agent_end") { if (event.willRetry === true) return; // 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; const tools = [...run.tools.values()].filter((t) => t.done).map((t) => t.record); const text = assistantText(message); const stopReason = message && message.stopReason; if (stopReason === "error" || stopReason === "aborted") { failTurn(run, `engine-${stopReason}`, `engine turn ended with ${stopReason}: ${(message && message.errorMessage) || ""}`.trim()); return; } settleTurn(run, { text, message, tools, turns: run.turns, usage: (message && message.usage) || null, model: message ? message.model : null, provider: message ? message.provider : null }); return; } if (event.type === "agent_settled") { state.busy = false; // A settle means pi has nothing queued. A turn that was accepted before // 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) (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(); } } // Send one prompt to pi and queue it as pending. Only called when pi is // idle from our point of view, so no streamingBehavior is ever needed. function send(turn, command) { state.pending.push(turn); request(command).then(() => { turn.accepted = true; }, (err) => { // 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 || 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 || state.wedged) throw new DiscordError("engine is not running", 1, { code: "engine-down" }); state.child.stdin.write(JSON.stringify(command) + "\n"); } function request(command) { const id = `r${state.nextId++}`; return new Promise((resolve, reject) => { state.responses.set(id, { resolve, reject }); try { write({ ...command, id }); } catch (err) { state.responses.delete(id); reject(err); } }); } 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); const child = spawn(command, args, { cwd, env, stdio: ["pipe", "pipe", "pipe"] }); state.child = child; child.stdout.setEncoding("utf8"); child.stdout.on("data", (chunk) => { state.buffer += chunk; let idx; while ((idx = state.buffer.indexOf("\n")) !== -1) { let line = state.buffer.slice(0, idx); state.buffer = state.buffer.slice(idx + 1); if (line.endsWith("\r")) line = line.slice(0, -1); if (line.length > 0) handleLine(line); } }); child.stderr.setEncoding("utf8"); child.stderr.on("data", (chunk) => log(`pi: ${chunk.trimEnd()}`)); child.on("error", (err) => { log(`engine spawn error: ${err.message}`); state.exited = { code: null, signal: null, error: err.message }; failAll("engine-down", `engine failed: ${err.message}`); onExit(state.exited); }); child.on("exit", (code, signal) => { state.exited = { code, signal }; failAll("engine-down", `engine exited (code ${code}, signal ${signal})`); onExit(state.exited); }); return child; }, // Resolves {text, message, tools, turns, usage, model, provider}. Rejects // 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, 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; }); const command = { type: "prompt", message: text }; // The clock starts on arrival, held time included: a message that // waits behind a long turn still fails after timeoutMs, and a held // turn that times out is simply never sent. turn.timer = setTimeoutImpl(() => { if (turn.done) return; const heldAt = state.held.findIndex((h) => h.turn === turn); if (heldAt !== -1) { state.held.splice(heldAt, 1); failTurn(turn, "timeout", `turn timed out after ${timeoutMs} ms while waiting for the engine`); return; } log(`engine: turn timed out after ${timeoutMs} ms, aborting`); try { write({ type: "abort" }); } catch (err) { log(`engine: abort failed: ${err.message}`); } failTurn(turn, "timeout", `turn timed out after ${timeoutMs} ms`); }, timeoutMs); if (state.exited !== null || state.wedged) { failTurn(turn, "engine-down", "engine is not running"); return done; } if (engineBusy() || state.held.length > 0) state.held.push({ turn, command }); else send(turn, command); return done; }, get busy() { return engineBusy() || state.held.length > 0; }, get pendingCount() { return state.pending.filter((t) => !t.done).length + state.held.length; }, stop(options) { return stopChild(options); }, }; }