events: an e2e test — an emitted event reaches the audit trail over the mesh's broker #2
Executable
+57
@@ -0,0 +1,57 @@
|
||||
#!/usr/bin/env bash
|
||||
# Build the runtime+audit-logger image the events test runs, and save it to a tar.
|
||||
#
|
||||
# The image is the tier-3 tool runtime (mesh-tools) carrying one tier-4 module (audit-logger) and
|
||||
# the sdk it imports. It is what MESH_LAB_RUNTIME points at:
|
||||
#
|
||||
# scripts/build-runtime-image.sh /tmp/mesh-runtime-audit.tar
|
||||
# MESH_LAB_RUNTIME=/tmp/mesh-runtime-audit.tar node --test test/integration/events.test.ts
|
||||
#
|
||||
# The sdk is vendored (dereferenced), not npm-installed: the sdk is not published, and the lab
|
||||
# machine has no route out anyway — the image must be self-contained. Sibling repositories are
|
||||
# assumed alongside this one; override with MESH_TOOLS / MESH_SDK / MESH_CATALOG.
|
||||
set -euo pipefail
|
||||
|
||||
OUT="${1:?usage: build-runtime-image.sh <output.tar>}"
|
||||
HERE="$(cd "$(dirname "$0")/.." && pwd)"
|
||||
ROOT="$(cd "$HERE/.." && pwd)"
|
||||
MESH_TOOLS="${MESH_TOOLS:-$ROOT/mesh-tools}"
|
||||
MESH_SDK="${MESH_SDK:-$ROOT/mesh-sdk}"
|
||||
MESH_CATALOG="${MESH_CATALOG:-$ROOT/mesh-catalog}"
|
||||
AUDIT="$MESH_CATALOG/modules/audit-logger"
|
||||
TAG="${RUNTIME_TAG:-mesh-runtime-audit:lab}"
|
||||
BASE="${RUNTIME_BASE:-node:22-bookworm-slim}"
|
||||
|
||||
echo "building $TAG from:"
|
||||
echo " runtime $MESH_TOOLS"
|
||||
echo " sdk $MESH_SDK"
|
||||
echo " module $AUDIT"
|
||||
|
||||
# Compile the three, so the image carries current dist. The sdk first — the others import it.
|
||||
( cd "$MESH_SDK" && npm run build >/dev/null )
|
||||
( cd "$MESH_TOOLS" && npm run build >/dev/null )
|
||||
( cd "$AUDIT" && npx tsc audit.ts index.ts --module NodeNext --moduleResolution NodeNext \
|
||||
--target ES2022 --outDir dist >/dev/null )
|
||||
|
||||
STAGE="$(mktemp -d)"
|
||||
trap 'rm -rf "$STAGE"' EXIT
|
||||
cp -r "$MESH_TOOLS/dist" "$STAGE/dist"
|
||||
cp -rL "$MESH_TOOLS/node_modules" "$STAGE/node_modules" # -L materialises the @novox/mesh-sdk symlink
|
||||
mkdir -p "$STAGE/modules/audit-logger"
|
||||
cp -r "$AUDIT/dist" "$STAGE/modules/audit-logger/dist"
|
||||
cp "$MESH_TOOLS/package.json" "$STAGE/package.json"
|
||||
|
||||
cat > "$STAGE/Dockerfile" <<DOCKER
|
||||
FROM $BASE
|
||||
WORKDIR /app
|
||||
COPY package.json ./
|
||||
COPY node_modules ./node_modules
|
||||
COPY dist ./dist
|
||||
COPY modules ./modules
|
||||
ENV MESH_TOOL_MODULES=/app/modules/audit-logger/dist/index.js
|
||||
ENTRYPOINT ["node", "dist/main.js"]
|
||||
DOCKER
|
||||
|
||||
docker build -t "$TAG" "$STAGE"
|
||||
docker save -o "$OUT" "$TAG"
|
||||
echo "saved $TAG -> $OUT"
|
||||
@@ -0,0 +1,221 @@
|
||||
/**
|
||||
* An event a module emits reaches an audit trail, over the broker the mesh raised.
|
||||
*
|
||||
* Everything else here proves the broker carries *commands* — a declaration crosses it, a
|
||||
* credential is delivered over it. This proves the other half of the bus (novox/hq ADR 0046): the
|
||||
* events exchange, where a module emits and any number listen, and the audit logger consumes `#`
|
||||
* and writes down what happened. It runs against the `mesh-broker` this scenario's own host raised
|
||||
* from the substrate bundle — not a broker a test stood up — because "the mesh hosts the broker"
|
||||
* (tier-1 substrate) is the thing being relied on.
|
||||
*
|
||||
* The wire shape it asserts is ADR 0047: metadata rides as headers so the body is only the
|
||||
* payload, a consumer gets a durable per-consumer queue `<node>.<module>.events`, and a
|
||||
* dead-letter exchange `mesh.events.dead` stands behind it. The queue and that exchange existing
|
||||
* on the raised broker is the check that the runtime provisioned the contract, not just that a
|
||||
* message happened to arrive.
|
||||
*
|
||||
* It needs the host binary and the substrate bundle, like the mesh walk, plus a runtime image:
|
||||
*
|
||||
* MESH_LAB_HOST_BINARY=.../mesh-host
|
||||
* MESH_LAB_BUNDLE=.../examples/substrate-first-node.lock
|
||||
* MESH_LAB_RUNTIME=.../mesh-runtime-audit.tar (docker save of the runtime+audit-logger image;
|
||||
* built by scripts/build-runtime-image.sh)
|
||||
*
|
||||
* The runtime is run as a container against the broker here, rather than assigned through the
|
||||
* control plane. Assigning it — so the mesh delivers its broker credential the way it does the
|
||||
* builder's — is the next step; this proves the events path itself first.
|
||||
*/
|
||||
|
||||
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";
|
||||
import { incus } from "../../src/incus/client.ts";
|
||||
import { machineName } from "../../src/lifecycle/names.ts";
|
||||
|
||||
const capability = await labIsUsable();
|
||||
const binary = hostBinaryPath();
|
||||
const bundle = process.env["MESH_LAB_BUNDLE"] ?? "";
|
||||
const runtime = process.env["MESH_LAB_RUNTIME"] ?? "";
|
||||
|
||||
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/)"
|
||||
: !runtime || !existsSync(runtime)
|
||||
? "MESH_LAB_RUNTIME is not set to a runtime image tar (scripts/build-runtime-image.sh)"
|
||||
: false;
|
||||
|
||||
const SCENARIO = "first-node";
|
||||
const MACHINE = "anchor";
|
||||
/** Where the audit-logger container writes its trail, on the machine — mounted from a host dir. */
|
||||
const TRAIL_DIR = "/var/lib/mesh-audit";
|
||||
const TRAIL = `${TRAIL_DIR}/audit.log`;
|
||||
/** The broker the substrate raised, reachable on the node's loopback (novox/hq ADR 0001). */
|
||||
const BROKER = "amqp://guest:guest@127.0.0.1:5672/";
|
||||
|
||||
let instanceId = "";
|
||||
/** The tag `docker load` reported for the runtime image, so the container names what was loaded. */
|
||||
let runtimeImage = "";
|
||||
|
||||
function quote(s: string): string {
|
||||
return `'${s.replaceAll("'", `'\\''`)}'`;
|
||||
}
|
||||
|
||||
/**
|
||||
* A command on the machine, its exit read from a marker on its own line.
|
||||
*
|
||||
* `exec 2>&1` on its own first line and the marker on its own last line, so a command that carries
|
||||
* a heredoc — the substrate bundle is written with one — terminates where it says it does rather
|
||||
* than swallowing the marker (the fault mesh.test.ts documents).
|
||||
*/
|
||||
async function on(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(command: string, timeoutMs?: number): Promise<string> {
|
||||
const { out, ok } = await on(command, timeoutMs);
|
||||
if (!ok) throw new Error(`${MACHINE}: ${command}\n${out}`);
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* The substrate bundle, its image references pointed at this scenario's own registry.
|
||||
*
|
||||
* A digest belongs to whatever registry serves it, so the committed bundle names a registry that
|
||||
* is not this one; matching by repository and rewriting to the digest this registry assigned is
|
||||
* what makes it applicable (the same rewrite mesh.test.ts does).
|
||||
*/
|
||||
function bundleFor(images: string[]): string {
|
||||
let text = readFileSync(bundle, "utf8");
|
||||
for (const pinned of images) {
|
||||
const repository = pinned.slice(pinned.indexOf("/") + 1, pinned.indexOf("@"));
|
||||
const escaped = repository.replaceAll("/", "\\/").replaceAll(".", "\\.");
|
||||
text = text.replaceAll(
|
||||
new RegExp(`[A-Za-z0-9_.:-]+\\/${escaped}@sha256:[0-9a-f]+`, "g"),
|
||||
pinned,
|
||||
);
|
||||
}
|
||||
return text;
|
||||
}
|
||||
|
||||
/** Read the trail back as parsed JSON lines. */
|
||||
async function trail(): Promise<Record<string, unknown>[]> {
|
||||
const raw = await must(`cat ${TRAIL} 2>/dev/null || true`);
|
||||
return raw
|
||||
.split("\n")
|
||||
.map((l) => l.trim())
|
||||
.filter(Boolean)
|
||||
.map((l) => JSON.parse(l) as Record<string, unknown>);
|
||||
}
|
||||
|
||||
before(async () => {
|
||||
if (skip) return;
|
||||
|
||||
const raised = await raise(loadScenario(`scenarios/${SCENARIO}.yml`), {
|
||||
onProgress: (m) => console.log(`raise: ${m}`),
|
||||
});
|
||||
instanceId = raised.instanceId;
|
||||
|
||||
// The node raises its substrate — store, broker and the rest — from the bundle, applied from a
|
||||
// file because the digests are this registry's and are not known until it is up.
|
||||
await must(`cat > /tmp/substrate.lock <<'MESHBUNDLE'\n${bundleFor(raised.images)}\nMESHBUNDLE`);
|
||||
await must(`${HOST_PATH} apply /tmp/substrate.lock`, 600_000);
|
||||
|
||||
const running = await must(`docker ps --format '{{.Names}}'`);
|
||||
assert.match(running, /mesh-broker/, `the substrate did not raise a broker:\n${running}`);
|
||||
|
||||
// Bring the runtime+audit-logger image onto the machine. Loaded, not pulled: the machine has no
|
||||
// route out (novox/hq the lab is a closed address space), so the image arrives as a tar the way
|
||||
// the builder binary does, and `docker load` names what it loaded.
|
||||
await incus([
|
||||
"file", "push", runtime, `${machineName(instanceId, MACHINE)}/tmp/runtime.tar`, "--mode", "0644",
|
||||
], 300_000);
|
||||
const loaded = await must(`docker load < /tmp/runtime.tar`);
|
||||
const named = loaded.match(/Loaded image:\s*(\S+)/)?.[1];
|
||||
assert.ok(named, `docker load did not name the image:\n${loaded}`);
|
||||
runtimeImage = named;
|
||||
|
||||
// Start the audit-logger: the runtime bound to the mesh's broker, consuming `#`, its trail on a
|
||||
// mounted directory so the assertions read what it actually wrote.
|
||||
await must(`mkdir -p ${TRAIL_DIR}`);
|
||||
await must(
|
||||
`docker run -d --name mesh-audit --network host ` +
|
||||
`-e MESH_BROKER_URL=${quote(BROKER)} -e MESH_MODULE=audit-logger -e MESH_NODE=${MACHINE} ` +
|
||||
`-e AUDIT_LOG=/trail/audit.log -v ${TRAIL_DIR}:/trail ` +
|
||||
`${runtimeImage}`,
|
||||
);
|
||||
// Give the subscription a moment to bind before anything is emitted at it.
|
||||
await new Promise((r) => setTimeout(r, 4000));
|
||||
}, { timeout: 1_800_000 });
|
||||
|
||||
after(async () => {
|
||||
if (instanceId) await destroy(instanceId);
|
||||
await destroyAll(`${SCENARIO}-`);
|
||||
}, { timeout: 600_000 });
|
||||
|
||||
test("an emitted event reaches the audit trail with its metadata in headers", {
|
||||
skip, timeout: 300_000,
|
||||
}, async () => {
|
||||
// Emitted as a different module (source=probe), so the trail's source is the emitter's, read
|
||||
// back from the header — proof it did not come from the body.
|
||||
await must(
|
||||
`docker exec -e MESH_MODULE=probe -e MESH_NODE=${MACHINE} mesh-audit ` +
|
||||
`node dist/main.js emit module.probe.site.created '{"domain":"my-app"}'`,
|
||||
);
|
||||
// A node-origin event too (novox/hq ADR 0047 reserves module.* / mesh.* / node.*), to show the
|
||||
// audit sink takes all of them, not only a module's.
|
||||
await must(
|
||||
`docker exec -e MESH_MODULE=probe -e MESH_NODE=${MACHINE} mesh-audit ` +
|
||||
`node dist/main.js emit node.${MACHINE}.tick '{"n":1}'`,
|
||||
);
|
||||
|
||||
// Poll: the trail is written by a separate container reacting to the broker, so it is not there
|
||||
// the instant emit returns.
|
||||
let lines: Record<string, unknown>[] = [];
|
||||
const until = Date.now() + 30_000;
|
||||
while (Date.now() < until) {
|
||||
lines = await trail();
|
||||
if (lines.length >= 2) break;
|
||||
await new Promise((r) => setTimeout(r, 2000));
|
||||
}
|
||||
|
||||
const types = lines.map((l) => l.type);
|
||||
assert.ok(
|
||||
types.includes("module.probe.site.created") && types.includes(`node.${MACHINE}.tick`),
|
||||
`both events did not reach the trail; it holds ${JSON.stringify(types)}\n` +
|
||||
`${(await on(`docker logs mesh-audit 2>&1 | tail -20`)).out}`,
|
||||
);
|
||||
|
||||
const site = lines.find((l) => l.type === "module.probe.site.created")!;
|
||||
assert.equal(site.source, "probe", "x-source did not survive as the trail's source");
|
||||
assert.equal(site.node, MACHINE, "x-node did not survive");
|
||||
assert.ok(typeof site.id === "string" && site.id.length > 0, "no x-event-id was recorded");
|
||||
assert.deepEqual(site.body, { domain: "my-app" }, "the body was not exactly the payload");
|
||||
});
|
||||
|
||||
test("the runtime provisioned the ADR 0047 queue and dead-letter on the raised broker", {
|
||||
skip, timeout: 120_000,
|
||||
}, async () => {
|
||||
// Asked of the broker itself, so this is the shape that actually exists on the bus, not the
|
||||
// shape the code intends. A durable per-consumer queue and a dead-letter home are what make the
|
||||
// trail survive a restart and set a poison event aside — assert they are there, not assumed.
|
||||
const queues = await must(`docker exec mesh-broker lavinmqctl list_queues name durable`);
|
||||
assert.match(queues, new RegExp(`${MACHINE}\\.audit-logger\\.events`),
|
||||
`the durable per-consumer queue is missing:\n${queues}`);
|
||||
|
||||
const exchanges = await must(`docker exec mesh-broker lavinmqctl list_exchanges name`);
|
||||
assert.match(exchanges, /mesh\.events(\s|$)/m, `the events exchange is missing:\n${exchanges}`);
|
||||
assert.match(exchanges, /mesh\.events\.dead/, `the dead-letter exchange is missing:\n${exchanges}`);
|
||||
});
|
||||
Reference in New Issue
Block a user