import { test } from "node:test"; import assert from "node:assert/strict"; import { spawn, spawnSync } from "node:child_process"; import { subscribe, unsubscribe } from "node:diagnostics_channel"; import { createServer } from "node:http"; import { once } from "node:events"; import { existsSync, readFileSync, statSync, writeFileSync, mkdirSync } from "node:fs"; import { join } from "node:path"; import { fileURLToPath } from "node:url"; import { Client } from "../../bus/src/client.mjs"; import { bootConfig, loadSystem, REPO } from "../src/config.mjs"; import { hostDir, hostFile, hostStatus, startHost, startTimeOf, stopHost, watchChildren } from "../src/host.mjs"; import { journalPath } from "../src/notifier.mjs"; import { IDS, makeDeployment } from "../../discord/tests/helpers.mjs"; import { fixture, notifyConfig, OPTIONS, tmp } from "./helpers.mjs"; const CLI = fileURLToPath(new URL("../src/cli.mjs", import.meta.url)); const CHANNEL = "100000000000000900"; const TOKEN = "MTAw.abcdefghijklmnopqrstuvwxyz0123456789"; // A fake Discord REST on 127.0.0.1: opens one DM channel, accepts messages. async function fakeDiscord(t) { const requests = []; let n = 0; const server = createServer((req, res) => { let body = ""; req.on("data", (b) => (body += b)); req.on("end", () => { requests.push({ method: req.method, url: req.url, body: body ? JSON.parse(body) : null, authorized: req.headers.authorization === `Bot ${TOKEN}` }); res.setHeader("content-type", "application/json"); if (req.url === "/users/@me/channels") return res.end(JSON.stringify({ id: CHANNEL, type: 1 })); if (req.url === `/channels/${CHANNEL}/messages`) return res.end(JSON.stringify({ id: `30000000000000${String(++n).padStart(4, "0")}` })); res.statusCode = 404; res.end("{}"); }); }); server.listen(0, "127.0.0.1"); await once(server, "listening"); t.after(() => server.close()); return { base: `http://127.0.0.1:${server.address().port}`, requests, dms: () => requests.filter((r) => r.url.endsWith("/messages") && r.body.nonce.startsWith("dm")) }; } async function until(fn, ms = 8000) { const end = Date.now() + ms; while (Date.now() < end) { if (fn()) return; await new Promise((r) => setTimeout(r, 50)); } throw new Error("timed out waiting"); } // Records every IPC message this process sends to a child it creates while // the spy is on, by trapping the `send` that node installs on a new child. // The host never exposes the notifier's reader capability; this is how a // test sees it. function spySends(t, onSend = () => {}) { const sent = []; const onChild = ({ process: child }) => { let send; Object.defineProperty(child, "send", { configurable: true, get: () => send, set(fn) { send = function (m, ...rest) { sent.push(m); onSend(m); return fn.call(this, m, ...rest); }; }, }); }; subscribe("child_process", onChild); t.after(() => unsubscribe("child_process", onChild)); return sent; } // Fails this process's sends to new children the way node fails a write to a // child that has just died: `fake(child, m)` returns an error to fail that // send (to its callback if it has one, otherwise as an 'error' event on the // next tick), or nothing to send it for real. Deterministic where the real // window after a SIGKILL is a millisecond wide (Darkwing R1, lead decision 73). // Every child it saw is killed after the test, so a host left half closed // fails the test rather than hang it. function failSends(t, fake) { const children = []; const onChild = ({ process: child }) => { children.push(child); let send; Object.defineProperty(child, "send", { configurable: true, get: () => send, set(fn) { send = function (m, ...rest) { const err = fake(this, m); if (!err) return fn.call(this, m, ...rest); const callback = rest.find((a) => typeof a === "function"); if (callback) process.nextTick(callback, err); else process.nextTick(() => this.emit("error", err)); return false; }; }, }); }; subscribe("child_process", onChild); t.after(() => { unsubscribe("child_process", onChild); for (const child of children) child.kill("SIGKILL"); }); } const epipe = (child) => { child.kill("SIGKILL"); return Object.assign(new Error("write EPIPE"), { code: "EPIPE", errno: -32, syscall: "write" }); }; const procText = (pid, what) => { try { return readFileSync(`/proc/${pid}/${what}`, "utf8"); } catch { return ""; } }; test("the host boots the broker, binds a launch in process, and the notifier DMs a blocking decision exactly once", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root, { dmRecipient: IDS.owner }); const discord = await fakeDiscord(t); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); assert.equal("trackers" in boot, false); const logs = []; const sends = spySends(t); const host = await startHost({ boot, business: "acme", notifier: { binding: "test-seat", base: discord.base, pollMs: 100 }, log: (l) => logs.push(l) }); t.after(() => host.close(0)); const start = sends.find((m) => m?.op === "start"); assert.equal(typeof start?.cap, "string", "the spy saw the notifier's start message"); assert.ok(start.cap.length >= 16); const state = JSON.parse(readFileSync(hostFile(f.dataRoot), "utf8")); assert.equal(statSync(hostFile(f.dataRoot)).mode & 0o777, 0o600); assert.deepEqual(Object.keys(state).sort(), ["business", "hostVersion", "notifier", "pid", "startTime", "startedAt"]); assert.equal(hostStatus(f.dataRoot).host.live, true); const launch = await host.bindLaunch({ business: "acme", role: "coder", run: "coder-run", harness: "pi", pid: process.pid, startTime: startTimeOf(process.pid) }); assert.equal(launch.run, "coder-run"); const coder = new Client({ path: host.path, cap: launch.cap }); await coder.call("role.claim"); const d = await coder.call("decision.raise", { action: "git.push.protected", target: "refactor", question: "Push?", options: OPTIONS, recommendation: "no", blocking: true, task_ref: "vikunja:1/7" }); await until(() => discord.dms().length === 1); await new Promise((r) => setTimeout(r, 500)); assert.equal(discord.dms().length, 1, "five more polls send nothing new"); assert.ok(discord.requests.every((r) => r.authorized)); assert.match(discord.dms()[0].body.content, new RegExp(`mosaic decide ${d.id.slice(0, 8)}`)); // No capability, the launch's or the notifier's reader, in a child's argv // or environment, or in the state file. for (const cap of [launch.cap, start.cap]) { for (const pid of Object.values(host.pids)) { assert.notEqual(procText(pid, "cmdline"), "", `pid ${pid} is readable`); assert.ok(!procText(pid, "cmdline").includes(cap)); assert.ok(!procText(pid, "environ").includes(cap)); } assert.ok(!readFileSync(hostFile(f.dataRoot), "utf8").includes(cap)); } assert.equal(await host.close(0), 0); const journal = readFileSync(journalPath(f.dataRoot, "acme"), "utf8"); for (const id of [IDS.owner, CHANNEL, TOKEN]) assert.ok(!journal.includes(id)); assert.equal(journal.trim().split("\n").filter((l) => JSON.parse(l).kind === "dm").length, 1); assert.equal(existsSync(hostFile(f.dataRoot)), false); assert.equal(existsSync(join(f.dataRoot, "bus", "writer.lock")), false); }); test("a notifier that dies takes the host down with exit 1, so the unit restarts the pair", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root, { dmRecipient: IDS.owner }); const discord = await fakeDiscord(t); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); const logs = []; const host = await startHost({ boot, business: "acme", notifier: { binding: "test-seat", base: discord.base, pollMs: 100 }, log: (l) => logs.push(l) }); t.after(() => host.close(0)); process.kill(host.pids.notifier, "SIGKILL"); assert.equal(await host.done, 1); assert.match(logs.join("\n"), /notifier exited \(SIGKILL\); stopping the host/); assert.equal(existsSync(join(f.dataRoot, "bus", "writer.lock")), false); assert.equal(existsSync(hostFile(f.dataRoot)), false); }); test("a second host for the same data root refuses with exit 3 while the first runs", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); const host = await startHost({ boot, business: "acme", log: () => {} }); t.after(() => host.close(0)); const second = startHost({ boot, business: "acme", log: () => {} }); // If the refusal regresses and a second host starts, close it too. t.after(async () => (await second.catch(() => null))?.close(0)); await assert.rejects(second, (e) => e.exitCode === 3 && e.message === `a bus host already runs for acme (pid ${process.pid})`); assert.equal(hostStatus(f.dataRoot).host.live, true, "the first host still runs"); assert.equal(await host.close(0), 0); }); test("a notifier that refuses stops the broker and the host refuses with exit 3", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); const started = startHost({ boot, business: "acme", notifier: { binding: "test-seat" }, log: () => {} }); // If the refusal regresses, the host starts; close it so the file still ends. t.after(async () => (await started.catch(() => null))?.close(0)); await assert.rejects(started, (e) => e.exitCode === 3 && /no dmRecipient/.test(e.message)); assert.equal(existsSync(join(f.dataRoot, "bus", "writer.lock")), false); assert.equal(existsSync(hostFile(f.dataRoot)), false); }); test("a notifier that refuses after the broker died still refuses with exit 3, without a send to the dead broker", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); const children = []; const onChild = ({ process: child }) => children.push(child); subscribe("child_process", onChild); t.after(() => unsubscribe("child_process", onChild)); const errors = []; // The broker dies as the host sends the notifier its start message. const sends = spySends(t, (m) => { if (m?.op !== "start") return; children[0].on("error", (e) => errors.push(e.code)); children[0].kill("SIGKILL"); }); const started = startHost({ boot, business: "acme", notifier: { binding: "test-seat" }, log: () => {} }); t.after(async () => (await started.catch(() => null))?.close(0)); await assert.rejects(started, (e) => e.exitCode === 3 && /no dmRecipient/.test(e.message)); await new Promise((r) => setImmediate(r)); assert.deepEqual(errors, [], "no close was sent over the closed channel"); assert.equal(sends.filter((m) => m?.op === "close").length, 0); }); test("a notifier that dies before it replies, after the broker died, still refuses, without a send to the dead broker", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); const children = []; const onChild = ({ process: child }) => children.push(child); subscribe("child_process", onChild); t.after(() => unsubscribe("child_process", onChild)); const errors = []; // The broker dies as the host sends the start message. The notifier is // stopped so it cannot reply, and killed once the host has seen the broker // disconnect: the host takes the no-reply path, not the refusal path. const sends = spySends(t, (m) => { if (m?.op !== "start") return; children[0].on("error", (e) => errors.push(e.code)); children[1].kill("SIGSTOP"); children[0].once("disconnect", () => children[1].kill("SIGKILL")); children[0].kill("SIGKILL"); }); const started = startHost({ boot, business: "acme", notifier: { binding: "test-seat" }, log: () => {} }); t.after(async () => (await started.catch(() => null))?.close(0)); await assert.rejects(started, (e) => e.exitCode === 1 && /notifier exited \(null\) before it replied/.test(e.message)); await new Promise((r) => setImmediate(r)); assert.deepEqual(errors, [], "no close was sent over the closed channel"); assert.equal(sends.filter((m) => m?.op === "close").length, 0); }); test("a close send that fails with EPIPE after the notifier refuses still gives the notifier's refusal, exit 3", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); failSends(t, (child, m) => (m?.op === "close" ? epipe(child) : null)); const started = startHost({ boot, business: "acme", notifier: { binding: "test-seat" }, log: () => {} }); t.after(async () => (await started.catch(() => null))?.close(0)); await assert.rejects(started, (e) => e.exitCode === 3 && /^notifier refused to start: .*no dmRecipient/.test(e.message)); }); test("a close send that fails with EPIPE after the notifier dies unanswered still gives the notifier's error", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); failSends(t, (child, m) => { if (m?.op === "close") return epipe(child); if (m?.op === "start") { // Stopped, it cannot reply; then it dies. child.kill("SIGSTOP"); setImmediate(() => child.kill("SIGKILL")); } return null; }); const started = startHost({ boot, business: "acme", notifier: { binding: "test-seat" }, log: () => {} }); t.after(async () => (await started.catch(() => null))?.close(0)); await assert.rejects(started, (e) => e.exitCode === 1 && /^notifier exited \(null\) before it replied/.test(e.message)); }); test("close() whose stop and close sends fail with EPIPE still finishes, with exit 1", async (t) => { const root = tmp(t); const f = fixture(root); makeDeployment(root, { dmRecipient: IDS.owner }); const discord = await fakeDiscord(t); const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); failSends(t, (child, m) => (m?.op === "stop" || m?.op === "close" ? epipe(child) : null)); const host = await startHost({ boot, business: "acme", notifier: { binding: "test-seat", base: discord.base, pollMs: 100 }, log: () => {} }); // No t.after close: a close() that rejected leaves `done` pending forever. assert.equal(await host.close(0), 1, "both children died by signal"); assert.equal(existsSync(hostFile(f.dataRoot)), false); }); test("watchChildren reports a child that died before it was called, and one that dies later", async (t) => { const early = spawn(process.execPath, ["-e", "process.exit(7)"], { stdio: "ignore" }); await once(early, "exit"); const killed = spawn(process.execPath, ["-e", "setTimeout(() => {}, 60000)"], { stdio: "ignore" }); await once(killed, "spawn"); killed.kill("SIGKILL"); await once(killed, "exit"); const signalled = []; watchChildren({ broker: killed }, (...d) => signalled.push(d)); assert.deepEqual(signalled, [["broker", null, "SIGKILL"]], "a death by signal before the watch is not lost either"); const late = spawn(process.execPath, ["-e", "setTimeout(() => {}, 60000)"], { stdio: "ignore" }); t.after(() => late.kill("SIGKILL")); await once(late, "spawn"); const deaths = []; watchChildren({ broker: early, notifier: late, none: null }, (...d) => deaths.push(d)); assert.deepEqual(deaths, [["broker", 7, null]], "the exit before the watch is not lost"); late.kill("SIGTERM"); await once(late, "exit"); assert.deepEqual(deaths, [["broker", 7, null], ["notifier", null, "SIGTERM"]]); }); test("bus stop refuses to signal a live pid that is not a bus host", async (t) => { const dataRoot = tmp(t); const child = spawn(process.execPath, ["-e", "setTimeout(() => {}, 60000)"], { stdio: "ignore" }); t.after(() => child.kill("SIGKILL")); await once(child, "spawn"); mkdirSync(hostDir(dataRoot), { recursive: true, mode: 0o700 }); writeFileSync(hostFile(dataRoot), JSON.stringify({ pid: child.pid, startTime: startTimeOf(child.pid), business: "acme" }), { mode: 0o600 }); await assert.rejects(stopHost(dataRoot, { timeoutMs: 1000 }), (e) => e.exitCode === 3 && /is not a bus host; refusing to signal it/.test(e.message)); await new Promise((r) => setTimeout(r, 200)); assert.equal(child.exitCode, null); assert.equal(child.signalCode, null, "the child was not signalled"); }); test("bus start refuses with exit 3 and the code when the broker refuses to boot; bus status names the lock", (t) => { const f = fixture(tmp(t)); notifyConfig(f.dataRoot, "acme", null); mkdirSync(join(f.dataRoot, "bus"), { recursive: true, mode: 0o700 }); writeFileSync(join(f.dataRoot, "bus", "writer.lock"), JSON.stringify({ pid: 999999999, at: "2026-10-08T00:00:00Z" }), { mode: 0o600 }); const r = spawnSync(process.execPath, [CLI, "bus", "start", "acme"], { env: f.env, encoding: "utf8", timeout: 40000 }); assert.equal(r.status, 3, r.stderr); assert.match(r.stderr, /broker refused to start: startup-refused/); const s = spawnSync(process.execPath, [CLI, "bus", "status"], { env: f.env, encoding: "utf8" }); assert.equal(s.status, 0, s.stderr); assert.match(s.stdout, /host: none/); assert.match(s.stdout, /writer\.lock: pid 999999999 \(not running/); }); test("bus start refuses with exit 3 without a notifier config", (t) => { const f = fixture(tmp(t)); const r = spawnSync(process.execPath, [CLI, "bus", "start", "acme"], { env: f.env, encoding: "utf8", timeout: 40000 }); assert.equal(r.status, 3); assert.match(r.stderr, /no notifier config/); assert.equal(existsSync(join(f.dataRoot, "bus", "writer.lock")), false); }); test("bus start runs until bus stop; status reports it while it runs", async (t) => { const f = fixture(tmp(t)); notifyConfig(f.dataRoot, "acme", null); const child = spawn(process.execPath, [CLI, "bus", "start", "acme"], { env: f.env, stdio: ["ignore", "pipe", "pipe"] }); t.after(() => child.exitCode === null && child.kill("SIGKILL")); let out = ""; child.stdout.on("data", (b) => (out += b)); child.stderr.on("data", (b) => (out += b)); await until(() => /bus host up: business acme/.test(out), 30000); const status = spawnSync(process.execPath, [CLI, "bus", "status", "--json"], { env: f.env, encoding: "utf8" }); const s = JSON.parse(status.stdout); assert.equal(s.host.live, true); assert.equal(s.host.pid, child.pid); assert.equal(s.host.notifier, null); assert.equal(s.socket, true); assert.equal(s.writerLock.live, true); const exit = once(child, "exit"); const stop = spawnSync(process.execPath, [CLI, "bus", "stop"], { env: f.env, encoding: "utf8", timeout: 70000 }); assert.equal(stop.status, 0, stop.stderr); assert.match(stop.stdout, new RegExp(`stopped bus host for acme \\(pid ${child.pid}\\)`)); assert.equal((await exit)[0], 0, out); assert.equal(existsSync(join(f.dataRoot, "bus", "writer.lock")), false); const again = spawnSync(process.execPath, [CLI, "bus", "stop"], { env: f.env, encoding: "utf8" }); assert.match(again.stdout, /no host runs/); }); test("bus-service.sh renders the unit and installs it into a given directory", (t) => { const dir = tmp(t); const script = join(REPO, "scripts", "bus-service.sh"); const render = spawnSync(script, ["render"], { encoding: "utf8" }); assert.equal(render.status, 0, render.stderr); assert.match(render.stdout, new RegExp(`ExecStart=${REPO.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")}/scripts/mosaic bus start %i`)); assert.match(render.stdout, /RestartPreventExitStatus=2 3 4/); assert.doesNotMatch(render.stdout, /@REPO@|@PATH@/); // A user unit cannot order on a system target (Darkwing F3). assert.doesNotMatch(render.stdout, /network-online/); const first = spawnSync(script, ["install", "--dir", dir, "--no-reload"], { encoding: "utf8" }); assert.equal(first.status, 0, first.stderr); assert.match(first.stdout, /written: /); assert.match(first.stdout, /^ {2}mkdir -m 0700 -p \/notify\/ +the notifier refuses a looser directory$/m); assert.equal(readFileSync(join(dir, "mosaic-bus@.service"), "utf8"), render.stdout); assert.match(spawnSync(script, ["install", "--dir", dir, "--no-reload"], { encoding: "utf8" }).stdout, /unchanged: /); assert.match(spawnSync(script, ["uninstall", "--dir", dir, "--no-reload"], { encoding: "utf8" }).stdout, /removed: /); assert.equal(spawnSync(script, ["bogus"], { encoding: "utf8" }).status, 4); });