Merge pull request 'lavinmq-bed: prove the AMQP provider provisions a require-only consumer' (#13) from feat/lavinmq-bed into main
This commit was merged in pull request #13.
This commit is contained in:
@@ -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]
|
||||
@@ -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_<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 };
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user