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]>
162 lines
6.9 KiB
JavaScript
162 lines
6.9 KiB
JavaScript
// A stand-in for `pi --mode rpc`. Reads JSONL commands on stdin, writes
|
|
// events on stdout. Behaviour is scripted per prompt text:
|
|
// "slow <ms>" answer "slow reply" after <ms>
|
|
// "garbage" emit one malformed line
|
|
// "error" end the turn with stopReason error
|
|
// "tools <n>" a first turn that calls <n> tools (read_file, with a
|
|
// tool_execution_start/end pair each, the last one refused),
|
|
// 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
|
|
// loop as real pi does: answered inside the same run, one agent_end for
|
|
// both, no agent_start of its own. The engine must therefore never send
|
|
// one; this fake makes that visible. Every command is mirrored to
|
|
// FAKE_PI_LOG when set.
|
|
import { appendFileSync } from "node:fs";
|
|
|
|
const logPath = process.env.FAKE_PI_LOG;
|
|
const out = (o) => process.stdout.write(JSON.stringify(o) + "\n");
|
|
let busy = false;
|
|
const queue = [];
|
|
|
|
function assistant(text, stopReason = "stop") {
|
|
return { role: "assistant", content: text === null ? [] : [{ type: "text", text }], stopReason, usage: { input: 3, output: 2 }, model: "fake", provider: "fake" };
|
|
}
|
|
|
|
function run(text) {
|
|
if (text === "mute") return;
|
|
busy = true;
|
|
out({ type: "agent_start" });
|
|
out({ type: "turn_start" });
|
|
const finish = (message) => {
|
|
out({ type: "turn_end", message, toolResults: [] });
|
|
out({ type: "agent_end", messages: [message] });
|
|
if (queue.length > 0) {
|
|
// Real pi: the follow-up continues this run; both answers are inside
|
|
// it and agent_end above already covered them. Just settle.
|
|
queue.length = 0;
|
|
}
|
|
busy = false;
|
|
out({ type: "agent_settled" });
|
|
};
|
|
const tm = /^tools (\d+)$/.exec(text);
|
|
if (tm || text === "toolonly") {
|
|
const n = tm ? Number(tm[1]) : 1;
|
|
const calls = [];
|
|
for (let i = 1; i <= n; i += 1) {
|
|
const id = `call_${i}`;
|
|
const last = i === n && n > 1;
|
|
calls.push({ type: "toolCall", id, name: "read_file", arguments: { root: "docs", path: `f${i}.md` } });
|
|
out({ type: "tool_execution_start", toolCallId: id, toolName: "read_file", args: { root: "docs", path: `f${i}.md` } });
|
|
out({ type: "tool_execution_end", toolCallId: id, toolName: "read_file", isError: false, result: { content: [{ type: "text", text: last ? "refused: budget" : "1: hello" }], details: last ? { tool: "read_file", root: "docs", path: `f${i}.md`, ok: false, reason: "tool budget for this message is used up", ms: 1 } : { tool: "read_file", root: "docs", path: `f${i}.md`, ok: true, bytes: 9, ms: 2 } } });
|
|
}
|
|
const toolTurn = { role: "assistant", content: calls, stopReason: "toolUse", usage: { input: 3, output: 2 }, model: "fake", provider: "fake" };
|
|
out({ type: "turn_end", message: toolTurn, toolResults: [] });
|
|
if (text === "toolonly") {
|
|
out({ type: "agent_end", messages: [toolTurn] });
|
|
busy = false;
|
|
out({ type: "agent_settled" });
|
|
return;
|
|
}
|
|
out({ type: "turn_start" });
|
|
const answer = assistant(`read ${n} file(s)`);
|
|
out({ type: "turn_end", message: answer, toolResults: [] });
|
|
out({ type: "agent_end", messages: [toolTurn, answer] });
|
|
queue.length = 0;
|
|
busy = false;
|
|
out({ type: "agent_settled" });
|
|
return;
|
|
}
|
|
const lm = /^late (\d+)$/.exec(text);
|
|
if (lm) {
|
|
setTimeout(() => {
|
|
const args = { root: "docs", path: "late.md" };
|
|
out({ type: "tool_execution_start", toolCallId: "call_late", toolName: "read_file", args });
|
|
out({ type: "tool_execution_end", toolCallId: "call_late", toolName: "read_file", isError: false, result: { content: [{ type: "text", text: "1: late" }], details: { tool: "read_file", ...args, ok: true, bytes: 5, ms: 1 } } });
|
|
out({ type: "turn_end", message: { role: "assistant", content: [{ type: "toolCall", id: "call_late", name: "read_file", arguments: args }], stopReason: "toolUse", usage: { input: 3, output: 2 }, model: "fake", provider: "fake" }, toolResults: [] });
|
|
out({ type: "turn_start" });
|
|
finish(assistant("late reply"));
|
|
}, Number(lm[1]));
|
|
return;
|
|
}
|
|
if (text === "retry") {
|
|
out({ type: "agent_end", messages: [], willRetry: true });
|
|
finish(assistant("after retry"));
|
|
return;
|
|
}
|
|
const m = /^slow (\d+)$/.exec(text);
|
|
if (m) {
|
|
const timer = setTimeout(() => finish(assistant("slow reply")), Number(m[1]));
|
|
current = { timer, finish };
|
|
return;
|
|
}
|
|
if (text === "garbage") {
|
|
process.stdout.write("this is not json\n");
|
|
finish(assistant("after garbage"));
|
|
return;
|
|
}
|
|
if (text === "error") {
|
|
finish({ ...assistant(null, "error"), errorMessage: "fake provider error" });
|
|
return;
|
|
}
|
|
finish(assistant(`echo: ${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 (!stalled && (idx = buffer.indexOf("\n")) !== -1) {
|
|
const line = buffer.slice(0, idx);
|
|
buffer = buffer.slice(idx + 1);
|
|
if (!line) continue;
|
|
const cmd = JSON.parse(line);
|
|
if (logPath) appendFileSync(logPath, line + "\n");
|
|
if (cmd.type === "prompt") {
|
|
if (busy && !cmd.streamingBehavior) {
|
|
out({ id: cmd.id, type: "response", command: "prompt", success: false, error: "agent is streaming; specify streamingBehavior" });
|
|
continue;
|
|
}
|
|
out({ id: cmd.id, type: "response", command: "prompt", success: true });
|
|
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 });
|
|
if (current) {
|
|
clearTimeout(current.timer);
|
|
const f = current.finish;
|
|
current = null;
|
|
f(assistant("", "aborted"));
|
|
}
|
|
} else {
|
|
out({ id: cmd.id, type: "response", command: cmd.type, success: true, data: {} });
|
|
}
|
|
}
|
|
}
|
|
process.stdin.on("end", () => process.exit(0));
|