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
138 lines
5.4 KiB
TypeScript
138 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 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, 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`;
|
|
}
|
|
|
|
/** 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 via the sibling mesh-tools `emit`, which wires a broker a run step has none. */
|
|
async function emitUsage(body: Record<string, unknown>): Promise<void> {
|
|
const main = process.env.MESH_TOOLS_MAIN ?? "/app/dist/main.js";
|
|
const { spawn } = await import("node:child_process");
|
|
await new Promise<void>((resolve) => {
|
|
const child = spawn(
|
|
process.execPath,
|
|
[main, "emit", "module.anthropic-consumer.usage.session", JSON.stringify(body)],
|
|
{ stdio: "inherit" },
|
|
);
|
|
child.on("exit", () => resolve());
|
|
child.on("error", (err) => {
|
|
console.error(`[anthropic-consumer] could not emit usage: ${err}`);
|
|
resolve();
|
|
});
|
|
});
|
|
}
|
|
|
|
await main();
|