model-usage: the vendor-neutral usage store (ADR 0054)
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
This commit is contained in:
@@ -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<void> {
|
||||
}
|
||||
}
|
||||
|
||||
// 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) {
|
||||
|
||||
@@ -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<void> {
|
||||
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}`);
|
||||
|
||||
@@ -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");
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
Vendored
+15
@@ -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<void>;
|
||||
}
|
||||
}
|
||||
@@ -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<typeof Pool>;
|
||||
|
||||
/** 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<void> {
|
||||
await this.pool.query(DDL);
|
||||
}
|
||||
|
||||
/** Upsert a reading, keeping the latest per `(licence, consumer, period, metric)`. */
|
||||
async upsert(row: UsageRow): Promise<void> {
|
||||
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<UsageRow[]> {
|
||||
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<void> {
|
||||
await this.pool.end();
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
Reference in New Issue
Block a user