// Session-grain usage emission (novox/hq ADR 0054). On a schedule, read every transcript under // `~/.claude/projects/*/.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 `.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 { 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): Promise { 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);