diff --git a/modules/anthropic-consumer/apply/index.ts b/modules/anthropic-consumer/apply/index.ts new file mode 100644 index 0000000..c0ab37a --- /dev/null +++ b/modules/anthropic-consumer/apply/index.ts @@ -0,0 +1,75 @@ +// The consumer's scheduled run: take the ACCESS token the mesh delivered and write it where the +// Claude CLI reads it, access-token-only (novox/hq ADR 0050). The refresh token is never here to +// strip — the manager holds it, and a holder's delivery has only ever been the access token. +// +// What the host delivers, per the manifest: +// secrets.model-access -> a file holding the sealed-then-unsealed ACCESS token (the host opened it +// with this node's private key; this process reads plaintext). +// binds.model-access -> a JSON file of the non-secret facts the licence serves (which licence, +// model, and — when the control plane carries them — grant expiry/scopes). +// +// Runs as `mesh-tools run` (no broker) on a schedule, so it is idempotent: same token in, same file +// out. + +import { readFileSync } from "node:fs"; + +import { deliver, type DeliveredGrant } from "../credentials.js"; +import { readAccountUuid, check } from "../identity.js"; + +function required(name: string): string { + const v = process.env[name]; + if (!v) throw new Error(`${name} is not set — the consumer runtime was deployed without it`); + return v; +} + +/** Read optional non-secret grant metadata (expiry, scopes, subscription) from the bound facts file. */ +function readBoundMeta(path: string | undefined): Partial { + if (!path) return {}; + try { + const raw = JSON.parse(readFileSync(path, "utf8")) as Record; + return { + expiresAt: typeof raw.expiresAt === "number" ? raw.expiresAt : null, + refreshTokenExpiresAt: typeof raw.refreshTokenExpiresAt === "number" ? raw.refreshTokenExpiresAt : null, + scopes: Array.isArray(raw.scopes) ? (raw.scopes as string[]) : null, + subscriptionType: typeof raw.subscriptionType === "string" ? raw.subscriptionType : null, + }; + } catch { + return {}; + } +} + +function main(): void { + const accessToken = readFileSync(required("MESH_MODEL_ACCESS_SECRET_FILE"), "utf8").trim(); + if (!accessToken) { + // Nothing was delivered — which reads exactly like a credential that never arrived, so it is + // said rather than written as an empty file the CLI would take for a login it should not do. + throw new Error("[anthropic-consumer] the delivered access token is empty; nothing was written"); + } + + const meta = readBoundMeta(process.env.MESH_MODEL_ACCESS_BIND_FILE); + const grant: DeliveredGrant = { accessToken, ...meta }; + + const target = process.env.MESH_CLAUDE_CREDENTIALS_FILE ?? `${homedir()}/.claude/.credentials.json`; + deliver(target, grant); + console.error(`[anthropic-consumer] wrote an access-token-only credential to ${target}`); + + // The mis-binding guard, best-effort and fail-closed. The expected account uuid is not yet plumbed + // (identity.ts TODO), so this reports what it can see rather than acting on it — it never delivers + // to a wrong account because it never learns one to deliver to. + const identityFile = process.env.MESH_CLAUDE_IDENTITY_FILE ?? `${homedir()}/.claude.json`; + const found = readAccountUuid(identityFile); + const expected = process.env.MESH_MODEL_ACCESS_ACCOUNT_UUID ?? null; + const verdict = check(found, expected); + if (verdict.state === "wrong-account") { + throw new Error( + `[anthropic-consumer] the CLI is logged in as ${verdict.found}, not the licensed ${verdict.expected}; refusing`, + ); + } + console.error(`[anthropic-consumer] identity check: ${verdict.state}`); +} + +function homedir(): string { + return process.env.HOME ?? "/root"; +} + +main(); diff --git a/modules/anthropic-consumer/credentials.ts b/modules/anthropic-consumer/credentials.ts new file mode 100644 index 0000000..574c9f5 --- /dev/null +++ b/modules/anthropic-consumer/credentials.ts @@ -0,0 +1,82 @@ +// Writing the access token where the Claude CLI reads it — the consumer half of model-access +// (novox/hq ADR 0050). A node holds an ACCESS token and nothing else: it cannot rotate, so it is +// never given a refresh token, and this enforces that on every write. +// +// The file shape and the strip are ported byte-exact from the mature implementation (see the port +// map): `~/.claude/.credentials.json` → `{ claudeAiOauth: { accessToken, expiresAt, +// refreshTokenExpiresAt?, scopes?, subscriptionType? } }`, and the refresh token is deleted, not +// merely omitted, so a full grant left by an interactive login is stripped back to access-only. + +import { readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; +import { dirname } from "node:path"; + +/** The access-token-only grant the mesh delivered — what the manager submitted, minus the refresh. */ +export interface DeliveredGrant { + readonly accessToken: string; + readonly expiresAt?: number | null; + readonly refreshTokenExpiresAt?: number | null; + readonly scopes?: string[] | null; + readonly subscriptionType?: string | null; +} + +interface ClaudeOauth { + accessToken?: string; + expiresAt?: number; + refreshTokenExpiresAt?: number; + scopes?: string[]; + subscriptionType?: string; + refreshToken?: string; +} + +interface Credentials { + claudeAiOauth?: ClaudeOauth; + [key: string]: unknown; +} + +/** Read the existing credentials file, or an empty object if there is none or it is unreadable. */ +function readLocal(path: string): Credentials { + try { + return JSON.parse(readFileSync(path, "utf8")) as Credentials; + } catch { + return {}; + } +} + +/** + * Overlay the delivered grant onto whatever is on disk, then STRIP the refresh token — the node + * carve-out. Returns the object to write, so the strip is testable without touching a file. + */ +export function applyGrant(local: Credentials, grant: DeliveredGrant): Credentials { + const oauth = local.claudeAiOauth ?? {}; + const next: Credentials = { + ...local, + claudeAiOauth: { + ...oauth, + accessToken: grant.accessToken, + ...(grant.expiresAt != null ? { expiresAt: grant.expiresAt } : {}), + ...(grant.refreshTokenExpiresAt != null + ? { refreshTokenExpiresAt: grant.refreshTokenExpiresAt } + : {}), + ...(grant.scopes ? { scopes: grant.scopes } : {}), + ...(grant.subscriptionType ? { subscriptionType: grant.subscriptionType } : {}), + }, + }; + // A node NEVER holds a refresh token: delete it, so a full grant on disk is reduced to access-only. + delete next.claudeAiOauth!.refreshToken; + return next; +} + +/** Atomic write-then-rename at 0600 — a partial credentials file must never be read as a whole one. */ +export function writeCredentials(path: string, creds: Credentials): void { + mkdirSync(dirname(path), { recursive: true }); + const tmp = `${path}.tmp`; + writeFileSync(tmp, JSON.stringify(creds, null, 2), { mode: 0o600 }); + renameSync(tmp, path); +} + +/** Read, overlay, strip, write — the whole consumer credential update, in one call. */ +export function deliver(path: string, grant: DeliveredGrant): Credentials { + const next = applyGrant(readLocal(path), grant); + writeCredentials(path, next); + return next; +} diff --git a/modules/anthropic-consumer/identity.ts b/modules/anthropic-consumer/identity.ts new file mode 100644 index 0000000..e3af1c0 --- /dev/null +++ b/modules/anthropic-consumer/identity.ts @@ -0,0 +1,49 @@ +// The mis-binding guard (novox/hq ADR 0050, port map §identity). Account identity is NOT in the +// token or any API — it lives in a sibling CLI state file, `~/.claude.json` → +// `oauthAccount.accountUuid`. The guard compares the account the CLI is actually logged in as to the +// account the licence was recorded against, and FAILS CLOSED: an absent file or an unrecorded licence +// account refuses rather than guesses, because delivering an access token to the wrong account is the +// exact fault this exists to catch. +// +// **Partial first cut, FLAGGED.** Reading the sibling file is implemented; the licence's recorded +// account uuid is not yet plumbed from the control plane to the consumer (the bound `model.json` does +// not carry it today). So `check` returns `licence-not-adopted` when no expected uuid is supplied, +// which is the fail-closed answer, and the wiring of the expected uuid is a TODO below. + +import { readFileSync } from "node:fs"; + +export type IdentityVerdict = + | { state: "verified"; accountUuid: string } + | { state: "no-identity-file" } + | { state: "licence-not-adopted" } + | { state: "wrong-account"; found: string; expected: string }; + +interface ClaudeJson { + oauthAccount?: { accountUuid?: string; emailAddress?: string; organizationUuid?: string }; +} + +/** Read `oauthAccount.accountUuid` from `~/.claude.json`, or null if the file or field is absent. */ +export function readAccountUuid(path: string): string | null { + try { + const raw = JSON.parse(readFileSync(path, "utf8")) as ClaudeJson; + return raw.oauthAccount?.accountUuid ?? null; + } catch { + return null; + } +} + +/** + * Compare the CLI's logged-in account to the one the licence was recorded against. Pure over its + * inputs so the fail-closed logic is tested without a filesystem. + * + * TODO(novox/hq ADR 0050, Phase C): plumb `expected` — the licence's recorded account uuid — from the + * control plane into the consumer's bound `model.json`, then adopt-on-first-sight or refuse per the + * port map's five states. Until then only the two safe verdicts are reachable: verified when an + * expected uuid is provided and matches, refuse otherwise. + */ +export function check(found: string | null, expected: string | null): IdentityVerdict { + if (found === null) return { state: "no-identity-file" }; + if (!expected) return { state: "licence-not-adopted" }; + if (found === expected) return { state: "verified", accountUuid: found }; + return { state: "wrong-account", found, expected }; +} diff --git a/modules/anthropic-consumer/module.json b/modules/anthropic-consumer/module.json new file mode 100644 index 0000000..0b5e52c --- /dev/null +++ b/modules/anthropic-consumer/module.json @@ -0,0 +1,91 @@ +{ + "module": "anthropic-consumer", + "version": "1", + "capabilities": [ + "container-runtime" + ], + "requires": [ + "model-access" + ], + "binds": { + "model-access": "/var/lib/anthropic-consumer/model.json" + }, + "secrets": { + "model-access": "/var/lib/anthropic-consumer/access-token" + }, + "own-secrets": { + "broker": "/var/lib/mesh/anthropic-consumer/broker" + }, + "emits": [ + "module.anthropic-consumer.usage.session" + ], + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/anthropic-consumer", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/anthropic-consumer", + "mode": "0700" + }, + { + "id": "claude-home", + "type": "directory", + "path": "/var/lib/anthropic-consumer/claude", + "mode": "0700" + }, + { + "id": "out", + "type": "directory", + "path": "/var/lib/anthropic-consumer/out", + "mode": "0700" + }, + { + "id": "apply", + "type": "container", + "name": "mesh-anthropic-consumer-apply", + "image": "mesh-runtime-anthropic-consumer@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "host", + "schedule": "*/5 * * * *", + "args": [ + "run", + "/app/modules/anthropic-consumer/dist/apply/index.js" + ], + "volumes": [ + "/var/lib/anthropic-consumer:/run/state" + ], + "env": { + "MESH_MODEL_ACCESS_SECRET_FILE": "/run/state/access-token", + "MESH_MODEL_ACCESS_BIND_FILE": "/run/state/model.json", + "MESH_CLAUDE_CREDENTIALS_FILE": "/run/state/claude/.credentials.json", + "MESH_CLAUDE_IDENTITY_FILE": "/run/state/claude/.claude.json" + } + }, + { + "id": "usage", + "type": "container", + "name": "mesh-anthropic-consumer-usage", + "image": "mesh-runtime-anthropic-consumer@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "host", + "schedule": "*/5 * * * *", + "args": [ + "run", + "/app/modules/anthropic-consumer/dist/usage/index.js" + ], + "volumes": [ + "/var/lib/mesh/anthropic-consumer/broker:/run/secrets/broker:ro", + "/var/lib/anthropic-consumer:/run/state" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_CLAUDE_PROJECTS_DIR": "/run/state/claude/projects", + "MESH_ANTHROPIC_USAGE_OUT": "/run/state/out/session-usage.json", + "MESH_TOOLS_MAIN": "/app/dist/main.js" + } + } + ] +} diff --git a/modules/anthropic-consumer/package.json b/modules/anthropic-consumer/package.json new file mode 100644 index 0000000..ef3f418 --- /dev/null +++ b/modules/anthropic-consumer/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-anthropic-consumer", + "version": "0.1.0", + "description": "anthropic-consumer — the consumer side of model-access (ADR 0050): writes the delivered access token to ~/.claude/.credentials.json (access-token-only) and reports session-grain usage from the CLI transcripts (ADR 0054).", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/anthropic-consumer/test/credentials.test.ts b/modules/anthropic-consumer/test/credentials.test.ts new file mode 100644 index 0000000..72ddcd7 --- /dev/null +++ b/modules/anthropic-consumer/test/credentials.test.ts @@ -0,0 +1,32 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtempSync, readFileSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import { applyGrant, deliver } from "../credentials.ts"; + +test("applyGrant strips the refresh token a full grant on disk left behind", () => { + const local = { claudeAiOauth: { accessToken: "at-old", refreshToken: "rt-must-not-survive" } }; + const next = applyGrant(local, { accessToken: "at-new", expiresAt: 123 }); + assert.equal(next.claudeAiOauth!.accessToken, "at-new"); + assert.equal(next.claudeAiOauth!.expiresAt, 123); + assert.ok(!("refreshToken" in next.claudeAiOauth!), "a node held onto a refresh token"); +}); + +test("deliver writes the port-map shape, access-token-only, and never a refresh token", () => { + const dir = mkdtempSync(join(tmpdir(), "anthropic-consumer-")); + const path = join(dir, ".credentials.json"); + // A prior interactive login left a full grant on disk. + writeFileSync(path, JSON.stringify({ claudeAiOauth: { accessToken: "at-old", refreshToken: "rt-login" } })); + + deliver(path, { accessToken: "at-delivered", expiresAt: 999, subscriptionType: "max" }); + + const raw = readFileSync(path, "utf8"); + const creds = JSON.parse(raw); + assert.equal(creds.claudeAiOauth.accessToken, "at-delivered"); + assert.equal(creds.claudeAiOauth.expiresAt, 999); + assert.equal(creds.claudeAiOauth.subscriptionType, "max"); + assert.doesNotMatch(raw, /rt-login/, "the refresh token is still on disk"); + assert.ok(!("refreshToken" in creds.claudeAiOauth)); +}); diff --git a/modules/anthropic-consumer/test/transcript.test.ts b/modules/anthropic-consumer/test/transcript.test.ts new file mode 100644 index 0000000..5048b5a --- /dev/null +++ b/modules/anthropic-consumer/test/transcript.test.ts @@ -0,0 +1,48 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { foldTranscript } from "../transcript.ts"; + +// A captured-shape transcript: two assistant turns and a user line, exactly the fields the port map +// names. Not imagined — the field names match the mature implementation's parse. +const TRANSCRIPT = [ + JSON.stringify({ type: "user", timestamp: "2026-01-01T00:00:00Z", cwd: "/work/app", gitBranch: "main" }), + JSON.stringify({ + type: "assistant", + timestamp: "2026-01-01T00:00:01Z", + costUSD: 0.01, + message: { + model: "claude-opus-4-8", + usage: { input_tokens: 100, cache_creation_input_tokens: 20, cache_read_input_tokens: 5, output_tokens: 40 }, + }, + }), + JSON.stringify({ + type: "assistant", + timestamp: "2026-01-01T00:00:02Z", + costUSD: 0.02, + message: { model: "claude-opus-4-8", usage: { input_tokens: 200, output_tokens: 60 } }, + }), + "", // a half-written trailing line is ordinary and must not be fatal. +].join("\n"); + +test("a transcript sums per-session token counts, cost, and metadata", () => { + const s = foldTranscript("session-abc", TRANSCRIPT); + assert.equal(s.sessionId, "session-abc"); + assert.equal(s.turns, 2); + assert.equal(s.inputTokens, 300); + assert.equal(s.cacheCreationTokens, 20); + assert.equal(s.cacheReadTokens, 5); + assert.equal(s.outputTokens, 100); + assert.equal(Math.round(s.costUSD * 100) / 100, 0.03); + assert.equal(s.model, "claude-opus-4-8"); + assert.equal(s.gitBranch, "main"); + assert.equal(s.cwd, "/work/app"); + assert.equal(s.startedAt, "2026-01-01T00:00:00Z"); + assert.equal(s.lastActive, "2026-01-01T00:00:02Z"); +}); + +test("a malformed line is skipped, not fatal", () => { + const s = foldTranscript("s", 'not json\n{"type":"assistant","message":{"usage":{"output_tokens":7}}}'); + assert.equal(s.outputTokens, 7); + assert.equal(s.turns, 1); +}); diff --git a/modules/anthropic-consumer/transcript.ts b/modules/anthropic-consumer/transcript.ts new file mode 100644 index 0000000..bb6149b --- /dev/null +++ b/modules/anthropic-consumer/transcript.ts @@ -0,0 +1,113 @@ +// Session-grain usage from the CLI's own transcripts (novox/hq ADR 0054). The mature implementation +// reads `~/.claude/projects//.jsonl` and sums the token counts each assistant +// message reports; this ports the token extraction and DROPS the per-message account-attribution +// timeline — the nox (node,module) session has a fixed licence binding (port map "don't-map" #3), so +// there is nothing to attribute per message. +// +// The fields are ported from the port map: assistant lines carry +// `message.usage.{input_tokens,cache_creation_input_tokens,cache_read_input_tokens,output_tokens}`, +// `costUSD`, `message.model`, `timestamp`; user lines carry `cwd`, `gitBranch`. + +import { createInterface } from "node:readline"; +import { createReadStream } from "node:fs"; + +/** One session's totals — the session-grain usage row ADR 0054 fixes. */ +export interface SessionUsage { + sessionId: string; + model: string | null; + gitBranch: string | null; + cwd: string | null; + turns: number; + inputTokens: number; + cacheCreationTokens: number; + cacheReadTokens: number; + outputTokens: number; + costUSD: number; + startedAt: string | null; + lastActive: string | null; +} + +interface Line { + type?: string; + timestamp?: string; + cwd?: string; + gitBranch?: string; + costUSD?: number; + message?: { + model?: string; + usage?: { + input_tokens?: number; + cache_creation_input_tokens?: number; + cache_read_input_tokens?: number; + output_tokens?: number; + }; + }; +} + +function empty(sessionId: string): SessionUsage { + return { + sessionId, + model: null, + gitBranch: null, + cwd: null, + turns: 0, + inputTokens: 0, + cacheCreationTokens: 0, + cacheReadTokens: 0, + outputTokens: 0, + costUSD: 0, + startedAt: null, + lastActive: null, + }; +} + +/** Fold one transcript line into a session's running totals. Pure, so it is tested on fixtures. */ +export function foldLine(acc: SessionUsage, raw: string): SessionUsage { + const line = parse(raw); + if (!line) return acc; + + if (line.timestamp) { + if (!acc.startedAt || line.timestamp < acc.startedAt) acc.startedAt = line.timestamp; + if (!acc.lastActive || line.timestamp > acc.lastActive) acc.lastActive = line.timestamp; + } + if (line.type === "user") { + if (line.cwd) acc.cwd = line.cwd; + if (line.gitBranch) acc.gitBranch = line.gitBranch; + } + if (line.type === "assistant") { + acc.turns += 1; + const u = line.message?.usage ?? {}; + acc.inputTokens += u.input_tokens ?? 0; + acc.cacheCreationTokens += u.cache_creation_input_tokens ?? 0; + acc.cacheReadTokens += u.cache_read_input_tokens ?? 0; + acc.outputTokens += u.output_tokens ?? 0; + acc.costUSD += line.costUSD ?? 0; + if (!acc.model && line.message?.model) acc.model = line.message.model; + } + return acc; +} + +function parse(raw: string): Line | null { + const trimmed = raw.trim(); + if (!trimmed) return null; + try { + return JSON.parse(trimmed) as Line; + } catch { + // A malformed line is skipped, never fatal: a transcript is an append-only log the CLI owns, and + // a half-written last line is ordinary. + return null; + } +} + +/** Sum a whole transcript string into one session's usage — the tested core of the streaming read. */ +export function foldTranscript(sessionId: string, text: string): SessionUsage { + return text.split("\n").reduce(foldLine, empty(sessionId)); +} + +/** Stream one `.jsonl` file line by line, so a large transcript never loads whole. */ +export async function readSessionFile(path: string, sessionId: string): Promise { + const acc = empty(sessionId); + const rl = createInterface({ input: createReadStream(path), crlfDelay: Infinity }); + for await (const line of rl) foldLine(acc, line); + return acc; +} diff --git a/modules/anthropic-consumer/tsconfig.json b/modules/anthropic-consumer/tsconfig.json new file mode 100644 index 0000000..ab08067 --- /dev/null +++ b/modules/anthropic-consumer/tsconfig.json @@ -0,0 +1,18 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": [ + "credentials.ts", + "transcript.ts", + "identity.ts", + "apply/index.ts", + "usage/index.ts" + ] +} diff --git a/modules/anthropic-consumer/usage/index.ts b/modules/anthropic-consumer/usage/index.ts new file mode 100644 index 0000000..7f31b11 --- /dev/null +++ b/modules/anthropic-consumer/usage/index.ts @@ -0,0 +1,92 @@ +// 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 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 { join, dirname } from "node:path"; + +import { readSessionFile, type SessionUsage } from "../transcript.js"; + +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}`); + } + } + + for (const r of readings) { + await emitUsage({ grain: "session", node, module, ...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): Promise { + const main = process.env.MESH_TOOLS_MAIN ?? "/app/dist/main.js"; + const { spawn } = await import("node:child_process"); + await new Promise((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(); diff --git a/modules/anthropic-manager/adopt/index.ts b/modules/anthropic-manager/adopt/index.ts new file mode 100644 index 0000000..794cd7f --- /dev/null +++ b/modules/anthropic-manager/adopt/index.ts @@ -0,0 +1,39 @@ +// Adoption: the ONE time an operator's refresh token enters the mesh, and it enters already sealed. +// +// The refresh token is read here, on the MANAGER NODE, sealed at rest to that node's own key, and +// only the sealed envelope leaves this process (novox/hq ADR 0050, Phase C). The control plane stores +// that envelope via `licence set-grant` without ever seeing the refresh token in the clear — the same +// bound every refresh keeps. This is the counterpart to `refresh/index.js`: adoption seals the first +// envelope, refresh opens and re-seals it. +// +// MESH_ANTHROPIC_REFRESH_TOKEN_FILE the operator's refresh token, read once and never written out +// MESH_NODE_SEALING_PUBLIC_FILE the manager node's public sealing key (base64 raw X25519) +// MESH_ANTHROPIC_GRANT_OUT where the sealed envelope is written, for `licence set-grant` + +import { readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; +import { dirname } from "node:path"; + +import { sealAtRest } from "../atrest.js"; + +function required(name: string): string { + const v = process.env[name]; + if (!v) throw new Error(`${name} is not set — adoption needs it`); + return v; +} + +const refreshToken = readFileSync(required("MESH_ANTHROPIC_REFRESH_TOKEN_FILE"), "utf8").trim(); +if (!refreshToken) throw new Error("[anthropic-manager] there is no refresh token to adopt"); + +const nodePub = readFileSync(required("MESH_NODE_SEALING_PUBLIC_FILE"), "utf8").trim(); +const envelope = sealAtRest(refreshToken, nodePub); + +const out = required("MESH_ANTHROPIC_GRANT_OUT"); +mkdirSync(dirname(out), { recursive: true }); +const tmp = `${out}.tmp`; +writeFileSync( + tmp, + JSON.stringify({ token: envelope.token, wrapped_key: envelope.wrappedKey, manager_key: envelope.managerKey }), + { mode: 0o600 }, +); +renameSync(tmp, out); +console.error("[anthropic-manager] sealed the refresh token at rest; only this node's key opens it"); diff --git a/modules/anthropic-manager/atrest.ts b/modules/anthropic-manager/atrest.ts new file mode 100644 index 0000000..759ab6c --- /dev/null +++ b/modules/anthropic-manager/atrest.ts @@ -0,0 +1,174 @@ +// The refresh token, encrypted at rest so ONE node — the manager — can read it back, and nothing +// else can: not the control plane, not a copy of its database, not another node. +// +// **Why this file exists at all.** novox/hq ADR 0050 draws one bounded carve-out in the mesh's "the +// control plane cannot read what it stores" guarantee: a refreshable-grant credential (Anthropic's +// subscription OAuth) must be rotated centrally, and rotating it means SOME node reads the refresh +// token back, every cycle. The ADR names exactly one such node — the *manager* — and this is the +// mechanism by which it, and only it, reads that token. mesh-control (the control plane) holds the +// output of this as three opaque strings and never runs the open: it has no key that could. +// +// **The construction (ECIES over the node's own sealing key).** Envelope encryption: +// - a fresh random 32-byte data key encrypts the refresh token with AES-256-GCM (`token`); +// - that data key is wrapped to the manager node's X25519 sealing key — the same key pair the +// host already holds for the node — via an ephemeral-static ECDH → HKDF-SHA256 → AES-256-GCM +// (`wrappedKey`, carrying the ephemeral public key in front); +// - `managerKey` is the node public key the data key was wrapped to, kept so a node that has since +// rotated its key learns it can no longer open this, rather than discovering it as a decrypt +// that fails. +// Recovering the refresh token needs the node's X25519 *private* half, which never leaves that +// machine. A copy of the control plane's database is a directory of ciphertexts and wrapped keys +// with nothing to open either. +// +// **On format.** This is the manager module's own at-rest format, distinct from mesh-control's Go +// `secrets.AtRest` (which is NaCl secretbox + sealed box). That is deliberate and safe: on the +// Anthropic path the manager module is the ONLY component that seals or opens the envelope — it +// seals at adoption, it opens and re-seals every refresh — and mesh-control stores the three parts +// as opaque strings it never interprets. The two never have to agree byte-for-byte because the +// bytes never cross the language boundary in an opened form. (If an operator-facing adopt path in +// mesh-control ever needed to produce the first envelope, the two would have to be unified — a NaCl +// port in TS, or a Go manager runtime. Flagged, not silently assumed.) + +import { + createCipheriv, + createDecipheriv, + createPrivateKey, + createPublicKey, + diffieHellman, + generateKeyPairSync, + hkdfSync, + randomBytes, + type KeyObject, +} from "node:crypto"; + +/** The three opaque parts mesh-control stores and forwards, and nothing else. */ +export interface Envelope { + /** base64( iv ‖ tag ‖ AES-256-GCM(dataKey, refreshToken) ). */ + readonly token: string; + /** base64( ephemeralPub(32) ‖ iv ‖ tag ‖ AES-256-GCM(kek, dataKey) ). */ + readonly wrappedKey: string; + /** base64 of the node's raw 32-byte X25519 public key the data key was wrapped to. */ + readonly managerKey: string; +} + +const INFO = Buffer.from("mesh-atrest-v1"); +const IV_LEN = 12; +const TAG_LEN = 16; +const RAW_KEY_LEN = 32; + +/** Import a node's raw 32-byte X25519 public key (standard base64, as the mesh records it). */ +function importPublic(rawBase64: string): KeyObject { + const raw = Buffer.from(rawBase64, "base64"); + if (raw.length !== RAW_KEY_LEN) { + throw new Error(`a sealing public key is 32 bytes, not ${raw.length}`); + } + return createPublicKey({ + key: { kty: "OKP", crv: "X25519", x: raw.toString("base64url") }, + format: "jwk", + }); +} + +/** Import a node's raw 32-byte X25519 private key together with its public half. */ +function importPrivate(rawPrivB64: string, rawPubB64: string): KeyObject { + const priv = Buffer.from(rawPrivB64, "base64"); + const pub = Buffer.from(rawPubB64, "base64"); + if (priv.length !== RAW_KEY_LEN) { + throw new Error(`a sealing private key is 32 bytes, not ${priv.length}`); + } + return createPrivateKey({ + key: { kty: "OKP", crv: "X25519", x: pub.toString("base64url"), d: priv.toString("base64url") }, + format: "jwk", + }); +} + +/** The raw 32-byte public key of an X25519 KeyObject. */ +function rawPublic(key: KeyObject): Buffer { + const jwk = key.export({ format: "jwk" }) as { x?: string }; + if (!jwk.x) throw new Error("a public key had no point"); + return Buffer.from(jwk.x, "base64url"); +} + +/** Bind the wrapping key to both the ephemeral and the recipient public key, as a sealed box does. */ +function deriveKek(shared: Buffer, ephemeralPub: Buffer, recipientPub: Buffer): Buffer { + const salt = Buffer.concat([ephemeralPub, recipientPub]); + return Buffer.from(hkdfSync("sha256", shared, salt, INFO, 32)); +} + +/** + * Seal a refresh token so only the holder of managerPublicKey's private half can read it. + * + * A fresh data key and ephemeral key each time, so two envelopes of the same token look nothing + * alike — a rotation that changed nothing is indistinguishable from one that changed everything. + */ +export function sealAtRest(refreshToken: string, managerPublicKeyB64: string): Envelope { + if (!refreshToken) throw new Error("there is nothing to seal"); + const recipient = importPublic(managerPublicKeyB64); + const recipientPub = rawPublic(recipient); + + // Ephemeral-static ECDH: a throwaway key pair whose public half rides in the envelope. + const eph = generateKeyPairSync("x25519"); + const ephemeralPub = rawPublic(eph.publicKey); + const shared = diffieHellman({ privateKey: eph.privateKey, publicKey: recipient }); + const kek = deriveKek(shared, ephemeralPub, recipientPub); + + const dataKey = randomBytes(32); + + const wrappedKey = Buffer.concat([ephemeralPub, aesSeal(kek, dataKey)]); + const token = aesSeal(dataKey, Buffer.from(refreshToken, "utf8")); + + return { + token: token.toString("base64"), + wrappedKey: wrappedKey.toString("base64"), + managerKey: managerPublicKeyB64, + }; +} + +/** + * Recover the refresh token, given the manager node's own key pair. This is the one place a refresh + * token is in the clear, and it runs only on the manager node. + */ +export function openAtRest(env: Envelope, managerPublicKeyB64: string, managerPrivateKeyB64: string): string { + const priv = importPrivate(managerPrivateKeyB64, managerPublicKeyB64); + const recipientPub = Buffer.from(managerPublicKeyB64, "base64"); + + const wrapped = Buffer.from(env.wrappedKey, "base64"); + if (wrapped.length < RAW_KEY_LEN + IV_LEN + TAG_LEN) { + throw new Error("the wrapped key is too short to hold what it must"); + } + const ephemeralPub = wrapped.subarray(0, RAW_KEY_LEN); + const wrappedRest = wrapped.subarray(RAW_KEY_LEN); + + const ephemeralKey = importPublic(ephemeralPub.toString("base64")); + const shared = diffieHellman({ privateKey: priv, publicKey: ephemeralKey }); + const kek = deriveKek(shared, ephemeralPub, recipientPub); + + let dataKey: Buffer; + try { + dataKey = aesOpen(kek, wrappedRest); + } catch { + throw new Error("this refresh token was not wrapped to this manager's key"); + } + if (dataKey.length !== 32) throw new Error("the wrapped data key is the wrong length"); + + const refresh = aesOpen(dataKey, Buffer.from(env.token, "base64")); + return refresh.toString("utf8"); +} + +// --- AES-256-GCM helpers: output/consume iv ‖ tag ‖ ciphertext --- + +function aesSeal(key: Buffer, plaintext: Buffer): Buffer { + const iv = randomBytes(IV_LEN); + const cipher = createCipheriv("aes-256-gcm", key, iv); + const ct = Buffer.concat([cipher.update(plaintext), cipher.final()]); + const tag = cipher.getAuthTag(); + return Buffer.concat([iv, tag, ct]); +} + +function aesOpen(key: Buffer, blob: Buffer): Buffer { + const iv = blob.subarray(0, IV_LEN); + const tag = blob.subarray(IV_LEN, IV_LEN + TAG_LEN); + const ct = blob.subarray(IV_LEN + TAG_LEN); + const decipher = createDecipheriv("aes-256-gcm", key, iv); + decipher.setAuthTag(tag); + return Buffer.concat([decipher.update(ct), decipher.final()]); +} diff --git a/modules/anthropic-manager/client.ts b/modules/anthropic-manager/client.ts new file mode 100644 index 0000000..5b37387 --- /dev/null +++ b/modules/anthropic-manager/client.ts @@ -0,0 +1,135 @@ +// The only file that talks to Anthropic — the vendor half of the refreshable-grant adapter +// (novox/hq ADR 0050). Isolated exactly as cloudflare-dns isolates its registrar call, so the +// vendor is swappable and the one place a token endpoint is reached is auditable. +// +// Two endpoints, and they are different hosts (port-map "don't-map" #1): the TOKEN host mints a new +// access token from the refresh token; the USAGE host reports utilisation against an access token. + +/** The token endpoint, overridable so the lab can point the whole flow at a stub without a vendor. */ +export function tokenEndpoint(env = process.env): string { + return env.MESH_ANTHROPIC_TOKEN_ENDPOINT ?? "https://platform.claude.com/v1/oauth/token"; +} + +/** The usage endpoint, likewise overridable for the lab. */ +export function usageEndpoint(env = process.env): string { + return env.MESH_ANTHROPIC_USAGE_ENDPOINT ?? "https://api.anthropic.com/api/oauth/usage"; +} + +// The OAuth client id is a hard-won constant, ported byte-exact from the mature implementation: a +// metadata URL in its place yields 400. It is not a secret (it identifies the public Claude Code +// client), so it lives in code. +const CLIENT_ID = "9d1c250a-e61b-44d9-88ed-5944d1962f5e"; + +/** The vendor's token response, snake_case as the wire has it. */ +export interface RefreshedGrant { + readonly access_token?: string; + readonly refresh_token?: string; + readonly expires_in?: number; + readonly refresh_token_expires_in?: number; + readonly scopes?: string[]; + readonly subscription_type?: string; +} + +/** + * Exchange a refresh token for a fresh grant. Returns null on any non-ok response, surfacing the + * OAuth error body (invalid_grant/invalid_client/…) — the difference between "the token is dead" and + * "the endpoint was unreachable", which a bare status hides. + */ +export async function refreshGrant( + refreshToken: string, + env = process.env, +): Promise { + const resp = await fetch(tokenEndpoint(env), { + method: "POST", + headers: { "content-type": "application/x-www-form-urlencoded" }, + body: new URLSearchParams({ + grant_type: "refresh_token", + refresh_token: refreshToken, + client_id: CLIENT_ID, + }), + }); + if (!resp.ok) { + const body = await resp.text().catch(() => ""); + console.error( + `[anthropic-manager] token refresh failed: ${resp.status} ${resp.statusText} — ${body.slice(0, 400)}`, + ); + return null; + } + return (await resp.json()) as RefreshedGrant; +} + +/** The vendor's usage response — utilisation percentages against several windows. */ +export interface UsageLimits { + readonly five_hour?: { utilization: number; resets_at?: string }; + readonly seven_day?: { utilization: number; resets_at?: string }; + readonly seven_day_sonnet?: { utilization: number; resets_at?: string }; + readonly seven_day_opus?: { utilization: number; resets_at?: string }; + readonly extra_usage?: { utilization: number }; + readonly [key: string]: unknown; +} + +/** + * Read utilisation for an access token. Never refreshes here (a 401 is just reported): a second + * refresh source racing the first is the fault the mature implementation warns against. + */ +export async function readUsage(accessToken: string, env = process.env): Promise { + const resp = await fetch(usageEndpoint(env), { + headers: { authorization: `Bearer ${accessToken}` }, + }); + if (!resp.ok) { + const body = await resp.text().catch(() => ""); + console.error(`[anthropic-manager] usage endpoint returned ${resp.status}: ${body.slice(0, 200)}`); + return null; + } + return (await resp.json()) as UsageLimits; +} + +/** The licence-grain reading ADR 0054 fixes, flattened from the vendor's windows. */ +export interface UsageReading { + readonly sessionPct: number | null; + readonly sessionResetsAt: string | null; + readonly weeklyPct: number | null; + readonly sonnetPct: number | null; + readonly extraPct: number | null; + readonly raw: UsageLimits; +} + +export function flattenUsage(u: UsageLimits): UsageReading { + return { + sessionPct: u.five_hour?.utilization ?? null, + sessionResetsAt: u.five_hour?.resets_at ?? null, + weeklyPct: u.seven_day?.utilization ?? null, + sonnetPct: u.seven_day_sonnet?.utilization ?? null, + extraPct: u.extra_usage?.utilization ?? null, + raw: u, + }; +} + +/** The access-token-only grant a holder is delivered — the port-map credential-file shape's fields. */ +export interface AccessGrant { + readonly accessToken: string; + readonly expiresAt: number | null; + readonly refreshTokenExpiresAt: number | null; + readonly scopes: string[] | null; + readonly subscriptionType: string | null; +} + +/** + * Turn a vendor refresh into what the manager submits: the access-token-only grant for holders, and + * the rotated refresh token if the vendor sent one. Never clobbers a good grant from an empty + * response — no access_token means the caller keeps what it had. + */ +export function grantFromRefresh(r: RefreshedGrant, nowMs: number): { access: AccessGrant; rotatedRefresh: string | null } | null { + if (!r.access_token) return null; + return { + access: { + accessToken: r.access_token, + expiresAt: typeof r.expires_in === "number" ? nowMs + r.expires_in * 1000 : null, + refreshTokenExpiresAt: + typeof r.refresh_token_expires_in === "number" ? nowMs + r.refresh_token_expires_in * 1000 : null, + scopes: r.scopes ?? null, + subscriptionType: r.subscription_type ?? null, + }, + rotatedRefresh: r.refresh_token ?? null, + }; +} diff --git a/modules/anthropic-manager/module.json b/modules/anthropic-manager/module.json new file mode 100644 index 0000000..b33c78c --- /dev/null +++ b/modules/anthropic-manager/module.json @@ -0,0 +1,68 @@ +{ + "module": "anthropic-manager", + "version": "1", + "capabilities": [ + "container-runtime" + ], + "own-secrets": { + "broker": "/var/lib/mesh/anthropic-manager/broker" + }, + "emits": [ + "module.anthropic-manager.usage.read" + ], + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/anthropic-manager", + "mode": "0700" + }, + { + "id": "keys", + "type": "directory", + "path": "/var/lib/mesh/anthropic-manager/keys", + "mode": "0700" + }, + { + "id": "out", + "type": "directory", + "path": "/var/lib/mesh/anthropic-manager/out", + "mode": "0700" + }, + { + "id": "config", + "type": "file", + "path": "/var/lib/mesh/anthropic-manager/config.json", + "merge": "json", + "content": "{}", + "mode": "0600" + }, + { + "id": "refresh", + "type": "container", + "name": "mesh-anthropic-manager-refresh", + "image": "mesh-runtime-anthropic-manager@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "host", + "schedule": "*/5 * * * *", + "args": [ + "run", + "/app/modules/anthropic-manager/dist/refresh/index.js" + ], + "volumes": [ + "/var/lib/mesh/anthropic-manager/broker:/run/secrets/broker:ro", + "/var/lib/mesh/anthropic-manager:/run/state" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_ANTHROPIC_LICENCE": "personal", + "MESH_ANTHROPIC_GRANT_FILE": "/run/state/grant.json", + "MESH_NODE_SEALING_PUBLIC_FILE": "/run/state/keys/sealing.pub", + "MESH_NODE_SEALING_PRIVATE_FILE": "/run/state/keys/sealing.priv", + "MESH_ANTHROPIC_ACCESS_OUT": "/run/state/out/access-token", + "MESH_ANTHROPIC_GRANT_OUT": "/run/state/out/grant.json", + "MESH_ANTHROPIC_USAGE_OUT": "/run/state/out/usage.json", + "MESH_TOOLS_MAIN": "/app/dist/main.js" + } + } + ] +} diff --git a/modules/anthropic-manager/package.json b/modules/anthropic-manager/package.json new file mode 100644 index 0000000..90d518a --- /dev/null +++ b/modules/anthropic-manager/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-anthropic-manager", + "version": "0.1.0", + "description": "anthropic-manager — the manager side of the model-access refreshable-grant (ADR 0050): opens the refresh token on the manager node alone, refreshes it against Anthropic's OAuth endpoint, and submits back only the access token and the re-sealed refresh envelope.", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/anthropic-manager/refresh/index.ts b/modules/anthropic-manager/refresh/index.ts new file mode 100644 index 0000000..245057f --- /dev/null +++ b/modules/anthropic-manager/refresh/index.ts @@ -0,0 +1,142 @@ +// The manager's scheduled run (novox/hq ADR 0050/0053). It is the whole of the carve-out in one +// place, and it runs on the MANAGER NODE, never in the control plane: +// +// 1. read the opaque refresh-token envelope the control plane forwarded (it cannot open it); +// 2. open it HERE with the node's own sealing key — the one moment a refresh token is in the clear, +// on the one node the ADR permits it; +// 3. call the vendor's OAuth token endpoint to mint a fresh access token (and maybe a rotated +// refresh token); +// 4. re-seal the rotated refresh token at rest (still openable by this node alone); +// 5. hand the control plane back ONLY the access token in the clear + the opaque re-sealed +// envelope — never the refresh token — which it seals per holder and stores; +// 6. poll usage with the fresh access token and record the licence-grain reading. +// +// mesh-control receives the products of steps 5–6 through `licence submit-refresh` (access token + +// opaque envelope). The refresh token never leaves this process except as ciphertext. +// +// This runs as `mesh-tools run`, which connects no broker, so the outputs are written to files the +// host mounts; the submit itself (the transport to mesh-control) is done by the caller invoking +// `mesh-control licence submit-refresh`. In the lab that caller is the scenario; in production it is +// an authenticated call the manager node makes. The transport is the one part stubbed here — FLAGGED +// — because a cross-node authenticated command surface is out of this module's scope. + +import { readFileSync, writeFileSync, renameSync, mkdirSync } from "node:fs"; +import { dirname } from "node:path"; + +import { openAtRest, sealAtRest, type Envelope } from "../atrest.js"; +import { refreshGrant, grantFromRefresh, readUsage, flattenUsage } from "../client.js"; + +function required(name: string): string { + const v = process.env[name]; + if (!v) throw new Error(`${name} is not set — the manager runtime was deployed without it`); + return v; +} + +function readTrimmed(path: string): string { + return readFileSync(path, "utf8").trim(); +} + +/** Accept an envelope in either the wire (snake_case) or internal (camelCase) shape. */ +function readEnvelope(path: string): Envelope { + const raw = JSON.parse(readFileSync(path, "utf8")) as Record; + const token = raw.token ?? ""; + const wrappedKey = raw.wrappedKey ?? raw.wrapped_key ?? ""; + const managerKey = raw.managerKey ?? raw.manager_key ?? ""; + if (!token || !wrappedKey || !managerKey) { + throw new Error("the refresh-token envelope is missing one of token/wrapped_key/manager_key"); + } + return { token, wrappedKey, managerKey }; +} + +/** Write the envelope in the wire (snake_case) shape mesh-control's `submit-refresh` reads. */ +function writeEnvelope(path: string, env: Envelope): void { + atomicWrite( + path, + JSON.stringify({ token: env.token, wrapped_key: env.wrappedKey, manager_key: env.managerKey }), + ); +} + +function atomicWrite(path: string, content: string): void { + mkdirSync(dirname(path), { recursive: true }); + const tmp = `${path}.tmp`; + writeFileSync(tmp, content, { mode: 0o600 }); + renameSync(tmp, path); +} + +async function main(): Promise { + const licence = process.env.MESH_ANTHROPIC_LICENCE ?? "unknown"; + + const envelope = readEnvelope(required("MESH_ANTHROPIC_GRANT_FILE")); + const nodePub = readTrimmed(required("MESH_NODE_SEALING_PUBLIC_FILE")); + const nodePriv = readTrimmed(required("MESH_NODE_SEALING_PRIVATE_FILE")); + + // Step 2: the one open, on the manager node. + const refreshToken = openAtRest(envelope, nodePub, nodePriv); + + // Step 3: the vendor call. + const refreshed = await refreshGrant(refreshToken); + if (!refreshed) { + // A dead endpoint or a rejected token: nothing to publish, and we do not clobber a good grant. + throw new Error(`[anthropic-manager] the refresh of ${licence} produced no grant`); + } + const grant = grantFromRefresh(refreshed, Date.now()); + if (!grant) { + throw new Error(`[anthropic-manager] the refresh of ${licence} returned no access token`); + } + + // Step 4: re-seal the rotated refresh token, if the vendor rotated it. Nothing to store otherwise. + if (grant.rotatedRefresh) { + const rotated = sealAtRest(grant.rotatedRefresh, nodePub); + if (process.env.MESH_ANTHROPIC_GRANT_OUT) { + writeEnvelope(process.env.MESH_ANTHROPIC_GRANT_OUT, rotated); + } + } + + // Step 5: the access token in the clear, for the control plane to seal per holder. This is all it + // ever receives that is not ciphertext. + atomicWrite(required("MESH_ANTHROPIC_ACCESS_OUT"), grant.access.accessToken); + + console.error( + `[anthropic-manager] refreshed ${licence}: access token minted` + + (grant.rotatedRefresh ? ", refresh token rotated and re-sealed" : ", refresh token unchanged"), + ); + + // Step 6: licence-grain usage, best-effort — a usage read failing must not fail the refresh. + try { + const usage = await readUsage(grant.access.accessToken); + if (usage) { + const reading = flattenUsage(usage); + if (process.env.MESH_ANTHROPIC_USAGE_OUT) { + atomicWrite( + process.env.MESH_ANTHROPIC_USAGE_OUT, + JSON.stringify({ licence, grain: "licence", ...reading }), + ); + } + await emitUsage({ licence, grain: "licence", ...reading }); + } + } catch (err) { + console.error(`[anthropic-manager] usage poll for ${licence} failed: ${err}`); + } +} + +/** + * Emit a usage event best-effort by shelling out to the sibling mesh-tools `emit` primitive, which + * is the one path that wires a broker from a run-once/scheduled step (which itself connects none). + * A broker hiccup must never fail a refresh that already happened. + */ +async function emitUsage(body: Record): Promise { + const main = process.env.MESH_TOOLS_MAIN ?? "/app/dist/main.js"; + const { spawn } = await import("node:child_process"); + await new Promise((resolve) => { + const child = spawn(process.execPath, [main, "emit", "module.anthropic-manager.usage.read", JSON.stringify(body)], { + stdio: "inherit", + }); + child.on("exit", () => resolve()); + child.on("error", (err) => { + console.error(`[anthropic-manager] could not emit usage: ${err}`); + resolve(); + }); + }); +} + +await main(); diff --git a/modules/anthropic-manager/test/atrest.test.ts b/modules/anthropic-manager/test/atrest.test.ts new file mode 100644 index 0000000..595bf63 --- /dev/null +++ b/modules/anthropic-manager/test/atrest.test.ts @@ -0,0 +1,47 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { generateKeyPairSync } from "node:crypto"; + +import { sealAtRest, openAtRest, type Envelope } from "../atrest.ts"; + +/** A node key pair as the mesh records it: raw 32-byte X25519 keys, standard base64. */ +function nodeKeys(): { pub: string; priv: string } { + const kp = generateKeyPairSync("x25519"); + const pub = (kp.publicKey.export({ format: "jwk" }) as { x: string }).x; + const priv = (kp.privateKey.export({ format: "jwk" }) as { d: string }).d; + // JWK is base64url; the mesh records standard base64 of the same 32 bytes. + const std = (b64url: string) => Buffer.from(b64url, "base64url").toString("base64"); + return { pub: std(pub), priv: std(priv) }; +} + +test("the manager seals a refresh token and reads it back with its own key", () => { + const { pub, priv } = nodeKeys(); + const env = sealAtRest("rt-the-refresh-token", pub); + assert.equal(env.managerKey, pub); + // Nothing in the envelope is the refresh token in the clear. + assert.doesNotMatch(env.token, /rt-the-refresh-token/); + assert.doesNotMatch(env.wrappedKey, /rt-the-refresh-token/); + assert.equal(openAtRest(env, pub, priv), "rt-the-refresh-token"); +}); + +test("a node that is not the manager cannot open the envelope", () => { + const manager = nodeKeys(); + const other = nodeKeys(); + const env = sealAtRest("rt-secret", manager.pub); + assert.throws(() => openAtRest(env, other.pub, other.priv)); +}); + +test("two seals of the same token look nothing alike", () => { + const { pub } = nodeKeys(); + const a = sealAtRest("rt-secret", pub); + const b = sealAtRest("rt-secret", pub); + assert.notEqual(a.token, b.token); + assert.notEqual(a.wrappedKey, b.wrappedKey); +}); + +test("a tampered envelope is refused, not silently mis-opened", () => { + const { pub, priv } = nodeKeys(); + const env = sealAtRest("rt-secret", pub); + const flipped: Envelope = { ...env, token: Buffer.from(env.token, "base64").reverse().toString("base64") }; + assert.throws(() => openAtRest(flipped, pub, priv)); +}); diff --git a/modules/anthropic-manager/test/client.test.ts b/modules/anthropic-manager/test/client.test.ts new file mode 100644 index 0000000..a28cd5f --- /dev/null +++ b/modules/anthropic-manager/test/client.test.ts @@ -0,0 +1,43 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { grantFromRefresh, flattenUsage } from "../client.ts"; + +test("an empty refresh response never clobbers a good grant", () => { + assert.equal(grantFromRefresh({}, 1000), null); +}); + +test("a refresh with an access token yields an access-token-only grant and epoch expiry", () => { + const out = grantFromRefresh( + { access_token: "at-new", expires_in: 3600, refresh_token: "rt-rotated", subscription_type: "pro" }, + 1_000_000, + ); + assert.ok(out); + assert.equal(out!.access.accessToken, "at-new"); + assert.equal(out!.access.expiresAt, 1_000_000 + 3600 * 1000); + assert.equal(out!.access.subscriptionType, "pro"); + // The rotated refresh token is reported separately, for the manager to re-seal — never put in the + // holder grant. + assert.equal(out!.rotatedRefresh, "rt-rotated"); + assert.ok(!("refreshToken" in (out!.access as object))); +}); + +test("a refresh that did not rotate the refresh token reports none to re-seal", () => { + const out = grantFromRefresh({ access_token: "at-new" }, 0); + assert.ok(out); + assert.equal(out!.rotatedRefresh, null); +}); + +test("usage flattens the vendor windows to the ADR 0054 grain", () => { + const r = flattenUsage({ + five_hour: { utilization: 42, resets_at: "2026-01-01T00:00:00Z" }, + seven_day: { utilization: 10 }, + seven_day_sonnet: { utilization: 5 }, + extra_usage: { utilization: 1 }, + }); + assert.equal(r.sessionPct, 42); + assert.equal(r.sessionResetsAt, "2026-01-01T00:00:00Z"); + assert.equal(r.weeklyPct, 10); + assert.equal(r.sonnetPct, 5); + assert.equal(r.extraPct, 1); +}); diff --git a/modules/anthropic-manager/tsconfig.json b/modules/anthropic-manager/tsconfig.json new file mode 100644 index 0000000..3a6da13 --- /dev/null +++ b/modules/anthropic-manager/tsconfig.json @@ -0,0 +1,17 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": [ + "atrest.ts", + "client.ts", + "adopt/index.ts", + "refresh/index.ts" + ] +}