From e8c3ea89799098bd2cade91d141f741cf86989ed Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 7 Sep 2026 04:07:26 +0200 Subject: [PATCH] model-usage: the vendor-neutral usage store (ADR 0054) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The home ADR 0050 left open for a usage reading. mesh-control is a CLI and cannot consume events, so the store that keeps the current usage picture is a MODULE — the audit-logger's sibling: it consumes `module.*.usage.*` and upserts each reading into its own provisioned postgres store, latest per (licence, consumer, period, metric), in the clear. One vendor-neutral table holds BOTH grains; they differ only in `consumer` (the holding module for the licence grain, the session for the finer one). The consumer creates its table on startup and, as a restart-until-ready service, self-heals rather than gating the apply on a run-once that must reach a provider over the overlay. The vendor->row normalisation moves into the adapter, as ADR 0054 requires: anthropic-manager (licence grain, utilization%) and anthropic-consumer (session grain, token/cost) now emit already-normalised { rows: UsageRow[], raw } on their existing keys, so the store stays vendor-blind. Proven end to end by the mesh-lab model-usage bed (green): a usage event emitted into the mesh is upserted at both grains, latest-per-key, in the clear. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- modules/anthropic-consumer/usage/index.ts | 49 +++++++++++- modules/anthropic-manager/refresh/index.ts | 36 ++++++++- modules/model-usage/index.ts | 38 ++++++++++ modules/model-usage/module.json | 66 +++++++++++++++++ modules/model-usage/package.json | 15 ++++ modules/model-usage/pg.d.ts | 15 ++++ modules/model-usage/store.ts | 86 ++++++++++++++++++++++ modules/model-usage/tools/index.ts | 32 ++++++++ modules/model-usage/tsconfig.json | 18 +++++ 9 files changed, 351 insertions(+), 4 deletions(-) create mode 100644 modules/model-usage/index.ts create mode 100644 modules/model-usage/module.json create mode 100644 modules/model-usage/package.json create mode 100644 modules/model-usage/pg.d.ts create mode 100644 modules/model-usage/store.ts create mode 100644 modules/model-usage/tools/index.ts create mode 100644 modules/model-usage/tsconfig.json diff --git a/modules/anthropic-consumer/usage/index.ts b/modules/anthropic-consumer/usage/index.ts index 7f31b11..f874547 100644 --- a/modules/anthropic-consumer/usage/index.ts +++ b/modules/anthropic-consumer/usage/index.ts @@ -6,11 +6,52 @@ // Runs as `mesh-tools run` (no broker), so events are emitted best-effort via the sibling mesh-tools // `emit` primitive; the totals are also written to a file so the reading is observable without one. -import { readdirSync, statSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; +import { readdirSync, statSync, readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; import { join, dirname } from "node:path"; import { readSessionFile, type SessionUsage } from "../transcript.js"; +/** The vendor-neutral usage row ADR 0054 fixes — the shape the model-usage store upserts. Kept local + * to the producer (the normalisation lives in the adapter), so nothing here couples to the SDK. */ +interface UsageRow { + licence: string; + consumer: string; + period: string; + metric: string; + value: number; +} + +/** The licence this session's usage is charged to. The model-access binding names it; failing that, + * the deployed env; failing that, "unknown" — a reading is never dropped for want of a licence. */ +function boundLicence(): string { + const bindFile = process.env.MESH_MODEL_ACCESS_BIND_FILE; + if (bindFile) { + try { + const raw = JSON.parse(readFileSync(bindFile, "utf8")) as { licence?: unknown }; + if (typeof raw.licence === "string" && raw.licence) return raw.licence; + } catch { + // A missing or unreadable bind file is not fatal — fall through to the env, then to "unknown". + } + } + return process.env.MESH_ANTHROPIC_LICENCE || "unknown"; +} + +/** Normalise one session reading into ADR 0054 rows: one row per metric, a row omitted when its + * number is not finite. `consumer` is node/module/session — the session grain. */ +function sessionRows(licence: string, node: string, module: string, r: SessionUsage): UsageRow[] { + const consumer = `${node}/${module}/${r.sessionId}`; + const rows: UsageRow[] = []; + const add = (metric: string, value: number): void => { + if (Number.isFinite(value)) rows.push({ licence, consumer, period: "session", metric, value }); + }; + add("input_tokens", r.inputTokens); + add("output_tokens", r.outputTokens); + add("cache_creation_tokens", r.cacheCreationTokens); + add("cache_read_tokens", r.cacheReadTokens); + add("cost_usd", r.costUSD); + return rows; +} + function projectsDir(): string { return process.env.MESH_CLAUDE_PROJECTS_DIR ?? `${process.env.HOME ?? "/root"}/.claude/projects`; } @@ -54,8 +95,12 @@ async function main(): Promise { } } + // Emit ADR-0054 rows, not a vendor-shaped body: the consumer of module.*.usage.* is the + // vendor-neutral model-usage store, so the session→row normalisation is done HERE. The full + // SessionUsage rides as `raw`, so model, branch, cwd and timestamps are not lost. + const licence = boundLicence(); for (const r of readings) { - await emitUsage({ grain: "session", node, module, ...r }); + await emitUsage({ rows: sessionRows(licence, node, module, r), raw: r }); } if (process.env.MESH_ANTHROPIC_USAGE_OUT) { diff --git a/modules/anthropic-manager/refresh/index.ts b/modules/anthropic-manager/refresh/index.ts index 25df864..e2f99de 100644 --- a/modules/anthropic-manager/refresh/index.ts +++ b/modules/anthropic-manager/refresh/index.ts @@ -25,9 +25,36 @@ import { readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; import { dirname } from "node:path"; +import { readEnv } from "@novox/mesh-sdk/primitives"; + import { seal } from "../sealedbox.js"; import { managerPublicKey, writeSealedGrant } from "../grantfile.js"; -import { refreshGrant, grantFromRefresh, readUsage, flattenUsage } from "../client.js"; +import { refreshGrant, grantFromRefresh, readUsage, flattenUsage, type UsageReading } from "../client.js"; + +/** The vendor-neutral usage row ADR 0054 fixes — the shape the model-usage store upserts. Kept local + * to the producer (the normalisation lives in the adapter), so nothing here couples to the SDK. */ +interface UsageRow { + licence: string; + consumer: string; + period: string; + metric: string; + value: number; +} + +/** Normalise a licence-grain reading into ADR 0054 rows: one utilization row per window, a row + * omitted when its percentage is absent. `consumer` is the holding module — the licence grain. */ +function licenceRows(licence: string, consumer: string, reading: UsageReading): UsageRow[] { + const rows: UsageRow[] = []; + const add = (period: string, pct: number | null): void => { + if (pct !== null && pct !== undefined && Number.isFinite(pct)) { + rows.push({ licence, consumer, period, metric: "utilization", value: pct }); + } + }; + add("5h", reading.sessionPct); + add("7d", reading.weeklyPct); + add("extra", reading.extraPct); + return rows; +} function required(name: string): string { const v = process.env[name]; @@ -95,7 +122,12 @@ async function main(): Promise { JSON.stringify({ licence, grain: "licence", ...reading }), ); } - await emitUsage({ licence, grain: "licence", ...reading }); + // Emit ADR-0054 rows, not a vendor-shaped body: the consumer of module.*.usage.* is the + // vendor-neutral model-usage store, so the normalisation is done HERE. The node names the + // holding module; with MESH_NODE unset the consumer is the module alone. + const node = readEnv("MESH_NODE", ""); + const consumer = node ? `${node}/anthropic-manager` : "anthropic-manager"; + await emitUsage({ rows: licenceRows(licence, consumer, reading), raw: usage }); } } catch (err) { console.error(`[anthropic-manager] usage poll for ${licence} failed: ${err}`); diff --git a/modules/model-usage/index.ts b/modules/model-usage/index.ts new file mode 100644 index 0000000..18b5ec4 --- /dev/null +++ b/modules/model-usage/index.ts @@ -0,0 +1,38 @@ +// model-usage's entrypoint — the usage context store's consumer (novox/hq ADR 0054). mesh-control 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. +// +// 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. + +import { on } from "@novox/mesh-sdk/events"; +import { UsageStore, type UsageRow } from "./store.js"; + +const store = UsageStore.fromEnv(); + +// 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(); + +await on("module.*.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}`); + } + } +}); + +console.log("[model-usage] recording model usage to its store"); diff --git a/modules/model-usage/module.json b/modules/model-usage/module.json new file mode 100644 index 0000000..d9cdba6 --- /dev/null +++ b/modules/model-usage/module.json @@ -0,0 +1,66 @@ +{ + "module": "model-usage", + "version": "1", + "slug": "usage", + "capabilities": [ + "container-runtime" + ], + "requires": [ + "postgres-database" + ], + "contributes": { + "postgres-database": { + "name": "model_usage" + } + }, + "binds": { + "postgres-database": "/var/lib/model-usage/database.json" + }, + "secrets": { + "postgres-database": "/var/lib/model-usage/database.secret" + }, + "consumes": [ + "module.*.usage.*" + ], + "own-secrets": { + "broker": "/var/lib/mesh/model-usage/broker" + }, + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/model-usage", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/model-usage", + "mode": "0700" + }, + { + "id": "db-env", + "type": "file", + "path": "/var/lib/model-usage/db.env", + "mode": "0600", + "content": "DATABASE_URL=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": [ + "/var/lib/mesh/model-usage/broker:/run/secrets/broker:ro", + "/var/lib/model-usage:/run/state" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker" + }, + "env-file": [ + "/var/lib/model-usage/db.env" + ] + } + ] +} diff --git a/modules/model-usage/package.json b/modules/model-usage/package.json new file mode 100644 index 0000000..916d04c --- /dev/null +++ b/modules/model-usage/package.json @@ -0,0 +1,15 @@ +{ + "name": "@novox/module-model-usage", + "version": "0.1.0", + "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, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0", + "pg": "^8" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/model-usage/pg.d.ts b/modules/model-usage/pg.d.ts new file mode 100644 index 0000000..d11d56e --- /dev/null +++ b/modules/model-usage/pg.d.ts @@ -0,0 +1,15 @@ +// Ambient types for `pg` (node-postgres), which ships its types only via the separate `@types/pg` +// package. Rather than pull that in at tsc time, this declares the exact slice model-usage uses — +// the same precedent anthropic-manager sets for `tweetnacl-sealedbox-js` (a local ambient .d.ts, +// listed in tsconfig `include`, default-imported). The real `pg` is installed into the module's +// 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. */ + export class Pool { + constructor(config?: { connectionString?: string }); + query(text: string, params?: unknown[]): Promise<{ rows: any[]; rowCount: number }>; + end(): Promise; + } +} diff --git a/modules/model-usage/store.ts b/modules/model-usage/store.ts new file mode 100644 index 0000000..4f96563 --- /dev/null +++ b/modules/model-usage/store.ts @@ -0,0 +1,86 @@ +// 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(); + } +} diff --git a/modules/model-usage/tools/index.ts b/modules/model-usage/tools/index.ts new file mode 100644 index 0000000..02571c0 --- /dev/null +++ b/modules/model-usage/tools/index.ts @@ -0,0 +1,32 @@ +// 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. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { UsageStore } from "../store.js"; + +export function getUsageTools(store: UsageStore): ToolDefinition[] { + return [ + { + name: "model_usage_current", + description: + "The latest model-usage reading per (licence, consumer, period, metric), across both grains, optionally restricted to one licence.", + input: { + licence: { type: "string", description: "restrict to one licence (optional)" }, + }, + run: async (args) => { + const licence = args.licence ? String(args.licence) : undefined; + return { usage: await store.current(licence) }; + }, + }, + ]; +} + +registerModuleTools("model-usage", (env) => { + try { + return getUsageTools(UsageStore.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/model-usage/tsconfig.json b/modules/model-usage/tsconfig.json new file mode 100644 index 0000000..5dcbb4e --- /dev/null +++ b/modules/model-usage/tsconfig.json @@ -0,0 +1,18 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": [ + "pg.d.ts", + "store.ts", + "index.ts", + "tools/index.ts", + "migrate/index.ts" + ] +} -- 2.54.0