One bus: the runtime pins the certificate after the server speaks, and the old transport goes #15
Generated
+76
@@ -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"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -14,11 +14,9 @@
|
|||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@novox/mesh-sdk": "^0.1.0",
|
"@novox/mesh-sdk": "^0.1.0",
|
||||||
"amqplib": "^0.10.9",
|
|
||||||
"nats": "^2.29.0"
|
"nats": "^2.29.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/amqplib": "^0.10.8",
|
|
||||||
"@types/node": "^22.20.1",
|
"@types/node": "^22.20.1",
|
||||||
"typescript": "^5.9.3"
|
"typescript": "^5.9.3"
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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;
|
|
||||||
}
|
|
||||||
+26
-9
@@ -17,6 +17,7 @@
|
|||||||
// correct.
|
// correct.
|
||||||
|
|
||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
|
import net from "node:net";
|
||||||
import tls from "node:tls";
|
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 } from "nats";
|
||||||
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
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 }> {
|
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> {
|
||||||
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
||||||
const port = url.port ? Number(url.port) : 4222;
|
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 certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
|
||||||
const socket = tls.connect(
|
const plain = net.connect({ host: url.hostname, port }, () => {});
|
||||||
{ host: url.hostname, port, rejectUnauthorized: false, servername: url.hostname },
|
let seenInfo = false;
|
||||||
() => {
|
let buffered = "";
|
||||||
const peer = socket.getPeerCertificate(true);
|
plain.on("error", reject);
|
||||||
socket.end();
|
plain.on("data", (chunk: Buffer) => {
|
||||||
resolve(peer);
|
if (seenInfo) return;
|
||||||
},
|
buffered += chunk.toString("utf8");
|
||||||
);
|
if (!buffered.includes("\r\n")) return;
|
||||||
socket.on("error", reject);
|
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");
|
const seen = createHash("sha256").update(certificate.raw).digest("hex");
|
||||||
if (seen !== normalizeFingerprint(fingerprint)) {
|
if (seen !== normalizeFingerprint(fingerprint)) {
|
||||||
|
|||||||
+5
-9
@@ -22,9 +22,7 @@
|
|||||||
|
|
||||||
import { readFileSync } from "node:fs";
|
import { readFileSync } from "node:fs";
|
||||||
import { pathToFileURL } from "node:url";
|
import { pathToFileURL } from "node:url";
|
||||||
import { connectAmqp, fatalBrokerReason as fatalAmqpReason } from "./broker-amqp.js";
|
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js";
|
||||||
import { connectNats, fatalBrokerReason as fatalNatsReason } from "./broker-nats.js";
|
|
||||||
import type { Credential } from "./broker-amqp.js";
|
|
||||||
import { runTools } from "./runtime.js";
|
import { runTools } from "./runtime.js";
|
||||||
import { invokeTool } from "@novox/mesh-sdk/tools";
|
import { invokeTool } from "@novox/mesh-sdk/tools";
|
||||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
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
|
// 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.
|
// whichever bus this runtime is on — each transport knows its own refusals.
|
||||||
function fatalBrokerReasonFor(err: unknown): string | null {
|
function fatalBrokerReasonFor(err: unknown): string | null {
|
||||||
return fatalNatsReason(err) ?? fatalAmqpReason(err);
|
return fatalNatsReason(err);
|
||||||
}
|
}
|
||||||
|
|
||||||
async function connectBroker(): Promise<Broker> {
|
async function connectBroker(): Promise<Broker> {
|
||||||
@@ -66,10 +64,8 @@ async function connectBroker(): Promise<Broker> {
|
|||||||
// nothing else in its environment changed (design 25; novox/hq design 28 task 5.2). The scheme
|
// 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
|
// is enough to know which bus to speak; a runtime that always dialled the old one would keep
|
||||||
// serving and answer nobody.
|
// serving and answer nobody.
|
||||||
if (credential.url.startsWith("nats://")) {
|
// One bus (novox/hq ADR 0131, design 28 task 5.5): the credential names it, and it is this.
|
||||||
return connectNats(credential);
|
return connectNats(credential);
|
||||||
}
|
|
||||||
return connectAmqp(credential, { assumeExchanges: true });
|
|
||||||
}
|
}
|
||||||
const url = process.env.MESH_BROKER_URL;
|
const url = process.env.MESH_BROKER_URL;
|
||||||
if (!url) {
|
if (!url) {
|
||||||
@@ -78,7 +74,7 @@ async function connectBroker(): Promise<Broker> {
|
|||||||
);
|
);
|
||||||
process.exit(1);
|
process.exit(1);
|
||||||
}
|
}
|
||||||
return connectAmqp(url);
|
return connectNats({ url });
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -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();
|
|
||||||
});
|
|
||||||
@@ -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);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
import { test } from "node:test";
|
import { test } from "node:test";
|
||||||
import assert from "node:assert/strict";
|
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
|
// 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
|
// 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", () => {
|
test("a refused login is fatal — a wrong or revoked credential, not an absent broker", () => {
|
||||||
for (const msg of [
|
// The bus refuses a login in its own words; each is final, because the next try says the same.
|
||||||
"Handshake terminated by server: 403 (ACCESS-REFUSED) with message \"ACCESS_REFUSED - Login was refused\"",
|
for (const msg of ["Authorization Violation", "nats: user authentication expired", "Permissions Violation for Subscription to \"x\""]) {
|
||||||
"Login was refused using authentication mechanism PLAIN",
|
|
||||||
"ACCESS_REFUSED",
|
|
||||||
]) {
|
|
||||||
assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`);
|
assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -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");
|
|
||||||
});
|
|
||||||
Reference in New Issue
Block a user