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.
49 lines
2.5 KiB
TypeScript
49 lines
2.5 KiB
TypeScript
// 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 });
|
|
};
|
|
}
|