Files
mesh-catalog/modules/anthropic-consumer/usage/index.ts
T
jochen 4aacea5453 anthropic-consumer: its usage runs in the node's runtime, and its apply is a scheduled process (hq ADR 0198)
Both containers go with the Dockerfile, build bases, bus credential and state directory. apply needs no bus and runs every five minutes as a process on the machine at the host paths the container mounted. usage emitted by spawning the runtime image's own emit command with the module's credential, which exists nowhere now, so it is loaded by the node's runtime instead: it emits through the SDK as this module and reads on the cadence the schedule gave it, once at start and every five minutes. That is the one code change.
2026-10-04 00:37:49 +02:00

139 lines
5.4 KiB
TypeScript

// Session-grain usage emission (novox/hq ADR 0054). On a schedule, read every transcript under
// `~/.claude/projects/*/<sessionId>.jsonl`, sum its tokens, and emit one session-grain usage event
// per session. The consumer IS the (node,module) session's fixed binding, so no per-message account
// attribution is done — just the totals (port map "don't-map" #3).
//
// Runs in the node's runtime (novox/hq ADR 0198), every five minutes, so events are emitted through
// the runtime as this module; the totals are also written to a file so the reading is observable
// without one.
import { readdirSync, statSync, readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs";
import { join, dirname } from "node:path";
import { emit } from "@novox/mesh-sdk/events";
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`;
}
/** Every `<sessionId>.jsonl` under the projects tree, with the project directory it sits in. */
function transcripts(root: string): { path: string; sessionId: string }[] {
const found: { path: string; sessionId: string }[] = [];
let projects: string[];
try {
projects = readdirSync(root);
} catch {
return found; // no projects yet is not a failure — there is simply nothing to report.
}
for (const proj of projects) {
const dir = join(root, proj);
let entries: string[];
try {
if (!statSync(dir).isDirectory()) continue;
entries = readdirSync(dir);
} catch {
continue;
}
for (const file of entries) {
if (!file.endsWith(".jsonl")) continue;
found.push({ path: join(dir, file), sessionId: file.replace(/\.jsonl$/, "") });
}
}
return found;
}
async function main(): Promise<void> {
const module = process.env.MESH_MODULE ?? "anthropic-consumer";
const node = process.env.MESH_NODE ?? "unknown";
const readings: SessionUsage[] = [];
for (const t of transcripts(projectsDir())) {
try {
readings.push(await readSessionFile(t.path, t.sessionId));
} catch (err) {
console.error(`[anthropic-consumer] could not read ${t.path}: ${err}`);
}
}
// 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({ rows: sessionRows(licence, node, module, r), raw: r });
}
if (process.env.MESH_ANTHROPIC_USAGE_OUT) {
atomicWrite(process.env.MESH_ANTHROPIC_USAGE_OUT, JSON.stringify(readings, null, 2));
}
console.error(`[anthropic-consumer] reported ${readings.length} session(s)`);
}
function atomicWrite(path: string, content: string): void {
mkdirSync(dirname(path), { recursive: true });
const tmp = `${path}.tmp`;
writeFileSync(tmp, content, { mode: 0o600 });
renameSync(tmp, path);
}
/** Emit best-effort through the runtime: a reading that could not be announced is still in the file. */
async function emitUsage(body: Record<string, unknown>): Promise<void> {
try {
await emit("usage.session", body);
} catch (err) {
console.error(`[anthropic-consumer] could not emit usage: ${err}`);
}
}
// The cadence the scheduled container had: once at start, then every five minutes. Not awaited, so the
// runtime's handshake is answered while a long first reading is still under way.
const EVERY_MS = 5 * 60 * 1000;
const tick = (): void => {
void main().catch((err) => console.error(`[anthropic-consumer] usage reading failed: ${err}`));
};
tick();
setInterval(tick, EVERY_MS);