The AMQP adapter now honours the full contract: events published persistent with metadata in headers; a durable per-consumer queue (<node>.<module>.events) with prefetch and a dead-letter exchange (mesh.events.dead); manual ack for at-least-once. Failure paths, not just the happy one: - a confirm channel, so a publish the broker never accepted fails the emit rather than vanishing — at-least-once starts at the emitter; - a handler that keeps failing is requeued once, then dead-lettered (poison set aside, never looping); - an undecodable body is dead-lettered at once — it never decodes on redelivery, and must not wedge the queue. Binding-conformance tests against a disposable broker (a stand-in for the mesh-hosted broker, ADR 0001): headers on the wire with a pure body, the redelivery-limit dead-letter, and the poison-body dead-letter. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
143 lines
5.7 KiB
TypeScript
143 lines
5.7 KiB
TypeScript
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 (<node>.<module>.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<void> {
|
|
return new Promise((r) => setTimeout(r, ms));
|
|
}
|
|
|
|
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
|
const start = Date.now();
|
|
while (!cond()) {
|
|
if (Date.now() - start > ms) throw new Error("condition not met in time");
|
|
await delay(25);
|
|
}
|
|
}
|