Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f04519e3e1 | ||
|
|
10e8191717 | ||
|
|
d703cebff4 | ||
|
|
46b56d53a6 | ||
|
|
d4a2802342 | ||
|
|
8acfa7a07d |
+55
-21
@@ -19,7 +19,7 @@
|
|||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
import net from "node:net";
|
import net from "node:net";
|
||||||
import tls from "node:tls";
|
import tls from "node:tls";
|
||||||
import { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription } from "nats";
|
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";
|
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
const sc = StringCodec();
|
const sc = StringCodec();
|
||||||
@@ -85,6 +85,10 @@ export async function connectNats(
|
|||||||
pass: cred.password,
|
pass: cred.password,
|
||||||
name: `${cred.node ?? "?"}.${self}`,
|
name: `${cred.node ?? "?"}.${self}`,
|
||||||
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
|
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.<user>.>` — 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
|
// 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.
|
// the mesh to exit. `close()` stays the only thing that ends the connection.
|
||||||
maxReconnectAttempts: -1,
|
maxReconnectAttempts: -1,
|
||||||
@@ -92,6 +96,10 @@ export async function connectNats(
|
|||||||
const js = conn.jetstream();
|
const js = conn.jetstream();
|
||||||
|
|
||||||
const subs: Subscription[] = [];
|
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<unknown>) => Promise<void> }[] = [];
|
||||||
|
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
|
||||||
let closed = false;
|
let closed = false;
|
||||||
|
|
||||||
return {
|
return {
|
||||||
@@ -176,21 +184,38 @@ export async function connectNats(
|
|||||||
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
|
* 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
|
* its own would be a module deciding its own delivery semantics, and its account cannot
|
||||||
* reach the JetStream API to do it anyway.
|
* 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<T>(
|
async subscribe<T>(
|
||||||
pattern: string,
|
pattern: string,
|
||||||
handler: (env: Envelope<T>) => Promise<void>,
|
handler: (env: Envelope<T>) => Promise<void>,
|
||||||
): Promise<() => void> {
|
): Promise<() => void> {
|
||||||
|
const listener = { pattern, handler: handler as (env: Envelope<unknown>) => Promise<void> };
|
||||||
|
listeners.push(listener);
|
||||||
|
if (!reading) {
|
||||||
const durable = `${cred.node ?? "?"}_${self}`;
|
const durable = `${cred.node ?? "?"}_${self}`;
|
||||||
const consumer = await js.consumers.get("EVENTS", durable);
|
const consumer = await js.consumers.get("EVENTS", durable);
|
||||||
const messages = await consumer.consume();
|
const messages = await consumer.consume();
|
||||||
|
reading = messages;
|
||||||
void (async () => {
|
void (async () => {
|
||||||
for await (const msg of messages) {
|
for await (const msg of messages) {
|
||||||
await deliver(msg, pattern, handler);
|
await deliver(msg, listeners);
|
||||||
}
|
}
|
||||||
})();
|
})();
|
||||||
|
}
|
||||||
return () => {
|
return () => {
|
||||||
void messages.close();
|
const at = listeners.indexOf(listener);
|
||||||
|
if (at >= 0) listeners.splice(at, 1);
|
||||||
|
if (listeners.length === 0 && reading) {
|
||||||
|
void reading.close();
|
||||||
|
reading = undefined;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
},
|
},
|
||||||
|
|
||||||
@@ -205,15 +230,20 @@ export async function connectNats(
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Deliver one event, acknowledging only once a handler has taken it. */
|
/**
|
||||||
async function deliver<T>(
|
* 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,
|
msg: JsMsg,
|
||||||
pattern: string,
|
listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[],
|
||||||
handler: (env: Envelope<T>) => Promise<void>,
|
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
let env: Envelope<T>;
|
let env: Envelope<unknown>;
|
||||||
try {
|
try {
|
||||||
env = toEnvelope<T>(msg);
|
env = toEnvelope<unknown>(msg);
|
||||||
} catch {
|
} catch {
|
||||||
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
|
// 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
|
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
|
||||||
@@ -221,15 +251,16 @@ async function deliver<T>(
|
|||||||
msg.term();
|
msg.term();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (!topicMatches(pattern, env.key)) {
|
const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key));
|
||||||
// The consumer's filters are the controller's, and may be wider than one subscription's
|
if (forThis.length === 0) {
|
||||||
// pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
|
// 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.
|
// would be redelivered until it expired.
|
||||||
msg.ack();
|
msg.ack();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
await handler(env);
|
for (const l of forThis) await l.handler(env);
|
||||||
msg.ack();
|
msg.ack();
|
||||||
} catch {
|
} catch {
|
||||||
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
||||||
@@ -298,14 +329,15 @@ function normalizeFingerprint(fingerprint: string): string {
|
|||||||
* A certificate authority is not consulted: the mesh issued this and knows its fingerprint,
|
* 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.
|
* 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
|
* **The pin is the only check.** What comes back is handed to the client as its TLS options, and
|
||||||
* hostname verification redundant in principle, but the NATS client exposes no hook to replace
|
* the client's transport spreads them into Node's own `tls.connect` — so the pinned certificate
|
||||||
* it — its TLS options are file paths and PEM strings, with no verify callback. So the
|
* is the one authority the handshake accepts, and the hostname check beside it is replaced with
|
||||||
* certificate the mesh issues the bus **must carry a subject-alternative name matching the
|
* one that always passes. Pinning the exact certificate makes verifying its name redundant, and
|
||||||
* address nodes dial it by**. The fingerprint check below still happens and is still the real
|
* the bus's certificate names the seat (`mesh-broker`), not the address a machine happens to
|
||||||
* guarantee; what cannot be switched off is the check *beside* it.
|
* 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<{ ca: string }> {
|
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<TlsOptions> {
|
||||||
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
||||||
const port = url.port ? Number(url.port) : 4222;
|
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,
|
// **The bus speaks first, in the clear.** A NATS server sends its INFO line before TLS begins,
|
||||||
@@ -342,7 +374,9 @@ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: str
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
|
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
|
||||||
return { ca: pem };
|
// 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 —
|
/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary —
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import { test } from "node:test";
|
import { test } from "node:test";
|
||||||
import assert from "node:assert/strict";
|
import assert from "node:assert/strict";
|
||||||
import { fatalBrokerReason, PinMismatchError } from "../src/broker-nats.ts";
|
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
|
||||||
|
|
||||||
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
||||||
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
||||||
@@ -41,3 +41,21 @@ test("a non-Error value does not crash the classifier", () => {
|
|||||||
assert.equal(fatalBrokerReason("just a string"), null);
|
assert.equal(fatalBrokerReason("just a string"), null);
|
||||||
assert.equal(fatalBrokerReason(undefined), null);
|
assert.equal(fatalBrokerReason(undefined), null);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// **One durable consumer feeds one reader, however many patterns a module registers.**
|
||||||
|
//
|
||||||
|
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
|
||||||
|
// reader that received one its own pattern does not match acknowledges it, which is right for a
|
||||||
|
// filter wider than anything registered and silent loss when it is another handler's. The matching is
|
||||||
|
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not
|
||||||
|
// by which reader happened to fetch it.
|
||||||
|
test("a message is for every pattern that matches it, and nothing else", () => {
|
||||||
|
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"];
|
||||||
|
const matched = (key: string) => registered.filter((p) => topicMatches(p, key));
|
||||||
|
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]);
|
||||||
|
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]);
|
||||||
|
// Nothing registered for it: the consumer's filter is the controller's and may be wider.
|
||||||
|
assert.deepEqual(matched("mesh-catalog.upgraded"), []);
|
||||||
|
// And a handler that asked for everything gets both, which is what the audit logger does.
|
||||||
|
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]);
|
||||||
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user