// The spool — where an event this module could not write waits until it can (novox/hq issue 276). // // **A handler takes an event by returning and asks for it again by throwing** (mesh-sdk README, // "Taking an event"). A write that failed is therefore thrown, and the bus offers the event again five // seconds later, up to the consumer's maximum deliveries — after which the bus gives it up. For a // module whose whole job is to keep what it is sent, giving up is losing, so: // // - **The first failed write puts the event in the spool**, one file per event, written and synced // before the handler throws. From then on the event is on this machine's disk, whatever the bus does. // - **The spool counts the failed deliveries**, because the runtime does not hand a handler the bus's // delivery count. On the last one the event is **taken**: it is in the spool, and the bus giving it // up would only raise a condition about an event that is not lost. // - **A background pass replays the spool** into the store once writes work again, oldest first, and // stops at the first failure. A write that succeeds — by a redelivery or by a replay — removes the // event from the spool. The module's writes are idempotent by the event's id, so a replay and a // redelivery of the same event write it once. // - **Over its bound** — too many events, or one waiting too long — the spool still keeps every event, // but the last delivery is no longer taken: the bus gives it up and the controller raises // `max-deliveries` for this module's consumer (to-be 45, S9). That is the one condition the mesh // has today that names a consumer which cannot keep up; a module has no standing of its own to report // (ADR 0224's is a provider's), so the bound borrows that one rather than inventing a second. // // The same file lives in audit-logger and model-usage; a third copy belongs in the SDK instead. import { mkdir, open, readdir, readFile, rename, unlink } from "node:fs/promises"; import { createHash } from "node:crypto"; import { join } from "node:path"; import type { Event } from "@novox/mesh-sdk/events"; /** The controller's maximum deliveries for a module's consumer (mesh-controller * internal/broker/derived.go). The spool takes an event on this, its last, failed delivery. */ export const MAX_DELIVERIES = 5; /** More events spooled than this is over the bound. */ export const BOUND_COUNT = 1000; /** An event spooled longer than this is over the bound. */ export const BOUND_AGE_MS = 30 * 60 * 1000; /** One event waiting in the spool. */ export interface SpoolEntry { key: string; event: Event; /** Failed deliveries seen so far. */ attempts: number; /** When it first failed, and last. */ first: string; last: string; /** The last failure's words. */ error: string; /** Taken on its last delivery: only the spool holds it now. */ held: boolean; } /** What the spool holds, for the module's status tool and its log. */ export interface Standing { dir: string; spooled: number; held: number; oldest?: string; lastError?: string; overBound: boolean; why?: string; bound: { count: number; ageMinutes: number }; } export interface SpoolOptions { maxDeliveries?: number; boundCount?: number; boundAgeMs?: number; now?: () => Date; } /** What a failed write became: asked for again, held and taken, or asked for again past the bound. */ export type Verdict = "retry" | "held" | "over-bound"; /** The key an event is spooled and deduplicated by: its x-event-id when it is a safe file name, else a * hash of what identifies it (an event without an id is still one event). */ export function keyOf(event: Pick): string { if (event.id && /^[A-Za-z0-9_-][A-Za-z0-9._-]{0,127}$/.test(event.id)) return event.id; const h = createHash("sha256"); for (const part of [event.id ?? "", event.type, event.source, event.node, event.at, JSON.stringify(event.body ?? null)]) { h.update(String(part)); h.update("\0"); } return "h-" + h.digest("hex"); } export class Spool { private readonly entries = new Map(); private chain: Promise = Promise.resolve(); private readonly max: number; private readonly boundCount: number; private readonly boundAgeMs: number; private readonly now: () => Date; private constructor(readonly dir: string, opts: SpoolOptions) { this.max = opts.maxDeliveries ?? MAX_DELIVERIES; this.boundCount = opts.boundCount ?? BOUND_COUNT; this.boundAgeMs = opts.boundAgeMs ?? BOUND_AGE_MS; this.now = opts.now ?? (() => new Date()); } /** Open the spool in `dir`, creating it, and read what it already holds. */ static async open(dir: string, opts: SpoolOptions = {}): Promise { const s = new Spool(dir, opts); await mkdir(dir, { recursive: true, mode: 0o700 }); await s.load(true); return s; } /** Read a spool another process writes, without changing it — for the status tool. */ static async read(dir: string, opts: SpoolOptions = {}): Promise { const s = new Spool(dir, opts); await s.load(false); return s; } get size(): number { return this.entries.size; } has(key: string): boolean { return this.entries.has(key); } list(): SpoolEntry[] { return [...this.entries.values()].sort((a, b) => a.first.localeCompare(b.first)); } /** Record a failed write of `event`. Synced to disk before it returns; throws if it could not be. */ failed(event: Event, err: unknown): Promise { return this.serial(async () => { const key = keyOf(event); const prev = this.entries.get(key); const at = this.now().toISOString(); const entry: SpoolEntry = { key, event, attempts: (prev?.attempts ?? 0) + 1, first: prev?.first ?? at, last: at, error: String(err instanceof Error ? err.message : err).slice(0, 500), held: false, }; const last = entry.attempts >= this.max; let oldest = entry.first; for (const e of this.entries.values()) if (e.first < oldest) oldest = e.first; const over = this.over(prev ? this.entries.size : this.entries.size + 1, oldest); entry.held = last && !over; await this.persist(entry); this.entries.set(key, entry); return last ? (over ? "over-bound" : "held") : "retry"; }); } /** The event was written: it need not wait any longer. */ done(key: string): Promise { if (!this.entries.has(key)) return Promise.resolve(); return this.serial(async () => { if (!this.entries.has(key)) return; await unlink(this.file(key)).catch((e: NodeJS.ErrnoException) => { if (e.code !== "ENOENT") throw e; }); this.entries.delete(key); }); } /** Write every spooled event again, oldest first, stopping at the first that still fails. */ async replay(write: (event: Event) => Promise): Promise<{ replayed: number; left: number; error?: string }> { let replayed = 0; for (const entry of this.list()) { try { await write(entry.event); } catch (err) { return { replayed, left: this.entries.size, error: String(err instanceof Error ? err.message : err) }; } await this.done(entry.key); replayed++; } return { replayed, left: this.entries.size }; } standing(): Standing { const list = this.list(); const oldest = list[0]?.first; const newest = list.reduce((n, e) => (!n || e.last > n.last ? e : n), undefined); const over = this.over(list.length, oldest); return { dir: this.dir, spooled: list.length, held: list.filter((e) => e.held).length, ...(oldest ? { oldest } : {}), ...(newest ? { lastError: newest.error } : {}), overBound: over, ...(over ? { why: this.why(list.length, oldest) } : {}), bound: { count: this.boundCount, ageMinutes: Math.round(this.boundAgeMs / 60000) }, }; } private over(count: number, oldest: string | undefined): boolean { return this.why(count, oldest) !== undefined; } private why(count: number, oldest: string | undefined): string | undefined { if (count > this.boundCount) return `${count} events spooled, more than ${this.boundCount}`; if (oldest && this.now().getTime() - Date.parse(oldest) > this.boundAgeMs) { return `an event has waited since ${oldest}, longer than ${Math.round(this.boundAgeMs / 60000)} minutes`; } return undefined; } private file(key: string): string { return join(this.dir, `${key}.json`); } private async load(tidy: boolean): Promise { let names: string[]; try { names = await readdir(this.dir); } catch (e) { if ((e as NodeJS.ErrnoException).code === "ENOENT") return; throw e; } for (const name of names) { if (name.endsWith(".tmp")) { // A write cut short: the event it was for was never answered as taken, so the bus has it. if (tidy) await unlink(join(this.dir, name)).catch(() => {}); continue; } if (!name.endsWith(".json")) continue; try { const entry = JSON.parse(await readFile(join(this.dir, name), "utf8")) as SpoolEntry; if (entry?.key && entry.event) this.entries.set(entry.key, entry); } catch (err) { // Kept on disk for a person, never deleted: it may be the only copy of an event. console.error(`[spool] ${join(this.dir, name)} cannot be read and is left as it is: ${err}`); } } } /** Written to a temporary file, synced, renamed over the entry, and the directory synced. */ private async persist(entry: SpoolEntry): Promise { const final = this.file(entry.key); const tmp = `${final}.tmp`; const fh = await open(tmp, "w", 0o600); try { await fh.writeFile(JSON.stringify(entry) + "\n"); await fh.sync(); } finally { await fh.close(); } await rename(tmp, final); const dh = await open(this.dir, "r"); try { await dh.sync(); } catch { // Not every filesystem syncs a directory; the rename stands either way. } finally { await dh.close(); } } private serial(fn: () => Promise): Promise { const run = this.chain.then(fn, fn); this.chain = run.catch(() => {}); return run; } } /** * The handler's whole discipline: write the event; if that fails, spool it and ask for it again — * or, on its last delivery, take it, because the spool now holds it. A write that cannot even be * spooled is asked for again, which is all that is left. */ export async function takeOrSpool( event: Event, write: (event: Event) => Promise, spool: Spool, say: (line: string) => void = (l) => console.error(l), ): Promise { const what = `${event.type} (${event.id || "no id"})`; try { await write(event); } catch (err) { let verdict: Verdict; try { verdict = await spool.failed(event, err); } catch (spoolErr) { say(`could not write ${what}: ${err}; and could not spool it either: ${spoolErr}; asking for it again`); throw err; } if (verdict === "held") { say(`could not write ${what} on its last delivery: ${err}; spooled and taken, written when writing works again`); return; } say( verdict === "over-bound" ? `could not write ${what} on its last delivery: ${err}; spooled, and NOT taken because the spool is over its bound (${spool.standing().why}) — the bus gives it up and the controller raises max-deliveries; the spool still writes it when it can` : `could not write ${what}: ${err}; spooled, asking for it again`, ); throw err; } // Written: a copy spooled by an earlier failed delivery is no longer needed. Failing to remove it // costs one idempotent write on the next replay, so it is said, never thrown. await spool.done(keyOf(event)).catch((e) => say(`wrote ${what} but could not remove its spooled copy: ${e}`)); } /** Replay the spool every `everyMs` while it holds anything, and say loudly — every fifteen minutes — * while it is over its bound. Returns a stop. */ export function replayEvery( spool: Spool, write: (event: Event) => Promise, everyMs = 30_000, say: (line: string) => void = (l) => console.error(l), ): () => void { let running = false; let loudAt = 0; const timer = setInterval(async () => { if (running || spool.size === 0) return; running = true; try { const r = await spool.replay(write); if (r.replayed > 0) say(`replayed ${r.replayed} spooled event(s); ${r.left} left`); const st = spool.standing(); if (st.overBound && Date.now() - loudAt > 15 * 60_000) { loudAt = Date.now(); say(`THE SPOOL IS OVER ITS BOUND: ${st.why}; ${st.spooled} event(s) wait in ${st.dir}, last error: ${r.error ?? st.lastError}`); } } catch (err) { say(`replaying the spool failed: ${err}`); } finally { running = false; } }, everyMs); timer.unref?.(); return () => clearInterval(timer); }