diff --git a/package-lock.json b/package-lock.json new file mode 100644 index 0000000..311b62d --- /dev/null +++ b/package-lock.json @@ -0,0 +1,76 @@ +{ + "name": "@novox/mesh-tools", + "version": "0.1.0", + "lockfileVersion": 3, + "requires": true, + "packages": { + "": { + "name": "@novox/mesh-tools", + "version": "0.1.0", + "dependencies": { + "@novox/mesh-sdk": "^0.1.0", + "nats": "^2.29.0" + }, + "bin": { + "mesh": "dist/mesh.js", + "mesh-tools": "dist/main.js" + }, + "devDependencies": { + "@types/node": "^22.20.1", + "typescript": "^5.9.3" + } + }, + "node_modules/@novox/mesh-sdk": { + "version": "0.1.1" + }, + "node_modules/@types/node": { + "version": "22.20.1", + "dev": true, + "license": "MIT", + "dependencies": { + "undici-types": "~6.21.0" + } + }, + "node_modules/nats": { + "version": "2.29.3", + "license": "Apache-2.0", + "dependencies": { + "nkeys.js": "1.1.0" + }, + "engines": { + "node": ">= 14.0.0" + } + }, + "node_modules/nkeys.js": { + "version": "1.1.0", + "license": "Apache-2.0", + "dependencies": { + "tweetnacl": "1.0.3" + }, + "engines": { + "node": ">=10.0.0" + } + }, + "node_modules/tweetnacl": { + "version": "1.0.3", + "license": "Unlicense" + }, + "node_modules/typescript": { + "version": "5.9.3", + "dev": true, + "license": "Apache-2.0", + "bin": { + "tsc": "bin/tsc", + "tsserver": "bin/tsserver" + }, + "engines": { + "node": ">=14.17" + } + }, + "node_modules/undici-types": { + "version": "6.21.0", + "dev": true, + "license": "MIT" + } + } +} diff --git a/package.json b/package.json index 6e420bb..9831df0 100644 --- a/package.json +++ b/package.json @@ -14,11 +14,9 @@ }, "dependencies": { "@novox/mesh-sdk": "^0.1.0", - "amqplib": "^0.10.9", "nats": "^2.29.0" }, "devDependencies": { - "@types/amqplib": "^0.10.8", "@types/node": "^22.20.1", "typescript": "^5.9.3" } diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts deleted file mode 100644 index 9b8db7c..0000000 --- a/src/broker-amqp.ts +++ /dev/null @@ -1,422 +0,0 @@ -// 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. */ -export 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. */ -export 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. */ -export 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. */ -export function localKeyFor(routingKey: string): string { - return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey; -} diff --git a/src/broker-nats.ts b/src/broker-nats.ts index f6f4ff7..148715c 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -17,6 +17,7 @@ // correct. import { createHash } from "node:crypto"; +import net from "node:net"; import tls from "node:tls"; import { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription } from "nats"; import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; @@ -307,16 +308,32 @@ function normalizeFingerprint(fingerprint: string): string { async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> { const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`); const port = url.port ? Number(url.port) : 4222; + // **The bus speaks first, in the clear.** A NATS server sends its INFO line before TLS begins, + // and only then expects the client to start the handshake; a raw TLS connect to that port reads + // the INFO line as a TLS record and fails with "wrong version number" — which is what every + // module met the first time it dialled the bus being built (2026-09-28). So: connect, wait for + // INFO, then start TLS on the same socket, and read the certificate the server presents. const certificate = await new Promise((resolve, reject) => { - const socket = tls.connect( - { host: url.hostname, port, rejectUnauthorized: false, servername: url.hostname }, - () => { - const peer = socket.getPeerCertificate(true); - socket.end(); - resolve(peer); - }, - ); - socket.on("error", reject); + const plain = net.connect({ host: url.hostname, port }, () => {}); + let seenInfo = false; + let buffered = ""; + plain.on("error", reject); + plain.on("data", (chunk: Buffer) => { + if (seenInfo) return; + buffered += chunk.toString("utf8"); + if (!buffered.includes("\r\n")) return; + seenInfo = true; + plain.removeAllListeners("data"); + const secure = tls.connect( + { socket: plain, rejectUnauthorized: false, servername: url.hostname }, + () => { + const peer = secure.getPeerCertificate(true); + secure.end(); + resolve(peer); + }, + ); + secure.on("error", reject); + }); }); const seen = createHash("sha256").update(certificate.raw).digest("hex"); if (seen !== normalizeFingerprint(fingerprint)) { diff --git a/src/main.ts b/src/main.ts index 21f609f..e128825 100644 --- a/src/main.ts +++ b/src/main.ts @@ -22,9 +22,7 @@ import { readFileSync } from "node:fs"; import { pathToFileURL } from "node:url"; -import { connectAmqp, fatalBrokerReason as fatalAmqpReason } from "./broker-amqp.js"; -import { connectNats, fatalBrokerReason as fatalNatsReason } from "./broker-nats.js"; -import type { Credential } from "./broker-amqp.js"; +import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js"; import { runTools } from "./runtime.js"; import { invokeTool } from "@novox/mesh-sdk/tools"; import { useBroker } from "@novox/mesh-sdk/messaging"; @@ -39,7 +37,7 @@ import { emit } from "@novox/mesh-sdk/events"; // fatalBrokerReasonFor is the reason a connection failure is final rather than "not yet", for // whichever bus this runtime is on — each transport knows its own refusals. function fatalBrokerReasonFor(err: unknown): string | null { - return fatalNatsReason(err) ?? fatalAmqpReason(err); + return fatalNatsReason(err); } async function connectBroker(): Promise { @@ -66,10 +64,8 @@ async function connectBroker(): Promise { // nothing else in its environment changed (design 25; novox/hq design 28 task 5.2). The scheme // is enough to know which bus to speak; a runtime that always dialled the old one would keep // serving and answer nobody. - if (credential.url.startsWith("nats://")) { - return connectNats(credential); - } - return connectAmqp(credential, { assumeExchanges: true }); + // One bus (novox/hq ADR 0131, design 28 task 5.5): the credential names it, and it is this. + return connectNats(credential); } const url = process.env.MESH_BROKER_URL; if (!url) { @@ -78,7 +74,7 @@ async function connectBroker(): Promise { ); process.exit(1); } - return connectAmqp(url); + return connectNats({ url }); } /** diff --git a/test/amqp.test.ts b/test/amqp.test.ts deleted file mode 100644 index b20be06..0000000 --- a/test/amqp.test.ts +++ /dev/null @@ -1,39 +0,0 @@ -import { test } from "node:test"; -import assert from "node:assert/strict"; - -import { registerModuleTools, resetTools, invokeTool } from "@novox/mesh-sdk/tools"; -import { connectAmqp } from "../dist/broker-amqp.js"; -import { runTools } from "../dist/runtime.js"; - -// Requires a real broker at $MESH_BROKER_URL. The test harness spins LavinMQ (the mesh's broker) -// and points this at it; skipped if it is not set, never failed for the environment. -const url = process.env.MESH_BROKER_URL; - -test("a module's tool serves and is invoked over a real AMQP broker", { skip: !url }, async () => { - resetTools(); - - // A module registers a real tool, exactly as umami does. - registerModuleTools("demo", () => [ - { - name: "greet", - description: "return a greeting", - input: { who: { type: "string" } }, - run: async (args) => ({ hello: String(args.who), from: "the tool runtime" }), - }, - ]); - - // The runtime binds the broker and serves the module (no module entrypoints to import here — the - // tool is registered in-process — but this is the exact runtime that serves on a node). - const serverBroker = await connectAmqp(url!); - const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] }); - - // A separate connection — a caller, like mesh-controller's command API — invokes over the broker, - // by module and tool (novox/hq ADR 0047: served on serve.demo.greet, invoked as demo.greet). - const caller = await connectAmqp(url!); - const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string }; - assert.equal(result.hello, "mesh"); - - stop(); - await caller.close(); - await serverBroker.close(); -}); diff --git a/test/events-wire.test.ts b/test/events-wire.test.ts deleted file mode 100644 index 6cceeb4..0000000 --- a/test/events-wire.test.ts +++ /dev/null @@ -1,142 +0,0 @@ -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 0042 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 foundation. 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 0042 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); - } -} diff --git a/test/patient-connect.test.ts b/test/patient-connect.test.ts index ff8db7f..a2a687c 100644 --- a/test/patient-connect.test.ts +++ b/test/patient-connect.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { fatalBrokerReason, PinMismatchError } from "../src/broker-amqp.ts"; +import { fatalBrokerReason, PinMismatchError } from "../src/broker-nats.ts"; // novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must // give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever @@ -31,11 +31,8 @@ test("a malformed broker URL is fatal — it never parses on the next try", () = }); test("a refused login is fatal — a wrong or revoked credential, not an absent broker", () => { - for (const msg of [ - "Handshake terminated by server: 403 (ACCESS-REFUSED) with message \"ACCESS_REFUSED - Login was refused\"", - "Login was refused using authentication mechanism PLAIN", - "ACCESS_REFUSED", - ]) { + // The bus refuses a login in its own words; each is final, because the next try says the same. + for (const msg of ["Authorization Violation", "nats: user authentication expired", "Permissions Violation for Subscription to \"x\""]) { assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`); } }); diff --git a/test/wire-unchanged.test.ts b/test/wire-unchanged.test.ts deleted file mode 100644 index 2d80144..0000000 --- a/test/wire-unchanged.test.ts +++ /dev/null @@ -1,68 +0,0 @@ -/** - * **The old bus's wire is byte-identical after the rename, and this is the test that lets the change - * be merged to a running mesh.** - * - * Every module's event names were converted from the old bus's routing keys to local names - * (novox/hq 04-ISSUES/127), and the old bus's client maps them back. If that mapping is wrong - * anywhere, a live mesh's events stop being delivered — silently, because a binding that matches - * nothing is not an error. - * - * So this pins the mapping against the literal routing keys the mesh used before, taken from the - * manifests as they were. It needs no bus: it is about a string. - */ -import assert from "node:assert/strict"; -import { test } from "node:test"; - -import { routingKeyFor, bindingFor, localKeyFor, topicMatches } from "../dist/broker-amqp.js"; - -test("a converted emit produces the routing key the mesh published before", () => { - // left: what the module's code says now. right: what went on the wire before, unchanged. - const same: [string, string, string][] = [ - ["plex", "playback.started", "module.plex.playback.started"], - ["sonarr", "download.completed", "module.sonarr.download.completed"], - ["builder", "built", "module.builder.built"], - ["mesh-catalog", "upgraded", "module.mesh-catalog.upgraded"], - ["keycloak", "user.created", "module.keycloak.user.created"], - ["mesh-vault", "secret.rotated", "module.mesh-vault.secret.rotated"], - ]; - for (const [self, local, before] of same) { - assert.equal(routingKeyFor(local, self), before, `${self} emitting ${local}`); - } -}); - -test("a converted subscription binds what it bound before", () => { - const same: [string, string][] = [ - ["builder.built", "module.builder.built"], - ["*.download.completed", "module.*.download.completed"], - ["*.usage.*", "module.*.usage.*"], - // The audit logger's "everything": `#` on this bus, and it must stay `#`. - ["**", "#"], - ]; - for (const [declared, before] of same) { - assert.equal(bindingFor(declared), before, `consuming ${declared}`); - } -}); - -test("a handler still matches what the bus delivers", () => { - // The key a handler is given is the local one now, and the pattern it compares against is local - // too — so the pair must still meet for every case the mesh actually has. - const pairs: [string, string][] = [ - ["builder.built", "module.builder.built"], - ["*.download.completed", "module.sonarr.download.completed"], - ["*.usage.*", "module.anthropic-consumer.usage.session"], - ["**", "module.anything.at.all"], - ]; - for (const [pattern, delivered] of pairs) { - assert.ok( - topicMatches(pattern, localKeyFor(delivered)), - `${pattern} no longer matches ${delivered}, so a running module would stop reacting`, - ); - } -}); - -test("a routing key already in the old form is left alone", () => { - // Belt for the transition: anything not yet converted still goes out as it did, so a module built - // from an older manifest keeps working beside one built from a current manifest. - assert.equal(routingKeyFor("module.plex.playback.started", "plex"), "module.plex.playback.started"); - assert.equal(bindingFor("module.builder.built"), "module.builder.built"); -});