diff --git a/scenarios/lavinmq-bed.yml b/scenarios/lavinmq-bed.yml new file mode 100644 index 0000000..ed9d8fd --- /dev/null +++ b/scenarios/lavinmq-bed.yml @@ -0,0 +1,59 @@ +# Two machines: one is the substrate, the other carries a lavinmq PROVIDER and a consumer of it. +# +# lavinmq is the mesh's own control-plane broker (mesh-broker), and it is ALSO offered as a +# user-facing capability: a module that needs a message queue gets its OWN broker — a scoped vhost and +# user on a lavinmq PROVIDER — not an account on the control broker. That makes lavinmq the sharp +# two-node case, for the same reason two-node-db is: the substrate's broker publishes 5672 on the node +# it runs on, and a lavinmq provider must publish 5672 too for its consumers to reach it — so the two +# cannot share a machine. The moment a real amqp provider must own a node's 5672, it collides with the +# control broker already there, and the chain is blocked single-node. +# +# This is the split that unblocks it. `anchor` runs the substrate (store, broker, control) and NOTHING +# else. `laptop` runs the whole chain: the lavinmq PROVIDER and the amqp-ping CONSUMER that requires +# it. Provider and consumer are co-located on laptop, so the grant never crosses a node boundary; only +# enrolment crosses to anchor, over the underlay both machines share. And because the substrate broker +# is on the OTHER node, the provider owns laptop's 5672 uncontested. +# +# 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 lavinmq /tmp/lavinmq.tar +# scripts/build-module-runtime.sh amqp-ping /tmp/amqp-ping.tar +# The service image cloudamqp/lavinmq:latest must be in the local daemon too — it is already stocked as +# the substrate's own broker image; each node pulls what it runs from the scenario's own registry by +# digest. +scenario: lavinmq-bed + +segments: + hosting: + kind: public + cidr: [192.0.2.0/24] + +machines: + anchor: + at: { segment: hosting, address: [192.0.2.10] } + inbound: allow + memory: 3GiB + cpus: 2 + laptop: + at: { segment: hosting, address: [192.0.2.11] } + inbound: allow + memory: 3GiB + cpus: 2 + +images: + # The first-node substrate: store, broker (itself lavinmq), control. + - postgres:17-alpine + - cloudamqp/lavinmq:latest + - mesh-control:development + # The lavinmq provider's runtime (reused for its run-once bootstrap and its provisioner) and the + # amqp-ping consumer's runtime, both built by scripts/build-module-runtime.sh into the local daemon + # and stocked into the scenario's own registry, which is where the hosts pull them from by digest. + # The lavinmq SERVICE image is cloudamqp/lavinmq:latest, already stocked above. + - mesh-runtime-lavinmq:development + - mesh-runtime-amqp-ping:development + +place: + all: [host, runtime] diff --git a/test/integration/lavinmq-bed.test.ts b/test/integration/lavinmq-bed.test.ts new file mode 100644 index 0000000..44a2ab8 --- /dev/null +++ b/test/integration/lavinmq-bed.test.ts @@ -0,0 +1,430 @@ +/** + * 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. + // ================================================================================================ + 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 }; + } +});