/** * A lavinmq PROVIDER and a consumer of it ride one node while the substrate'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 substrate'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 substrate 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 substrate 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/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 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 substrate'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 } from "./harness.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 = "lavinmq-bed"; /** The node that carries the provider and its consumer. anchor carries only the substrate. */ const NODE = "laptop"; /** The login the mesh derives for the consumer: mesh__ (slug "ping"; ADR 0049). */ const CONSUMER_LOGIN = "mesh_laptop_ping"; let instanceId = ""; let stocked: string[] = []; 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-control /mesh-control ${command}`, timeoutMs); } /** The pinned reference for one of the scenario's images, by repository. */ function pinned(repository: string): string { const found = stocked.find((r) => r.slice(r.indexOf("/") + 1, r.indexOf("@")) === repository); assert.ok(found, `the scenario stocks no ${repository}; it serves ${stocked.join(", ")}`); return found; } /** The substrate bundle, its image references pointed at this scenario's own registry. */ function bundleFor(images: string[]): string { let text = readFileSync(bundle, "utf8"); for (const ref of images) { const repository = ref.slice(ref.indexOf("/") + 1, ref.indexOf("@")); const escaped = repository.replaceAll("/", "\\/").replaceAll(".", "\\."); text = text.replaceAll(new RegExp(`[A-Za-z0-9_.:-]+\\/${escaped}@sha256:[0-9a-f]+`, "g"), ref); } return text; } 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 { 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; stocked = 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("the lavinmq provider and its consumer ride laptop while the substrate 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 { 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`); 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`); 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 substrate 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 substrate broker is not on the first node"); assert.doesNotMatch(onAnchor, /(^|\n)lavinmq(\n|$)/, "the lavinmq provider landed on the substrate node — the 5672 collision this bed exists to avoid"); assert.doesNotMatch(onLaptop, /(^|\n)mesh-store(\n|$)/, "the substrate 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(); 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 substrate 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. // ================================================================================================ // DIAG: the provisioner logged no create — is the grant in its receives file, and is the management // API it must reach actually up? (A stuck waitReady logs neither a create nor an error.) { const grants = (await on(NODE, `cat /var/lib/lavinmq-module/grants/mesh.json 2>&1 || echo MISSING`)).out; const api = (await on(NODE, `docker exec lavinmq sh -c 'wget -qO- http://127.0.0.1:15672/api/overview 2>&1 | head -c 200 || curl -s http://127.0.0.1:15672/api/overview 2>&1 | head -c 200' 2>&1`)).out; const full = (await on(NODE, `docker logs mesh-lavinmq 2>&1`)).out; console.log(`\n=== LAVINMQ DIAG ===\ngrants mesh.json:\n${grants}\n--- mgmt API (lavinmq:15672/api/overview) ---\n${api}\n--- mesh-lavinmq full log ---\n${full}\n=== END LAVINMQ DIAG ===\n`); } 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 }; } });