From f85d1dc303f5117058a893fbe9352d73adc5d1c4 Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 00:26:09 +0200 Subject: [PATCH] broker: the ADR 0047 wire shape, and the event failure paths MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The AMQP adapter now honours the full contract: events published persistent with metadata in headers; a durable per-consumer queue (..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 --- src/broker-amqp.ts | 168 +++++++++++++++++++++++++++++++++++++-- test/events-wire.test.ts | 142 +++++++++++++++++++++++++++++++++ 2 files changed, 302 insertions(+), 8 deletions(-) create mode 100644 test/events-wire.test.ts diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index b31a5eb..3fc9be3 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -1,16 +1,22 @@ // A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker // (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime -// provides the binding — so a broker-client change never rebuilds the modules. +// provides the binding — so a broker-client change never rebuilds the modules. This is where the +// ADR 0047 wire shape lives: the two exchanges, persistent events, per-consumer durable queues, +// prefetch, dead-letter — none of which a module ever sees. import amqp from "amqplib"; import { randomUUID } from "node:crypto"; -import type { Broker, Envelope } from "@novox/mesh-sdk/messaging"; +import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; -// Two topic exchanges, kept apart on purpose: tool invocations are request/reply and are not -// events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node +// Two topic exchanges, kept apart on purpose (ADR 0047): tool invocations are request/reply and are +// not events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node // events — never the RPC traffic. const RPC_EXCHANGE = "mesh.rpc"; const EVENTS_EXCHANGE = "mesh.events"; +// Where an event rejected past its redelivery limit is set aside for inspection. +const DEAD_EXCHANGE = "mesh.events.dead"; +// Bound in-flight events so one slow consumer can't pull the whole backlog into memory (ADR 0047). +const EVENT_PREFETCH = 32; interface Reply { result?: unknown; @@ -20,10 +26,22 @@ interface Reply { /** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */ export async function connectAmqp(url: string): Promise { const conn = await amqp.connect(url); - const ch = await conn.createChannel(); + // A confirm channel, so an event publish awaits the broker's ack: a publish the broker never + // accepted (it was mid-restart, the connection dropped) fails the emit rather than vanishing — + // at-least-once starts at the emitter, not only the consumer (ADR 0047). + const ch = await conn.createConfirmChannel(); await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true }); await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true }); + // The dead-letter home for poison events. A durable queue bound to `#` retains them for + // inspection — a dead-letter exchange with no queue behind it would drop them silently, which is + // exactly the loss the audit trail exists to prevent. + await ch.assertExchange(DEAD_EXCHANGE, "topic", { durable: true }); + await ch.assertQueue(DEAD_EXCHANGE, { durable: true }); + await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#"); + + await ch.prefetch(EVENT_PREFETCH); + // Request/reply: one exclusive reply queue, correlationId → resolver. const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true }); const pending = new Map void>(); @@ -40,6 +58,40 @@ export async function connectAmqp(url: string): Promise { { noAck: true }, ); + // One durable event queue per consumer (ADR 0047: ..events), with many bindings and + // a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message + // once however many bindings match, so the local match is what keeps a two-pattern module from + // running the wrong handler. + type EventSub = { pattern: string; handler: (env: Envelope) => Promise }; + const eventSubs: EventSub[] = []; + let eventQueue: string | undefined; + let eventConsumerTag: string | undefined; + + async function dispatchEvent(msg: amqp.ConsumeMessage): Promise { + const key = msg.fields.routingKey; + let env: Envelope; + try { + env = toEnvelope(msg); + } catch { + // An undecodable body will never decode on redelivery — dead-letter it at once rather than + // wedge the queue or loop. Decoding sits before the handler try on purpose: a poison message + // is a different failure from a handler that threw, and gets no retry. + ch.nack(msg, false, false); + return; + } + try { + for (const s of eventSubs) { + if (topicMatches(s.pattern, key)) await s.handler(env); + } + ch.ack(msg); + } catch { + // First handler failure: requeue once. A second (already redelivered) dead-letters it, so a + // poison event is set aside rather than looping forever or vanishing (ADR 0047). Redelivery + // re-runs every matching handler, so a consumer must be idempotent — which the ADR requires. + ch.nack(msg, false, !msg.fields.redelivered); + } + } + return { async request(key: string, body: Req): Promise { const id = randomUUID(); @@ -61,6 +113,8 @@ export async function connectAmqp(url: string): Promise { }, async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { + // A durable, shared serve queue (ADR 0047): several runtimes serving one tool key compete for + // invocations rather than each answering the same call. const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true }); await ch.bindQueue(queue, RPC_EXCHANGE, key); const consumer = await ch.consume(queue, (msg) => { @@ -84,17 +138,72 @@ export async function connectAmqp(url: string): Promise { }, async publish(env: Envelope): Promise { - ch.publish(EVENTS_EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body))); + const headers = env.headers ?? ({} as EventHeaders); + // Events are persistent (delivery-mode 2): an audit trail that loses events on a broker + // restart is not one (ADR 0047). Metadata rides as headers; the body is only the payload. + // The publish is awaited to the broker's confirm — an unaccepted publish rejects here. + await new Promise((resolve, reject) => { + ch.publish( + EVENTS_EXCHANGE, + env.key, + Buffer.from(JSON.stringify(env.body)), + { + persistent: true, + contentType: + typeof headers["content-type"] === "string" ? headers["content-type"] : "application/json", + messageId: headers["x-event-id"], + headers: { ...headers }, + }, + (err) => (err ? reject(err instanceof Error ? err : new Error(String(err))) : resolve()), + ); + }); }, async subscribe(pattern: string, handler: (env: Envelope) => Promise): Promise<() => void> { + const node = process.env.MESH_NODE; + const mod = process.env.MESH_MODULE; + + // A module we can name gets its ADR 0047 durable queue; an anonymous subscriber (a test, an + // ad-hoc listener) gets a transient exclusive one that dies with the connection. + if (node && mod) { + if (!eventQueue) { + const name = `${node}.${mod}.events`; + await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE }); + eventQueue = name; + } + await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern); + const sub: EventSub = { pattern, handler: handler as EventSub["handler"] }; + eventSubs.push(sub); + if (!eventConsumerTag) { + const consumer = await ch.consume(eventQueue, (msg) => { + if (msg) void dispatchEvent(msg); + }); + eventConsumerTag = consumer.consumerTag; + } + return () => { + const i = eventSubs.indexOf(sub); + if (i >= 0) eventSubs.splice(i, 1); + }; + } + const { queue } = await ch.assertQueue("", { exclusive: true }); await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern); const consumer = await ch.consume(queue, (msg) => { if (!msg) return; void (async () => { - await handler({ key: msg.fields.routingKey, node: "", body: JSON.parse(msg.content.toString()) as T }); - ch.ack(msg); + let env: Envelope; + try { + env = toEnvelope(msg); + } catch { + ch.nack(msg, false, false); // undecodable — drop, never retry + return; + } + try { + await handler(env); + ch.ack(msg); + } catch { + ch.nack(msg, false, !msg.fields.redelivered); + } })(); }); return () => void ch.cancel(consumer.consumerTag); @@ -106,3 +215,46 @@ export async function connectAmqp(url: string): Promise { }, }; } + +/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */ +function toEnvelope(msg: amqp.ConsumeMessage): Envelope { + const raw = msg.properties.headers ?? {}; + const headers: Record = {}; + for (const [k, v] of Object.entries(raw)) { + if (v == null) continue; + headers[k] = typeof v === "string" ? v : String(v); + } + if (!headers["content-type"] && msg.properties.contentType) { + headers["content-type"] = msg.properties.contentType; + } + return { + key: msg.fields.routingKey, + node: headers["x-node"] ?? "", + body: JSON.parse(msg.content.toString()) as T, + headers: headers as EventHeaders, + }; +} + +/** AMQP topic matching: `*` matches one word, `#` zero or more. Used to fan a shared queue's + * deliveries out to the handlers whose pattern actually matches the routing key. */ +function topicMatches(pattern: string, key: string): boolean { + return matchFrom(pattern.split("."), 0, key.split("."), 0); +} + +function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { + while (pi < p.length) { + const tok = p[pi]; + if (tok === "#") { + if (pi === p.length - 1) return true; // trailing # swallows the rest, including nothing + for (let skip = ki; skip <= k.length; skip++) { + if (matchFrom(p, pi + 1, k, skip)) return true; + } + return false; + } + if (ki >= k.length) return false; + if (tok !== "*" && tok !== k[ki]) return false; + pi++; + ki++; + } + return ki === k.length; +} diff --git a/test/events-wire.test.ts b/test/events-wire.test.ts new file mode 100644 index 0000000..8e21c00 --- /dev/null +++ b/test/events-wire.test.ts @@ -0,0 +1,142 @@ +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); + } +}