Files
mesh-lab/test/integration/lavinmq-bed.test.ts
T
jschoubben 8e494c0e78 Add lavinmq-bed: a two-node amqp provider + consumer bed
A lavinmq provider and an amqp-ping consumer ride laptop while the
substrate's own broker owns 5672 on anchor — the twin of two-node-db.
lavinmq is the mesh's control broker AND a user-facing capability, so a
provider must publish 5672 for its consumers and cannot share a node
with the control broker that already owns it; the split unblocks the
chain single-node.

The bed proves, layered: the run-once bootstrap computed the admin hash
and wrote the broker config before the broker started (ADR 0052); the
service and both runtimes are up and stable; each module got its scoped
broker account on the substrate broker; the provisioner created the
consumer's vhost AND user, both named for the derived login; and the
consumer connected to that vhost with the mesh-minted password and
round-tripped a message. The consumer uses ${bound:amqp:as} for user
and vhost both, and the provider's serves carries the port.

Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
2026-09-06 15:50:33 +02:00

431 lines
22 KiB
TypeScript

/**
* 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_<node>_<slug> (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<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 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<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;
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<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`);
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<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 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.
// ================================================================================================
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 };
}
});