feat(cli): the mosaic CLI, broker host and decision notifier (row 39, S4, rocko)
packages/cli adds mosaic inbox, decide, tasks, agents and trail over the human-cli transport, and mosaic bus start, stop and status as the trusted host (unit mosaic-bus@<business>, scripts/bus-service.sh). The host boots packages/bus/src/process.mjs, passes config.trackers from the tracker.* variables (lead decision 70), and runs a notifier child. The notifier DMs each open blocking decision once and sends an 08:00 America/Chicago digest, journaled in notify/<business>/sent.jsonl at 0600 with no Discord ids. A torn journal tail is copied aside and truncated; a malformed line, a directory looser than 0700 or a symlinked journal refuses (lead decision 71). packages/discord gains dmRecipient, createDm and notify.mjs. Candidate agents/rocko/work/slice1-s4, baseb9b6cf00, build.patch b52f7d68, manifest e858504e (29 files). Darkwing approved round 2 on #1521 (comment 26855), Filbert approved round 2 (comment 26856). The packet's mutant table lists M28 as killed; it survived, and BUILD-LOG records the correction. Integration gate in a worktree on2557e29dwith the patch applied: bus 67, business 60, cli 49, control-board 124, discord 178, ledger 78, mosaic 69, queue 148, seat 19, tasks 51 and webui 14, all with no failures. Conversation is 149/3, the same K1, K3 and K10 cases that fail on the base; S4 doesn't touch the package. Every scripts/test-*.sh is green, with test-release 14/14 and test-task 98/98 on the existing gate2 compose network. A scratch test, not in this commit, booted the real host with trackers against S3's fake Vikunja: the adapter went ready and a task.close on a missing task answered task-not-found after a Vikunja read. Co-Authored-By: Claude Opus 5.5 <[email protected]>
This commit is contained in:
@@ -0,0 +1,211 @@
|
||||
#!/usr/bin/env node
|
||||
// mosaic inbox | decide | tasks | agents | trail, and mosaic bus start|stop|status.
|
||||
// See packages/cli/README.md. Exit codes are in src/errors.mjs.
|
||||
//
|
||||
// The human commands read and resolve through the broker's human transport
|
||||
// (src/transport.mjs), never the SQLite file. `bus start` is the trusted
|
||||
// host (src/host.mjs) the systemd unit mosaic-bus@<business> runs.
|
||||
|
||||
import { lstatSync, readFileSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
import { createInterface } from "node:readline/promises";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { BINDING_NAME } from "../../discord/src/binding.mjs";
|
||||
import { CliError } from "./errors.mjs";
|
||||
import { bootConfig, loadSystem, socketPath } from "./config.mjs";
|
||||
import { humanTransport, refuseInsideAgent } from "./transport.mjs";
|
||||
import { formatAgents, formatDecision, formatInbox, formatTasks, formatTrail, shortId } from "./format.mjs";
|
||||
import { hostStatus, readHostState, startHost, stopHost } from "./host.mjs";
|
||||
|
||||
export const USAGE = `usage:
|
||||
mosaic inbox [--business <id>] [--json]
|
||||
mosaic decide <decision> <option> [--note <text>] [--yes] [--business <id>]
|
||||
mosaic tasks [--business <id>] [--json]
|
||||
mosaic agents [--business <id>] [--json]
|
||||
mosaic trail <task|decision> [--business <id>] [--json]
|
||||
mosaic bus start <business>
|
||||
mosaic bus stop
|
||||
mosaic bus status [--json]`;
|
||||
|
||||
const usage = (why) => new CliError(why ? `${why}\n${USAGE}` : USAGE, 4);
|
||||
|
||||
function parse(args, { flags = [], values = [] }) {
|
||||
const out = { positional: [], flags: new Set(), values: {} };
|
||||
for (let i = 0; i < args.length; i++) {
|
||||
const a = args[i];
|
||||
if (flags.includes(a)) out.flags.add(a);
|
||||
else if (values.includes(a)) {
|
||||
if (i + 1 >= args.length || Object.hasOwn(out.values, a)) throw usage(`${a} needs one value`);
|
||||
out.values[a] = args[++i];
|
||||
} else if (a.startsWith("--")) throw usage(`unknown option ${a}`);
|
||||
else out.positional.push(a);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// The notifier's private config: `<dataRoot>/notify/<business>/notify.json`,
|
||||
// 0600, `{"notifyVersion": 1, "binding": "<discord binding>" | null}`.
|
||||
// Missing refuses, so a host never starts without someone deciding whether
|
||||
// it DMs; null runs the host without a notifier.
|
||||
export function readNotifyConfig(dataRoot, business) {
|
||||
const file = join(dataRoot, "notify", business, "notify.json");
|
||||
let st;
|
||||
try {
|
||||
st = lstatSync(file);
|
||||
} catch {
|
||||
throw new CliError(`no notifier config at ${file}; write {"notifyVersion": 1, "binding": "<discord binding>"} there, mode 0600, or "binding": null to run without DMs`, 3);
|
||||
}
|
||||
if (!st.isFile() || st.uid !== process.getuid() || (st.mode & 0o777) !== 0o600) {
|
||||
throw new CliError(`notifier config must be a regular file, mode 0600, owned by this user: ${file}`, 3);
|
||||
}
|
||||
let doc;
|
||||
try {
|
||||
doc = JSON.parse(readFileSync(file, "utf8"));
|
||||
} catch {
|
||||
throw new CliError(`notifier config is not JSON: ${file}`, 3);
|
||||
}
|
||||
const keys = doc && typeof doc === "object" && !Array.isArray(doc) ? Object.keys(doc).sort().join(",") : "";
|
||||
if (keys !== "binding,notifyVersion" || doc.notifyVersion !== 1 || (doc.binding !== null && !(typeof doc.binding === "string" && BINDING_NAME.test(doc.binding)))) {
|
||||
throw new CliError(`notifier config must be exactly {"notifyVersion": 1, "binding": <binding name> | null}: ${file}`, 3);
|
||||
}
|
||||
return doc.binding;
|
||||
}
|
||||
|
||||
// --business, else the running host's business.
|
||||
function businessFor(parsed, dataRoot) {
|
||||
if (parsed.values["--business"]) return parsed.values["--business"];
|
||||
const state = readHostState(dataRoot);
|
||||
if (state?.live) return state.business;
|
||||
throw usage("no bus host runs here, so name the business with --business");
|
||||
}
|
||||
|
||||
function findDecision(list, ref) {
|
||||
const exact = list.find((d) => d.id === ref);
|
||||
if (exact) return exact;
|
||||
if (ref.length < 8) throw new CliError(`decision reference ${JSON.stringify(ref)} is too short; use at least 8 characters of the id`, 2);
|
||||
const hits = list.filter((d) => d.id.startsWith(ref));
|
||||
if (hits.length === 1) return hits[0];
|
||||
if (hits.length > 1) throw new CliError(`${ref} matches ${hits.length} open decisions; use more of the id`, 2);
|
||||
throw new CliError(`no open decision for you matches ${ref}; see mosaic inbox`, 2);
|
||||
}
|
||||
|
||||
async function confirm(io, question) {
|
||||
const rl = createInterface({ input: io.stdin, output: io.stdout });
|
||||
try {
|
||||
return /^(y|yes)$/i.test((await rl.question(question)).trim());
|
||||
} finally {
|
||||
rl.close();
|
||||
}
|
||||
}
|
||||
|
||||
async function decide(parsed, call, io) {
|
||||
if (parsed.positional.length !== 2) throw usage("decide takes a decision and an option");
|
||||
const [ref, choice] = parsed.positional;
|
||||
const d = findDecision(await call("inbox"), ref);
|
||||
const option = d.options.find((o) => o.key === choice);
|
||||
if (!option) throw new CliError(`decision ${shortId(d.id)} has no option ${JSON.stringify(choice)}; its options are ${d.options.map((o) => o.key).join(", ")}`, 2);
|
||||
io.stdout.write(formatDecision(d));
|
||||
const effect = d.authorization ? (choice === d.authorization.approvalChoice ? "this authorizes the action" : "this declines the action") : null;
|
||||
io.stdout.write(`your choice: ${choice} (${option.text})${effect ? `; ${effect}` : ""}\n`);
|
||||
if (!parsed.flags.has("--yes")) {
|
||||
if (!io.stdin.isTTY) throw usage("stdin is not a terminal; pass --yes to resolve without the prompt");
|
||||
if (!(await confirm(io, `resolve ${shortId(d.id)} with ${choice}? [y/N] `))) throw new CliError("not resolved", 1);
|
||||
}
|
||||
const args = { id: d.id, choice };
|
||||
if (parsed.values["--note"] !== undefined) args.note = parsed.values["--note"];
|
||||
try {
|
||||
await call("decision.resolve", args);
|
||||
} catch (e) {
|
||||
if (e.code === "outcome-unknown") {
|
||||
throw new CliError(`outcome unknown: the broker may have recorded it. Check mosaic inbox or mosaic trail ${d.id} before trying again`, 1);
|
||||
}
|
||||
if (e.code === "decision-closed") throw new CliError(`decision ${shortId(d.id)} was closed before your answer arrived; see mosaic trail ${d.id}`, 2);
|
||||
throw e;
|
||||
}
|
||||
io.stdout.write(`resolved ${d.id}: ${choice}\n`);
|
||||
}
|
||||
|
||||
async function human(verb, rest, io, deps) {
|
||||
const spec = verb === "decide" ? { flags: ["--yes"], values: ["--business", "--note"] } : { flags: ["--json"], values: ["--business"] };
|
||||
const parsed = parse(rest, spec);
|
||||
refuseInsideAgent(io.env);
|
||||
const system = deps.system ?? loadSystem({ env: io.env });
|
||||
const business = businessFor(parsed, system.dataRoot);
|
||||
const call = deps.transport ? deps.transport(business) : humanTransport({ socket: socketPath(system.dataRoot), business, env: io.env });
|
||||
if (verb === "decide") return decide(parsed, call, io);
|
||||
const json = parsed.flags.has("--json");
|
||||
const print = (value, text) => io.stdout.write(json ? `${JSON.stringify(value, null, 2)}\n` : text(value));
|
||||
if (verb === "trail") {
|
||||
if (parsed.positional.length !== 1) throw usage("trail takes one task or decision");
|
||||
const subject = parsed.positional[0];
|
||||
return print(await call("trail", { subject }), (rows) => formatTrail(subject, rows));
|
||||
}
|
||||
if (parsed.positional.length !== 0) throw usage(`${verb} takes no arguments`);
|
||||
if (verb === "inbox") return print(await call("inbox"), formatInbox);
|
||||
if (verb === "tasks") return print(await call("tasks"), formatTasks);
|
||||
return print(await call("agents"), formatAgents);
|
||||
}
|
||||
|
||||
async function bus(rest, io, deps) {
|
||||
const [sub, ...args] = rest;
|
||||
if (sub === "start") {
|
||||
const parsed = parse(args, {});
|
||||
if (parsed.positional.length !== 1) throw usage("bus start takes one business");
|
||||
const businessId = parsed.positional[0];
|
||||
const system = deps.system ?? loadSystem({ env: io.env });
|
||||
const boot = bootConfig({ system, businessId, env: io.env, warn: (w) => io.stderr.write(`mosaic-bus: warning: ${w}\n`) });
|
||||
const binding = readNotifyConfig(system.dataRoot, businessId);
|
||||
const host = await startHost({ boot, business: businessId, notifier: binding ? { binding } : null });
|
||||
io.stdout.write(`bus host up: business ${businessId}, socket ${host.path}, notifier ${binding ?? "off"}\n`);
|
||||
const stop = () => host.close(0);
|
||||
process.on("SIGTERM", stop);
|
||||
process.on("SIGINT", stop);
|
||||
const code = await host.done;
|
||||
process.off("SIGTERM", stop);
|
||||
process.off("SIGINT", stop);
|
||||
io.stdout.write(`bus host stopped (${code})\n`);
|
||||
if (code !== 0) throw new CliError("bus host stopped after a child exited", code);
|
||||
return;
|
||||
}
|
||||
if (sub === "stop") {
|
||||
if (args.length !== 0) throw usage("bus stop takes no arguments");
|
||||
const system = deps.system ?? loadSystem({ env: io.env });
|
||||
const r = await stopHost(system.dataRoot);
|
||||
io.stdout.write(r.stopped ? `stopped bus host for ${r.business} (pid ${r.pid})\n` : `${r.reason}\n`);
|
||||
return;
|
||||
}
|
||||
if (sub === "status") {
|
||||
const parsed = parse(args, { flags: ["--json"] });
|
||||
if (parsed.positional.length !== 0) throw usage("bus status takes no arguments");
|
||||
const system = deps.system ?? loadSystem({ env: io.env });
|
||||
const s = hostStatus(system.dataRoot);
|
||||
if (parsed.flags.has("--json")) return io.stdout.write(`${JSON.stringify(s, null, 2)}\n`);
|
||||
const lines = [
|
||||
s.host ? `host: ${s.host.business}, pid ${s.host.pid}, ${s.host.live ? "running" : "not running (stale state file)"}, since ${s.host.startedAt}, notifier ${s.host.notifier ?? "off"}` : "host: none",
|
||||
`socket: ${s.socket ? "present" : "absent"}`,
|
||||
s.writerLock
|
||||
? `writer.lock: pid ${s.writerLock.pid ?? "?"}${s.writerLock.live ? "" : " (not running: a broker died hard; remove the lock by hand once you've checked, see packages/bus/README.md)"}`
|
||||
: "writer.lock: absent",
|
||||
];
|
||||
return io.stdout.write(`${lines.join("\n")}\n`);
|
||||
}
|
||||
throw usage();
|
||||
}
|
||||
|
||||
export async function main(argv, io = { env: process.env, stdin: process.stdin, stdout: process.stdout, stderr: process.stderr }, deps = {}) {
|
||||
const [verb, ...rest] = argv;
|
||||
if (["inbox", "decide", "tasks", "agents", "trail"].includes(verb)) return human(verb, rest, io, deps);
|
||||
if (verb === "bus") return bus(rest, io, deps);
|
||||
if (verb === "-h" || verb === "--help") return io.stdout.write(`${USAGE}\n`);
|
||||
throw usage();
|
||||
}
|
||||
|
||||
if (process.argv[1] === fileURLToPath(import.meta.url)) {
|
||||
try {
|
||||
await main(process.argv.slice(2));
|
||||
} catch (error) {
|
||||
if (!(error instanceof CliError)) throw error;
|
||||
process.stderr.write(`mosaic: ${error.message}\n`);
|
||||
process.exitCode = error.exitCode;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
// What the host and the human commands read before they touch the bus:
|
||||
// the system config (through scripts/mosaic-config.mjs, as `mosaic business`
|
||||
// does), the business file, and the boot message for the broker process.
|
||||
// Everything here reads and refuses; nothing writes.
|
||||
|
||||
import { spawnSync } from "node:child_process";
|
||||
import { existsSync } from "node:fs";
|
||||
import { dirname, join, resolve } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { BusinessError, configDir, loadBusiness, loadProject, projectFilePath, resolveInstance, systemVars } from "../../business/src/index.mjs";
|
||||
import { busBusiness } from "../../bus/src/business.mjs";
|
||||
import { BusError } from "../../bus/src/broker.mjs";
|
||||
import { CliError } from "./errors.mjs";
|
||||
|
||||
export const REPO = resolve(dirname(fileURLToPath(import.meta.url)), "..", "..", "..");
|
||||
|
||||
// The broker's socket. One broker per data root: the bus store takes
|
||||
// `<dataRoot>/bus/writer.lock` (packages/bus/src/store.mjs).
|
||||
export const socketPath = (dataRoot) => join(dataRoot, "bus", "broker.sock");
|
||||
|
||||
export function loadSystem({ env = process.env } = {}) {
|
||||
const proc = spawnSync(process.execPath, [join(REPO, "scripts", "mosaic-config.mjs"), "validate"], {
|
||||
encoding: "utf8",
|
||||
maxBuffer: 1024 * 1024,
|
||||
env,
|
||||
});
|
||||
if (proc.status !== 0) throw new CliError(`system config problem (mosaic-config exit ${proc.status}): ${proc.stderr.trim()}`, 3);
|
||||
try {
|
||||
return JSON.parse(proc.stdout);
|
||||
} catch {
|
||||
throw new CliError("system config: mosaic-config printed something that isn't JSON", 3);
|
||||
}
|
||||
}
|
||||
|
||||
export function rolesDir(env = process.env) {
|
||||
return resolve(env.MOSAIC_ROLES_DIR || join(REPO, "roles"));
|
||||
}
|
||||
|
||||
// A business-package or adapter refusal is a config problem: exit 3.
|
||||
function business(fn) {
|
||||
try {
|
||||
return fn();
|
||||
} catch (error) {
|
||||
if (error instanceof BusinessError) throw new CliError(error.message, 3);
|
||||
if (error instanceof BusError) throw new CliError(`business adapter refused: ${error.code}`, 3);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
// config.trackers[<business>] (lead decision 70): plain data from the
|
||||
// tracker.* variables. tracker.project is a project-layer variable, so the
|
||||
// entry comes from the one declared project whose file sets it. No
|
||||
// tracker.baseUrl, or no project naming a tracker project, means no entry
|
||||
// and no task verbs; two projects naming one each refuse, since the boot
|
||||
// shape holds one tracker project per business.
|
||||
function trackerFor(system, b, warn) {
|
||||
const first = Object.keys(b.roles)[0];
|
||||
const base = resolveInstance({ system, business: b, instance: first }).vars;
|
||||
if (base["tracker.baseUrl"] === undefined) return null;
|
||||
const named = [];
|
||||
for (const [id, entry] of Object.entries(b.projects)) {
|
||||
if (!existsSync(projectFilePath(entry.root))) continue;
|
||||
const project = loadProject(entry.root);
|
||||
if (project.id !== id) throw new CliError(`project file ${project.file} has id ${project.id}, but business ${b.id} declares it as ${id}`, 3);
|
||||
const vars = resolveInstance({ system, business: b, project, instance: first }).vars;
|
||||
if (vars["tracker.project"] !== undefined) named.push({ id, vars });
|
||||
}
|
||||
if (named.length > 1) {
|
||||
throw new CliError(`business ${b.id}: projects ${named.map((n) => n.id).join(", ")} each set tracker.project; the broker takes one tracker project per business`, 3);
|
||||
}
|
||||
if (named.length === 0) {
|
||||
warn(`business ${b.id} sets tracker.baseUrl, but no project file sets tracker.project; the broker gets no task verbs`);
|
||||
return null;
|
||||
}
|
||||
const v = named[0].vars;
|
||||
return {
|
||||
baseUrl: v["tracker.baseUrl"],
|
||||
project: v["tracker.project"],
|
||||
pollSeconds: v["tracker.pollSeconds"],
|
||||
reconcileMinutes: v["tracker.reconcileMinutes"],
|
||||
};
|
||||
}
|
||||
|
||||
// The `{op: 'boot'}` config for packages/bus/src/process.mjs, for one
|
||||
// business. The host holds a reader capability for the notifier and binds
|
||||
// launches later through bindLaunch, so `launches` starts empty.
|
||||
export function bootConfig({ system, businessId, env = process.env, warn = () => {} }) {
|
||||
return business(() => {
|
||||
const vars = systemVars(system);
|
||||
const b = loadBusiness(businessId, { dir: configDir(env), rolesDir: rolesDir(env) });
|
||||
const resolved = {};
|
||||
for (const instance of Object.keys(b.roles)) resolved[instance] = resolveInstance({ system: vars, business: b, instance });
|
||||
const config = {
|
||||
dataRoot: system.dataRoot,
|
||||
businesses: { [b.id]: busBusiness(b, resolved) },
|
||||
launches: [],
|
||||
readers: [b.id],
|
||||
};
|
||||
const tracker = trackerFor(vars, b, warn);
|
||||
if (tracker) config.trackers = { [b.id]: tracker };
|
||||
return config;
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
// Exit codes for every `mosaic` command in this package: 0 ok; 1 failed
|
||||
// (a child died, a boot timed out, an outcome is unknown); 2 invalid (bad
|
||||
// input, a decision or option that doesn't exist); 3 refused (system or
|
||||
// business config problem, the broker or notifier refused, the human proof
|
||||
// refused, a host already running); 4 usage.
|
||||
export class CliError extends Error {
|
||||
constructor(message, exitCode = 2) {
|
||||
super(message);
|
||||
this.name = "CliError";
|
||||
this.exitCode = exitCode;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,97 @@
|
||||
// Plain-text views of broker results. Every function takes what the broker
|
||||
// returned and keeps its order: the broker sorts, this file never re-sorts.
|
||||
|
||||
const one = (s, max = 160) => {
|
||||
const flat = String(s ?? "").replace(/\s+/g, " ").trim();
|
||||
return flat.length > max ? `${flat.slice(0, max - 1)}…` : flat;
|
||||
};
|
||||
|
||||
export const shortId = (id) => String(id).slice(0, 8);
|
||||
|
||||
export function optionsLine(d) {
|
||||
return d.options.map((o) => `${o.key} = ${one(o.text, 80)}`).join("; ");
|
||||
}
|
||||
|
||||
// What approving means, when the decision carries authorization context.
|
||||
export function authorizationLines(d) {
|
||||
const a = d.authorization;
|
||||
if (!a) return [];
|
||||
const lines = [`action: ${a.action}${a.target ? ` on ${a.target}` : ""}`];
|
||||
lines.push(`choosing "${a.approvalChoice}" authorizes it; any other choice declines`);
|
||||
return lines;
|
||||
}
|
||||
|
||||
export function formatInbox(list) {
|
||||
if (list.length === 0) return "inbox: empty\n";
|
||||
const out = [`inbox: ${list.length} open decision(s)`];
|
||||
for (const d of list) {
|
||||
const flags = [d.class, d.blocking ? "blocking" : null].filter(Boolean).join(", ");
|
||||
out.push("", `${shortId(d.id)} ${d.action} (${flags}) raised ${d.at} by ${d.raised_by_role}`);
|
||||
out.push(` ${one(d.question, 400)}`);
|
||||
out.push(` options: ${optionsLine(d)} (recommended: ${d.recommendation})`);
|
||||
for (const l of authorizationLines(d)) out.push(` ${l}`);
|
||||
if (d.task_ref) out.push(` task: ${d.task_ref}`);
|
||||
out.push(` decide: mosaic decide ${shortId(d.id)} <option>`);
|
||||
}
|
||||
return `${out.join("\n")}\n`;
|
||||
}
|
||||
|
||||
export function formatDecision(d) {
|
||||
const out = [
|
||||
`decision ${d.id}`,
|
||||
` ${d.class}${d.blocking ? ", blocking" : ""}, raised ${d.at} by ${d.raised_by_role} (${d.raised_by_run})`,
|
||||
` question: ${d.question}`,
|
||||
` options: ${optionsLine(d)}`,
|
||||
` recommended: ${d.recommendation}`,
|
||||
];
|
||||
for (const l of authorizationLines(d)) out.push(` ${l}`);
|
||||
if (d.task_ref) out.push(` task: ${d.task_ref}`);
|
||||
return `${out.join("\n")}\n`;
|
||||
}
|
||||
|
||||
export function formatAgents(list) {
|
||||
if (list.length === 0) return "agents: no role is claimed\n";
|
||||
return `${list.map((c) => `${c.role} ${c.holder_run} ${c.harness ?? "-"} since ${c.at}`).join("\n")}\n`;
|
||||
}
|
||||
|
||||
export function formatTasks(list) {
|
||||
if (list.length === 0) return "tasks: none\n";
|
||||
return `${list
|
||||
.map((t) => {
|
||||
const f = t.fields ?? {};
|
||||
const title = f.title ?? f.name ?? "";
|
||||
const state = f.status ?? f.state ?? (f.done === true ? "done" : f.done === false ? "open" : "");
|
||||
return [t.task_ref, state, one(title, 100)].filter(Boolean).join(" ");
|
||||
})
|
||||
.join("\n")}\n`;
|
||||
}
|
||||
|
||||
function trailSummary(r) {
|
||||
switch (r.table) {
|
||||
case "events":
|
||||
return [r.kind, r.actor_role ?? null, r.body?.operation ?? null, r.body?.decision ? `decision ${shortId(r.body.decision)}` : null].filter(Boolean).join(" ");
|
||||
case "decisions":
|
||||
return `${r.class} ${r.action} → ${r.route_to}${r.blocking ? " (blocking)" : ""}: ${one(r.question, 120)}`;
|
||||
case "decision_events":
|
||||
return [r.op, r.choice ? `choice ${r.choice}` : null, `by ${r.by}`, r.via ? `via ${r.via}` : null].filter(Boolean).join(" ");
|
||||
case "messages":
|
||||
return `${r.from_role} → ${r.to_role}: ${one(r.body, 120)}`;
|
||||
case "deliveries":
|
||||
return [r.op, r.transport, r.holder_run].filter(Boolean).join(" ");
|
||||
case "role_claims":
|
||||
return `${r.op} ${r.role} by ${r.holder_run}`;
|
||||
case "task_snapshots":
|
||||
return `snapshot ${r.task_ref ?? ""}`.trim();
|
||||
default:
|
||||
return "";
|
||||
}
|
||||
}
|
||||
|
||||
// Rows print in the order the broker returned them (at, table, seq).
|
||||
export function formatTrail(subject, rows) {
|
||||
const out = [`trail ${subject}: ${rows.length} row(s)`];
|
||||
for (const r of rows) out.push(`${r.at} ${r.table}#${r.seq} ${trailSummary(r)}`.trimEnd());
|
||||
const decision = rows.find((r) => r.table === "decisions" && r.id === subject);
|
||||
if (decision?.task_ref) out.push("", `task: ${decision.task_ref}`, `follow with: mosaic trail ${decision.task_ref}`);
|
||||
return `${out.join("\n")}\n`;
|
||||
}
|
||||
@@ -0,0 +1,247 @@
|
||||
// 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, `<dataRoot>/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/<pid>/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);
|
||||
broker.send({ op: "close" });
|
||||
await ended(broker, CLOSE_TIMEOUT_MS);
|
||||
throw e;
|
||||
}
|
||||
if (ok?.ok !== true) {
|
||||
await ended(notify, 5000);
|
||||
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,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
// @mosaic/cli: the human commands and the bus host. See README.md.
|
||||
export { CliError } from "./errors.mjs";
|
||||
export { bootConfig, loadSystem, rolesDir, socketPath } from "./config.mjs";
|
||||
export { AGENT_MARKERS, HUMAN_CLI, busExit, humanTransport, refuseInsideAgent } from "./transport.mjs";
|
||||
export { authorizationLines, formatAgents, formatDecision, formatInbox, formatTasks, formatTrail, optionsLine, shortId } from "./format.mjs";
|
||||
export { DIGEST_HOUR, POLL_MS, ZONE, createNotifier, digestContent, digestNonce, dmContent, dmNonce, journalPath, openJournal, runLoop, zoned } from "./notifier.mjs";
|
||||
export { BOOT_TIMEOUT_MS, hostFile, hostStatus, readHostState, startHost, stopHost } from "./host.mjs";
|
||||
export { USAGE, main, readNotifyConfig } from "./cli.mjs";
|
||||
@@ -0,0 +1,52 @@
|
||||
// The notifier as a child of the bus host. Started with fork(); the reader
|
||||
// capability arrives over IPC in `{op: 'start'}` and never touches argv,
|
||||
// stdout or the environment. Replies `{ok: true}` once the binding, the
|
||||
// token file and the journal check out, or `{ok: false, error}` and exits
|
||||
// 3. `{op: 'stop'}`, SIGTERM or the host going away stop it after the poll
|
||||
// in flight.
|
||||
import { Client } from "../../bus/src/client.mjs";
|
||||
import { openDirect } from "../../discord/src/notify.mjs";
|
||||
import { createNotifier, runLoop } from "./notifier.mjs";
|
||||
|
||||
const log = (line) => process.stderr.write(`mosaic-notify: ${line}\n`);
|
||||
let loop = null;
|
||||
let started = false;
|
||||
let stopping = false;
|
||||
|
||||
async function stop(code) {
|
||||
if (stopping) return;
|
||||
stopping = true;
|
||||
try {
|
||||
await loop?.stop();
|
||||
} catch {
|
||||
code = 1;
|
||||
}
|
||||
process.exitCode = code;
|
||||
if (process.connected) process.disconnect();
|
||||
}
|
||||
|
||||
if (!process.send) {
|
||||
process.stderr.write("trusted-host-required\n");
|
||||
process.exitCode = 2;
|
||||
} else {
|
||||
const timer = setTimeout(() => stop(2), 10000);
|
||||
process.on("message", (m) => {
|
||||
if (m?.op === "stop") return stop(0);
|
||||
if (m?.op !== "start" || started) return;
|
||||
started = true;
|
||||
clearTimeout(timer);
|
||||
try {
|
||||
const client = new Client({ path: m.path, cap: m.cap });
|
||||
const direct = openDirect({ dataRoot: m.dataRoot, name: m.binding, ...(m.base ? { base: m.base } : {}), log });
|
||||
const notifier = createNotifier({ business: m.business, dataRoot: m.dataRoot, inbox: () => client.call("inbox"), direct, log });
|
||||
loop = runLoop(notifier, { pollMs: m.pollMs ?? undefined, log });
|
||||
process.send({ ok: true });
|
||||
} catch (e) {
|
||||
// DiscordError and CliError messages name files, never a token or id.
|
||||
process.send({ ok: false, error: e?.message ?? "notifier-refused" }, () => stop(3));
|
||||
}
|
||||
});
|
||||
process.on("disconnect", () => stop(stopping ? process.exitCode : 2));
|
||||
process.on("SIGTERM", () => stop(0));
|
||||
process.on("SIGINT", () => {});
|
||||
}
|
||||
@@ -0,0 +1,276 @@
|
||||
// The notifier (lead decision 70): every poll it reads the human inbox
|
||||
// through a reader capability, DMs each open blocking decision once, and
|
||||
// sends one digest a day at 08:00 America/Chicago. It never writes to the
|
||||
// bus; its memory is the journal `<dataRoot>/notify/<business>/sent.jsonl`
|
||||
// (0600, in a 0700 directory), one line per send attempt:
|
||||
//
|
||||
// {at, kind: "dm"|"digest", decision, day?, outcome: "confirmed"|"refused"|"unknown", messageId, status?}
|
||||
//
|
||||
// A decision counts as sent once it has a confirmed line; a day's digest
|
||||
// likewise. A refused or unknown send is retried with backoff (30 s
|
||||
// doubling to 30 min): a duplicate costs less than a miss, and Discord's
|
||||
// nonce folds a retry inside its dedupe window into the first message.
|
||||
// No Discord channel or user id goes in the journal, a log line or an
|
||||
// error; the Discord side (packages/discord/src/notify.mjs) keeps them.
|
||||
|
||||
import { createHash } from "node:crypto";
|
||||
import { closeSync, constants, fstatSync, fsyncSync, ftruncateSync, lstatSync, mkdirSync, openSync, readFileSync, writeSync } from "node:fs";
|
||||
import { dirname, join } from "node:path";
|
||||
import { CliError } from "./errors.mjs";
|
||||
import { authorizationLines, optionsLine, shortId } from "./format.mjs";
|
||||
|
||||
export const ZONE = "America/Chicago";
|
||||
export const DIGEST_HOUR = 8;
|
||||
export const POLL_MS = 30000;
|
||||
const BACKOFF_MS = 30000;
|
||||
const BACKOFF_MAX_MS = 30 * 60 * 1000;
|
||||
const LIMIT = 2000;
|
||||
|
||||
export const journalPath = (dataRoot, business) => join(dataRoot, "notify", business, "sent.jsonl");
|
||||
|
||||
// Local date and hour in the IANA zone. An unknown zone throws RangeError,
|
||||
// so a broken time-zone database refuses rather than guessing.
|
||||
export function zoned(date, zone = ZONE) {
|
||||
const parts = Object.fromEntries(
|
||||
new Intl.DateTimeFormat("en-US", { timeZone: zone, year: "numeric", month: "2-digit", day: "2-digit", hour: "2-digit", hourCycle: "h23" })
|
||||
.formatToParts(date)
|
||||
.map((p) => [p.type, p.value]),
|
||||
);
|
||||
return { day: `${parts.year}-${parts.month}-${parts.day}`, hour: Number(parts.hour) };
|
||||
}
|
||||
|
||||
// Discord nonces are at most 25 characters. The digest nonce carries the
|
||||
// business, so two businesses sharing a bot and a recipient never send the
|
||||
// same nonce on the same day.
|
||||
const hash23 = (text) => createHash("sha256").update(text).digest("hex").slice(0, 23);
|
||||
export const dmNonce = (id) => `dm${hash23(String(id))}`;
|
||||
export const digestNonce = (business, day) => `dg${hash23(`${business}\n${day}`)}`;
|
||||
|
||||
function clip(text, max) {
|
||||
return text.length > max ? `${text.slice(0, max - 1)}…` : text;
|
||||
}
|
||||
|
||||
export function dmContent(business, d) {
|
||||
const lines = [
|
||||
`Mosaic (${business}): a blocking decision needs you.`,
|
||||
clip(d.question, 1000),
|
||||
`options: ${optionsLine(d)} (recommended: ${d.recommendation})`,
|
||||
...authorizationLines(d),
|
||||
`raised by ${d.raised_by_role}${d.task_ref ? `, task ${d.task_ref}` : ""}`,
|
||||
`decide: mosaic decide ${shortId(d.id)} <option>`,
|
||||
];
|
||||
return clip(lines.join("\n"), LIMIT);
|
||||
}
|
||||
|
||||
export function digestContent(business, day, inbox, dmSent) {
|
||||
if (inbox.length === 0) return `Mosaic digest (${business}, ${day}): your inbox is empty.`;
|
||||
const head = `Mosaic digest (${business}, ${day}): ${inbox.length} open decision(s).`;
|
||||
const tail = "Run mosaic inbox for the full list.";
|
||||
const lines = [head];
|
||||
let shown = 0;
|
||||
for (const d of inbox) {
|
||||
const mark = d.blocking ? (dmSent(d.id) ? "[blocking, DM sent] " : "[blocking, DM pending] ") : "";
|
||||
const line = `- ${mark}${shortId(d.id)} ${d.action}: ${clip(d.question.replace(/\s+/g, " "), 160)}`;
|
||||
const more = inbox.length - shown - 1;
|
||||
const reserve = more > 0 ? `\n… and ${more} more.`.length : 0;
|
||||
if ([...lines, line].join("\n").length + reserve + tail.length + 1 > LIMIT) break;
|
||||
lines.push(line);
|
||||
shown++;
|
||||
}
|
||||
if (shown < inbox.length) lines.push(`… and ${inbox.length - shown} more.`);
|
||||
lines.push(tail);
|
||||
return lines.join("\n");
|
||||
}
|
||||
|
||||
const { O_APPEND, O_CREAT, O_NOFOLLOW, O_RDWR, O_WRONLY } = constants;
|
||||
|
||||
// UTC stamp for a torn-tail copy, e.g. 20261008T235212345Z.
|
||||
const stamp = (date) => date.toISOString().replace(/[-:.]/g, "");
|
||||
|
||||
// Step 1 of a torn-tail repair (lead decision 71): the torn bytes go to a
|
||||
// new file, torn-<UTC stamp>.bin (0600), fsynced with its directory. A name
|
||||
// that exists already gets a counter, so a repeat never overwrites a copy.
|
||||
function copyTorn(dir, bytes, date) {
|
||||
for (let n = 0; ; n++) {
|
||||
const name = `torn-${stamp(date)}${n ? `-${n}` : ""}.bin`;
|
||||
let fd;
|
||||
try {
|
||||
fd = openSync(join(dir, name), O_WRONLY | O_CREAT | constants.O_EXCL | O_NOFOLLOW, 0o600);
|
||||
} catch (e) {
|
||||
if (e.code === "EEXIST") continue;
|
||||
throw e;
|
||||
}
|
||||
try {
|
||||
writeSync(fd, bytes);
|
||||
fsyncSync(fd);
|
||||
} finally {
|
||||
closeSync(fd);
|
||||
}
|
||||
const dfd = openSync(dir, constants.O_RDONLY);
|
||||
try {
|
||||
fsyncSync(dfd);
|
||||
} finally {
|
||||
closeSync(dfd);
|
||||
}
|
||||
return name;
|
||||
}
|
||||
}
|
||||
|
||||
// Opens (creating if needed) the journal. The directory must be 0700 or
|
||||
// tighter and the file 0600, both owned by this user; a symlinked journal
|
||||
// refuses. Every complete line must parse, or the open refuses with exit 3.
|
||||
// A final line without its newline is a write that never finished (lead
|
||||
// decision 71): its bytes are copied to torn-<stamp>.bin, then the journal
|
||||
// is truncated to its last newline and fsynced, and both steps are logged.
|
||||
// A crash between the two leaves the tail torn, and the next open repeats
|
||||
// both; the second copy is harmless.
|
||||
export function openJournal(file, { log = () => {}, now = () => new Date() } = {}) {
|
||||
const dir = dirname(file);
|
||||
mkdirSync(dir, { recursive: true, mode: 0o700 });
|
||||
const ds = lstatSync(dir);
|
||||
if (!ds.isDirectory() || ds.uid !== process.getuid() || (ds.mode & 0o077) !== 0) {
|
||||
throw new CliError(`notify journal directory must be mode 0700 and owned by this user: ${dir}`, 3);
|
||||
}
|
||||
let fd;
|
||||
try {
|
||||
fd = openSync(file, O_RDWR | O_APPEND | O_CREAT | O_NOFOLLOW, 0o600);
|
||||
} catch (e) {
|
||||
if (e.code === "ELOOP") throw new CliError(`notify journal must not be a symlink: ${file}`, 3);
|
||||
throw e;
|
||||
}
|
||||
const sent = new Set();
|
||||
const days = new Set();
|
||||
try {
|
||||
const st = fstatSync(fd);
|
||||
if (!st.isFile() || st.uid !== process.getuid() || (st.mode & 0o777) !== 0o600) {
|
||||
throw new CliError(`notify journal must be a regular file, mode 0600, owned by this user: ${file}`, 3);
|
||||
}
|
||||
const bytes = readFileSync(fd);
|
||||
const end = bytes.lastIndexOf(0x0a) + 1;
|
||||
const lines = bytes.subarray(0, end).toString("utf8").split("\n");
|
||||
lines.pop();
|
||||
lines.forEach((line, i) => {
|
||||
let r;
|
||||
try {
|
||||
r = JSON.parse(line);
|
||||
} catch {
|
||||
r = null;
|
||||
}
|
||||
if (!r || typeof r !== "object" || !["dm", "digest"].includes(r.kind)) {
|
||||
throw new CliError(`notify journal line ${i + 1} is malformed: ${file}`, 3);
|
||||
}
|
||||
if (r.outcome !== "confirmed") return;
|
||||
if (r.kind === "dm") sent.add(r.decision);
|
||||
else days.add(r.day);
|
||||
});
|
||||
if (end < bytes.length) {
|
||||
const name = copyTorn(dir, bytes.subarray(end), now());
|
||||
log(`notify journal: copied a torn final line (${bytes.length - end} bytes) to ${name}`);
|
||||
ftruncateSync(fd, end);
|
||||
fsyncSync(fd);
|
||||
log(`notify journal: truncated ${file} to its last newline (${end} bytes)`);
|
||||
}
|
||||
} finally {
|
||||
closeSync(fd);
|
||||
}
|
||||
return {
|
||||
sent,
|
||||
days,
|
||||
append(record) {
|
||||
const afd = openSync(file, O_WRONLY | O_APPEND | O_NOFOLLOW);
|
||||
try {
|
||||
writeSync(afd, `${JSON.stringify(record)}\n`);
|
||||
} finally {
|
||||
closeSync(afd);
|
||||
}
|
||||
if (record.outcome !== "confirmed") return;
|
||||
if (record.kind === "dm") sent.add(record.decision);
|
||||
else days.add(record.day);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
// inbox(): the broker's inbox through the reader capability.
|
||||
// direct.send({content, nonce}) resolves {messageId} or throws a RestOutcome.
|
||||
export function createNotifier({ business, dataRoot, inbox, direct, now = () => new Date(), zone = ZONE, log = () => {} }) {
|
||||
zoned(now(), zone);
|
||||
const journal = openJournal(journalPath(dataRoot, business), { log, now });
|
||||
const backoff = new Map();
|
||||
|
||||
function waiting(key, t) {
|
||||
const b = backoff.get(key);
|
||||
return b !== undefined && t < b.next;
|
||||
}
|
||||
|
||||
async function attempt(key, record, message) {
|
||||
const t = now().getTime();
|
||||
try {
|
||||
const { messageId } = await direct.send(message);
|
||||
backoff.delete(key);
|
||||
journal.append({ at: new Date(t).toISOString(), ...record, outcome: "confirmed", messageId });
|
||||
return true;
|
||||
} catch (err) {
|
||||
const kind = err?.kind === "refused" ? "refused" : "unknown";
|
||||
const status = Number.isInteger(err?.details?.status) ? err.details.status : null;
|
||||
const n = (backoff.get(key)?.n ?? -1) + 1;
|
||||
backoff.set(key, { n, next: t + Math.min(BACKOFF_MS * 2 ** n, BACKOFF_MAX_MS) });
|
||||
journal.append({ at: new Date(t).toISOString(), ...record, outcome: kind, messageId: null, ...(status !== null ? { status } : {}) });
|
||||
log(`notify: ${record.kind} ${kind}${status !== null ? ` (HTTP ${status})` : ""}; retry in ${Math.round(Math.min(BACKOFF_MS * 2 ** n, BACKOFF_MAX_MS) / 1000)} s`);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// One poll. Returns what it did, for the tests and the log.
|
||||
async function tick() {
|
||||
const done = { dms: 0, digest: false, failed: 0 };
|
||||
let list;
|
||||
try {
|
||||
list = await inbox();
|
||||
} catch (err) {
|
||||
log(`notify: inbox read failed (${err?.code ?? err?.message ?? "error"}); next poll retries`);
|
||||
return { ...done, inboxError: true };
|
||||
}
|
||||
const t = now();
|
||||
for (const d of list) {
|
||||
if (!d.blocking || journal.sent.has(d.id)) continue;
|
||||
const key = `dm:${d.id}`;
|
||||
if (waiting(key, t.getTime())) continue;
|
||||
if (await attempt(key, { kind: "dm", decision: d.id }, { content: dmContent(business, d), nonce: dmNonce(d.id) })) done.dms++;
|
||||
else done.failed++;
|
||||
}
|
||||
const { day, hour } = zoned(t, zone);
|
||||
const key = `digest:${day}`;
|
||||
if (hour >= DIGEST_HOUR && !journal.days.has(day) && !waiting(key, t.getTime())) {
|
||||
const content = digestContent(business, day, list, (id) => journal.sent.has(id));
|
||||
if (await attempt(key, { kind: "digest", decision: null, day }, { content, nonce: digestNonce(business, day) })) done.digest = true;
|
||||
else done.failed++;
|
||||
}
|
||||
return done;
|
||||
}
|
||||
|
||||
return Object.freeze({ tick, journal: journalPath(dataRoot, business) });
|
||||
}
|
||||
|
||||
// Runs tick() every pollMs until stop(). Ticks never overlap.
|
||||
export function runLoop(notifier, { pollMs = POLL_MS, log = () => {} } = {}) {
|
||||
let timer = null;
|
||||
let stopped = false;
|
||||
let running = Promise.resolve();
|
||||
const loop = () => {
|
||||
if (stopped) return;
|
||||
running = notifier
|
||||
.tick()
|
||||
.catch((err) => log(`notify: poll failed: ${err?.message ?? err}`))
|
||||
.finally(() => {
|
||||
if (!stopped) timer = setTimeout(loop, pollMs);
|
||||
});
|
||||
};
|
||||
loop();
|
||||
return {
|
||||
async stop() {
|
||||
stopped = true;
|
||||
clearTimeout(timer);
|
||||
await running;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
// The human transport (lead decision 70): run
|
||||
// `packages/bus/src/human-cli.mjs <socket>` as a child, write
|
||||
// `{business, verb, args}` on its stdin, read the JSON reply on its stdout,
|
||||
// or one bus error code on its stderr. The bus does the proof: the child
|
||||
// re-executes itself with a nonce and the broker checks the process and its
|
||||
// ancestry for agent markers. Nothing here carries a capability.
|
||||
|
||||
import { spawnSync } from "node:child_process";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { CliError } from "./errors.mjs";
|
||||
|
||||
export const HUMAN_CLI = fileURLToPath(new URL("../../bus/src/human-cli.mjs", import.meta.url));
|
||||
|
||||
// The broker's markers (packages/bus/src/human.mjs). The broker checks the
|
||||
// whole ancestry; this check is only the clear early message for the case
|
||||
// it would refuse anyway.
|
||||
export const AGENT_MARKERS = Object.freeze([
|
||||
"MOSAIC_BUS_CAP",
|
||||
"MOSAIC_RUN_ID",
|
||||
"MOSAIC_AGENT_RUN",
|
||||
"CLAUDECODE",
|
||||
"CLAUDE_CODE_ENTRYPOINT",
|
||||
"CODEX_THREAD_ID",
|
||||
"PI_AGENT_DIR",
|
||||
]);
|
||||
|
||||
export function refuseInsideAgent(env = process.env) {
|
||||
const found = AGENT_MARKERS.filter((k) => env[k]);
|
||||
if (found.length > 0) {
|
||||
throw new CliError(`refused: this runs only from a human shell, outside any agent run (found ${found.join(", ")} in the environment)`, 3);
|
||||
}
|
||||
}
|
||||
|
||||
// Exit codes by bus error code; anything else is invalid input (2).
|
||||
const EXIT = {
|
||||
"human-required": 3,
|
||||
unauthenticated: 3,
|
||||
"unknown-business": 3,
|
||||
"read-only": 3,
|
||||
"outcome-unknown": 1,
|
||||
"response-too-large": 1,
|
||||
"invalid-response": 1,
|
||||
};
|
||||
|
||||
export function busExit(code) {
|
||||
return EXIT[code] ?? 2;
|
||||
}
|
||||
|
||||
// Returns a call(verb, args) for one business. `cli` and `spawn` exist for
|
||||
// the tests; the command line never sets them.
|
||||
export function humanTransport({ socket, business, env = process.env, cli = HUMAN_CLI, spawn = spawnSync, timeoutMs = 30000 }) {
|
||||
return async function call(verb, args = {}) {
|
||||
const proc = spawn(process.execPath, [cli, socket], {
|
||||
input: `${JSON.stringify({ business, verb, args })}\n`,
|
||||
encoding: "utf8",
|
||||
env,
|
||||
timeout: timeoutMs,
|
||||
maxBuffer: 8 * 1024 * 1024,
|
||||
});
|
||||
if (proc.error || proc.status === null) throw new CliError("outcome-unknown: the bus transport did not finish", 1);
|
||||
if (proc.status !== 0) {
|
||||
const code = String(proc.stderr ?? "").trim().split("\n").reverse().find((l) => /^[a-z-]{1,64}$/.test(l)) ?? "invalid-response";
|
||||
const error = new CliError(code, busExit(code));
|
||||
error.code = code;
|
||||
throw error;
|
||||
}
|
||||
try {
|
||||
return JSON.parse(proc.stdout);
|
||||
} catch {
|
||||
const error = new CliError("invalid-response", 1);
|
||||
error.code = "invalid-response";
|
||||
throw error;
|
||||
}
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user