From c1517c39a0878b5ceaa80971674ee38e4bc58528 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 23:29:17 +0200 Subject: [PATCH 1/6] The tool runtime's client on NATS, behind the unchanged contract MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Task 3.6 of novox/hq ADR 0116. A module is still written against request, handle, publish, subscribe, close; only what is underneath changes. main.ts still selects the AMQP client — steps 1 to 4 leave every node on AMQP, so this ships beside it and is selected at the rollout. Round-tripped against a real server (test/roundtrip.mjs): a tool answered across two connections, a throwing handler reaching the caller as an error rather than a timeout, an event delivered once with its key, body, node and event id intact, and an event landing under its emitter's own namespace. Three things the compiler and the server corrected: - the envelope's field is `key`, not `type`, and the payload is `env.body` with metadata in headers — not the whole envelope re-encoded. An implementation that nested the envelope would pass all its own tests and agree with nobody, which is what the conformance suite exists to stop. - the NATS client's TLS options are PEM strings with no verify hook, so the AMQP client's `checkServerIdentity: () => undefined` has no equivalent. The fingerprint check still happens and is still the guarantee, but the bus's certificate must now carry a SAN matching the address nodes dial. That is a constraint on the mesh's certificates, recorded where it bites. - a durable consumer is bound, never created: a module's account cannot reach the JetStream API, and a runtime creating its own would be a module choosing its own delivery semantics. --- package.json | 3 +- src/broker-nats.ts | 341 +++++++++++++++++++++++++++++++++++++++++++++ test/roundtrip.mjs | 57 ++++++++ 3 files changed, 400 insertions(+), 1 deletion(-) create mode 100644 src/broker-nats.ts create mode 100644 test/roundtrip.mjs diff --git a/package.json b/package.json index 3ed7292..7c04aa6 100644 --- a/package.json +++ b/package.json @@ -12,7 +12,8 @@ }, "dependencies": { "@novox/mesh-sdk": "^0.1.0", - "amqplib": "^0.10.9" + "amqplib": "^0.10.9", + "nats": "^2.29.0" }, "devDependencies": { "@types/amqplib": "^0.10.8", diff --git a/src/broker-nats.ts b/src/broker-nats.ts new file mode 100644 index 0000000..caae49f --- /dev/null +++ b/src/broker-nats.ts @@ -0,0 +1,341 @@ +// The tool runtime's broker client, on NATS. +// +// **The sdk's contract does not change** (novox/hq ADR 0106, ADR 0039): a module is written +// against `request`, `handle`, `publish`, `subscribe`, `close`, and the runtime implements them. +// That is why a module built before any of this runs on the new runtime without a rebuild, and +// why the sdk's own diff for the whole bus change is three comments. +// +// What changes is underneath: exchanges and per-tool queues become subjects, and durability +// becomes JetStream (novox/hq design 25, design 29). +// +// mesh.mod..event. an event this module emits +// mesh.mod..tool. a tool this module serves +// mesh.seat..accept. work submitted to a role +// +// The module never writes one of those: it names its events and tools locally and the mesh +// derives the subject (design 29 §1), so reorganising the subject space leaves every module +// correct. + +import { createHash } from "node:crypto"; +import tls from "node:tls"; +import { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription } from "nats"; +import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; + +const sc = StringCodec(); + +/** Requests wait this long for an answer before failing, matching the AMQP client's behaviour so + * a module's timeout handling does not change with the transport. */ +const REQUEST_TIMEOUT_MS = 30_000; + +export class PinMismatchError extends Error {} + +/** A broker credential as the mesh delivers it (novox/hq ADR 0120): the bus's address, the + * fingerprint of the certificate it must present, and the node and module the account is scoped + * to — the runtime derives its subjects from those rather than being told them. */ +export interface Credential { + url: string; + fingerprint?: string; + node?: string; + module?: string; + user?: string; + password?: string; +} + +/** Whether a connection failure is worth retrying, or is a fact about this configuration that + * retrying cannot change. Mirrors the AMQP client's judgement so the runtime's supervisor does + * not have to know which transport it is on. */ +export function fatalBrokerReason(err: unknown): string | null { + if (err instanceof PinMismatchError) return "the bus's certificate does not match the pin"; + const e = err as { code?: string; message?: string }; + const message = typeof e?.message === "string" ? e.message : String(err); + if (e?.code === "ERR_INVALID_URL" || /invalid url/i.test(message)) { + return "the bus address is not a usable URL"; + } + if (/authorization violation|user authentication expired|permissions violation/i.test(message)) { + return "the bus refused this account"; + } + return null; +} + +/** + * Connect to the mesh bus and return a Broker. + * + * **A module's subjects come from its credential, not from its calls.** `node` and `module` name + * the account the mesh issued, and every subject this client publishes or subscribes is derived + * from them — so a module cannot name another's namespace even by mistake, and what it emits + * matches what the mesh authorised (ADR 0074's identity rule). + */ +export async function connectNats( + target: string | Credential, + opts: { module?: string } = {}, +): Promise { + const cred: Credential = typeof target === "string" ? { url: target } : target; + const self = cred.module ?? opts.module; + if (!self) { + throw new Error( + "a broker credential with no module: the runtime derives its subjects from the account " + + "the mesh issued, and cannot guess which module it is", + ); + } + + const conn = await natsConnect({ + servers: cred.url, + user: cred.user, + pass: cred.password, + name: `${cred.node ?? "?"}.${self}`, + tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined, + // Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on + // the mesh to exit. The AMQP client's supervisor did this a level up; here the library does + // it, and `close()` is still the only thing that ends the connection. + maxReconnectAttempts: -1, + }); + const js = conn.jetstream(); + + const subs: Subscription[] = []; + let closed = false; + + return { + /** + * Ask one question and await one answer. + * + * Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3), + * and a lost one is a timeout the caller already handles. The reply travels on the inbox the + * request carries, which the responder may answer because its account has `allow_responses` + * — one reply to a message it actually received, and nothing wider. + */ + async request(key: string, body: Req): Promise { + const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), { + timeout: REQUEST_TIMEOUT_MS, + }); + const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string }; + if (reply.error) throw new Error(reply.error); + return reply.result as Res; + }, + + /** + * Answer a question. + * + * A queue group, so several nodes may serve one tool and exactly one of them answers each + * call — the same property the AMQP client got from a shared durable queue. + */ + async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { + const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` }); + subs.push(sub); + void (async () => { + for await (const msg of sub) { + let reply: { result?: Res; error?: string }; + try { + reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) }; + } catch (err) { + // The caller is told, rather than left to time out: a handler that threw is a + // different failure from a tool nobody serves, and only one of them is worth retrying. + reply = { error: err instanceof Error ? err.message : String(err) }; + } + msg.respond(sc.encode(JSON.stringify(reply))); + } + })(); + return () => { + sub.unsubscribe(); + }; + }, + + /** + * Emit an event. + * + * Published into JetStream and awaited, so a publish the bus never accepted fails the emit + * rather than vanishing — at-least-once starts at the emitter, not only the consumer + * (ADR 0042), which is what the AMQP client's confirm channel was for. + * + * `msgID` is the event's own id, so a redelivery after a crash between publishing and + * acknowledging is de-duplicated by the server inside its window rather than seen twice. + */ + async publish(env: Envelope): Promise { + // **The body is the payload and the metadata rides as headers**, exactly as on AMQP + // (ADR 0042). NATS has headers of its own, so the envelope's shape on the wire is + // preserved rather than re-encoded — which matters because that shape is what the + // conformance suite pins, and an implementation that nested the whole envelope in the + // body would pass every one of its own tests and agree with nobody. + const meta = (env.headers ?? {}) as Record; + 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( + pattern: string, + handler: (env: Envelope) => Promise, + ): 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 { + 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( + msg: JsMsg, + pattern: string, + handler: (env: Envelope) => Promise, +): Promise { + let env: Envelope; + try { + env = toEnvelope(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(msg: JsMsg): Envelope { + const headers: Record = {}; + 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..event.`. The + // module never sees the subject, only the key it declared. + key: keyFromSubject(msg.subject), + node: headers["x-node"] ?? "", + body: JSON.parse(sc.decode(msg.data)) as T, + headers: headers as EventHeaders, + }; +} + +/** The event key inside a module's event subject. */ +function keyFromSubject(subject: string): string { + const marker = ".event."; + const at = subject.indexOf(marker); + return at < 0 ? subject : subject.slice(at + marker.length); +} + +/** A module's own event subject. Derived, never taken from the caller: the module names its + * event and the mesh decides where it lands (design 29 §1). */ +function eventSubject(type: string, self: string): string { + return `mesh.mod.${self}.event.${type}`; +} + +/** A tool's subject. A bare name is this module's own tool; `.` addresses + * another's, which is how a request reaches a module that is not this one. */ +function toolSubject(key: string, self: string): string { + const dot = key.indexOf("."); + if (dot < 0) return `mesh.mod.${self}.tool.${key}`; + return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`; +} + +function normalizeFingerprint(fingerprint: string): string { + return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase(); +} + +/** + * Dial once to see the certificate, and refuse unless it is exactly the one the mesh pinned. + * A certificate authority is not consulted: the mesh issued this and knows its fingerprint, + * which is stronger than trusting whoever a machine's trust store happens to contain. + * + * **One behaviour differs from the AMQP client, and it is a constraint on the mesh rather than a + * detail of this file.** That client passed `checkServerIdentity: () => undefined`, because + * pinning the exact certificate makes hostname verification redundant. The NATS client exposes no + * such hook — its TLS options are file paths and PEM strings, with no verify callback — so the + * certificate the mesh issues the bus **must carry a subject-alternative name matching the address + * nodes dial it by**. The fingerprint check below still happens and is still the real guarantee; + * what cannot be switched off is the check *beside* it. + */ +async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> { + const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`); + const port = url.port ? Number(url.port) : 4222; + const certificate = await new Promise((resolve, reject) => { + const socket = tls.connect( + { host: url.hostname, port, rejectUnauthorized: false, servername: url.hostname }, + () => { + const peer = socket.getPeerCertificate(true); + socket.end(); + resolve(peer); + }, + ); + socket.on("error", reject); + }); + const seen = createHash("sha256").update(certificate.raw).digest("hex"); + if (seen !== normalizeFingerprint(fingerprint)) { + throw new PinMismatchError( + `the bus at ${url.hostname}:${port} presented ${seen}, not the pinned ${normalizeFingerprint(fingerprint)}`, + ); + } + const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`; + return { ca: pem }; +} + +/** The mesh's topic matching, unchanged from AMQP: `*` is one token, `#` the rest. Kept because + * it is the module's vocabulary — a module's `consumes` pattern reads the same as it always did, + * and the subject it becomes is the mesh's business. */ +export function topicMatches(pattern: string, key: string): boolean { + return matchFrom(pattern.split("."), 0, key.split("."), 0); +} + +function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { + if (pi === p.length) return ki === k.length; + if (p[pi] === "#") { + for (let skip = ki; skip <= k.length; skip++) { + if (matchFrom(p, pi + 1, k, skip)) return true; + } + return false; + } + if (ki === k.length) return false; + if (p[pi] !== "*" && p[pi] !== k[ki]) return false; + return matchFrom(p, pi + 1, k, ki + 1); +} diff --git a/test/roundtrip.mjs b/test/roundtrip.mjs new file mode 100644 index 0000000..702fef8 --- /dev/null +++ b/test/roundtrip.mjs @@ -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); From 19560ca6a705c1f58e8c2662ad4af0664e525da9 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 23:40:59 +0200 Subject: [PATCH 2/6] Hold the runtime's NATS client to the shared fixtures Read back from the stream rather than from the client that wrote it, so the check is what reached the wire. The runner lives with the implementation; the fixture stays in one place. --- test/conformance.mjs | 54 ++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 54 insertions(+) create mode 100644 test/conformance.mjs diff --git a/test/conformance.mjs b/test/conformance.mjs new file mode 100644 index 0000000..0d8517c --- /dev/null +++ b/test/conformance.mjs @@ -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); From 2197c36fefbcfd40a6dc5edda2bfa2c8b3344f59 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 23:51:00 +0200 Subject: [PATCH 3/6] Describe the client on its own terms MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Same cleanup: the comments explained each decision by contrast with what came before instead of stating it. The certificate constraint stays — it is a fact about the mesh's certificates, not a comparison. --- src/broker-nats.ts | 46 +++++++++++++++++++++------------------------- 1 file changed, 21 insertions(+), 25 deletions(-) diff --git a/src/broker-nats.ts b/src/broker-nats.ts index caae49f..bd82bbb 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -5,8 +5,8 @@ // That is why a module built before any of this runs on the new runtime without a rebuild, and // why the sdk's own diff for the whole bus change is three comments. // -// What changes is underneath: exchanges and per-tool queues become subjects, and durability -// becomes JetStream (novox/hq design 25, design 29). +// Underneath, everything is a subject and durability is JetStream (novox/hq design 25, +// design 29). // // mesh.mod..event. an event this module emits // mesh.mod..tool. a tool this module serves @@ -23,8 +23,8 @@ import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; const sc = StringCodec(); -/** Requests wait this long for an answer before failing, matching the AMQP client's behaviour so - * a module's timeout handling does not change with the transport. */ +/** Requests wait this long for an answer before failing. Unchanged from what modules already + * expect, so a module's timeout handling is not something the bus quietly redefines. */ const REQUEST_TIMEOUT_MS = 30_000; export class PinMismatchError extends Error {} @@ -42,8 +42,8 @@ export interface Credential { } /** Whether a connection failure is worth retrying, or is a fact about this configuration that - * retrying cannot change. Mirrors the AMQP client's judgement so the runtime's supervisor does - * not have to know which transport it is on. */ + * retrying cannot change. The runtime's supervisor asks this and does not need to know what it + * is connected to. */ export function fatalBrokerReason(err: unknown): string | null { if (err instanceof PinMismatchError) return "the bus's certificate does not match the pin"; const e = err as { code?: string; message?: string }; @@ -85,8 +85,7 @@ export async function connectNats( name: `${cred.node ?? "?"}.${self}`, tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined, // Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on - // the mesh to exit. The AMQP client's supervisor did this a level up; here the library does - // it, and `close()` is still the only thing that ends the connection. + // the mesh to exit. `close()` stays the only thing that ends the connection. maxReconnectAttempts: -1, }); const js = conn.jetstream(); @@ -116,7 +115,7 @@ export async function connectNats( * Answer a question. * * A queue group, so several nodes may serve one tool and exactly one of them answers each - * call — the same property the AMQP client got from a shared durable queue. + * call. */ async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` }); @@ -144,17 +143,15 @@ export async function connectNats( * * Published into JetStream and awaited, so a publish the bus never accepted fails the emit * rather than vanishing — at-least-once starts at the emitter, not only the consumer - * (ADR 0042), which is what the AMQP client's confirm channel was for. + * (ADR 0042). * * `msgID` is the event's own id, so a redelivery after a crash between publishing and * acknowledging is de-duplicated by the server inside its window rather than seen twice. */ async publish(env: Envelope): Promise { - // **The body is the payload and the metadata rides as headers**, exactly as on AMQP - // (ADR 0042). NATS has headers of its own, so the envelope's shape on the wire is - // preserved rather than re-encoded — which matters because that shape is what the - // conformance suite pins, and an implementation that nested the whole envelope in the - // body would pass every one of its own tests and agree with nobody. + // **The body is the payload and the metadata rides as headers** (ADR 0042). That shape + // is what the conformance suite pins: an implementation that nested the whole envelope in + // the body would pass every one of its own tests and agree with nobody. const meta = (env.headers ?? {}) as Record; const h = natsHeaders(); for (const [k, v] of Object.entries(meta)) { @@ -288,13 +285,12 @@ function normalizeFingerprint(fingerprint: string): string { * A certificate authority is not consulted: the mesh issued this and knows its fingerprint, * which is stronger than trusting whoever a machine's trust store happens to contain. * - * **One behaviour differs from the AMQP client, and it is a constraint on the mesh rather than a - * detail of this file.** That client passed `checkServerIdentity: () => undefined`, because - * pinning the exact certificate makes hostname verification redundant. The NATS client exposes no - * such hook — its TLS options are file paths and PEM strings, with no verify callback — so the - * certificate the mesh issues the bus **must carry a subject-alternative name matching the address - * nodes dial it by**. The fingerprint check below still happens and is still the real guarantee; - * what cannot be switched off is the check *beside* it. + * **A constraint on the mesh, not a detail of this file.** Pinning the exact certificate makes + * hostname verification redundant in principle, but the NATS client exposes no hook to replace + * it — its TLS options are file paths and PEM strings, with no verify callback. So the + * certificate the mesh issues the bus **must carry a subject-alternative name matching the + * address nodes dial it by**. The fingerprint check below still happens and is still the real + * guarantee; what cannot be switched off is the check *beside* it. */ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> { const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`); @@ -320,9 +316,9 @@ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: str return { ca: pem }; } -/** The mesh's topic matching, unchanged from AMQP: `*` is one token, `#` the rest. Kept because - * it is the module's vocabulary — a module's `consumes` pattern reads the same as it always did, - * and the subject it becomes is the mesh's business. */ +/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary — + * a module's `consumes` pattern is matched here, and the subject it becomes is the mesh's + * business, not the module's. */ export function topicMatches(pattern: string, key: string): boolean { return matchFrom(pattern.split("."), 0, key.split("."), 0); } From fbeb373d1a66cae51360ec232183e262b06fcba7 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 27 Sep 2026 14:42:40 +0200 Subject: [PATCH 4/6] Both clients map local event names to their own wire MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A module names its events locally and each transport works out where they land. That is what design 29 says and what neither client did: both passed the name straight through, which happened to be right on the old bus because modules were writing routing keys, and wrong on the new one (novox/hq 04-ISSUES/127). The old bus's client now turns a local name into `module..` on the way out and back on the way in. Without that, converting the modules to local names would have broken the mesh that is actually running. **A handler and a manifest now say the same thing.** The key a module sees was the event name alone, so a manifest declaring `consumes: builder.built` produced a pattern that could never match what it was compared against — and a module consuming one event from two emitters could only tell them apart by reading a header. The subject already carries the emitter, so naming it in the key makes a mismatch between manifest and code a typo instead of a category error. Both matchers accept `**` for the rest of a name, which is how a manifest spells it; the old bus's `#` still works, because both buses ship until the rollout. --- src/broker-amqp.ts | 40 +++++++++++++++++++++++++++++++++++----- src/broker-nats.ts | 20 +++++++++++++++++--- 2 files changed, 52 insertions(+), 8 deletions(-) diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index eee7b03..62c5e51 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -81,6 +81,10 @@ export async function connectAmqp( opts: { assumeExchanges?: boolean } = {}, ): Promise { 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 ? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint)) : await amqp.connect(cred.url); @@ -140,7 +144,7 @@ export async function connectAmqp( let eventConsumerTag: string | undefined; async function dispatchEvent(msg: amqp.ConsumeMessage): Promise { - const key = msg.fields.routingKey; + const key = localKeyFor(msg.fields.routingKey); let env: Envelope; try { env = toEnvelope(msg); @@ -220,7 +224,7 @@ export async function connectAmqp( await new Promise((resolve, reject) => { ch.publish( EVENTS_EXCHANGE, - env.key, + routingKeyFor(env.key, self), Buffer.from(JSON.stringify(env.body)), { persistent: true, @@ -253,7 +257,7 @@ export async function connectAmqp( } 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"] }; eventSubs.push(sub); if (!eventConsumerTag) { @@ -269,7 +273,7 @@ export async function connectAmqp( } 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) => { if (!msg) return; void (async () => { @@ -376,7 +380,9 @@ function topicMatches(pattern: string, key: string): boolean { function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { while (pi < p.length) { 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 for (let skip = ki; skip <= k.length; skip++) { 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; } + +/** This bus spells an event as a routing key that repeats the emitter's name: `module..`. + * 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. */ +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. */ +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. */ +function localKeyFor(routingKey: string): string { + return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey; +} diff --git a/src/broker-nats.ts b/src/broker-nats.ts index bd82bbb..f6f4ff7 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -255,11 +255,23 @@ function toEnvelope(msg: JsMsg): Envelope { }; } -/** The event key inside a module's event subject. */ +/** 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); - return at < 0 ? subject : subject.slice(at + marker.length); + if (at < 0) return subject; + const event = subject.slice(at + marker.length); + // `mesh.mod..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 @@ -325,7 +337,9 @@ export function topicMatches(pattern: string, key: string): boolean { function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { if (pi === p.length) return ki === k.length; - if (p[pi] === "#") { + // `**` 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; } From 9acc40145a27588cd761d58ed711b1ab5abeee71 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 27 Sep 2026 17:03:13 +0200 Subject: [PATCH 5/6] A person's client: the mesh's tools from a workstation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Design 25 §7's second item. Two surfaces over one thing — a command line for somebody at a terminal, an MCP server for an agent — and both are adapters over the same three calls: what tools are there, what does this one take, call it. A second way of reaching a tool would be a second thing to keep correct. It uses the client a module's runtime uses. Not a bridge and not a second protocol: a person connects as their own bus user and publishes on the tool subjects their account permits, so "what may this person do" is answered by the same permission list that answers it for a module, and an audit has nothing separate to read. `mesh tools` lists what the *catalogue* has, not what this credential may call. The two differ and the difference is the point: somebody seeing only their own tools cannot tell "not installed" from "not yours", and those need different people to fix them. A failed call says which of three things happened, because the remedies are in three different places: nobody serves that tool, this credential may not call it, or the tool itself was slow. Without that they are one timeout and a stack trace. The MCP surface decides nothing. The tool names are the ones a person types, the schemas are the modules' own, and an answer is passed through unshaped — an adapter that summarised somebody else's answer would be deciding what matters in it. A tool that fails comes back as a tool error rather than a protocol error, because the request was well-formed and the mesh answered it. Written against the protocol directly: it is three methods and one framing, and a dependency here would be a dependency on every workstation. Tests drive both surfaces against a real bus, including that a host's notification is answered with nothing and an unknown method is refused. They run one file at a time, because each stands up a module serving the same tool subjects and run together their requests get split between them — which showed up as one test reading another's answer. --- package.json | 6 +- src/client.ts | 125 ++++++++++++++++++++++++++++++++++++++ src/mcp.ts | 139 ++++++++++++++++++++++++++++++++++++++++++ src/mesh.ts | 145 ++++++++++++++++++++++++++++++++++++++++++++ test/client.test.ts | 105 ++++++++++++++++++++++++++++++++ test/mcp.test.ts | 124 +++++++++++++++++++++++++++++++++++++ 6 files changed, 642 insertions(+), 2 deletions(-) create mode 100644 src/client.ts create mode 100644 src/mcp.ts create mode 100644 src/mesh.ts create mode 100644 test/client.test.ts create mode 100644 test/mcp.test.ts diff --git a/package.json b/package.json index 7c04aa6..6e420bb 100644 --- a/package.json +++ b/package.json @@ -4,11 +4,13 @@ "description": "The Novox Mesh tool runtime \u2014 binds the mesh broker and serves the assigned modules' tools.", "type": "module", "bin": { - "mesh-tools": "./dist/main.js" + "mesh-tools": "./dist/main.js", + "mesh": "./dist/mesh.js" }, "scripts": { "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": { "@novox/mesh-sdk": "^0.1.0", diff --git a/src/client.ts b/src/client.ts new file mode 100644 index 0000000..cbbdf43 --- /dev/null +++ b/src/client.ts @@ -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 { + 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 { + 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 { + const answered = await bus.request, { 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 `.`, 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 { + if (!key.includes(".")) { + throw new Error( + `"${key}" does not name a tool: write ., as \`mesh tools\` lists them`, + ); + } + return bus.request(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}`; +} diff --git a/src/mcp.ts b/src/mcp.ts new file mode 100644 index 0000000..f512753 --- /dev/null +++ b/src/mcp.ts @@ -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; +} + +/** + * 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 { + 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 { + 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(); +} diff --git a/src/mesh.ts b/src/mesh.ts new file mode 100644 index 0000000..013df0b --- /dev/null +++ b/src/mesh.ts @@ -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 . [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 . [json] +mesh mcp + + --credential the JSON \`operator issue\` printed; default $MESH_CREDENTIAL`; + +async function main(argv: string[]): Promise { + 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 . 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>, who?: string): Promise { + 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>, + args: string[], +): Promise { + const key = args.shift(); + if (!key) { + console.error("mesh call . [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 { + 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; diff --git a/test/client.test.ts b/test/client.test.ts new file mode 100644 index 0000000..ca2e4aa --- /dev/null +++ b/test/client.test.ts @@ -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 . 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/); +}); diff --git a/test/mcp.test.ts b/test/mcp.test.ts new file mode 100644 index 0000000..fca4b51 --- /dev/null +++ b/test/mcp.test.ts @@ -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 }) { + 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[]> { + 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); + 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/); +}); From 3d55aeb1d8d281404aea0d919f7d946026c6bef4 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 27 Sep 2026 17:37:19 +0200 Subject: [PATCH 6/6] The old bus's wire is pinned unchanged, because this has to merge to a running mesh MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every module's event names were converted from the old bus's routing keys to local names, and this 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 the mapping is pinned against the literal routing keys the mesh published before, taken from the manifests as they were: what each module now emits, what each now binds, and that a handler still matches what the bus delivers. Including the audit logger's "everything", which must stay `#` on this bus. And a key already in the old form is left alone, so a module built from an older manifest keeps working beside one built from a current manifest — which is the state the mesh will actually be in between deployments. --- src/broker-amqp.ts | 8 ++--- test/wire-unchanged.test.ts | 68 +++++++++++++++++++++++++++++++++++++ 2 files changed, 72 insertions(+), 4 deletions(-) create mode 100644 test/wire-unchanged.test.ts diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index 62c5e51..9b8db7c 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -373,7 +373,7 @@ function toEnvelope(msg: amqp.ConsumeMessage): Envelope { /** 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. */ -function topicMatches(pattern: string, key: string): boolean { +export function topicMatches(pattern: string, key: string): boolean { return matchFrom(pattern.split("."), 0, key.split("."), 0); } @@ -404,19 +404,19 @@ function matchFrom(p: string[], pi: number, k: string[], ki: number): boolean { * **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. */ -function routingKeyFor(key: string, self: string): string { +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. */ -function bindingFor(pattern: string): string { +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. */ -function localKeyFor(routingKey: string): string { +export function localKeyFor(routingKey: string): string { return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey; } diff --git a/test/wire-unchanged.test.ts b/test/wire-unchanged.test.ts new file mode 100644 index 0000000..2d80144 --- /dev/null +++ b/test/wire-unchanged.test.ts @@ -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"); +});