Compare commits
13
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dea98e509a | ||
|
|
80b02740ab | ||
|
|
621d033d53 | ||
|
|
38831c5c56 | ||
|
|
fdad2f3268 | ||
|
|
81972a4995 | ||
|
|
10e8191717 | ||
|
|
d703cebff4 | ||
|
|
46b56d53a6 | ||
|
|
d4a2802342 | ||
|
|
8acfa7a07d | ||
|
|
f0104b7846 | ||
|
|
cd26131c61 |
@@ -25,6 +25,20 @@ MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool
|
|||||||
`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables
|
`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables
|
||||||
and starts it like any other supervised workload.
|
and starts it like any other supervised workload.
|
||||||
|
|
||||||
|
## `mesh` — the tools for whoever is on a machine
|
||||||
|
|
||||||
|
The same package carries the client (novox/hq design 25 §7, design 34): `mesh tools`, `mesh call
|
||||||
|
<module>.<tool> [json]`, `mesh mcp` (an MCP server over stdio for a program a person starts) and
|
||||||
|
`mesh serve` (the **console**: MCP over HTTP on a machine's loopback, started by the mesh as the
|
||||||
|
`mesh-console` module on the credential in `MESH_BROKER_FILE` — novox/hq ADR 0152). `mesh serve`
|
||||||
|
refuses to bind anything but loopback. With `--console <url>`, `tools` and `call` go through a console
|
||||||
|
already on the machine and need no credential.
|
||||||
|
|
||||||
|
Discovery asks the modules: every runtime answers a `tools` verb for each module it serves, with names,
|
||||||
|
descriptions and schemas from the code that answers them, and the console asks the catalogue which
|
||||||
|
modules the mesh holds and each module what it serves. A module that does not answer is named, never
|
||||||
|
dropped. A module may not name a tool of its own `tools`; the runtime refuses it at load.
|
||||||
|
|
||||||
## Verified
|
## Verified
|
||||||
|
|
||||||
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the
|
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the
|
||||||
|
|||||||
Generated
+76
@@ -0,0 +1,76 @@
|
|||||||
|
{
|
||||||
|
"name": "@novox/mesh-tools",
|
||||||
|
"version": "0.1.0",
|
||||||
|
"lockfileVersion": 3,
|
||||||
|
"requires": true,
|
||||||
|
"packages": {
|
||||||
|
"": {
|
||||||
|
"name": "@novox/mesh-tools",
|
||||||
|
"version": "0.1.0",
|
||||||
|
"dependencies": {
|
||||||
|
"@novox/mesh-sdk": "^0.1.0",
|
||||||
|
"nats": "^2.29.0"
|
||||||
|
},
|
||||||
|
"bin": {
|
||||||
|
"mesh": "dist/mesh.js",
|
||||||
|
"mesh-tools": "dist/main.js"
|
||||||
|
},
|
||||||
|
"devDependencies": {
|
||||||
|
"@types/node": "^22.20.1",
|
||||||
|
"typescript": "^5.9.3"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/@novox/mesh-sdk": {
|
||||||
|
"version": "0.1.1"
|
||||||
|
},
|
||||||
|
"node_modules/@types/node": {
|
||||||
|
"version": "22.20.1",
|
||||||
|
"dev": true,
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"undici-types": "~6.21.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/nats": {
|
||||||
|
"version": "2.29.3",
|
||||||
|
"license": "Apache-2.0",
|
||||||
|
"dependencies": {
|
||||||
|
"nkeys.js": "1.1.0"
|
||||||
|
},
|
||||||
|
"engines": {
|
||||||
|
"node": ">= 14.0.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/nkeys.js": {
|
||||||
|
"version": "1.1.0",
|
||||||
|
"license": "Apache-2.0",
|
||||||
|
"dependencies": {
|
||||||
|
"tweetnacl": "1.0.3"
|
||||||
|
},
|
||||||
|
"engines": {
|
||||||
|
"node": ">=10.0.0"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/tweetnacl": {
|
||||||
|
"version": "1.0.3",
|
||||||
|
"license": "Unlicense"
|
||||||
|
},
|
||||||
|
"node_modules/typescript": {
|
||||||
|
"version": "5.9.3",
|
||||||
|
"dev": true,
|
||||||
|
"license": "Apache-2.0",
|
||||||
|
"bin": {
|
||||||
|
"tsc": "bin/tsc",
|
||||||
|
"tsserver": "bin/tsserver"
|
||||||
|
},
|
||||||
|
"engines": {
|
||||||
|
"node": ">=14.17"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"node_modules/undici-types": {
|
||||||
|
"version": "6.21.0",
|
||||||
|
"dev": true,
|
||||||
|
"license": "MIT"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -14,11 +14,9 @@
|
|||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@novox/mesh-sdk": "^0.1.0",
|
"@novox/mesh-sdk": "^0.1.0",
|
||||||
"amqplib": "^0.10.9",
|
|
||||||
"nats": "^2.29.0"
|
"nats": "^2.29.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/amqplib": "^0.10.8",
|
|
||||||
"@types/node": "^22.20.1",
|
"@types/node": "^22.20.1",
|
||||||
"typescript": "^5.9.3"
|
"typescript": "^5.9.3"
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,422 +0,0 @@
|
|||||||
// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker
|
|
||||||
// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime
|
|
||||||
// provides the binding — so a broker-client change never rebuilds the modules. This is where the
|
|
||||||
// ADR 0042 wire shape lives: the two exchanges, persistent events, per-consumer durable queues,
|
|
||||||
// prefetch, dead-letter — none of which a module ever sees.
|
|
||||||
|
|
||||||
import amqp from "amqplib";
|
|
||||||
import * as tls from "node:tls";
|
|
||||||
import { randomUUID, createHash } from "node:crypto";
|
|
||||||
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
|
||||||
|
|
||||||
// Two topic exchanges, kept apart on purpose (ADR 0042): tool invocations are request/reply and are
|
|
||||||
// not events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node
|
|
||||||
// events — never the RPC traffic.
|
|
||||||
const RPC_EXCHANGE = "mesh.rpc";
|
|
||||||
const EVENTS_EXCHANGE = "mesh.events";
|
|
||||||
// Where an event rejected past its redelivery limit is set aside for inspection.
|
|
||||||
const DEAD_EXCHANGE = "mesh.events.dead";
|
|
||||||
// Bound in-flight events so one slow consumer can't pull the whole backlog into memory (ADR 0042).
|
|
||||||
const EVENT_PREFETCH = 32;
|
|
||||||
|
|
||||||
interface Reply {
|
|
||||||
result?: unknown;
|
|
||||||
error?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The broker presented a certificate whose fingerprint is not the one the mesh pinned. A distinct
|
|
||||||
* type rather than a message to grep, so a caller deciding "wait or refuse" (serve mode's patient
|
|
||||||
* reconnect, novox/hq issue 058) tells this apart from an absent broker by `instanceof`, not by a
|
|
||||||
* prose string that a later reword would silently turn back into an infinite retry against an
|
|
||||||
* impostor. */
|
|
||||||
export class PinMismatchError extends Error {}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Why a broker connection failed in a way no amount of waiting will fix — or null when it is worth
|
|
||||||
* retrying. Serve mode's patient reconnect (novox/hq issue 058) uses this to tell a permanent
|
|
||||||
* fault from a broker that is merely not up yet. Three failures are permanent:
|
|
||||||
*
|
|
||||||
* - the certificate does not match the pin — an impostor does not become the broker by being
|
|
||||||
* asked again (typed, so a reworded message cannot silently turn this back into a retry);
|
|
||||||
* - the broker URL is not a URL — a malformed address never parses on the next try;
|
|
||||||
* - the broker answered and refused the login — a wrong or revoked credential, not an absent
|
|
||||||
* broker, and it will refuse the next attempt identically.
|
|
||||||
*
|
|
||||||
* Everything else — connection refused, timeout, DNS not resolving yet — is the overlay still
|
|
||||||
* coming up, and is retried.
|
|
||||||
*/
|
|
||||||
export function fatalBrokerReason(err: unknown): string | null {
|
|
||||||
if (err instanceof PinMismatchError) return "the broker's certificate does not match the pin";
|
|
||||||
const e = err as { code?: unknown; message?: unknown };
|
|
||||||
const code = typeof e?.code === "string" ? e.code : "";
|
|
||||||
const message = typeof e?.message === "string" ? e.message : String(err);
|
|
||||||
if (code === "ERR_INVALID_URL" || /invalid url/i.test(message)) {
|
|
||||||
return `the broker URL is not a URL (${message})`;
|
|
||||||
}
|
|
||||||
if (/access[-_ ]?refused|login was refused|handshake terminated|\b403\b/i.test(message)) {
|
|
||||||
return `the broker refused the login (${message})`;
|
|
||||||
}
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A broker credential as the mesh delivers it (novox/hq ADR 0043): an amqps URL, the fingerprint
|
|
||||||
* of the certificate the broker must present, and the node and module the account is scoped to (so
|
|
||||||
* the runtime names its queue as the mesh did). A plain string is a bootstrap URL. */
|
|
||||||
export interface Credential {
|
|
||||||
url: string;
|
|
||||||
fingerprint?: string;
|
|
||||||
node?: string;
|
|
||||||
module?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down.
|
|
||||||
*
|
|
||||||
* A scoped module (novox/hq ADR 0043) passes `assumeExchanges: true`: its account may not declare
|
|
||||||
* an exchange, and the foundation already owns them, so it declares only its own queue. A credential
|
|
||||||
* carrying a fingerprint is dialled over amqps, pinned to exactly that certificate.
|
|
||||||
*/
|
|
||||||
export async function connectAmqp(
|
|
||||||
target: string | Credential,
|
|
||||||
opts: { assumeExchanges?: boolean } = {},
|
|
||||||
): Promise<Broker> {
|
|
||||||
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);
|
|
||||||
|
|
||||||
// A confirm channel, so an event publish awaits the broker's ack: a publish the broker never
|
|
||||||
// accepted (it was mid-restart, the connection dropped) fails the emit rather than vanishing —
|
|
||||||
// at-least-once starts at the emitter, not only the consumer (ADR 0042).
|
|
||||||
const ch = await conn.createConfirmChannel();
|
|
||||||
|
|
||||||
// The foundation owns the exchanges (ADR 0043). A bootstrap/admin connection declares them; a
|
|
||||||
// scoped module assumes they exist and never tries — its account could not, and the dead-letter
|
|
||||||
// queue behind the exchange is the foundation's to keep, not a module's.
|
|
||||||
if (!opts.assumeExchanges) {
|
|
||||||
await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true });
|
|
||||||
await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true });
|
|
||||||
await ch.assertExchange(DEAD_EXCHANGE, "topic", { durable: true });
|
|
||||||
await ch.assertQueue(DEAD_EXCHANGE, { durable: true });
|
|
||||||
await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#");
|
|
||||||
}
|
|
||||||
|
|
||||||
await ch.prefetch(EVENT_PREFETCH);
|
|
||||||
|
|
||||||
// Request/reply is set up lazily: a consumer-only module (the audit logger) never calls a tool,
|
|
||||||
// and its scoped account may not declare the exclusive reply queue this would otherwise need.
|
|
||||||
const pending = new Map<string, (r: Reply) => void>();
|
|
||||||
let replyQueue: string | undefined;
|
|
||||||
async function ensureReply(): Promise<string> {
|
|
||||||
if (replyQueue) return replyQueue;
|
|
||||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
|
||||||
// Replies come back through the RPC exchange keyed by this queue's own name, not the default
|
|
||||||
// exchange (novox/hq ADR 0047): a serving module's scoped account may write to mesh.rpc but not
|
|
||||||
// the default exchange, which would let it publish into any queue on the broker.
|
|
||||||
await ch.bindQueue(queue, RPC_EXCHANGE, queue);
|
|
||||||
replyQueue = queue;
|
|
||||||
await ch.consume(
|
|
||||||
queue,
|
|
||||||
(msg) => {
|
|
||||||
if (!msg) return;
|
|
||||||
const resolve = pending.get(msg.properties.correlationId);
|
|
||||||
if (resolve) {
|
|
||||||
pending.delete(msg.properties.correlationId);
|
|
||||||
resolve(JSON.parse(msg.content.toString()) as Reply);
|
|
||||||
}
|
|
||||||
},
|
|
||||||
{ noAck: true },
|
|
||||||
);
|
|
||||||
return queue;
|
|
||||||
}
|
|
||||||
|
|
||||||
// One durable event queue per consumer (ADR 0042: <node>.<module>.events), with many bindings and
|
|
||||||
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
|
|
||||||
// once however many bindings match, so the local match is what keeps a two-pattern module from
|
|
||||||
// running the wrong handler.
|
|
||||||
type EventSub = { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> };
|
|
||||||
const eventSubs: EventSub[] = [];
|
|
||||||
let eventQueue: string | undefined;
|
|
||||||
let eventConsumerTag: string | undefined;
|
|
||||||
|
|
||||||
async function dispatchEvent(msg: amqp.ConsumeMessage): Promise<void> {
|
|
||||||
const key = localKeyFor(msg.fields.routingKey);
|
|
||||||
let env: Envelope<unknown>;
|
|
||||||
try {
|
|
||||||
env = toEnvelope(msg);
|
|
||||||
} catch {
|
|
||||||
// An undecodable body will never decode on redelivery — dead-letter it at once rather than
|
|
||||||
// wedge the queue or loop. Decoding sits before the handler try on purpose: a poison message
|
|
||||||
// is a different failure from a handler that threw, and gets no retry.
|
|
||||||
ch.nack(msg, false, false);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
try {
|
|
||||||
for (const s of eventSubs) {
|
|
||||||
if (topicMatches(s.pattern, key)) await s.handler(env);
|
|
||||||
}
|
|
||||||
ch.ack(msg);
|
|
||||||
} catch {
|
|
||||||
// First handler failure: requeue once. A second (already redelivered) dead-letters it, so a
|
|
||||||
// poison event is set aside rather than looping forever or vanishing (ADR 0042). Redelivery
|
|
||||||
// re-runs every matching handler, so a consumer must be idempotent — which the ADR requires.
|
|
||||||
ch.nack(msg, false, !msg.fields.redelivered);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
|
||||||
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
|
||||||
const reply = await ensureReply();
|
|
||||||
const id = randomUUID();
|
|
||||||
const answered = new Promise<Res>((resolve, reject) => {
|
|
||||||
const timer = setTimeout(() => {
|
|
||||||
if (pending.delete(id)) reject(new Error(`request ${key} timed out`));
|
|
||||||
}, 30_000);
|
|
||||||
pending.set(id, (r) => {
|
|
||||||
clearTimeout(timer);
|
|
||||||
if (r.error) reject(new Error(r.error));
|
|
||||||
else resolve(r.result as Res);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
|
||||||
correlationId: id,
|
|
||||||
replyTo: reply,
|
|
||||||
});
|
|
||||||
return answered;
|
|
||||||
},
|
|
||||||
|
|
||||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
|
||||||
// A durable, shared serve queue (ADR 0042): several runtimes serving one tool key compete for
|
|
||||||
// invocations rather than each answering the same call.
|
|
||||||
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
|
|
||||||
await ch.bindQueue(queue, RPC_EXCHANGE, key);
|
|
||||||
const consumer = await ch.consume(queue, (msg) => {
|
|
||||||
if (!msg) return;
|
|
||||||
void (async () => {
|
|
||||||
let reply: Reply;
|
|
||||||
try {
|
|
||||||
reply = { result: await handler(JSON.parse(msg.content.toString()) as Req) };
|
|
||||||
} catch (err) {
|
|
||||||
reply = { error: err instanceof Error ? err.message : String(err) };
|
|
||||||
}
|
|
||||||
if (msg.properties.replyTo) {
|
|
||||||
// Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped
|
|
||||||
// account answers with write on mesh.rpc alone — never the default exchange (ADR 0047).
|
|
||||||
ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), {
|
|
||||||
correlationId: msg.properties.correlationId,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
ch.ack(msg);
|
|
||||||
})();
|
|
||||||
});
|
|
||||||
return () => void ch.cancel(consumer.consumerTag);
|
|
||||||
},
|
|
||||||
|
|
||||||
async publish<T>(env: Envelope<T>): Promise<void> {
|
|
||||||
const headers = env.headers ?? ({} as EventHeaders);
|
|
||||||
// Events are persistent (delivery-mode 2): an audit trail that loses events on a broker
|
|
||||||
// restart is not one (ADR 0042). Metadata rides as headers; the body is only the payload.
|
|
||||||
// The publish is awaited to the broker's confirm — an unaccepted publish rejects here.
|
|
||||||
await new Promise<void>((resolve, reject) => {
|
|
||||||
ch.publish(
|
|
||||||
EVENTS_EXCHANGE,
|
|
||||||
routingKeyFor(env.key, self),
|
|
||||||
Buffer.from(JSON.stringify(env.body)),
|
|
||||||
{
|
|
||||||
persistent: true,
|
|
||||||
contentType:
|
|
||||||
typeof headers["content-type"] === "string" ? headers["content-type"] : "application/json",
|
|
||||||
messageId: headers["x-event-id"],
|
|
||||||
headers: { ...headers },
|
|
||||||
},
|
|
||||||
(err) => (err ? reject(err instanceof Error ? err : new Error(String(err))) : resolve()),
|
|
||||||
);
|
|
||||||
});
|
|
||||||
},
|
|
||||||
|
|
||||||
async subscribe<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
|
||||||
const node = process.env.MESH_NODE;
|
|
||||||
const mod = process.env.MESH_MODULE;
|
|
||||||
|
|
||||||
// A module we can name gets its ADR 0042 durable queue; an anonymous subscriber (a test, an
|
|
||||||
// ad-hoc listener) gets a transient exclusive one that dies with the connection.
|
|
||||||
if (node && mod) {
|
|
||||||
if (!eventQueue) {
|
|
||||||
const name = `${node}.${mod}.events`;
|
|
||||||
if (opts.assumeExchanges) {
|
|
||||||
// The mesh pre-declared this queue with its dead-letter when it issued the account: a
|
|
||||||
// scoped account may not declare a dead-lettered queue itself (the broker refuses that
|
|
||||||
// to a non-administrator). Passively check it is there, then bind and consume.
|
|
||||||
await ch.checkQueue(name);
|
|
||||||
} else {
|
|
||||||
await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE });
|
|
||||||
}
|
|
||||||
eventQueue = name;
|
|
||||||
}
|
|
||||||
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, bindingFor(pattern));
|
|
||||||
const sub: EventSub = { pattern, handler: handler as EventSub["handler"] };
|
|
||||||
eventSubs.push(sub);
|
|
||||||
if (!eventConsumerTag) {
|
|
||||||
const consumer = await ch.consume(eventQueue, (msg) => {
|
|
||||||
if (msg) void dispatchEvent(msg);
|
|
||||||
});
|
|
||||||
eventConsumerTag = consumer.consumerTag;
|
|
||||||
}
|
|
||||||
return () => {
|
|
||||||
const i = eventSubs.indexOf(sub);
|
|
||||||
if (i >= 0) eventSubs.splice(i, 1);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
|
||||||
await ch.bindQueue(queue, EVENTS_EXCHANGE, bindingFor(pattern));
|
|
||||||
const consumer = await ch.consume(queue, (msg) => {
|
|
||||||
if (!msg) return;
|
|
||||||
void (async () => {
|
|
||||||
let env: Envelope<T>;
|
|
||||||
try {
|
|
||||||
env = toEnvelope<T>(msg);
|
|
||||||
} catch {
|
|
||||||
ch.nack(msg, false, false); // undecodable — drop, never retry
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
try {
|
|
||||||
await handler(env);
|
|
||||||
ch.ack(msg);
|
|
||||||
} catch {
|
|
||||||
ch.nack(msg, false, !msg.fields.redelivered);
|
|
||||||
}
|
|
||||||
})();
|
|
||||||
});
|
|
||||||
return () => void ch.cancel(consumer.consumerTag);
|
|
||||||
},
|
|
||||||
|
|
||||||
async close(): Promise<void> {
|
|
||||||
await ch.close();
|
|
||||||
await conn.close();
|
|
||||||
},
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Normalise a certificate fingerprint to bare lower-case hex, dropping an `sha256:` prefix and
|
|
||||||
* any colon grouping, so two spellings of the same fingerprint compare equal. */
|
|
||||||
function normalizeFingerprint(fingerprint: string): string {
|
|
||||||
return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase();
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Socket options that pin the broker to exactly the certificate whose fingerprint the mesh
|
|
||||||
* delivered (novox/hq ADR 0043, as the builder does). Done in two phases so a credential never
|
|
||||||
* reaches an impostor: first a bare TLS connection that sends nothing fetches the certificate and
|
|
||||||
* the fingerprint is checked; only then does the real connection trust *that* certificate as its
|
|
||||||
* own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key.
|
|
||||||
* Node's `checkServerIdentity` does not run under `rejectUnauthorized: false`, so a one-phase
|
|
||||||
* "connect then check" would have already sent the password to whoever answered.
|
|
||||||
*/
|
|
||||||
async function pinnedOptions(rawUrl: string, fingerprint: string): Promise<tls.ConnectionOptions> {
|
|
||||||
const url = new URL(rawUrl);
|
|
||||||
const host = url.hostname;
|
|
||||||
const port = url.port ? Number(url.port) : 5671;
|
|
||||||
|
|
||||||
const certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
|
|
||||||
const socket = tls.connect({ host, port, servername: host, rejectUnauthorized: false }, () => {
|
|
||||||
const peer = socket.getPeerCertificate(true);
|
|
||||||
socket.destroy();
|
|
||||||
if (!peer || !peer.raw) reject(new Error("the broker presented no certificate to pin"));
|
|
||||||
else resolve(peer);
|
|
||||||
});
|
|
||||||
socket.setTimeout(15_000, () => {
|
|
||||||
socket.destroy();
|
|
||||||
reject(new Error("timed out fetching the broker's certificate"));
|
|
||||||
});
|
|
||||||
socket.on("error", reject);
|
|
||||||
});
|
|
||||||
|
|
||||||
const seen = createHash("sha256").update(certificate.raw).digest("hex");
|
|
||||||
if (seen !== normalizeFingerprint(fingerprint)) {
|
|
||||||
throw new PinMismatchError(
|
|
||||||
`the broker's certificate (sha256:${seen}) does not match the pinned ${fingerprint} — refusing`,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
const pem =
|
|
||||||
"-----BEGIN CERTIFICATE-----\n" +
|
|
||||||
(certificate.raw.toString("base64").match(/.{1,64}/g) ?? []).join("\n") +
|
|
||||||
"\n-----END CERTIFICATE-----\n";
|
|
||||||
// Trust that one certificate and nothing else; the mesh's own name is not in any public store,
|
|
||||||
// so identity is the pin, not the hostname — checkServerIdentity is satisfied deliberately.
|
|
||||||
return { ca: [pem], checkServerIdentity: () => undefined };
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */
|
|
||||||
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
|
||||||
const raw = msg.properties.headers ?? {};
|
|
||||||
const headers: Record<string, string> = {};
|
|
||||||
for (const [k, v] of Object.entries(raw)) {
|
|
||||||
if (v == null) continue;
|
|
||||||
headers[k] = typeof v === "string" ? v : String(v);
|
|
||||||
}
|
|
||||||
if (!headers["content-type"] && msg.properties.contentType) {
|
|
||||||
headers["content-type"] = msg.properties.contentType;
|
|
||||||
}
|
|
||||||
return {
|
|
||||||
key: msg.fields.routingKey,
|
|
||||||
node: headers["x-node"] ?? "",
|
|
||||||
body: JSON.parse(msg.content.toString()) as T,
|
|
||||||
headers: headers as EventHeaders,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/** 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. */
|
|
||||||
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];
|
|
||||||
// `**` 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;
|
|
||||||
}
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
if (ki >= k.length) return false;
|
|
||||||
if (tok !== "*" && tok !== k[ki]) return false;
|
|
||||||
pi++;
|
|
||||||
ki++;
|
|
||||||
}
|
|
||||||
return ki === k.length;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** This bus spells an event as a routing key that repeats the emitter's name: `module.<emitter>.<event>`.
|
|
||||||
* A module names its events locally and the mesh derives where they land (novox/hq design 29 §1), so
|
|
||||||
* the mapping lives here rather than in every module.
|
|
||||||
*
|
|
||||||
* **Why it exists at all.** Until 04-ISSUES/127 every module passed the routing key itself, which
|
|
||||||
* worked on this bus and derived into a namespace nobody owns on the one being built. Converting the
|
|
||||||
* modules to local names without this would have broken the mesh that is actually running. */
|
|
||||||
export function routingKeyFor(key: string, self: string): string {
|
|
||||||
return key.startsWith("module.") ? key : `module.${self}.${key}`;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A local pattern as this bus's binding. `**` is the mesh's wildcard for the rest of a name; here
|
|
||||||
* that is `#`, and on the bus being built it is `>`. Neither spelling appears in a manifest. */
|
|
||||||
export function bindingFor(pattern: string): string {
|
|
||||||
const here = pattern.split(".").map((part) => (part === "**" ? "#" : part)).join(".");
|
|
||||||
if (here === "#") return "#";
|
|
||||||
return here.startsWith("module.") ? here : `module.${here}`;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A routing key as the local name a handler and a manifest both use: the emitter and the event. */
|
|
||||||
export function localKeyFor(routingKey: string): string {
|
|
||||||
return routingKey.startsWith("module.") ? routingKey.slice("module.".length) : routingKey;
|
|
||||||
}
|
|
||||||
+88
-37
@@ -17,8 +17,9 @@
|
|||||||
// correct.
|
// correct.
|
||||||
|
|
||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
|
import net from "node:net";
|
||||||
import tls from "node:tls";
|
import tls from "node:tls";
|
||||||
import { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription } from "nats";
|
import { connect as natsConnect, headers as natsHeaders, StringCodec, type JsMsg, type Subscription, type TlsOptions } from "nats";
|
||||||
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
const sc = StringCodec();
|
const sc = StringCodec();
|
||||||
@@ -84,6 +85,10 @@ export async function connectNats(
|
|||||||
pass: cred.password,
|
pass: cred.password,
|
||||||
name: `${cred.node ?? "?"}.${self}`,
|
name: `${cred.node ?? "?"}.${self}`,
|
||||||
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
|
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
|
||||||
|
// Its own inbox, not a random one: every user's inbox is private to it (design 25 §4), and the
|
||||||
|
// grant names `_INBOX.<user>.>` — a reply space the client invented would be refused, and with
|
||||||
|
// it every pull for the next message and every answer to a tool call.
|
||||||
|
inboxPrefix: cred.user ? `_INBOX.${cred.user}` : undefined,
|
||||||
// Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on
|
// Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on
|
||||||
// the mesh to exit. `close()` stays the only thing that ends the connection.
|
// the mesh to exit. `close()` stays the only thing that ends the connection.
|
||||||
maxReconnectAttempts: -1,
|
maxReconnectAttempts: -1,
|
||||||
@@ -91,6 +96,10 @@ export async function connectNats(
|
|||||||
const js = conn.jetstream();
|
const js = conn.jetstream();
|
||||||
|
|
||||||
const subs: Subscription[] = [];
|
const subs: Subscription[] = [];
|
||||||
|
// Every registration, and the one reader that dispatches to them. A module has one durable
|
||||||
|
// consumer; the loop belongs to the connection rather than to a subscription.
|
||||||
|
const listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[] = [];
|
||||||
|
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
|
||||||
let closed = false;
|
let closed = false;
|
||||||
|
|
||||||
return {
|
return {
|
||||||
@@ -175,21 +184,38 @@ export async function connectNats(
|
|||||||
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
|
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
|
||||||
* its own would be a module deciding its own delivery semantics, and its account cannot
|
* its own would be a module deciding its own delivery semantics, and its account cannot
|
||||||
* reach the JetStream API to do it anyway.
|
* reach the JetStream API to do it anyway.
|
||||||
|
*
|
||||||
|
* **One consumer, one loop, however many patterns a module registers.** A module has exactly one
|
||||||
|
* durable consumer, so two loops reading it would each take half the messages — and a loop that
|
||||||
|
* received one its own pattern does not match acknowledges it, which is the right answer for a
|
||||||
|
* filter wider than anything registered and silent loss when it is another handler's. Every
|
||||||
|
* registration is therefore dispatched from one reader, and a message is acknowledged once every
|
||||||
|
* handler it is for has taken it.
|
||||||
*/
|
*/
|
||||||
async subscribe<T>(
|
async subscribe<T>(
|
||||||
pattern: string,
|
pattern: string,
|
||||||
handler: (env: Envelope<T>) => Promise<void>,
|
handler: (env: Envelope<T>) => Promise<void>,
|
||||||
): Promise<() => void> {
|
): Promise<() => void> {
|
||||||
const durable = `${cred.node ?? "?"}_${self}`;
|
const listener = { pattern, handler: handler as (env: Envelope<unknown>) => Promise<void> };
|
||||||
const consumer = await js.consumers.get("EVENTS", durable);
|
listeners.push(listener);
|
||||||
const messages = await consumer.consume();
|
if (!reading) {
|
||||||
void (async () => {
|
const durable = `${cred.node ?? "?"}_${self}`;
|
||||||
for await (const msg of messages) {
|
const consumer = await js.consumers.get("EVENTS", durable);
|
||||||
await deliver(msg, pattern, handler);
|
const messages = await consumer.consume();
|
||||||
}
|
reading = messages;
|
||||||
})();
|
void (async () => {
|
||||||
|
for await (const msg of messages) {
|
||||||
|
await deliver(msg, listeners);
|
||||||
|
}
|
||||||
|
})();
|
||||||
|
}
|
||||||
return () => {
|
return () => {
|
||||||
void messages.close();
|
const at = listeners.indexOf(listener);
|
||||||
|
if (at >= 0) listeners.splice(at, 1);
|
||||||
|
if (listeners.length === 0 && reading) {
|
||||||
|
void reading.close();
|
||||||
|
reading = undefined;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
},
|
},
|
||||||
|
|
||||||
@@ -204,15 +230,20 @@ export async function connectNats(
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Deliver one event, acknowledging only once a handler has taken it. */
|
/**
|
||||||
async function deliver<T>(
|
* Deliver one event to every handler it is for, acknowledging only once each has taken it.
|
||||||
|
*
|
||||||
|
* Several registrations share one durable consumer, so matching happens here rather than by having
|
||||||
|
* each registration read the stream: two readers of one consumer would split it between them, and a
|
||||||
|
* message that reached the wrong one would be acknowledged as not-for-me and lost.
|
||||||
|
*/
|
||||||
|
async function deliver(
|
||||||
msg: JsMsg,
|
msg: JsMsg,
|
||||||
pattern: string,
|
listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[],
|
||||||
handler: (env: Envelope<T>) => Promise<void>,
|
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
let env: Envelope<T>;
|
let env: Envelope<unknown>;
|
||||||
try {
|
try {
|
||||||
env = toEnvelope<T>(msg);
|
env = toEnvelope<unknown>(msg);
|
||||||
} catch {
|
} catch {
|
||||||
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
|
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
|
||||||
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
|
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
|
||||||
@@ -220,15 +251,16 @@ async function deliver<T>(
|
|||||||
msg.term();
|
msg.term();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (!topicMatches(pattern, env.key)) {
|
const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key));
|
||||||
// The consumer's filters are the controller's, and may be wider than one subscription's
|
if (forThis.length === 0) {
|
||||||
// pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
|
// The consumer's filters are the controller's, derived from what the module declared it
|
||||||
|
// consumes, and may be wider than anything it registered a handler for. Acknowledge it, or it
|
||||||
// would be redelivered until it expired.
|
// would be redelivered until it expired.
|
||||||
msg.ack();
|
msg.ack();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
await handler(env);
|
for (const l of forThis) await l.handler(env);
|
||||||
msg.ack();
|
msg.ack();
|
||||||
} catch {
|
} catch {
|
||||||
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
||||||
@@ -297,26 +329,43 @@ function normalizeFingerprint(fingerprint: string): string {
|
|||||||
* A certificate authority is not consulted: the mesh issued this and knows its fingerprint,
|
* A certificate authority is not consulted: the mesh issued this and knows its fingerprint,
|
||||||
* which is stronger than trusting whoever a machine's trust store happens to contain.
|
* which is stronger than trusting whoever a machine's trust store happens to contain.
|
||||||
*
|
*
|
||||||
* **A constraint on the mesh, not a detail of this file.** Pinning the exact certificate makes
|
* **The pin is the only check.** What comes back is handed to the client as its TLS options, and
|
||||||
* hostname verification redundant in principle, but the NATS client exposes no hook to replace
|
* the client's transport spreads them into Node's own `tls.connect` — so the pinned certificate
|
||||||
* it — its TLS options are file paths and PEM strings, with no verify callback. So the
|
* is the one authority the handshake accepts, and the hostname check beside it is replaced with
|
||||||
* certificate the mesh issues the bus **must carry a subject-alternative name matching the
|
* one that always passes. Pinning the exact certificate makes verifying its name redundant, and
|
||||||
* address nodes dial it by**. The fingerprint check below still happens and is still the real
|
* the bus's certificate names the seat (`mesh-broker`), not the address a machine happens to
|
||||||
* guarantee; what cannot be switched off is the check *beside* it.
|
* dial it by: every module on the mesh met "does not match certificate's altnames" the first time
|
||||||
|
* it reached the handshake (2026-09-28).
|
||||||
*/
|
*/
|
||||||
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: string }> {
|
async function pinnedTls(rawUrl: string, fingerprint: string): Promise<TlsOptions> {
|
||||||
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
const url = new URL(rawUrl.includes("://") ? rawUrl : `nats://${rawUrl}`);
|
||||||
const port = url.port ? Number(url.port) : 4222;
|
const port = url.port ? Number(url.port) : 4222;
|
||||||
|
// **The bus speaks first, in the clear.** A NATS server sends its INFO line before TLS begins,
|
||||||
|
// and only then expects the client to start the handshake; a raw TLS connect to that port reads
|
||||||
|
// the INFO line as a TLS record and fails with "wrong version number" — which is what every
|
||||||
|
// module met the first time it dialled the bus being built (2026-09-28). So: connect, wait for
|
||||||
|
// INFO, then start TLS on the same socket, and read the certificate the server presents.
|
||||||
const certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
|
const certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
|
||||||
const socket = tls.connect(
|
const plain = net.connect({ host: url.hostname, port }, () => {});
|
||||||
{ host: url.hostname, port, rejectUnauthorized: false, servername: url.hostname },
|
let seenInfo = false;
|
||||||
() => {
|
let buffered = "";
|
||||||
const peer = socket.getPeerCertificate(true);
|
plain.on("error", reject);
|
||||||
socket.end();
|
plain.on("data", (chunk: Buffer) => {
|
||||||
resolve(peer);
|
if (seenInfo) return;
|
||||||
},
|
buffered += chunk.toString("utf8");
|
||||||
);
|
if (!buffered.includes("\r\n")) return;
|
||||||
socket.on("error", reject);
|
seenInfo = true;
|
||||||
|
plain.removeAllListeners("data");
|
||||||
|
const secure = tls.connect(
|
||||||
|
{ socket: plain, rejectUnauthorized: false, servername: url.hostname },
|
||||||
|
() => {
|
||||||
|
const peer = secure.getPeerCertificate(true);
|
||||||
|
secure.end();
|
||||||
|
resolve(peer);
|
||||||
|
},
|
||||||
|
);
|
||||||
|
secure.on("error", reject);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
const seen = createHash("sha256").update(certificate.raw).digest("hex");
|
const seen = createHash("sha256").update(certificate.raw).digest("hex");
|
||||||
if (seen !== normalizeFingerprint(fingerprint)) {
|
if (seen !== normalizeFingerprint(fingerprint)) {
|
||||||
@@ -325,7 +374,9 @@ async function pinnedTls(rawUrl: string, fingerprint: string): Promise<{ ca: str
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
|
const pem = `-----BEGIN CERTIFICATE-----\n${certificate.raw.toString("base64").replace(/(.{64})/g, "$1\n")}\n-----END CERTIFICATE-----\n`;
|
||||||
return { ca: pem };
|
// Node's option, not the client's: the transport passes the whole object on. `undefined` from
|
||||||
|
// checkServerIdentity is "the name is fine"; the pin above already decided the rest.
|
||||||
|
return { ca: pem, checkServerIdentity: () => undefined } as TlsOptions;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary —
|
/** The mesh's topic matching: `*` is one token, `#` the rest. This is the module's vocabulary —
|
||||||
|
|||||||
+86
-31
@@ -1,30 +1,32 @@
|
|||||||
/**
|
/**
|
||||||
* A person's client: the mesh's tools from a workstation (novox/hq design 25 §7).
|
* The mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
|
||||||
*
|
*
|
||||||
* Two surfaces over one thing. A command line, for somebody at a terminal; an MCP server, for an
|
* 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,
|
* 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.
|
* 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
|
* **It uses the same client a module's runtime uses.** Not a second protocol and not a bridge: the
|
||||||
* person connects as their own bus user, publishes on the tool subjects their account permits, and the
|
* caller connects as its own bus user — a person's, or the console's — publishes on the tool subjects
|
||||||
* server refuses anything else. So "what may this person do" is answered by the same permission list
|
* that account permits, and the server refuses anything else. So "what may this ask" is answered by the
|
||||||
* that answers it for a module, and there is nothing here for an audit to read separately.
|
* 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
|
* What the caller 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;
|
* account (design 25 §4): it cannot publish an event, so it cannot claim a module said something; it
|
||||||
* they have no consumer, so there is no delivery to acknowledge; and they cannot answer a request, so
|
* has no consumer, so there is no delivery to acknowledge; and it cannot answer a request, so it cannot
|
||||||
* they cannot impersonate a module on a bus where anyone may serve a tool.
|
* impersonate a module on a bus where anyone may serve a tool.
|
||||||
*/
|
*/
|
||||||
import { readFile } from "node:fs/promises";
|
import { readFile } from "node:fs/promises";
|
||||||
|
|
||||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
import { connectNats, type Credential } from "./broker-nats.js";
|
import { connectNats, type Credential } from "./broker-nats.js";
|
||||||
|
import { TOOLS_VERB, type ToolsAnswer } from "./runtime.js";
|
||||||
|
|
||||||
/** Where the catalogue answers what tools the mesh has. */
|
/** Where the catalogue answers which modules the mesh holds. */
|
||||||
const CATALOGUE_TOOLS = "mesh-catalog.catalog_tools";
|
const CATALOGUE_MODULES = "mesh-catalog.catalog_modules";
|
||||||
|
|
||||||
/** A tool as the catalogue describes one. */
|
/** A tool as its module describes it. */
|
||||||
export interface Tool {
|
export interface Tool {
|
||||||
module: string;
|
module: string;
|
||||||
name: string;
|
name: string;
|
||||||
@@ -33,6 +35,20 @@ export interface Tool {
|
|||||||
input?: unknown;
|
input?: unknown;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* What the mesh could say about its tools when asked (design 34 §3).
|
||||||
|
*
|
||||||
|
* **Silence is named, never dropped.** A module the catalogue holds and nothing answered for is in
|
||||||
|
* `notAnswering`, because a tool that is not offered looks exactly like a tool that does not exist,
|
||||||
|
* and those need different people to fix them.
|
||||||
|
*/
|
||||||
|
export interface Listing {
|
||||||
|
tools: Tool[];
|
||||||
|
/** Modules the catalogue holds whose runtime did not answer `tools`: not assigned, not up, or built
|
||||||
|
* before the runtime answered it. Each may still be called by name. */
|
||||||
|
notAnswering: string[];
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* A person's credential, as `operator issue` prints it.
|
* A person's credential, as `operator issue` prints it.
|
||||||
*
|
*
|
||||||
@@ -72,24 +88,63 @@ export async function connectAs(held: PersonCredential): Promise<Broker> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* What tools the mesh has, asked of the catalogue.
|
* Connect as the console: the module credential the mesh delivered (novox/hq ADR 0152), read from
|
||||||
*
|
* the same variable every runtime reads. It names the node and the module, so the account's inbox
|
||||||
* **Asked, not configured.** The catalogue is the only thing that knows what is installed, and a
|
* and subjects derive from what the mesh authorised and from nothing in this process's environment.
|
||||||
* client carrying its own list would be a list that goes stale the first time a module is assigned —
|
|
||||||
* silently, because a tool that is not offered looks exactly like a tool that does not exist.
|
|
||||||
*/
|
*/
|
||||||
export async function toolsOn(bus: Broker): Promise<Tool[]> {
|
export async function connectAsTheConsole(path: string): Promise<{ bus: Broker; who: string }> {
|
||||||
const answered = await bus.request<Record<string, never>, { tools?: Tool[] } | Tool[]>(
|
const raw = await readFile(path, "utf8");
|
||||||
CATALOGUE_TOOLS,
|
let held: Credential;
|
||||||
{},
|
try {
|
||||||
);
|
held = JSON.parse(raw) as Credential;
|
||||||
const tools = Array.isArray(answered) ? answered : (answered.tools ?? []);
|
} catch (e) {
|
||||||
return tools
|
throw new Error(`${path} is not a broker credential: ${(e as Error).message}`);
|
||||||
.slice()
|
}
|
||||||
.sort((a: Tool, b: Tool) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
|
if (!held.url || !held.module || !held.user) {
|
||||||
|
throw new Error(
|
||||||
|
`${path} names no bus, module or user: the console runs on the credential the mesh sealed to ` +
|
||||||
|
"this machine for it, and nothing else",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return { bus: await connectNats(held), who: `${held.node ?? "?"}.${held.module}` };
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what their account
|
/**
|
||||||
|
* What tools the mesh has, asked of the modules (design 34 §3).
|
||||||
|
*
|
||||||
|
* The catalogue says which modules the mesh holds; each module says what it serves, through the one
|
||||||
|
* verb its runtime answers for it. **Asked, not configured**: a client carrying its own list would be a
|
||||||
|
* list that goes stale the first time a module is assigned. Every module is asked at once, and the bus
|
||||||
|
* refuses at once a request nothing serves, so the cost is bounded by the modules that are up.
|
||||||
|
*/
|
||||||
|
export async function toolsOn(bus: Broker): Promise<Listing> {
|
||||||
|
const answered = await bus.request<Record<string, never>, { modules?: { module: string }[] }>(
|
||||||
|
CATALOGUE_MODULES,
|
||||||
|
{},
|
||||||
|
);
|
||||||
|
const names = (answered.modules ?? []).map((m) => m.module).filter((m) => typeof m === "string");
|
||||||
|
|
||||||
|
const asked = await Promise.allSettled(
|
||||||
|
names.map((module) => bus.request<Record<string, never>, ToolsAnswer>(`${module}.${TOOLS_VERB}`, {})),
|
||||||
|
);
|
||||||
|
const tools: Tool[] = [];
|
||||||
|
const notAnswering: string[] = [];
|
||||||
|
asked.forEach((outcome, i) => {
|
||||||
|
const module = names[i]!;
|
||||||
|
if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
|
||||||
|
for (const t of outcome.value.tools) {
|
||||||
|
tools.push({ module, name: t.name, description: t.description, input: t.input });
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
notAnswering.push(module);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
tools.sort((a, b) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
|
||||||
|
notAnswering.sort();
|
||||||
|
return { tools, notAnswering };
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what the account
|
||||||
* permits — one vocabulary, so a refusal names the thing they asked for. */
|
* permits — one vocabulary, so a refusal names the thing they asked for. */
|
||||||
export async function callTool(bus: Broker, key: string, args: unknown): Promise<unknown> {
|
export async function callTool(bus: Broker, key: string, args: unknown): Promise<unknown> {
|
||||||
if (!key.includes(".")) {
|
if (!key.includes(".")) {
|
||||||
@@ -104,18 +159,18 @@ export async function callTool(bus: Broker, key: string, args: unknown): Promise
|
|||||||
* Why a call failed, said so that the remedy is in the words.
|
* 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,
|
* 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
|
* the mesh refused this account, or the tool itself failed. Without this they are one timeout and a
|
||||||
* stack trace.
|
* stack trace.
|
||||||
*/
|
*/
|
||||||
export function whyItFailed(key: string, err: unknown): string {
|
export function whyItFailed(key: string, err: unknown): string {
|
||||||
const message = err instanceof Error ? err.message : String(err);
|
const message = err instanceof Error ? err.message : String(err);
|
||||||
if (/no responders|503/i.test(message)) {
|
if (/no responders|503/i.test(message)) {
|
||||||
return `nothing serves ${key}. The module may not be assigned to any machine, or it is down — ` +
|
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.";
|
"`mesh tools` lists what answered.";
|
||||||
}
|
}
|
||||||
if (/permissions violation|authorization/i.test(message)) {
|
if (/permissions violation|authorization/i.test(message)) {
|
||||||
return `this credential may not call ${key}. What it may call was fixed when it was issued; ` +
|
return `this account may not call ${key}. What it may call was fixed when it was issued — a ` +
|
||||||
"`operator issue` again with the tool named, or ask somebody who can.";
|
"person's by `operator issue`, the console's by its manifest.";
|
||||||
}
|
}
|
||||||
if (/timeout/i.test(message)) {
|
if (/timeout/i.test(message)) {
|
||||||
return `${key} did not answer in time. Something is serving it, so this is the tool being slow ` +
|
return `${key} did not answer in time. Something is serving it, so this is the tool being slow ` +
|
||||||
|
|||||||
+165
@@ -0,0 +1,165 @@
|
|||||||
|
/**
|
||||||
|
* The console's endpoint: MCP over HTTP, on a machine's loopback (novox/hq ADR 0152, design 34 §2).
|
||||||
|
*
|
||||||
|
* **Loopback is the authority boundary.** Whoever can connect is on the machine, and whoever is on the
|
||||||
|
* machine is the account that owns the mesh there (ADR 0034, ADR 0144). So there is no token and no
|
||||||
|
* login here, and the one thing this file enforces is that it binds nothing else: a console reachable
|
||||||
|
* from another machine would be authority over the mesh handed to whoever finds the port.
|
||||||
|
*
|
||||||
|
* The transport is the streamable-HTTP shape an agent host speaks: `POST /mcp` with one JSON-RPC
|
||||||
|
* message, answered with one JSON body. No session, because the surface holds nothing per caller; no
|
||||||
|
* event stream, because nothing here has anything to say unasked.
|
||||||
|
*/
|
||||||
|
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
|
||||||
|
|
||||||
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
|
import { mcpSurface, type Reply, type Request } from "./mcp.js";
|
||||||
|
|
||||||
|
/** The most a request body may be. A tool's arguments are small; a megabyte is somebody else's file. */
|
||||||
|
const BODY_LIMIT = 1 << 20;
|
||||||
|
|
||||||
|
export interface Listening {
|
||||||
|
/** Where it listens, as `host:port`, with the port the machine actually gave. */
|
||||||
|
address: string;
|
||||||
|
close(): Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Hosts that are this machine and no other. */
|
||||||
|
const loopback = new Set(["127.0.0.1", "::1", "localhost", "[::1]"]);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Listen on `host:port`. Refused unless the host is loopback — said before binding, so a manifest or
|
||||||
|
* a flag that would open the console to a network is a startup failure rather than something
|
||||||
|
* discovered by whoever finds it.
|
||||||
|
*/
|
||||||
|
export async function serveMcpHttp(bus: Broker, who: string, listen: string): Promise<Listening> {
|
||||||
|
const at = listen.lastIndexOf(":");
|
||||||
|
if (at < 0) {
|
||||||
|
throw new Error(`"${listen}" is not host:port`);
|
||||||
|
}
|
||||||
|
const host = listen.slice(0, at);
|
||||||
|
const port = Number(listen.slice(at + 1));
|
||||||
|
if (!loopback.has(host)) {
|
||||||
|
throw new Error(
|
||||||
|
`the console listens on loopback and nowhere else (novox/hq ADR 0152): "${host}" is not this ` +
|
||||||
|
"machine's own address — whoever is on the machine owns the mesh there, and nobody else may reach this",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
if (!Number.isInteger(port) || port < 0 || port > 65535) {
|
||||||
|
throw new Error(`"${listen.slice(at + 1)}" is not a port`);
|
||||||
|
}
|
||||||
|
|
||||||
|
const surface = mcpSurface(bus, who);
|
||||||
|
const server = createServer((req, res) => {
|
||||||
|
void route(req, res, surface.handle).catch((e) => {
|
||||||
|
json(res, 500, { jsonrpc: "2.0", id: null, error: { code: -32603, message: String(e) } });
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
await new Promise<void>((resolve, reject) => {
|
||||||
|
server.once("error", reject);
|
||||||
|
server.listen(port, host.replace(/^\[|\]$/g, ""), () => resolve());
|
||||||
|
});
|
||||||
|
const bound = server.address();
|
||||||
|
const address = typeof bound === "object" && bound ? `${host}:${bound.port}` : listen;
|
||||||
|
return {
|
||||||
|
address,
|
||||||
|
close: () =>
|
||||||
|
new Promise<void>((resolve) => {
|
||||||
|
server.close(() => resolve());
|
||||||
|
}),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
async function route(
|
||||||
|
req: IncomingMessage,
|
||||||
|
res: ServerResponse,
|
||||||
|
handle: (r: Request) => Promise<Reply | undefined>,
|
||||||
|
): Promise<void> {
|
||||||
|
const path = (req.url ?? "/").split("?")[0];
|
||||||
|
if (path === "/") {
|
||||||
|
res.writeHead(200, { "content-type": "text/plain; charset=utf-8" });
|
||||||
|
res.end("the mesh's console: MCP over HTTP at POST /mcp (novox/hq design 34)\n");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (path !== "/mcp") {
|
||||||
|
json(res, 404, { error: "the console serves /mcp and nothing else" });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
switch (req.method) {
|
||||||
|
case "POST":
|
||||||
|
break;
|
||||||
|
case "DELETE":
|
||||||
|
// A host ending a session. There is no session to end; saying so is the truthful answer.
|
||||||
|
res.writeHead(204).end();
|
||||||
|
return;
|
||||||
|
case "GET":
|
||||||
|
// A host opening an event stream. The console has nothing to say unasked.
|
||||||
|
res.writeHead(405, { allow: "POST, DELETE" }).end();
|
||||||
|
return;
|
||||||
|
default:
|
||||||
|
res.writeHead(405, { allow: "POST, DELETE" }).end();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let body: string;
|
||||||
|
try {
|
||||||
|
body = await read(req);
|
||||||
|
} catch (e) {
|
||||||
|
json(res, 413, { jsonrpc: "2.0", id: null, error: { code: -32600, message: String(e) } });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let parsed: unknown;
|
||||||
|
try {
|
||||||
|
parsed = JSON.parse(body);
|
||||||
|
} catch {
|
||||||
|
json(res, 400, { jsonrpc: "2.0", id: null, error: { code: -32700, message: "the body is not JSON" } });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// One message, or a batch of them; a batch is answered as a batch. A notification gets no reply
|
||||||
|
// and, alone, no body: 202 is how the transport says "heard".
|
||||||
|
if (Array.isArray(parsed)) {
|
||||||
|
const replies = (await Promise.all(parsed.map((r) => handle(r as Request)))).filter(Boolean);
|
||||||
|
if (replies.length === 0) {
|
||||||
|
res.writeHead(202).end();
|
||||||
|
} else {
|
||||||
|
json(res, 200, replies);
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const reply = await handle(parsed as Request);
|
||||||
|
if (!reply) {
|
||||||
|
res.writeHead(202).end();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
json(res, 200, reply);
|
||||||
|
}
|
||||||
|
|
||||||
|
function read(req: IncomingMessage): Promise<string> {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
let size = 0;
|
||||||
|
const chunks: Buffer[] = [];
|
||||||
|
req.on("data", (chunk: Buffer) => {
|
||||||
|
size += chunk.length;
|
||||||
|
if (size > BODY_LIMIT) {
|
||||||
|
reject(new Error(`the request is larger than ${BODY_LIMIT} bytes`));
|
||||||
|
req.destroy();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
chunks.push(chunk);
|
||||||
|
});
|
||||||
|
req.on("end", () => resolve(Buffer.concat(chunks).toString("utf8")));
|
||||||
|
req.on("error", reject);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
function json(res: ServerResponse, status: number, body: unknown): void {
|
||||||
|
const text = JSON.stringify(body);
|
||||||
|
res.writeHead(status, {
|
||||||
|
"content-type": "application/json; charset=utf-8",
|
||||||
|
"content-length": Buffer.byteLength(text),
|
||||||
|
});
|
||||||
|
res.end(text);
|
||||||
|
}
|
||||||
+47
-9
@@ -5,6 +5,10 @@
|
|||||||
// starts consuming as it is imported, so this also runs consumers.
|
// starts consuming as it is imported, so this also runs consumers.
|
||||||
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
|
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
|
||||||
// and what an events test uses to put a message on the wire.
|
// and what an events test uses to put a message on the wire.
|
||||||
|
// mesh-tools prepare bring this module's state to the shape this version needs and exit
|
||||||
|
// — the runtime's answer to the word the mesh asks every module
|
||||||
|
// (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which
|
||||||
|
// the module's own image names beside MESH_TOOL_MODULES.
|
||||||
// mesh-tools run ENTRYPOINT run one compiled module entrypoint to completion and exit — the
|
// mesh-tools run ENTRYPOINT run one compiled module entrypoint to completion and exit — the
|
||||||
// runtime side of a run-once step (novox/hq ADR 0052). It imports
|
// runtime side of a run-once step (novox/hq ADR 0052). It imports
|
||||||
// the given entrypoint, whose top-level code does its work — seed a
|
// the given entrypoint, whose top-level code does its work — seed a
|
||||||
@@ -18,13 +22,12 @@
|
|||||||
// amqps account scoped to this module. Preferred: a module holds its own.
|
// amqps account scoped to this module. Preferred: a module holds its own.
|
||||||
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
|
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
|
||||||
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
|
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
|
||||||
|
// MESH_PREPARE /path/a,/path/b,… compiled entrypoints that prepare this module's state
|
||||||
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
|
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
|
||||||
|
|
||||||
import { readFileSync } from "node:fs";
|
import { readFileSync } from "node:fs";
|
||||||
import { pathToFileURL } from "node:url";
|
import { pathToFileURL } from "node:url";
|
||||||
import { connectAmqp, fatalBrokerReason as fatalAmqpReason } from "./broker-amqp.js";
|
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js";
|
||||||
import { connectNats, fatalBrokerReason as fatalNatsReason } from "./broker-nats.js";
|
|
||||||
import type { Credential } from "./broker-amqp.js";
|
|
||||||
import { runTools } from "./runtime.js";
|
import { runTools } from "./runtime.js";
|
||||||
import { invokeTool } from "@novox/mesh-sdk/tools";
|
import { invokeTool } from "@novox/mesh-sdk/tools";
|
||||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||||
@@ -39,7 +42,7 @@ import { emit } from "@novox/mesh-sdk/events";
|
|||||||
// fatalBrokerReasonFor is the reason a connection failure is final rather than "not yet", for
|
// fatalBrokerReasonFor is the reason a connection failure is final rather than "not yet", for
|
||||||
// whichever bus this runtime is on — each transport knows its own refusals.
|
// whichever bus this runtime is on — each transport knows its own refusals.
|
||||||
function fatalBrokerReasonFor(err: unknown): string | null {
|
function fatalBrokerReasonFor(err: unknown): string | null {
|
||||||
return fatalNatsReason(err) ?? fatalAmqpReason(err);
|
return fatalNatsReason(err);
|
||||||
}
|
}
|
||||||
|
|
||||||
async function connectBroker(): Promise<Broker> {
|
async function connectBroker(): Promise<Broker> {
|
||||||
@@ -66,10 +69,8 @@ async function connectBroker(): Promise<Broker> {
|
|||||||
// nothing else in its environment changed (design 25; novox/hq design 28 task 5.2). The scheme
|
// nothing else in its environment changed (design 25; novox/hq design 28 task 5.2). The scheme
|
||||||
// is enough to know which bus to speak; a runtime that always dialled the old one would keep
|
// is enough to know which bus to speak; a runtime that always dialled the old one would keep
|
||||||
// serving and answer nobody.
|
// serving and answer nobody.
|
||||||
if (credential.url.startsWith("nats://")) {
|
// One bus (novox/hq ADR 0131, design 28 task 5.5): the credential names it, and it is this.
|
||||||
return connectNats(credential);
|
return connectNats(credential);
|
||||||
}
|
|
||||||
return connectAmqp(credential, { assumeExchanges: true });
|
|
||||||
}
|
}
|
||||||
const url = process.env.MESH_BROKER_URL;
|
const url = process.env.MESH_BROKER_URL;
|
||||||
if (!url) {
|
if (!url) {
|
||||||
@@ -78,7 +79,7 @@ async function connectBroker(): Promise<Broker> {
|
|||||||
);
|
);
|
||||||
process.exit(1);
|
process.exit(1);
|
||||||
}
|
}
|
||||||
return connectAmqp(url);
|
return connectNats({ url });
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -181,12 +182,49 @@ async function runEntry(entrypoint: string): Promise<void> {
|
|||||||
await import(pathToFileURL(entrypoint).href);
|
await import(pathToFileURL(entrypoint).href);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Bring this module's state to the shape this version needs, and exit — the runtime's answer to the
|
||||||
|
* one word the mesh asks every module (novox/hq ADR 0135).
|
||||||
|
*
|
||||||
|
* The entrypoints come from `MESH_PREPARE`, which a module's own image sets beside the entrypoints it
|
||||||
|
* already lists there: the module knows which of its files prepares its state, and nothing else could.
|
||||||
|
* Each is imported in the order given, to completion, with no broker — preparation runs before the
|
||||||
|
* version that would use it, so there is nothing yet to talk to.
|
||||||
|
*
|
||||||
|
* **An empty list is a failure, not a no-op.** The mesh only asks this of a module whose manifest says
|
||||||
|
* it prepares something; a module that says so and names nothing has been built wrong, and exiting 0
|
||||||
|
* would let that version serve against a state nobody shaped.
|
||||||
|
*/
|
||||||
|
async function prepareState(): Promise<void> {
|
||||||
|
const named = (process.env.MESH_PREPARE ?? "")
|
||||||
|
.split(",")
|
||||||
|
.map((entry) => entry.trim())
|
||||||
|
.filter((entry) => entry !== "");
|
||||||
|
if (named.length === 0) {
|
||||||
|
console.error(
|
||||||
|
"mesh-tools prepare: this module was asked to prepare its state and its image names nothing " +
|
||||||
|
"to do it with — set MESH_PREPARE to the compiled entrypoint(s) that prepare it, the way " +
|
||||||
|
"MESH_TOOL_MODULES names the ones it serves",
|
||||||
|
);
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
for (const entrypoint of named) {
|
||||||
|
console.log(`[mesh-tools] preparing with ${entrypoint}`);
|
||||||
|
await runEntry(entrypoint);
|
||||||
|
}
|
||||||
|
console.log(`[mesh-tools] prepared: ${named.length} entrypoint(s) ran to completion`);
|
||||||
|
}
|
||||||
|
|
||||||
async function main(): Promise<void> {
|
async function main(): Promise<void> {
|
||||||
const [command, ...rest] = process.argv.slice(2);
|
const [command, ...rest] = process.argv.slice(2);
|
||||||
if (command === "run") {
|
if (command === "run") {
|
||||||
await runEntry(rest[0] ?? "");
|
await runEntry(rest[0] ?? "");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
if (command === "prepare") {
|
||||||
|
await prepareState();
|
||||||
|
return;
|
||||||
|
}
|
||||||
if (command === "invoke") {
|
if (command === "invoke") {
|
||||||
const [module, tool] = rest;
|
const [module, tool] = rest;
|
||||||
if (!module || !tool) {
|
if (!module || !tool) {
|
||||||
|
|||||||
+154
-89
@@ -1,47 +1,179 @@
|
|||||||
/**
|
/**
|
||||||
* The mesh's tools as an MCP server, over stdio (novox/hq design 25 §7).
|
* The mesh's tools as an MCP server (novox/hq design 25 §7, design 34).
|
||||||
*
|
*
|
||||||
* **A thin adapter and nothing more.** Every tool an agent sees is one the catalogue listed and one
|
* **A thin adapter and nothing more.** Every tool an agent sees is one a module answered for and one
|
||||||
* this credential may call; the schema is the module's own; the answer is the module's own. Nothing
|
* this account may call; the schema is the module's own; the answer is the module's own. Nothing here
|
||||||
* here decides anything, which is why it is short — an MCP surface that reshaped arguments or
|
* decides anything, which is why it is short — an MCP surface that reshaped arguments or summarised
|
||||||
* summarised answers would be a second definition of what a tool is, and the module's manifest is the
|
* answers would be a second definition of what a tool is, and the module's code is the first.
|
||||||
* first.
|
|
||||||
*
|
*
|
||||||
* Implemented against the protocol directly rather than through a library: the surface is three
|
* 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.
|
* methods and one framing, and a dependency here would be a dependency on every machine.
|
||||||
|
*
|
||||||
|
* One handler, two transports. Over stdio for a program a person starts (`mesh mcp`), over HTTP on a
|
||||||
|
* machine's loopback for the console the mesh assigns there (`mesh serve`, http.ts). The handler does
|
||||||
|
* not know which asked.
|
||||||
*/
|
*/
|
||||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
import { callTool, toolsOn, whyItFailed, type Tool } from "./client.js";
|
import { callTool, toolsOn, whyItFailed, type Listing } from "./client.js";
|
||||||
|
|
||||||
/** The protocol version this speaks. Stated, because a host that wants another should be told so
|
/** 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. */
|
* rather than discovering it through a shape it did not expect. */
|
||||||
const PROTOCOL = "2024-11-05";
|
export const PROTOCOL = "2025-03-26";
|
||||||
|
|
||||||
interface Request {
|
export interface Request {
|
||||||
jsonrpc: string;
|
jsonrpc: string;
|
||||||
id?: number | string | null;
|
id?: number | string | null;
|
||||||
method: string;
|
method: string;
|
||||||
params?: Record<string, unknown>;
|
params?: Record<string, unknown>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface Reply {
|
||||||
|
jsonrpc: "2.0";
|
||||||
|
id: Request["id"];
|
||||||
|
result?: unknown;
|
||||||
|
error?: { code: number; message: string };
|
||||||
|
}
|
||||||
|
|
||||||
|
/** How long a fetched tool list is kept before the modules are asked again. An agent asks on every
|
||||||
|
* turn; the mesh changes on the order of minutes. */
|
||||||
|
export const LISTING_KEPT_MS = 30_000;
|
||||||
|
|
||||||
|
export interface Surface {
|
||||||
|
/** Answer one request; undefined for a notification, which expects none. */
|
||||||
|
handle(request: Request): Promise<Reply | undefined>;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Serve until stdin closes, which is how a host ends a session.
|
* The surface over one bus connection, as one account.
|
||||||
*
|
*
|
||||||
* The tool list is fetched once, on the first `tools/list`, and kept. An agent asks for it repeatedly
|
* The tool list is fetched when first asked and kept for a short while (design 34 §3): asking every
|
||||||
* and the catalogue's answer does not change mid-session; refetching would make every turn cost a
|
* module on every `tools/list` would fan out on every agent turn for something nobody changed, and
|
||||||
* round trip to a module for something nobody changed.
|
* never refreshing would hide a module assigned a moment ago.
|
||||||
|
*/
|
||||||
|
export function mcpSurface(bus: Broker, who: string): Surface {
|
||||||
|
let known: { listing: Listing; at: number } | undefined;
|
||||||
|
|
||||||
|
const answer = (id: Request["id"], result: unknown): Reply => ({ jsonrpc: "2.0", id, result });
|
||||||
|
const refuse = (id: Request["id"], code: number, message: string): Reply => ({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id,
|
||||||
|
error: { code, message },
|
||||||
|
});
|
||||||
|
|
||||||
|
const listing = async (): Promise<Listing> => {
|
||||||
|
if (!known || Date.now() - known.at > LISTING_KEPT_MS) {
|
||||||
|
known = { listing: await toolsOn(bus), at: Date.now() };
|
||||||
|
}
|
||||||
|
return known.listing;
|
||||||
|
};
|
||||||
|
|
||||||
|
return {
|
||||||
|
async handle(request) {
|
||||||
|
// 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":
|
||||||
|
return 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 account was issued, so a ` +
|
||||||
|
`refusal means the account, not the tool. The list is what the running modules ` +
|
||||||
|
`answered; a module that did not answer is named in the list's _meta and can still be ` +
|
||||||
|
`called by <module>.<tool>. The mesh's own verbs (status, push, assign) are not served ` +
|
||||||
|
`on the bus yet.`,
|
||||||
|
});
|
||||||
|
|
||||||
|
case "notifications/initialized":
|
||||||
|
return undefined;
|
||||||
|
|
||||||
|
case "ping":
|
||||||
|
return notification ? undefined : answer(request.id, {});
|
||||||
|
|
||||||
|
case "tools/list": {
|
||||||
|
let have: Listing;
|
||||||
|
try {
|
||||||
|
have = await listing();
|
||||||
|
} catch (e) {
|
||||||
|
return refuse(request.id, -32603, whyItFailed("mesh-catalog.catalog_modules", e));
|
||||||
|
}
|
||||||
|
return answer(request.id, {
|
||||||
|
tools: have.tools.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: asSchema(t.input),
|
||||||
|
})),
|
||||||
|
// Silence, named (design 34 §3): the modules the catalogue holds and nothing answered
|
||||||
|
// for. Not a tool, so not in `tools`; not dropped either.
|
||||||
|
_meta: { notAnswering: have.notAnswering },
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
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.
|
||||||
|
return 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.
|
||||||
|
return answer(request.id, {
|
||||||
|
content: [{ type: "text", text: whyItFailed(name, e) }],
|
||||||
|
isError: true,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
default:
|
||||||
|
return notification ? undefined : refuse(request.id, -32601, `mesh's MCP surface has no ${request.method}`);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A module's declared input as a JSON schema an agent can read.
|
||||||
|
*
|
||||||
|
* The sdk keeps a tool's input opaque, and the catalogue's modules write it as a bare map of
|
||||||
|
* property to description — `{ module: { type, description } }` — which is the `properties` of a
|
||||||
|
* schema rather than a schema. Wrapped here when that is what arrived; passed through when a module
|
||||||
|
* already wrote a schema; an empty object when it declared nothing. The module's words are kept
|
||||||
|
* either way.
|
||||||
|
*/
|
||||||
|
export function asSchema(input: unknown): Record<string, unknown> {
|
||||||
|
if (!input || typeof input !== "object" || Array.isArray(input)) {
|
||||||
|
return { type: "object", properties: {} };
|
||||||
|
}
|
||||||
|
const given = input as Record<string, unknown>;
|
||||||
|
if (given.type === "object" || "properties" in given) return given;
|
||||||
|
if (Object.keys(given).length === 0) return { type: "object", properties: {} };
|
||||||
|
return { type: "object", properties: given };
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Serve over stdio until stdin closes, which is how a host ends a session.
|
||||||
*/
|
*/
|
||||||
export async function serveMcp(bus: Broker, who: string): Promise<void> {
|
export async function serveMcp(bus: Broker, who: string): Promise<void> {
|
||||||
let known: Tool[] | undefined;
|
const surface = mcpSurface(bus, who);
|
||||||
|
|
||||||
const say = (message: unknown) => {
|
const say = (message: unknown) => {
|
||||||
process.stdout.write(`${JSON.stringify(message)}\n`);
|
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()) {
|
for await (const line of lines()) {
|
||||||
let request: Request;
|
let request: Request;
|
||||||
try {
|
try {
|
||||||
@@ -51,75 +183,8 @@ export async function serveMcp(bus: Broker, who: string): Promise<void> {
|
|||||||
// reply to a request that was never framed is noise on the same channel.
|
// reply to a request that was never framed is noise on the same channel.
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
// A notification has no id and expects no answer; `initialized` is the one every host sends.
|
const reply = await surface.handle(request);
|
||||||
const notification = request.id === undefined || request.id === null;
|
if (reply) say(reply);
|
||||||
|
|
||||||
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}`);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+179
-52
@@ -1,63 +1,100 @@
|
|||||||
#!/usr/bin/env node
|
#!/usr/bin/env node
|
||||||
/**
|
/**
|
||||||
* `mesh` — the mesh's tools from a workstation, for a person (novox/hq design 25 §7).
|
* `mesh` — the mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
|
||||||
*
|
*
|
||||||
* Three verbs and nothing else. What tools are there, call one, and serve the same two to an agent
|
* Four verbs and nothing else. What tools are there, call one, serve the same two to an agent over
|
||||||
* over MCP. Deliberately thin: everything that could be a decision is one the mesh already made, and a
|
* stdio, and serve them on a machine's loopback as the console the mesh assigns. Deliberately thin:
|
||||||
* client that grew opinions would be a second place the mesh's behaviour is defined.
|
* 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 tools what the running modules answer
|
||||||
* mesh call <module>.<tool> [json] call one, arguments as JSON on the command line or on stdin
|
* mesh call <module>.<tool> [json] call one, arguments as JSON on the command line or on stdin
|
||||||
* mesh mcp the same, as an MCP server over stdio
|
* mesh mcp the same two, as an MCP server over stdio
|
||||||
|
* mesh serve [--listen host:port] the console: MCP over HTTP on this machine's loopback
|
||||||
*
|
*
|
||||||
* The credential comes from MESH_CREDENTIAL, or --credential. It is the JSON `operator issue` printed.
|
* Who it speaks as, in order of preference:
|
||||||
|
* MESH_BROKER_FILE the module credential the mesh delivered — the console's (ADR 0152)
|
||||||
|
* MESH_CREDENTIAL / --credential <file> a person's, as `operator issue` printed it (design 25 §7)
|
||||||
|
* MESH_CONSOLE / --console <url> no credential: `tools` and `call` go through a console
|
||||||
|
* already running on this machine, over loopback HTTP
|
||||||
*/
|
*/
|
||||||
import { readFile } from "node:fs/promises";
|
import { readFile } from "node:fs/promises";
|
||||||
|
|
||||||
import { callTool, connectAs, credentialFrom, toolsOn, whyItFailed, type Tool } from "./client.js";
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
|
import {
|
||||||
|
callTool,
|
||||||
|
connectAs,
|
||||||
|
connectAsTheConsole,
|
||||||
|
credentialFrom,
|
||||||
|
toolsOn,
|
||||||
|
whyItFailed,
|
||||||
|
type Listing,
|
||||||
|
} from "./client.js";
|
||||||
|
import { serveMcpHttp } from "./http.js";
|
||||||
import { serveMcp } from "./mcp.js";
|
import { serveMcp } from "./mcp.js";
|
||||||
|
|
||||||
const usage = `mesh tools
|
const usage = `mesh tools
|
||||||
mesh call <module>.<tool> [json]
|
mesh call <module>.<tool> [json]
|
||||||
mesh mcp
|
mesh mcp
|
||||||
|
mesh serve [--listen host:port]
|
||||||
|
|
||||||
--credential <file> the JSON \`operator issue\` printed; default $MESH_CREDENTIAL`;
|
--credential <file> a person's credential, the JSON \`operator issue\` printed; default $MESH_CREDENTIAL
|
||||||
|
--console <url> a console on this machine to ask through instead; default $MESH_CONSOLE
|
||||||
|
--listen <host:port> where \`serve\` listens; loopback only; default $MESH_CONSOLE_LISTEN or 127.0.0.1:4270
|
||||||
|
MESH_BROKER_FILE the module credential the mesh delivered, which \`serve\` runs on`;
|
||||||
|
|
||||||
|
/** What the console listens on when nothing says otherwise. */
|
||||||
|
const DEFAULT_LISTEN = "127.0.0.1:4270";
|
||||||
|
|
||||||
async function main(argv: string[]): Promise<number> {
|
async function main(argv: string[]): Promise<number> {
|
||||||
const args = [...argv];
|
const args = [...argv];
|
||||||
let credentialPath = process.env.MESH_CREDENTIAL ?? "";
|
let credentialPath = process.env.MESH_CREDENTIAL ?? "";
|
||||||
|
let consoleUrl = process.env.MESH_CONSOLE ?? "";
|
||||||
|
let listen = process.env.MESH_CONSOLE_LISTEN ?? DEFAULT_LISTEN;
|
||||||
for (let i = 0; i < args.length; i++) {
|
for (let i = 0; i < args.length; i++) {
|
||||||
if (args[i] === "--credential") {
|
const take = () => {
|
||||||
credentialPath = args[i + 1] ?? "";
|
const v = args[i + 1] ?? "";
|
||||||
args.splice(i, 2);
|
args.splice(i, 2);
|
||||||
i--;
|
i--;
|
||||||
}
|
return v;
|
||||||
|
};
|
||||||
|
if (args[i] === "--credential") credentialPath = take();
|
||||||
|
else if (args[i] === "--console") consoleUrl = take();
|
||||||
|
else if (args[i] === "--listen") listen = take();
|
||||||
}
|
}
|
||||||
const verb = args.shift();
|
const verb = args.shift();
|
||||||
if (!verb || verb === "help" || verb === "--help") {
|
if (!verb || verb === "help" || verb === "--help") {
|
||||||
console.log(usage);
|
console.log(usage);
|
||||||
return verb ? 0 : 1;
|
return verb ? 0 : 1;
|
||||||
}
|
}
|
||||||
if (!credentialPath) {
|
|
||||||
console.error(
|
// Through a console already on this machine: no credential to hold, which is the point of one.
|
||||||
"no credential: set MESH_CREDENTIAL or pass --credential <file>. It is the JSON " +
|
if (consoleUrl && (verb === "tools" || verb === "call")) {
|
||||||
"`operator issue` printed, saved verbatim.",
|
return verb === "tools" ? listingVia(consoleUrl) : callingVia(consoleUrl, args);
|
||||||
);
|
|
||||||
return 1;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
const held = await credentialFrom(credentialPath);
|
const { bus, who } = await connecting(credentialPath, verb);
|
||||||
const bus = await connectAs(held);
|
|
||||||
try {
|
try {
|
||||||
switch (verb) {
|
switch (verb) {
|
||||||
case "tools":
|
case "tools":
|
||||||
return await listing(bus, held.person);
|
return await listing(bus, who);
|
||||||
case "call":
|
case "call":
|
||||||
return await calling(bus, args);
|
return await calling(bus, args);
|
||||||
case "mcp":
|
case "mcp":
|
||||||
// Serves until stdin closes, which is how an MCP host ends a session.
|
// Serves until stdin closes, which is how an MCP host ends a session.
|
||||||
await serveMcp(bus, held.person ?? held.user ?? "somebody");
|
await serveMcp(bus, who);
|
||||||
return 0;
|
return 0;
|
||||||
|
case "serve": {
|
||||||
|
const up = await serveMcpHttp(bus, who, listen);
|
||||||
|
console.log(`mesh console listening on http://${up.address}/mcp as ${who}`);
|
||||||
|
await new Promise<void>((resolve) => {
|
||||||
|
process.once("SIGTERM", () => resolve());
|
||||||
|
process.once("SIGINT", () => resolve());
|
||||||
|
});
|
||||||
|
await up.close();
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
default:
|
default:
|
||||||
console.error(`mesh has no "${verb}".\n\n${usage}`);
|
console.error(`mesh has no "${verb}".\n\n${usage}`);
|
||||||
return 1;
|
return 1;
|
||||||
@@ -67,50 +104,70 @@ async function main(argv: string[]): Promise<number> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async function listing(bus: Awaited<ReturnType<typeof connectAs>>, who?: string): Promise<number> {
|
/**
|
||||||
let tools: Tool[];
|
* Who this process is on the bus. The module credential first: a console is started by the mesh with
|
||||||
try {
|
* MESH_BROKER_FILE and nothing else, and must not fall back to a person's file lying around.
|
||||||
tools = await toolsOn(bus);
|
*/
|
||||||
} catch (e) {
|
async function connecting(credentialPath: string, verb: string): Promise<{ bus: Broker; who: string }> {
|
||||||
console.error(whyItFailed("mesh-catalog.catalog_tools", e));
|
const delivered = process.env.MESH_BROKER_FILE;
|
||||||
return 1;
|
if (delivered) {
|
||||||
|
return connectAsTheConsole(delivered);
|
||||||
}
|
}
|
||||||
if (tools.length === 0) {
|
if (!credentialPath) {
|
||||||
console.log("the catalogue lists no tools; nothing on this mesh serves any");
|
throw new Error(
|
||||||
return 0;
|
verb === "serve"
|
||||||
|
? "no credential: the console runs on MESH_BROKER_FILE, the module credential the mesh " +
|
||||||
|
"delivered; to run it by hand, pass --credential <file> with a person's credential"
|
||||||
|
: "no credential: set MESH_CREDENTIAL or pass --credential <file> (the JSON `operator issue` " +
|
||||||
|
"printed, saved verbatim), or --console <url> to ask through a console on this machine",
|
||||||
|
);
|
||||||
}
|
}
|
||||||
// **What the catalogue has, not what this credential may call.** The two differ and the difference
|
const held = await credentialFrom(credentialPath);
|
||||||
// is the point: a person seeing only their own tools cannot tell "not installed" from "not yours",
|
return { bus: await connectAs(held), who: held.person ?? held.user ?? "somebody" };
|
||||||
// and those need different people to fix them.
|
}
|
||||||
for (const t of tools) {
|
|
||||||
|
function printListing(have: Listing, who?: string): void {
|
||||||
|
if (have.tools.length === 0) {
|
||||||
|
console.log("no running module answered with any tool");
|
||||||
|
}
|
||||||
|
// **What the modules answered, not what this account may call.** The two differ and the
|
||||||
|
// difference is the point: an account seeing only its own tools cannot tell "not installed" from
|
||||||
|
// "not yours", and those need different people to fix them.
|
||||||
|
for (const t of have.tools) {
|
||||||
const name = `${t.module}.${t.name}`;
|
const name = `${t.module}.${t.name}`;
|
||||||
console.log(t.description ? `${name.padEnd(36)} ${t.description}` : name);
|
console.log(t.description ? `${name.padEnd(36)} ${t.description}` : name);
|
||||||
}
|
}
|
||||||
if (who) {
|
if (have.notAnswering.length > 0) {
|
||||||
console.log(`\nthis is what the mesh has. What ${who} may call was fixed when the credential was issued.`);
|
console.log(
|
||||||
|
`\nheld by the mesh and not answering: ${have.notAnswering.join(", ")} — not assigned, not up, ` +
|
||||||
|
"or built before the runtime answered `tools`; each can still be called by name",
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
if (who) {
|
||||||
|
console.log(`\nasked as ${who}; what ${who} may call was fixed when the account was issued.`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function listing(bus: Broker, who?: string): Promise<number> {
|
||||||
|
let have: Listing;
|
||||||
|
try {
|
||||||
|
have = await toolsOn(bus);
|
||||||
|
} catch (e) {
|
||||||
|
console.error(whyItFailed("mesh-catalog.catalog_modules", e));
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
printListing(have, who);
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
async function calling(
|
async function calling(bus: Broker, args: string[]): Promise<number> {
|
||||||
bus: Awaited<ReturnType<typeof connectAs>>,
|
|
||||||
args: string[],
|
|
||||||
): Promise<number> {
|
|
||||||
const key = args.shift();
|
const key = args.shift();
|
||||||
if (!key) {
|
if (!key) {
|
||||||
console.error("mesh call <module>.<tool> [json]");
|
console.error("mesh call <module>.<tool> [json]");
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
|
const parsed = await argumentsFrom(args);
|
||||||
let parsed: unknown = {};
|
if (parsed === undefined) return 1;
|
||||||
if (raw.trim() !== "") {
|
|
||||||
try {
|
|
||||||
parsed = JSON.parse(raw);
|
|
||||||
} catch (e) {
|
|
||||||
console.error(`the arguments are not JSON: ${(e as Error).message}`);
|
|
||||||
return 1;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
try {
|
try {
|
||||||
const answer = await callTool(bus, key, parsed);
|
const answer = await callTool(bus, key, parsed);
|
||||||
console.log(JSON.stringify(answer, null, 2));
|
console.log(JSON.stringify(answer, null, 2));
|
||||||
@@ -121,6 +178,76 @@ async function calling(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** The console's answer to one MCP request, over loopback HTTP. */
|
||||||
|
async function viaConsole(consoleUrl: string, method: string, params?: unknown): Promise<any> {
|
||||||
|
const endpoint = consoleUrl.endsWith("/mcp") ? consoleUrl : `${consoleUrl.replace(/\/$/, "")}/mcp`;
|
||||||
|
const res = await fetch(endpoint, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", accept: "application/json" },
|
||||||
|
body: JSON.stringify({ jsonrpc: "2.0", id: 1, method, params }),
|
||||||
|
});
|
||||||
|
if (!res.ok) {
|
||||||
|
throw new Error(`the console at ${endpoint} answered ${res.status}`);
|
||||||
|
}
|
||||||
|
const reply = (await res.json()) as { result?: any; error?: { message: string } };
|
||||||
|
if (reply.error) throw new Error(reply.error.message);
|
||||||
|
return reply.result;
|
||||||
|
}
|
||||||
|
|
||||||
|
async function listingVia(consoleUrl: string): Promise<number> {
|
||||||
|
try {
|
||||||
|
const result = await viaConsole(consoleUrl, "tools/list");
|
||||||
|
const have: Listing = {
|
||||||
|
tools: (result.tools ?? []).map((t: { name: string; description?: string; inputSchema?: unknown }) => {
|
||||||
|
const at = t.name.indexOf(".");
|
||||||
|
return { module: t.name.slice(0, at), name: t.name.slice(at + 1), description: t.description, input: t.inputSchema };
|
||||||
|
}),
|
||||||
|
notAnswering: result._meta?.notAnswering ?? [],
|
||||||
|
};
|
||||||
|
printListing(have);
|
||||||
|
return 0;
|
||||||
|
} catch (e) {
|
||||||
|
console.error(e instanceof Error ? e.message : String(e));
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function callingVia(consoleUrl: string, args: string[]): Promise<number> {
|
||||||
|
const key = args.shift();
|
||||||
|
if (!key) {
|
||||||
|
console.error("mesh call <module>.<tool> [json]");
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
const parsed = await argumentsFrom(args);
|
||||||
|
if (parsed === undefined) return 1;
|
||||||
|
try {
|
||||||
|
const result = await viaConsole(consoleUrl, "tools/call", { name: key, arguments: parsed });
|
||||||
|
const text = result?.content?.[0]?.text ?? JSON.stringify(result);
|
||||||
|
if (result?.isError) {
|
||||||
|
console.error(text);
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
console.log(text);
|
||||||
|
return 0;
|
||||||
|
} catch (e) {
|
||||||
|
console.error(e instanceof Error ? e.message : String(e));
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A call's arguments: JSON on the command line, else on stdin, else nothing. Undefined when what
|
||||||
|
* was given is not JSON, after saying so. */
|
||||||
|
async function argumentsFrom(args: string[]): Promise<unknown> {
|
||||||
|
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
|
||||||
|
if (raw.trim() === "") return {};
|
||||||
|
try {
|
||||||
|
return JSON.parse(raw);
|
||||||
|
} catch (e) {
|
||||||
|
console.error(`the arguments are not JSON: ${(e as Error).message}`);
|
||||||
|
return undefined;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Arguments on stdin, for a call whose JSON is too long or too quoted to type. Empty when stdin is a
|
/** 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
|
* terminal, so `mesh call x.y` with no arguments does not hang waiting for something nobody is
|
||||||
* typing. */
|
* typing. */
|
||||||
|
|||||||
+38
-2
@@ -6,9 +6,22 @@
|
|||||||
import { pathToFileURL } from "node:url";
|
import { pathToFileURL } from "node:url";
|
||||||
import { resolve } from "node:path";
|
import { resolve } from "node:path";
|
||||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||||
import { serveTools, listTools } from "@novox/mesh-sdk/tools";
|
import { collectTools, serveTools, listTools, toolKey } from "@novox/mesh-sdk/tools";
|
||||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
|
||||||
|
* module's tool names, descriptions and argument schemas, from the code that answers them and
|
||||||
|
* from nowhere else. Discovery asks the module, because a copy kept anywhere else drifts.
|
||||||
|
*/
|
||||||
|
export const TOOLS_VERB = "tools";
|
||||||
|
|
||||||
|
/** What `tools` answers for one module. */
|
||||||
|
export interface ToolsAnswer {
|
||||||
|
module: string;
|
||||||
|
tools: { name: string; description: string; input: Readonly<Record<string, unknown>> }[];
|
||||||
|
}
|
||||||
|
|
||||||
export interface RuntimeOptions {
|
export interface RuntimeOptions {
|
||||||
/** The mesh broker to serve over. */
|
/** The mesh broker to serve over. */
|
||||||
broker: Broker;
|
broker: Broker;
|
||||||
@@ -30,6 +43,29 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
|||||||
// runtime that always served would fail for exactly the modules that never needed it.
|
// runtime that always served would fail for exactly the modules that never needed it.
|
||||||
const tools = listTools();
|
const tools = listTools();
|
||||||
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
|
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
|
||||||
|
|
||||||
|
// And, for every module that serves any, the verb that says what it serves. Refused before
|
||||||
|
// anything is bound if a module named a tool of its own `tools`: one name answering two things
|
||||||
|
// is the fault nobody can diagnose afterwards, and the runtime is the only place that sees both.
|
||||||
|
const stops: Array<() => void> = [stop];
|
||||||
|
for (const { module, tools: own } of collectTools()) {
|
||||||
|
if (own.length === 0) continue;
|
||||||
|
if (own.some((t) => t.name === TOOLS_VERB)) {
|
||||||
|
stop();
|
||||||
|
throw new Error(
|
||||||
|
`${module} names a tool "${TOOLS_VERB}", which is the verb the runtime answers for every ` +
|
||||||
|
"module with what it serves (novox/hq ADR 0152) — refused, rename it",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
const answer: ToolsAnswer = {
|
||||||
|
module,
|
||||||
|
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input })),
|
||||||
|
};
|
||||||
|
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async () => answer));
|
||||||
|
}
|
||||||
|
|
||||||
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
||||||
return stop;
|
return () => {
|
||||||
|
for (const s of stops) s();
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,39 +0,0 @@
|
|||||||
import { test } from "node:test";
|
|
||||||
import assert from "node:assert/strict";
|
|
||||||
|
|
||||||
import { registerModuleTools, resetTools, invokeTool } from "@novox/mesh-sdk/tools";
|
|
||||||
import { connectAmqp } from "../dist/broker-amqp.js";
|
|
||||||
import { runTools } from "../dist/runtime.js";
|
|
||||||
|
|
||||||
// Requires a real broker at $MESH_BROKER_URL. The test harness spins LavinMQ (the mesh's broker)
|
|
||||||
// and points this at it; skipped if it is not set, never failed for the environment.
|
|
||||||
const url = process.env.MESH_BROKER_URL;
|
|
||||||
|
|
||||||
test("a module's tool serves and is invoked over a real AMQP broker", { skip: !url }, async () => {
|
|
||||||
resetTools();
|
|
||||||
|
|
||||||
// A module registers a real tool, exactly as umami does.
|
|
||||||
registerModuleTools("demo", () => [
|
|
||||||
{
|
|
||||||
name: "greet",
|
|
||||||
description: "return a greeting",
|
|
||||||
input: { who: { type: "string" } },
|
|
||||||
run: async (args) => ({ hello: String(args.who), from: "the tool runtime" }),
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
|
|
||||||
// The runtime binds the broker and serves the module (no module entrypoints to import here — the
|
|
||||||
// tool is registered in-process — but this is the exact runtime that serves on a node).
|
|
||||||
const serverBroker = await connectAmqp(url!);
|
|
||||||
const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] });
|
|
||||||
|
|
||||||
// A separate connection — a caller, like mesh-controller's command API — invokes over the broker,
|
|
||||||
// by module and tool (novox/hq ADR 0047: served on serve.demo.greet, invoked as demo.greet).
|
|
||||||
const caller = await connectAmqp(url!);
|
|
||||||
const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string };
|
|
||||||
assert.equal(result.hello, "mesh");
|
|
||||||
|
|
||||||
stop();
|
|
||||||
await caller.close();
|
|
||||||
await serverBroker.close();
|
|
||||||
});
|
|
||||||
+23
-12
@@ -18,16 +18,22 @@ import { callTool, toolsOn, whyItFailed } from "../dist/client.js";
|
|||||||
|
|
||||||
const url = process.env.MESH_TEST_NATS;
|
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
|
/** A mesh as discovery sees it (design 34 §3): the catalogue holding three modules, two of them up
|
||||||
* talk to. Two connections, because a person and a module are different users even in a test. */
|
* and answering `tools`, one held and not running. Two connections, because a person and a module
|
||||||
|
* are different users even in a test. */
|
||||||
async function aMeshWithTools() {
|
async function aMeshWithTools() {
|
||||||
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||||
const shop = await connectNats({ url: url!, module: "shop" });
|
const shop = await connectNats({ url: url!, module: "shop" });
|
||||||
await catalogue.handle("catalog_tools", async () => ({
|
await catalogue.handle("catalog_modules", async () => ({
|
||||||
tools: [
|
modules: [{ module: "shop" }, { module: "mesh-catalog" }, { module: "ghost" }],
|
||||||
{ module: "shop", name: "price", description: "what something costs", input: { type: "object" } },
|
}));
|
||||||
{ module: "mesh-catalog", name: "catalog_tools", description: "what tools the mesh has" },
|
await catalogue.handle("tools", async () => ({
|
||||||
],
|
module: "mesh-catalog",
|
||||||
|
tools: [{ name: "catalog_modules", description: "what modules the mesh has", input: {} }],
|
||||||
|
}));
|
||||||
|
await shop.handle("tools", async () => ({
|
||||||
|
module: "shop",
|
||||||
|
tools: [{ name: "price", description: "what something costs", input: { type: "object" } }],
|
||||||
}));
|
}));
|
||||||
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||||
return {
|
return {
|
||||||
@@ -38,17 +44,22 @@ async function aMeshWithTools() {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
test("a person sees what the catalogue says the mesh has, sorted", async (t) => {
|
test("a person sees what the running modules answer, sorted, and who did not answer", async (t) => {
|
||||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
const mesh = await aMeshWithTools();
|
const mesh = await aMeshWithTools();
|
||||||
const person = await connectNats({ url, module: "person.ada" });
|
const person = await connectNats({ url, module: "person.ada" });
|
||||||
try {
|
try {
|
||||||
const tools = await toolsOn(person);
|
const began = Date.now();
|
||||||
|
const have = await toolsOn(person);
|
||||||
assert.deepEqual(
|
assert.deepEqual(
|
||||||
tools.map((x) => `${x.module}.${x.name}`),
|
have.tools.map((x) => `${x.module}.${x.name}`),
|
||||||
["mesh-catalog.catalog_tools", "shop.price"],
|
["mesh-catalog.catalog_modules", "shop.price"],
|
||||||
"the list is what the catalogue answered, in a stable order",
|
"the list is what the modules answered, in a stable order",
|
||||||
);
|
);
|
||||||
|
// Silence is named, never dropped: a module the catalogue holds and nothing answered for.
|
||||||
|
assert.deepEqual(have.notAnswering, ["ghost"]);
|
||||||
|
// And at once: a module that is not running costs nothing, or the list is unusable.
|
||||||
|
assert.ok(Date.now() - began < 5_000, "an absent module waited out the timeout");
|
||||||
} finally {
|
} finally {
|
||||||
await person.close();
|
await person.close();
|
||||||
await mesh.close();
|
await mesh.close();
|
||||||
|
|||||||
@@ -1,142 +0,0 @@
|
|||||||
import { test } from "node:test";
|
|
||||||
import assert from "node:assert/strict";
|
|
||||||
import amqp from "amqplib";
|
|
||||||
|
|
||||||
import { connectAmqp } from "../dist/broker-amqp.js";
|
|
||||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
|
||||||
import { emit, on, type Event } from "@novox/mesh-sdk/events";
|
|
||||||
|
|
||||||
// Binding conformance: does the mesh-tools AMQP *adapter* honour the ADR 0042 wire contract —
|
|
||||||
// headers on the wire, a body that is only the payload, persistent messages, a durable per-consumer
|
|
||||||
// queue, dead-letter, poison handling? This tests the adapter in isolation, against a disposable
|
|
||||||
// broker that stands in for the one the mesh hosts (ADR 0001). It is NOT an event test of the mesh:
|
|
||||||
// that lives in a full lab scenario where the mesh raises the broker as foundation. Requires
|
|
||||||
// $MESH_BROKER_URL (a throwaway broker); skipped, never failed, when it is not set.
|
|
||||||
const url = process.env.MESH_BROKER_URL;
|
|
||||||
|
|
||||||
test("an event rides the wire with ADR 0042 headers and a body that is only the payload", { skip: !url }, async () => {
|
|
||||||
const sub = await connectAmqp(url!);
|
|
||||||
const pub = await connectAmqp(url!);
|
|
||||||
|
|
||||||
// The consumer's identity names its durable queue (<node>.<module>.events).
|
|
||||||
process.env.MESH_NODE = "lab";
|
|
||||||
process.env.MESH_MODULE = "audit-logger";
|
|
||||||
useBroker(() => sub);
|
|
||||||
const got: Event[] = [];
|
|
||||||
await on("#", async (e) => void got.push(e));
|
|
||||||
await delay(200); // let the binding settle before publishing
|
|
||||||
|
|
||||||
// The emitter is a different module on the same node.
|
|
||||||
process.env.MESH_MODULE = "umami";
|
|
||||||
useBroker(() => pub);
|
|
||||||
await emit("module.umami.site.created", { domain: "my-app" }, { causationId: "cmd-1" });
|
|
||||||
|
|
||||||
await waitFor(() => got.length > 0, 4000);
|
|
||||||
const e = got[0];
|
|
||||||
assert.equal(e.type, "module.umami.site.created");
|
|
||||||
assert.equal(e.source, "umami"); // x-source — read from a header, not the body
|
|
||||||
assert.equal(e.node, "lab"); // x-node
|
|
||||||
assert.ok(e.id, "x-event-id present"); // the handle a consumer dedups on
|
|
||||||
assert.equal(e.causationId, "cmd-1"); // x-causation-id round-trips
|
|
||||||
assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/); // x-time, RFC-3339
|
|
||||||
assert.deepEqual(e.body, { domain: "my-app" }); // provenance never leaked into the body
|
|
||||||
|
|
||||||
// The durable per-consumer queue exists and is bound — a passive assert throws if it does not.
|
|
||||||
const probe = await amqp.connect(url!);
|
|
||||||
const pch = await probe.createChannel();
|
|
||||||
await pch.checkQueue("lab.audit-logger.events");
|
|
||||||
await probe.close();
|
|
||||||
|
|
||||||
await sub.close();
|
|
||||||
await pub.close();
|
|
||||||
delete process.env.MESH_NODE;
|
|
||||||
delete process.env.MESH_MODULE;
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a handler that keeps failing dead-letters the event past the redelivery limit", { skip: !url }, async () => {
|
|
||||||
const module = `flaky-${Date.now()}`;
|
|
||||||
const c = await connectAmqp(url!);
|
|
||||||
process.env.MESH_NODE = "lab";
|
|
||||||
process.env.MESH_MODULE = module;
|
|
||||||
useBroker(() => c);
|
|
||||||
|
|
||||||
let attempts = 0;
|
|
||||||
await on("module.test.boom", async () => {
|
|
||||||
attempts++;
|
|
||||||
throw new Error("boom");
|
|
||||||
});
|
|
||||||
await delay(200);
|
|
||||||
|
|
||||||
await emit("module.test.boom", { n: 1 });
|
|
||||||
|
|
||||||
// First delivery requeues once; the redelivered copy is dead-lettered — two attempts, then it
|
|
||||||
// leaves the consumer queue for good.
|
|
||||||
await waitFor(() => attempts >= 2, 5000);
|
|
||||||
await delay(300);
|
|
||||||
assert.equal(attempts, 2, "attempted twice, not looping forever");
|
|
||||||
|
|
||||||
// The poison event is retained on the dead-letter queue for inspection, not vanished.
|
|
||||||
const probe = await amqp.connect(url!);
|
|
||||||
const pch = await probe.createChannel();
|
|
||||||
const dead = await pch.get("mesh.events.dead", { noAck: true });
|
|
||||||
assert.ok(dead, "the rejected event is on mesh.events.dead");
|
|
||||||
assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.test.boom");
|
|
||||||
await probe.close();
|
|
||||||
|
|
||||||
await c.close();
|
|
||||||
delete process.env.MESH_NODE;
|
|
||||||
delete process.env.MESH_MODULE;
|
|
||||||
});
|
|
||||||
|
|
||||||
test("an undecodable event body is dead-lettered, not looped and not silently swallowed", { skip: !url }, async () => {
|
|
||||||
const module = `poison-${Date.now()}`;
|
|
||||||
const c = await connectAmqp(url!);
|
|
||||||
process.env.MESH_NODE = "lab";
|
|
||||||
process.env.MESH_MODULE = module;
|
|
||||||
useBroker(() => c);
|
|
||||||
|
|
||||||
// Drain any earlier dead events so the assert below sees only this test's.
|
|
||||||
const drain = await amqp.connect(url!);
|
|
||||||
const dch = await drain.createChannel();
|
|
||||||
await dch.purgeQueue("mesh.events.dead");
|
|
||||||
|
|
||||||
let handlerRuns = 0;
|
|
||||||
await on("module.poison.raw", async () => void handlerRuns++);
|
|
||||||
await delay(200);
|
|
||||||
|
|
||||||
// Publish a body that is not JSON straight onto the events exchange — a malformed emitter. A
|
|
||||||
// confirm channel, waited on, so the broker has the message before the connection closes.
|
|
||||||
const raw = await amqp.connect(url!);
|
|
||||||
const rch = await raw.createConfirmChannel();
|
|
||||||
rch.publish("mesh.events", "module.poison.raw", Buffer.from("this is not json{"), {
|
|
||||||
persistent: true,
|
|
||||||
headers: { "x-event-id": "poison-1", "x-source": "bad", "x-node": "lab" },
|
|
||||||
});
|
|
||||||
await rch.waitForConfirms();
|
|
||||||
await raw.close();
|
|
||||||
|
|
||||||
// The handler never ran (the body never decoded), and the message is on the dead queue — set
|
|
||||||
// aside for inspection, not stuck redelivering forever.
|
|
||||||
await delay(600);
|
|
||||||
assert.equal(handlerRuns, 0, "a body that never decodes never reaches the handler");
|
|
||||||
const dead = await dch.get("mesh.events.dead", { noAck: true });
|
|
||||||
assert.ok(dead, "the poison event is retained on mesh.events.dead");
|
|
||||||
assert.equal((dead as amqp.GetMessage).fields.routingKey, "module.poison.raw");
|
|
||||||
await drain.close();
|
|
||||||
|
|
||||||
await c.close();
|
|
||||||
delete process.env.MESH_NODE;
|
|
||||||
delete process.env.MESH_MODULE;
|
|
||||||
});
|
|
||||||
|
|
||||||
function delay(ms: number): Promise<void> {
|
|
||||||
return new Promise((r) => setTimeout(r, ms));
|
|
||||||
}
|
|
||||||
|
|
||||||
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
|
||||||
const start = Date.now();
|
|
||||||
while (!cond()) {
|
|
||||||
if (Date.now() - start > ms) throw new Error("condition not met in time");
|
|
||||||
await delay(25);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Vendored
+6
@@ -0,0 +1,6 @@
|
|||||||
|
// A module naming a tool after the verb the runtime answers for every module — refused at load.
|
||||||
|
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||||
|
|
||||||
|
registerModuleTools("clash", () => [
|
||||||
|
{ name: "tools", description: "mine, not the runtime's", input: {}, run: async () => ({}) },
|
||||||
|
]);
|
||||||
Vendored
+12
@@ -0,0 +1,12 @@
|
|||||||
|
// A module's tool entrypoint, as the runtime imports one: registers and returns.
|
||||||
|
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||||
|
|
||||||
|
registerModuleTools("shop", () => [
|
||||||
|
{
|
||||||
|
name: "price",
|
||||||
|
description: "what something costs",
|
||||||
|
input: { of: { type: "string", description: "the thing" } },
|
||||||
|
run: async (args) => ({ of: args.of ?? "nothing", cost: 12 }),
|
||||||
|
},
|
||||||
|
{ name: "refund", description: "give it back", input: {}, run: async () => ({ done: true }) },
|
||||||
|
]);
|
||||||
@@ -0,0 +1,113 @@
|
|||||||
|
/**
|
||||||
|
* The console: the same surface over HTTP on loopback, started the way the mesh starts it — on the
|
||||||
|
* module credential in MESH_BROKER_FILE (novox/hq ADR 0152, design 34 §2).
|
||||||
|
*
|
||||||
|
* 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/http.test.ts
|
||||||
|
*/
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
import { spawn, type ChildProcess } from "node:child_process";
|
||||||
|
import { mkdtemp, writeFile } from "node:fs/promises";
|
||||||
|
import { join } from "node:path";
|
||||||
|
|
||||||
|
import { connectNats } from "../dist/broker-nats.js";
|
||||||
|
import { serveMcpHttp } from "../dist/http.js";
|
||||||
|
|
||||||
|
const url = process.env.MESH_TEST_NATS;
|
||||||
|
|
||||||
|
async function aMesh(t: { after: (fn: () => Promise<void> | void) => void }) {
|
||||||
|
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||||
|
const shop = await connectNats({ url: url!, module: "shop" });
|
||||||
|
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "shop" }] }));
|
||||||
|
await shop.handle("tools", async () => ({
|
||||||
|
module: "shop",
|
||||||
|
tools: [{ name: "price", description: "what something costs", input: {} }],
|
||||||
|
}));
|
||||||
|
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||||
|
t.after(async () => {
|
||||||
|
await catalogue.close();
|
||||||
|
await shop.close();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
/** `mesh serve` as the mesh runs it: MESH_BROKER_FILE, a listen address, nothing else. */
|
||||||
|
async function aConsole(t: { after: (fn: () => Promise<void> | void) => void }): Promise<string> {
|
||||||
|
const dir = await mkdtemp("/tmp/mesh-console-");
|
||||||
|
const credential = join(dir, "broker");
|
||||||
|
await writeFile(credential, JSON.stringify({ url, node: "desk", module: "mesh-console", user: "desk.mesh-console", password: "x" }));
|
||||||
|
const child: ChildProcess = spawn(process.execPath, ["dist/mesh.js", "serve", "--listen", "127.0.0.1:0"], {
|
||||||
|
env: { ...process.env, MESH_BROKER_FILE: credential, MESH_CREDENTIAL: "" },
|
||||||
|
stdio: ["ignore", "pipe", "pipe"],
|
||||||
|
});
|
||||||
|
t.after(() => {
|
||||||
|
child.kill("SIGTERM");
|
||||||
|
});
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
let out = "";
|
||||||
|
let err = "";
|
||||||
|
child.stdout!.on("data", (d) => {
|
||||||
|
out += d.toString();
|
||||||
|
const m = /listening on (http:\/\/[^/]+\/mcp)/.exec(out);
|
||||||
|
if (m) resolve(m[1]!);
|
||||||
|
});
|
||||||
|
child.stderr!.on("data", (d) => (err += d.toString()));
|
||||||
|
child.on("exit", (code) => reject(new Error(`serve exited ${code}: ${err}`)));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async function post(endpoint: string, body: unknown): Promise<{ status: number; json?: any }> {
|
||||||
|
const res = await fetch(endpoint, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", accept: "application/json" },
|
||||||
|
body: JSON.stringify(body),
|
||||||
|
});
|
||||||
|
const text = await res.text();
|
||||||
|
return { status: res.status, json: text ? JSON.parse(text) : undefined };
|
||||||
|
}
|
||||||
|
|
||||||
|
test("the console answers a host on loopback, as the account the mesh gave it", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
await aMesh(t);
|
||||||
|
const endpoint = await aConsole(t);
|
||||||
|
|
||||||
|
const hello = await post(endpoint, { jsonrpc: "2.0", id: 1, method: "initialize", params: {} });
|
||||||
|
assert.equal(hello.status, 200);
|
||||||
|
assert.match(hello.json.result.instructions, /desk\.mesh-console/, "the handshake names the console's account");
|
||||||
|
|
||||||
|
const heard = await post(endpoint, { jsonrpc: "2.0", method: "notifications/initialized" });
|
||||||
|
assert.equal(heard.status, 202, "a notification is heard and not answered");
|
||||||
|
|
||||||
|
const listed = await post(endpoint, { jsonrpc: "2.0", id: 2, method: "tools/list" });
|
||||||
|
assert.deepEqual(listed.json.result.tools.map((x: { name: string }) => x.name), ["shop.price"]);
|
||||||
|
|
||||||
|
const called = await post(endpoint, {
|
||||||
|
jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "shop.price", arguments: { of: "a hat" } },
|
||||||
|
});
|
||||||
|
assert.deepEqual(JSON.parse(called.json.result.content[0].text), { of: "a hat", cost: 12 });
|
||||||
|
|
||||||
|
// A person's client through the same endpoint, with no credential of its own.
|
||||||
|
const { main } = await import("../dist/mesh.js");
|
||||||
|
const logged: string[] = [];
|
||||||
|
const was = console.log;
|
||||||
|
console.log = (line: string) => logged.push(String(line));
|
||||||
|
try {
|
||||||
|
assert.equal(await main(["tools", "--console", endpoint]), 0);
|
||||||
|
} finally {
|
||||||
|
console.log = was;
|
||||||
|
}
|
||||||
|
assert.ok(logged.some((l) => l.startsWith("shop.price")), `the client did not list through the console: ${logged}`);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("the console binds loopback and nowhere else", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
const bus = await connectNats({ url, module: "mesh-console", node: "desk" });
|
||||||
|
try {
|
||||||
|
await assert.rejects(() => serveMcpHttp(bus, "desk.mesh-console", "0.0.0.0:0"), /loopback and nowhere else/);
|
||||||
|
const up = await serveMcpHttp(bus, "desk.mesh-console", "127.0.0.1:0");
|
||||||
|
assert.match(up.address, /^127\.0\.0\.1:\d+$/);
|
||||||
|
await up.close();
|
||||||
|
} finally {
|
||||||
|
await bus.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
+11
-3
@@ -20,8 +20,12 @@ const url = process.env.MESH_TEST_NATS;
|
|||||||
async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void) => void }) {
|
async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void) => void }) {
|
||||||
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||||
const shop = await connectNats({ url: url!, module: "shop" });
|
const shop = await connectNats({ url: url!, module: "shop" });
|
||||||
await catalogue.handle("catalog_tools", async () => ({
|
await catalogue.handle("catalog_modules", async () => ({
|
||||||
tools: [{ module: "shop", name: "price", description: "what something costs" }],
|
modules: [{ module: "shop" }, { module: "ghost" }],
|
||||||
|
}));
|
||||||
|
await shop.handle("tools", async () => ({
|
||||||
|
module: "shop",
|
||||||
|
tools: [{ name: "price", description: "what something costs", input: { of: { type: "string" } } }],
|
||||||
}));
|
}));
|
||||||
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||||
t.after(async () => {
|
t.after(async () => {
|
||||||
@@ -80,7 +84,7 @@ test("a host initialises, lists the mesh's tools and calls one", async (t) => {
|
|||||||
assert.equal(replies.length, 3, `expected three replies, got ${JSON.stringify(replies)}`);
|
assert.equal(replies.length, 3, `expected three replies, got ${JSON.stringify(replies)}`);
|
||||||
|
|
||||||
const hello = byId.get(1)!.result;
|
const hello = byId.get(1)!.result;
|
||||||
assert.equal(hello.protocolVersion, "2024-11-05");
|
assert.equal(hello.protocolVersion, "2025-03-26");
|
||||||
assert.ok(hello.capabilities.tools, "a server offering no tools is not this one");
|
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");
|
assert.match(hello.instructions, /ada/, "the handshake says whose authority a call is made under");
|
||||||
|
|
||||||
@@ -88,6 +92,10 @@ test("a host initialises, lists the mesh's tools and calls one", async (t) => {
|
|||||||
assert.equal(listed.length, 1);
|
assert.equal(listed.length, 1);
|
||||||
assert.equal(listed[0].name, "shop.price", "a tool is named the way a person names it");
|
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");
|
assert.ok(listed[0].inputSchema, "a tool with no schema is one an agent cannot call");
|
||||||
|
// A module's bare property map arrives as a schema an agent can read, its words kept.
|
||||||
|
assert.deepEqual(listed[0].inputSchema, { type: "object", properties: { of: { type: "string" } } });
|
||||||
|
// Silence is named: the module the catalogue holds and nothing answered for.
|
||||||
|
assert.deepEqual(byId.get(2)!.result._meta.notAnswering, ["ghost"]);
|
||||||
|
|
||||||
const called = byId.get(3)!.result;
|
const called = byId.get(3)!.result;
|
||||||
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
|
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
|
import { spawn } from "node:child_process";
|
||||||
import { test } from "node:test";
|
import { test } from "node:test";
|
||||||
import assert from "node:assert/strict";
|
import assert from "node:assert/strict";
|
||||||
import { fatalBrokerReason, PinMismatchError } from "../src/broker-amqp.ts";
|
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
|
||||||
|
|
||||||
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
||||||
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
||||||
@@ -31,11 +32,8 @@ test("a malformed broker URL is fatal — it never parses on the next try", () =
|
|||||||
});
|
});
|
||||||
|
|
||||||
test("a refused login is fatal — a wrong or revoked credential, not an absent broker", () => {
|
test("a refused login is fatal — a wrong or revoked credential, not an absent broker", () => {
|
||||||
for (const msg of [
|
// The bus refuses a login in its own words; each is final, because the next try says the same.
|
||||||
"Handshake terminated by server: 403 (ACCESS-REFUSED) with message \"ACCESS_REFUSED - Login was refused\"",
|
for (const msg of ["Authorization Violation", "nats: user authentication expired", "Permissions Violation for Subscription to \"x\""]) {
|
||||||
"Login was refused using authentication mechanism PLAIN",
|
|
||||||
"ACCESS_REFUSED",
|
|
||||||
]) {
|
|
||||||
assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`);
|
assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
@@ -44,3 +42,39 @@ test("a non-Error value does not crash the classifier", () => {
|
|||||||
assert.equal(fatalBrokerReason("just a string"), null);
|
assert.equal(fatalBrokerReason("just a string"), null);
|
||||||
assert.equal(fatalBrokerReason(undefined), null);
|
assert.equal(fatalBrokerReason(undefined), null);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// **A module asked to prepare its state and naming nothing is a failure, not a no-op** (novox/hq
|
||||||
|
// ADR 0135). The mesh asks this only of a module whose manifest says it prepares something, so an
|
||||||
|
// image that names nothing was built wrong, and exiting 0 would let that version serve against a
|
||||||
|
// state nobody shaped.
|
||||||
|
test("preparing with nothing named fails rather than passing quietly", async () => {
|
||||||
|
const runtime = new URL("../dist/main.js", import.meta.url).pathname;
|
||||||
|
const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
|
||||||
|
const child = spawn(process.execPath, [runtime, "prepare"], {
|
||||||
|
env: { ...process.env, MESH_PREPARE: "" },
|
||||||
|
});
|
||||||
|
let said = "";
|
||||||
|
child.stderr.on("data", (chunk) => (said += String(chunk)));
|
||||||
|
child.on("close", (code) => resolve({ code, said }));
|
||||||
|
});
|
||||||
|
assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
|
||||||
|
assert.match(ran.said, /MESH_PREPARE/);
|
||||||
|
});
|
||||||
|
|
||||||
|
// **One durable consumer feeds one reader, however many patterns a module registers.**
|
||||||
|
//
|
||||||
|
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
|
||||||
|
// reader that received one its own pattern does not match acknowledges it, which is right for a
|
||||||
|
// filter wider than anything registered and silent loss when it is another handler's. The matching is
|
||||||
|
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not
|
||||||
|
// by which reader happened to fetch it.
|
||||||
|
test("a message is for every pattern that matches it, and nothing else", () => {
|
||||||
|
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"];
|
||||||
|
const matched = (key: string) => registered.filter((p) => topicMatches(p, key));
|
||||||
|
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]);
|
||||||
|
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]);
|
||||||
|
// Nothing registered for it: the consumer's filter is the controller's and may be wider.
|
||||||
|
assert.deepEqual(matched("mesh-catalog.upgraded"), []);
|
||||||
|
// And a handler that asked for everything gets both, which is what the audit logger does.
|
||||||
|
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]);
|
||||||
|
});
|
||||||
|
|||||||
@@ -0,0 +1,60 @@
|
|||||||
|
/**
|
||||||
|
* Every module's runtime answers `tools` for it (novox/hq ADR 0152, design 34 §3): the names,
|
||||||
|
* descriptions and schemas from the code that answers them. Against a real bus, because the claim is
|
||||||
|
* what a second connection gets back.
|
||||||
|
*
|
||||||
|
* 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/runtime-tools.test.ts
|
||||||
|
*/
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import { test } from "node:test";
|
||||||
|
import { fileURLToPath } from "node:url";
|
||||||
|
|
||||||
|
import { resetTools } from "@novox/mesh-sdk/tools";
|
||||||
|
|
||||||
|
import { connectNats } from "../dist/broker-nats.js";
|
||||||
|
import { runTools } from "../dist/runtime.js";
|
||||||
|
|
||||||
|
const url = process.env.MESH_TEST_NATS;
|
||||||
|
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
||||||
|
|
||||||
|
test("a module registering two tools answers three names, the third being what it serves", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
resetTools();
|
||||||
|
const shop = await connectNats({ url, node: "one", module: "shop" });
|
||||||
|
const asker = await connectNats({ url, module: "person.ada" });
|
||||||
|
const stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
|
||||||
|
try {
|
||||||
|
const answer = await asker.request<Record<string, never>, { module: string; tools: { name: string; input: unknown }[] }>(
|
||||||
|
"shop.tools",
|
||||||
|
{},
|
||||||
|
);
|
||||||
|
assert.equal(answer.module, "shop");
|
||||||
|
assert.deepEqual(answer.tools.map((x) => x.name), ["price", "refund"]);
|
||||||
|
// The schema travels with the name: a name alone is not callable by something that has never
|
||||||
|
// seen the mesh before.
|
||||||
|
assert.deepEqual(answer.tools[0].input, { of: { type: "string", description: "the thing" } });
|
||||||
|
// And the tools themselves still answer beside it.
|
||||||
|
const priced = await asker.request<{ of: string }, { cost: number }>("shop.price", { of: "a hat" });
|
||||||
|
assert.equal(priced.cost, 12);
|
||||||
|
} finally {
|
||||||
|
stop();
|
||||||
|
await asker.close();
|
||||||
|
await shop.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("a module naming a tool of its own `tools` is refused at load", async (t) => {
|
||||||
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||||
|
resetTools();
|
||||||
|
const clash = await connectNats({ url, node: "one", module: "clash" });
|
||||||
|
try {
|
||||||
|
await assert.rejects(
|
||||||
|
() => runTools({ broker: clash, moduleEntrypoints: [fixture("clash-tools.mjs")] }),
|
||||||
|
/names a tool "tools"/,
|
||||||
|
);
|
||||||
|
} finally {
|
||||||
|
resetTools();
|
||||||
|
await clash.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -1,68 +0,0 @@
|
|||||||
/**
|
|
||||||
* **The old bus's wire is byte-identical after the rename, and this is the test that lets the change
|
|
||||||
* be merged to a running mesh.**
|
|
||||||
*
|
|
||||||
* Every module's event names were converted from the old bus's routing keys to local names
|
|
||||||
* (novox/hq 04-ISSUES/127), and the old bus's client maps them back. If that mapping is wrong
|
|
||||||
* anywhere, a live mesh's events stop being delivered — silently, because a binding that matches
|
|
||||||
* nothing is not an error.
|
|
||||||
*
|
|
||||||
* So this pins the mapping against the literal routing keys the mesh used before, taken from the
|
|
||||||
* manifests as they were. It needs no bus: it is about a string.
|
|
||||||
*/
|
|
||||||
import assert from "node:assert/strict";
|
|
||||||
import { test } from "node:test";
|
|
||||||
|
|
||||||
import { routingKeyFor, bindingFor, localKeyFor, topicMatches } from "../dist/broker-amqp.js";
|
|
||||||
|
|
||||||
test("a converted emit produces the routing key the mesh published before", () => {
|
|
||||||
// left: what the module's code says now. right: what went on the wire before, unchanged.
|
|
||||||
const same: [string, string, string][] = [
|
|
||||||
["plex", "playback.started", "module.plex.playback.started"],
|
|
||||||
["sonarr", "download.completed", "module.sonarr.download.completed"],
|
|
||||||
["builder", "built", "module.builder.built"],
|
|
||||||
["mesh-catalog", "upgraded", "module.mesh-catalog.upgraded"],
|
|
||||||
["keycloak", "user.created", "module.keycloak.user.created"],
|
|
||||||
["mesh-vault", "secret.rotated", "module.mesh-vault.secret.rotated"],
|
|
||||||
];
|
|
||||||
for (const [self, local, before] of same) {
|
|
||||||
assert.equal(routingKeyFor(local, self), before, `${self} emitting ${local}`);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a converted subscription binds what it bound before", () => {
|
|
||||||
const same: [string, string][] = [
|
|
||||||
["builder.built", "module.builder.built"],
|
|
||||||
["*.download.completed", "module.*.download.completed"],
|
|
||||||
["*.usage.*", "module.*.usage.*"],
|
|
||||||
// The audit logger's "everything": `#` on this bus, and it must stay `#`.
|
|
||||||
["**", "#"],
|
|
||||||
];
|
|
||||||
for (const [declared, before] of same) {
|
|
||||||
assert.equal(bindingFor(declared), before, `consuming ${declared}`);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a handler still matches what the bus delivers", () => {
|
|
||||||
// The key a handler is given is the local one now, and the pattern it compares against is local
|
|
||||||
// too — so the pair must still meet for every case the mesh actually has.
|
|
||||||
const pairs: [string, string][] = [
|
|
||||||
["builder.built", "module.builder.built"],
|
|
||||||
["*.download.completed", "module.sonarr.download.completed"],
|
|
||||||
["*.usage.*", "module.anthropic-consumer.usage.session"],
|
|
||||||
["**", "module.anything.at.all"],
|
|
||||||
];
|
|
||||||
for (const [pattern, delivered] of pairs) {
|
|
||||||
assert.ok(
|
|
||||||
topicMatches(pattern, localKeyFor(delivered)),
|
|
||||||
`${pattern} no longer matches ${delivered}, so a running module would stop reacting`,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a routing key already in the old form is left alone", () => {
|
|
||||||
// Belt for the transition: anything not yet converted still goes out as it did, so a module built
|
|
||||||
// from an older manifest keeps working beside one built from a current manifest.
|
|
||||||
assert.equal(routingKeyFor("module.plex.playback.started", "plex"), "module.plex.playback.started");
|
|
||||||
assert.equal(bindingFor("module.builder.built"), "module.builder.built");
|
|
||||||
});
|
|
||||||
Reference in New Issue
Block a user