One name per thing, per the HQ glossary: the module/container/image/binary/repo becomes mesh-controller, the seat the-controller, and the store+broker pair the foundation (embedded base bundles, default template and example lock renamed with their go:embed directives). No behaviour change — a pure vocabulary rename. Claude-Session: https://claude.ai/code/session_01D6qtiYU3P9jk3pnAXyAFyx
396 lines
21 KiB
TypeScript
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-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<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 (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-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})`);
|
|
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-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 };
|
|
}
|
|
});
|