import { test } from "node:test"; import assert from "node:assert/strict"; import amqp from "amqplib"; import { connectAmqp } from "../dist/broker-amqp.js"; import { useBroker } from "@novox/mesh-sdk/messaging"; import { emit, on, type Event } from "@novox/mesh-sdk/events"; // Binding conformance: does the mesh-tools AMQP *adapter* honour the ADR 0047 wire contract — // headers on the wire, a body that is only the payload, persistent messages, a durable per-consumer // queue, dead-letter, poison handling? This tests the adapter in isolation, against a disposable // broker that stands in for the one the mesh hosts (ADR 0001). It is NOT an event test of the mesh: // that lives in a full lab scenario where the mesh raises the broker as substrate. Requires // $MESH_BROKER_URL (a throwaway broker); skipped, never failed, when it is not set. const url = process.env.MESH_BROKER_URL; test("an event rides the wire with ADR 0047 headers and a body that is only the payload", { skip: !url }, async () => { const sub = await connectAmqp(url!); const pub = await connectAmqp(url!); // The consumer's identity names its durable queue (..events). process.env.MESH_NODE = "lab"; process.env.MESH_MODULE = "audit-logger"; useBroker(() => sub); const got: Event[] = []; await on("#", async (e) => void got.push(e)); await delay(200); // let the binding settle before publishing // The emitter is a different module on the same node. process.env.MESH_MODULE = "umami"; useBroker(() => pub); await emit("module.umami.site.created", { domain: "my-app" }, { causationId: "cmd-1" }); await waitFor(() => got.length > 0, 4000); const e = got[0]; assert.equal(e.type, "module.umami.site.created"); assert.equal(e.source, "umami"); // x-source — read from a header, not the body assert.equal(e.node, "lab"); // x-node assert.ok(e.id, "x-event-id present"); // the handle a consumer dedups on assert.equal(e.causationId, "cmd-1"); // x-causation-id round-trips assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/); // x-time, RFC-3339 assert.deepEqual(e.body, { domain: "my-app" }); // provenance never leaked into the body // The durable per-consumer queue exists and is bound — a passive assert throws if it does not. const probe = await amqp.connect(url!); const pch = await probe.createChannel(); await pch.checkQueue("lab.audit-logger.events"); await probe.close(); await sub.close(); await pub.close(); delete process.env.MESH_NODE; delete process.env.MESH_MODULE; }); test("a handler that keeps failing dead-letters the event past the redelivery limit", { skip: !url }, async () => { const module = `flaky-${Date.now()}`; const c = await connectAmqp(url!); process.env.MESH_NODE = "lab"; process.env.MESH_MODULE = module; useBroker(() => c); let attempts = 0; await on("module.test.boom", async () => { attempts++; throw new Error("boom"); }); await delay(200); await emit("module.test.boom", { n: 1 }); // First delivery requeues once; the redelivered copy is dead-lettered — two attempts, then it // leaves the consumer queue for good. await waitFor(() => attempts >= 2, 5000); await delay(300); assert.equal(attempts, 2, "attempted twice, not looping forever"); // The poison event is retained on the dead-letter queue for inspection, not vanished. const probe = await amqp.connect(url!); const pch = await probe.createChannel(); const dead = await pch.get("mesh.events.dead", { noAck: true }); assert.ok(dead, "the rejected event is on mesh.events.dead"); assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.test.boom"); await probe.close(); await c.close(); delete process.env.MESH_NODE; delete process.env.MESH_MODULE; }); test("an undecodable event body is dead-lettered, not looped and not silently swallowed", { skip: !url }, async () => { const module = `poison-${Date.now()}`; const c = await connectAmqp(url!); process.env.MESH_NODE = "lab"; process.env.MESH_MODULE = module; useBroker(() => c); // Drain any earlier dead events so the assert below sees only this test's. const drain = await amqp.connect(url!); const dch = await drain.createChannel(); await dch.purgeQueue("mesh.events.dead"); let handlerRuns = 0; await on("module.poison.raw", async () => void handlerRuns++); await delay(200); // Publish a body that is not JSON straight onto the events exchange — a malformed emitter. A // confirm channel, waited on, so the broker has the message before the connection closes. const raw = await amqp.connect(url!); const rch = await raw.createConfirmChannel(); rch.publish("mesh.events", "module.poison.raw", Buffer.from("this is not json{"), { persistent: true, headers: { "x-event-id": "poison-1", "x-source": "bad", "x-node": "lab" }, }); await rch.waitForConfirms(); await raw.close(); // The handler never ran (the body never decoded), and the message is on the dead queue — set // aside for inspection, not stuck redelivering forever. await delay(600); assert.equal(handlerRuns, 0, "a body that never decodes never reaches the handler"); const dead = await dch.get("mesh.events.dead", { noAck: true }); assert.ok(dead, "the poison event is retained on mesh.events.dead"); assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.poison.raw"); await drain.close(); await c.close(); delete process.env.MESH_NODE; delete process.env.MESH_MODULE; }); function delay(ms: number): Promise { return new Promise((r) => setTimeout(r, ms)); } async function waitFor(cond: () => boolean, ms: number): Promise { const start = Date.now(); while (!cond()) { if (Date.now() - start > ms) throw new Error("condition not met in time"); await delay(25); } }