diff --git a/packages/cli/README.md b/packages/cli/README.md index 589078c3..0eb1423c 100644 --- a/packages/cli/README.md +++ b/packages/cli/README.md @@ -145,9 +145,10 @@ inbox through the reader capability and does two things: - **Digest.** It sends one digest a day at 08:00 America/Chicago. The hour comes from the IANA zone, so daylight saving time is handled. If the host starts after 08:00 and the day has no digest yet, the digest goes at once. - The digest lists the inbox and marks each blocking decision as DM sent or - DM pending. An empty inbox gets one line. Each message stays within - Discord's 2000 characters. + The digest lists the inbox. It marks each blocking decision with one of + three states: "DM sent", "DM pending" or "DM refused, not retried". An + empty inbox gets one line. Each message stays within Discord's 2000 + characters. The notifier only reads the bus. Its memory is the journal `/notify//sent.jsonl` (directory 0700, file 0600), with @@ -160,21 +161,41 @@ one line per send attempt: - A decision, or a day's digest, counts as sent once it has a `confirmed` line. -- A `refused` or `unknown` send is retried. The wait starts at 30 s and - doubles up to 30 min. A duplicate costs less than a miss. Every retry +- A `refused` or `unknown` send is retried. After an `unknown` outcome the + wait starts at 30 s and doubles up to 30 min. A duplicate costs less than + a miss. Every retry carries the same Discord nonce, so a retry inside Discord's dedupe window returns the first message. A DM's nonce comes from the decision id. A digest's comes from the business and the day. +- A definite refusal is a DM refused with an HTTP 4xx other than 429. The + next attempt waits the full 30 min, counted from the refusal's `at` in + the journal, so a restart does not shorten it. After 5 of them for one + decision, at least 2 hours apart end to end, the notifier appends one + `gave-up` line, logs it once and never sends that DM again. The digest + marks the decision "DM refused, not retried". The count comes from the + journal, so a restart keeps it. An `unknown` outcome (network, 5xx or + 429) retries without a limit (lead decisions 72 and 73). +- Recovery after a give-up is manual. Fixing the binding does not resend + the DM, which stays given up. The decision stays open in `mosaic inbox`, + the digest lists it, and the operator decides it with `mosaic decide`. - On open, a final line without its newline is a torn write. The notifier copies those bytes to `torn-.bin` in the same directory (0600, a new file, synced), then truncates `sent.jsonl` to its last newline and syncs it. It logs both steps. A crash between the two leaves the torn tail in place, and the next open repairs it with a second copy. The next append therefore starts on a line of its own. -- A malformed complete line refuses with exit 3 and changes nothing. -- The journal refuses with exit 3 when its directory is looser than 0700 or - not yours, when `sent.jsonl` is a symlink, or when the file is not a - regular 0600 file you own. +- A malformed complete line refuses with exit 3 and changes nothing. Each + line is type-checked: `at` an ISO timestamp, `kind` `dm` or `digest`, + `decision` a string (required for `dm`, null for `digest`), `day` a real + `YYYY-MM-DD` date (required for `digest`), `outcome` one of `confirmed`, + `refused`, `unknown` or `gave-up` (`gave-up` only for `dm`), `messageId` + a string (required when confirmed) or null, `status` an integer when + present. +- The journal refuses with exit 3 when its directory is looser than 0700, + not yours, a symlink (dangling or not) or not writable, or cannot be + created (for example, a parent is a file). It also refuses when + `sent.jsonl` is a symlink, a directory or not writable, or is not a + regular 0600 file you own. The message names the path. - No Discord channel or user id goes in the journal, a log line or an error. diff --git a/packages/cli/src/host.mjs b/packages/cli/src/host.mjs index 159160c0..083f8f1d 100644 --- a/packages/cli/src/host.mjs +++ b/packages/cli/src/host.mjs @@ -141,13 +141,13 @@ export async function startHost({ boot, business, notifier = null, bootTimeoutMs } catch (e) { notify.kill("SIGTERM"); await ended(notify, 5000); - broker.send({ op: "close" }); + if (broker.connected) broker.send({ op: "close" }, () => {}); await ended(broker, CLOSE_TIMEOUT_MS); throw e; } if (ok?.ok !== true) { await ended(notify, 5000); - broker.send({ op: "close" }); + 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); } @@ -181,11 +181,11 @@ export async function startHost({ boot, business, notifier = null, bootTimeoutMs closing = true; let result = code; if (notify) { - if (notify.connected) notify.send({ op: "stop" }); + 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 (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 }); diff --git a/packages/cli/src/notifier.mjs b/packages/cli/src/notifier.mjs index a4bb78d2..2a392b83 100644 --- a/packages/cli/src/notifier.mjs +++ b/packages/cli/src/notifier.mjs @@ -4,17 +4,25 @@ // bus; its memory is the journal `/notify//sent.jsonl` // (0600, in a 0700 directory), one line per send attempt: // -// {at, kind: "dm"|"digest", decision, day?, outcome: "confirmed"|"refused"|"unknown", messageId, status?} +// {at, kind: "dm"|"digest", decision, day?, outcome: "confirmed"|"refused"|"unknown"|"gave-up", 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. +// The exception is a definite refusal, an HTTP 4xx other than 429 (lead +// decision 72): the next try waits the full 30 min, counted from the +// refusal's line so a restart cannot shorten it (lead decision 73), and +// after five for one decision, counted from the journal, the notifier +// appends one `gave-up` line, logs once and stops sending that DM. Five +// tries span two hours, long enough to fix a binding or token typo. The +// digest marks it "DM refused, not retried"; recovery is the operator +// deciding it. Unknown outcomes (network, 5xx, 429) retry without a limit. // 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 { accessSync, 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"; @@ -25,6 +33,7 @@ export const POLL_MS = 30000; const BACKOFF_MS = 30000; const BACKOFF_MAX_MS = 30 * 60 * 1000; const LIMIT = 2000; +export const REFUSAL_LIMIT = 5; export const journalPath = (dataRoot, business) => join(dataRoot, "notify", business, "sent.jsonl"); @@ -62,14 +71,14 @@ export function dmContent(business, d) { return clip(lines.join("\n"), LIMIT); } -export function digestContent(business, day, inbox, dmSent) { +export function digestContent(business, day, inbox, dmSent, dmGaveUp = () => false) { 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 mark = d.blocking ? (dmSent(d.id) ? "[blocking, DM sent] " : dmGaveUp(d.id) ? "[blocking, DM refused, not retried] " : "[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; @@ -116,9 +125,32 @@ function copyTorn(dir, bytes, date) { } } +const OUTCOMES = ["confirmed", "refused", "unknown", "gave-up"]; +const DAY = /^\d{4}-\d{2}-\d{2}$/; +const isDay = (d) => typeof d === "string" && DAY.test(d) && !Number.isNaN(Date.parse(d)) && new Date(d).toISOString().startsWith(d); +const UNWRITABLE = ["EACCES", "EPERM", "EROFS"]; + +// The first field of a journal record that has the wrong type, or null +// (lead decision 72, J4). A confirmed line carries the message id. +function badField(r) { + if (!r || typeof r !== "object" || Array.isArray(r)) return "record"; + if (typeof r.at !== "string" || Number.isNaN(Date.parse(r.at)) || new Date(r.at).toISOString() !== r.at) return "at"; + if (!["dm", "digest"].includes(r.kind)) return "kind"; + if (r.kind === "dm" ? typeof r.decision !== "string" || r.decision === "" : r.decision !== null && r.decision !== undefined) return "decision"; + if (r.kind === "digest" ? !isDay(r.day) : r.day !== undefined && !isDay(r.day)) return "day"; + if (!OUTCOMES.includes(r.outcome) || (r.outcome === "gave-up" && r.kind !== "dm")) return "outcome"; + if (r.outcome === "confirmed" ? typeof r.messageId !== "string" : r.messageId !== null && typeof r.messageId !== "string") return "messageId"; + if (r.status !== undefined && !Number.isInteger(r.status)) return "status"; + return null; +} + +// A DM refused with an HTTP 4xx other than 429. +const definite = (r) => r.kind === "dm" && r.outcome === "refused" && Number.isInteger(r.status) && r.status >= 400 && r.status < 500 && r.status !== 429; + // 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. +// tighter, writable and not a symlink, and the file 0600, both owned by +// this user; a symlinked journal refuses. Every complete line must parse +// and type-check, 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-.bin, then the journal // is truncated to its last newline and fsynced, and both steps are logged. @@ -126,20 +158,49 @@ function copyTorn(dir, bytes, date) { // both; the second copy is harmless. export function openJournal(file, { log = () => {}, now = () => new Date() } = {}) { const dir = dirname(file); - mkdirSync(dir, { recursive: true, mode: 0o700 }); + try { + mkdirSync(dir, { recursive: true, mode: 0o700 }); + } catch (e) { + // A dangling link in place of the directory fails mkdir with ENOENT. + if (e.code === "ENOENT" && lstatSync(dir, { throwIfNoEntry: false })?.isSymbolicLink()) throw new CliError(`notify journal directory must not be a symlink: ${dir}`, 3); + if ([...UNWRITABLE, "ENOTDIR"].includes(e.code)) throw new CliError(`notify journal directory cannot be created (${e.code}): ${dir}`, 3); + if (e.code !== "EEXIST") throw e; + } const ds = lstatSync(dir); + if (ds.isSymbolicLink()) throw new CliError(`notify journal directory must not be a symlink: ${dir}`, 3); 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); } + try { + accessSync(dir, constants.W_OK | constants.X_OK); + } catch (e) { + if (UNWRITABLE.includes(e.code)) throw new CliError(`notify journal directory is not writable (${e.code}): ${dir}`, 3); + throw e; + } 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); + if (UNWRITABLE.includes(e.code)) throw new CliError(`notify journal is not writable (${e.code}): ${file}`, 3); + if (e.code === "EISDIR") throw new CliError(`notify journal must be a regular file, mode 0600, owned by this user: ${file}`, 3); throw e; } const sent = new Set(); const days = new Set(); + const refusals = new Map(); + const gaveUp = new Set(); + const refusedAt = new Map(); + const count = (r) => { + if (r.kind === "dm" && r.outcome === "gave-up") gaveUp.add(r.decision); + if (definite(r)) { + refusals.set(r.decision, (refusals.get(r.decision) ?? 0) + 1); + refusedAt.set(r.decision, Date.parse(r.at)); + } + if (r.outcome !== "confirmed") return; + if (r.kind === "dm") sent.add(r.decision); + else days.add(r.day); + }; try { const st = fstatSync(fd); if (!st.isFile() || st.uid !== process.getuid() || (st.mode & 0o777) !== 0o600) { @@ -156,12 +217,9 @@ export function openJournal(file, { log = () => {}, now = () => new Date() } = { } 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); + const bad = r === null ? "json" : badField(r); + if (bad) throw new CliError(`notify journal line ${i + 1} is malformed (${bad}): ${file}`, 3); + count(r); }); if (end < bytes.length) { const name = copyTorn(dir, bytes.subarray(end), now()); @@ -176,6 +234,9 @@ export function openJournal(file, { log = () => {}, now = () => new Date() } = { return { sent, days, + refusals, + gaveUp, + refusedAt, append(record) { const afd = openSync(file, O_WRONLY | O_APPEND | O_NOFOLLOW); try { @@ -183,9 +244,7 @@ export function openJournal(file, { log = () => {}, now = () => new Date() } = { } finally { closeSync(afd); } - if (record.outcome !== "confirmed") return; - if (record.kind === "dm") sent.add(record.decision); - else days.add(record.day); + count(record); }, }; } @@ -202,6 +261,13 @@ export function createNotifier({ business, dataRoot, inbox, direct, now = () => return b !== undefined && t < b.next; } + // One `gave-up` line and one log line; the decision is never sent again. + function giveUp(decision, t) { + journal.append({ at: new Date(t).toISOString(), kind: "dm", decision, outcome: "gave-up", messageId: null }); + backoff.delete(`dm:${decision}`); + log(`notify: dm ${shortId(decision)} refused ${journal.refusals.get(decision)} times; gave up, not retried`); + } + async function attempt(key, record, message) { const t = now().getTime(); try { @@ -212,10 +278,16 @@ export function createNotifier({ business, dataRoot, inbox, direct, now = () => } catch (err) { const kind = err?.kind === "refused" ? "refused" : "unknown"; const status = Number.isInteger(err?.details?.status) ? err.details.status : null; + const line = { at: new Date(t).toISOString(), ...record, outcome: kind, messageId: null, ...(status !== null ? { status } : {}) }; 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`); + const wait = definite(line) ? BACKOFF_MAX_MS : Math.min(BACKOFF_MS * 2 ** n, BACKOFF_MAX_MS); + backoff.set(key, { n, next: t + wait }); + journal.append(line); + if (record.kind === "dm" && (journal.refusals.get(record.decision) ?? 0) >= REFUSAL_LIMIT) { + giveUp(record.decision, t); + return false; + } + log(`notify: ${record.kind} ${kind}${status !== null ? ` (HTTP ${status})` : ""}; retry in ${Math.round(wait / 1000)} s`); return false; } } @@ -232,7 +304,15 @@ export function createNotifier({ business, dataRoot, inbox, direct, now = () => } const t = now(); for (const d of list) { - if (!d.blocking || journal.sent.has(d.id)) continue; + if (!d.blocking || journal.sent.has(d.id) || journal.gaveUp.has(d.id)) continue; + // A crash between the fifth refusal and its gave-up line. + if ((journal.refusals.get(d.id) ?? 0) >= REFUSAL_LIMIT) { + giveUp(d.id, t.getTime()); + continue; + } + // A definite refusal waits the full cap from its journal line, which + // the in-memory backoff loses on a restart. + if (t.getTime() < (journal.refusedAt.get(d.id) ?? -Infinity) + BACKOFF_MAX_MS) 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++; @@ -241,7 +321,7 @@ export function createNotifier({ business, dataRoot, inbox, direct, now = () => 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)); + const content = digestContent(business, day, list, (id) => journal.sent.has(id), (id) => journal.gaveUp.has(id)); if (await attempt(key, { kind: "digest", decision: null, day }, { content, nonce: digestNonce(business, day) })) done.digest = true; else done.failed++; } diff --git a/packages/cli/tests/host.test.mjs b/packages/cli/tests/host.test.mjs index 8d20b47f..f8f11ec9 100644 --- a/packages/cli/tests/host.test.mjs +++ b/packages/cli/tests/host.test.mjs @@ -1,6 +1,7 @@ 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"; @@ -48,6 +49,70 @@ async function until(fn, ms = 8000) { 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"); @@ -64,8 +129,12 @@ test("the host boots the broker, binds a launch in process, and the notifier DMs 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); @@ -84,12 +153,16 @@ test("the host boots the broker, binds a launch in process, and the notifier DMs 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 in a child's argv or environment, or in the state file. - for (const pid of Object.values(host.pids)) { - assert.ok(!procText(pid, "cmdline").includes(launch.cap)); - assert.ok(!procText(pid, "environ").includes(launch.cap)); + // 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.ok(!readFileSync(hostFile(f.dataRoot), "utf8").includes(launch.cap)); assert.equal(await host.close(0), 0); const journal = readFileSync(journalPath(f.dataRoot, "acme"), "utf8"); @@ -107,6 +180,7 @@ test("a notifier that dies takes the host down with exit 1, so the unit restarts 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/); @@ -114,6 +188,21 @@ test("a notifier that dies takes the host down with exit 1, so the unit restarts 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); @@ -127,9 +216,111 @@ test("a notifier that refuses stops the broker and the host refuses with exit 3" 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"); @@ -215,6 +406,7 @@ test("bus-service.sh renders the unit and installs it into a given directory", ( 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: /); diff --git a/packages/cli/tests/notifier.test.mjs b/packages/cli/tests/notifier.test.mjs index 99127662..3b65af47 100644 --- a/packages/cli/tests/notifier.test.mjs +++ b/packages/cli/tests/notifier.test.mjs @@ -1,12 +1,15 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { appendFileSync, chmodSync, mkdirSync, readdirSync, readFileSync, statSync, symlinkSync, writeFileSync } from "node:fs"; +import { appendFileSync, chmodSync, mkdirSync, readdirSync, readFileSync, rmSync, statSync, symlinkSync, writeFileSync } from "node:fs"; import { dirname, join } from "node:path"; -import { createNotifier, digestContent, digestNonce, dmNonce, journalPath, openJournal, runLoop, zoned } from "../src/notifier.mjs"; +import { createNotifier, digestContent, digestNonce, dmNonce, journalPath, openJournal, POLL_MS, REFUSAL_LIMIT, runLoop, zoned } from "../src/notifier.mjs"; import { RestOutcome } from "../../discord/src/rest.mjs"; import { broker, tmp } from "./helpers.mjs"; +const AT = "2026-10-08T12:00:00.000Z"; + // Discord-side fake: records sends, answers from a script (default: ok). +// An entry is "ok", "refused" (HTTP 403), "unknown", or {kind, status}. function fakeDirect(script = []) { const sends = []; let n = 0; @@ -16,7 +19,8 @@ function fakeDirect(script = []) { sends.push(m); const next = script.shift() ?? "ok"; if (next === "ok") return { messageId: `30000000000000${String(++n).padStart(4, "0")}` }; - throw new RestOutcome(next, `dm: ${next}`, { status: next === "refused" ? 403 : null }); + const { kind, status } = typeof next === "string" ? { kind: next, status: next === "refused" ? 403 : null } : next; + throw new RestOutcome(kind, `dm: ${kind}`, { status }); }, }; } @@ -91,14 +95,117 @@ test("a failed DM is journaled, backs off, and is retried until it lands", async assert.equal((await s.notifier.tick()).failed, 0, "inside the first 30 s backoff"); s.time.advance(25_000); assert.equal((await s.notifier.tick()).failed, 1, "second attempt refused"); - s.time.advance(45_000); - assert.equal((await s.notifier.tick()).dms, 0, "inside the 60 s backoff"); - s.time.advance(20_000); + s.time.advance(29 * 60_000); + assert.equal((await s.notifier.tick()).dms, 0, "a definite refusal waits the full 30 min, not the 60 s step"); + s.time.advance(60_000); assert.equal((await s.notifier.tick()).dms, 1); assert.deepEqual(s.journal().map((r) => r.outcome), ["unknown", "refused", "confirmed"]); assert.equal(s.journal()[1].status, 403); assert.equal(new Set(s.direct.sends.map((m) => m.nonce)).size, 1, "every retry reuses the nonce"); - assert.ok(s.logs.some((l) => /retry in 30 s/.test(l))); + assert.ok(s.logs.some((l) => /unknown; retry in 30 s/.test(l))); + assert.ok(s.logs.some((l) => /refused \(HTTP 403\); retry in 1800 s/.test(l))); +}); + +// 05:00Z is 00:00 Chicago: eight hours of polls before the digest is due. +const MIDNIGHT = "2026-10-08T05:00:00Z"; +const PAST_BACKOFF = 31 * 60_000; + +test("five definite refusals stop a DM: one gave-up line, one log line, and a restart keeps the count", async (t) => { + assert.equal(REFUSAL_LIMIT, 5); + const s = setup(t, MIDNIGHT, ["refused", "refused", "refused", { kind: "refused", status: 404 }, { kind: "refused", status: 400 }]); + const d = s.raise("git.push.protected", { target: "refactor", blocking: true, task_ref: "vikunja:1/7" }); + for (let i = 0; i < 3; i++) { + assert.equal((await s.notifier.tick()).failed, 1); + s.time.advance(PAST_BACKOFF); + } + const after = s.make(); + for (let i = 0; i < 2; i++) { + assert.equal((await after.tick()).failed, 1, "the restart did not reset the count"); + s.time.advance(PAST_BACKOFF); + } + assert.equal(s.direct.sends.length, 5); + assert.deepEqual(s.journal().map((r) => r.outcome), ["refused", "refused", "refused", "refused", "refused", "gave-up"]); + assert.deepEqual(s.journal().at(-1), { at: s.journal().at(-1).at, kind: "dm", decision: d.id, outcome: "gave-up", messageId: null }); + assert.deepEqual(await after.tick(), { dms: 0, digest: false, failed: 0 }); + assert.deepEqual(await s.make().tick(), { dms: 0, digest: false, failed: 0 }); + assert.equal(s.direct.sends.length, 5, "never sent again, before or after a restart"); + assert.equal(s.journal().length, 6); + const gave = s.logs.filter((l) => /gave up/.test(l)); + assert.deepEqual(gave, [`notify: dm ${d.id.slice(0, 8)} refused 5 times; gave up, not retried`]); + assert.equal(s.logs.filter((l) => /retry in/.test(l)).length, 4, "the fifth refusal logs the give-up, not a retry"); + s.time.clock.t = new Date("2026-10-08T13:00:00Z"); + assert.equal((await s.make().tick()).digest, true); + assert.match(s.direct.sends.at(-1).content, new RegExp(`\\[blocking, DM refused, not retried\\] ${d.id.slice(0, 8)} git\\.push\\.protected`)); +}); + +test("429s, 5xx-style unknowns and refusals without a status never count toward the limit", async (t) => { + const script = [ + ...Array(6).fill({ kind: "refused", status: 429 }), + ...Array(3).fill("unknown"), + { kind: "refused", status: null }, + ...Array(REFUSAL_LIMIT - 1).fill("refused"), + "ok", + ]; + const s = setup(t, MIDNIGHT, [...script]); + s.raise("git.push.protected", { target: "refactor", blocking: true, task_ref: "vikunja:1/7" }); + for (let i = 0; i < script.length; i++) { + await s.notifier.tick(); + s.time.advance(PAST_BACKOFF); + } + assert.equal(s.direct.sends.length, script.length); + assert.equal(s.journal().at(-1).outcome, "confirmed"); + assert.ok(!s.journal().some((r) => r.outcome === "gave-up")); + assert.ok(!s.logs.some((l) => /gave up/.test(l))); +}); + +test("a crash between the fifth refusal and its gave-up line: the next poll appends it and sends nothing", async (t) => { + const s = setup(t, MIDNIGHT); + const d = s.raise("git.push.protected", { target: "refactor", blocking: true, task_ref: "vikunja:1/7" }); + const file = journalPath(s.dataRoot, "demo"); + const line = { at: AT, kind: "dm", decision: d.id, outcome: "refused", messageId: null, status: 403 }; + appendFileSync(file, `${JSON.stringify(line)}\n`.repeat(REFUSAL_LIMIT)); + assert.deepEqual(await s.make().tick(), { dms: 0, digest: false, failed: 0 }); + assert.equal(s.direct.sends.length, 0); + assert.equal(s.journal().at(-1).outcome, "gave-up"); + assert.equal(s.logs.filter((l) => /gave up/.test(l)).length, 1); + await s.make().tick(); + assert.equal(s.journal().length, REFUSAL_LIMIT + 1, "one gave-up line only"); +}); + +// Polls every POLL_MS from the current clock up to `until` ms later; returns +// the minute (from `from`) of each new send. +async function pollUntil(s, notifier, from, until, onPoll = () => {}) { + const minutes = []; + for (let elapsed = s.time.clock.t.getTime() - from; elapsed <= until; elapsed += POLL_MS) { + const before = s.direct.sends.length; + await notifier.tick(); + if (s.direct.sends.length > before) minutes.push(elapsed / 60_000); + onPoll(elapsed); + s.time.advance(POLL_MS); + } + return minutes; +} + +test("polled every POLL_MS against a permanent 403, a DM is sent at 0, 30, 60, 90 and 120 min and gives up only then", async (t) => { + const s = setup(t, MIDNIGHT, Array(10).fill("refused")); + s.raise("git.push.protected", { target: "refactor", blocking: true, task_ref: "vikunja:1/7" }); + const from = s.time.clock.t.getTime(); + const gaveUp = () => s.journal().some((r) => r.outcome === "gave-up"); + const minutes = await pollUntil(s, s.notifier, from, 120 * 60_000 + POLL_MS, (elapsed) => { + assert.equal(gaveUp(), elapsed >= 120 * 60_000, `gave-up line at +${elapsed / 1000} s`); + }); + assert.deepEqual(minutes, [0, 30, 60, 90, 120]); + assert.equal(s.logs.filter((l) => /refused \(HTTP 403\); retry in 1800 s/.test(l)).length, 4); +}); + +test("a restart after the second refusal does not send before that refusal's 30 min are up", async (t) => { + const s = setup(t, MIDNIGHT, Array(10).fill("refused")); + s.raise("git.push.protected", { target: "refactor", blocking: true, task_ref: "vikunja:1/7" }); + const from = s.time.clock.t.getTime(); + assert.deepEqual(await pollUntil(s, s.notifier, from, 30 * 60_000), [0, 30]); + s.time.advance(60_000); + const after = s.make(); + assert.deepEqual(await pollUntil(s, after, from, 60 * 60_000), [60], "nothing between the restart and +60 min"); }); test("the digest goes at 08:00 Chicago once a day, with blocking ones marked as DM'd", async (t) => { @@ -150,7 +257,7 @@ const FRAGMENT = '{"at":"x","kind":"dm","dec'; test("the journal: a torn tail is copied out and truncated, so an append after it reopens cleanly", (t) => { const file = journalPath(tmp(t), "demo"); - openJournal(file).append({ at: "x", kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }); + openJournal(file).append({ at: AT, kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }); const good = readFileSync(file); appendFileSync(file, FRAGMENT); const logs = []; @@ -164,7 +271,7 @@ test("the journal: a torn tail is copied out and truncated, so an append after i assert.equal(logs.length, 2); assert.match(logs[0], /copied a torn final line \(26 bytes\) to torn-20261008T235212345Z\.bin/); assert.match(logs[1], /truncated .* to its last newline/); - j.append({ at: "y", kind: "dm", decision: "b", outcome: "confirmed", messageId: "2" }); + j.append({ at: AT, kind: "dm", decision: "b", outcome: "confirmed", messageId: "2" }); const again = []; assert.deepEqual([...openJournal(file, { log: (l) => again.push(l) }).sent], ["a", "b"]); assert.deepEqual(again, [], "nothing torn the second time"); @@ -172,7 +279,7 @@ test("the journal: a torn tail is copied out and truncated, so an append after i test("the journal: a crash between the copy and the truncate leaves a tail the next open repairs", (t) => { const file = journalPath(tmp(t), "demo"); - openJournal(file).append({ at: "x", kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }); + openJournal(file).append({ at: AT, kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }); appendFileSync(file, FRAGMENT); const now = () => new Date("2026-10-08T23:52:12.345Z"); // The log after step 1 throws: the process dies before step 2. @@ -182,7 +289,7 @@ test("the journal: a crash between the copy and the truncate leaves a tail the n const j = openJournal(file, { now }); assert.deepEqual(tornFiles(file), ["torn-20261008T235212345Z-1.bin", "torn-20261008T235212345Z.bin"], "a second copy, the first kept"); for (const n of tornFiles(file)) assert.equal(readFileSync(join(dirname(file), n), "utf8"), FRAGMENT); - j.append({ at: "y", kind: "dm", decision: "b", outcome: "confirmed", messageId: "2" }); + j.append({ at: AT, kind: "dm", decision: "b", outcome: "confirmed", messageId: "2" }); assert.deepEqual([...openJournal(file).sent], ["a", "b"]); }); @@ -222,6 +329,108 @@ test("the journal: a loose file mode, a loose directory or a symlinked journal r assert.throws(() => openJournal(linked), (e) => e.exitCode === 3 && /must not be a symlink/.test(e.message)); }); +test("the journal: a line with a wrong type refuses with exit 3 and names the field", (t) => { + const file = journalPath(tmp(t), "demo"); + openJournal(file); + const dm = { at: AT, kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }; + const digest = { at: AT, kind: "digest", decision: null, day: "2026-10-08", outcome: "confirmed", messageId: "2" }; + const bad = [ + ["record", [1]], + ["at", { ...dm, at: "x" }], + ["at", { ...dm, at: "2026-10-08" }], + ["at", { ...dm, at: 1 }], + ["kind", { ...dm, kind: "dg" }], + ["decision", { ...dm, decision: 7 }], + ["decision", { ...dm, decision: "" }], + ["decision", { ...digest, decision: "a" }], + ["day", { ...digest, day: undefined }], + ["day", { ...digest, day: "2026-10-8" }], + ["day", { ...digest, day: "2026-13-45" }], + ["day", { ...digest, day: "2026-02-30" }], + ["day", { ...dm, day: 20261008 }], + ["outcome", { ...dm, outcome: "sent" }], + ["outcome", { ...digest, outcome: "gave-up", messageId: null }], + ["messageId", { ...dm, messageId: null }], + ["messageId", { ...dm, outcome: "refused", messageId: 1 }], + ["status", { ...dm, outcome: "refused", messageId: null, status: "403" }], + ["status", { ...dm, outcome: "refused", messageId: null, status: 403.5 }], + ]; + for (const [field, r] of bad) { + const text = `${JSON.stringify(dm)}\n${JSON.stringify(r)}\n`; + writeFileSync(file, text); + assert.throws(() => openJournal(file), (e) => e.exitCode === 3 && e.message === `notify journal line 2 is malformed (${field}): ${file}`, `${field}: ${JSON.stringify(r)}`); + assert.equal(readFileSync(file, "utf8"), text, "a refusal changes nothing"); + } + const good = [dm, digest, { ...dm, outcome: "refused", messageId: null, status: 403 }, { ...dm, outcome: "gave-up", messageId: null }, { ...digest, outcome: "unknown", messageId: null }]; + writeFileSync(file, good.map((r) => `${JSON.stringify(r)}\n`).join("")); + const j = openJournal(file); + assert.deepEqual([...j.sent], ["a"]); + assert.deepEqual([...j.days], ["2026-10-08"]); + assert.deepEqual([...j.gaveUp], ["a"]); + assert.deepEqual([...j.refusals], [["a", 1]]); +}); + +test("the journal: a symlinked directory refuses and says it is a link", (t) => { + const root = tmp(t); + const real = join(root, "real"); + mkdirSync(real, { mode: 0o700 }); + const file = journalPath(root, "demo"); + mkdirSync(dirname(dirname(file)), { mode: 0o700 }); + symlinkSync(real, dirname(file)); + assert.throws(() => openJournal(file), (e) => e.exitCode === 3 && e.message === `notify journal directory must not be a symlink: ${dirname(file)}`); + assert.deepEqual(readdirSync(real), [], "nothing was created through the link"); +}); + +test("the journal: a dangling directory link, a parent that is a file and a journal that is a directory each refuse with exit 3", (t) => { + const root = tmp(t); + const dangling = journalPath(root, "demo"); + mkdirSync(dirname(dirname(dangling)), { mode: 0o700 }); + symlinkSync(join(root, "missing"), dirname(dangling)); + assert.throws(() => openJournal(dangling), (e) => e.exitCode === 3 && e.message === `notify journal directory must not be a symlink: ${dirname(dangling)}`); + const under = join(root, "plain", "notify", "demo", "sent.jsonl"); + writeFileSync(join(root, "plain"), ""); + assert.throws(() => openJournal(under), (e) => e.exitCode === 3 && e.message === `notify journal directory cannot be created (ENOTDIR): ${dirname(under)}`); + const asDir = journalPath(root, "acme"); + mkdirSync(asDir, { recursive: true, mode: 0o700 }); + assert.throws(() => openJournal(asDir), (e) => e.exitCode === 3 && e.message === `notify journal must be a regular file, mode 0600, owned by this user: ${asDir}`); +}); + +test("the journal: an append after the file was swapped for a symlink refuses and writes nothing through it", (t) => { + const root = tmp(t); + const file = journalPath(root, "demo"); + const j = openJournal(file); + const other = join(root, "elsewhere.jsonl"); + writeFileSync(other, "", { mode: 0o600 }); + rmSync(file); + symlinkSync(other, file); + assert.throws(() => j.append({ at: AT, kind: "dm", decision: "a", outcome: "confirmed", messageId: "1" }), (e) => e.code === "ELOOP"); + assert.equal(readFileSync(other, "utf8"), ""); +}); + +test("the journal: a directory it cannot write or create refuses with exit 3 and names the path", { skip: process.getuid() === 0 && "root ignores the modes" }, (t) => { + const root = tmp(t); + const file = journalPath(root, "demo"); + mkdirSync(dirname(file), { recursive: true, mode: 0o700 }); + // Modes are restored in finally: tmp()'s cleanup runs first among the after hooks. + chmodSync(dirname(file), 0o500); + try { + assert.throws(() => openJournal(file), (e) => e.exitCode === 3 && e.message === `notify journal directory is not writable (EACCES): ${dirname(file)}`); + } finally { + chmodSync(dirname(file), 0o700); + } + openJournal(file); + chmodSync(file, 0o400); + assert.throws(() => openJournal(file), (e) => e.exitCode === 3 && e.message === `notify journal is not writable (EACCES): ${file}`); + const parent = join(root, "notify"); + const other = journalPath(root, "acme"); + chmodSync(parent, 0o500); + try { + assert.throws(() => openJournal(other), (e) => e.exitCode === 3 && e.message === `notify journal directory cannot be created (EACCES): ${dirname(other)}`); + } finally { + chmodSync(parent, 0o700); + } +}); + test("digest content stays within Discord's 2000 characters", () => { const inbox = Array.from({ length: 60 }, (_, i) => ({ id: `${String(i).padStart(8, "0")}-x`, action: "deploy", question: "q".repeat(300), blocking: i % 2 === 0 })); const text = digestContent("demo", "2026-10-08", inbox, () => true); diff --git a/packages/cli/tests/trackers-boot.test.mjs b/packages/cli/tests/trackers-boot.test.mjs new file mode 100644 index 00000000..b7710013 --- /dev/null +++ b/packages/cli/tests/trackers-boot.test.mjs @@ -0,0 +1,66 @@ +// Sage's row 39 gate check, not part of the candidate: the S4 host boots the +// S3 adapter through the real process.mjs when bootConfig emits trackers. +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { chmodSync, mkdirSync, writeFileSync } from "node:fs"; +import { join } from "node:path"; +import { Client } from "../../bus/src/client.mjs"; +import { bootConfig, loadSystem } from "../src/config.mjs"; +import { startHost, startTimeOf } from "../src/host.mjs"; +import { writeJson } from "../../business/tests/helpers.mjs"; +import { FakeVikunja } from "../../tasks/src/fake.mjs"; +import { SCOPES } from "../../tasks/tests/world.mjs"; +import { fixture, tmp } from "./helpers.mjs"; + +test("bootConfig trackers reach the S3 adapter in the real broker child, which goes ready against a fake Vikunja", async (t) => { + const root = tmp(t); + const fake = new FakeVikunja(); + const srv = await fake.listen(); + t.after(() => srv.close()); + const owner = fake.user("svc-acme"); + const bots = {}; + for (const r of ["sync", "pm", "cto", "coder", "reviewer"]) bots[r] = fake.user(`bot-acme-${r}`, { owner }); + const project = fake.project("Acme", owner); + for (const r of ["pm", "cto", "coder", "reviewer"]) fake.share(project, bots[r], 1); + fake.share(project, bots.sync, 0); + fake.install(project); + + const f = fixture(root, "acme", (doc) => { + doc.vars["tracker.baseUrl"] = srv.url; + doc.tracker.sync.botId = bots.sync; + for (const r of Object.keys(doc.roles)) doc.roles[r].tracker.botId = bots[r]; + return doc; + }); + const put = (file, token) => { + writeFileSync(file, token + "\n"); + chmodSync(file, 0o600); + }; + put(f.doc.tracker.sync.credentials.vikunja.file, fake.token(bots.sync, SCOPES.sync)); + for (const r of Object.keys(f.doc.roles)) put(f.doc.roles[r].credentials.vikunja.file, fake.token(bots[r], r === "pm" ? SCOPES.pm : SCOPES.worker)); + writeJson(join(root, "project", ".mosaic", "project.json"), { projectVersion: 1, id: "stack", vars: { "tracker.project": project } }); + + mkdirSync(f.dataRoot, { recursive: true, mode: 0o700 }); + const boot = bootConfig({ system: loadSystem({ env: f.env }), businessId: "acme", env: f.env }); + assert.deepEqual(boot.trackers, { acme: { baseUrl: srv.url, project, pollSeconds: 60, reconcileMinutes: 60 } }); + const host = await startHost({ boot, business: "acme", log: () => {} }); + t.after(() => host.close(0)); + + const launch = await host.bindLaunch({ business: "acme", role: "pm", run: "pm-run", harness: "pi", pid: process.pid, startTime: startTimeOf(process.pid) }); + const pm = new Client({ path: host.path, cap: launch.cap }); + await pm.call("role.claim"); + let code; + const end = Date.now() + 15000; + do { + code = await pm.call("task.close", {}).then(() => "ok", (e) => e.code); + if (code !== "tracker-starting") break; + await new Promise((r) => setTimeout(r, 100)); + } while (Date.now() < end); + console.log(`task.close {} answered: ${code}; fake saw ${fake.requests.length} requests, first ${fake.requests.slice(0, 3).map((r) => `${r.method} ${r.path} ${r.status}`).join(", ")}`); + assert.ok(fake.requests.some((r) => r.path === "/info"), "the adapter called the fake"); + assert.doesNotMatch(code, /^(tracker-|credential-|scope-too-broad)/, "the adapter is ready, so the verb fails on its own arguments"); + const before = fake.requests.length; + const missing = await pm.call("task.close", { task_ref: `vikunja:${project}/999`, verdict: "gate check" }).then(() => "ok", (e) => e.code); + console.log(`task.close on a missing task answered: ${missing}; it made ${fake.requests.slice(before).map((r) => `${r.method} ${r.path} ${r.status}`).join(", ")}`); + assert.ok(fake.requests.length > before, "a well-formed verb reached Vikunja through the adapter"); + assert.equal(await host.close(0), 0); +}); diff --git a/packages/discord/tests/journal.test.mjs b/packages/discord/tests/journal.test.mjs index e3443510..681d8b83 100644 --- a/packages/discord/tests/journal.test.mjs +++ b/packages/discord/tests/journal.test.mjs @@ -17,6 +17,15 @@ function journal() { return dir; } +// A pid that is provably dead: a child spawned and reaped here, then +// confirmed gone with kill(pid, 0). A fixed number can belong to a live process. +function deadPid() { + for (;;) { + const { pid } = spawnSync(process.execPath, ["-e", ""], { stdio: "ignore" }); + if (Number.isSafeInteger(pid) && !pidAlive(pid)) return pid; + } +} + function publish(dir, rec) { mkdirSync(lockPath(dir), { recursive: true }); writeFileSync(ownerPath(dir), JSON.stringify(rec) + "\n", { mode: 0o600 }); @@ -76,7 +85,7 @@ test("lock: the claim is exclusive; a second start against a live owner refuses" test("lock: a stale lock (dead owner, reused pid, or record without start) refuses run and is never signaled; only unlock clears it", () => { const dir = journal(); - const dead = { pid: 2 ** 22 - 7, start: "1", boot: bootId() }; + const dead = { pid: deadPid(), start: "1", boot: bootId() }; publish(dir, dead); assert.equal(stopTarget(dir), null); assert.throws(() => writePid(dir, process.pid), (err) => err instanceof DiscordError && /which is gone/.test(err.message) && /unlock/.test(err.message)); @@ -116,7 +125,7 @@ test("lock: a stale lock (dead owner, reused pid, or record without start) refus assert.throws(() => unlock(dir), /cannot be verified; nothing removed/); rmSync(stopPath(dir)); rmSync(lockPath(dir), { recursive: true }); - publish(dir, { pid: 2 ** 22 - 7, start: "1" }); + publish(dir, { pid: deadPid(), start: "1" }); assert.equal(ownerState(readPid(dir)), "dead"); unlock(dir); rmSync(stopPath(dir)); @@ -227,8 +236,9 @@ test("lock: identity syntax; only canonical unsigned decimal start ticks and low test("lock: a process whose start marker or boot id cannot be read refuses to claim", () => { const dir = journal(); - assert.equal(processStart(2 ** 22 - 7), null); - assert.throws(() => writePid(dir, 2 ** 22 - 7), (err) => err instanceof DiscordError && /start time or the boot id/.test(err.message)); + const gone = deadPid(); + assert.equal(processStart(gone), null); + assert.throws(() => writePid(dir, gone), (err) => err instanceof DiscordError && /start time or the boot id/.test(err.message)); assert.equal(existsSync(lockPath(dir)), false, "nothing was left behind"); const noBoot = (pid) => ({ start: processStart(pid), boot: null }); assert.throws(() => writePid(dir, process.pid, { identity: noBoot }), /start time or the boot id/); @@ -276,7 +286,7 @@ test("lock: four processes racing for the same binding; exactly one claims it an test("lock: stale handoff; concurrent starts over a stale lock all refuse, nothing reclaims, one unlock then exactly one live owner", async () => { const dir = journal(); - const stale = { pid: 2 ** 22 - 7, start: "1", boot: bootId() }; + const stale = { pid: deadPid(), start: "1", boot: bootId() }; publish(dir, stale); // Several starts race over the stale lock: none may reclaim it. const r1 = await race(dir, 4, "stale"); @@ -315,7 +325,7 @@ test("lock: four-party schedule; claims landing inside an unlock's gap never sur const worker = fileURLToPath(new URL("../fixtures/claim-worker.mjs", import.meta.url)); const go = join(dir, "go"); writeFileSync(go, ""); - const stale = { pid: 2 ** 22 - 7, start: "1", boot: bootId() }; + const stale = { pid: deadPid(), start: "1", boot: bootId() }; publish(dir, stale); const claims = []; const claim = (tag) => { diff --git a/scripts/bus-service.sh b/scripts/bus-service.sh index 4d8c4b95..a60b3610 100755 --- a/scripts/bus-service.sh +++ b/scripts/bus-service.sh @@ -68,6 +68,7 @@ install_unit() { fi cat </notify/ the notifier refuses a looser directory write /notify//notify.json, mode 0600: {"notifyVersion": 1, "binding": ""} or "binding": null for no DMs systemctl --user enable --now mosaic-bus@ start now and at login