Files
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

428 lines
22 KiB
TypeScript

/**
* A lavinmq PROVIDER and a consumer of it ride one node while the foundation's own broker owns 5672 on
* another — the end-to-end proof that a module which needs a message queue gets its OWN broker.
*
* lavinmq is the mesh's control-plane broker AND a user-facing capability: a consumer that requires
* `amqp` is given a scoped vhost + user on a lavinmq PROVIDER, not an account on the control broker
* (novox/hq ADR 0048). That makes it the two-node case, the twin of two-node-db: the foundation's
* broker publishes 5672 on anchor, and a lavinmq provider must publish 5672 for its consumers to
* reach it — so the two cannot share a machine. `anchor` runs the foundation and nothing else; `laptop`
* runs the provider AND the amqp-ping consumer, co-located, and owns laptop's 5672 uncontested.
*
* The proof is layered:
* - the provider's run-once bootstrap computed the admin password hash lavinmq needs and wrote its
* config BEFORE the broker started (ADR 0052) — proven because the broker came up at all and is
* gone from `docker ps -a` (a step, not a service);
* - the lavinmq service and BOTH runtimes (bootstrap done, provisioner running) are up and stable;
* - each module got its own scoped broker account on the foundation broker (anchor), named for the
* node that runs it;
* - the provisioner (in mesh-lavinmq on laptop) created the consumer's vhost AND user, both named
* for the login the mesh derived — checked with `lavinmqctl` inside the broker;
* - the consumer was bound its scoped login, and its container connected to that vhost with the
* mesh-minted password and round-tripped a message (published and consumed it back).
*
* The db-name lesson, applied to AMQP: the consumer uses `${bound:amqp:as}` for BOTH its username and
* its vhost — the provider named the vhost after the login — and nothing is hardcoded; the provider's
* `serves` carries the port so the consumer references `${bound:amqp:port}`.
*
* 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 lavinmq /tmp/lavinmq.tar
* scripts/build-module-runtime.sh amqp-ping /tmp/amqp-ping.tar
* cloudamqp/lavinmq:latest must be in the local daemon too (it is the foundation's own broker image);
* scenarios/lavinmq-bed.yml stocks all of them, and each node pulls what it runs 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 = "lavinmq-bed";
/** The node that carries the provider and its consumer. anchor carries only the foundation. */
const NODE = "laptop";
/** The login the mesh derives for the consumer: mesh_<node>_<slug> (slug "ping"; ADR 0049). */
const CONSUMER_LOGIN = "mesh_laptop_ping";
let instanceId = "";
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;
}
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("the lavinmq provider and its consumer ride laptop while the foundation broker owns 5672 on anchor", {
skip, timeout: 1_500_000,
}, async () => {
// ================================================================================================
// THE MANIFESTS — the committed catalogue shapes, verbatim from mesh-catalog. Only the images (to
// what this scenario serves by digest) change.
// ================================================================================================
// --- lavinmq: a provider of the `amqp` interface. Its server publishes 5672 so consumers reach it
// at the node address the mesh writes into their grant; the management API (15672) stays inside the
// module's private network. A run-once bootstrap computes the admin password HASH lavinmq needs
// (the mesh delivers a plain secret; lavinmq's config wants a hash) and writes the config the broker
// reads, BEFORE the broker starts. The provisioner then reaches the management API as that admin.
const lavinmqManifest = JSON.stringify({
module: "lavinmq",
version: "1",
provides: [{ name: "amqp", scope: "mesh" }],
serves: { amqp: { port: 5672 } },
emits: ["module.lavinmq.amqp.provisioned", "module.lavinmq.amqp.deprovisioned"],
consumes: ["module.lavinmq.amqp.provisioned", "module.lavinmq.amqp.deprovisioned"],
receives: { amqp: "/var/lib/lavinmq-module/grants/mesh.json" },
grants: { amqp: "/var/lib/lavinmq-module/grants" },
"own-secrets": { default: "/var/lib/lavinmq-module/default.secret", broker: "/var/lib/mesh/lavinmq/broker" },
resources: [
{ id: "mesh-state", type: "directory", path: "/var/lib/mesh/lavinmq", mode: "0700" },
{ id: "state", type: "directory", path: "/var/lib/lavinmq-module", mode: "0700" },
{ id: "grants-dir", type: "directory", path: "/var/lib/lavinmq-module/grants", mode: "0700" },
{ id: "data", type: "directory", path: "/services/lavinmq/data", mode: "0700" },
{ id: "net", type: "network", name: "lavinmq" },
{
id: "bootstrap", type: "container", name: "lavinmq-bootstrap",
image: pinned("mesh-runtime-lavinmq"), "run-once": true,
volumes: [
"/var/lib/lavinmq-module:/var/lib/lavinmq-module",
"/var/lib/lavinmq-module/default.secret:/run/secrets/default:ro",
],
env: {
MESH_PROVISION_ADMIN_USER: "mesh-admin",
MESH_PROVISION_PASSWORD_FILE: "/run/secrets/default",
MESH_LAVINMQ_CONFIG_OUT: "/var/lib/lavinmq-module/lavinmq.ini",
MESH_LAVINMQ_DATA_DIR: "/var/lib/lavinmq",
},
args: ["run", "/app/modules/lavinmq/dist/bootstrap/index.js"],
},
{
id: "server", type: "container", name: "lavinmq", image: pinned("cloudamqp/lavinmq"),
network: "lavinmq", ports: ["5672"],
volumes: [
"/services/lavinmq/data:/var/lib/lavinmq",
"/var/lib/lavinmq-module/lavinmq.ini:/etc/lavinmq/lavinmq.ini:ro",
],
args: ["--config", "/etc/lavinmq/lavinmq.ini"],
},
{
id: "runtime", type: "container", name: "mesh-lavinmq", image: pinned("mesh-runtime-lavinmq"),
network: "lavinmq",
volumes: [
"/var/lib/mesh/lavinmq/broker:/run/secrets/broker:ro",
"/var/lib/lavinmq-module/grants:/var/lib/lavinmq-module/grants:ro",
"/var/lib/lavinmq-module/default.secret:/run/secrets/default:ro",
],
env: {
MESH_BROKER_FILE: "/run/secrets/broker",
MESH_RECEIVES: "/var/lib/lavinmq-module/grants/mesh.json",
MESH_PROVISION_LAVINMQ: "http://lavinmq:15672",
MESH_PROVISION_ADMIN_USER: "mesh-admin",
MESH_PROVISION_PASSWORD_FILE: "/run/secrets/default",
},
},
],
});
// --- amqp-ping: a consumer that requires `amqp`. It contributes nothing (the vhost is the login),
// binds the grant, and reads its connection from an env-file the mesh fills — USER and VHOST both
// ${bound:amqp:as} (the db-name lesson), PORT ${bound:amqp:port}, PASSWORD ${secret:amqp}. Its
// container connects and round-trips a message, then stays up. It carries a `slug` so its identity
// mesh_laptop_ping fits the 20-char backend bound (ADR 0049).
const amqpPingManifest = JSON.stringify({
module: "amqp-ping",
slug: "ping",
version: "1",
requires: ["amqp"],
contributes: {},
binds: { amqp: "/var/lib/amqp-ping/amqp.json" },
secrets: { amqp: "/var/lib/amqp-ping/amqp.secret" },
"own-secrets": { broker: "/var/lib/mesh/amqp-ping/broker" },
resources: [
{ id: "mesh-state", type: "directory", path: "/var/lib/mesh/amqp-ping", mode: "0700" },
{ id: "state", type: "directory", path: "/var/lib/amqp-ping", mode: "0700" },
{
id: "amqp-env", type: "file", path: "/var/lib/amqp-ping/amqp.env", mode: "0600",
content:
"MESH_AMQP_HOST=${bound:amqp:at}\nMESH_AMQP_PORT=${bound:amqp:port}\n" +
"MESH_AMQP_USER=${bound:amqp:as}\nMESH_AMQP_VHOST=${bound:amqp:as}\n" +
"MESH_AMQP_PASSWORD=${secret:amqp}\n",
},
{ id: "net", type: "network", name: "amqp-ping" },
{
id: "runtime", type: "container", name: "amqp-ping", image: pinned("mesh-runtime-amqp-ping"),
network: "amqp-ping",
"env-file": ["/var/lib/amqp-ping/amqp.env"],
args: ["run", "/app/modules/amqp-ping/dist/index.js"],
"restart-on": ["amqp-env"],
},
],
});
// --- add, issue a scoped broker account, assign to laptop --------------------------------------
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`);
const issued = await mesh(`module issue ${name} --node ${NODE}`);
assert.match(issued, /broker account/, issued);
await mesh(`assign ${NODE} ${name}`);
}
// The consumer connects to the provider by the binding's `at` — the provider's private-network
// address the mesh fills. So the overlay is assigned first, to both nodes, even though provider and
// consumer are co-located (the same reason two-node-db does it).
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("lavinmq", lavinmqManifest);
await addIssueAssign("amqp-ping", amqpPingManifest);
// ONE push, ONE convergence — provider and consumer resolved and applied together on laptop, past
// the run-once bootstrap that gates the broker.
await mesh(`push ${NODE}`);
await settled(NODE);
// ================================================================================================
// THE two-node split — the foundation broker owns 5672 on anchor, the provider owns it on laptop.
// ================================================================================================
const onAnchor = await must("anchor", `docker ps --format '{{.Names}}'`);
const onLaptop = await must(NODE, `docker ps --format '{{.Names}}'`);
assert.match(onAnchor, /(^|\n)mesh-broker(\n|$)/, "the foundation broker is not on the first node");
assert.doesNotMatch(onAnchor, /(^|\n)lavinmq(\n|$)/,
"the lavinmq provider landed on the foundation node — the 5672 collision this bed exists to avoid");
assert.doesNotMatch(onLaptop, /(^|\n)mesh-store(\n|$)/, "the foundation store leaked onto the second node");
assert.match(onLaptop, /(^|\n)lavinmq(\n|$)/, "the lavinmq provider is not on the second node");
// ================================================================================================
// THE run-once bootstrap ran to completion and is gone (a step, not a service). The broker being up
// at all is the proof it ran at the right phase — an un-bootstrapped broker has no config to read.
// ================================================================================================
const allOnLaptop = await must(NODE, `docker ps -a --format '{{.Names}}'`);
assert.doesNotMatch(allOnLaptop, /(^|\n)lavinmq-bootstrap(\n|$)/,
`the run-once bootstrap was left as a container instead of run to completion:\n${allOnLaptop}`);
const iniExists = await on(NODE, `test -f /var/lib/lavinmq-module/lavinmq.ini`);
assert.ok(iniExists.ok, "the bootstrap did not write the broker config");
// ================================================================================================
// THE service and both runtimes are up on laptop.
// ================================================================================================
for (const name of ["lavinmq", "mesh-lavinmq", "amqp-ping"]) {
assert.match(onLaptop, new RegExp(`(^|\\n)${name}(\\n|$)`),
`${name} is not running on the second node after the push:\n${onLaptop}\n---host log---\n${(await on(NODE, `tail -60 /var/log/mesh-host.log`)).out}`);
}
// Everything up and STABLE — no crash-loop. Sample restart counts, wait, require they did not rise.
const stable = ["lavinmq", "mesh-lavinmq", "amqp-ping"];
const restarts = new Map<string, number>();
for (const name of stable) {
const [running, count] = (await must(NODE, `docker inspect -f '{{.State.Running}} {{.RestartCount}}' ${name}`)).trim().split(" ");
assert.equal(running, "true",
`${name} is not running after the push:\n${(await on(NODE, `docker logs ${name} 2>&1 | tail -40`)).out}`);
restarts.set(name, Number(count));
}
await new Promise((r) => setTimeout(r, 20000));
for (const name of stable) {
const [running, count] = (await must(NODE, `docker inspect -f '{{.State.Running}} {{.RestartCount}}' ${name}`)).trim().split(" ");
assert.equal(running, "true",
`${name} fell over after the push:\n${(await on(NODE, `docker logs ${name} 2>&1 | tail -40`)).out}`);
assert.ok(Number(count) <= (restarts.get(name) ?? 0),
`${name} is crash-looping (restart count rose ${restarts.get(name)} -> ${count}):\n` +
`${(await on(NODE, `docker logs ${name} 2>&1 | tail -40`)).out}`);
}
// ================================================================================================
// Each module got its own scoped broker account on the foundation broker (anchor), named for the
// node that runs it and the module.
// ================================================================================================
const users = await must("anchor", `docker exec mesh-broker lavinmqctl list_users 2>&1`);
for (const acct of ["laptop-lavinmq", "laptop-amqp-ping"]) {
assert.match(users, new RegExp(acct), `the scoped account ${acct} is not on the broker:\n${users}`);
}
// ================================================================================================
// THE provider path — the provisioner (in mesh-lavinmq on laptop) created the consumer's vhost AND
// user, both named for the login the mesh derived. Checked with lavinmqctl inside the broker (its
// local control socket needs no auth); the admin authenticating to the management API is implicit,
// because the vhost could not exist otherwise.
// ================================================================================================
let vhosts = { out: "", ok: false };
const untilVhost = Date.now() + 90_000;
while (Date.now() < untilVhost) {
vhosts = await on(NODE, `docker exec lavinmq lavinmqctl list_vhosts 2>&1`);
if (vhosts.ok && new RegExp(CONSUMER_LOGIN).test(vhosts.out)) break;
await new Promise((r) => setTimeout(r, 3000));
}
assert.match(vhosts.out, new RegExp(CONSUMER_LOGIN),
`the provisioner never created the consumer's vhost ${CONSUMER_LOGIN}:\n` +
`${(await on(NODE, `docker logs mesh-lavinmq 2>&1 | tail -40`)).out}\n---\n${vhosts.out}`);
const provUsers = await must(NODE, `docker exec lavinmq lavinmqctl list_users 2>&1`);
assert.match(provUsers, new RegExp(CONSUMER_LOGIN),
`the provisioner created the vhost but not the user ${CONSUMER_LOGIN}:\n${provUsers}`);
// ================================================================================================
// THE consumer was bound its scoped login, and its container round-tripped a message. The binding
// proves the mesh derived and delivered the login; the log line proves the container connected to
// its vhost with the mesh-minted password and published+consumed a message back.
// ================================================================================================
const bound = await waitForBinding("/var/lib/amqp-ping/amqp.json");
assert.equal(bound.provision, "amqp", `amqp-ping was bound the wrong provision: ${bound.provision}`);
assert.equal(bound.as, CONSUMER_LOGIN,
`amqp-ping's login is ${bound.as}, not the derived ${CONSUMER_LOGIN}`);
let log = "";
const untilPing = Date.now() + 150_000;
while (Date.now() < untilPing) {
log = (await on(NODE, `docker logs amqp-ping 2>&1 | tail -40`)).out;
if (/round-trip ok/.test(log)) break;
await new Promise((r) => setTimeout(r, 3000));
}
assert.match(log, /round-trip ok/,
`the consumer never round-tripped a message over its granted broker:\n${log}`);
assert.match(log, new RegExp(`vhost ${CONSUMER_LOGIN}`),
`the consumer connected to a vhost other than its own login:\n${log}`);
// Helper: wait for the mesh to write the consumer's binding 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 };
}
});