import { readFileSync } from "node:fs"; // The vendor-neutral usage store (novox/hq ADR 0054). ONE table holds BOTH grains — licence and // 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 — 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. import pg from "pg"; import { requireEnv } from "@novox/mesh-sdk/primitives"; // `pg` is CommonJS: its default export is the module object, so the Pool class is a property of it. // Destructuring the default is the form that resolves at runtime under Node's ESM loader (the same // shape anthropic-manager uses for tweetnacl-sealedbox-js). const { Pool } = pg; type PgPool = InstanceType; /** One usage reading — the vendor-neutral row ADR 0054 fixes. `raw` carries the vendor payload the * reading was normalised from (model, windows, timestamps), so nothing is lost by normalising. */ export interface UsageRow { licence: string; consumer: string; period: string; metric: string; value: number; raw?: unknown; } /** * 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, period text NOT NULL, metric text NOT NULL, value double precision NOT NULL, 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) {} /** Build a store from the resolved environment — DATABASE_URL is the granted postgres connection, * templated into the module's env-file from the mesh's binding (umami's DATABASE_URL precedent). */ static fromEnv(env: NodeJS.ProcessEnv = process.env): 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"); // 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 })); } /** 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 { await this.pool.query( "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 { 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. */ async current(licence?: string): Promise { const where = licence ? "WHERE licence = $1" : ""; const params = licence ? [licence] : []; const { rows } = await this.pool.query( `SELECT licence, consumer, period, metric, value, raw FROM usage ${where} ORDER BY licence, consumer, period, metric`, params, ); return rows as UsageRow[]; } async close(): Promise { await this.pool.end(); } } /** The content of a file the environment names, its line ending gone — or undefined when it names none. */ function readMaybe(path: string | undefined): string | undefined { if (!path) return undefined; try { return readFileSync(path, "utf8").replace(/\r?\n$/, ""); } catch { return undefined; } }