From f85d1dc303f5117058a893fbe9352d73adc5d1c4 Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 00:26:09 +0200 Subject: [PATCH 1/5] 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); + } +} From 99ce1e12523633bd8717b3a3e3740f8435831144 Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 00:56:37 +0200 Subject: [PATCH 2/5] runtime: an emit primitive, so an events test can put a message on the wire 'mesh-tools emit [json]' connects, emits one ADR 0047 event (awaiting the publish confirm), and exits. The serve path already runs a module's on('#') subscription as an import side effect, so the runtime hosts both an emitter and the audit-logger consumer. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- src/main.ts | 50 +++++++++++++++++++++++++++++++++++++++++++++----- 1 file changed, 45 insertions(+), 5 deletions(-) diff --git a/src/main.ts b/src/main.ts index 4c17be6..c3bcc4c 100644 --- a/src/main.ts +++ b/src/main.ts @@ -1,13 +1,21 @@ -// The runnable entrypoint. Reads its configuration from the environment the host resolved for it, -// connects the mesh broker, and serves the assigned modules' tools until stopped. +// The runnable entrypoint. Two modes: // -// MESH_BROKER_URL amqp://… the mesh broker -// MESH_TOOL_MODULES /path/a,/path/b,… compiled tool entrypoints of the assigned modules +// mesh-tools serve — bind the broker and serve the assigned modules until +// stopped. A module entrypoint that subscribes to events (on("#")) +// starts consuming as it is imported, so this also runs consumers. +// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive, +// and what an events test uses to put a message on the wire. +// +// MESH_BROKER_URL amqp://… the mesh broker (both modes) +// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode) +// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047) import { connectAmqp } from "./broker-amqp.js"; import { runTools } from "./runtime.js"; +import { useBroker } from "@novox/mesh-sdk/messaging"; +import { emit } from "@novox/mesh-sdk/events"; -async function main(): Promise { +async function serve(): Promise { const url = requireEnv("MESH_BROKER_URL"); const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "") .split(",") @@ -26,6 +34,38 @@ async function main(): Promise { process.on("SIGINT", () => void shutdown()); } +async function emitOnce(type: string, bodyJson: string): Promise { + const url = requireEnv("MESH_BROKER_URL"); + let body: unknown = {}; + if (bodyJson) { + try { + body = JSON.parse(bodyJson); + } catch { + console.error(`mesh-tools emit: body is not JSON: ${bodyJson}`); + process.exit(1); + } + } + const broker = await connectAmqp(url); + useBroker(() => broker); + // emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close. + await emit(type, body); + await broker.close(); +} + +async function main(): Promise { + const [command, ...rest] = process.argv.slice(2); + if (command === "emit") { + const type = rest[0]; + if (!type) { + console.error("mesh-tools emit [json-body] — a routing key is required"); + process.exit(1); + } + await emitOnce(type, rest[1] ?? ""); + return; + } + await serve(); +} + function requireEnv(name: string): string { const v = process.env[name]; if (!v) { From 04a689e008dc05d090bf010a241c89cd081cb1a1 Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 01:51:17 +0200 Subject: [PATCH 3/5] runtime: connect with a sealed credential, scoped, over pinned amqps (ADR 0048) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A module reads its broker credential from MESH_BROKER_FILE — the sealed {url,fingerprint} the mesh delivered — and connects over amqps pinned to exactly that certificate. The pin is two-phase (fetch cert, verify, then trust only it), because Node's checkServerIdentity does not run under rejectUnauthorized:false, so a naive connect-then-check would already have sent the password to whoever answered. A scoped module (assumeExchanges) never declares the exchanges (its account may not) nor its own queue with a dead-letter (the broker refuses that to a non-administrator) — the mesh pre-declared the queue, so it passively checks it, binds and consumes. The RPC reply queue is lazy, and a module that registered no tools serves none: a pure-events consumer touches only what its account allows. Verified end-to-end against a real broker as the scoped account: the audit logger consumes # and records events, over an account that is not the broker's own. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- dlxprobe.mjs | 48 +++++++++++++++ src/broker-amqp.ts | 144 ++++++++++++++++++++++++++++++++++++--------- src/main.ts | 58 +++++++++++++----- src/runtime.ts | 5 +- 4 files changed, 210 insertions(+), 45 deletions(-) create mode 100644 dlxprobe.mjs diff --git a/dlxprobe.mjs b/dlxprobe.mjs new file mode 100644 index 0000000..fd68229 --- /dev/null +++ b/dlxprobe.mjs @@ -0,0 +1,48 @@ +import amqp from "amqplib"; + +const PORT = process.argv[2]; +const MPORT = process.argv[3]; +const B = `http://127.0.0.1:${MPORT}`; +const AUTH = "Basic " + Buffer.from("guest:guest").toString("base64"); + +async function api(method, path, body) { + const r = await fetch(B + path, { + method, + headers: { "content-type": "application/json", authorization: AUTH }, + body: body ? JSON.stringify(body) : undefined, + }); + if (r.status >= 300 && r.status !== 404) throw new Error(`${method} ${path} -> ${r.status}`); +} + +await api("PUT", "/api/exchanges/%2f/mesh.events.dead", { type: "topic", durable: true }); +await api("PUT", "/api/users/al", { password: "s", tags: "" }); + +const Q = "anchor.al.events"; +const D = "mesh.events.dead"; +const q = Q.replace(/\./g, "\\."); +const d = D.replace(/\./g, "\\."); + +// configure, write, read patterns per grant on the dead exchange +const combos = { + "none": { configure: `^${q}$`, write: `^${q}$`, read: `^${q}$` }, + "read-dead": { configure: `^${q}$`, write: `^${q}$`, read: `^(${q}|${d})$` }, + "write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^${q}$` }, + "configure-dead": { configure: `^(${q}|${d})$`, write: `^${q}$`, read: `^${q}$` }, + "read+write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^(${q}|${d})$` }, +}; + +let i = 0; +for (const [label, perms] of Object.entries(combos)) { + await api("PUT", "/api/permissions/%2f/al", perms); + const queue = `${Q}.${i++}`; // fresh each time + try { + const c = await amqp.connect(`amqp://al:s@127.0.0.1:${PORT}/`); + const ch = await c.createChannel(); + ch.on("error", () => {}); + await ch.assertQueue(queue, { durable: true, deadLetterExchange: D }); + console.log(`${label}: declare-with-DLX OK`); + await c.close(); + } catch (e) { + console.log(`${label}: FAIL - ${String(e.message).slice(0, 70)}`); + } +} diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index 3fc9be3..ef3c183 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -5,7 +5,8 @@ // prefetch, dead-letter — none of which a module ever sees. import amqp from "amqplib"; -import { randomUUID } from "node:crypto"; +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 0047): tool invocations are request/reply and are @@ -23,40 +24,69 @@ interface Reply { error?: string; } -/** 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); +/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL and the + * fingerprint of the certificate the broker must present. A plain string is a bootstrap URL. */ +export interface Credential { + url: string; + fingerprint?: string; +} + +/** + * Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. + * + * A scoped module (novox/hq ADR 0048) passes `assumeExchanges: true`: its account may not declare + * an exchange, and the substrate 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; + 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 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, "#"); + // The substrate owns the exchanges (ADR 0048). 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 substrate'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: one exclusive reply queue, correlationId → resolver. - const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true }); + // 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>(); - await ch.consume( - replyQueue, - (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 }, - ); + let replyQueue: string | undefined; + async function ensureReply(): Promise { + if (replyQueue) return replyQueue; + const { queue } = await ch.assertQueue("", { exclusive: true }); + 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 0047: ..events), with many bindings and // a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message @@ -94,6 +124,7 @@ export async function connectAmqp(url: string): Promise { return { async request(key: string, body: Req): Promise { + const reply = await ensureReply(); const id = randomUUID(); const answered = new Promise((resolve, reject) => { const timer = setTimeout(() => { @@ -107,7 +138,7 @@ export async function connectAmqp(url: string): Promise { }); ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), { correlationId: id, - replyTo: replyQueue, + replyTo: reply, }); return answered; }, @@ -168,7 +199,14 @@ export async function connectAmqp(url: string): Promise { if (node && mod) { if (!eventQueue) { const name = `${node}.${mod}.events`; - await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE }); + 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, pattern); @@ -216,6 +254,56 @@ export async function connectAmqp(url: string): Promise { }; } +/** 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 0048, 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 Error( + `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 ?? {}; diff --git a/src/main.ts b/src/main.ts index c3bcc4c..1e98b6d 100644 --- a/src/main.ts +++ b/src/main.ts @@ -6,23 +6,59 @@ // mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive, // and what an events test uses to put a message on the wire. // -// MESH_BROKER_URL amqp://… the mesh broker (both modes) -// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode) -// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047) +// The broker, in order of preference: +// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0048) — an +// amqps account scoped to this module. Preferred: a module holds its own. +// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account. +// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode) +// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047) +import { readFileSync } from "node:fs"; import { connectAmqp } from "./broker-amqp.js"; +import type { Credential } from "./broker-amqp.js"; import { runTools } from "./runtime.js"; import { useBroker } from "@novox/mesh-sdk/messaging"; +import type { Broker } from "@novox/mesh-sdk/messaging"; import { emit } from "@novox/mesh-sdk/events"; +/** + * Connect the way this process is meant to: with its sealed credential if the mesh gave it one, and + * over the plain bootstrap URL otherwise. A scoped module assumes the substrate's exchanges exist — + * its account may not declare them (ADR 0048). + */ +async function connectBroker(): Promise { + const file = process.env.MESH_BROKER_FILE; + if (file) { + let credential: Credential; + try { + credential = JSON.parse(readFileSync(file, "utf8")) as Credential; + } catch (err) { + console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`); + process.exit(1); + } + if (!credential.url) { + console.error(`mesh-tools: ${file} carries no url — it is not a broker credential`); + process.exit(1); + } + return connectAmqp(credential, { assumeExchanges: true }); + } + const url = process.env.MESH_BROKER_URL; + if (!url) { + console.error( + "mesh-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach", + ); + process.exit(1); + } + return connectAmqp(url); +} + async function serve(): Promise { - const url = requireEnv("MESH_BROKER_URL"); const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "") .split(",") .map((s) => s.trim()) .filter(Boolean); - const broker = await connectAmqp(url); + const broker = await connectBroker(); const stop = await runTools({ broker, moduleEntrypoints }); const shutdown = async (): Promise => { @@ -35,7 +71,6 @@ async function serve(): Promise { } async function emitOnce(type: string, bodyJson: string): Promise { - const url = requireEnv("MESH_BROKER_URL"); let body: unknown = {}; if (bodyJson) { try { @@ -45,7 +80,7 @@ async function emitOnce(type: string, bodyJson: string): Promise { process.exit(1); } } - const broker = await connectAmqp(url); + const broker = await connectBroker(); useBroker(() => broker); // emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close. await emit(type, body); @@ -66,13 +101,4 @@ async function main(): Promise { await serve(); } -function requireEnv(name: string): string { - const v = process.env[name]; - if (!v) { - console.error(`mesh-tools: ${name} is not set — the runtime cannot serve without it`); - process.exit(1); - } - return v; -} - void main(); diff --git a/src/runtime.ts b/src/runtime.ts index d43b33e..80d9e04 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -25,8 +25,11 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> { await import(pathToFileURL(resolve(entry)).href); } - const stop = await serveTools(opts.broker); + // Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the + // audit logger) registers none, and its scoped account may not declare the serve queue — so a + // runtime that always served would fail for exactly the modules that never needed it. const tools = listTools(); + const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {}; console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`); return stop; } From bf5339cec2cd1ccd2048e61b56e902c94844ceff Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 01:53:23 +0200 Subject: [PATCH 4/5] runtime: take node and module identity from the sealed credential The mesh scoped the account to a node and module; the credential now carries both, so the runtime names its queue and stamps its events as the mesh authorised without a manifest interpolating a node the vocabulary has no token for. Verified: with only MESH_BROKER_FILE, the audit logger consumed as anchor/audit-logger. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- src/broker-amqp.ts | 7 +++++-- src/main.ts | 5 +++++ 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index ef3c183..78c57dd 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -24,11 +24,14 @@ interface Reply { error?: string; } -/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL and the - * fingerprint of the certificate the broker must present. A plain string is a bootstrap URL. */ +/** A broker credential as the mesh delivers it (novox/hq ADR 0048): 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; } /** diff --git a/src/main.ts b/src/main.ts index 1e98b6d..bea038a 100644 --- a/src/main.ts +++ b/src/main.ts @@ -40,6 +40,11 @@ async function connectBroker(): Promise { console.error(`mesh-tools: ${file} carries no url — it is not a broker credential`); process.exit(1); } + // The mesh scoped this account to a node and module; take the runtime's identity from the + // credential so its queue and the events it emits match what the mesh authorised, no matter + // what the environment says. + if (credential.node) process.env.MESH_NODE = credential.node; + if (credential.module) process.env.MESH_MODULE = credential.module; return connectAmqp(credential, { assumeExchanges: true }); } const url = process.env.MESH_BROKER_URL; From 4558248f166804474f732753445bfe7d52ba7b2d Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 4 Sep 2026 21:56:04 +0200 Subject: [PATCH 5/5] runtime: RPC replies ride mesh.rpc; an invoke subcommand (ADR 0052) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replies go through the RPC exchange keyed by the caller's reply-queue name, not the default exchange — so a serving module's scoped account answers with write on mesh.rpc alone, never the default exchange (which would let it publish into any queue). 'mesh-tools invoke [args]' is the caller's side, the sibling of emit. Verified against a real broker: a scoped account serves its tool and is refused another module's serve queue. --- src/broker-amqp.ts | 8 +++++++- src/main.ts | 26 ++++++++++++++++++++++++++ test/amqp.test.ts | 13 ++++--------- 3 files changed, 37 insertions(+), 10 deletions(-) diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index 78c57dd..3bd8e76 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -75,6 +75,10 @@ export async function connectAmqp( 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 0052): 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, @@ -161,7 +165,9 @@ export async function connectAmqp( reply = { error: err instanceof Error ? err.message : String(err) }; } if (msg.properties.replyTo) { - ch.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { + // 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 0052). + ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { correlationId: msg.properties.correlationId, }); } diff --git a/src/main.ts b/src/main.ts index bea038a..769b6a5 100644 --- a/src/main.ts +++ b/src/main.ts @@ -17,6 +17,7 @@ import { readFileSync } from "node:fs"; import { connectAmqp } from "./broker-amqp.js"; import type { Credential } from "./broker-amqp.js"; import { runTools } from "./runtime.js"; +import { invokeTool } from "@novox/mesh-sdk/tools"; import { useBroker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging"; import { emit } from "@novox/mesh-sdk/events"; @@ -92,8 +93,33 @@ async function emitOnce(type: string, bodyJson: string): Promise { await broker.close(); } +async function invokeOnce(module: string, tool: string, argsJson: string): Promise { + let args: Record = {}; + if (argsJson) { + try { + args = JSON.parse(argsJson) as Record; + } catch { + console.error(`mesh-tools invoke: args are not JSON: ${argsJson}`); + process.exit(1); + } + } + const broker = await connectBroker(); + const result = await invokeTool(broker, module, tool, args); + process.stdout.write(JSON.stringify(result) + "\n"); + await broker.close(); +} + async function main(): Promise { const [command, ...rest] = process.argv.slice(2); + if (command === "invoke") { + const [module, tool] = rest; + if (!module || !tool) { + console.error("mesh-tools invoke [json-args] — a module and tool are required"); + process.exit(1); + } + await invokeOnce(module, tool, rest[2] ?? ""); + return; + } if (command === "emit") { const type = rest[0]; if (!type) { diff --git a/test/amqp.test.ts b/test/amqp.test.ts index 205eafb..8a406e4 100644 --- a/test/amqp.test.ts +++ b/test/amqp.test.ts @@ -1,7 +1,7 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { registerModuleTools, resetTools } from "@novox/mesh-sdk/tools"; +import { registerModuleTools, resetTools, invokeTool } from "@novox/mesh-sdk/tools"; import { connectAmqp } from "../dist/broker-amqp.js"; import { runTools } from "../dist/runtime.js"; @@ -27,17 +27,12 @@ test("a module's tool serves and is invoked over a real AMQP broker", { skip: !u const serverBroker = await connectAmqp(url!); const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] }); - // A separate connection — a caller, like mesh-control's command API — invokes over the broker. + // A separate connection — a caller, like mesh-control's command API — invokes over the broker, + // by module and tool (novox/hq ADR 0052: served on serve.demo.greet, invoked as demo.greet). const caller = await connectAmqp(url!); - const result = await caller.request<{ tool: string; args: Record }, { hello: string }>( - "tools.invoke", - { tool: "greet", args: { who: "mesh" } }, - ); + const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string }; assert.equal(result.hello, "mesh"); - // An unknown tool is refused over the wire, not silently dropped. - await assert.rejects(caller.request("tools.invoke", { tool: "nope", args: {} }), /no such tool/); - stop(); await caller.close(); await serverBroker.close();