Files
mesh-tools/src/broker-nats.ts
T
jschoubben fbeb373d1a Both clients map local event names to their own wire
A module names its events locally and each transport works out where they land.
That is what design 29 says and what neither client did: both passed the name
straight through, which happened to be right on the old bus because modules were
writing routing keys, and wrong on the new one (novox/hq 04-ISSUES/127).

The old bus's client now turns a local name into `module.<emitter>.<event>` on the
way out and back on the way in. Without that, converting the modules to local
names would have broken the mesh that is actually running.

**A handler and a manifest now say the same thing.** The key a module sees was the
event name alone, so a manifest declaring `consumes: builder.built` produced a
pattern that could never match what it was compared against — and a module
consuming one event from two emitters could only tell them apart by reading a
header. The subject already carries the emitter, so naming it in the key makes a
mismatch between manifest and code a typo instead of a category error.

Both matchers accept `**` for the rest of a name, which is how a manifest spells
it; the old bus's `#` still works, because both buses ship until the rollout.
2026-09-27 14:42:40 +02:00

352 lines
15 KiB
TypeScript

// The tool runtime's broker client, on NATS.
//
// **The sdk's contract does not change** (novox/hq ADR 0106, ADR 0039): a module is written
// against `request`, `handle`, `publish`, `subscribe`, `close`, and the runtime implements them.
// That is why a module built before any of this runs on the new runtime without a rebuild, and
// why the sdk's own diff for the whole bus change is three comments.
//
// Underneath, everything is a subject and durability is JetStream (novox/hq design 25,
// design 29).
//
// mesh.mod.<module>.event.<type> an event this module emits
// mesh.mod.<module>.tool.<tool> a tool this module serves
// mesh.seat.<seat>.accept.<verb> work submitted to a role
//
// The module never writes one of those: it names its events and tools locally and the mesh
// derives the subject (design 29 §1), so reorganising the subject space leaves every module
// correct.
import { createHash } from "node:crypto";
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";
const sc = StringCodec();
/** Requests wait this long for an answer before failing. Unchanged from what modules already
* expect, so a module's timeout handling is not something the bus quietly redefines. */
const REQUEST_TIMEOUT_MS = 30_000;
export class PinMismatchError extends Error {}
/** A broker credential as the mesh delivers it (novox/hq ADR 0120): the bus's address, the
* fingerprint of the certificate it must present, and the node and module the account is scoped
* to — the runtime derives its subjects from those rather than being told them. */
export interface Credential {
url: string;
fingerprint?: string;
node?: string;
module?: string;
user?: string;
password?: string;
}
/** Whether a connection failure is worth retrying, or is a fact about this configuration that
* retrying cannot change. The runtime's supervisor asks this and does not need to know what it
* is connected to. */
export function fatalBrokerReason(err: unknown): string | null {
if (err instanceof PinMismatchError) return "the bus's certificate does not match the pin";
const e = err as { code?: string; message?: string };
const message = typeof e?.message === "string" ? e.message : String(err);
if (e?.code === "ERR_INVALID_URL" || /invalid url/i.test(message)) {
return "the bus address is not a usable URL";
}
if (/authorization violation|user authentication expired|permissions violation/i.test(message)) {
return "the bus refused this account";
}
return null;
}
/**
* Connect to the mesh bus and return a Broker.
*
* **A module's subjects come from its credential, not from its calls.** `node` and `module` name
* the account the mesh issued, and every subject this client publishes or subscribes is derived
* from them — so a module cannot name another's namespace even by mistake, and what it emits
* matches what the mesh authorised (ADR 0074's identity rule).
*/
export async function connectNats(
target: string | Credential,
opts: { module?: string } = {},
): Promise<Broker> {
const cred: Credential = typeof target === "string" ? { url: target } : target;
const self = cred.module ?? opts.module;
if (!self) {
throw new Error(
"a broker credential with no module: the runtime derives its subjects from the account " +
"the mesh issued, and cannot guess which module it is",
);
}
const conn = await natsConnect({
servers: cred.url,
user: cred.user,
pass: cred.password,
name: `${cred.node ?? "?"}.${self}`,
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : 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,
});
const js = conn.jetstream();
const subs: Subscription[] = [];
let closed = false;
return {
/**
* Ask one question and await one answer.
*
* Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3),
* and a lost one is a timeout the caller already handles. The reply travels on the inbox the
* request carries, which the responder may answer because its account has `allow_responses`
* — one reply to a message it actually received, and nothing wider.
*/
async request<Req, Res>(key: string, body: Req): Promise<Res> {
const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), {
timeout: REQUEST_TIMEOUT_MS,
});
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string };
if (reply.error) throw new Error(reply.error);
return reply.result as Res;
},
/**
* Answer a question.
*
* A queue group, so several nodes may serve one tool and exactly one of them answers each
* call.
*/
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` });
subs.push(sub);
void (async () => {
for await (const msg of sub) {
let reply: { result?: Res; error?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
msg.respond(sc.encode(JSON.stringify(reply)));
}
})();
return () => {
sub.unsubscribe();
};
},
/**
* Emit an event.
*
* Published into JetStream and awaited, so a publish the bus never accepted fails the emit
* rather than vanishing — at-least-once starts at the emitter, not only the consumer
* (ADR 0042).
*
* `msgID` is the event's own id, so a redelivery after a crash between publishing and
* acknowledging is de-duplicated by the server inside its window rather than seen twice.
*/
async publish<T>(env: Envelope<T>): Promise<void> {
// **The body is the payload and the metadata rides as headers** (ADR 0042). That shape
// is what the conformance suite pins: an implementation that nested the whole envelope in
// the body would pass every one of its own tests and agree with nobody.
const meta = (env.headers ?? {}) as Record<string, string>;
const h = natsHeaders();
for (const [k, v] of Object.entries(meta)) {
if (v != null) h.set(k, String(v));
}
if (!meta["content-type"]) h.set("content-type", "application/json");
if (env.node) h.set("x-node", env.node);
await js.publish(eventSubject(env.key, self), sc.encode(JSON.stringify(env.body)), {
headers: h,
// De-duplicated by the server inside its window, so a redelivery after a crash between
// publishing and acknowledging is not seen twice. Only the emitter can make this id.
msgID: meta["x-event-id"],
});
},
/**
* React to events.
*
* The durable consumer is the **controller's** to create, from what this module declared it
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
* its own would be a module deciding its own delivery semantics, and its account cannot
* reach the JetStream API to do it anyway.
*/
async subscribe<T>(
pattern: string,
handler: (env: Envelope<T>) => Promise<void>,
): Promise<() => void> {
const durable = `${cred.node ?? "?"}_${self}`;
const consumer = await js.consumers.get("EVENTS", durable);
const messages = await consumer.consume();
void (async () => {
for await (const msg of messages) {
await deliver(msg, pattern, handler);
}
})();
return () => {
void messages.close();
};
},
async close(): Promise<void> {
if (closed) return;
closed = true;
for (const sub of subs) sub.unsubscribe();
// Drain rather than close: an in-flight reply is finished instead of dropped, which for a
// tool call is the difference between an answer and an unexplained timeout at the caller.
await conn.drain();
},
};
}
/** Deliver one event, acknowledging only once a handler has taken it. */
async function deliver<T>(
msg: JsMsg,
pattern: string,
handler: (env: Envelope<T>) => Promise<void>,
): Promise<void> {
let env: Envelope<T>;
try {
env = toEnvelope<T>(msg);
} catch {
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
// that were never an envelope.
msg.term();
return;
}
if (!topicMatches(pattern, env.key)) {
// The consumer's filters are the controller's, and may be wider than one subscription's
// pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
// would be redelivered until it expired.
msg.ack();
return;
}
try {
await handler(env);
msg.ack();
} catch {
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
// attempt, and one failing permanently exhausts max-deliver and dead-letters rather than
// spinning. The consumer's limits are the controller's; this only says "not done".
msg.nak(5_000);
}
}
/** Rebuild the envelope a module sees, from the subject, the headers and the payload — the
* mirror of publish, and the reason both live beside each other. */
function toEnvelope<T>(msg: JsMsg): Envelope<T> {
const headers: Record<string, string> = {};
if (msg.headers) {
for (const k of msg.headers.keys()) headers[k] = msg.headers.get(k);
}
return {
// The event's own key, recovered from the subject: `mesh.mod.<module>.event.<key>`. The
// module never sees the subject, only the key it declared.
key: keyFromSubject(msg.subject),
node: headers["x-node"] ?? "",
body: JSON.parse(sc.decode(msg.data)) as T,
headers: headers as EventHeaders,
};
}
/** The event key a module sees: the emitter and the event, which is exactly how its manifest names
* what it consumes (novox/hq design 29 §1, 04-ISSUES/127).
*
* **One vocabulary for the declaration and the handler.** This returned the event name alone, so a
* manifest declaring `consumes: builder.built` produced a handler pattern that could never match
* the key it was compared against — and a module consuming the same event from two emitters could
* not tell them apart except by reading a header. The subject already carries the emitter; naming it
* here makes a mismatch between manifest and code a typo rather than a category error. */
function keyFromSubject(subject: string): string {
const marker = ".event.";
const at = subject.indexOf(marker);
if (at < 0) return subject;
const event = subject.slice(at + marker.length);
// `mesh.mod.<emitter>.event.…` — the emitter is the token before the marker.
const before = subject.slice(0, at).split(".");
const emitter = before[before.length - 1];
return emitter ? `${emitter}.${event}` : event;
}
/** A module's own event subject. Derived, never taken from the caller: the module names its
* event and the mesh decides where it lands (design 29 §1). */
function eventSubject(type: string, self: string): string {
return `mesh.mod.${self}.event.${type}`;
}
/** A tool's subject. A bare name is this module's own tool; `<module>.<tool>` addresses
* another's, which is how a request reaches a module that is not this one. */
function toolSubject(key: string, self: string): string {
const dot = key.indexOf(".");
if (dot < 0) return `mesh.mod.${self}.tool.${key}`;
return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`;
}
function normalizeFingerprint(fingerprint: string): string {
return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase();
}
/**
* Dial once to see the certificate, and refuse unless it is exactly the one the mesh pinned.
* 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.
*/
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;
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 seen = createHash("sha256").update(certificate.raw).digest("hex");
if (seen !== normalizeFingerprint(fingerprint)) {
throw new PinMismatchError(
`the bus at ${url.hostname}:${port} presented ${seen}, not the pinned ${normalizeFingerprint(fingerprint)}`,
);
}
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
return { ca: pem };
}
/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary —
* a module's `consumes` pattern is matched here, and the subject it becomes is the mesh's
* business, not the module's. */
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 {
if (pi === p.length) return ki === k.length;
// `**` is the mesh's wildcard for the rest of a name; `#` is the old bus's, accepted so a pattern
// written either way behaves the same while both buses ship (novox/hq design 29 §1).
if (p[pi] === "#" || p[pi] === "**") {
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 (p[pi] !== "*" && p[pi] !== k[ki]) return false;
return matchFrom(p, pi + 1, k, ki + 1);
}