The one address a runtime derives for itself is mesh.assignment.<node>.<module>. It reads the
membership there with a direct get on the ASSIGNMENTS stream, serves each tool exactly where the
membership says — the plain subject in the module's queue when the mesh issued one, this machine's
beside it — and follows the subject, re-serving when a new membership arrives. A mesh that has issued
nothing yet gets the shape it always derived, and the log says so.
A seat's verbs are the role's, not the software's (ADR 0159): a module implements them with
registerModuleTools("<seat>", …), the runtime serves that on the seat's subjects when the credential
claims the seat, and never lists it among the module's own tools. A module named like its seat
registers once and is both.
The tools answer carries each tool's subjects, and the console and CLI call the subject the listing
gave them instead of composing one.
592 lines
27 KiB
TypeScript
592 lines
27 KiB
TypeScript
// 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.<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 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.<node>.<module>` — 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<string, string[]>;
|
|
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<Res> {
|
|
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<Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>>;
|
|
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): 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<RuntimeBroker> {
|
|
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.<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
|
|
// 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<unknown>) => Promise<void> }[] = [];
|
|
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["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 = <Req, Res>(
|
|
subject: string,
|
|
queue: string | undefined,
|
|
handler: (body: Req) => Promise<Res>,
|
|
): (() => 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 <Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>> => {
|
|
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<Req, Res>(key: string, body: Req): Promise<Res> {
|
|
return (await ask<Req, Res>(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<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): 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<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): 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<T>(env: Envelope<T>): Promise<void> {
|
|
// **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)) {
|
|
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<T>(
|
|
pattern: string,
|
|
handler: (env: Envelope<T>) => 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 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<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 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<unknown>) => Promise<void> }[],
|
|
): Promise<void> {
|
|
let env: Envelope<unknown>;
|
|
try {
|
|
env = toEnvelope<unknown>(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<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 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.<emitter>.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; `<module>.<tool>` addresses
|
|
* another's, which is how a request reaches a module that is not this one; `seat:<seat>.<verb>`
|
|
* addresses a role's tool, answered by whoever holds the seat (novox/hq ADR 0132) — with
|
|
* `seat:<seat>.<verb>@<node>` 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:<seat>.<verb>`);
|
|
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}`;
|
|
}
|
|
// `<module>.<tool>@<node>` 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<TlsOptions> {
|
|
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<tls.DetailedPeerCertificate>((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);
|
|
}
|