Author SHA1 Message Date
mesh-admin 10e8191717 Merge pull request 'The runtime hears on its own inbox' (#17) from fix/the-runtime-hears-on-its-own-inbox into main 2026-09-28 02:30:24 +00:00
jschoubben d703cebff4 The runtime hears on its own inbox
Every user's inbox is private to it and the grant names it; a reply space the client invented was
refused, and with it every pull for the next message and every answer to a tool call.
2026-09-28 04:30:22 +02:00
mesh-admin 46b56d53a6 Merge pull request 'The pin is the only check: the runtime stops verifying the bus's name' (#16) from fix/the-pin-is-the-only-check into main 2026-09-28 02:03:46 +00:00
jschoubben d4a2802342 The pin is the only check: the runtime stops verifying the bus's name
Every module on the new runtime reached the handshake and failed on 'does not match
certificate's altnames': the bus's certificate names the seat, not the address a machine dials
it by, and pinning the exact certificate already decides everything a name check could. The
client's transport spreads the TLS options into Node's tls.connect, so the hostname check is
replaced with one that passes and the pinned certificate is the one authority accepted.
2026-09-28 04:03:43 +02:00
mesh-admin 8acfa7a07d Merge pull request 'One bus: the runtime pins the certificate after the server speaks, and the old transport goes' (#15) from feat/one-bus into main 2026-09-28 01:09:01 +00:00
jschoubben f0104b7846 One bus: the runtime pins the certificate after the server speaks, and the old transport goes
Every module that dialled the new bus failed its handshake with "wrong version
number": the runtime pinned the server's certificate by a raw TLS connection to a
port on which the server speaks first, in the clear. The pin is taken after the
INFO line now, on the same socket, and then the real connection verifies against
exactly that certificate.

And the old transport is deleted — its client, its tests, its dependency — with
the wire-compatibility pins that only existed for the move (novox/hq ADR 0131,
design 28 task 5.5). A credential names the bus, and there is one.
2026-09-28 03:08:59 +02:00
mesh-admin cd26131c61 Merge pull request 'The runtime speaks the bus its credential names' (#14) from feat/the-runtime-speaks-the-bus-its-credential-names into main 2026-09-28 00:42:29 +00:00
jschoubben 9e3ff6fa45 The runtime speaks the bus its credential names
A module moved to the bus being built was handed a credential for it — nats://
with user, password and fingerprint beside the address — and nothing else in its
environment changed. The runtime always dialled the old bus, so every moved module
kept serving and answered nobody. The scheme in the credential is enough to know
which bus to speak; the nats broker was already written and never chosen.
2026-09-28 02:42:26 +02:00
jschoubben a4447f1251 Merge pull request 'A person's own client, and pins that the wire did not change' (#13) from feat/nats-genesis into main 2026-09-27 17:20:07 +00:00
9 changed files with 137 additions and 702 deletions
+76
View File
@@ -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"
}
}
}
-2
View File
@@ -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"
}
-422
View File
@@ -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<Broker> {
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<string, (r: Reply) => void>();
let replyQueue: string | undefined;
async function ensureReply(): Promise<string> {
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: <node>.<module>.events), with many bindings and
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
// once however many bindings match, so the local match is what keeps a two-pattern module from
// running the wrong handler.
type EventSub = { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> };
const eventSubs: EventSub[] = [];
let eventQueue: string | undefined;
let eventConsumerTag: string | undefined;
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
const key = localKeyFor(msg.fields.routingKey);
let env: Envelope<unknown>;
try {
env = toEnvelope(msg);
} catch {
// An undecodable body will never decode on redelivery — dead-letter it at once rather than
// wedge the queue or loop. Decoding sits before the handler try on purpose: a poison message
// is a different failure from a handler that threw, and gets no retry.
ch.nack(msg, false, false);
return;
}
try {
for (const s of eventSubs) {
if (topicMatches(s.pattern, key)) await s.handler(env);
}
ch.ack(msg);
} catch {
// First handler failure: requeue once. A second (already redelivered) dead-letters it, so a
// poison event is set aside rather than looping forever or vanishing (ADR 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<Req, Res>(key: string, body: Req): Promise<Res> {
const reply = await ensureReply();
const id = randomUUID();
const answered = new Promise<Res>((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<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): 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<T>(env: Envelope<T>): Promise<void> {
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<void>((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<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
const node = process.env.MESH_NODE;
const mod = process.env.MESH_MODULE;
// A module we can name gets its ADR 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<T>;
try {
env = toEnvelope<T>(msg);
} catch {
ch.nack(msg, false, false); // undecodable — drop, never retry
return;
}
try {
await handler(env);
ch.ack(msg);
} catch {
ch.nack(msg, false, !msg.fields.redelivered);
}
})();
});
return () => void ch.cancel(consumer.consumerTag);
},
async close(): Promise<void> {
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<tls.ConnectionOptions> {
const url = new URL(rawUrl);
const host = url.hostname;
const port = url.port ? Number(url.port) : 5671;
const certificate = await new Promise<tls.DetailedPeerCertificate>((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<T>(msg: amqp.ConsumeMessage): Envelope<T> {
const raw = msg.properties.headers ?? {};
const headers: Record<string, string> = {};
for (const [k, v] of Object.entries(raw)) {
if (v == null) continue;
headers[k] = typeof v === "string" ? v : String(v);
}
if (!headers["content-type"] && msg.properties.contentType) {
headers["content-type"] = msg.properties.contentType;
}
return {
key: msg.fields.routingKey,
node: headers["x-node"] ?? "",
body: JSON.parse(msg.content.toString()) as T,
headers: headers as EventHeaders,
};
}
/** AMQP topic matching: `*` matches one word, `#` zero or more. Used to fan a shared queue's
* deliveries out to the handlers whose pattern actually matches the routing key. */
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.<emitter>.<event>`.
* 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;
}
+42 -18
View File
@@ -17,8 +17,9 @@
// 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 { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription, type TlsOptions } from "nats";
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
const sc = StringCodec();
@@ -84,6 +85,10 @@ export async function connectNats(
pass: cred.password,
name: `${cred.node ?? "?"}.${self}`,
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
// Its own inbox, not a random one: every user's inbox is private to it (design 25 §4), and the
// grant names `_INBOX.<user>.>` — a reply space the client invented would be refused, and with
// it every pull for the next message and every answer to a tool call.
inboxPrefix: cred.user ? `_INBOX.${cred.user}` : undefined,
// Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on
// the mesh to exit. `close()` stays the only thing that ends the connection.
maxReconnectAttempts: -1,
@@ -297,26 +302,43 @@ function normalizeFingerprint(fingerprint: string): string {
* A certificate authority is not consulted: the mesh issued this and knows its fingerprint,
* which is stronger than trusting whoever a machine's trust store happens to contain.
*
* **A constraint on the mesh, not a detail of this file.** Pinning the exact certificate makes
* hostname verification redundant in principle, but the NATS client exposes no hook to replace
* it — its TLS options are file paths and PEM strings, with no verify callback. So the
* certificate the mesh issues the bus **must carry a subject-alternative name matching the
* address nodes dial it by**. The fingerprint check below still happens and is still the real
* guarantee; what cannot be switched off is the check *beside* it.
* **The pin is the only check.** What comes back is handed to the client as its TLS options, and
* the client's transport spreads them into Node's own `tls.connect` — so the pinned certificate
* is the one authority the handshake accepts, and the hostname check beside it is replaced with
* one that always passes. Pinning the exact certificate makes verifying its name redundant, and
* the bus's certificate names the seat (`mesh-broker`), not the address a machine happens to
* dial it by: every module on the mesh met "does not match certificate's altnames" the first time
* it reached the handshake (2026-09-28).
*/
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> {
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<TlsOptions> {
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<tls.DetailedPeerCertificate>((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)) {
@@ -325,7 +347,9 @@ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: str
);
}
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
return { ca: pem };
// Node's option, not the client's: the transport passes the whole object on. `undefined` from
// checkServerIdentity is "the name is fine"; the pin above already decided the rest.
return { ca: pem, checkServerIdentity: () => undefined } as TlsOptions;
}
/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary —
+16 -5
View File
@@ -22,8 +22,7 @@
import { readFileSync } from "node:fs";
import { pathToFileURL } from "node:url";
import { connectAmqp, fatalBrokerReason } from "./broker-amqp.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";
@@ -35,6 +34,12 @@ import { emit } from "@novox/mesh-sdk/events";
* over the plain bootstrap URL otherwise. A scoped module assumes the foundation's exchanges exist —
* its account may not declare them (ADR 0043).
*/
// 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);
}
async function connectBroker(): Promise<Broker> {
const file = process.env.MESH_BROKER_FILE;
if (file) {
@@ -54,7 +59,13 @@ async function connectBroker(): Promise<Broker> {
// 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 });
// **The credential names the bus.** A module moved to the bus being built was handed a
// credential for it — `nats://…` with user, password and fingerprint beside the address — and
// 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.
// 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) {
@@ -63,7 +74,7 @@ async function connectBroker(): Promise<Broker> {
);
process.exit(1);
}
return connectAmqp(url);
return connectNats({ url });
}
/**
@@ -83,7 +94,7 @@ async function connectBrokerPatiently(): Promise<Broker> {
try {
return await connectBroker();
} catch (err) {
const fatal = fatalBrokerReason(err);
const fatal = fatalBrokerReasonFor(err);
if (fatal !== null) {
console.error(`mesh-tools: ${fatal} — waiting will not fix this; giving up`);
throw err;
-39
View File
@@ -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();
});
-142
View File
@@ -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 (<node>.<module>.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<void> {
return new Promise((r) => setTimeout(r, ms));
}
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
const start = Date.now();
while (!cond()) {
if (Date.now() - start > ms) throw new Error("condition not met in time");
await delay(25);
}
}
+3 -6
View File
@@ -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}`);
}
});
-68
View File
@@ -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");
});