Files
mesh-lab/test/integration/assigned-model-usage.test.ts
T
jschoubben e085e31f95 A bed that raises the foundation from the bundle derives the anchor's filter before it relies on the hub
The base ruleset (ADR 0088) admits ssh, the bus and the registry and nothing else until the mesh
derives one, and the mesh derives one only where the filter module is assigned — which genesis
does and these beds did not. Without it the hub's WireGuard port stayed closed, no joined node's
tunnel formed, and every module dialling the anchor by its overlay name timed out fetching the
broker's certificate; the model-usage bed showed it as a login that failed for a role never made.
2026-09-21 13:30:26 +02:00

417 lines
23 KiB
TypeScript

/**
* 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, deriveTheFilterOn } 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<string> {
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<string> {
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<void> {
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 (process.env["MESH_LAB_KEEP"]) {
console.log(`\nLEFT STANDING: ${instanceId} — not destroyed (MESH_LAB_KEEP).`);
return;
}
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<void> {
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`);
// The anchor's derived filter, admitting the hub's port — what genesis does on the control-node,
// and what a bed raised from the bundle must do itself (ADR 0088; see the harness).
await deriveTheFilterOn({ machine: "anchor", node: "anchor", hubPort: 51820, must, mesh, on });
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,
"--- containers:", (await on(NODE, `docker ps -a --format '{{.Names}} {{.Status}}'`)).out,
"--- postgres server:", (await on(NODE, `docker logs --tail 15 postgres 2>&1; ls -la /var/lib/postgres/`)).out,
"--- laptop filter:", (await on(NODE, `nft list ruleset 2>&1 | head -60`)).out,
"--- laptop → broker from a container:", (await on(NODE, `docker run --rm --network postgres alpine sh -c 'nc -zvw3 192.0.2.10 5671' 2>&1`)).out,
"--- anchor filter counters:", (await on("anchor", `nft -a list table inet mesh 2>&1 | head -60`)).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<string> {
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<void> {
const source = key.split(".")[1] ?? "usage-probe"; // module.<producer>.usage.<grain>
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 };
}
});