// 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. Delivery is at-least-once, so an upsert is idempotent-latest by design — a duplicate // event writes the same row again, never a second. // // 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; } const DDL = ` 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) ); `; 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 { return new UsageStore(new Pool({ connectionString: requireEnv("DATABASE_URL", env) })); } /** 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 { await this.pool.query(DDL); } /** Upsert a reading, keeping the latest per `(licence, consumer, period, metric)`. */ async upsert(row: UsageRow): Promise { 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 ?? {})], ); } /** 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(); } }