Files
mesh-lab/test/integration/assigned-model-usage.test.ts
T
jschoubben 675facdb0d The beds name images the way a machine would find them
Twenty-eight integration tests each carried their own copy of the same two helpers,
which pointed a manifest and the substrate bundle at whatever the lab's registry had
assigned. They now share two in the harness, and the difference is the point: ours is
rewritten to the ID the machine holds it under, and everything else is left exactly as
written so the machine pulls it.

**The substrate bundle is where the fiction was most load-bearing.** mesh-host's
`examples/substrate-first-node.lock` pins all three of its images at
`192.0.2.250:5000/…`, which is the address the lab's registry served from — it was
written for a target, and the target was the lab. Two of those are ordinary third-party
images and become the digests mesh-catalog's own postgres and lavinmq modules pin, so
the substrate's store and broker are literally the images the mesh runs. mesh-control
exists in no registry at all and becomes the ID the machine was handed. **The bundle
itself should be fixed in mesh-host and this substitution deleted with it.**

Beds that wrote a manifest by hand named an image by repository and let the rewrite
supply a digest. There is nothing to supply one now, so `onTheMachine` refuses an
unpinned reference and hands back the digest the catalogue pins — a bed runs the image
the mesh ships, and a bed that drifts from the catalogue is testing a different
postgres.

Three beds took a third-party image out of the raised list, which no longer contains
one: certificates (pebble), objectstore (minio and its client) and provisioner
(postgres) now name theirs and pull it. builds and mesh publish into the MESH's own
artifact store — the `registry` module's image, on the node, on 5000 — rather than into
scenery the lab raised. That is a different claim, and only one of them exists in
production.

New unit tests cover what a full raise would otherwise be the only way to check: the
routes an egress machine gets (that its gateway is still the path to the rest of the
scenario, that a range with no path is unreachable rather than leaked to the uplink,
that each family gets its own next hop), which machine is handed which of our images,
and the `images:` rule that refuses a third-party entry. The "shipped scenarios are
valid" test now loads every scenario rather than two of them.

Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
2026-09-10 23:16:41 +02:00

396 lines
21 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-control 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 substrate store both
* want host port 5432, so the substrate (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 substrate bundle:
*
* MESH_LAB_HOST_BINARY=.../mesh-host
* MESH_LAB_BUNDLE=.../examples/substrate-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, substrateBundle, 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 substrate 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
* substrate. */
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-control /mesh-control ${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 substrate bundle: ours by the ID the machine holds, everything else upstream. */
function bundleFor(images: HeldImage[]): string {
return substrateBundle(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-control /mesh-control 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/substrate.lock <<'MESHBUNDLE'\n${bundleFor(raised.images)}\nMESHBUNDLE`);
await must("anchor", `${HOST_PATH} apply /tmp/substrate.lock`, 600_000);
const up = await must("anchor", `docker ps --format '{{.Names}}'`);
for (const c of ["mesh-store", "mesh-broker", "mesh-control"]) {
assert.match(up, new RegExp(c), `the substrate 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: "superuser-env", type: "file", path: "/var/lib/postgres/superuser.env", mode: "0600", content: "POSTGRES_PASSWORD=${secret:superuser}\n" },
{ 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",
env: { POSTGRES_USER: "postgres", POSTGRES_DB: "postgres" },
"env-file": ["/var/lib/postgres/superuser.env"],
ports: ["5432"],
volumes: ["/services/postgres/db-data:/var/lib/postgresql/data"],
},
{
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" },
{
id: "db-env", type: "file", path: "/var/lib/model-usage/db.env", mode: "0600",
content:
"DATABASE_URL=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",
],
env: { MESH_BROKER_FILE: "/run/secrets/broker" },
"env-file": ["/var/lib/model-usage/db.env"],
},
],
});
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-control:/${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 substrate 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})`);
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}`);
// Inject a usage event into the mesh. model-usage is a PURE CONSUMER, so its own broker account has
// no publish right (mesh-control 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 };
}
});