/** * The usage context store, proved end to end (novox/hq ADR 0054). * * model-usage is a MODULE, not a control-plane feature, because mesh-controller is a CLI and cannot * consume events: the store that keeps the latest usage reading has to be something that subscribes * to `module.*.usage.*` and upserts. This bed raises a real mesh, provisions model-usage its own * postgres store, emits usage events into the mesh, and reads the store back to prove the three * properties the ADR fixes: * * 1. UPSERT-LATEST — a second reading on the same (licence, consumer, period, metric) REPLACES the * first; the store keeps one row, the latest. At-least-once delivery makes the upsert idempotent. * 2. TWO GRAINS — the licence grain and the session grain are ONE table, differing only in * `consumer`; both are stored and both are selectable, told apart by `consumer`. * 3. IN THE CLEAR — the value and its raw payload are ordinary columns an ordinary select returns; * nothing about usage is sealed. * * The topology mirrors two-node-db: the app-postgres provider and the mesh's own foundation store both * want host port 5432, so the foundation (store, broker, control) owns `anchor` and NOTHING else, and * `laptop` runs the postgres PROVIDER and the model-usage CONSUMER co-located. Only enrolment crosses * to anchor, over the underlay both machines share. * * It needs a host binary and the foundation bundle: * * MESH_LAB_HOST_BINARY=.../mesh-host * MESH_LAB_BUNDLE=.../examples/foundation-first-node.lock * * HELPER — stock the two runtimes into the local daemon before the run (some may already be there): * scripts/build-module-runtime.sh postgres /tmp/postgres.tar * scripts/build-module-runtime.sh model-usage /tmp/model-usage.tar * The service image (postgres:17-alpine) must be in the local daemon too; scenarios/model-usage-bed.yml * stocks it, and each node pulls what it runs from the scenario's own registry by digest. */ import { test, before, after } from "node:test"; import assert from "node:assert/strict"; import { existsSync, readFileSync } from "node:fs"; import { loadScenario } from "../../src/declaration/parse.ts"; import { raise } from "../../src/lifecycle/raise.ts"; import { destroy, exec } from "../../src/lifecycle/operate.ts"; import { hostBinaryPath, HOST_PATH } from "../../src/lifecycle/place.ts"; import { labIsUsable, destroyAll, foundationBundle, onTheMachine } from "./harness.ts"; import type { HeldImage } from "../../src/pinning.ts"; const capability = await labIsUsable(); const binary = hostBinaryPath(); const bundle = process.env["MESH_LAB_BUNDLE"] ?? ""; const skip = !capability.usable ? `lab not usable: ${capability.why}` : !binary || !existsSync(binary) ? "MESH_LAB_HOST_BINARY is not set to a built mesh-host" : !bundle || !existsSync(bundle) ? "MESH_LAB_BUNDLE is not set to a foundation bundle (mesh-host examples/)" : false; const SCENARIO = "model-usage-bed"; /** The node that carries the postgres provider and the model-usage consumer. anchor carries only the * foundation. */ const NODE = "laptop"; let instanceId = ""; /** The mesh's own images, as the machines hold them. */ let held: HeldImage[] = []; function quote(s: string): string { return `'${s.replaceAll("'", `'\\''`)}'`; } async function on(machine: string, command: string, timeoutMs?: number): Promise<{ out: string; ok: boolean }> { const { stdout } = await exec(instanceId, machine, [ "sh", "-c", `exec 2>&1\n${command}\necho "__exit=$?"`, ], timeoutMs); const marker = stdout.lastIndexOf("__exit="); if (marker < 0) return { out: stdout, ok: false }; return { out: stdout.slice(0, marker), ok: stdout.slice(marker + 7).trim() === "0" }; } async function must(machine: string, command: string, timeoutMs?: number): Promise { const { out, ok } = await on(machine, command, timeoutMs); if (!ok) throw new Error(`${machine}: ${command}\n${out}`); return out; } /** The control plane, a container on the first node. */ async function mesh(command: string, timeoutMs?: number): Promise { return must("anchor", `docker exec mesh-controller /mesh-controller ${command}`, timeoutMs); } /** The reference a manifest should carry, once this scenario has been raised. */ /** What a manifest's image reference becomes on the machine — ours by ID, everything else as written. */ function pinned(reference: string): string { return onTheMachine(reference, held); } /** The foundation bundle: ours by the ID the machine holds, everything else upstream. */ function bundleFor(images: HeldImage[]): string { return foundationBundle(bundle, images); } function tokenFrom(said: string): string { const found = said.split("\n").map((l) => l.trim()).find((l) => l.length > 100 && !l.includes(" ")); assert.ok(found, `no token in:\n${said}`); return found; } /** * Wait until a node has actually applied what it was last sent. `push` sends and returns; the node * applies afterwards, so asserting immediately after a push is a race. The mesh is asked in its own * terms — settled when it is neither waiting for what it was sent nor wrong about what it applied. */ async function settled(node: string, withinMs = 1_200_000): Promise { const until = Date.now() + withinMs; let last = ""; while (Date.now() < until) { let state: { wrong: { node: string; outcome: string; refused?: string; failed?: { id: string; error: string }[] }[]; waiting: { node: string; never: boolean }[]; reported: { node: string; outcome: string; current: boolean }[]; } | undefined; let said = ""; try { const asked = await on("anchor", `docker exec mesh-controller /mesh-controller status --json`); said = asked.out; if (asked.ok) state = JSON.parse(said); } catch (err) { said = (err as Error).message; } if (!state) { last = said; await new Promise((r) => setTimeout(r, 5000)); continue; } const bad = state.wrong.find((w) => w.node === node); if (bad) { const why = [bad.refused, ...(bad.failed ?? []).map((f) => `${f.id}: ${f.error}`)] .filter(Boolean).join("\n "); throw new Error(`${node} did not apply what it was sent (${bad.outcome}):\n ${why}`); } const word = state.reported.find((r) => r.node === node); const acted = word?.outcome === "applied" && word.current; if (!state.waiting.some((w) => w.node === node) && acted) return; last = said; await new Promise((r) => setTimeout(r, 5000)); } const ps = (await on(NODE, `docker ps -a --format '{{.Names}}\t{{.Status}}'`)).out; const hostLog = (await on(NODE, `tail -80 /var/log/mesh-host.log`)).out; throw new Error( `${node} never caught up within ${Math.round(withinMs / 1000)}s.\nLast status:\n${last}\n` + `--- ${NODE} docker ps -a ---\n${ps}\n--- ${NODE} mesh-host.log tail ---\n${hostLog}`); } before(async () => { if (skip) return; const raised = await raise(loadScenario(`scenarios/${SCENARIO}.yml`), { onProgress: (m) => console.log(`raise: ${m}`), }); instanceId = raised.instanceId; held = raised.images; await must("anchor", `cat > /tmp/foundation.lock <<'MESHBUNDLE'\n${bundleFor(raised.images)}\nMESHBUNDLE`); await must("anchor", `${HOST_PATH} apply /tmp/foundation.lock`, 600_000); const up = await must("anchor", `docker ps --format '{{.Names}}'`); for (const c of ["mesh-store", "mesh-broker", "mesh-controller"]) { assert.match(up, new RegExp(c), `the foundation did not raise ${c}:\n${up}`); } for (const [machine, node] of [["anchor", "anchor"], ["laptop", "laptop"]] as const) { await mesh(`node add ${node}`); const token = tokenFrom(await mesh(`token issue --node ${node}`)); const said = await must(machine, `${HOST_PATH} enrol --token ${quote(token)}`); assert.match(said, new RegExp(`enrolled as ${node}`), said); await must(machine, `nohup ${HOST_PATH} run > /var/log/mesh-host.log 2>&1 & sleep 3`); } }, { timeout: 1_800_000 }); after(async () => { if (instanceId) await destroy(instanceId); await destroyAll(`${SCENARIO}-`); }, { timeout: 600_000 }); test("usage events are upserted into model-usage's store — latest-per-key, both grains, in the clear", { skip, timeout: 1_500_000, }, async () => { // ================================================================================================ // THE MANIFESTS — postgres verbatim from two-node-db (its server publishes 5432 so the consumer // reaches it), and model-usage the committed catalogue shape with its images pinned. // ================================================================================================ const postgresManifest = JSON.stringify({ module: "postgres", version: "1", provides: [{ name: "postgres-database", scope: "mesh" }], serves: { "postgres-database": { port: 5432 } }, emits: ["module.postgres.database.provisioned", "module.postgres.database.deprovisioned"], consumes: ["module.postgres.database.provisioned", "module.postgres.database.deprovisioned"], receives: { "postgres-database": "/var/lib/postgres/grants/mesh.json" }, grants: { "postgres-database": "/var/lib/postgres/grants" }, "own-secrets": { superuser: "/var/lib/postgres/superuser.secret", broker: "/var/lib/mesh/postgres/broker" }, resources: [ { id: "mesh-state", type: "directory", path: "/var/lib/mesh/postgres", mode: "0700" }, { id: "state", type: "directory", path: "/var/lib/postgres", mode: "0700" }, { id: "grants", type: "directory", path: "/var/lib/postgres/grants", mode: "0700" }, { id: "data", type: "directory", path: "/services/postgres/db-data", mode: "0700" }, { id: "net", type: "network", name: "postgres" }, { id: "server", type: "container", name: "postgres", image: pinned("postgres"), network: "postgres", // The superuser reaches postgres as a file (novox/hq ADR 0086), the shape the catalogue's manifest has. env: { POSTGRES_USER: "postgres", POSTGRES_DB: "postgres", POSTGRES_PASSWORD_FILE: "/run/secrets/superuser" }, ports: ["5432"], volumes: ["/services/postgres/db-data:/var/lib/postgresql/data", "/var/lib/postgres/superuser.secret:/run/secrets/superuser:ro"], }, { id: "runtime", type: "container", name: "mesh-postgres", image: pinned("mesh-runtime-postgres"), network: "postgres", volumes: [ "/var/lib/mesh/postgres/broker:/run/secrets/broker:ro", "/var/lib/postgres/grants:/var/lib/postgres/grants:ro", "/var/lib/postgres/superuser.secret:/run/secrets/superuser:ro", ], env: { MESH_BROKER_FILE: "/run/secrets/broker", MESH_RECEIVES: "/var/lib/postgres/grants/mesh.json", MESH_PROVISION_POSTGRES: "postgres://postgres@postgres:5432/postgres?sslmode=disable", MESH_PROVISION_PASSWORD_FILE: "/run/secrets/superuser", }, }, ], }); // --- model-usage: requires postgres-database, owns a provisioned store, consumes module.*.usage.*, // runs a run-once migrate then the long-lived event consumer. Both containers on the host network so // they reach the granted postgres (at the provider's address the mesh writes) and the broker. ------ const modelUsageManifest = JSON.stringify({ module: "model-usage", version: "1", // `mesh_laptop_model-usage` is 23 chars, over the 20 an S3 access key keeps (ADR 0049); a short // slug makes the consumer identity `mesh_laptop_usage` (17). db/role/`as` all derive from it. 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" }, // The connection string carries the password, so it reaches the runtime as a file the mesh // templates (novox/hq ADR 0086), the shape the catalogue's manifest has. { id: "database-url", type: "file", path: "/var/lib/model-usage/database.url", mode: "0600", content: "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: pinned("mesh-runtime-model-usage"), network: "host", volumes: [ "/var/lib/mesh/model-usage/broker:/run/secrets/broker:ro", "/var/lib/model-usage:/run/state", "/var/lib/model-usage/database.url:/run/secrets/database-url:ro", ], env: { MESH_BROKER_FILE: "/run/secrets/broker", DATABASE_URL_FILE: "/run/secrets/database-url" }, }, ], }); async function addIssueAssign(name: string, manifest: string): Promise { await must("anchor", `printf %s ${quote(manifest)} > /tmp/${name}.json && docker cp /tmp/${name}.json mesh-controller:/${name}.json`); await mesh(`module add /${name}.json`); await mesh(`module issue ${name} --node ${NODE}`); await mesh(`assign ${NODE} ${name}`); } // The consumer connects to its provider by the binding's `at` — the overlay address the mesh fills. // The overlay networking is assigned first, to both nodes, so the address is non-empty. await mesh("overlay place anchor --hub --endpoint 192.0.2.10:51820 --site lab"); await mesh(`overlay place ${NODE} --site lab`); await mesh("assign anchor networking"); await mesh(`assign ${NODE} networking`); await addIssueAssign("postgres", postgresManifest); await addIssueAssign("model-usage", modelUsageManifest); // ONE push, ONE convergence — provider and consumer resolved and applied together on the node. await mesh(`push ${NODE}`); await settled(NODE); // The provider owns 5432 on laptop; the foundation store owns it on anchor. const onLaptop = await must(NODE, `docker ps --format '{{.Names}}'`); for (const name of ["postgres", "mesh-postgres", "mesh-model-usage"]) { assert.match(onLaptop, new RegExp(`(^|\\n)${name}(\\n|$)`), `${name} is not running on ${NODE} after the push:\n${onLaptop}\n---host log---\n${(await on(NODE, `tail -60 /var/log/mesh-host.log`)).out}`); } // ================================================================================================ // The store was provisioned — model-usage got its postgres binding. Build the connection the way // two-node-db does: the login and password the mesh minted, against the provider on its own network. // ================================================================================================ const bound = await waitForBinding("/var/lib/model-usage/database.json"); assert.equal(bound.provision, "postgres-database", `model-usage was bound the wrong provision: ${bound.provision}`); const pw = (await must(NODE, `cat /var/lib/model-usage/database.secret`)).trim(); assert.ok(bound.as && pw, `model-usage's login or password was empty (as=${bound.as})`); // What the machine can say when the login below fails — asked now, so the failure carries it. const account = async () => [ "--- provisioner (mesh-postgres):", (await on(NODE, `docker logs --tail 25 mesh-postgres 2>&1`)).out, "--- runtime (mesh-model-usage):", (await on(NODE, `docker logs --tail 25 mesh-model-usage 2>&1`)).out, "--- grants:", (await on(NODE, `ls -la /var/lib/postgres/grants/; cat /var/lib/postgres/grants/mesh.json 2>&1 | head -30`)).out, "--- roles:", (await on(NODE, `docker exec postgres psql -U postgres -tAc "select rolname from pg_roles where rolname like 'mesh%'" 2>&1`)).out, "--- host log:", (await on(NODE, `tail -40 /var/log/mesh-host.log 2>&1`)).out, ].join("\n"); const conn = `postgresql://${bound.as}:${encodeURIComponent(pw)}@postgres:5432/${bound.as}?sslmode=disable`; async function usageQuery(sql: string): Promise<{ out: string; ok: boolean }> { return on(NODE, `docker exec mesh-postgres psql ${quote(conn)} -tAc ${quote(sql)} 2>&1`); } /** Poll a query until its output matches `want` (the store is written asynchronously by the * consumer, and the run-once migrate creates the table shortly after the binding lands). */ async function waitFor(sql: string, want: RegExp, withinMs = 120_000): Promise { const until = Date.now() + withinMs; let last = ""; while (Date.now() < until) { const r = await usageQuery(sql); last = r.out; if (r.ok && want.test(r.out)) return r.out; await new Promise((res) => setTimeout(res, 3000)); } return last; } // The consumer creates the `usage` table on startup (idempotent CREATE TABLE IF NOT EXISTS); wait // for it before emitting, so the consumer's upsert has a table to write (a reading that arrives // before the table would be logged and lost). const tableReady = await waitFor("select to_regclass('usage') is not null", /^t$/m); assert.match(tableReady, /^t$/m, `the usage table was never created (the consumer did not migrate):\n${tableReady}\n${await account()}`); // Inject a usage event into the mesh. model-usage is a PURE CONSUMER, so its own broker account has // no publish right (mesh-controller grants write to mesh.events only to a module that declares `emits`). // So publish from the postgres provider's container — a publisher already on the node — overriding // the source header to the real producer the key names, exactly as events.test.ts injects with // `-e MESH_MODULE=...`. Write permission is per-exchange, not per-key, so postgres may carry any // routing key; the store reads the body, not the source, so model-usage's consume→upsert path is // exercised exactly as it is when the real anthropic producer emits. async function emit(key: string, body: unknown): Promise { const source = key.split(".")[1] ?? "usage-probe"; // module..usage. await must(NODE, `docker exec -e MESH_MODULE=${source} mesh-postgres node dist/main.js emit ${key} ${quote(JSON.stringify(body))}`); } const S = "module.anthropic-consumer.usage.session"; const L = "module.anthropic-manager.usage.read"; // ------------------------------------------------------------------------------------------------ // (1) UPSERT-LATEST — the same session-grain key emitted twice; the second value REPLACES the first. // ------------------------------------------------------------------------------------------------ await emit(S, { rows: [{ licence: "personal", consumer: "n/anthropic-consumer/s1", period: "session", metric: "input_tokens", value: 100 }], raw: {} }); await waitFor("select value from usage where consumer='n/anthropic-consumer/s1' and metric='input_tokens'", /^100$/m); await emit(S, { rows: [{ licence: "personal", consumer: "n/anthropic-consumer/s1", period: "session", metric: "input_tokens", value: 250 }], raw: {} }); const latest = await waitFor("select value from usage where consumer='n/anthropic-consumer/s1' and metric='input_tokens'", /^250$/m); assert.match(latest, /^250$/m, `the upsert did not keep the latest reading (expected 250):\n${latest}`); const count = (await usageQuery("select count(*) from usage where consumer='n/anthropic-consumer/s1' and metric='input_tokens'")).out; assert.match(count, /^1$/m, `the upsert kept more than one row for the key (expected exactly 1):\n${count}`); // ------------------------------------------------------------------------------------------------ // (2) TWO GRAINS — a licence-grain reading in the SAME table, told apart from the session grain only // by `consumer`. After it, at least two distinct consumers exist under the one licence. // ------------------------------------------------------------------------------------------------ await emit(L, { rows: [{ licence: "personal", consumer: "n/anthropic-manager", period: "5h", metric: "utilization", value: 12 }], raw: {} }); const consumers = await waitFor("select count(distinct consumer) from usage where licence='personal'", /^[2-9]\d*$/m); assert.match(consumers, /^[2-9]\d*$/m, `the two grains did not both land under the licence (expected >= 2 consumers):\n${consumers}`); const licenceGrain = await waitFor("select value from usage where consumer='n/anthropic-manager' and period='5h' and metric='utilization'", /^12$/m); assert.match(licenceGrain, /^12$/m, `the licence-grain row is not selectable:\n${licenceGrain}`); // Both grains distinguishable by consumer: the session grain's consumer carries a session segment, // the licence grain's is the holding module alone. const bothGrains = (await usageQuery("select consumer from usage where licence='personal' order by consumer")).out; assert.match(bothGrains, /n\/anthropic-consumer\/s1/, `the session grain is missing:\n${bothGrains}`); assert.match(bothGrains, /n\/anthropic-manager/, `the licence grain is missing:\n${bothGrains}`); // ------------------------------------------------------------------------------------------------ // (3) IN THE CLEAR — the value is returned by an ordinary select as plaintext (250 above was read // straight out, not unsealed); the raw payload is likewise an ordinary jsonb column. // ------------------------------------------------------------------------------------------------ const clear = (await usageQuery("select value, raw from usage where consumer='n/anthropic-consumer/s1' and metric='input_tokens'")).out; assert.match(clear, /250/, `the value is not returned as plaintext by an ordinary select:\n${clear}`); // Helper: wait for the mesh to write model-usage's bound file with an `as`, and parse it. async function waitForBinding(path: string): Promise<{ as: string; provision: string }> { let raw = ""; const until = Date.now() + 90_000; while (Date.now() < until) { const got = await on(NODE, `cat ${path} 2>/dev/null`); if (got.ok && /"as"/.test(got.out)) { raw = got.out; break; } await new Promise((r) => setTimeout(r, 3000)); } assert.match(raw, /"as"/, `the mesh never wrote a binding with a login to ${path}:\n${raw}`); return JSON.parse(raw) as { as: string; provision: string }; } });