Files
stack/packages/conversation/src/cohort.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

236 lines
12 KiB
JavaScript

// Engine launch, cohort observation and force stop (#1507, CHAT-03 §6).
//
// ScopeLauncher starts the supervisor shim (shim.mjs) in a delegated
// systemd user scope named after the claim. The cohort is the scope's
// `engine` cgroup, and the scope's invocation ID is the membership epoch. A
// unit with the right name but another invocation ID is a different cohort:
// it gets no signal, and its evidence is unavailable.
//
// PgroupLauncher is the fallback with no scope: the engine leads its own
// process group. A process group can't be enumerated completely, so a force
// stop on it always ends `uncertain`.
//
// Force stop: check the invocation ID against systemd and the shim; TERM
// every member and wait a bounded grace; freeze `engine` and wait for
// `frozen 1`; enumerate every member with pid and start time; write
// `cgroup.kill`; wait for `populated 0`. Only when all of that succeeded is
// membership complete. Anything unavailable ends the stop `uncertain`.
import { spawn, spawnSync } from "node:child_process";
import { existsSync, readFileSync } from "node:fs";
import { connect } from "node:net";
import { dirname, join } from "node:path";
import { fileURLToPath } from "node:url";
import { processStart } from "../../discord/src/journal.mjs";
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
import { hash, newId, record, sealProof } from "./records.mjs";
export const AUTHORITY = "mosaic-conversation-shim-fixture";
export const SHIM_PATH = join(dirname(fileURLToPath(import.meta.url)), "shim.mjs");
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
export function cohortRefOf({ machineId, bootId, unitName, invocationId }) {
return "cohort-" + hash({ machineId, bootId, unitName, invocationId }).slice(0, 40);
}
// One request per connection keeps a dead shim from wedging a caller.
export function shimRequest(socketPath, op, extra = {}, timeoutMs = 3000) {
return new Promise((resolve) => {
if (!socketPath || !existsSync(socketPath)) return resolve({ ok: false, unavailable: "the shim socket is gone" });
const sock = connect(socketPath);
let done = false;
const finish = (v) => {
if (done) return;
done = true;
clearTimeout(timer);
sock.destroy();
resolve(v);
};
const timer = setTimeout(() => finish({ ok: false, unavailable: `the shim did not answer ${op}` }), timeoutMs);
const splitter = new LineSplitter((line) => {
const p = parseLine(line);
finish(p.error ? { ok: false, unavailable: `shim answer ${p.error}` } : p.value);
});
sock.on("data", (c) => splitter.push(c));
sock.on("error", () => finish({ ok: false, unavailable: "the shim is unreachable" }));
sock.on("close", () => finish({ ok: false, unavailable: "the shim closed the connection" }));
sock.on("connect", () => sock.write(encodeLine({ id: 1, op, ...extra })));
});
}
export function systemctlShow(unitName) {
const r = spawnSync("systemctl", ["--user", "show", "-p", "LoadState,ActiveState,InvocationID,ControlGroup", `${unitName}.scope`], { encoding: "utf8", timeout: 5000 });
if (r.status !== 0) return null;
const out = {};
for (const line of r.stdout.split("\n")) {
const i = line.indexOf("=");
if (i > 0) out[line.slice(0, i)] = line.slice(i + 1);
}
return { loadState: out.LoadState ?? null, activeState: out.ActiveState ?? null, invocationId: out.InvocationID || null, controlGroup: out.ControlGroup || null };
}
// claim.classify's unit lookup. A unit systemd has collected is absent; one
// that is loaded and active is alive; anything unreadable is unknown.
export const systemdUnits = Object.freeze({
lookup(rec) {
if (!rec?.unitName) return { state: "unknown" };
const s = systemctlShow(rec.unitName);
if (!s) return { state: "unknown" };
if (s.loadState === "not-found" || (s.activeState === "inactive" && !s.controlGroup)) return { state: "absent" };
if (["active", "activating", "deactivating", "reloading"].includes(s.activeState)) return { state: "alive", invocationId: s.invocationId };
return { state: "unknown", detail: s };
},
});
export function scopeAvailable() {
const r = spawnSync("systemd-run", ["--user", "--scope", "--quiet", "-p", "Delegate=yes", "--", "true"], { timeout: 10000 });
return r.status === 0;
}
export class ScopeLauncher {
constructor({ shimPath = SHIM_PATH, startTimeoutMs = 10000 } = {}) {
this.kind = "scope";
this.shimPath = shimPath;
this.startTimeoutMs = startTimeoutMs;
}
async launch({ unitName, socketPath, command, args, cwd, env }) {
const proc = spawn("systemd-run", ["--user", "--scope", "-p", "Delegate=yes", `--unit=${unitName}`, "--quiet", "--", process.execPath, this.shimPath, "--socket", socketPath, "--", command, ...args], {
cwd, env, stdio: ["pipe", "pipe", "pipe"],
});
const exited = new Promise((r) => proc.on("exit", (code, signal) => r({ code, signal })));
const end = Date.now() + this.startTimeoutMs;
let hello = null;
while (Date.now() < end) {
const race = await Promise.race([exited.then((e) => ({ exited: e })), sleep(20).then(() => null)]);
if (race?.exited) throw Object.assign(new Error(`the scope exited before the shim answered (${JSON.stringify(race.exited)})`), { code: "launch-failed" });
if (existsSync(socketPath)) {
hello = await shimRequest(socketPath, "hello");
if (hello.ok) break;
}
}
if (!hello?.ok) throw Object.assign(new Error("the shim never answered"), { code: "launch-failed", proc });
const show = systemctlShow(unitName);
return {
kind: "scope", proc, stdin: proc.stdin, stdout: proc.stdout, stderr: proc.stderr,
pid: hello.enginePid, start: hello.engineStart === null ? null : String(hello.engineStart),
invocationId: hello.invocationId, scope: hello.scope, shimSocket: socketPath,
systemd: show, exited,
};
}
}
export class PgroupLauncher {
constructor() {
this.kind = "pgroup";
}
async launch({ command, args, cwd, env }) {
const proc = spawn(command, args, { cwd, env, stdio: ["pipe", "pipe", "pipe"], detached: true });
const exited = new Promise((r) => proc.on("exit", (code, signal) => r({ code, signal })));
await new Promise((resolve, reject) => {
proc.once("spawn", resolve);
proc.once("error", reject);
});
return { kind: "pgroup", proc, stdin: proc.stdin, stdout: proc.stdout, stderr: proc.stderr, pid: proc.pid, start: processStart(proc.pid), invocationId: null, scope: null, shimSocket: null, exited };
}
}
// Runs the force-stop escalation from the TERM phase. `onPhase(name)` is
// awaited before each phase begins, so the caller records the phase before
// any signal. Returns { outcome: "proven" | "unavailable", ... }. Nothing is
// recorded as done unless it was observed.
export async function forceStopCohort({ kind, unitName, invocationId, shimSocket, pid, graceMs = 1000, onPhase = async () => {} }) {
if (kind === "pgroup") {
await onPhase("term");
try {
process.kill(-pid, "SIGTERM");
} catch {
// the group is gone
}
await sleep(graceMs);
await onPhase("kill");
try {
process.kill(-pid, "SIGKILL");
} catch {
// the group is gone
}
return { outcome: "unavailable", reason: "process-group fallback: membership can't be enumerated completely" };
}
const show = systemctlShow(unitName);
if (!show || !invocationId || show.invocationId !== invocationId) {
return { outcome: "unavailable", reason: `invocation ID mismatch or unreadable (recorded ${invocationId}, found ${show?.invocationId ?? "none"}); no signal sent` };
}
const hello = await shimRequest(shimSocket, "hello");
if (!hello.ok) return { outcome: "unavailable", reason: hello.unavailable };
if (hello.invocationId !== invocationId || hello.scope !== show.controlGroup) {
return { outcome: "unavailable", reason: "the shim does not answer for the recorded scope; no signal sent" };
}
await onPhase("term");
const term = await shimRequest(shimSocket, "term");
if (!term.ok) return { outcome: "unavailable", reason: term.unavailable, phase: "term" };
const end = Date.now() + graceMs;
for (;;) {
const e = await shimRequest(shimSocket, "events");
if (!e.ok) return { outcome: "unavailable", reason: e.unavailable, phase: "term" };
if (e.populated === 0 || Date.now() >= end) break;
await sleep(20);
}
await onPhase("kill");
const frozen = await shimRequest(shimSocket, "freeze", { timeoutMs: 3000 }, 6000);
if (!frozen.ok) return { outcome: "unavailable", reason: frozen.unavailable, phase: "kill" };
if (frozen.timedOut) return { outcome: "unavailable", reason: "engine never reported frozen 1", phase: "kill" };
const listed = await shimRequest(shimSocket, "members");
if (!listed.ok) return { outcome: "unavailable", reason: listed.unavailable, phase: "kill" };
const killed = await shimRequest(shimSocket, "kill", { timeoutMs: 5000 }, 8000);
if (!killed.ok) return { outcome: "unavailable", reason: killed.unavailable, phase: "kill" };
if (killed.timedOut || killed.populated !== 0) return { outcome: "unavailable", reason: "engine never reported populated 0", phase: "kill" };
const observedAt = new Date().toISOString();
if (listed.members.some((m) => !Number.isInteger(m.startTicks) || m.startTicks < 1)) {
return { outcome: "unavailable", reason: "a member's start time was unreadable at enumeration", phase: "kill" };
}
return {
outcome: "proven", membershipComplete: true, epoch: invocationId, observedAt, boot: listed.boot,
members: listed.members.map((m) => ({ pid: m.pid, boot: listed.boot, startTicks: m.startTicks, terminatedAt: observedAt })),
};
}
export function cohortProof({ binding, stop, result }) {
return sealProof(record("cohortProof", {
id: newId("cohort-proof"), authority: AUTHORITY, conversation: binding.scope.conversation, execution: binding.execution,
cohortRef: binding.cohortRef, membershipEpoch: result.epoch, membershipComplete: result.membershipComplete === true,
members: result.members, observedAt: result.observedAt, verificationDigest: "", stop,
}));
}
// A tool-start with no tool-end at stop time is an `uncertain` effect.
// Killing never counts as rollback.
export function effectReport({ binding, stop, tools, observedAt }) {
return sealProof(record("effectReport", {
id: newId("effects"), authority: AUTHORITY, conversation: binding.scope.conversation, execution: binding.execution,
cohortRef: binding.cohortRef,
invocations: [...tools.values()].map((t) => ({ id: t.call, disposition: t.end ? "completed" : "uncertain", evidence: t.end ?? t.start })),
observedAt, verificationDigest: "", stop,
}));
}
// The current boot's start time, from /proc/stat btime. Every process of an
// earlier boot ended before it.
export function bootTime() {
const line = readFileSync("/proc/stat", "utf8").split("\n").find((l) => l.startsWith("btime "));
return new Date(Number(line.slice(6)) * 1000).toISOString();
}
// A boot proof for a claim recorded under an earlier boot of this host. Its
// evidence is the boot change; it goes through the same verifier.
export function bootProof({ claim, conversation, execution, cohortRef, stop, now = () => new Date() }) {
const terminatedAt = bootTime();
const members = claim.engine?.pid && claim.engine?.start ? [{ pid: claim.engine.pid, boot: claim.host.bootId, startTicks: Number(claim.engine.start), terminatedAt }] : [];
return sealProof(record("cohortProof", {
id: newId("boot-proof"), authority: AUTHORITY, conversation, execution, cohortRef,
membershipEpoch: claim.invocationId ?? `boot-${claim.host.bootId}`, membershipComplete: true, members,
observedAt: now().toISOString(), verificationDigest: "", stop,
}));
}