Controller, claim store, live-session guard, engine link and seal, turn tracker, cohort force stop and recovery, client library, transcript and mediated terminal, with the fake engine and tests. Fixtures only; no live cutover. Dewey built it. Darkwing (comment 26690) and Filbert (comment 26694) approved round 2. Manifest I1-r2-manifest.sha256 (2b48e333, 27 files). Suites on an export: conversation 152/152, control-board 124, webui 14, seat 19, chat-00/01/01c checks, and all nine scripts/test-*.sh green. Follow-ups for I3 are in DEFERRED. Gate E stays with Jason. Co-Authored-By: Claude Opus 5.5 <[email protected]>
695 lines
27 KiB
JavaScript
695 lines
27 KiB
JavaScript
#!/usr/bin/env node
|
||
// The fake engine for CHAT-03 I1 (#1507, brief §3 N10).
|
||
//
|
||
// It models pinned Pi 0.85.1's RPC prompt path where CHAT-03 depends on it
|
||
// (agent-session.js 821–949, rpc-mode.js 298–335):
|
||
// - `prompt` awaits the input handlers (line 843), checks `isStreaming` and
|
||
// throws with no streamingBehavior (line 860), awaits the auth and
|
||
// compaction checks (895) and before_agent_start (915), acks, and then
|
||
// starts the run. A Mosaic prompt is never queued.
|
||
// - a prompt that collides with a running agent is acked, its throw is
|
||
// swallowed, it settles with no agent_start, and it leaves isStreaming
|
||
// false while the other run goes on;
|
||
// - `agent_settled` comes from a `finally`, and one run can hold several
|
||
// agent_start … agent_end pairs;
|
||
// - `clear_queue` emits an empty `queue_update` before its response,
|
||
// returns only the steer and follow-up items, drops agent-level custom
|
||
// messages without returning them, and leaves nextTurn messages;
|
||
// - `abort` waits for idle, and anything still queued then runs inside the
|
||
// same run.
|
||
// Unsealed-extension effects (queueing, extension prompts, triggerTurn,
|
||
// input handlers) are simulated here; no real extension ever loads.
|
||
//
|
||
// Two ways to run it. In process, a test builds `FakePi` on a pair of streams
|
||
// through `FakeLauncher` and drives it directly. As a process
|
||
// (`node fake-pi.mjs --mode rpc ...`), it is driven over the Unix socket in
|
||
// FAKE_PI_CONTROL and logs argv, every byte received and every line sent to
|
||
// FAKE_PI_LOG. The process form exists for the cohort fixtures (K-series),
|
||
// which need a real process tree.
|
||
|
||
import { appendFileSync, readFileSync, rmSync, writeFileSync } from "node:fs";
|
||
import { spawn } from "node:child_process";
|
||
import { createServer, connect } from "node:net";
|
||
import { PassThrough, Writable } from "node:stream";
|
||
import { randomBytes } from "node:crypto";
|
||
import { pathToFileURL } from "node:url";
|
||
import { LineSplitter, encodeLine, parseLine } from "../src/framing.mjs";
|
||
import { parseSnapshot } from "../src/pi.mjs";
|
||
|
||
export const BUSY_ERROR = "Agent is already processing. Specify streamingBehavior ('steer' or 'followUp') to queue the message.";
|
||
export const DEFAULT_SCRIPT = Object.freeze([{ text: "ok", stop: "stop" }]);
|
||
|
||
const entryId = () => randomBytes(4).toString("hex");
|
||
|
||
export class FakePi {
|
||
constructor({ input, output, argv = [], leaf, append = true, log = null, now = () => new Date() }) {
|
||
this.input = input;
|
||
this.output = output;
|
||
this.argv = argv;
|
||
this.log = log;
|
||
this.now = now;
|
||
this.append = append;
|
||
const i = argv.indexOf("--session");
|
||
this.sessionFile = i >= 0 ? argv[i + 1] : null;
|
||
this.sessionId = null;
|
||
this.leaf = null;
|
||
if (this.sessionFile) {
|
||
const parsed = parseSnapshot(readFileSync(this.sessionFile, "utf8"));
|
||
this.sessionId = parsed.header.id;
|
||
this.leaf = parsed.defaultLeaf?.entry.id ?? null;
|
||
// Pi's startup append (sdk.js 240–252): a session with no messages, or
|
||
// whose branch has no thinking_level_change, gains one when Pi starts,
|
||
// so the loaded leaf moves. (A new session with a model also gains a
|
||
// model_change; the fake has no model.)
|
||
const branch = [];
|
||
for (let r = parsed.defaultLeaf; r; r = typeof r.entry.parentId === "string" ? parsed.byId.get(r.entry.parentId) : null) branch.push(r.entry);
|
||
if (!branch.some((e) => e.type === "message") || !branch.some((e) => e.type === "thinking_level_change")) this.startupAppend = true;
|
||
}
|
||
this.leafOverride = leaf;
|
||
this.streaming = false;
|
||
this.run = null;
|
||
this.runs = [];
|
||
this.steering = [];
|
||
this.followUp = [];
|
||
this.agentQueue = [];
|
||
this.nextTurn = [];
|
||
this.scripts = [];
|
||
this.plans = [];
|
||
this.extPlans = [];
|
||
this.drop = new Map();
|
||
this.fail = new Map();
|
||
this.armed = new Map();
|
||
this.paused = new Map();
|
||
this.pauseWaiters = [];
|
||
this.idleWaiters = [];
|
||
this.received = [];
|
||
this.commands = [];
|
||
this.sent = [];
|
||
this.appends = [];
|
||
this.uiSerial = 0;
|
||
this.persistHeld = false;
|
||
this.unpersisted = [];
|
||
if (this.startupAppend && append) this.#appendEntry({ type: "thinking_level_change", thinkingLevel: "off" });
|
||
this.onPaused = null;
|
||
this.closed = false;
|
||
this.log?.({ t: "argv", argv });
|
||
const splitter = new LineSplitter((line) => this.#command(line));
|
||
input.on("data", (c) => {
|
||
this.received.push(Buffer.from(c));
|
||
this.log?.({ t: "in", b64: Buffer.from(c).toString("base64") });
|
||
splitter.push(c);
|
||
});
|
||
input.on("error", () => {});
|
||
output.on?.("error", () => {});
|
||
}
|
||
|
||
// ---- test surface -----------------------------------------------------
|
||
|
||
bytes() {
|
||
return Buffer.concat(this.received);
|
||
}
|
||
|
||
// The script for the next run a prompt starts (FIFO). Steps: {text, stop,
|
||
// errorMessage}, {pause: name}, {tool: {id, name, args, result, isError,
|
||
// hold}}, {emit: object}, {raw: string}, {continue: true} (another
|
||
// agent_start … agent_end pair, as a retry would), {hang: true} (ignores
|
||
// abort). A first step {failBeforeUser: text, raw?} ends the run in a
|
||
// failure message before any user message.
|
||
//
|
||
// Pause points (`arm`): "843", "895", "915", "run-start" (after the ack),
|
||
// "after-start" (after agent_start), "clear", "clear-response", "abort",
|
||
// "state" (holds a get_state reply), a step's {pause} name and
|
||
// `tool:<id>`. An extension prompt's points carry an `ext:` prefix.
|
||
script(steps) {
|
||
this.scripts.push(steps);
|
||
}
|
||
|
||
// How the next RPC prompt's preflight goes: {handled: true} (an input
|
||
// handler takes it: ack, no run), {error: message, at: "843"|"895"|"915"}.
|
||
plan(p) {
|
||
this.plans.push(p);
|
||
}
|
||
|
||
// Leaves the next `n` commands of `type` unanswered.
|
||
dropResponse(type, n = 1) {
|
||
this.drop.set(type, (this.drop.get(type) ?? 0) + n);
|
||
}
|
||
|
||
// Answers the next `n` commands of `type` with `success: false`, as Pi's
|
||
// RPC loop does when a handler throws.
|
||
failResponse(type, n = 1) {
|
||
this.fail.set(type, (this.fail.get(type) ?? 0) + n);
|
||
}
|
||
|
||
arm(point, times = 1) {
|
||
this.armed.set(point, (this.armed.get(point) ?? 0) + times);
|
||
}
|
||
|
||
resume(point) {
|
||
const p = this.paused.get(point);
|
||
if (!p) throw new Error(`fake-pi: not paused at ${point}`);
|
||
this.paused.delete(point);
|
||
p.resolve();
|
||
}
|
||
|
||
waitPaused(point) {
|
||
if (this.paused.has(point)) return Promise.resolve();
|
||
return new Promise((resolve) => this.pauseWaiters.push({ point, resolve }));
|
||
}
|
||
|
||
waitIdle() {
|
||
if (!this.run) return Promise.resolve();
|
||
return new Promise((resolve) => this.idleWaiters.push(resolve));
|
||
}
|
||
|
||
// An unsealed extension's queueing. steer and followUp go into Pi's
|
||
// session queues and emit queue_update; "agent" is a custom message queued
|
||
// straight into the agent (no event, never returned by clear); "nextTurn"
|
||
// waits for the next prompt.
|
||
queue(kind, text) {
|
||
const msg = { role: kind === "steer" || kind === "followUp" ? "user" : "custom", text };
|
||
if (kind === "steer") this.steering.push(msg);
|
||
else if (kind === "followUp") this.followUp.push(msg);
|
||
else if (kind === "agent") this.agentQueue.push(msg);
|
||
else if (kind === "nextTurn") this.nextTurn.push(msg);
|
||
else throw new Error(`fake-pi: unknown queue ${kind}`);
|
||
if (kind === "steer" || kind === "followUp") this.#queueUpdate();
|
||
}
|
||
|
||
// An extension prompt. With preflight it goes through the prompt path (as
|
||
// pi.sendUserMessage would) with points named `ext:<point>`; without, it is
|
||
// a sendCustomMessage triggerTurn, which starts a run with no preflight
|
||
// (agent-session.js 1120–1121).
|
||
// `steps` is the extension run's script (default DEFAULT_SCRIPT). A
|
||
// triggerTurn while a run streams is queued straight into the agent
|
||
// (lines 1112–1118) and gives no event.
|
||
extensionPrompt({ text = "extension", preflight = true, plan = {}, steps = DEFAULT_SCRIPT } = {}) {
|
||
if (preflight) {
|
||
this.extPlans.push(plan);
|
||
return this.#prompt(null, text, "ext", steps);
|
||
}
|
||
if (this.run) {
|
||
this.agentQueue.push({ role: "custom", text });
|
||
return Promise.resolve({ queued: true });
|
||
}
|
||
return this.#runAgent({ role: "custom", text }, "ext", steps);
|
||
}
|
||
|
||
// Holds session appends after message_end until `flushPersist()`, so a
|
||
// page can be read between an event and its entry (E4).
|
||
holdPersist() {
|
||
this.persistHeld = true;
|
||
}
|
||
|
||
flushPersist() {
|
||
this.persistHeld = false;
|
||
for (const m of this.unpersisted.splice(0)) this.#persist(m);
|
||
}
|
||
|
||
dialog(method) {
|
||
this.#out({ type: "extension_ui_request", id: `ui-${++this.uiSerial}`, method, title: `fake ${method}` });
|
||
}
|
||
|
||
emit(value) {
|
||
this.#out(value);
|
||
}
|
||
|
||
raw(text) {
|
||
this.sent.push(text);
|
||
this.log?.({ t: "out", line: text });
|
||
if (!this.closed) this.output.write(text);
|
||
}
|
||
|
||
end() {
|
||
this.closed = true;
|
||
this.output.end?.();
|
||
}
|
||
|
||
// ---- protocol ----------------------------------------------------------
|
||
|
||
#out(value) {
|
||
const line = encodeLine(value);
|
||
this.raw(line);
|
||
}
|
||
|
||
#respond(id, command, success, data, error) {
|
||
const v = { id, type: "response", command, success };
|
||
if (success && data !== undefined) v.data = data;
|
||
if (!success) v.error = error;
|
||
this.#out(v);
|
||
}
|
||
|
||
async #point(name, run = null) {
|
||
const n = this.armed.get(name) ?? 0;
|
||
if (n <= 0) return;
|
||
this.armed.set(name, n - 1);
|
||
await new Promise((resolve) => {
|
||
this.paused.set(name, { resolve, run });
|
||
for (const w of this.pauseWaiters.filter((x) => x.point === name)) w.resolve();
|
||
this.pauseWaiters = this.pauseWaiters.filter((x) => x.point !== name);
|
||
this.onPaused?.(name);
|
||
});
|
||
}
|
||
|
||
#command(line) {
|
||
const parsed = parseLine(line);
|
||
if (parsed.error) return;
|
||
const cmd = parsed.value;
|
||
this.commands.push(cmd);
|
||
this.log?.({ t: "cmd", cmd: { id: cmd.id, type: cmd.type } });
|
||
const drop = this.drop.get(cmd.type) ?? 0;
|
||
if (drop > 0) {
|
||
this.drop.set(cmd.type, drop - 1);
|
||
return;
|
||
}
|
||
const fail = this.fail.get(cmd.type) ?? 0;
|
||
if (fail > 0) {
|
||
this.fail.set(cmd.type, fail - 1);
|
||
return this.#respond(cmd.id, cmd.type, false, undefined, "fake-pi: scripted failure");
|
||
}
|
||
switch (cmd.type) {
|
||
case "prompt":
|
||
if (typeof cmd.message !== "string") return this.#respond(cmd.id, "prompt", false, undefined, "message required");
|
||
void this.#prompt(cmd.id, cmd.message, null);
|
||
return;
|
||
case "abort":
|
||
void this.#abort().then(() => this.#respond(cmd.id, "abort", true));
|
||
return;
|
||
case "clear_queue":
|
||
void this.#clear(cmd.id);
|
||
return;
|
||
case "get_state": {
|
||
// Armed "state" holds the reply; the state is read when it goes out.
|
||
const reply = () => this.#respond(cmd.id, "get_state", true, {
|
||
isStreaming: this.streaming, isCompacting: false, sessionFile: this.sessionFile, sessionId: this.sessionId,
|
||
pendingMessageCount: this.steering.length + this.followUp.length, messageCount: this.appends.length,
|
||
});
|
||
if ((this.armed.get("state") ?? 0) > 0) {
|
||
void this.#point("state").then(reply);
|
||
return undefined;
|
||
}
|
||
return reply();
|
||
}
|
||
case "get_tree":
|
||
return this.#respond(cmd.id, "get_tree", true, { tree: [], leafId: this.leafOverride !== undefined ? this.leafOverride : this.leaf });
|
||
default:
|
||
return this.#respond(cmd.id, cmd.type, false, undefined, `Unknown command: ${cmd.type}`);
|
||
}
|
||
}
|
||
|
||
async #prompt(id, text, tag, steps) {
|
||
const plan = (tag ? this.extPlans : this.plans).shift() ?? {};
|
||
const P = (n) => this.#point(tag ? `${tag}:${n}` : n);
|
||
try {
|
||
await P("843");
|
||
if (plan.at === "843" && plan.error) throw new Error(plan.error);
|
||
if (plan.handled) {
|
||
if (id) this.#respond(id, "prompt", true);
|
||
return;
|
||
}
|
||
if (this.streaming) throw new Error(BUSY_ERROR);
|
||
await P("895");
|
||
if (plan.at === "895" && plan.error) throw new Error(plan.error);
|
||
await P("915");
|
||
if ((plan.at ?? "915") === "915" && plan.error) throw new Error(plan.error);
|
||
} catch (err) {
|
||
if (id) this.#respond(id, "prompt", false, undefined, err.message);
|
||
return;
|
||
}
|
||
if (id) this.#respond(id, "prompt", true);
|
||
await P("run-start");
|
||
await this.#runAgent({ role: "user", text }, tag, steps);
|
||
}
|
||
|
||
async #runAgent(first, tag = null, extSteps = DEFAULT_SCRIPT) {
|
||
if (this.run) {
|
||
// agent.prompt throws on a running agent; the throw is swallowed by
|
||
// the RPC prompt's catch, and the finally settles (agent-session.js
|
||
// 772–785).
|
||
this.streaming = false;
|
||
this.#out({ type: "agent_settled" });
|
||
return { collided: true };
|
||
}
|
||
const steps = tag ? extSteps : this.scripts.shift() ?? DEFAULT_SCRIPT;
|
||
const run = { id: this.runs.length + 1, aborted: false, hang: false, first: first.text };
|
||
this.runs.push(run);
|
||
this.run = run;
|
||
this.streaming = true;
|
||
try {
|
||
if (steps[0]?.failBeforeUser !== undefined) {
|
||
this.#out({ type: "agent_start" });
|
||
if (steps[0].raw !== undefined) this.raw(steps[0].raw);
|
||
else this.#assistantEnd({ text: "", stop: "error", errorMessage: steps[0].failBeforeUser });
|
||
this.#out({ type: "agent_end", messages: [] });
|
||
return run;
|
||
}
|
||
this.#out({ type: "agent_start" });
|
||
// Between agent_start and the first message (N19's window).
|
||
await this.#point(tag ? `${tag}:after-start` : "after-start", run);
|
||
const pending = [first, ...(first.role === "user" ? this.nextTurn.splice(0) : [])];
|
||
let current = steps;
|
||
for (;;) {
|
||
for (const m of pending.splice(0)) this.#message(m);
|
||
await this.#steps(run, current);
|
||
const next = this.steering.shift() ?? this.followUp.shift() ?? this.agentQueue.shift();
|
||
if (!next) break;
|
||
if (next.role === "user") this.#queueUpdate();
|
||
run.aborted = false;
|
||
pending.push(next);
|
||
current = DEFAULT_SCRIPT;
|
||
}
|
||
this.#out({ type: "agent_end", messages: [] });
|
||
return run;
|
||
} finally {
|
||
this.run = null;
|
||
this.streaming = false;
|
||
this.#out({ type: "agent_settled" });
|
||
for (const w of this.idleWaiters.splice(0)) w();
|
||
}
|
||
}
|
||
|
||
async #steps(run, steps) {
|
||
for (const s of steps) {
|
||
if (run.aborted) break;
|
||
if (s.pause) await this.#point(s.pause, run);
|
||
else if (s.hang) {
|
||
run.hang = true;
|
||
await new Promise(() => {});
|
||
} else if (s.emit) this.#out(s.emit);
|
||
else if (s.raw !== undefined) this.raw(s.raw);
|
||
else if (s.continue) {
|
||
this.#out({ type: "agent_end", messages: [] });
|
||
this.#out({ type: "agent_start" });
|
||
} else if (s.tool) await this.#tool(run, s.tool);
|
||
else if (s.text !== undefined || s.stop) {
|
||
this.#out({ type: "message_start", message: { role: "assistant", content: [] } });
|
||
if (s.text) this.#out({ type: "message_update", assistantMessageEvent: { type: "text_delta", contentIndex: 0, delta: s.text }, message: { role: "assistant" } });
|
||
if (s.final !== false) this.#assistantEnd(s);
|
||
}
|
||
}
|
||
if (run.aborted) {
|
||
this.#out({ type: "message_start", message: { role: "assistant", content: [] } });
|
||
this.#assistantEnd({ text: "", stop: "aborted", errorMessage: "Request was aborted" });
|
||
}
|
||
}
|
||
|
||
async #tool(run, t) {
|
||
const call = { type: "toolCall", id: t.id, name: t.name ?? "bash", arguments: t.args ?? {} };
|
||
this.#out({ type: "message_start", message: { role: "assistant", content: [] } });
|
||
const m = { role: "assistant", content: [call], stopReason: "toolUse" };
|
||
this.#out({ type: "message_end", message: m });
|
||
this.#persist(m);
|
||
this.#out({ type: "tool_execution_start", toolCallId: t.id, toolName: call.name, args: call.arguments });
|
||
if (t.hold) await this.#point(`tool:${t.id}`, run);
|
||
if (run.aborted) return;
|
||
const result = { content: [{ type: "text", text: t.result ?? "done" }] };
|
||
this.#out({ type: "tool_execution_end", toolCallId: t.id, toolName: call.name, result, isError: t.isError === true });
|
||
const r = { role: "toolResult", toolCallId: t.id, toolName: call.name, content: result.content, isError: t.isError === true };
|
||
this.#out({ type: "message_start", message: r });
|
||
this.#out({ type: "message_end", message: r });
|
||
this.#persist(r);
|
||
}
|
||
|
||
#message(m) {
|
||
const msg = m.role === "user" ? { role: "user", content: [{ type: "text", text: m.text }] } : { role: "custom", customType: "fake", content: [{ type: "text", text: m.text }], display: true };
|
||
this.#out({ type: "message_start", message: msg });
|
||
this.#out({ type: "message_end", message: msg });
|
||
this.#persist(msg);
|
||
}
|
||
|
||
#assistantEnd({ text = "", stop = "stop", errorMessage }) {
|
||
const m = { role: "assistant", content: text ? [{ type: "text", text }] : [], stopReason: stop };
|
||
if (errorMessage) m.errorMessage = errorMessage;
|
||
this.#out({ type: "message_end", message: m });
|
||
this.#persist(m);
|
||
}
|
||
|
||
// Pi persists on message_end (agent-session.js 386–398). These appends are
|
||
// the engine's own and are recorded so W11 can exclude them.
|
||
#persist(message) {
|
||
if (!this.append || !this.sessionFile || message.role === "custom") return;
|
||
if (this.persistHeld) return void this.unpersisted.push(message);
|
||
this.#appendEntry({ type: "message", message });
|
||
}
|
||
|
||
#appendEntry(fields) {
|
||
const id = entryId();
|
||
const { type, ...rest } = fields;
|
||
const entry = { type, id, parentId: this.leaf, timestamp: this.now().toISOString(), ...rest };
|
||
appendFileSync(this.sessionFile, JSON.stringify(entry) + "\n");
|
||
this.appends.push(id);
|
||
this.leaf = id;
|
||
}
|
||
|
||
#queueUpdate() {
|
||
this.#out({ type: "queue_update", steering: this.steering.map((m) => m.text), followUp: this.followUp.map((m) => m.text) });
|
||
}
|
||
|
||
async #clear(id) {
|
||
await this.#point("clear");
|
||
const steering = this.steering.splice(0).map((m) => m.text);
|
||
const followUp = this.followUp.splice(0).map((m) => m.text);
|
||
this.agentQueue.length = 0;
|
||
this.#queueUpdate();
|
||
await this.#point("clear-response");
|
||
this.#respond(id, "clear_queue", true, { steering, followUp });
|
||
}
|
||
|
||
async #abort() {
|
||
await this.#point("abort");
|
||
const run = this.run;
|
||
if (!run) return;
|
||
if (!run.hang) {
|
||
run.aborted = true;
|
||
for (const [name, p] of [...this.paused]) {
|
||
if (p.run === run) {
|
||
this.paused.delete(name);
|
||
p.resolve();
|
||
}
|
||
}
|
||
}
|
||
await this.waitIdle();
|
||
}
|
||
}
|
||
|
||
// In-process launcher: the controller gets the fake's streams. Its kind is
|
||
// "fake", so no signal ever reaches a real process. Its `forceStop` is a
|
||
// fixture stand-in for the shim: the K fixtures prove real scopes.
|
||
export class FakeLauncher {
|
||
constructor({ leaf, append = true, stdinFailAfter = null, onEngine = null, stopOutcome = "proven" } = {}) {
|
||
this.kind = "fake";
|
||
this.stopOutcome = stopOutcome;
|
||
this.opts = { leaf, append };
|
||
this.stdinFailAfter = stdinFailAfter;
|
||
this.onEngine = onEngine;
|
||
this.engines = [];
|
||
this.launches = [];
|
||
}
|
||
|
||
async launch({ unitName, command, args, cwd }) {
|
||
this.launches.push({ unitName, command, args, cwd });
|
||
const toEngine = new PassThrough();
|
||
const fromEngine = new PassThrough();
|
||
const stderr = new PassThrough();
|
||
// H19: the pipe takes `gate.left` more bytes, then fails mid-line. A
|
||
// test may set `engine.stdinGate.left` after the preflight commands.
|
||
const gate = { left: this.stdinFailAfter };
|
||
const stdin = new Writable({
|
||
write(chunk, _enc, cb) {
|
||
if (gate.left === null) {
|
||
toEngine.write(chunk);
|
||
return cb();
|
||
}
|
||
if (gate.left <= 0) return cb(Object.assign(new Error("EPIPE"), { code: "EPIPE" }));
|
||
const take = chunk.subarray(0, gate.left);
|
||
gate.left -= take.length;
|
||
toEngine.write(take);
|
||
if (take.length < chunk.length) return cb(Object.assign(new Error("EPIPE"), { code: "EPIPE" }));
|
||
return cb();
|
||
},
|
||
final(cb) {
|
||
toEngine.end();
|
||
cb();
|
||
},
|
||
});
|
||
const piArgs = args.slice(args.indexOf("--mode"));
|
||
const engine = new FakePi({ input: toEngine, output: fromEngine, argv: piArgs, ...this.opts });
|
||
engine.stdinGate = gate;
|
||
this.engines.push(engine);
|
||
this.onEngine?.(engine);
|
||
let exit;
|
||
const exited = new Promise((r) => (exit = r));
|
||
engine.kill = () => {
|
||
engine.end();
|
||
exit({ code: null, signal: "SIGKILL" });
|
||
};
|
||
const pid = 4194304 + this.engines.length;
|
||
const invocationId = `fake-inv-${this.engines.length}`;
|
||
engine.identity = { pid, invocationId, unitName };
|
||
return { kind: "fake", proc: null, stdin, stdout: fromEngine, stderr, pid, start: "1", invocationId, scope: null, shimSocket: null, exited, kill: () => engine.kill() };
|
||
}
|
||
|
||
// The fixture cohort stop. `stopOutcome` "proven" kills the fake and
|
||
// reports a complete one-member cohort; anything else is unavailable.
|
||
async forceStop({ unitName, invocationId, onPhase = async () => {} }) {
|
||
const engine = this.engines.find((e) => e.identity?.unitName === unitName && e.identity?.invocationId === invocationId);
|
||
this.stops = (this.stops ?? 0) + 1;
|
||
if (!engine) return { outcome: "unavailable", reason: "no fake engine for this unit and invocation" };
|
||
await onPhase("term");
|
||
await onPhase("kill");
|
||
if (this.stopOutcome !== "proven") return { outcome: "unavailable", reason: "fixture: cohort evidence unavailable" };
|
||
engine.kill();
|
||
const observedAt = new Date().toISOString();
|
||
return { outcome: "proven", membershipComplete: true, epoch: invocationId, observedAt, boot: "fake-boot", members: [{ pid: engine.identity.pid, boot: "fake-boot", startTicks: 1, terminatedAt: observedAt }] };
|
||
}
|
||
|
||
get last() {
|
||
return this.engines[this.engines.length - 1];
|
||
}
|
||
}
|
||
|
||
// ---- process form --------------------------------------------------------
|
||
|
||
export class ControlClient {
|
||
constructor(path) {
|
||
this.path = path;
|
||
this.serial = 0;
|
||
this.pending = new Map();
|
||
this.events = [];
|
||
this.eventWaiters = [];
|
||
}
|
||
|
||
async connect(timeoutMs = 10000) {
|
||
const end = Date.now() + timeoutMs;
|
||
for (;;) {
|
||
try {
|
||
await new Promise((resolve, reject) => {
|
||
this.sock = connect(this.path);
|
||
this.sock.once("connect", resolve);
|
||
this.sock.once("error", reject);
|
||
});
|
||
break;
|
||
} catch (err) {
|
||
if (Date.now() > end) throw err;
|
||
await new Promise((r) => setTimeout(r, 25));
|
||
}
|
||
}
|
||
const splitter = new LineSplitter((line) => {
|
||
const v = JSON.parse(line);
|
||
if (v.id !== undefined && this.pending.has(v.id)) {
|
||
this.pending.get(v.id)(v);
|
||
this.pending.delete(v.id);
|
||
} else {
|
||
this.events.push(v);
|
||
for (const w of this.eventWaiters.filter((x) => x.pred(v))) w.resolve(v);
|
||
this.eventWaiters = this.eventWaiters.filter((x) => !x.pred(v));
|
||
}
|
||
});
|
||
this.sock.on("data", (c) => splitter.push(c));
|
||
this.sock.on("error", () => {});
|
||
return this;
|
||
}
|
||
|
||
call(op, args = {}) {
|
||
const id = ++this.serial;
|
||
return new Promise((resolve) => {
|
||
this.pending.set(id, resolve);
|
||
this.sock.write(encodeLine({ id, op, ...args }));
|
||
});
|
||
}
|
||
|
||
waitEvent(pred) {
|
||
const hit = this.events.find(pred);
|
||
if (hit) return Promise.resolve(hit);
|
||
return new Promise((resolve) => this.eventWaiters.push({ pred, resolve }));
|
||
}
|
||
|
||
close() {
|
||
this.sock?.destroy();
|
||
}
|
||
}
|
||
|
||
const children = [];
|
||
|
||
// A tool child. `setsid` leaves the engine's session and process group (K1,
|
||
// K2); `forkLoop` forks every 5 ms (K12); `ignoreTerm` survives SIGTERM, so
|
||
// only the kill phase ends it (K3, K10, K11). With `pidLog`, the fork loop
|
||
// appends each child's pid and a `term` line when it gets SIGTERM (K12).
|
||
function spawnChild({ setsid = false, forkLoop = false, ignoreTerm = false, pidLog = null } = {}) {
|
||
const note = pidLog ? `const note=(s)=>require('node:fs').appendFileSync(${JSON.stringify(pidLog)},s+'\\n');` : "const note=()=>{};";
|
||
const code = note + (ignoreTerm ? "process.on('SIGTERM',()=>note('term'));" : "") + (forkLoop
|
||
? "const {spawn}=require('node:child_process');setInterval(()=>{try{const c=spawn('sleep',['1000'],{stdio:'ignore'});if(c.pid)note(String(c.pid))}catch{}},5);setInterval(()=>{},1e9)"
|
||
: "setInterval(()=>{},1e9)");
|
||
const child = spawn(process.execPath, ["-e", code], { stdio: "ignore", detached: setsid });
|
||
children.push(child.pid);
|
||
return child.pid;
|
||
}
|
||
|
||
// K13: a member writes its own pid to another cgroup's cgroup.procs.
|
||
function escape(target) {
|
||
try {
|
||
writeFileSync(target, String(process.pid));
|
||
return { escaped: true };
|
||
} catch (err) {
|
||
return { escaped: false, code: err.code ?? String(err.message) };
|
||
}
|
||
}
|
||
|
||
async function main() {
|
||
const argv = process.argv.slice(2);
|
||
const logPath = process.env.FAKE_PI_LOG ?? null;
|
||
const log = logPath ? (v) => appendFileSync(logPath, JSON.stringify({ ...v, pid: process.pid }) + "\n") : null;
|
||
const leaf = process.env.FAKE_PI_LEAF !== undefined ? (process.env.FAKE_PI_LEAF === "null" ? null : process.env.FAKE_PI_LEAF) : undefined;
|
||
const fake = new FakePi({ input: process.stdin, output: process.stdout, argv, leaf, log });
|
||
if (process.env.FAKE_PI_SCRIPT) for (const s of JSON.parse(process.env.FAKE_PI_SCRIPT)) fake.script(s);
|
||
process.stdout.on("error", () => {});
|
||
process.stdin.on("end", () => log?.({ t: "stdin-end" }));
|
||
const controlPath = process.env.FAKE_PI_CONTROL;
|
||
if (!controlPath) return;
|
||
// A force-stopped predecessor (SIGKILL) leaves its socket file behind; the path is per-fixture.
|
||
rmSync(controlPath, { force: true });
|
||
const clients = new Set();
|
||
fake.onPaused = (point) => {
|
||
for (const c of clients) c.write(encodeLine({ event: "paused", point }));
|
||
};
|
||
createServer((sock) => {
|
||
clients.add(sock);
|
||
sock.on("close", () => clients.delete(sock));
|
||
sock.on("error", () => {});
|
||
const splitter = new LineSplitter((line) => {
|
||
const req = JSON.parse(line);
|
||
const reply = (result) => sock.write(encodeLine({ id: req.id, ok: true, result }));
|
||
const ops = {
|
||
script: () => fake.script(req.steps),
|
||
plan: () => fake.plan(req.plan),
|
||
arm: () => fake.arm(req.point, req.times ?? 1),
|
||
resume: () => fake.resume(req.point),
|
||
queue: () => fake.queue(req.kind, req.text),
|
||
dialog: () => fake.dialog(req.method),
|
||
emit: () => fake.emit(req.value),
|
||
raw: () => fake.raw(req.text),
|
||
drop: () => fake.dropResponse(req.type, req.n ?? 1),
|
||
extension: () => void fake.extensionPrompt(req.args ?? {}),
|
||
state: () => ({ streaming: fake.streaming, runs: fake.runs.length, commands: fake.commands, pid: process.pid, children, appends: fake.appends }),
|
||
child: () => ({ pid: spawnChild(req.args ?? {}) }),
|
||
escape: () => escape(req.target),
|
||
cgroup: () => readFileSync(`/proc/${req.pid ?? process.pid}/cgroup`, "utf8"),
|
||
waitPaused: () => fake.waitPaused(req.point),
|
||
// H19: stop reading stdin so the controller's write fills the pipe.
|
||
stall: () => void process.stdin.pause(),
|
||
};
|
||
if (!ops[req.op]) return sock.write(encodeLine({ id: req.id, ok: false, error: `unknown op ${req.op}` }));
|
||
try {
|
||
const out = ops[req.op]();
|
||
if (out && typeof out.then === "function") out.then(reply);
|
||
else reply(out ?? null);
|
||
} catch (err) {
|
||
sock.write(encodeLine({ id: req.id, ok: false, error: String(err.message) }));
|
||
}
|
||
});
|
||
sock.on("data", (c) => splitter.push(c));
|
||
}).listen(controlPath);
|
||
}
|
||
|
||
if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) await main();
|