// 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. This is where the // ADR 0042 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 * as tls from "node:tls"; import { randomUUID, createHash } from "node:crypto"; import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; // Two topic exchanges, kept apart on purpose (ADR 0042): 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 0042). const EVENT_PREFETCH = 32; interface Reply { result?: unknown; error?: string; } /** The broker presented a certificate whose fingerprint is not the one the mesh pinned. A distinct * type rather than a message to grep, so a caller deciding "wait or refuse" (serve mode's patient * reconnect, novox/hq issue 058) tells this apart from an absent broker by `instanceof`, not by a * prose string that a later reword would silently turn back into an infinite retry against an * impostor. */ export class PinMismatchError extends Error {} /** * Why a broker connection failed in a way no amount of waiting will fix — or null when it is worth * retrying. Serve mode's patient reconnect (novox/hq issue 058) uses this to tell a permanent * fault from a broker that is merely not up yet. Three failures are permanent: * * - the certificate does not match the pin — an impostor does not become the broker by being * asked again (typed, so a reworded message cannot silently turn this back into a retry); * - the broker URL is not a URL — a malformed address never parses on the next try; * - the broker answered and refused the login — a wrong or revoked credential, not an absent * broker, and it will refuse the next attempt identically. * * Everything else — connection refused, timeout, DNS not resolving yet — is the overlay still * coming up, and is retried. */ export function fatalBrokerReason(err: unknown): string | null { if (err instanceof PinMismatchError) return "the broker's certificate does not match the pin"; const e = err as { code?: unknown; message?: unknown }; const code = typeof e?.code === "string" ? e.code : ""; const message = typeof e?.message === "string" ? e.message : String(err); if (code === "ERR_INVALID_URL" || /invalid url/i.test(message)) { return `the broker URL is not a URL (${message})`; } if (/access[-_ ]?refused|login was refused|handshake terminated|\b403\b/i.test(message)) { return `the broker refused the login (${message})`; } return null; } /** A broker credential as the mesh delivers it (novox/hq ADR 0043): an amqps URL, the fingerprint * of the certificate the broker must present, and the node and module the account is scoped to (so * the runtime names its queue as the mesh did). A plain string is a bootstrap URL. */ export interface Credential { url: string; fingerprint?: string; node?: string; module?: string; } /** * Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. * * A scoped module (novox/hq ADR 0043) passes `assumeExchanges: true`: its account may not declare * an exchange, and the foundation already owns them, so it declares only its own queue. A credential * carrying a fingerprint is dialled over amqps, pinned to exactly that certificate. */ export async function connectAmqp( target: string | Credential, opts: { assumeExchanges?: boolean } = {}, ): Promise { const cred: Credential = typeof target === "string" ? { url: target } : target; // This module's own name, for turning a local event name into this bus's routing key. From the // credential where the mesh issued one, and from the environment for an ad-hoc client that has no // credential of its own — the same two places subscribe already looks. const self = cred.module ?? process.env.MESH_MODULE ?? ""; const conn = cred.fingerprint ? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint)) : await amqp.connect(cred.url); // 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 0042). const ch = await conn.createConfirmChannel(); // The foundation owns the exchanges (ADR 0043). A bootstrap/admin connection declares them; a // scoped module assumes they exist and never tries — its account could not, and the dead-letter // queue behind the exchange is the foundation's to keep, not a module's. if (!opts.assumeExchanges) { await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true }); await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true }); 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 is set up lazily: a consumer-only module (the audit logger) never calls a tool, // and its scoped account may not declare the exclusive reply queue this would otherwise need. const pending = new Map void>(); let replyQueue: string | undefined; async function ensureReply(): Promise { if (replyQueue) return replyQueue; const { queue } = await ch.assertQueue("", { exclusive: true }); // Replies come back through the RPC exchange keyed by this queue's own name, not the default // exchange (novox/hq ADR 0047): a serving module's scoped account may write to mesh.rpc but not // the default exchange, which would let it publish into any queue on the broker. await ch.bindQueue(queue, RPC_EXCHANGE, queue); replyQueue = queue; await ch.consume( queue, (msg) => { if (!msg) return; const resolve = pending.get(msg.properties.correlationId); if (resolve) { pending.delete(msg.properties.correlationId); resolve(JSON.parse(msg.content.toString()) as Reply); } }, { noAck: true }, ); return queue; } // One durable event queue per consumer (ADR 0042: ..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 = localKeyFor(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 0042). 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 reply = await ensureReply(); const id = randomUUID(); const answered = new Promise((resolve, reject) => { const timer = setTimeout(() => { if (pending.delete(id)) reject(new Error(`request ${key} timed out`)); }, 30_000); pending.set(id, (r) => { clearTimeout(timer); if (r.error) reject(new Error(r.error)); else resolve(r.result as Res); }); }); ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), { correlationId: id, replyTo: reply, }); return answered; }, async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { // A durable, shared serve queue (ADR 0042): 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) => { if (!msg) return; void (async () => { let reply: Reply; try { reply = { result: await handler(JSON.parse(msg.content.toString()) as Req) }; } catch (err) { reply = { error: err instanceof Error ? err.message : String(err) }; } if (msg.properties.replyTo) { // Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped // account answers with write on mesh.rpc alone — never the default exchange (ADR 0047). ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { correlationId: msg.properties.correlationId, }); } ch.ack(msg); })(); }); return () => void ch.cancel(consumer.consumerTag); }, async publish(env: Envelope): Promise { 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 0042). 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, routingKeyFor(env.key, self), 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 0042 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`; if (opts.assumeExchanges) { // The mesh pre-declared this queue with its dead-letter when it issued the account: a // scoped account may not declare a dead-lettered queue itself (the broker refuses that // to a non-administrator). Passively check it is there, then bind and consume. await ch.checkQueue(name); } else { await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE }); } eventQueue = name; } await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, bindingFor(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, bindingFor(pattern)); const consumer = await ch.consume(queue, (msg) => { if (!msg) return; void (async () => { 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); }, async close(): Promise { await ch.close(); await conn.close(); }, }; } /** Normalise a certificate fingerprint to bare lower-case hex, dropping an `sha256:` prefix and * any colon grouping, so two spellings of the same fingerprint compare equal. */ function normalizeFingerprint(fingerprint: string): string { return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase(); } /** * Socket options that pin the broker to exactly the certificate whose fingerprint the mesh * delivered (novox/hq ADR 0043, as the builder does). Done in two phases so a credential never * reaches an impostor: first a bare TLS connection that sends nothing fetches the certificate and * the fingerprint is checked; only then does the real connection trust *that* certificate as its * own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key. * Node's `checkServerIdentity` does not run under `rejectUnauthorized: false`, so a one-phase * "connect then check" would have already sent the password to whoever answered. */ async function pinnedOptions(rawUrl: string, fingerprint: string): Promise { const url = new URL(rawUrl); const host = url.hostname; const port = url.port ? Number(url.port) : 5671; const certificate = await new Promise((resolve, reject) => { const socket = tls.connect({ host, port, servername: host, rejectUnauthorized: false }, () => { const peer = socket.getPeerCertificate(true); socket.destroy(); if (!peer || !peer.raw) reject(new Error("the broker presented no certificate to pin")); else resolve(peer); }); socket.setTimeout(15_000, () => { socket.destroy(); reject(new Error("timed out fetching the broker's certificate")); }); socket.on("error", reject); }); const seen = createHash("sha256").update(certificate.raw).digest("hex"); if (seen !== normalizeFingerprint(fingerprint)) { throw new PinMismatchError( `the broker's certificate (sha256:${seen}) does not match the pinned ${fingerprint} — refusing`, ); } const pem = "-----BEGIN CERTIFICATE-----\n" + (certificate.raw.toString("base64").match(/.{1,64}/g) ?? []).join("\n") + "\n-----END CERTIFICATE-----\n"; // Trust that one certificate and nothing else; the mesh's own name is not in any public store, // so identity is the pin, not the hostname — checkServerIdentity is satisfied deliberately. return { ca: [pem], checkServerIdentity: () => undefined }; } /** 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]; // `**` is the mesh's wildcard for the rest of a name; `#` is this bus's, accepted so a pattern // written either way behaves the same while both buses ship (novox/hq design 29 §1). if (tok === "#" || 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; } /** This bus spells an event as a routing key that repeats the emitter's name: `module..`. * A module names its events locally and the mesh derives where they land (novox/hq design 29 §1), so * the mapping lives here rather than in every module. * * **Why it exists at all.** Until 04-ISSUES/127 every module passed the routing key itself, which * worked on this bus and derived into a namespace nobody owns on the one being built. Converting the * modules to local names without this would have broken the mesh that is actually running. */ function routingKeyFor(key: string, self: string): string { return key.startsWith("module.") ? key : `module.${self}.${key}`; } /** A local pattern as this bus's binding. `**` is the mesh's wildcard for the rest of a name; here * that is `#`, and on the bus being built it is `>`. Neither spelling appears in a manifest. */ function bindingFor(pattern: string): string { const here = pattern.split(".").map((part) => (part === "**" ? "#" : part)).join("."); if (here === "#") return "#"; return here.startsWith("module.") ? here : `module.${here}`; } /** A routing key as the local name a handler and a manifest both use: the emitter and the event. */ function localKeyFor(routingKey: string): string { return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey; }