audit-logger, model-usage: retry a failed write and never lose the event (hq issue 276)
Both caught a failed write and took the event, losing it silently; the SDK's rule is to throw when the work was not done. A failed write now throws so the bus offers the event again, and is spooled on disk at once; on its last delivery the spooled event is taken, and a background pass replays the spool once writing works. The runtime does not pass the delivery count, so the spool counts failed deliveries itself, across restarts. Over its bound (1000 events or 30 minutes) the last delivery is no longer taken, so the bus gives it up and the controller raises max-deliveries - the one existing condition that names a consumer which cannot keep up - while the spool still holds it. Writes are idempotent by event: the trail skips an id it already wrote; the usage upsert keeps the reading observed latest (migration 2), so a late replay never overwrites a newer one. Each module has a status tool for the spool, declared as valuable data (ADR 0233). model-usage moves to the bundle shape (ADR 0198) with its schema in a prepare step and numbered migrations; its old container shape had no image. Both on mesh-sdk 0.1.13. The log-only handlers of redis, mssql, mosquitto, mongodb, mesh-vault, showcase and the catalogue no longer throw a TypeError on an event without a body.
This commit is contained in:
@@ -1,20 +1,33 @@
|
||||
// The audit handler — appends one line per event to an append-only audit log. It is the whole of
|
||||
// the module's code: an audit logger is not a privileged component, only a module that listens to
|
||||
// everything and writes it down (novox/hq ADR 0041).
|
||||
//
|
||||
// **Written once per event.** Delivery is at-least-once, and since novox/hq issue 276 a failed write
|
||||
// is asked for again and an event that keeps failing is replayed from the spool — so the same event can
|
||||
// reach the trail more than once. An append-only file has no ON CONFLICT; its equivalent is the SDK's
|
||||
// other answer ("Taking an event"): skip an event id already written, remembering it only once the line
|
||||
// is on disk. The ids remembered are the most recent ones, seeded from the end of the trail on start,
|
||||
// which covers every redelivery (seconds apart) and a restart between a write and its answer.
|
||||
|
||||
import { appendFile, mkdir } from "node:fs/promises";
|
||||
import { mkdir, open, stat } from "node:fs/promises";
|
||||
import { dirname } from "node:path";
|
||||
import type { Event } from "@novox/mesh-sdk/events";
|
||||
import { keyOf } from "./spool.js";
|
||||
|
||||
/** Where the trail is written. A directory the host applies; one file per node. */
|
||||
export function auditLogPath(env: NodeJS.ProcessEnv = process.env): string {
|
||||
return env.AUDIT_LOG ?? "/var/lib/audit-logger/audit.log";
|
||||
}
|
||||
|
||||
/** Append an event to the trail as one JSON line, keeping the metadata an audit needs first. The id
|
||||
* is the event's own x-event-id (ADR 0042) — the handle a reader dedups the at-least-once trail on. */
|
||||
export async function record(event: Event, path: string): Promise<void> {
|
||||
const line =
|
||||
/** Where an event that could not be written waits (see spool.ts). Beside the trail unless said. */
|
||||
export function auditSpoolPath(env: NodeJS.ProcessEnv = process.env): string {
|
||||
return env.AUDIT_SPOOL ?? `${dirname(auditLogPath(env))}/spool`;
|
||||
}
|
||||
|
||||
/** One event as the trail's line, keeping the metadata an audit needs first. The id is the event's own
|
||||
* x-event-id (ADR 0042) — the handle a reader dedups the at-least-once trail on. */
|
||||
export function lineOf(event: Event): string {
|
||||
return (
|
||||
JSON.stringify({
|
||||
id: event.id,
|
||||
type: event.type,
|
||||
@@ -23,8 +36,82 @@ export async function record(event: Event, path: string): Promise<void> {
|
||||
at: event.at,
|
||||
...(event.causationId ? { causationId: event.causationId } : {}),
|
||||
...(event.schema ? { schema: event.schema } : {}),
|
||||
body: event.body,
|
||||
}) + "\n";
|
||||
await mkdir(dirname(path), { recursive: true }).catch(() => {});
|
||||
await appendFile(path, line, { mode: 0o600 });
|
||||
body: event.body ?? null,
|
||||
}) + "\n"
|
||||
);
|
||||
}
|
||||
|
||||
/** How many recent event ids a trail remembers, and how much of its end it reads to seed them. */
|
||||
const REMEMBERED = 50_000;
|
||||
const SEED_BYTES = 8 * 1024 * 1024;
|
||||
|
||||
export class Trail {
|
||||
private readonly seen = new Map<string, true>();
|
||||
|
||||
private constructor(readonly path: string) {}
|
||||
|
||||
/** Open the trail at `path`, remembering the ids at its end. */
|
||||
static async open(path: string): Promise<Trail> {
|
||||
const t = new Trail(path);
|
||||
await t.seed();
|
||||
return t;
|
||||
}
|
||||
|
||||
has(event: Event): boolean {
|
||||
return this.seen.has(keyOf(event));
|
||||
}
|
||||
|
||||
/** Append the event as one line and sync it, unless this trail already holds it. Throws when the
|
||||
* line could not be written — the handler's cue to ask for the event again. */
|
||||
async record(event: Event): Promise<void> {
|
||||
const key = keyOf(event);
|
||||
if (this.seen.has(key)) return;
|
||||
await mkdir(dirname(this.path), { recursive: true }).catch(() => {});
|
||||
const fh = await open(this.path, "a", 0o600);
|
||||
try {
|
||||
await fh.appendFile(lineOf(event));
|
||||
await fh.datasync();
|
||||
} finally {
|
||||
await fh.close();
|
||||
}
|
||||
this.remember(key);
|
||||
}
|
||||
|
||||
private remember(key: string): void {
|
||||
this.seen.set(key, true);
|
||||
if (this.seen.size > REMEMBERED) {
|
||||
const first = this.seen.keys().next().value;
|
||||
if (first !== undefined) this.seen.delete(first);
|
||||
}
|
||||
}
|
||||
|
||||
private async seed(): Promise<void> {
|
||||
let size: number;
|
||||
try {
|
||||
size = (await stat(this.path)).size;
|
||||
} catch {
|
||||
return; // no trail yet
|
||||
}
|
||||
const from = Math.max(0, size - SEED_BYTES);
|
||||
const fh = await open(this.path, "r");
|
||||
let text: string;
|
||||
try {
|
||||
const buf = Buffer.alloc(size - from);
|
||||
await fh.read(buf, 0, buf.length, from);
|
||||
text = buf.toString("utf8");
|
||||
} finally {
|
||||
await fh.close();
|
||||
}
|
||||
const lines = text.split("\n");
|
||||
if (from > 0) lines.shift(); // the first is cut
|
||||
for (const line of lines) {
|
||||
if (!line) continue;
|
||||
try {
|
||||
const l = JSON.parse(line);
|
||||
this.remember(keyOf({ id: l.id ?? "", type: l.type, source: l.source, node: l.node, at: l.at, body: l.body }));
|
||||
} catch {
|
||||
// a line cut short by a failed write: its event was asked for again
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,20 +1,26 @@
|
||||
// The audit-logger's entrypoint. The per-node module runtime imports this once the broker is
|
||||
// bound; the on("#") subscription is the whole handshake — it consumes every event on the mesh
|
||||
// (module.*, mesh.*, node.*) and writes each to the trail.
|
||||
//
|
||||
// **No event is lost** (novox/hq issue 276). A line that cannot be written is thrown, so the bus offers
|
||||
// the event again; the first failure spools it on disk, its last delivery is taken from the spool, and
|
||||
// the spool is replayed into the trail once writing works again (spool.ts). The trail skips an event id
|
||||
// it already holds, so a redelivery or a replay never writes it twice (audit.ts).
|
||||
|
||||
import { on } from "@novox/mesh-sdk/events";
|
||||
import { auditLogPath, record } from "./audit.js";
|
||||
import { auditLogPath, auditSpoolPath, Trail } from "./audit.js";
|
||||
import { replayEvery, Spool, takeOrSpool } from "./spool.js";
|
||||
|
||||
const path = auditLogPath();
|
||||
const say = (line: string) => console.error(`[audit-logger] ${line}`);
|
||||
const trail = await Trail.open(auditLogPath());
|
||||
const spool = await Spool.open(auditSpoolPath());
|
||||
const write = (event: Parameters<Trail["record"]>[0]) => trail.record(event);
|
||||
|
||||
await on("#", async (event) => {
|
||||
try {
|
||||
await record(event, path);
|
||||
} catch (err) {
|
||||
// A failure to write the audit trail is worth a loud line, never a swallowed one — but it must
|
||||
// not throw back into the broker and wedge the subscription.
|
||||
console.error(`[audit-logger] could not record ${event.type}: ${err}`);
|
||||
}
|
||||
});
|
||||
replayEvery(spool, write, 30_000, say);
|
||||
|
||||
console.log(`[audit-logger] recording all mesh events to ${path}`);
|
||||
await on("#", (event) => takeOrSpool(event, write, spool, say));
|
||||
|
||||
console.log(
|
||||
`[audit-logger] recording all mesh events to ${trail.path}` +
|
||||
(spool.size > 0 ? `; ${spool.size} spooled event(s) wait in ${spool.dir}` : ""),
|
||||
);
|
||||
|
||||
@@ -12,13 +12,16 @@
|
||||
"kind": "bundle",
|
||||
"language": "typescript",
|
||||
"entrypoints": [
|
||||
"index.js"
|
||||
"index.js",
|
||||
"tools/index.js"
|
||||
],
|
||||
"loads": [
|
||||
"index.js"
|
||||
"index.js",
|
||||
"tools/index.js"
|
||||
],
|
||||
"env": {
|
||||
"AUDIT_LOG": "${dir:trail}/audit.log"
|
||||
"AUDIT_LOG": "${dir:trail}/audit.log",
|
||||
"AUDIT_SPOOL": "${dir:spool}"
|
||||
}
|
||||
}
|
||||
]
|
||||
@@ -30,6 +33,12 @@
|
||||
"path": "${dir:trail}",
|
||||
"class": "valuable",
|
||||
"why": "the audit trail, written nowhere else"
|
||||
},
|
||||
{
|
||||
"id": "spool",
|
||||
"path": "${dir:spool}",
|
||||
"class": "valuable",
|
||||
"why": "events the trail could not take yet; after the bus gives one up, the only copy"
|
||||
}
|
||||
]
|
||||
},
|
||||
@@ -44,6 +53,11 @@
|
||||
"id": "trail",
|
||||
"type": "directory",
|
||||
"mode": "0700"
|
||||
},
|
||||
{
|
||||
"id": "spool",
|
||||
"type": "directory",
|
||||
"mode": "0700"
|
||||
}
|
||||
],
|
||||
"capabilities": [
|
||||
|
||||
@@ -4,8 +4,12 @@
|
||||
"description": "audit-logger — records every event on the mesh (module, mesh and node) to an append-only trail.",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"scripts": {
|
||||
"build": "tsc audit.ts spool.ts index.ts tools/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --strict --skipLibCheck --outDir dist",
|
||||
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
|
||||
},
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0"
|
||||
"@novox/mesh-sdk": "^0.1.13"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
|
||||
@@ -0,0 +1,336 @@
|
||||
// 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<Event, "id" | "type" | "source" | "node" | "at" | "body">): 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<string, SpoolEntry>();
|
||||
private chain: Promise<unknown> = 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<Spool> {
|
||||
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<Spool> {
|
||||
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<Verdict> {
|
||||
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<void> {
|
||||
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<void>): 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<SpoolEntry | undefined>((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<void> {
|
||||
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<void> {
|
||||
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<T>(fn: () => Promise<T>): Promise<T> {
|
||||
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<void>,
|
||||
spool: Spool,
|
||||
say: (line: string) => void = (l) => console.error(l),
|
||||
): Promise<void> {
|
||||
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<void>,
|
||||
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);
|
||||
}
|
||||
@@ -1,40 +1,188 @@
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { mkdtemp, readFile } from "node:fs/promises";
|
||||
import { mkdtemp, readFile, readdir } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
|
||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||
import { emit, on } from "@novox/mesh-sdk/events";
|
||||
import { record } from "../audit.ts";
|
||||
import { emit, on, type Event } from "@novox/mesh-sdk/events";
|
||||
import { Trail } from "../dist/audit.js";
|
||||
import { Spool, takeOrSpool } from "../dist/spool.js";
|
||||
import { getAuditTools } from "../dist/tools/index.js";
|
||||
|
||||
async function lines(path: string): Promise<Record<string, any>[]> {
|
||||
const text = await readFile(path, "utf8").catch(() => "");
|
||||
return text.trim() === "" ? [] : text.trim().split("\n").map((l) => JSON.parse(l));
|
||||
}
|
||||
|
||||
function ev(id: string, body: unknown = { n: 1 }): Event {
|
||||
return { type: "umami.site.created", id, source: "umami", node: "anchor", at: "2026-10-06T10:00:00.000Z", body };
|
||||
}
|
||||
|
||||
async function setup() {
|
||||
const dir = await mkdtemp(join(tmpdir(), "audit-"));
|
||||
const trail = await Trail.open(join(dir, "audit.log"));
|
||||
const spool = await Spool.open(join(dir, "spool"));
|
||||
const store = { broken: false };
|
||||
const write = async (e: Event) => {
|
||||
if (store.broken) throw new Error("ENOSPC: no space left on device");
|
||||
await trail.record(e);
|
||||
};
|
||||
const said: string[] = [];
|
||||
const handle = (e: Event) => takeOrSpool(e, write, spool, (l) => said.push(l));
|
||||
return { dir, trail, spool, store, write, handle, said, path: join(dir, "audit.log") };
|
||||
}
|
||||
|
||||
test("audit-logger records every event to the trail as one line each", async () => {
|
||||
const broker = memBroker();
|
||||
useBroker(() => broker);
|
||||
const dir = await mkdtemp(join(tmpdir(), "audit-"));
|
||||
const path = join(dir, "audit.log");
|
||||
const { handle, path } = await setup();
|
||||
|
||||
// The audit-logger's whole behaviour: consume everything, record it.
|
||||
await on("#", async (event) => record(event, path)); // the pattern index.ts subscribes
|
||||
await on("#", handle); // the pattern index.ts subscribes
|
||||
|
||||
process.env.MESH_MODULE = "umami";
|
||||
process.env.MESH_NODE = "anchor";
|
||||
await emit("site.created", { domain: "my-app" });
|
||||
await emit("node.anchor.joined", { role: "worker" }); // a node event, not a module one
|
||||
|
||||
const lines = (await readFile(path, "utf8")).trim().split("\n").map((l) => JSON.parse(l));
|
||||
assert.equal(lines.length, 2);
|
||||
const got = await lines(path);
|
||||
assert.equal(got.length, 2);
|
||||
// A module names its events locally (design 29); the module is the `source`, which together with
|
||||
// the type says whose event it was. This broker does no namespacing, so the type is as emitted.
|
||||
assert.deepEqual(lines.map((l) => l.type), ["site.created", "node.anchor.joined"]);
|
||||
assert.equal(lines[0].source, "umami");
|
||||
assert.equal(lines[0].node, "anchor");
|
||||
assert.equal(lines[0].body.domain, "my-app");
|
||||
assert.deepEqual(got.map((l) => l.type), ["site.created", "node.anchor.joined"]);
|
||||
assert.equal(got[0].source, "umami");
|
||||
assert.equal(got[0].node, "anchor");
|
||||
assert.equal(got[0].body.domain, "my-app");
|
||||
|
||||
delete process.env.MESH_MODULE;
|
||||
delete process.env.MESH_NODE;
|
||||
});
|
||||
|
||||
test("a write that fails is thrown, so the bus offers the event again — and it is spooled at once", async () => {
|
||||
const broker = memBroker();
|
||||
useBroker(() => broker);
|
||||
const { handle, store, spool, path } = await setup();
|
||||
await on("#", handle);
|
||||
store.broken = true;
|
||||
await assert.rejects(emit("site.created", { domain: "x" }), /ENOSPC/);
|
||||
assert.equal(spool.size, 1, "spooled on the first failure");
|
||||
assert.equal(spool.list()[0].attempts, 1);
|
||||
assert.equal(spool.list()[0].held, false);
|
||||
assert.deepEqual(await lines(path), []);
|
||||
});
|
||||
|
||||
test("a redelivered event is written once — also after a restart, and also when its earlier failure was spooled", async () => {
|
||||
const { handle, store, spool, path, dir } = await setup();
|
||||
await handle(ev("e1"));
|
||||
await handle(ev("e1")); // the answer was lost; the bus offers it again
|
||||
assert.equal((await lines(path)).length, 1);
|
||||
|
||||
// A restart between writing and answering: the trail remembers the ids at its end.
|
||||
const again = await Trail.open(path);
|
||||
await again.record(ev("e1"));
|
||||
assert.equal((await lines(path)).length, 1);
|
||||
|
||||
// Failed once, then written: the spooled copy goes.
|
||||
store.broken = true;
|
||||
await assert.rejects(handle(ev("e2")));
|
||||
assert.equal(spool.size, 1);
|
||||
store.broken = false;
|
||||
await handle(ev("e2"));
|
||||
assert.equal(spool.size, 0);
|
||||
assert.deepEqual(await readdir(join(dir, "spool")), []);
|
||||
assert.deepEqual((await lines(path)).map((l) => l.id), ["e1", "e2"]);
|
||||
});
|
||||
|
||||
test("on its last delivery a failed event is spooled and taken", async () => {
|
||||
const { handle, store, spool, said } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 1; i < 5; i++) await assert.rejects(handle(ev("e3")), /ENOSPC/, `delivery ${i} asks again`);
|
||||
await handle(ev("e3")); // the fifth — the controller's max-deliver — returns: taken
|
||||
const [entry] = spool.list();
|
||||
assert.equal(entry.attempts, 5);
|
||||
assert.equal(entry.held, true);
|
||||
assert.match(said.at(-1)!, /spooled and taken/);
|
||||
});
|
||||
|
||||
test("the delivery count survives a restart, because it is the spool's", async () => {
|
||||
const { handle, store, dir } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 0; i < 3; i++) await assert.rejects(handle(ev("e4")));
|
||||
const reopened = await Spool.open(join(dir, "spool"));
|
||||
assert.equal(reopened.list()[0].attempts, 3);
|
||||
});
|
||||
|
||||
test("the spool replays into the trail once writing works, once each", async () => {
|
||||
const { handle, store, spool, write, path } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 0; i < 5; i++) await handle(ev("e5")).catch(() => {});
|
||||
for (let i = 0; i < 5; i++) await handle(ev("e6")).catch(() => {});
|
||||
assert.equal(spool.size, 2);
|
||||
|
||||
// Still broken: the pass stops at the first failure and keeps everything.
|
||||
let r = await spool.replay(write);
|
||||
assert.equal(r.replayed, 0);
|
||||
assert.equal(spool.size, 2);
|
||||
|
||||
store.broken = false;
|
||||
r = await spool.replay(write);
|
||||
assert.deepEqual(r, { replayed: 2, left: 0 });
|
||||
await handle(ev("e5")); // a late redelivery after the replay
|
||||
assert.deepEqual((await lines(path)).map((l) => l.id), ["e5", "e6"]);
|
||||
});
|
||||
|
||||
test("over its bound the spool keeps the event but does not take it, so the bus gives it up and says so", async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), "audit-"));
|
||||
const spool = await Spool.open(join(dir, "spool"), { boundCount: 1 });
|
||||
const failing = async () => {
|
||||
throw new Error("down");
|
||||
};
|
||||
for (let i = 0; i < 5; i++) await takeOrSpool(ev("a"), failing, spool, () => {}).catch(() => {});
|
||||
assert.equal(spool.list()[0].held, true, "within the bound: taken");
|
||||
for (let i = 1; i < 5; i++) await assert.rejects(takeOrSpool(ev("b"), failing, spool, () => {}));
|
||||
await assert.rejects(takeOrSpool(ev("b"), failing, spool, () => {}), /down/, "over the bound: asked again on its last");
|
||||
assert.equal(spool.size, 2, "and still spooled");
|
||||
const st = spool.standing();
|
||||
assert.equal(st.overBound, true);
|
||||
assert.equal(st.spooled, 2);
|
||||
assert.equal(st.held, 1);
|
||||
assert.match(st.why!, /more than 1/);
|
||||
});
|
||||
|
||||
test("an event without a body, or without an id, is recorded once", async () => {
|
||||
const { handle, path } = await setup();
|
||||
await handle(ev("e7", null));
|
||||
await handle(ev("", undefined));
|
||||
await handle(ev("", undefined));
|
||||
const got = await lines(path);
|
||||
assert.equal(got.length, 2);
|
||||
assert.equal(got[0].body, null);
|
||||
});
|
||||
|
||||
test("an event waiting longer than the bound puts the spool over it", async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), "audit-"));
|
||||
let now = new Date("2026-10-06T10:00:00Z");
|
||||
const spool = await Spool.open(join(dir, "spool"), { now: () => now });
|
||||
await takeOrSpool(ev("old"), async () => { throw new Error("down"); }, spool, () => {}).catch(() => {});
|
||||
assert.equal(spool.standing().overBound, false);
|
||||
now = new Date("2026-10-06T10:31:00Z");
|
||||
assert.match(spool.standing().why!, /longer than 30 minutes/);
|
||||
});
|
||||
|
||||
test("audit_status says what waits in the spool, read from disk by the tool's own process", async () => {
|
||||
const { handle, store, dir } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 0; i < 5; i++) await handle(ev("s1")).catch(() => {});
|
||||
await handle(ev("s2")).catch(() => {});
|
||||
const [tool] = getAuditTools({ AUDIT_LOG: join(dir, "audit.log"), AUDIT_SPOOL: join(dir, "spool") });
|
||||
assert.equal(tool.name, "audit_status");
|
||||
const got = (await tool.run({}, {} as never)) as { spool: { spooled: number; held: number; lastError: string } };
|
||||
assert.equal(got.spool.spooled, 2);
|
||||
assert.equal(got.spool.held, 1);
|
||||
assert.match(got.spool.lastError, /ENOSPC/);
|
||||
});
|
||||
|
||||
function memBroker() {
|
||||
const subs: { pattern: string; handler: (env: { key: string; node: string; body: unknown }) => Promise<void> }[] = [];
|
||||
return {
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
// The audit-logger's one tool — its standing: where the trail is, and what waits in the spool
|
||||
// (novox/hq issue 276). Launched as its own process beside the handler, so it reads the spool from
|
||||
// disk rather than asking the handler.
|
||||
|
||||
import { stat } from "node:fs/promises";
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { auditLogPath, auditSpoolPath } from "../audit.js";
|
||||
import { Spool } from "../spool.js";
|
||||
|
||||
export function getAuditTools(env: NodeJS.ProcessEnv = process.env): ToolDefinition[] {
|
||||
return [
|
||||
{
|
||||
name: "audit_status",
|
||||
description:
|
||||
"The audit trail's standing: its path and size, and the events that could not be written yet — how many wait in the spool, how many were taken from the bus and are held only there, the oldest, the last error, and whether the spool is over its bound (then the bus gives events up and max-deliveries is raised).",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const path = auditLogPath(env);
|
||||
const size = await stat(path).then((s) => s.size, () => null);
|
||||
const spool = await Spool.read(auditSpoolPath(env));
|
||||
return { trail: { path, bytes: size }, spool: spool.standing() };
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
registerModuleTools("audit-logger", (env) => getAuditTools(env));
|
||||
@@ -8,5 +8,10 @@
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true
|
||||
},
|
||||
"include": ["audit.ts", "index.ts"]
|
||||
"include": [
|
||||
"audit.ts",
|
||||
"spool.ts",
|
||||
"index.ts",
|
||||
"tools/index.ts"
|
||||
]
|
||||
}
|
||||
|
||||
@@ -58,7 +58,7 @@ interface Built {
|
||||
* months ago — see `replay` above.
|
||||
*/
|
||||
const placeTheBuild = async (event: { body: unknown }): Promise<void> => {
|
||||
const body = event.body as Built;
|
||||
const body = (event.body ?? {}) as Built;
|
||||
if (!body.module || !body.commit) {
|
||||
// Said rather than dropped: a build that announced itself without saying what it built is a
|
||||
// fault in the builder, and a silent discard here would make it look like a missing event.
|
||||
|
||||
@@ -17,15 +17,15 @@ interface SecretEvent {
|
||||
}
|
||||
|
||||
await on<SecretEvent>("provisioned", async (e) => {
|
||||
console.log(`[mesh-vault] secret provisioned for ${e.body.as} on ${e.body.consumer} (${e.body.fingerprint})`);
|
||||
console.log(`[mesh-vault] secret provisioned for ${e.body?.as} on ${e.body?.consumer} (${e.body?.fingerprint})`);
|
||||
});
|
||||
|
||||
await on<SecretEvent>("rotated", async (e) => {
|
||||
console.log(`[mesh-vault] secret rotated for ${e.body.as} — rotation ${e.body.rotations} (${e.body.fingerprint})`);
|
||||
console.log(`[mesh-vault] secret rotated for ${e.body?.as} — rotation ${e.body?.rotations} (${e.body?.fingerprint})`);
|
||||
});
|
||||
|
||||
await on<SecretEvent>("deprovisioned", async (e) => {
|
||||
console.log(`[mesh-vault] secret withdrawn from ${e.body.as}`);
|
||||
console.log(`[mesh-vault] secret withdrawn from ${e.body?.as}`);
|
||||
});
|
||||
|
||||
console.log("[mesh-vault] auditing secret lifecycle events");
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
// What model-usage does with one usage event: read its rows, and write each to the store. Kept apart
|
||||
// from index.ts so it is tested without a broker (novox/hq issue 276).
|
||||
//
|
||||
// The producers emit ALREADY-NORMALISED rows (ADR 0054). A body is `{ rows: UsageRow[], raw }`; each
|
||||
// row carries its own `raw` or the body's as a fallback. A row the store could never take — a missing
|
||||
// key, a value that is not a number — is said and skipped: asking for the event again would not change
|
||||
// it. A store that did not take a row throws, so the event is asked for again (spool.ts).
|
||||
|
||||
import type { Event } from "@novox/mesh-sdk/events";
|
||||
import type { Observed, UsageRow } from "./store.js";
|
||||
|
||||
/** Where a usage event that could not be written waits (spool.ts). */
|
||||
export function usageSpoolPath(env: NodeJS.ProcessEnv = process.env): string {
|
||||
return env.USAGE_SPOOL ?? "/var/lib/model-usage/spool";
|
||||
}
|
||||
|
||||
/** The slice of the store the consumer writes through. */
|
||||
export interface UsageWriter {
|
||||
upsert(row: UsageRow, observed?: Observed): Promise<void>;
|
||||
}
|
||||
|
||||
/** The rows of a usage event that can be stored, and the words for those that cannot. */
|
||||
export function rowsOf(event: Pick<Event, "body">): { rows: UsageRow[]; skipped: string[] } {
|
||||
const body = (event.body ?? {}) as { rows?: unknown; raw?: unknown };
|
||||
const rows: UsageRow[] = [];
|
||||
const skipped: string[] = [];
|
||||
if (!Array.isArray(body.rows)) return { rows, skipped: body.rows === undefined ? [] : ["rows is not a list"] };
|
||||
body.rows.forEach((r, i) => {
|
||||
const row = (r ?? {}) as Partial<UsageRow>;
|
||||
const missing = (["licence", "consumer", "period", "metric"] as const).filter(
|
||||
(k) => typeof row[k] !== "string" || row[k] === "",
|
||||
);
|
||||
if (missing.length > 0) return void skipped.push(`row ${i} has no ${missing.join(", ")}`);
|
||||
if (typeof row.value !== "number" || !Number.isFinite(row.value)) return void skipped.push(`row ${i}'s value is not a number`);
|
||||
rows.push({ ...(row as UsageRow), raw: row.raw ?? body.raw ?? {} });
|
||||
});
|
||||
return { rows, skipped };
|
||||
}
|
||||
|
||||
/** Write one event's rows. Every row is upserted idempotently, so a retry after a partial write
|
||||
* writes the rows already written to the same values. */
|
||||
export function writerFor(store: UsageWriter, say: (line: string) => void = (l) => console.error(l)) {
|
||||
return async (event: Event): Promise<void> => {
|
||||
const { rows, skipped } = rowsOf(event);
|
||||
for (const why of skipped) say(`${event.type} (${event.id || "no id"}): ${why}; that row is skipped`);
|
||||
for (const row of rows) await store.upsert(row, { at: event.at, eventId: event.id });
|
||||
};
|
||||
}
|
||||
@@ -1,38 +1,32 @@
|
||||
// model-usage's entrypoint — the usage context store's consumer (novox/hq ADR 0054). mesh-controller is
|
||||
// a CLI and cannot consume events, so the store that keeps the latest usage reading is a MODULE: it
|
||||
// subscribes to `module.*.usage.*` and upserts each row. Like the audit-logger, the on(...) IS the
|
||||
// whole handshake — the runtime imports this once the broker is bound, and every usage event any
|
||||
// producer emits lands here as well as on the audit trail.
|
||||
// subscribes to `*.usage.*` and upserts each row. Like the audit-logger, the on(...) IS the whole
|
||||
// handshake — the node's runtime launches this once the broker is bound (ADR 0198), and every usage
|
||||
// event any producer emits lands here as well as on the audit trail.
|
||||
//
|
||||
// The producers (the anthropic adapters) emit ALREADY-NORMALISED rows: the vendor→row normalisation
|
||||
// lives in the adapter, not here, so this consumer is vendor-neutral (ADR 0054). A body is
|
||||
// `{ rows: UsageRow[], raw }`; each row is upserted, carrying its own `raw` or the body's as a
|
||||
// fallback. Delivery is at-least-once, so a duplicate is fine — the upsert keeps the latest.
|
||||
// **No event is lost** (novox/hq issue 276). A row the store did not take is thrown, so the bus offers
|
||||
// the event again; the first failure spools it on disk, its last delivery is taken from the spool, and
|
||||
// the spool is replayed into the store once it answers again (spool.ts). The upsert keeps the reading
|
||||
// observed latest, so a redelivery or a late replay never writes an older one over a newer.
|
||||
//
|
||||
// The schema is not brought up here: the prepare step does it before this version starts
|
||||
// (prepare/index.ts; ADR 0135), as for the module catalogue.
|
||||
|
||||
import { on } from "@novox/mesh-sdk/events";
|
||||
import { UsageStore, type UsageRow } from "./store.js";
|
||||
import { UsageStore } from "./store.js";
|
||||
import { usageSpoolPath, writerFor } from "./consume.js";
|
||||
import { replayEvery, Spool, takeOrSpool } from "./spool.js";
|
||||
|
||||
const say = (line: string) => console.error(`[model-usage] ${line}`);
|
||||
const store = UsageStore.fromEnv();
|
||||
const spool = await Spool.open(usageSpoolPath());
|
||||
const write = writerFor(store, say);
|
||||
|
||||
// Create the store's one table before subscribing. The DDL is idempotent (CREATE TABLE IF NOT
|
||||
// EXISTS), so a restart re-runs it harmlessly. This is done here, in the long-lived consumer, rather
|
||||
// than as a gating run-once step: the consumer is a `--restart unless-stopped` service, so if the
|
||||
// provider is not yet reachable — its overlay address comes up as the same push settles — this exits
|
||||
// and is restarted until it can connect, without ever halting the apply. A run-once migrate that had
|
||||
// to reach the provider over the overlay would block the very apply that brings the overlay up.
|
||||
await store.migrate();
|
||||
replayEvery(spool, write, 30_000, say);
|
||||
|
||||
await on("*.usage.*", async (event) => {
|
||||
const body = event.body as { rows?: UsageRow[]; raw?: unknown };
|
||||
for (const row of body.rows ?? []) {
|
||||
try {
|
||||
await store.upsert({ ...row, raw: row.raw ?? body.raw ?? {} });
|
||||
} catch (err) {
|
||||
// A failed upsert is a loud line, never a throw back into the broker that would wedge the
|
||||
// subscription (the audit-logger's discipline).
|
||||
console.error(`[model-usage] could not upsert a row from ${event.type}: ${err}`);
|
||||
}
|
||||
}
|
||||
});
|
||||
await on("*.usage.*", (event) => takeOrSpool(event, write, spool, say));
|
||||
|
||||
console.log("[model-usage] recording model usage to its store");
|
||||
console.log(
|
||||
"[model-usage] recording model usage to its store" +
|
||||
(spool.size > 0 ? `; ${spool.size} spooled event(s) wait in ${spool.dir}` : ""),
|
||||
);
|
||||
|
||||
@@ -22,9 +22,6 @@
|
||||
"consumes": [
|
||||
"*.usage.*"
|
||||
],
|
||||
"own-secrets": {
|
||||
"broker": "${dir:mesh-state}/broker"
|
||||
},
|
||||
"data": {
|
||||
"own": [
|
||||
{
|
||||
@@ -32,16 +29,16 @@
|
||||
"path": "${dir:state}",
|
||||
"class": "cache",
|
||||
"why": "the mesh's own files for it; its records are in its database"
|
||||
},
|
||||
{
|
||||
"id": "spool",
|
||||
"path": "${dir:spool}",
|
||||
"class": "valuable",
|
||||
"why": "usage events the store could not take yet; after the bus gives one up, the only copy"
|
||||
}
|
||||
]
|
||||
},
|
||||
"resources": [
|
||||
{
|
||||
"id": "mesh-state",
|
||||
"type": "directory",
|
||||
"mode": "0700",
|
||||
"place": "mesh"
|
||||
},
|
||||
{
|
||||
"id": "state",
|
||||
"type": "directory",
|
||||
@@ -56,20 +53,48 @@
|
||||
"content": "postgresql://${bound:postgres-database:as}:${secret:postgres-database}@${bound:postgres-database:at}:${bound:postgres-database:port}/${bound:postgres-database:as}\n"
|
||||
},
|
||||
{
|
||||
"id": "runtime",
|
||||
"type": "container",
|
||||
"name": "mesh-model-usage",
|
||||
"image": "mesh-runtime-model-usage@sha256:0000000000000000000000000000000000000000000000000000000000000000",
|
||||
"network": "host",
|
||||
"volumes": [
|
||||
"${dir:mesh-state}/broker:/run/secrets/broker:ro",
|
||||
"${dir:state}:/run/state",
|
||||
"${dir:state}/database.url:/run/secrets/database-url:ro"
|
||||
"id": "spool",
|
||||
"type": "directory",
|
||||
"mode": "0700"
|
||||
},
|
||||
{
|
||||
"id": "prepare",
|
||||
"type": "process",
|
||||
"name": "model-usage-prepare",
|
||||
"artifact": "code",
|
||||
"run": [
|
||||
"node",
|
||||
"prepare/index.js"
|
||||
],
|
||||
"run-once": true,
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"DATABASE_URL_FILE": "/run/secrets/database-url"
|
||||
}
|
||||
"DATABASE_URL_FILE": "${dir:state}/database.url"
|
||||
},
|
||||
"restart-on": [
|
||||
"database-url"
|
||||
]
|
||||
}
|
||||
]
|
||||
],
|
||||
"build": {
|
||||
"artifacts": [
|
||||
{
|
||||
"name": "code",
|
||||
"kind": "bundle",
|
||||
"language": "typescript",
|
||||
"entrypoints": [
|
||||
"index.js",
|
||||
"tools/index.js",
|
||||
"prepare/index.js"
|
||||
],
|
||||
"loads": [
|
||||
"index.js",
|
||||
"tools/index.js"
|
||||
],
|
||||
"env": {
|
||||
"DATABASE_URL_FILE": "${dir:state}/database.url",
|
||||
"USAGE_SPOOL": "${dir:spool}"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,8 +4,12 @@
|
||||
"description": "model-usage — the usage context store (novox/hq ADR 0054): consumes module.*.usage.* and upserts the latest vendor-neutral reading per (licence, consumer, period, metric) into a provisioned postgres store, kept in the clear.",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"scripts": {
|
||||
"build": "tsc pg.d.ts store.ts spool.ts consume.ts index.ts tools/index.ts prepare/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --strict --skipLibCheck --outDir dist",
|
||||
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
|
||||
},
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0",
|
||||
"@novox/mesh-sdk": "^0.1.13",
|
||||
"pg": "^8"
|
||||
},
|
||||
"devDependencies": {
|
||||
|
||||
Vendored
+9
-3
@@ -5,10 +5,16 @@
|
||||
// runtime image (package.json `dependencies`; novox/hq ADR 0052), so this types the code without
|
||||
// deciding what runs.
|
||||
declare module "pg" {
|
||||
/** A lazily-connecting connection pool. Only the connection string, one query form, and end() are
|
||||
* used here. */
|
||||
/** A lazily-connecting connection pool. Only the connection string and two timeouts, one query form,
|
||||
* and end() are used here. */
|
||||
export class Pool {
|
||||
constructor(config?: { connectionString?: string });
|
||||
constructor(config?: {
|
||||
connectionString?: string;
|
||||
/** How long a connection may take before the query fails (ms). */
|
||||
connectionTimeoutMillis?: number;
|
||||
/** How long a query may take before it fails, client side (ms). */
|
||||
query_timeout?: number;
|
||||
});
|
||||
query(text: string, params?: unknown[]): Promise<{ rows: any[]; rowCount: number }>;
|
||||
end(): Promise<void>;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
// The usage store's schema, brought to the shape this version needs (novox/hq ADR 0135): the host runs
|
||||
// this as a run-once step before the version that needs it, with no bus, and does not start that
|
||||
// version if it fails. Migrations are numbered and recorded (store.ts).
|
||||
import { UsageStore } from "../store.js";
|
||||
|
||||
const store = UsageStore.fromEnv();
|
||||
const applied = await store.migrate();
|
||||
console.log(
|
||||
applied.length > 0
|
||||
? `[model-usage] applied migration(s) ${applied.join(", ")}; the usage store's schema is what this version needs`
|
||||
: "[model-usage] the usage store's schema is what this version needs",
|
||||
);
|
||||
await store.close();
|
||||
@@ -0,0 +1,336 @@
|
||||
// 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<Event, "id" | "type" | "source" | "node" | "at" | "body">): 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<string, SpoolEntry>();
|
||||
private chain: Promise<unknown> = 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<Spool> {
|
||||
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<Spool> {
|
||||
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<Verdict> {
|
||||
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<void> {
|
||||
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<void>): 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<SpoolEntry | undefined>((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<void> {
|
||||
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<void> {
|
||||
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<T>(fn: () => Promise<T>): Promise<T> {
|
||||
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<void>,
|
||||
spool: Spool,
|
||||
say: (line: string) => void = (l) => console.error(l),
|
||||
): Promise<void> {
|
||||
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<void>,
|
||||
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);
|
||||
}
|
||||
@@ -3,8 +3,9 @@ import { readFileSync } from "node:fs";
|
||||
// session — which differ only in `consumer`; a reading is one row `(licence, consumer, period,
|
||||
// metric, value)` plus its `raw` vendor payload. The store keeps the LATEST reading per
|
||||
// `(licence, consumer, period, metric)`: an event carries a fresh total, and the upsert replaces the
|
||||
// previous one. Delivery is at-least-once, so an upsert is idempotent-latest by design — a duplicate
|
||||
// event writes the same row again, never a second.
|
||||
// previous one — unless the previous one was observed later. Delivery is at-least-once, so an upsert
|
||||
// is idempotent-latest by design: a duplicate event writes the same row again, never a second, and an
|
||||
// older event replayed late changes nothing.
|
||||
//
|
||||
// Usage is stored IN THE CLEAR (ADR 0054): the value and its raw payload are ordinary columns, not
|
||||
// sealed. This is not a credential; it is a reading a query answers.
|
||||
@@ -29,7 +30,20 @@ export interface UsageRow {
|
||||
raw?: unknown;
|
||||
}
|
||||
|
||||
const DDL = `
|
||||
/**
|
||||
* The store's schema, as numbered migrations applied in order and recorded in `usage_schema`. Each is
|
||||
* written to be safe on a store that already has it, because the first version created the table
|
||||
* without recording anything.
|
||||
*
|
||||
* 2 (novox/hq issue 276): a reading now carries the time its event was emitted, and an upsert replaces
|
||||
* a reading only with one emitted no earlier. A failed event is asked for again and, past its last
|
||||
* delivery, replayed from the spool — possibly after a newer reading arrived; it must not overwrite it.
|
||||
*/
|
||||
export const MIGRATIONS: { version: number; why: string; sql: string }[] = [
|
||||
{
|
||||
version: 1,
|
||||
why: "the one table, both grains",
|
||||
sql: `
|
||||
CREATE TABLE IF NOT EXISTS usage (
|
||||
licence text NOT NULL,
|
||||
consumer text NOT NULL,
|
||||
@@ -39,8 +53,32 @@ CREATE TABLE IF NOT EXISTS usage (
|
||||
raw jsonb NOT NULL DEFAULT '{}'::jsonb,
|
||||
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||
PRIMARY KEY (licence, consumer, period, metric)
|
||||
);
|
||||
`;
|
||||
);`,
|
||||
},
|
||||
{
|
||||
version: 2,
|
||||
why: "a reading knows when it was observed and which event said it, so an older one never replaces it",
|
||||
sql: `
|
||||
ALTER TABLE usage ADD COLUMN IF NOT EXISTS observed_at timestamptz;
|
||||
ALTER TABLE usage ADD COLUMN IF NOT EXISTS event_id text;`,
|
||||
},
|
||||
];
|
||||
|
||||
/** The upsert: the latest reading per key, where latest is the event's emit time — so a redelivered or
|
||||
* replayed event writes the row it already wrote, and an older one changes nothing. */
|
||||
export const UPSERT = `
|
||||
INSERT INTO usage (licence, consumer, period, metric, value, raw, observed_at, event_id)
|
||||
VALUES ($1, $2, $3, $4, $5, $6::jsonb, COALESCE($7::timestamptz, now()), $8)
|
||||
ON CONFLICT (licence, consumer, period, metric)
|
||||
DO UPDATE SET value = excluded.value, raw = excluded.raw, observed_at = excluded.observed_at,
|
||||
event_id = excluded.event_id, updated_at = now()
|
||||
WHERE usage.observed_at IS NULL OR usage.observed_at <= excluded.observed_at`;
|
||||
|
||||
/** Who wrote a reading, and when it was true — the event it came from. */
|
||||
export interface Observed {
|
||||
at?: string;
|
||||
eventId?: string;
|
||||
}
|
||||
|
||||
export class UsageStore {
|
||||
private constructor(private readonly pool: PgPool) {}
|
||||
@@ -51,24 +89,42 @@ export class UsageStore {
|
||||
// As a file first (novox/hq ADR 0086): the connection string carries the password.
|
||||
const url = env["DATABASE_URL"] ?? readMaybe(env["DATABASE_URL_FILE"]);
|
||||
if (!url) throw new Error("DATABASE_URL_FILE (or DATABASE_URL) is not set — model-usage cannot reach its database");
|
||||
return new UsageStore(new Pool({ connectionString: url }));
|
||||
// Bounded, so a store that does not answer is a failed write the handler throws well within the
|
||||
// runtime's two-minute event timeout — not an event left hanging until the runtime gives up on it.
|
||||
return new UsageStore(new Pool({ connectionString: url, connectionTimeoutMillis: 10_000, query_timeout: 30_000 }));
|
||||
}
|
||||
|
||||
/** Create the one table if it is not there. Run once by the migrate entry before the consumer
|
||||
* starts; idempotent, so re-running is harmless. */
|
||||
async migrate(): Promise<void> {
|
||||
await this.pool.query(DDL);
|
||||
}
|
||||
|
||||
/** Upsert a reading, keeping the latest per `(licence, consumer, period, metric)`. */
|
||||
async upsert(row: UsageRow): Promise<void> {
|
||||
/** Bring the schema to the newest migration, recording each applied. Idempotent: run by the
|
||||
* prepare step before this version starts (novox/hq ADR 0135). */
|
||||
async migrate(): Promise<number[]> {
|
||||
await this.pool.query(
|
||||
`INSERT INTO usage (licence, consumer, period, metric, value, raw)
|
||||
VALUES ($1, $2, $3, $4, $5, $6::jsonb)
|
||||
ON CONFLICT (licence, consumer, period, metric)
|
||||
DO UPDATE SET value = excluded.value, raw = excluded.raw, updated_at = now()`,
|
||||
[row.licence, row.consumer, row.period, row.metric, row.value, JSON.stringify(row.raw ?? {})],
|
||||
"CREATE TABLE IF NOT EXISTS usage_schema (version integer PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())",
|
||||
);
|
||||
const { rows } = await this.pool.query("SELECT version FROM usage_schema");
|
||||
const have = new Set(rows.map((r) => Number(r.version)));
|
||||
const applied: number[] = [];
|
||||
for (const m of MIGRATIONS) {
|
||||
if (have.has(m.version)) continue;
|
||||
await this.pool.query(m.sql);
|
||||
await this.pool.query("INSERT INTO usage_schema (version) VALUES ($1) ON CONFLICT DO NOTHING", [m.version]);
|
||||
applied.push(m.version);
|
||||
}
|
||||
return applied;
|
||||
}
|
||||
|
||||
/** Upsert a reading, keeping the latest per `(licence, consumer, period, metric)` by when it was
|
||||
* observed. Throws when the store did not take it. */
|
||||
async upsert(row: UsageRow, observed: Observed = {}): Promise<void> {
|
||||
await this.pool.query(UPSERT, [
|
||||
row.licence,
|
||||
row.consumer,
|
||||
row.period,
|
||||
row.metric,
|
||||
row.value,
|
||||
JSON.stringify(row.raw ?? {}),
|
||||
observed.at || null,
|
||||
observed.eventId || null,
|
||||
]);
|
||||
}
|
||||
|
||||
/** The current reading per key, whole or filtered to one licence — for the read-only tool. */
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
// What holds model-usage to taking no usage event it did not write (novox/hq issue 276): a store that
|
||||
// refuses is thrown, so the bus offers the event again; a redelivery writes nothing twice; the last
|
||||
// delivery is spooled and taken; the spool replays once the store answers; an older event replayed
|
||||
// late does not overwrite a newer reading; and an event without a body is taken, not a TypeError.
|
||||
//
|
||||
// The store is a fake keyed exactly as the table is, with the same rule for "latest". When
|
||||
// TEST_DATABASE_URL names a throwaway postgres, the same rule is also checked against the real SQL.
|
||||
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { mkdtemp } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
|
||||
import type { Event } from "@novox/mesh-sdk/events";
|
||||
import { rowsOf, writerFor } from "../dist/consume.js";
|
||||
import { Spool, takeOrSpool } from "../dist/spool.js";
|
||||
import { UsageStore } from "../dist/store.js";
|
||||
import { getUsageTools } from "../dist/tools/index.js";
|
||||
|
||||
type Row = { licence: string; consumer: string; period: string; metric: string; value: number; raw?: unknown };
|
||||
|
||||
function fakeStore() {
|
||||
const rows = new Map<string, { row: Row; at: string; eventId?: string }>();
|
||||
let writes = 0;
|
||||
return {
|
||||
broken: false,
|
||||
rows,
|
||||
get writes() {
|
||||
return writes;
|
||||
},
|
||||
async upsert(row: Row, observed: { at?: string; eventId?: string } = {}) {
|
||||
if (this.broken) throw new Error("connect ECONNREFUSED 10.0.0.5:5432");
|
||||
writes++;
|
||||
const key = [row.licence, row.consumer, row.period, row.metric].join("|");
|
||||
const had = rows.get(key);
|
||||
const at = observed.at || new Date().toISOString();
|
||||
if (had && had.at > at) return; // the table's WHERE: an older reading changes nothing
|
||||
rows.set(key, { row, at, eventId: observed.eventId });
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function usage(id: string, at: string, value: number, body?: unknown): Event {
|
||||
return {
|
||||
type: "anthropic-licence-manager.usage.reported",
|
||||
id,
|
||||
source: "anthropic-licence-manager",
|
||||
node: "laptop",
|
||||
at,
|
||||
body: body === undefined ? { rows: [{ licence: "team", consumer: "*", period: "2026-10", metric: "tokens", value }], raw: { v: 1 } } : body,
|
||||
};
|
||||
}
|
||||
|
||||
async function setup() {
|
||||
const store = fakeStore();
|
||||
const dir = await mkdtemp(join(tmpdir(), "usage-"));
|
||||
const spool = await Spool.open(join(dir, "spool"));
|
||||
const write = writerFor(store, () => {});
|
||||
const handle = (e: Event) => takeOrSpool(e, write, spool, () => {});
|
||||
return { store, spool, write, handle, dir };
|
||||
}
|
||||
|
||||
test("a store that refuses a row is thrown, so the event is offered again, and it is spooled", async () => {
|
||||
const { store, spool, handle } = await setup();
|
||||
store.broken = true;
|
||||
await assert.rejects(handle(usage("u1", "2026-10-06T10:00:00Z", 5)), /ECONNREFUSED/);
|
||||
assert.equal(spool.size, 1);
|
||||
assert.equal(spool.list()[0].held, false);
|
||||
});
|
||||
|
||||
test("a redelivered event writes the same reading, never a second, and clears its spooled copy", async () => {
|
||||
const { store, spool, handle } = await setup();
|
||||
store.broken = true;
|
||||
await assert.rejects(handle(usage("u2", "2026-10-06T10:00:00Z", 5)));
|
||||
store.broken = false;
|
||||
await handle(usage("u2", "2026-10-06T10:00:00Z", 5));
|
||||
await handle(usage("u2", "2026-10-06T10:00:00Z", 5));
|
||||
assert.equal(store.rows.size, 1);
|
||||
assert.equal([...store.rows.values()][0].row.value, 5);
|
||||
assert.equal(spool.size, 0);
|
||||
});
|
||||
|
||||
test("on its last delivery a usage event is spooled and taken", async () => {
|
||||
const { store, spool, handle } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 1; i < 5; i++) await assert.rejects(handle(usage("u3", "2026-10-06T10:00:00Z", 5)));
|
||||
await handle(usage("u3", "2026-10-06T10:00:00Z", 5));
|
||||
assert.equal(spool.list()[0].held, true);
|
||||
assert.equal(spool.list()[0].attempts, 5);
|
||||
});
|
||||
|
||||
test("the spool replays once the store answers, and an older reading replayed late does not overwrite a newer", async () => {
|
||||
const { store, spool, write, handle } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 0; i < 5; i++) await handle(usage("old", "2026-10-06T10:00:00Z", 5)).catch(() => {});
|
||||
store.broken = false;
|
||||
await handle(usage("new", "2026-10-06T11:00:00Z", 9)); // a newer total arrives first
|
||||
const r = await spool.replay(write);
|
||||
assert.deepEqual(r, { replayed: 1, left: 0 });
|
||||
assert.equal([...store.rows.values()][0].row.value, 9);
|
||||
assert.equal([...store.rows.values()][0].eventId, "new");
|
||||
});
|
||||
|
||||
test("an event without a body is taken without a TypeError; a row the store could never take is skipped, not retried", async () => {
|
||||
const { store, handle } = await setup();
|
||||
await handle(usage("n1", "2026-10-06T10:00:00Z", 0, null));
|
||||
await handle(usage("n2", "2026-10-06T10:00:00Z", 0, { rows: "nonsense" }));
|
||||
await handle(usage("n3", "2026-10-06T10:00:00Z", 0, { rows: [null, { licence: "team", consumer: "*", period: "p", metric: "m", value: "x" }] }));
|
||||
assert.equal(store.writes, 0);
|
||||
assert.deepEqual(rowsOf({ body: null }), { rows: [], skipped: [] });
|
||||
assert.equal(rowsOf({ body: { rows: [null] } }).skipped.length, 1);
|
||||
});
|
||||
|
||||
test("model_usage_status reports the spool, and is offered even when the store cannot be opened", async () => {
|
||||
const { store, handle, dir } = await setup();
|
||||
store.broken = true;
|
||||
for (let i = 0; i < 5; i++) await handle(usage("s1", "2026-10-06T10:00:00Z", 1)).catch(() => {});
|
||||
const tools = getUsageTools(undefined, { USAGE_SPOOL: join(dir, "spool") });
|
||||
assert.deepEqual(tools.map((t) => t.name), ["model_usage_status"]);
|
||||
const got = (await tools[0].run({}, {} as never)) as { spool: { spooled: number; held: number } };
|
||||
assert.equal(got.spool.spooled, 1);
|
||||
assert.equal(got.spool.held, 1);
|
||||
});
|
||||
|
||||
const url = process.env.TEST_DATABASE_URL;
|
||||
test("against postgres: migrations are numbered and repeatable, and the upsert keeps the reading observed latest", { skip: !url && "TEST_DATABASE_URL is not set" }, async () => {
|
||||
const store = UsageStore.fromEnv({ DATABASE_URL: url });
|
||||
try {
|
||||
await store.migrate();
|
||||
assert.deepEqual(await store.migrate(), [], "a second run applies nothing");
|
||||
const row = { licence: "pg", consumer: "*", period: "2026-10", metric: "tokens", value: 5, raw: {} };
|
||||
await store.upsert(row, { at: "2026-10-06T10:00:00Z", eventId: "a" });
|
||||
await store.upsert({ ...row, value: 9 }, { at: "2026-10-06T11:00:00Z", eventId: "b" });
|
||||
await store.upsert(row, { at: "2026-10-06T10:00:00Z", eventId: "a" }); // the old one, replayed late
|
||||
await store.upsert({ ...row, value: 9 }, { at: "2026-10-06T11:00:00Z", eventId: "b" }); // redelivered
|
||||
const got = (await store.current("pg")).map((r) => r.value);
|
||||
assert.deepEqual(got, [9]);
|
||||
} finally {
|
||||
await store.close();
|
||||
}
|
||||
});
|
||||
@@ -1,14 +1,27 @@
|
||||
// model-usage's tools — its own read-only query surface over the usage store (novox/hq ADR 0039:
|
||||
// the tools and the client they call live in the module). The store is opened from the environment;
|
||||
// when it cannot be (no DATABASE_URL), the module contributes no tools rather than failing the whole
|
||||
// tool runtime — the postgres precedent.
|
||||
// the tools and the client they call live in the module), and its standing: what waits in the spool
|
||||
// (novox/hq issue 276). The store is opened from the environment; when it cannot be (no DATABASE_URL),
|
||||
// the query tool is not contributed rather than failing the whole tool runtime — the postgres
|
||||
// precedent — but the status tool is, since the spool is what matters most when the store is away.
|
||||
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { UsageStore } from "../store.js";
|
||||
import { Spool } from "../spool.js";
|
||||
import { usageSpoolPath } from "../consume.js";
|
||||
|
||||
export function getUsageTools(store: UsageStore): ToolDefinition[] {
|
||||
return [
|
||||
export function getUsageTools(store: UsageStore | undefined, env: NodeJS.ProcessEnv = process.env): ToolDefinition[] {
|
||||
const tools: ToolDefinition[] = [
|
||||
{
|
||||
name: "model_usage_status",
|
||||
description:
|
||||
"The usage store's standing: the usage events that could not be written yet — how many wait in the spool, how many were taken from the bus and are held only there, the oldest, the last error, and whether the spool is over its bound (then the bus gives events up and max-deliveries is raised).",
|
||||
input: {},
|
||||
// Launched as its own process beside the handler, so it reads the spool from disk.
|
||||
run: async () => ({ spool: (await Spool.read(usageSpoolPath(env))).standing() }),
|
||||
},
|
||||
];
|
||||
if (store) {
|
||||
tools.push({
|
||||
name: "model_usage_current",
|
||||
description:
|
||||
"The latest model-usage reading per (licence, consumer, period, metric), across both grains, optionally restricted to one licence.",
|
||||
@@ -19,14 +32,17 @@ export function getUsageTools(store: UsageStore): ToolDefinition[] {
|
||||
const licence = args.licence ? String(args.licence) : undefined;
|
||||
return { usage: await store.current(licence) };
|
||||
},
|
||||
},
|
||||
];
|
||||
});
|
||||
}
|
||||
return tools;
|
||||
}
|
||||
|
||||
registerModuleTools("model-usage", (env) => {
|
||||
let store: UsageStore | undefined;
|
||||
try {
|
||||
return getUsageTools(UsageStore.fromEnv(env));
|
||||
store = UsageStore.fromEnv(env);
|
||||
} catch {
|
||||
return [];
|
||||
store = undefined;
|
||||
}
|
||||
return getUsageTools(store, env);
|
||||
});
|
||||
|
||||
@@ -11,8 +11,10 @@
|
||||
"include": [
|
||||
"pg.d.ts",
|
||||
"store.ts",
|
||||
"spool.ts",
|
||||
"consume.ts",
|
||||
"index.ts",
|
||||
"tools/index.ts",
|
||||
"migrate/index.ts"
|
||||
"prepare/index.ts"
|
||||
]
|
||||
}
|
||||
|
||||
@@ -15,11 +15,11 @@ interface DatabaseEvent {
|
||||
}
|
||||
|
||||
await on<DatabaseEvent>("database.provisioned", async (e) => {
|
||||
console.log(`[mongodb] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
console.log(`[mongodb] database provisioned for ${e.body?.consumer} (db ${e.body?.database})`);
|
||||
});
|
||||
|
||||
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
|
||||
console.log(`[mongodb] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
console.log(`[mongodb] database deprovisioned for ${e.body?.consumer} (db ${e.body?.database})`);
|
||||
});
|
||||
|
||||
console.log("[mongodb] auditing database lifecycle events");
|
||||
|
||||
@@ -15,11 +15,11 @@ interface TopicEvent {
|
||||
}
|
||||
|
||||
await on<TopicEvent>("topic.provisioned", async (e) => {
|
||||
console.log(`[mosquitto] topic provisioned for ${e.body.consumer} (client ${e.body.username})`);
|
||||
console.log(`[mosquitto] topic provisioned for ${e.body?.consumer} (client ${e.body?.username})`);
|
||||
});
|
||||
|
||||
await on<TopicEvent>("topic.deprovisioned", async (e) => {
|
||||
console.log(`[mosquitto] topic deprovisioned for ${e.body.consumer} (client ${e.body.username})`);
|
||||
console.log(`[mosquitto] topic deprovisioned for ${e.body?.consumer} (client ${e.body?.username})`);
|
||||
});
|
||||
|
||||
console.log("[mosquitto] auditing topic lifecycle events");
|
||||
|
||||
@@ -15,11 +15,11 @@ interface DatabaseEvent {
|
||||
}
|
||||
|
||||
await on<DatabaseEvent>("database.provisioned", async (e) => {
|
||||
console.log(`[mssql] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
console.log(`[mssql] database provisioned for ${e.body?.consumer} (db ${e.body?.database})`);
|
||||
});
|
||||
|
||||
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
|
||||
console.log(`[mssql] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
console.log(`[mssql] database deprovisioned for ${e.body?.consumer} (db ${e.body?.database})`);
|
||||
});
|
||||
|
||||
console.log("[mssql] auditing database lifecycle events");
|
||||
|
||||
@@ -15,11 +15,11 @@ interface CacheEvent {
|
||||
}
|
||||
|
||||
await on<CacheEvent>("cache.provisioned", async (e) => {
|
||||
console.log(`[redis] cache provisioned for ${e.body.consumer} (user ${e.body.username})`);
|
||||
console.log(`[redis] cache provisioned for ${e.body?.consumer} (user ${e.body?.username})`);
|
||||
});
|
||||
|
||||
await on<CacheEvent>("cache.deprovisioned", async (e) => {
|
||||
console.log(`[redis] cache deprovisioned for ${e.body.consumer} (user ${e.body.username})`);
|
||||
console.log(`[redis] cache deprovisioned for ${e.body?.consumer} (user ${e.body?.username})`);
|
||||
});
|
||||
|
||||
console.log("[redis] auditing cache lifecycle events");
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
import { on, emit } from "@novox/mesh-sdk/events";
|
||||
|
||||
await on<{ who?: string }>("greeted", async (event) => {
|
||||
console.log(`[showcase] greeted ${event.body.who ?? "somebody"}`);
|
||||
console.log(`[showcase] greeted ${event.body?.who ?? "somebody"}`);
|
||||
// A consumer may emit, which is what makes an event graph rather than a list of sinks.
|
||||
await emit("acknowledged", { who: event.body.who ?? "somebody" });
|
||||
await emit("acknowledged", { who: event.body?.who ?? "somebody" });
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user