#!/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:`. 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:`; 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();