A person's own client, and pins that the wire did not change #13
+6
-3
@@ -4,15 +4,18 @@
|
|||||||
"description": "The Novox Mesh tool runtime \u2014 binds the mesh broker and serves the assigned modules' tools.",
|
"description": "The Novox Mesh tool runtime \u2014 binds the mesh broker and serves the assigned modules' tools.",
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"bin": {
|
"bin": {
|
||||||
"mesh-tools": "./dist/main.js"
|
"mesh-tools": "./dist/main.js",
|
||||||
|
"mesh": "./dist/mesh.js"
|
||||||
},
|
},
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"build": "tsc",
|
"build": "tsc",
|
||||||
"test": "node --test --experimental-strip-types 'test/*.test.ts'"
|
"pretest": "tsc",
|
||||||
|
"test": "node --test --test-concurrency=1 --experimental-strip-types 'test/*.test.ts'"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@novox/mesh-sdk": "^0.1.0",
|
"@novox/mesh-sdk": "^0.1.0",
|
||||||
"amqplib": "^0.10.9"
|
"amqplib": "^0.10.9",
|
||||||
|
"nats": "^2.29.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/amqplib": "^0.10.8",
|
"@types/amqplib": "^0.10.8",
|
||||||
|
|||||||
+36
-6
@@ -81,6 +81,10 @@ export async function connectAmqp(
|
|||||||
opts: { assumeExchanges?: boolean } = {},
|
opts: { assumeExchanges?: boolean } = {},
|
||||||
): Promise<Broker> {
|
): Promise<Broker> {
|
||||||
const cred: Credential = typeof target === "string" ? { url: target } : target;
|
const cred: Credential = typeof target === "string" ? { url: target } : target;
|
||||||
|
// This module's own name, for turning a local event name into this bus's routing key. From the
|
||||||
|
// credential where the mesh issued one, and from the environment for an ad-hoc client that has no
|
||||||
|
// credential of its own — the same two places subscribe already looks.
|
||||||
|
const self = cred.module ?? process.env.MESH_MODULE ?? "";
|
||||||
const conn = cred.fingerprint
|
const conn = cred.fingerprint
|
||||||
? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint))
|
? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint))
|
||||||
: await amqp.connect(cred.url);
|
: await amqp.connect(cred.url);
|
||||||
@@ -140,7 +144,7 @@ export async function connectAmqp(
|
|||||||
let eventConsumerTag: string | undefined;
|
let eventConsumerTag: string | undefined;
|
||||||
|
|
||||||
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
|
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
|
||||||
const key = msg.fields.routingKey;
|
const key = localKeyFor(msg.fields.routingKey);
|
||||||
let env: Envelope<unknown>;
|
let env: Envelope<unknown>;
|
||||||
try {
|
try {
|
||||||
env = toEnvelope(msg);
|
env = toEnvelope(msg);
|
||||||
@@ -220,7 +224,7 @@ export async function connectAmqp(
|
|||||||
await new Promise<void>((resolve, reject) => {
|
await new Promise<void>((resolve, reject) => {
|
||||||
ch.publish(
|
ch.publish(
|
||||||
EVENTS_EXCHANGE,
|
EVENTS_EXCHANGE,
|
||||||
env.key,
|
routingKeyFor(env.key, self),
|
||||||
Buffer.from(JSON.stringify(env.body)),
|
Buffer.from(JSON.stringify(env.body)),
|
||||||
{
|
{
|
||||||
persistent: true,
|
persistent: true,
|
||||||
@@ -253,7 +257,7 @@ export async function connectAmqp(
|
|||||||
}
|
}
|
||||||
eventQueue = name;
|
eventQueue = name;
|
||||||
}
|
}
|
||||||
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern);
|
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, bindingFor(pattern));
|
||||||
const sub: EventSub = { pattern, handler: handler as EventSub["handler"] };
|
const sub: EventSub = { pattern, handler: handler as EventSub["handler"] };
|
||||||
eventSubs.push(sub);
|
eventSubs.push(sub);
|
||||||
if (!eventConsumerTag) {
|
if (!eventConsumerTag) {
|
||||||
@@ -269,7 +273,7 @@ export async function connectAmqp(
|
|||||||
}
|
}
|
||||||
|
|
||||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
const { queue } = await ch.assertQueue("", { exclusive: true });
|
||||||
await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern);
|
await ch.bindQueue(queue, EVENTS_EXCHANGE, bindingFor(pattern));
|
||||||
const consumer = await ch.consume(queue, (msg) => {
|
const consumer = await ch.consume(queue, (msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
void (async () => {
|
void (async () => {
|
||||||
@@ -369,14 +373,16 @@ function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
|||||||
|
|
||||||
/** AMQP topic matching: `*` matches one word, `#` zero or more. Used to fan a shared queue's
|
/** AMQP topic matching: `*` matches one word, `#` zero or more. Used to fan a shared queue's
|
||||||
* deliveries out to the handlers whose pattern actually matches the routing key. */
|
* deliveries out to the handlers whose pattern actually matches the routing key. */
|
||||||
function topicMatches(pattern: string, key: string): boolean {
|
export function topicMatches(pattern: string, key: string): boolean {
|
||||||
return matchFrom(pattern.split("."), 0, key.split("."), 0);
|
return matchFrom(pattern.split("."), 0, key.split("."), 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean {
|
function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean {
|
||||||
while (pi < p.length) {
|
while (pi < p.length) {
|
||||||
const tok = p[pi];
|
const tok = p[pi];
|
||||||
if (tok === "#") {
|
// `**` is the mesh's wildcard for the rest of a name; `#` is this bus's, accepted so a pattern
|
||||||
|
// written either way behaves the same while both buses ship (novox/hq design 29 §1).
|
||||||
|
if (tok === "#" || tok === "**") {
|
||||||
if (pi === p.length - 1) return true; // trailing # swallows the rest, including nothing
|
if (pi === p.length - 1) return true; // trailing # swallows the rest, including nothing
|
||||||
for (let skip = ki; skip <= k.length; skip++) {
|
for (let skip = ki; skip <= k.length; skip++) {
|
||||||
if (matchFrom(p, pi + 1, k, skip)) return true;
|
if (matchFrom(p, pi + 1, k, skip)) return true;
|
||||||
@@ -390,3 +396,27 @@ function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean {
|
|||||||
}
|
}
|
||||||
return ki === k.length;
|
return ki === k.length;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** This bus spells an event as a routing key that repeats the emitter's name: `module.<emitter>.<event>`.
|
||||||
|
* A module names its events locally and the mesh derives where they land (novox/hq design 29 §1), so
|
||||||
|
* the mapping lives here rather than in every module.
|
||||||
|
*
|
||||||
|
* **Why it exists at all.** Until 04-ISSUES/127 every module passed the routing key itself, which
|
||||||
|
* worked on this bus and derived into a namespace nobody owns on the one being built. Converting the
|
||||||
|
* modules to local names without this would have broken the mesh that is actually running. */
|
||||||
|
export function routingKeyFor(key: string, self: string): string {
|
||||||
|
return key.startsWith("module.") ? key : `module.${self}.${key}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A local pattern as this bus's binding. `**` is the mesh's wildcard for the rest of a name; here
|
||||||
|
* that is `#`, and on the bus being built it is `>`. Neither spelling appears in a manifest. */
|
||||||
|
export function bindingFor(pattern: string): string {
|
||||||
|
const here = pattern.split(".").map((part) => (part === "**" ? "#" : part)).join(".");
|
||||||
|
if (here === "#") return "#";
|
||||||
|
return here.startsWith("module.") ? here : `module.${here}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A routing key as the local name a handler and a manifest both use: the emitter and the event. */
|
||||||
|
export function localKeyFor(routingKey: string): string {
|
||||||
|
return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey;
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,351 @@
|
|||||||
|
// 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 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. 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;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 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<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. `close()` stays 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.
|
||||||
|
*/
|
||||||
|
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).
|
||||||
|
*
|
||||||
|
* `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.
|
||||||
|
*/
|
||||||
|
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 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. */
|
||||||
|
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.
|
||||||
|
*
|
||||||
|
* **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}`);
|
||||||
|
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: `*` 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);
|
||||||
|
}
|
||||||
+125
@@ -0,0 +1,125 @@
|
|||||||
|
/**
|
||||||
|
* A person's client: the mesh's tools from a workstation (novox/hq design 25 §7).
|
||||||
|
*
|
||||||
|
* Two surfaces over one thing. A command line, for somebody at a terminal; an MCP server, for an
|
||||||
|
* agent. Both are adapters over the same three calls — what tools are there, what does this one take,
|
||||||
|
* call it — because a second way of reaching a tool is a second thing to keep correct.
|
||||||
|
*
|
||||||
|
* **It uses the same client a module's runtime uses.** Not a second protocol and not a bridge: a
|
||||||
|
* person connects as their own bus user, publishes on the tool subjects their account permits, and the
|
||||||
|
* server refuses anything else. So "what may this person do" is answered by the same permission list
|
||||||
|
* that answers it for a module, and there is nothing here for an audit to read separately.
|
||||||
|
*
|
||||||
|
* What a person may NOT do is the more interesting half, and none of it is enforced here — it is the
|
||||||
|
* account (design 25 §4): they cannot publish an event, so they cannot claim a module said something;
|
||||||
|
* they have no consumer, so there is no delivery to acknowledge; and they cannot answer a request, so
|
||||||
|
* they cannot impersonate a module on a bus where anyone may serve a tool.
|
||||||
|
*/
|
||||||
|
import { readFile } from "node:fs/promises";
|
||||||
|
|
||||||
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
|
import { connectNats, type Credential } from "./broker-nats.js";
|
||||||
|
|
||||||
|
/** Where the catalogue answers what tools the mesh has. */
|
||||||
|
const CATALOGUE_TOOLS = "mesh-catalog.catalog_tools";
|
||||||
|
|
||||||
|
/** A tool as the catalogue describes one. */
|
||||||
|
export interface Tool {
|
||||||
|
module: string;
|
||||||
|
name: string;
|
||||||
|
description?: string;
|
||||||
|
/** The JSON schema of what it takes, as the module declared it. */
|
||||||
|
input?: unknown;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A person's credential, as `operator issue` prints it.
|
||||||
|
*
|
||||||
|
* The same shape a module is handed, minus the parts a module needs and a person does not: no node,
|
||||||
|
* because a person is not on a machine, and no module, because they are not one.
|
||||||
|
*/
|
||||||
|
export interface PersonCredential extends Credential {
|
||||||
|
person?: string;
|
||||||
|
invokes?: string[];
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Read the credential from the file `operator issue` produced. */
|
||||||
|
export async function credentialFrom(path: string): Promise<PersonCredential> {
|
||||||
|
const raw = await readFile(path, "utf8");
|
||||||
|
let held: PersonCredential;
|
||||||
|
try {
|
||||||
|
held = JSON.parse(raw) as PersonCredential;
|
||||||
|
} catch (e) {
|
||||||
|
throw new Error(
|
||||||
|
`${path} is not a credential this mesh issued: ${(e as Error).message}. ` +
|
||||||
|
"It is the JSON `operator issue` printed, saved verbatim.",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
if (!held.url || !held.user || !held.password) {
|
||||||
|
throw new Error(
|
||||||
|
`${path} names no bus, user or password. It is the JSON \`operator issue\` printed, saved ` +
|
||||||
|
"verbatim — not an edited copy of it.",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return held;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Connect as this person. The module name the runtime wants is their own user, because every subject
|
||||||
|
* it derives is for a tool somebody else serves. */
|
||||||
|
export async function connectAs(held: PersonCredential): Promise<Broker> {
|
||||||
|
return connectNats({ ...held, module: held.user });
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* What tools the mesh has, asked of the catalogue.
|
||||||
|
*
|
||||||
|
* **Asked, not configured.** The catalogue is the only thing that knows what is installed, and a
|
||||||
|
* client carrying its own list would be a list that goes stale the first time a module is assigned —
|
||||||
|
* silently, because a tool that is not offered looks exactly like a tool that does not exist.
|
||||||
|
*/
|
||||||
|
export async function toolsOn(bus: Broker): Promise<Tool[]> {
|
||||||
|
const answered = await bus.request<Record<string, never>, { tools?: Tool[] } | Tool[]>(
|
||||||
|
CATALOGUE_TOOLS,
|
||||||
|
{},
|
||||||
|
);
|
||||||
|
const tools = Array.isArray(answered) ? answered : (answered.tools ?? []);
|
||||||
|
return tools
|
||||||
|
.slice()
|
||||||
|
.sort((a: Tool, b: Tool) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what their account
|
||||||
|
* permits — one vocabulary, so a refusal names the thing they asked for. */
|
||||||
|
export async function callTool(bus: Broker, key: string, args: unknown): Promise<unknown> {
|
||||||
|
if (!key.includes(".")) {
|
||||||
|
throw new Error(
|
||||||
|
`"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return bus.request<unknown, unknown>(key, args ?? {});
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Why a call failed, said so that the remedy is in the words.
|
||||||
|
*
|
||||||
|
* Three answers a person actually gets, and they need different things done: nobody serves that tool,
|
||||||
|
* the mesh refused this person, or the tool itself failed. Without this they are one timeout and a
|
||||||
|
* stack trace.
|
||||||
|
*/
|
||||||
|
export function whyItFailed(key: string, err: unknown): string {
|
||||||
|
const message = err instanceof Error ? err.message : String(err);
|
||||||
|
if (/no responders|503/i.test(message)) {
|
||||||
|
return `nothing serves ${key}. The module may not be assigned to any machine, or it is down — ` +
|
||||||
|
"`mesh tools` lists what the catalogue says is there.";
|
||||||
|
}
|
||||||
|
if (/permissions violation|authorization/i.test(message)) {
|
||||||
|
return `this credential may not call ${key}. What it may call was fixed when it was issued; ` +
|
||||||
|
"`operator issue` again with the tool named, or ask somebody who can.";
|
||||||
|
}
|
||||||
|
if (/timeout/i.test(message)) {
|
||||||
|
return `${key} did not answer in time. Something is serving it, so this is the tool being slow ` +
|
||||||
|
"rather than absent.";
|
||||||
|
}
|
||||||
|
return `${key} failed: ${message}`;
|
||||||
|
}
|
||||||
+139
@@ -0,0 +1,139 @@
|
|||||||
|
/**
|
||||||
|
* The mesh's tools as an MCP server, over stdio (novox/hq design 25 §7).
|
||||||
|
*
|
||||||
|
* **A thin adapter and nothing more.** Every tool an agent sees is one the catalogue listed and one
|
||||||
|
* this credential may call; the schema is the module's own; the answer is the module's own. Nothing
|
||||||
|
* here decides anything, which is why it is short — an MCP surface that reshaped arguments or
|
||||||
|
* summarised answers would be a second definition of what a tool is, and the module's manifest is the
|
||||||
|
* first.
|
||||||
|
*
|
||||||
|
* Implemented against the protocol directly rather than through a library: the surface is three
|
||||||
|
* methods and one framing, and a dependency here would be a dependency on every workstation.
|
||||||
|
*/
|
||||||
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
|
import { callTool, toolsOn, whyItFailed, type Tool } from "./client.js";
|
||||||
|
|
||||||
|
/** The protocol version this speaks. Stated, because a host that wants another should be told so
|
||||||
|
* rather than discovering it through a shape it did not expect. */
|
||||||
|
const PROTOCOL = "2024-11-05";
|
||||||
|
|
||||||
|
interface Request {
|
||||||
|
jsonrpc: string;
|
||||||
|
id?: number | string | null;
|
||||||
|
method: string;
|
||||||
|
params?: Record<string, unknown>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Serve until stdin closes, which is how a host ends a session.
|
||||||
|
*
|
||||||
|
* The tool list is fetched once, on the first `tools/list`, and kept. An agent asks for it repeatedly
|
||||||
|
* and the catalogue's answer does not change mid-session; refetching would make every turn cost a
|
||||||
|
* round trip to a module for something nobody changed.
|
||||||
|
*/
|
||||||
|
export async function serveMcp(bus: Broker, who: string): Promise<void> {
|
||||||
|
let known: Tool[] | undefined;
|
||||||
|
|
||||||
|
const say = (message: unknown) => {
|
||||||
|
process.stdout.write(`${JSON.stringify(message)}\n`);
|
||||||
|
};
|
||||||
|
const answer = (id: Request["id"], result: unknown) => say({ jsonrpc: "2.0", id, result });
|
||||||
|
const refuse = (id: Request["id"], code: number, message: string) =>
|
||||||
|
say({ jsonrpc: "2.0", id, error: { code, message } });
|
||||||
|
|
||||||
|
for await (const line of lines()) {
|
||||||
|
let request: Request;
|
||||||
|
try {
|
||||||
|
request = JSON.parse(line) as Request;
|
||||||
|
} catch {
|
||||||
|
// Unparseable, and with no id there is nobody to tell. Skipped rather than answered, because a
|
||||||
|
// reply to a request that was never framed is noise on the same channel.
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
// A notification has no id and expects no answer; `initialized` is the one every host sends.
|
||||||
|
const notification = request.id === undefined || request.id === null;
|
||||||
|
|
||||||
|
switch (request.method) {
|
||||||
|
case "initialize":
|
||||||
|
answer(request.id, {
|
||||||
|
protocolVersion: PROTOCOL,
|
||||||
|
capabilities: { tools: {} },
|
||||||
|
serverInfo: { name: "mesh", version: "1" },
|
||||||
|
// Said in the handshake, because an agent that knows whose authority it is acting under can
|
||||||
|
// say so when a call is refused — and a refusal is the one thing here that is not the
|
||||||
|
// mesh's fault or the tool's.
|
||||||
|
instructions:
|
||||||
|
`These are the tools of a Novox mesh, reached as ${who}. Every call goes to the module ` +
|
||||||
|
`that serves it; what may be called was fixed when this credential was issued, so a ` +
|
||||||
|
`refusal means the credential, not the tool.`,
|
||||||
|
});
|
||||||
|
break;
|
||||||
|
|
||||||
|
case "notifications/initialized":
|
||||||
|
break;
|
||||||
|
|
||||||
|
case "tools/list": {
|
||||||
|
try {
|
||||||
|
known ??= await toolsOn(bus);
|
||||||
|
} catch (e) {
|
||||||
|
refuse(request.id, -32603, whyItFailed("mesh-catalog.catalog_tools", e));
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
answer(request.id, {
|
||||||
|
tools: known.map((t) => ({
|
||||||
|
name: `${t.module}.${t.name}`,
|
||||||
|
description: t.description ?? `${t.name}, served by ${t.module}`,
|
||||||
|
// The module's own schema, passed through. An empty object is a tool that takes nothing,
|
||||||
|
// which is a real answer and not a missing one.
|
||||||
|
inputSchema: t.input ?? { type: "object", properties: {} },
|
||||||
|
})),
|
||||||
|
});
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
case "tools/call": {
|
||||||
|
const name = String(request.params?.name ?? "");
|
||||||
|
const args = request.params?.arguments ?? {};
|
||||||
|
try {
|
||||||
|
const result = await callTool(bus, name, args);
|
||||||
|
// Text, because that is what every host renders. The content is the module's answer as
|
||||||
|
// JSON, unshaped: an adapter that flattened it would be deciding what matters in somebody
|
||||||
|
// else's answer.
|
||||||
|
answer(request.id, {
|
||||||
|
content: [{ type: "text", text: JSON.stringify(result, null, 2) }],
|
||||||
|
});
|
||||||
|
} catch (e) {
|
||||||
|
// **An error the agent can act on, not a stack.** isError rather than a protocol failure,
|
||||||
|
// because the call was well-formed and the mesh answered it — with a refusal, an absence or
|
||||||
|
// a fault, and the words say which.
|
||||||
|
answer(request.id, {
|
||||||
|
content: [{ type: "text", text: whyItFailed(name, e) }],
|
||||||
|
isError: true,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
default:
|
||||||
|
if (!notification) {
|
||||||
|
refuse(request.id, -32601, `mesh's MCP surface has no ${request.method}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** stdin as newline-framed messages, which is what MCP over stdio is. */
|
||||||
|
async function* lines(): AsyncGenerator<string> {
|
||||||
|
let buffered = "";
|
||||||
|
for await (const chunk of process.stdin) {
|
||||||
|
buffered += (chunk as Buffer).toString("utf8");
|
||||||
|
let at: number;
|
||||||
|
while ((at = buffered.indexOf("\n")) >= 0) {
|
||||||
|
const line = buffered.slice(0, at).trim();
|
||||||
|
buffered = buffered.slice(at + 1);
|
||||||
|
if (line !== "") yield line;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (buffered.trim() !== "") yield buffered.trim();
|
||||||
|
}
|
||||||
+145
@@ -0,0 +1,145 @@
|
|||||||
|
#!/usr/bin/env node
|
||||||
|
/**
|
||||||
|
* `mesh` — the mesh's tools from a workstation, for a person (novox/hq design 25 §7).
|
||||||
|
*
|
||||||
|
* Three verbs and nothing else. What tools are there, call one, and serve the same two to an agent
|
||||||
|
* over MCP. Deliberately thin: everything that could be a decision is one the mesh already made, and a
|
||||||
|
* client that grew opinions would be a second place the mesh's behaviour is defined.
|
||||||
|
*
|
||||||
|
* mesh tools what this credential may call
|
||||||
|
* mesh call <module>.<tool> [json] call one, arguments as JSON on the command line or on stdin
|
||||||
|
* mesh mcp the same, as an MCP server over stdio
|
||||||
|
*
|
||||||
|
* The credential comes from MESH_CREDENTIAL, or --credential. It is the JSON `operator issue` printed.
|
||||||
|
*/
|
||||||
|
import { readFile } from "node:fs/promises";
|
||||||
|
|
||||||
|
import { callTool, connectAs, credentialFrom, toolsOn, whyItFailed, type Tool } from "./client.js";
|
||||||
|
import { serveMcp } from "./mcp.js";
|
||||||
|
|
||||||
|
const usage = `mesh tools
|
||||||
|
mesh call <module>.<tool> [json]
|
||||||
|
mesh mcp
|
||||||
|
|
||||||
|
--credential <file> the JSON \`operator issue\` printed; default $MESH_CREDENTIAL`;
|
||||||
|
|
||||||
|
async function main(argv: string[]): Promise<number> {
|
||||||
|
const args = [...argv];
|
||||||
|
let credentialPath = process.env.MESH_CREDENTIAL ?? "";
|
||||||
|
for (let i = 0; i < args.length; i++) {
|
||||||
|
if (args[i] === "--credential") {
|
||||||
|
credentialPath = args[i + 1] ?? "";
|
||||||
|
args.splice(i, 2);
|
||||||
|
i--;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
const verb = args.shift();
|
||||||
|
if (!verb || verb === "help" || verb === "--help") {
|
||||||
|
console.log(usage);
|
||||||
|
return verb ? 0 : 1;
|
||||||
|
}
|
||||||
|
if (!credentialPath) {
|
||||||
|
console.error(
|
||||||
|
"no credential: set MESH_CREDENTIAL or pass --credential <file>. It is the JSON " +
|
||||||
|
"`operator issue` printed, saved verbatim.",
|
||||||
|
);
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
const held = await credentialFrom(credentialPath);
|
||||||
|
const bus = await connectAs(held);
|
||||||
|
try {
|
||||||
|
switch (verb) {
|
||||||
|
case "tools":
|
||||||
|
return await listing(bus, held.person);
|
||||||
|
case "call":
|
||||||
|
return await calling(bus, args);
|
||||||
|
case "mcp":
|
||||||
|
// Serves until stdin closes, which is how an MCP host ends a session.
|
||||||
|
await serveMcp(bus, held.person ?? held.user ?? "somebody");
|
||||||
|
return 0;
|
||||||
|
default:
|
||||||
|
console.error(`mesh has no "${verb}".\n\n${usage}`);
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
await bus.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function listing(bus: Awaited<ReturnType<typeof connectAs>>, who?: string): Promise<number> {
|
||||||
|
let tools: Tool[];
|
||||||
|
try {
|
||||||
|
tools = await toolsOn(bus);
|
||||||
|
} catch (e) {
|
||||||
|
console.error(whyItFailed("mesh-catalog.catalog_tools", e));
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
if (tools.length === 0) {
|
||||||
|
console.log("the catalogue lists no tools; nothing on this mesh serves any");
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
// **What the catalogue has, not what this credential may call.** The two differ and the difference
|
||||||
|
// is the point: a person seeing only their own tools cannot tell "not installed" from "not yours",
|
||||||
|
// and those need different people to fix them.
|
||||||
|
for (const t of tools) {
|
||||||
|
const name = `${t.module}.${t.name}`;
|
||||||
|
console.log(t.description ? `${name.padEnd(36)} ${t.description}` : name);
|
||||||
|
}
|
||||||
|
if (who) {
|
||||||
|
console.log(`\nthis is what the mesh has. What ${who} may call was fixed when the credential was issued.`);
|
||||||
|
}
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function calling(
|
||||||
|
bus: Awaited<ReturnType<typeof connectAs>>,
|
||||||
|
args: string[],
|
||||||
|
): Promise<number> {
|
||||||
|
const key = args.shift();
|
||||||
|
if (!key) {
|
||||||
|
console.error("mesh call <module>.<tool> [json]");
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
|
||||||
|
let parsed: unknown = {};
|
||||||
|
if (raw.trim() !== "") {
|
||||||
|
try {
|
||||||
|
parsed = JSON.parse(raw);
|
||||||
|
} catch (e) {
|
||||||
|
console.error(`the arguments are not JSON: ${(e as Error).message}`);
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
const answer = await callTool(bus, key, parsed);
|
||||||
|
console.log(JSON.stringify(answer, null, 2));
|
||||||
|
return 0;
|
||||||
|
} catch (e) {
|
||||||
|
console.error(whyItFailed(key, e));
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Arguments on stdin, for a call whose JSON is too long or too quoted to type. Empty when stdin is a
|
||||||
|
* terminal, so `mesh call x.y` with no arguments does not hang waiting for something nobody is
|
||||||
|
* typing. */
|
||||||
|
async function maybeStdin(): Promise<string> {
|
||||||
|
if (process.stdin.isTTY) return "";
|
||||||
|
const chunks: Buffer[] = [];
|
||||||
|
for await (const chunk of process.stdin) chunks.push(chunk as Buffer);
|
||||||
|
return Buffer.concat(chunks).toString("utf8");
|
||||||
|
}
|
||||||
|
|
||||||
|
// Only when run, so a test can import the pieces.
|
||||||
|
if (process.argv[1] && import.meta.url === new URL(`file://${process.argv[1]}`).href) {
|
||||||
|
main(process.argv.slice(2))
|
||||||
|
.then((code) => process.exit(code))
|
||||||
|
.catch((e) => {
|
||||||
|
console.error(e instanceof Error ? e.message : String(e));
|
||||||
|
process.exit(1);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
export { main, usage };
|
||||||
|
export const _readFile = readFile;
|
||||||
@@ -0,0 +1,105 @@
|
|||||||
|
/**
|
||||||
|
* A person's client, against a real bus.
|
||||||
|
*
|
||||||
|
* What is worth checking is not that a request/reply works — the runtime's own tests cover that — but
|
||||||
|
* that the two surfaces are the same thing. An agent and a person must see the same tools and get the
|
||||||
|
* same answers, or the MCP surface becomes a second definition of what a tool is.
|
||||||
|
*
|
||||||
|
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||||
|
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/client.test.ts
|
||||||
|
*/
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
|
||||||
|
// The built output, not the source: the client imports its siblings as `.js`, which is what ships and
|
||||||
|
// what every other file here does, and cannot be loaded as TypeScript directly. `pretest` builds.
|
||||||
|
import { connectNats } from "../dist/broker-nats.js";
|
||||||
|
import { callTool, toolsOn, whyItFailed } from "../dist/client.js";
|
||||||
|
|
||||||
|
const url = process.env.MESH_TEST_NATS;
|
||||||
|
|
||||||
|
/** A module serving the catalogue's tool list and one tool of its own, so the client has a mesh to
|
||||||
|
* talk to. Two connections, because a person and a module are different users even in a test. */
|
||||||
|
async function aMeshWithTools() {
|
||||||
|
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||||
|
const shop = await connectNats({ url: url!, module: "shop" });
|
||||||
|
await catalogue.handle("catalog_tools", async () => ({
|
||||||
|
tools: [
|
||||||
|
{ module: "shop", name: "price", description: "what something costs", input: { type: "object" } },
|
||||||
|
{ module: "mesh-catalog", name: "catalog_tools", description: "what tools the mesh has" },
|
||||||
|
],
|
||||||
|
}));
|
||||||
|
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||||
|
return {
|
||||||
|
async close() {
|
||||||
|
await catalogue.close();
|
||||||
|
await shop.close();
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
test("a person sees what the catalogue says the mesh has, sorted", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const mesh = await aMeshWithTools();
|
||||||
|
const person = await connectNats({ url, module: "person.ada" });
|
||||||
|
try {
|
||||||
|
const tools = await toolsOn(person);
|
||||||
|
assert.deepEqual(
|
||||||
|
tools.map((x) => `${x.module}.${x.name}`),
|
||||||
|
["mesh-catalog.catalog_tools", "shop.price"],
|
||||||
|
"the list is what the catalogue answered, in a stable order",
|
||||||
|
);
|
||||||
|
} finally {
|
||||||
|
await person.close();
|
||||||
|
await mesh.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a person calls a tool and gets the module's own answer, unshaped", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const mesh = await aMeshWithTools();
|
||||||
|
const person = await connectNats({ url, module: "person.ada" });
|
||||||
|
try {
|
||||||
|
const answer = await callTool(person, "shop.price", { of: "a hat" });
|
||||||
|
assert.deepEqual(answer, { of: "a hat", cost: 12 });
|
||||||
|
} finally {
|
||||||
|
await person.close();
|
||||||
|
await mesh.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a tool nobody serves says so at once, and says what to do about it", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const person = await connectNats({ url: url!, module: "person.ada" });
|
||||||
|
try {
|
||||||
|
const began = Date.now();
|
||||||
|
await assert.rejects(() => callTool(person, "ghost.missing", {}));
|
||||||
|
// At once, not after the whole wait: "that module is down" and "that tool is slow" need
|
||||||
|
// different things done, and a timeout cannot tell them apart.
|
||||||
|
assert.ok(Date.now() - began < 5_000, "a tool nobody serves waited out the timeout");
|
||||||
|
} finally {
|
||||||
|
await person.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a name that is not <module>.<tool> is refused before anything is sent", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const person = await connectNats({ url: url!, module: "person.ada" });
|
||||||
|
try {
|
||||||
|
await assert.rejects(() => callTool(person, "price", {}), /does not name a tool/);
|
||||||
|
} finally {
|
||||||
|
await person.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("each way a call fails says what to do about it", () => {
|
||||||
|
// The three answers a person actually gets. Without this they are one timeout and a stack trace,
|
||||||
|
// and the remedies are in three different places.
|
||||||
|
assert.match(whyItFailed("shop.price", new Error("no responders")), /nothing serves shop\.price/);
|
||||||
|
assert.match(
|
||||||
|
whyItFailed("shop.price", new Error("Permissions Violation for Publish")),
|
||||||
|
/may not call shop\.price/,
|
||||||
|
);
|
||||||
|
assert.match(whyItFailed("shop.price", new Error("timeout")), /did not answer in time/);
|
||||||
|
assert.match(whyItFailed("shop.price", new Error("something else")), /something else/);
|
||||||
|
});
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
// The TypeScript implementation, held to the shared fixtures (novox/hq ADR 0074, design 19).
|
||||||
|
//
|
||||||
|
// Run against a NATS server, because the question is what actually reaches the wire:
|
||||||
|
//
|
||||||
|
// docker run -d --rm --name c -p 14222:4222 nats:2.10-alpine -js
|
||||||
|
// node test/conformance.mjs
|
||||||
|
//
|
||||||
|
// **The runner lives with the implementation it exercises; the fixture does not.** It is read
|
||||||
|
// from the sdk's conformance directory by sibling path — the same file the Go suite reads. A
|
||||||
|
// fixture copied into each implementation is two fixtures, and two fixtures drift, which is the
|
||||||
|
// failure the suite exists to prevent.
|
||||||
|
import { readFileSync } from "node:fs";
|
||||||
|
import { connect } from "nats";
|
||||||
|
|
||||||
|
const clientPath = process.argv[2] ?? "../dist/broker-nats.js";
|
||||||
|
const { connectNats } = await import(clientPath);
|
||||||
|
const f = JSON.parse(readFileSync(new URL("../../mesh-sdk/conformance/events/module-event.json", import.meta.url)));
|
||||||
|
const URL_ = process.env.MESH_TEST_NATS ?? "nats://127.0.0.1:14222";
|
||||||
|
|
||||||
|
let failed = 0;
|
||||||
|
const check = (ok, what) => { console.log(` ${ok ? "ok " : "FAIL"} ${what}`); if (!ok) failed++; };
|
||||||
|
|
||||||
|
const admin = await connect({ servers: URL_ });
|
||||||
|
const jsm = await admin.jetstreamManager();
|
||||||
|
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>", "mesh.seat.*.event.>"] });
|
||||||
|
|
||||||
|
const shop = await connectNats({ url: URL_, node: f.given.node, module: f.given.module });
|
||||||
|
await shop.publish({ key: f.given.key, node: f.given.node, body: f.given.body, headers: f.given.headers });
|
||||||
|
|
||||||
|
// What actually landed, read back from the stream rather than from the client that wrote it.
|
||||||
|
const msg = await jsm.streams.getMessage("EVENTS", { last_by_subj: f.wire.subject });
|
||||||
|
check(!!msg, `it lands on ${f.wire.subject}`);
|
||||||
|
|
||||||
|
if (msg) {
|
||||||
|
const got = {};
|
||||||
|
if (msg.header) for (const k of msg.header.keys()) got[k] = msg.header.get(k);
|
||||||
|
for (const h of f.wire.requiredHeaders) {
|
||||||
|
check(got[h] !== undefined && got[h] !== "", `${h} is set`);
|
||||||
|
}
|
||||||
|
check(got["content-type"] === f.wire.headerFormats["content-type"], "content-type is as pinned");
|
||||||
|
check(new RegExp(f.wire.headerFormats["x-event-id"]).test(got["x-event-id"]), "x-event-id is as pinned");
|
||||||
|
check(!Number.isNaN(Date.parse(got["x-time"])), "x-time parses as a date");
|
||||||
|
check(got["x-source"] === f.given.module, "x-source agrees with the subject's module");
|
||||||
|
|
||||||
|
const payload = JSON.parse(new TextDecoder().decode(msg.data));
|
||||||
|
check(payload.envelope === undefined && payload.key === undefined,
|
||||||
|
"the payload is the body alone, not the envelope (the fixture refuses nesting)");
|
||||||
|
check(JSON.stringify(payload) === JSON.stringify(f.given.body),
|
||||||
|
"the body round-trips — semantic, not byte-exact, per the README");
|
||||||
|
}
|
||||||
|
|
||||||
|
await shop.close(); await admin.close();
|
||||||
|
console.log(failed ? `\n${failed} failed` : "\nall passed");
|
||||||
|
process.exit(failed ? 1 : 0);
|
||||||
@@ -0,0 +1,124 @@
|
|||||||
|
/**
|
||||||
|
* The MCP surface, driven the way a host drives it.
|
||||||
|
*
|
||||||
|
* **The claim worth checking is that it is the same thing the command line is.** An agent and a
|
||||||
|
* person must see the same tools and get the same answers, or this becomes a second definition of what
|
||||||
|
* a tool is — which is exactly what a thin adapter is supposed to avoid.
|
||||||
|
*
|
||||||
|
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||||
|
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/mcp.test.ts
|
||||||
|
*/
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
import { spawn } from "node:child_process";
|
||||||
|
|
||||||
|
import { connectNats } from "../dist/broker-nats.js";
|
||||||
|
|
||||||
|
const url = process.env.MESH_TEST_NATS;
|
||||||
|
|
||||||
|
/** A module answering the catalogue's list and one tool, plus a credential file the client reads. */
|
||||||
|
async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void) => void }) {
|
||||||
|
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||||
|
const shop = await connectNats({ url: url!, module: "shop" });
|
||||||
|
await catalogue.handle("catalog_tools", async () => ({
|
||||||
|
tools: [{ module: "shop", name: "price", description: "what something costs" }],
|
||||||
|
}));
|
||||||
|
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||||
|
t.after(async () => {
|
||||||
|
await catalogue.close();
|
||||||
|
await shop.close();
|
||||||
|
});
|
||||||
|
|
||||||
|
const { mkdtemp, writeFile } = await import("node:fs/promises");
|
||||||
|
const { join } = await import("node:path");
|
||||||
|
const dir = await mkdtemp("/tmp/mesh-client-");
|
||||||
|
const path = join(dir, "credential.json");
|
||||||
|
await writeFile(
|
||||||
|
path,
|
||||||
|
JSON.stringify({ url, user: "person.ada", password: "x", person: "ada", invokes: ["shop.price"] }),
|
||||||
|
);
|
||||||
|
return path;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Drive `mesh mcp` over stdio and collect the replies, as a host would. */
|
||||||
|
function driving(credential: string, requests: unknown[]): Promise<Record<string, any>[]> {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
const child = spawn(process.execPath, ["dist/mesh.js", "mcp", "--credential", credential], {
|
||||||
|
stdio: ["pipe", "pipe", "pipe"],
|
||||||
|
});
|
||||||
|
let out = "";
|
||||||
|
let err = "";
|
||||||
|
child.stdout.on("data", (d) => (out += d.toString()));
|
||||||
|
child.stderr.on("data", (d) => (err += d.toString()));
|
||||||
|
child.on("error", reject);
|
||||||
|
child.on("close", () => {
|
||||||
|
const replies = out
|
||||||
|
.split("\n")
|
||||||
|
.filter((l) => l.trim() !== "")
|
||||||
|
.map((l) => JSON.parse(l) as Record<string, any>);
|
||||||
|
if (replies.length === 0 && err !== "") reject(new Error(err));
|
||||||
|
else resolve(replies);
|
||||||
|
});
|
||||||
|
for (const r of requests) child.stdin.write(`${JSON.stringify(r)}\n`);
|
||||||
|
child.stdin.end();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
test("a host initialises, lists the mesh's tools and calls one", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const credential = await aMeshAndACredential(t);
|
||||||
|
|
||||||
|
const replies = await driving(credential, [
|
||||||
|
{ jsonrpc: "2.0", id: 1, method: "initialize", params: {} },
|
||||||
|
{ jsonrpc: "2.0", method: "notifications/initialized" },
|
||||||
|
{ jsonrpc: "2.0", id: 2, method: "tools/list" },
|
||||||
|
{ jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "shop.price", arguments: { of: "a hat" } } },
|
||||||
|
]);
|
||||||
|
|
||||||
|
const byId = new Map(replies.map((r) => [r.id, r]));
|
||||||
|
// A notification is answered with nothing, or a host waiting on ids sees a reply it cannot match.
|
||||||
|
assert.equal(replies.length, 3, `expected three replies, got ${JSON.stringify(replies)}`);
|
||||||
|
|
||||||
|
const hello = byId.get(1)!.result;
|
||||||
|
assert.equal(hello.protocolVersion, "2024-11-05");
|
||||||
|
assert.ok(hello.capabilities.tools, "a server offering no tools is not this one");
|
||||||
|
assert.match(hello.instructions, /ada/, "the handshake says whose authority a call is made under");
|
||||||
|
|
||||||
|
const listed = byId.get(2)!.result.tools;
|
||||||
|
assert.equal(listed.length, 1);
|
||||||
|
assert.equal(listed[0].name, "shop.price", "a tool is named the way a person names it");
|
||||||
|
assert.ok(listed[0].inputSchema, "a tool with no schema is one an agent cannot call");
|
||||||
|
|
||||||
|
const called = byId.get(3)!.result;
|
||||||
|
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
|
||||||
|
// The module's own answer, unshaped. An adapter that summarised it would be deciding what matters
|
||||||
|
// in somebody else's answer.
|
||||||
|
assert.deepEqual(JSON.parse(called.content[0].text), { of: "a hat", cost: 12 });
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a tool nobody serves comes back as an error the agent can act on", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const credential = await aMeshAndACredential(t);
|
||||||
|
|
||||||
|
const replies = await driving(credential, [
|
||||||
|
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "ghost.missing", arguments: {} } },
|
||||||
|
]);
|
||||||
|
const result = replies[0].result;
|
||||||
|
// isError, not a protocol failure: the call was well-formed and the mesh answered it — with an
|
||||||
|
// absence. A JSON-RPC error would tell the agent its request was malformed, which it was not.
|
||||||
|
assert.ok(result?.isError, `expected a tool error, got ${JSON.stringify(replies[0])}`);
|
||||||
|
assert.match(result.content[0].text, /nothing serves ghost\.missing/);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a method this surface does not have is refused, and a notification is not", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const credential = await aMeshAndACredential(t);
|
||||||
|
|
||||||
|
const replies = await driving(credential, [
|
||||||
|
{ jsonrpc: "2.0", id: 1, method: "resources/list" },
|
||||||
|
{ jsonrpc: "2.0", method: "notifications/cancelled" },
|
||||||
|
]);
|
||||||
|
assert.equal(replies.length, 1, "a notification was answered");
|
||||||
|
assert.equal(replies[0].error.code, -32601);
|
||||||
|
assert.match(replies[0].error.message, /resources\/list/);
|
||||||
|
});
|
||||||
@@ -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);
|
||||||
@@ -0,0 +1,68 @@
|
|||||||
|
/**
|
||||||
|
* **The old bus's wire is byte-identical after the rename, and this is the test that lets the change
|
||||||
|
* be merged to a running mesh.**
|
||||||
|
*
|
||||||
|
* Every module's event names were converted from the old bus's routing keys to local names
|
||||||
|
* (novox/hq 04-ISSUES/127), and the old bus's client maps them back. If that mapping is wrong
|
||||||
|
* anywhere, a live mesh's events stop being delivered — silently, because a binding that matches
|
||||||
|
* nothing is not an error.
|
||||||
|
*
|
||||||
|
* So this pins the mapping against the literal routing keys the mesh used before, taken from the
|
||||||
|
* manifests as they were. It needs no bus: it is about a string.
|
||||||
|
*/
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
|
||||||
|
import { routingKeyFor, bindingFor, localKeyFor, topicMatches } from "../dist/broker-amqp.js";
|
||||||
|
|
||||||
|
test("a converted emit produces the routing key the mesh published before", () => {
|
||||||
|
// left: what the module's code says now. right: what went on the wire before, unchanged.
|
||||||
|
const same: [string, string, string][] = [
|
||||||
|
["plex", "playback.started", "module.plex.playback.started"],
|
||||||
|
["sonarr", "download.completed", "module.sonarr.download.completed"],
|
||||||
|
["builder", "built", "module.builder.built"],
|
||||||
|
["mesh-catalog", "upgraded", "module.mesh-catalog.upgraded"],
|
||||||
|
["keycloak", "user.created", "module.keycloak.user.created"],
|
||||||
|
["mesh-vault", "secret.rotated", "module.mesh-vault.secret.rotated"],
|
||||||
|
];
|
||||||
|
for (const [self, local, before] of same) {
|
||||||
|
assert.equal(routingKeyFor(local, self), before, `${self} emitting ${local}`);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a converted subscription binds what it bound before", () => {
|
||||||
|
const same: [string, string][] = [
|
||||||
|
["builder.built", "module.builder.built"],
|
||||||
|
["*.download.completed", "module.*.download.completed"],
|
||||||
|
["*.usage.*", "module.*.usage.*"],
|
||||||
|
// The audit logger's "everything": `#` on this bus, and it must stay `#`.
|
||||||
|
["**", "#"],
|
||||||
|
];
|
||||||
|
for (const [declared, before] of same) {
|
||||||
|
assert.equal(bindingFor(declared), before, `consuming ${declared}`);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a handler still matches what the bus delivers", () => {
|
||||||
|
// The key a handler is given is the local one now, and the pattern it compares against is local
|
||||||
|
// too — so the pair must still meet for every case the mesh actually has.
|
||||||
|
const pairs: [string, string][] = [
|
||||||
|
["builder.built", "module.builder.built"],
|
||||||
|
["*.download.completed", "module.sonarr.download.completed"],
|
||||||
|
["*.usage.*", "module.anthropic-consumer.usage.session"],
|
||||||
|
["**", "module.anything.at.all"],
|
||||||
|
];
|
||||||
|
for (const [pattern, delivered] of pairs) {
|
||||||
|
assert.ok(
|
||||||
|
topicMatches(pattern, localKeyFor(delivered)),
|
||||||
|
`${pattern} no longer matches ${delivered}, so a running module would stop reacting`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a routing key already in the old form is left alone", () => {
|
||||||
|
// Belt for the transition: anything not yet converted still goes out as it did, so a module built
|
||||||
|
// from an older manifest keeps working beside one built from a current manifest.
|
||||||
|
assert.equal(routingKeyFor("module.plex.playback.started", "plex"), "module.plex.playback.started");
|
||||||
|
assert.equal(bindingFor("module.builder.built"), "module.builder.built");
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user