The tool runtime's client on NATS, behind the unchanged contract

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.
This commit is contained in:
2026-09-26 23:29:17 +02:00
parent 6b380b67e0
commit c1517c39a0
3 changed files with 400 additions and 1 deletions
+341
View File
@@ -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.<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, 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<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. 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<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 — the same property the AMQP client got from a shared durable queue.
*/
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), 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<T>(env: Envelope<T>): Promise<void> {
// **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<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 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; `<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.
*
* **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<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, 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);
}