Files
stack/packages/conversation/src/shim.mjs
T
jason.woltjeandClaude Opus 5.5 243e153c8b feat(conversation): CHAT-03 I1, mediated control of a sealed headless Pi (#1507)
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]>
2026-10-04 15:47:53 -05:00

190 lines
7.1 KiB
JavaScript

#!/usr/bin/env node
// The supervisor shim (#1507, CHAT-03 §6). Not a library: `cohort.mjs` starts
// it as
//
// systemd-run --user --scope -p Delegate=yes --unit=<name> --quiet -- \
// node shim.mjs --socket <path> -- <engine argv...>
//
// Inside the delegated scope it moves itself into a `supervisor` child cgroup
// and starts the engine in an `engine` child cgroup. A small shell writes its
// own pid into engine/cgroup.procs before it execs `unshare -U
// --map-current-user --cgroup`, which execs the engine, so the engine is
// contained from its first instruction and its cgroup namespace is rooted at
// `engine`. The engine inherits the controller's stdin and stdout directly;
// the shim closes its own copies. The shim is not a cohort member. It holds
// the scope open, so systemd can't collect the cgroup before emptiness is
// read, and it outlives its controller so a restarted controller can reach
// the cohort again.
//
// Control is JSON lines over a Unix socket in the controller's 0700 socket
// directory. Every request names an op; every answer is {ok, ...} or
// {ok:false, unavailable}. Emptiness is a readable engine/cgroup.events with
// `populated 0`. A missing or unreadable file is unavailable, never empty.
import { spawn } from "node:child_process";
import { closeSync, mkdirSync, readdirSync, readFileSync, unlinkSync, writeFileSync } from "node:fs";
import { createServer } from "node:net";
import { join } from "node:path";
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
const args = process.argv.slice(2);
const sep = args.indexOf("--");
const opt = (name) => {
const i = args.indexOf(name);
return i >= 0 && i < sep ? args[i + 1] : null;
};
const socketPath = opt("--socket");
const engineArgv = sep >= 0 ? args.slice(sep + 1) : [];
if (!socketPath || engineArgv.length === 0) {
process.stderr.write("shim: usage: shim.mjs --socket <path> -- <engine argv...>\n");
process.exit(2);
}
const own = readFileSync("/proc/self/cgroup", "utf8").split("\n").find((l) => l.startsWith("0::"));
const scope = join("/sys/fs/cgroup", own.slice(3).trim());
const supervisor = join(scope, "supervisor");
const engine = join(scope, "engine");
const invocationId = process.env.INVOCATION_ID ?? null;
const bootId = readFileSync("/proc/sys/kernel/random/boot_id", "utf8").trim();
mkdirSync(supervisor, { recursive: true });
mkdirSync(engine, { recursive: true });
writeFileSync(join(supervisor, "cgroup.procs"), String(process.pid));
const startOf = (pid) => {
try {
const stat = readFileSync(`/proc/${pid}/stat`, "utf8");
const fields = stat.slice(stat.lastIndexOf(")") + 2).split(" ");
return Number(fields[19]);
} catch {
return null;
}
};
const child = spawn(
"/bin/sh",
["-c", 'echo $$ > "$1/cgroup.procs" || exit 97; shift; exec unshare -U --map-current-user --cgroup -- "$@"', "mosaic-engine", engine, ...engineArgv],
{ stdio: [0, 1, 2] },
);
const enginePid = child.pid;
let engineExit = null;
child.on("exit", (code, signal) => {
engineExit = { code, signal, at: new Date().toISOString() };
});
// The engine holds the controller's pipes; the shim's copies would hide EOF.
closeSync(0);
closeSync(1);
function events() {
try {
const text = readFileSync(join(engine, "cgroup.events"), "utf8");
const out = {};
for (const line of text.split("\n")) {
const [k, v] = line.split(" ");
if (k) out[k] = Number(v);
}
if (out.populated !== 0 && out.populated !== 1) return { ok: false, unavailable: "cgroup.events has no populated field" };
return { ok: true, populated: out.populated, frozen: out.frozen ?? null };
} catch (err) {
return { ok: false, unavailable: `engine/cgroup.events unreadable (${err.code ?? err.message})` };
}
}
// Every member of `engine` and its descendants: the engine's namespace can
// create child cgroups, so enumeration is recursive.
function members() {
const out = [];
const walk = (dir) => {
const procs = readFileSync(join(dir, "cgroup.procs"), "utf8").split("\n").filter(Boolean).map(Number);
for (const pid of procs) out.push({ pid, startTicks: startOf(pid), cgroup: dir.slice(scope.length) || "/" });
for (const name of readdirSync(dir, { withFileTypes: true })) if (name.isDirectory()) walk(join(dir, name.name));
};
try {
walk(engine);
return { ok: true, members: out, boot: bootId };
} catch (err) {
return { ok: false, unavailable: `engine enumeration failed (${err.code ?? err.message})` };
}
}
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
async function waitFor(pred, ms) {
const end = Date.now() + ms;
for (;;) {
const e = events();
if (!e.ok) return e;
if (pred(e)) return e;
if (Date.now() >= end) return { ...e, timedOut: true };
await sleep(10);
}
}
async function handle(req) {
switch (req.op) {
case "hello":
return { ok: true, invocationId, scope: scope.slice("/sys/fs/cgroup".length), enginePid, engineStart: startOf(enginePid), shimPid: process.pid, shimStart: startOf(process.pid), boot: bootId, engineExit };
case "events":
return events();
case "members":
return members();
case "term": {
const m = members();
if (!m.ok) return m;
const signalled = [];
for (const { pid, startTicks } of m.members) {
if (startOf(pid) !== startTicks) continue;
try {
process.kill(pid, "SIGTERM");
signalled.push(pid);
} catch {
// gone already
}
}
return { ok: true, signalled };
}
case "freeze":
try {
writeFileSync(join(engine, "cgroup.freeze"), "1");
} catch (err) {
return { ok: false, unavailable: `cgroup.freeze unwritable (${err.code ?? err.message})` };
}
return waitFor((e) => e.frozen === 1 || e.populated === 0, Number(req.timeoutMs) || 2000);
case "kill":
try {
writeFileSync(join(engine, "cgroup.kill"), "1");
} catch (err) {
return { ok: false, unavailable: `cgroup.kill unwritable (${err.code ?? err.message})` };
}
return waitFor((e) => e.populated === 0, Number(req.timeoutMs) || 5000);
case "release": {
const e = events();
if (!e.ok || e.populated !== 0) return { ok: false, unavailable: "the engine cgroup is not empty" };
setTimeout(() => {
try {
unlinkSync(socketPath);
} catch {
// already gone
}
process.exit(0);
}, 10);
return { ok: true };
}
default:
return { ok: false, unavailable: `unknown op ${String(req.op)}` };
}
}
const server = createServer((sock) => {
const splitter = new LineSplitter(async (line) => {
const parsed = parseLine(line);
const answer = parsed.error ? { ok: false, unavailable: parsed.error } : await handle(parsed.value);
if (!sock.destroyed) sock.write(encodeLine({ id: parsed.value?.id ?? null, ...answer }));
}, { maxBytes: 65536 });
sock.on("data", (chunk) => splitter.push(chunk));
sock.on("error", () => {});
});
server.listen(socketPath);
process.on("SIGTERM", () => {}); // a stray TERM never drops the scope's anchor
process.on("SIGHUP", () => {});