// The trusted bus host (lead decision 70). It forks // packages/bus/src/process.mjs, boots it over IPC, keeps the reader // capability it gets back, and hands it to the notifier child over IPC. // Capabilities never appear in argv, stdout or the environment. While the // host runs, `/bus-host/host.json` (0600) names it for `mosaic bus // stop|status` and for the human commands' default business. // // startHost() is also the in-process API S6 uses: bindLaunch(record) binds // a launched run's process identity and returns its capability. import { fork } from "node:child_process"; import { once } from "node:events"; import { existsSync, mkdirSync, readFileSync, renameSync, rmSync, writeFileSync } from "node:fs"; import { join } from "node:path"; import { fileURLToPath } from "node:url"; import { CliError } from "./errors.mjs"; const BROKER = fileURLToPath(new URL("../../bus/src/process.mjs", import.meta.url)); const NOTIFIER = fileURLToPath(new URL("./notifier-process.mjs", import.meta.url)); export const BOOT_TIMEOUT_MS = 30000; const START_TIMEOUT_MS = 15000; const CLOSE_TIMEOUT_MS = 20000; export const hostDir = (dataRoot) => join(dataRoot, "bus-host"); export const hostFile = (dataRoot) => join(hostDir(dataRoot), "host.json"); // `/proc//stat` field 22, the start time in clock ticks. Null when the // process is gone or has exited and waits to be reaped (state Z). export function startTimeOf(pid) { try { const raw = readFileSync(`/proc/${pid}/stat`, "utf8"); const fields = raw.slice(raw.lastIndexOf(")") + 2).split(" "); return fields[0] === "Z" ? null : fields[19]; } catch { return null; } } export function cmdlineOf(pid) { try { return readFileSync(`/proc/${pid}/cmdline`, "utf8").split("\0").filter(Boolean); } catch { return null; } } // The state file, or null. `live` is true only when the pid still runs with // the recorded start time, so a recycled pid never counts as the host. export function readHostState(dataRoot) { const file = hostFile(dataRoot); if (!existsSync(file)) return null; let state; try { state = JSON.parse(readFileSync(file, "utf8")); } catch { throw new CliError(`bus host state is unreadable: ${file}`, 3); } if (!state || !Number.isSafeInteger(state.pid) || typeof state.startTime !== "string" || typeof state.business !== "string") { throw new CliError(`bus host state is malformed: ${file}`, 3); } return { ...state, live: startTimeOf(state.pid) === state.startTime }; } function writeHostState(dataRoot, state) { mkdirSync(hostDir(dataRoot), { recursive: true, mode: 0o700 }); const file = hostFile(dataRoot); const tmp = `${file}.${process.pid}.tmp`; writeFileSync(tmp, `${JSON.stringify(state, null, 2)}\n`, { mode: 0o600, flag: "wx" }); renameSync(tmp, file); } // The first IPC message from child, or a refusal when it exits or the wait runs out. function firstReply(child, timeoutMs, what) { return new Promise((resolve, reject) => { const done = (fn, v) => { clearTimeout(timer); child.off("message", onMessage); child.off("exit", onExit); fn(v); }; const onMessage = (m) => done(resolve, m); const onExit = (code) => done(reject, new CliError(`${what} exited (${code}) before it replied`, 1)); const timer = setTimeout(() => done(reject, new CliError(`${what} did not reply within ${Math.round(timeoutMs / 1000)} s`, 1)), timeoutMs); child.on("message", onMessage); child.on("exit", onExit); }); } async function ended(child, timeoutMs) { if (child.exitCode !== null || child.signalCode !== null) return child.exitCode; const timer = setTimeout(() => child.kill("SIGKILL"), timeoutMs); const [code] = await once(child, "exit"); clearTimeout(timer); return code; } // Calls onDeath(name, code, signal) once for each child that exits. A child // that exited before this runs (the broker while the notifier starts, say) // has already emitted its exit event, so its exit state is checked here. export function watchChildren(children, onDeath) { for (const [name, child] of Object.entries(children)) { if (!child) continue; if (child.exitCode !== null || child.signalCode !== null) onDeath(name, child.exitCode, child.signalCode); else child.once("exit", (code, signal) => onDeath(name, code, signal)); } } // boot: the {op:'boot'} config from bootConfig(). notifier: null, or // {binding, base?, pollMs?}; base and pollMs exist for the tests, and the // command line never sets them. Resolves once both children are up. export async function startHost({ boot, business, notifier = null, bootTimeoutMs = BOOT_TIMEOUT_MS, log = (l) => process.stderr.write(`mosaic-bus: ${l}\n`) }) { const dataRoot = boot.dataRoot; const prior = readHostState(dataRoot); if (prior?.live) throw new CliError(`a bus host already runs for ${prior.business} (pid ${prior.pid})`, 3); const broker = fork(BROKER, [], { stdio: ["ignore", "inherit", "inherit", "ipc"] }); const reply = firstReply(broker, bootTimeoutMs, "broker"); broker.send({ op: "boot", config: boot }); let ready; try { ready = await reply; } catch (e) { broker.kill("SIGTERM"); await ended(broker, 5000); throw e; } if (ready?.ok !== true) { await ended(broker, 5000); throw new CliError(`broker refused to start: ${typeof ready?.error === "string" ? ready.error : "startup-refused"}`, 3); } const reader = ready.readers.find((r) => r.business === business); let notify = null; if (notifier) { notify = fork(NOTIFIER, [], { stdio: ["ignore", "inherit", "inherit", "ipc"] }); const started = firstReply(notify, START_TIMEOUT_MS, "notifier"); notify.send({ op: "start", path: ready.path, cap: reader.cap, business, dataRoot, binding: notifier.binding, base: notifier.base, pollMs: notifier.pollMs }); let ok; try { ok = await started; } catch (e) { notify.kill("SIGTERM"); await ended(notify, 5000); if (broker.connected) broker.send({ op: "close" }, () => {}); await ended(broker, CLOSE_TIMEOUT_MS); throw e; } if (ok?.ok !== true) { await ended(notify, 5000); if (broker.connected) broker.send({ op: "close" }, () => {}); await ended(broker, CLOSE_TIMEOUT_MS); throw new CliError(`notifier refused to start: ${typeof ok?.error === "string" ? ok.error : "notifier-refused"}`, 3); } } const state = { hostVersion: 1, pid: process.pid, startTime: startTimeOf(process.pid), business, startedAt: new Date().toISOString(), notifier: notifier ? notifier.binding : null }; rmSync(hostFile(dataRoot), { force: true }); writeHostState(dataRoot, state); let closing = false; let finish; const done = new Promise((r) => (finish = r)); // Replies from process.mjs carry no request id, so requests go one at a time. let queue = Promise.resolve(); function bindLaunch(record) { const run = queue.then(async () => { if (closing) throw new CliError("bus host is closing", 1); const r = firstReply(broker, 10000, "broker"); broker.send({ op: "bindLaunch", record }); const m = await r; if (m?.ok !== true) throw new CliError(`launch bind refused: ${m?.error ?? "bind-refused"}`, 3); return m.launch; }); queue = run.catch(() => {}); return run; } async function close(code = 0) { if (closing) return done; closing = true; let result = code; if (notify) { if (notify.connected) notify.send({ op: "stop" }, () => {}); if ((await ended(notify, CLOSE_TIMEOUT_MS)) !== 0 && result === 0) result = 1; } await queue; if (broker.connected) broker.send({ op: "close" }, () => {}); if ((await ended(broker, CLOSE_TIMEOUT_MS)) !== 0 && result === 0) result = 1; const now = readHostState(dataRoot); if (now && now.pid === state.pid && now.startTime === state.startTime) rmSync(hostFile(dataRoot), { force: true }); finish(result); return done; } // A child that dies on its own takes the host down with exit 1: the unit // restarts the pair rather than run a broker without its notifier. // Runs after `queue` exists, since close() awaits it. watchChildren({ broker, notifier: notify }, (name, code, signal) => { if (closing) return; log(`${name} exited (${signal ?? code}); stopping the host`); close(1); }); return Object.freeze({ path: ready.path, business, bindLaunch, close, done, pids: Object.freeze({ broker: broker.pid, notifier: notify?.pid ?? null }) }); } // `mosaic bus stop`: SIGTERM to the recorded host, after checking that the // pid still runs with the recorded start time and is a `bus start` process. export async function stopHost(dataRoot, { timeoutMs = 60000 } = {}) { const state = readHostState(dataRoot); if (!state || !state.live) return { stopped: false, reason: state ? "stale state file; no host runs" : "no host runs" }; const argv = cmdlineOf(state.pid) ?? []; const i = argv.findIndex((a) => a.endsWith("/packages/cli/src/cli.mjs")); if (i < 0 || argv[i + 1] !== "bus" || argv[i + 2] !== "start") { throw new CliError(`pid ${state.pid} is not a bus host; refusing to signal it`, 3); } process.kill(state.pid, "SIGTERM"); const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (startTimeOf(state.pid) !== state.startTime) return { stopped: true, business: state.business, pid: state.pid }; await new Promise((r) => setTimeout(r, 200)); } throw new CliError(`bus host pid ${state.pid} did not stop within ${Math.round(timeoutMs / 1000)} s`, 1); } // `mosaic bus status`: what the files and /proc say. Never removes a lock. export function hostStatus(dataRoot) { const state = readHostState(dataRoot); const lock = join(dataRoot, "bus", "writer.lock"); // The bus store writes {pid, at}. A dead pid means a broker was killed // hard; the bus README leaves removing the lock to the operator. let owner = null; if (existsSync(lock)) { try { owner = JSON.parse(readFileSync(lock, "utf8")); } catch { owner = {}; } } const lockPid = Number.isSafeInteger(owner?.pid) ? owner.pid : null; return { host: state ? { business: state.business, pid: state.pid, startedAt: state.startedAt, notifier: state.notifier ?? null, live: state.live } : null, socket: existsSync(join(dataRoot, "bus", "broker.sock")), writerLock: owner ? { pid: lockPid, at: typeof owner.at === "string" ? owner.at : null, live: lockPid !== null && startTimeOf(lockPid) !== null } : null, }; }