// 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..event. an event this module emits // mesh.mod..tool. a tool this module serves // mesh.seat..accept. 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 net from "node:net"; import tls from "node:tls"; 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(); /** 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; /** The seats this module claims, with the verbs each promises (novox/hq ADR 0159). The runtime * serves each claimed seat's verbs with its tools of the same name; the bus admits only the * holder's subscription, so claiming and not holding costs a refused subscription and nothing else. */ claims?: { seat: string; scope?: string; serves?: string[] }[]; } /** * What the mesh issued this assignment (novox/hq ADR 0160): where its tools are served, in which * queue, the verbs of the seats it holds, where its events land, what it may reach. Read from the * ASSIGNMENTS stream at `mesh.assignment..` — the one subject a runtime derives for * itself — and followed live. Absent for a mesh older than the membership, and then the runtime * serves the shape it always derived, and says so. */ export interface Membership { node: string; module: string; /** Addresses a tool is answered on; `{tool}` stands for the tool's name. */ serves: { subject: string; queue?: string }[]; seats?: { seat: string; verb: string; subject: string }[]; emits: string; reaches?: Record; tools: string; } /** What a tool call answers: the module's own result, and which machine answered it * (novox/hq ADR 0159) — a module on several machines is otherwise an answer from nowhere. */ export interface Answered { result: Res; node?: string; } /** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs — * an answer that says which machine gave it, and serving a subject that is not a module's own tool * (a seat's verb). */ export interface RuntimeBroker extends Broker { /** Call a tool by key, or — when `on` names a subject the mesh listed for it (ADR 0160) — there. */ ask(key: string, body: Req, on?: string): Promise>; handleSubject(subject: string, handler: (body: Req) => Promise): Promise<() => void>; /** What the mesh issued this assignment, or undefined when nothing has been issued yet. */ membership(): Membership | undefined; /** Called when the mesh issues a new membership; the runtime re-serves on it. */ onMembership(handler: (m: Membership) => void): void; } const ASSIGNMENTS_STREAM = "ASSIGNMENTS"; /** The one address a runtime derives for itself (ADR 0160). */ export function membershipSubject(node: string, module: string): string { return `mesh.assignment.${node}.${module}`; } /** 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 { 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, // Its own inbox, not a random one: every user's inbox is private to it (design 25 §4), and the // grant names `_INBOX..>` — 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, }); const js = conn.jetstream(); const subs: Subscription[] = []; // Every registration, and the one reader that dispatches to them. A module has one durable // consumer; the loop belongs to the connection rather than to a subscription. const listeners: { pattern: string; handler: (env: Envelope) => Promise }[] = []; let reading: Awaited>["consume"]>> | undefined; let closed = false; const node = cred.node; // The membership, read once at connect and followed. A direct get is one request on the // stream's API, which is the whole of what this account may ask JetStream for its own subject; // a 404 is a mesh that has not issued one, which is a fact to say and not an error to retry. let issued: Membership | undefined; const issuedHandlers: ((m: Membership) => void)[] = []; const subjectOfMine = node ? membershipSubject(node, self) : ""; if (subjectOfMine) { try { const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}`, sc.encode(JSON.stringify({ last_by_subj: subjectOfMine })), { timeout: 5_000 }); const status = got.headers?.code ?? 0; if (status === 0 && got.data.length > 0) { issued = JSON.parse(sc.decode(got.data)) as Membership; } } catch { // Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below. } if (!issued) { console.log(`[mesh-tools] no membership issued for ${self} on ${node} yet; serving the derived shape until one arrives`); } try { const live = conn.subscribe(subjectOfMine); subs.push(live); void (async () => { for await (const msg of live) { try { issued = JSON.parse(sc.decode(msg.data)) as Membership; console.log(`[mesh-tools] ${self} on ${node} was issued a new membership; re-serving on it`); for (const h of issuedHandlers) h(issued); } catch (err) { console.log(`[mesh-tools] a membership arrived that is not one: ${err}`); } } })(); } catch { // A subscription this account may not make is a mesh older than the membership. } } /** The subjects a tool of this module is served on: from the membership when issued, derived * otherwise (the shape the mesh issues on day one, so the two agree). */ const servedOn = (tool: string): { subject: string; queue?: string }[] => { if (issued) { const m = issued; const out = m.serves.map((s) => ({ subject: s.subject.replace("{tool}", tool), queue: s.queue })); // The verb that lists what this module serves is answered on the mesh's plain address for it // whatever the placement — one answer suffices, so a queue — and on this machine's beside it. if (tool === "tools" && m.tools && !out.some((s) => s.subject === m.tools)) { out.unshift({ subject: m.tools, queue: `serve.${self}` }); } return out; } const base = `mesh.mod.${self}.tool.${tool}`; const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }]; if (node) out.push({ subject: `${base}.${node}` }); return out; }; /** Where a call by key goes: a subject the membership says this module reaches, when it says * one — the machine's when named — else the derived shape. */ const reachedAt = (key: string): string => { const [name, wanted] = key.split("@", 2); const reach = issued?.reaches?.[name]; if (reach && reach.length > 0) { if (wanted) { const at = reach.find((s) => s.endsWith(`.${wanted}`)); if (at) return at; } else { return reach[0]; } } return toolSubject(key, self); }; /** Answer one subject with one handler, and say which machine answered (novox/hq ADR 0159). */ const answerOn = ( subject: string, queue: string | undefined, handler: (body: Req) => Promise, ): (() => void) => { const sub = conn.subscribe(subject, queue ? { queue } : {}); subs.push(sub); void (async () => { for await (const msg of sub) { let reply: { result?: Res; error?: string; node?: 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) }; } if (node) reply.node = node; msg.respond(sc.encode(JSON.stringify(reply))); } })(); return () => sub.unsubscribe(); }; /** * Ask one question and await one answer, with the machine that gave it. * * 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. */ const ask = async (key: string, body: Req, on?: string): Promise> => { const msg = await conn.request(on ?? reachedAt(key), sc.encode(JSON.stringify(body)), { timeout: REQUEST_TIMEOUT_MS, }); const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string; node?: string }; if (reply.error) throw new Error(reply.error); return { result: reply.result as Res, node: reply.node }; }; return { async request(key: string, body: Req): Promise { return (await ask(key, body)).result; }, ask, /** * Answer a question, two ways (novox/hq ADR 0159): on the module's subject in a queue group, * so several machines may serve one tool and exactly one of them answers each call; and on the * same subject with this machine as its last token, so a caller that names the machine reaches * this instance and no other. A runtime that does not know its machine serves only the first, * which is how it always behaved. */ async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { // A seat's verb named outright is served on the seat's subject as given, for a holder that // knows its role without a membership; everything else is this module's own tool, served // where the mesh issued it (ADR 0160). A key naming another module is not served here at all. if (key.startsWith("seat:")) { const stop = answerOn(toolSubject(key, self), undefined, handler); return () => stop(); } const dot = key.indexOf("."); if (dot >= 0 && key.slice(0, dot) !== self) { throw new Error(`${self} cannot serve ${key}: a module serves its own tools`); } const tool = dot < 0 ? key : key.slice(dot + 1); let stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler)); // When a new membership arrives, serve where it now says and stop serving where it no longer does. issuedHandlers.push(() => { stops.forEach((stop) => stop()); stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler)); }); return () => stops.forEach((stop) => stop()); }, membership: () => issued, onMembership: (handler: (m: Membership) => void) => { issuedHandlers.push(handler); }, async handleSubject(subject: string, handler: (body: Req) => Promise): Promise<() => void> { return answerOn(subject, undefined, handler); }, /** * 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(env: Envelope): Promise { // **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; 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. * * **One consumer, one loop, however many patterns a module registers.** A module has exactly one * durable consumer, so two loops reading it would each take half the messages — and a loop that * received one its own pattern does not match acknowledges it, which is the right answer for a * filter wider than anything registered and silent loss when it is another handler's. Every * registration is therefore dispatched from one reader, and a message is acknowledged once every * handler it is for has taken it. */ async subscribe( pattern: string, handler: (env: Envelope) => Promise, ): Promise<() => void> { const listener = { pattern, handler: handler as (env: Envelope) => Promise }; listeners.push(listener); if (!reading) { const durable = `${cred.node ?? "?"}_${self}`; const consumer = await js.consumers.get("EVENTS", durable); const messages = await consumer.consume(); reading = messages; void (async () => { for await (const msg of messages) { await deliver(msg, listeners); } })(); } return () => { const at = listeners.indexOf(listener); if (at >= 0) listeners.splice(at, 1); if (listeners.length === 0 && reading) { void reading.close(); reading = undefined; } }; }, async close(): Promise { 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 to every handler it is for, acknowledging only once each has taken it. * * Several registrations share one durable consumer, so matching happens here rather than by having * each registration read the stream: two readers of one consumer would split it between them, and a * message that reached the wrong one would be acknowledged as not-for-me and lost. */ async function deliver( msg: JsMsg, listeners: { pattern: string; handler: (env: Envelope) => Promise }[], ): Promise { let env: Envelope; try { env = toEnvelope(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; } const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key)); if (forThis.length === 0) { // The consumer's filters are the controller's, derived from what the module declared it // consumes, and may be wider than anything it registered a handler for. Acknowledge it, or it // would be redelivered until it expired. msg.ack(); return; } try { for (const l of forThis) await l.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(msg: JsMsg): Envelope { const headers: Record = {}; 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..event.`. 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..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; `.` addresses * another's, which is how a request reaches a module that is not this one; `seat:.` * addresses a role's tool, answered by whoever holds the seat (novox/hq ADR 0132) — with * `seat:.@` for a node-scoped seat, whose tool carries the machine (design 33 §4). */ function toolSubject(key: string, self: string): string { if (key.startsWith("seat:")) { const rest = key.slice("seat:".length); const dot = rest.indexOf("."); if (dot < 0) throw new Error(`"${key}" names a seat and no verb: seat:.`); const seat = rest.slice(0, dot); const [verb, node] = rest.slice(dot + 1).split("@", 2); return node ? `mesh.seat.${seat}.tool.${verb}.${node}` : `mesh.seat.${seat}.tool.${verb}`; } // `.@` names the machine (novox/hq ADR 0159): the same subject with the // machine as its last token, which is what that instance serves beside the queue. const [name, node] = key.split("@", 2); const dot = name.indexOf("."); const base = dot < 0 ? `mesh.mod.${self}.tool.${name}` : `mesh.mod.${name.slice(0, dot)}.tool.${name.slice(dot + 1)}`; return node ? `${base}.${node}` : base; } /** A seat's verb, as its holder serves it: flat for a mesh seat, carrying the machine for a * node-scoped one (design 33 §4) — the same shape the controller grants. */ export function seatToolSubject(seat: string, verb: string, scope: string | undefined, node: string | undefined): string { const base = `mesh.seat.${seat}.tool.${verb}`; return scope === "node" && node ? `${base}.${node}` : base; } 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. * * **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 { const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`); const port = url.port ? Number(url.port) : 4222; // **The bus speaks first, in the clear.** A NATS server sends its INFO line before TLS begins, // and only then expects the client to start the handshake; a raw TLS connect to that port reads // the INFO line as a TLS record and fails with "wrong version number" — which is what every // module met the first time it dialled the bus being built (2026-09-28). So: connect, wait for // INFO, then start TLS on the same socket, and read the certificate the server presents. const certificate = await new Promise((resolve, reject) => { const 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)) { 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`; // 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 — * 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); }