A person's own client, and pins that the wire did not change #13

Merged
jschoubben merged 6 commits from feat/nats-genesis into main 2026-09-27 17:20:07 +00:00
Showing only changes of commit 2197c36fef - Show all commits
+21 -25
View File
@@ -5,8 +5,8 @@
// 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).
// 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
@@ -23,8 +23,8 @@ 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. */
/** 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 {}
@@ -42,8 +42,8 @@ export interface Credential {
}
/** 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. */
* 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 };
@@ -85,8 +85,7 @@ export async function connectNats(
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.
// the mesh to exit. `close()` stays the only thing that ends the connection.
maxReconnectAttempts: -1,
});
const js = conn.jetstream();
@@ -116,7 +115,7 @@ export async function connectNats(
* 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.
* call.
*/
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` });
@@ -144,17 +143,15 @@ export async function connectNats(
*
* 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.
* (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**, 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.
// **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)) {
@@ -288,13 +285,12 @@ function normalizeFingerprint(fingerprint: string): string {
* 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.
* **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}`);
@@ -320,9 +316,9 @@ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: str
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. */
/** 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);
}