broker: the ADR 0047 wire shape, and the event failure paths
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
This commit is contained in:
+160
-8
@@ -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<Broker> {
|
||||
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<string, (r: Reply) => void>();
|
||||
@@ -40,6 +58,40 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
||||
{ noAck: true },
|
||||
);
|
||||
|
||||
// One durable event queue per consumer (ADR 0047: <node>.<module>.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<unknown>) => Promise<void> };
|
||||
const eventSubs: EventSub[] = [];
|
||||
let eventQueue: string | undefined;
|
||||
let eventConsumerTag: string | undefined;
|
||||
|
||||
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
|
||||
const key = msg.fields.routingKey;
|
||||
let env: Envelope<unknown>;
|
||||
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<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||
const id = randomUUID();
|
||||
@@ -61,6 +113,8 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
||||
},
|
||||
|
||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): 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<Broker> {
|
||||
},
|
||||
|
||||
async publish<T>(env: Envelope<T>): Promise<void> {
|
||||
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<void>((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<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): 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<T>;
|
||||
try {
|
||||
env = toEnvelope<T>(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<Broker> {
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */
|
||||
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
||||
const raw = msg.properties.headers ?? {};
|
||||
const headers: Record<string, string> = {};
|
||||
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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user