diff --git a/package.json b/package.json index 3ed7292..6e420bb 100644 --- a/package.json +++ b/package.json @@ -4,15 +4,18 @@ "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", - "amqplib": "^0.10.9" + "amqplib": "^0.10.9", + "nats": "^2.29.0" }, "devDependencies": { "@types/amqplib": "^0.10.8", diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index eee7b03..9b8db7c 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 () => { @@ -369,14 +373,16 @@ 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); } 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. */ +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; +} diff --git a/src/broker-nats.ts b/src/broker-nats.ts new file mode 100644 index 0000000..f6f4ff7 --- /dev/null +++ b/src/broker-nats.ts @@ -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..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. 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 { + 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(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. + */ + 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). + * + * `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** (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)) { + 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 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..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; `.` 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((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); +} 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/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); 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/); +}); 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); 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"); +});