From c1517c39a0878b5ceaa80971674ee38e4bc58528 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 23:29:17 +0200 Subject: [PATCH] The tool runtime's client on NATS, behind the unchanged contract MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Task 3.6 of novox/hq ADR 0116. A module is still written against request, handle, publish, subscribe, close; only what is underneath changes. main.ts still selects the AMQP client — steps 1 to 4 leave every node on AMQP, so this ships beside it and is selected at the rollout. Round-tripped against a real server (test/roundtrip.mjs): a tool answered across two connections, a throwing handler reaching the caller as an error rather than a timeout, an event delivered once with its key, body, node and event id intact, and an event landing under its emitter's own namespace. Three things the compiler and the server corrected: - the envelope's field is `key`, not `type`, and the payload is `env.body` with metadata in headers — not the whole envelope re-encoded. An implementation that nested the envelope would pass all its own tests and agree with nobody, which is what the conformance suite exists to stop. - the NATS client's TLS options are PEM strings with no verify hook, so the AMQP client's `checkServerIdentity: () => undefined` has no equivalent. The fingerprint check still happens and is still the guarantee, but the bus's certificate must now carry a SAN matching the address nodes dial. That is a constraint on the mesh's certificates, recorded where it bites. - a durable consumer is bound, never created: a module's account cannot reach the JetStream API, and a runtime creating its own would be a module choosing its own delivery semantics. --- package.json | 3 +- src/broker-nats.ts | 341 +++++++++++++++++++++++++++++++++++++++++++++ test/roundtrip.mjs | 57 ++++++++ 3 files changed, 400 insertions(+), 1 deletion(-) create mode 100644 src/broker-nats.ts create mode 100644 test/roundtrip.mjs diff --git a/package.json b/package.json index 3ed7292..7c04aa6 100644 --- a/package.json +++ b/package.json @@ -12,7 +12,8 @@ }, "dependencies": { "@novox/mesh-sdk": "^0.1.0", - "amqplib": "^0.10.9" + "amqplib": "^0.10.9", + "nats": "^2.29.0" }, "devDependencies": { "@types/amqplib": "^0.10.8", diff --git a/src/broker-nats.ts b/src/broker-nats.ts new file mode 100644 index 0000000..caae49f --- /dev/null +++ b/src/broker-nats.ts @@ -0,0 +1,341 @@ +// 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. +// +// What changes is underneath: exchanges and per-tool queues become subjects, and durability +// becomes 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 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, matching the AMQP client's behaviour so + * a module's timeout handling does not change with the transport. */ +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. Mirrors the AMQP client's judgement so the runtime's supervisor does + * not have to know which transport it is on. */ +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, + // Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on + // the mesh to exit. The AMQP client's supervisor did this a level up; here the library does + // it, and `close()` is still 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(key: string, body: Req): Promise { + 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 — the same property the AMQP client got from a shared durable queue. + */ + async handle(key: string, handler: (body: Req) => Promise): 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), which is what the AMQP client's confirm channel was for. + * + * `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**, exactly as on AMQP + // (ADR 0042). NATS has headers of its own, so the envelope's shape on the wire is + // preserved rather than re-encoded — which matters because that shape is what the + // conformance suite pins, and 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. + */ + async subscribe( + pattern: string, + handler: (env: Envelope) => Promise, + ): 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 { + 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( + msg: JsMsg, + 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; + } + 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(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 inside a module's event subject. */ +function keyFromSubject(subject: string): string { + const marker = ".event."; + const at = subject.indexOf(marker); + return at < 0 ? subject : subject.slice(at + marker.length); +} + +/** 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. */ +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. + * + * **One behaviour differs from the AMQP client, and it is a constraint on the mesh rather than a + * detail of this file.** That client passed `checkServerIdentity: () => undefined`, because + * pinning the exact certificate makes hostname verification redundant. The NATS client exposes no + * such hook — 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((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, unchanged from AMQP: `*` is one token, `#` the rest. Kept because + * it is the module's vocabulary — a module's `consumes` pattern reads the same as it always did, + * and the subject it becomes is the mesh's business. */ +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; + if (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); +} diff --git a/test/roundtrip.mjs b/test/roundtrip.mjs new file mode 100644 index 0000000..702fef8 --- /dev/null +++ b/test/roundtrip.mjs @@ -0,0 +1,57 @@ +// Round-trip the runtime's NATS client against a real server: a tool call answered, and an +// event emitted and received with its envelope intact. +import { connect } from "nats"; +import { connectNats } from "../dist/broker-nats.js"; + +const URL = "nats://127.0.0.1:14222"; + +// The controller's job, done by hand here: the stream and the module's durable consumer. +const admin = await connect({ servers: URL }); +const jsm = await admin.jetstreamManager(); +await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>", "mesh.seat.*.event.>"] }); +await jsm.consumers.add("EVENTS", { + durable_name: "one_audit", ack_policy: "explicit", + filter_subjects: ["mesh.mod.shop.event.order.placed"], +}); + +const shop = await connectNats({ url: URL, node: "one", module: "shop" }); +const audit = await connectNats({ url: URL, node: "one", module: "audit" }); + +let failures = 0; +const check = (ok, what) => { console.log(` ${ok ? "ok " : "FAIL"} ${what}`); if (!ok) failures++; }; + +// A tool, served and called. +await shop.handle("price", async (body) => ({ total: body.qty * 3 })); +const answer = await audit.request("shop.price", { qty: 4 }); +check(answer.total === 12, "a tool call is answered across two connections"); + +// A handler that throws reaches the caller as an error, not a timeout. +await shop.handle("boom", async () => { throw new Error("no"); }); +let threw = null; +try { await audit.request("shop.boom", {}); } catch (e) { threw = e.message; } +check(threw === "no", "a handler that throws answers the caller instead of timing out"); + +// An event, emitted and received with its envelope intact. +const seen = []; +await audit.subscribe("order.placed", async (env) => { seen.push(env); }); +await shop.publish({ + key: "order.placed", node: "one", body: { id: "a1" }, + headers: { "x-event-id": "e1", "x-node": "one", "content-type": "application/json" }, +}); +await new Promise((r) => setTimeout(r, 800)); +check(seen.length === 1, `exactly one delivery (saw ${seen.length})`); +if (seen[0]) { + check(seen[0].key === "order.placed", "the key survives the subject round trip"); + check(seen[0].body?.id === "a1", "the body is the payload, not the whole envelope"); + check(seen[0].node === "one", "the node comes back from the headers"); + check(seen[0].headers?.["x-event-id"] === "e1", "the event id survives as a header"); +} + +// A module cannot reach into another's namespace by naming its own event oddly. +await shop.publish({ key: "other", node: "one", body: {}, headers: { "x-event-id": "e2" } }); +const msg = await jsm.streams.getMessage("EVENTS", { last_by_subj: "mesh.mod.shop.event.other" }); +check(!!msg, "an event lands under the emitting module's own namespace"); + +await shop.close(); await audit.close(); await admin.close(); +console.log(failures ? `\n${failures} failed` : "\nall passed"); +process.exit(failures ? 1 : 0);