17 Commits
Author SHA1 Message Date
mesh-admin a3ace58362 Merge pull request 'A registration under a stranger's name is said and skipped, not fatal' (#26) from fix/a-registration-under-a-strangers-name-is-refused-not-fatal into main 2026-10-01 14:47:31 +00:00
jschoubben feb13c1de6 A registration under a stranger's name is said and skipped, not fatal
A module implementing a seat its credential does not (yet) claim registered tools under the seat's
name; the runtime treated them as the module's own, refused to serve another module's key, and the
whole runtime restarted in a loop — postgres on the control node, 2026-10-01, whose credential predated
the claims. Such a registration is now named in the log and left out; the module's own tools serve.
2026-10-01 16:47:03 +02:00
mesh-admin e6676068c7 Merge pull request 'The membership is read by its subject' (#25) from fix/the-membership-is-read-by-its-subject into main 2026-10-01 14:22:26 +00:00
jschoubben a01b5f5b78 The membership is read by its subject
The runtime asked the assignments stream's root for the last message by subject in the body, and the
mesh grants a module's account only the subject-addressed form of the direct get — the server refused
every read (Publish Violation on $JS.API.DIRECT.GET.ASSIGNMENTS for every module on the new runtime),
so every runtime kept the derived shape. The address is now the stream then the subject, nothing in the
body, which is the one address the account has.
2026-10-01 16:16:57 +02:00
mesh-admin ab8b13f512 Merge pull request 'A seat's verb keeps a node of its own; only a module's tool gives it to the subject' (#24) from fix/a-seat-verbs-node-is-its-own into main 2026-10-01 13:44:22 +00:00
jschoubben ce8c37a7fc The listing test tells the seat's two verbs apart: one keeps its own node, the other takes none 2026-10-01 15:43:36 +02:00
jschoubben 809e2b11ec The listing test indexes the module's tool after the seat's two verbs 2026-10-01 15:42:58 +02:00
jschoubben 670b486ac0 A seat's verb keeps a node of its own; only a module's tool gives it to the subject
The console moved every call's node into the subject (ADR 0159), so mesh-controller.push {node: x}
became a call to the seat's verb on machine x, which nothing serves — the mesh's own verbs could not
be given a machine from the console at all. A role's verb takes no machine from the console; its
arguments are its own.
2026-10-01 15:42:35 +02:00
mesh-admin 56f82588af Merge pull request 'A runtime serves what the mesh issued it, and a seat's verbs are implemented under the seat's name (hq ADR 0160)' (#23) from feat/a-runtime-serves-what-it-is-issued into main 2026-10-01 13:23:55 +00:00
jschoubben e7b98f1fbc A runtime serves what the mesh issued it, and a seat's verbs are implemented under the seat's name (hq ADR 0160)
The one address a runtime derives for itself is mesh.assignment.<node>.<module>. It reads the
membership there with a direct get on the ASSIGNMENTS stream, serves each tool exactly where the
membership says — the plain subject in the module's queue when the mesh issued one, this machine's
beside it — and follows the subject, re-serving when a new membership arrives. A mesh that has issued
nothing yet gets the shape it always derived, and the log says so.

A seat's verbs are the role's, not the software's (ADR 0159): a module implements them with
registerModuleTools("<seat>", …), the runtime serves that on the seat's subjects when the credential
claims the seat, and never lists it among the module's own tools. A module named like its seat
registers once and is both.

The tools answer carries each tool's subjects, and the console and CLI call the subject the listing
gave them instead of composing one.
2026-10-01 15:17:14 +02:00
mesh-admin 278a25b3e5 Merge pull request 'A tool call names the machine it is for, every answer says which machine answered, and a holder's runtime serves its seat's verbs (ADR 0159)' (#22) from feat/a-tool-call-names-the-machine into main 2026-10-01 12:01:40 +00:00
jschoubben c65f1993ba A tool call names the machine it is for, every answer says which machine answered, and a holder's
runtime serves its seat's verbs (novox/hq ADR 0159)

A module on several machines served one subject in one queue group, so a call reached whichever
instance answered first and nobody could ask one machine's instance. Now each instance also serves
its subject with its machine as the last token, `<module>.<tool>@<node>` addresses it, and every
answer carries the machine that gave it: the console lists `node` on every module tool, strips it
into the subject, and appends "answered by <node>" to the answer; `mesh call` prints it.

And a seat's verbs are served by whoever claims the seat: the credential names the claimed seats
and the verbs each promises, the runtime serves each verb with the module's tool of the same name
on the seat's own subject, and the bus — which admits only the holder's subscription — decides where
that serving is real. Design 33 §3 and §4, built.
2026-10-01 13:58:22 +02:00
jschoubben 6a91d144b3 Merge pull request 'The console lists and calls a role's tools' (#21) from feat/the-mesh-answers-for-itself into main
Reviewed-on: #21
2026-09-30 15:54:31 +00:00
jschoubben 71965ef958 The console lists and calls a role's tools
seat:<seat>.<verb> addresses a role's tool (with @<node> for a node-scoped seat); the listing asks the
mesh-controller seat's tools verb beside the modules and marks a role's tools; <seat>.<verb> resolves
to the seat when the seat declares that verb, a module's own name otherwise (novox/hq ADR 0154).
2026-09-30 17:41:56 +02:00
jschoubben dea98e509a Merge pull request 'The console: mesh serve on loopback, and every runtime answers tools' (#20) from feat/the-console into main
Reviewed-on: #20
2026-09-30 14:46:54 +00:00
jschoubben 80b02740ab The console: mesh serve on loopback, and every runtime answers tools
The runtime serves a tools verb per module with names, descriptions and schemas (design 34 §3), and
refuses a module naming its own tool tools. Discovery asks catalog_modules then each module, naming
what did not answer. One MCP handler over two transports: stdio (mesh mcp) and loopback HTTP (mesh
serve, the mesh-console module, novox/hq ADR 0152); serve refuses any bind but loopback. tools/call
may go through a running console with --console and no credential.
2026-09-30 16:19:22 +02:00
mesh-admin 621d033d53 Merge pull request 'One consumer, one reader, however many patterns a module registers' (#19) from fix/one-consumer-one-loop into main 2026-09-28 14:25:47 +00:00
18 changed files with 1713 additions and 248 deletions
+14
View File
@@ -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
+225 -33
View File
@@ -40,6 +40,55 @@ export interface Credential {
module?: string; module?: string;
user?: string; user?: string;
password?: string; password?: string;
/** The seats this module claims, with the verbs each promises (novox/hq ADR 0159). The runtime
* serves each claimed seat's verbs with its tools of the same name; the bus admits only the
* holder's subscription, so claiming and not holding costs a refused subscription and nothing else. */
claims?: { seat: string; scope?: string; serves?: string[] }[];
}
/**
* What the mesh issued this assignment (novox/hq ADR 0160): where its tools are served, in which
* queue, the verbs of the seats it holds, where its events land, what it may reach. Read from the
* ASSIGNMENTS stream at `mesh.assignment.<node>.<module>` — the one subject a runtime derives for
* itself — and followed live. Absent for a mesh older than the membership, and then the runtime
* serves the shape it always derived, and says so.
*/
export interface Membership {
node: string;
module: string;
/** Addresses a tool is answered on; `{tool}` stands for the tool's name. */
serves: { subject: string; queue?: string }[];
seats?: { seat: string; verb: string; subject: string }[];
emits: string;
reaches?: Record<string, string[]>;
tools: string;
}
/** What a tool call answers: the module's own result, and which machine answered it
* (novox/hq ADR 0159) — a module on several machines is otherwise an answer from nowhere. */
export interface Answered<Res> {
result: Res;
node?: string;
}
/** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs —
* an answer that says which machine gave it, and serving a subject that is not a module's own tool
* (a seat's verb). */
export interface RuntimeBroker extends Broker {
/** Call a tool by key, or — when `on` names a subject the mesh listed for it (ADR 0160) — there. */
ask<Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>>;
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void>;
/** What the mesh issued this assignment, or undefined when nothing has been issued yet. */
membership(): Membership | undefined;
/** Called when the mesh issues a new membership; the runtime re-serves on it. */
onMembership(handler: (m: Membership) => void): void;
}
const ASSIGNMENTS_STREAM = "ASSIGNMENTS";
/** The one address a runtime derives for itself (ADR 0160). */
export function membershipSubject(node: string, module: string): string {
return `mesh.assignment.${node}.${module}`;
} }
/** Whether a connection failure is worth retrying, or is a fact about this configuration that /** Whether a connection failure is worth retrying, or is a fact about this configuration that
@@ -69,7 +118,7 @@ export function fatalBrokerReason(err: unknown): string | null {
export async function connectNats( export async function connectNats(
target: string | Credential, target: string | Credential,
opts: { module?: string } = {}, opts: { module?: string } = {},
): Promise<Broker> { ): Promise<RuntimeBroker> {
const cred: Credential = typeof target === "string" ? { url: target } : target; const cred: Credential = typeof target === "string" ? { url: target } : target;
const self = cred.module ?? opts.module; const self = cred.module ?? opts.module;
if (!self) { if (!self) {
@@ -102,36 +151,96 @@ export async function connectNats(
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined; let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
let closed = false; let closed = false;
return { const node = cred.node;
/**
* Ask one question and await one answer.
*
* Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3),
* and a lost one is a timeout the caller already handles. The reply travels on the inbox the
* request carries, which the responder may answer because its account has `allow_responses`
* — one reply to a message it actually received, and nothing wider.
*/
async request<Req, Res>(key: string, body: Req): Promise<Res> {
const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), {
timeout: REQUEST_TIMEOUT_MS,
});
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string };
if (reply.error) throw new Error(reply.error);
return reply.result as Res;
},
/** // The membership, read once at connect and followed. A direct get is one request on the
* Answer a question. // stream's API, which is the whole of what this account may ask JetStream for its own subject;
* // a 404 is a mesh that has not issued one, which is a fact to say and not an error to retry.
* A queue group, so several nodes may serve one tool and exactly one of them answers each let issued: Membership | undefined;
* call. const issuedHandlers: ((m: Membership) => void)[] = [];
*/ const subjectOfMine = node ? membershipSubject(node, self) : "";
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> { if (subjectOfMine) {
const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` }); try {
// The subject-addressed form of a direct get — the stream, then the subject, nothing in
// the body — because that is the one address the mesh grants this account on the
// stream's API; the body form asks the stream's root, which it may not.
const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}.${subjectOfMine}`,
new Uint8Array(0), { timeout: 5_000 });
const status = got.headers?.code ?? 0;
if (status === 0 && got.data.length > 0) {
issued = JSON.parse(sc.decode(got.data)) as Membership;
}
} catch {
// Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below.
}
if (!issued) {
console.log(`[mesh-tools] no membership issued for ${self} on ${node} yet; serving the derived shape until one arrives`);
}
try {
const live = conn.subscribe(subjectOfMine);
subs.push(live);
void (async () => {
for await (const msg of live) {
try {
issued = JSON.parse(sc.decode(msg.data)) as Membership;
console.log(`[mesh-tools] ${self} on ${node} was issued a new membership; re-serving on it`);
for (const h of issuedHandlers) h(issued);
} catch (err) {
console.log(`[mesh-tools] a membership arrived that is not one: ${err}`);
}
}
})();
} catch {
// A subscription this account may not make is a mesh older than the membership.
}
}
/** The subjects a tool of this module is served on: from the membership when issued, derived
* otherwise (the shape the mesh issues on day one, so the two agree). */
const servedOn = (tool: string): { subject: string; queue?: string }[] => {
if (issued) {
const m = issued;
const out = m.serves.map((s) => ({ subject: s.subject.replace("{tool}", tool), queue: s.queue }));
// The verb that lists what this module serves is answered on the mesh's plain address for it
// whatever the placement — one answer suffices, so a queue — and on this machine's beside it.
if (tool === "tools" && m.tools && !out.some((s) => s.subject === m.tools)) {
out.unshift({ subject: m.tools, queue: `serve.${self}` });
}
return out;
}
const base = `mesh.mod.${self}.tool.${tool}`;
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }];
if (node) out.push({ subject: `${base}.${node}` });
return out;
};
/** Where a call by key goes: a subject the membership says this module reaches, when it says
* one — the machine's when named — else the derived shape. */
const reachedAt = (key: string): string => {
const [name, wanted] = key.split("@", 2);
const reach = issued?.reaches?.[name];
if (reach && reach.length > 0) {
if (wanted) {
const at = reach.find((s) => s.endsWith(`.${wanted}`));
if (at) return at;
} else {
return reach[0];
}
}
return toolSubject(key, self);
};
/** Answer one subject with one handler, and say which machine answered (novox/hq ADR 0159). */
const answerOn = <Req, Res>(
subject: string,
queue: string | undefined,
handler: (body: Req) => Promise<Res>,
): (() => void) => {
const sub = conn.subscribe(subject, queue ? { queue } : {});
subs.push(sub); subs.push(sub);
void (async () => { void (async () => {
for await (const msg of sub) { for await (const msg of sub) {
let reply: { result?: Res; error?: string }; let reply: { result?: Res; error?: string; node?: string };
try { try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) }; reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) { } catch (err) {
@@ -139,12 +248,73 @@ export async function connectNats(
// different failure from a tool nobody serves, and only one of them is worth retrying. // different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) }; reply = { error: err instanceof Error ? err.message : String(err) };
} }
if (node) reply.node = node;
msg.respond(sc.encode(JSON.stringify(reply))); msg.respond(sc.encode(JSON.stringify(reply)));
} }
})(); })();
return () => { return () => sub.unsubscribe();
sub.unsubscribe();
}; };
/**
* Ask one question and await one answer, with the machine that gave it.
*
* Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3),
* and a lost one is a timeout the caller already handles. The reply travels on the inbox the
* request carries, which the responder may answer because its account has `allow_responses`
* — one reply to a message it actually received, and nothing wider.
*/
const ask = async <Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>> => {
const msg = await conn.request(on ?? reachedAt(key), sc.encode(JSON.stringify(body)), {
timeout: REQUEST_TIMEOUT_MS,
});
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string; node?: string };
if (reply.error) throw new Error(reply.error);
return { result: reply.result as Res, node: reply.node };
};
return {
async request<Req, Res>(key: string, body: Req): Promise<Res> {
return (await ask<Req, Res>(key, body)).result;
},
ask,
/**
* Answer a question, two ways (novox/hq ADR 0159): on the module's subject in a queue group,
* so several machines may serve one tool and exactly one of them answers each call; and on the
* same subject with this machine as its last token, so a caller that names the machine reaches
* this instance and no other. A runtime that does not know its machine serves only the first,
* which is how it always behaved.
*/
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
// A seat's verb named outright is served on the seat's subject as given, for a holder that
// knows its role without a membership; everything else is this module's own tool, served
// where the mesh issued it (ADR 0160). A key naming another module is not served here at all.
if (key.startsWith("seat:")) {
const stop = answerOn(toolSubject(key, self), undefined, handler);
return () => stop();
}
const dot = key.indexOf(".");
if (dot >= 0 && key.slice(0, dot) !== self) {
throw new Error(`${self} cannot serve ${key}: a module serves its own tools`);
}
const tool = dot < 0 ? key : key.slice(dot + 1);
let stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
// When a new membership arrives, serve where it now says and stop serving where it no longer does.
issuedHandlers.push(() => {
stops.forEach((stop) => stop());
stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
});
return () => stops.forEach((stop) => stop());
},
membership: () => issued,
onMembership: (handler: (m: Membership) => void) => {
issuedHandlers.push(handler);
},
async handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
return answerOn(subject, undefined, handler);
}, },
/** /**
@@ -313,11 +483,33 @@ function eventSubject(type: string, self: string): string {
} }
/** A tool's subject. A bare name is this module's own tool; `<module>.<tool>` addresses /** A tool's subject. A bare name is this module's own tool; `<module>.<tool>` addresses
* another's, which is how a request reaches a module that is not this one. */ * another's, which is how a request reaches a module that is not this one; `seat:<seat>.<verb>`
* addresses a role's tool, answered by whoever holds the seat (novox/hq ADR 0132) — with
* `seat:<seat>.<verb>@<node>` for a node-scoped seat, whose tool carries the machine (design 33 §4). */
function toolSubject(key: string, self: string): string { function toolSubject(key: string, self: string): string {
const dot = key.indexOf("."); if (key.startsWith("seat:")) {
if (dot < 0) return `mesh.mod.${self}.tool.${key}`; const rest = key.slice("seat:".length);
return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`; const dot = rest.indexOf(".");
if (dot < 0) throw new Error(`"${key}" names a seat and no verb: seat:<seat>.<verb>`);
const seat = rest.slice(0, dot);
const [verb, node] = rest.slice(dot + 1).split("@", 2);
return node ? `mesh.seat.${seat}.tool.${verb}.${node}` : `mesh.seat.${seat}.tool.${verb}`;
}
// `<module>.<tool>@<node>` names the machine (novox/hq ADR 0159): the same subject with the
// machine as its last token, which is what that instance serves beside the queue.
const [name, node] = key.split("@", 2);
const dot = name.indexOf(".");
const base = dot < 0
? `mesh.mod.${self}.tool.${name}`
: `mesh.mod.${name.slice(0, dot)}.tool.${name.slice(dot + 1)}`;
return node ? `${base}.${node}` : base;
}
/** A seat's verb, as its holder serves it: flat for a mesh seat, carrying the machine for a
* node-scoped one (design 33 §4) — the same shape the controller grants. */
export function seatToolSubject(seat: string, verb: string, scope: string | undefined, node: string | undefined): string {
const base = `mesh.seat.${seat}.tool.${verb}`;
return scope === "node" && node ? `${base}.${node}` : base;
} }
function normalizeFingerprint(fingerprint: string): string { function normalizeFingerprint(fingerprint: string): string {
+183 -36
View File
@@ -1,36 +1,91 @@
/** /**
* 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 Answered, 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. */ /** Where the mesh answers every role's tools, from its records: the mesh-controller seat's own
* `tools` verb (novox/hq ADR 0154, design 33 §5). */
const SEAT_TOOLS = "seat:mesh-controller.tools";
/** A tool as its module describes it. */
export interface Tool { export interface Tool {
/** The module that serves it — or, for a role's tool, the seat. */
module: string; module: string;
name: string; name: string;
description?: string; description?: string;
/** The JSON schema of what it takes, as the module declared it. */ /** The JSON schema of what it takes, as the module declared it. */
input?: unknown; input?: unknown;
/** True for a role's tool: addressed to the seat, answered by whoever holds it (ADR 0132). */
seat?: boolean;
/** Where the tool is answered, as the mesh issued it (ADR 0160): the plain subject first when the
* module answers for itself anywhere, then one per machine. Absent for a runtime older than this. */
subjects?: string[];
}
/**
* 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. The mesh's own records are
* listed here as `mesh-controller (seat)` when the control plane did not answer. */
notAnswering: string[];
}
/** The seats and the verbs each declares, from the last listing, so a call can tell a role's tool
* from a module's when the two share a prefix (a module and a seat may share a name). */
export type Seats = Map<string, Set<string>>;
/** The roles' tools, keyed the way `toolKey` names them. */
export function seatsIn(have: Listing): Seats {
const seats: Seats = new Map();
for (const t of have.tools) {
if (!t.seat) continue;
if (!seats.has(t.module)) seats.set(t.module, new Set());
seats.get(t.module)!.add(t.name);
}
return seats;
}
/** The key a call uses for `<prefix>.<name>`: a role's when the prefix is a seat declaring that
* verb, a module's otherwise. Both names for one capability are deliberate and bounded (ADR 0132);
* the seat wins only for a verb it actually declares, so a module's own tool is never shadowed. */
export function toolKey(name: string, seats?: Seats): string {
if (name.startsWith("seat:")) return name;
const dot = name.indexOf(".");
if (dot < 0) return name;
const prefix = name.slice(0, dot);
const verb = name.slice(dot + 1);
if (seats?.get(prefix)?.has(verb)) return `seat:${prefix}.${verb}`;
return name;
} }
/** /**
@@ -72,50 +127,142 @@ 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;
} catch (e) {
throw new Error(`${path} is not a broker credential: ${(e as Error).message}`);
}
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",
); );
const tools = Array.isArray(answered) ? answered : (answered.tools ?? []); }
return tools return { bus: await connectNats(held), who: `${held.node ?? "?"}.${held.module}` };
.slice()
.sort((a: Tool, b: Tool) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
} }
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what their account /**
* 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, roles] = await Promise.all([
bus.request<Record<string, never>, { modules?: { module: string }[] }>(CATALOGUE_MODULES, {}),
// The roles' tools, from the mesh's records (design 33 §5). Asked beside the modules rather
// than first: a control plane that is restarting must not hide every module's tools with it.
bus
.request<Record<string, never>, { seats?: { seat: string; scope?: string; tools?: ToolsAnswer["tools"] }[] }>(
SEAT_TOOLS,
{},
)
.catch(() => undefined),
]);
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[] = [];
if (roles) {
for (const s of roles.seats ?? []) {
// A node-scoped seat's tool is asked of one machine, and the listing does not know which;
// those wait for a caller naming the node (`seat:<seat>.<verb>@<node>`).
if (s.scope === "node") continue;
for (const t of s.tools ?? []) {
tools.push({ module: s.seat, name: t.name, description: t.description, input: t.input, seat: true });
}
}
} else {
notAnswering.push("mesh-controller (seat)");
}
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, subjects: t.subjects });
}
} 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> { /** Call a tool and learn which machine answered (novox/hq ADR 0159). `<module>.<tool>@<node>` asks
if (!key.includes(".")) { * the instance on one machine; without it, whichever instance answers first does, and the answer
* says which. */
export async function callTool(
bus: Broker,
key: string,
args: unknown,
seats?: Seats,
listing?: Listing,
): Promise<Answered<unknown>> {
const at = key.indexOf("@");
const name = at < 0 ? key : key.slice(0, at);
const node = at < 0 ? "" : key.slice(at + 1);
if (!name.includes(".")) {
throw new Error( throw new Error(
`"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them`, `"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them, ` +
"or <module>.<tool>@<node> for the instance on one machine",
); );
} }
return bus.request<unknown, unknown>(key, args ?? {}); const resolved = toolKey(name, seats) + (node ? `@${node}` : "");
// Where the tool is answered is the module's to say and the mesh's to issue (ADR 0160): when the
// listing carried subjects for it, the call goes to one of those and composes nothing.
const on = subjectListed(name, node, listing);
const asking = bus as Broker & { ask?: <Req, Res>(k: string, b: Req, on?: string) => Promise<Answered<Res>> };
if (typeof asking.ask === "function") return asking.ask<unknown, unknown>(resolved, args ?? {}, on);
return { result: await bus.request<unknown, unknown>(resolved, args ?? {}) };
}
/** The subject the listing says answers `<module>.<tool>` — the machine's when one is named, else
* the plain one — or undefined when the listing said none, and the key is composed as before. */
export function subjectListed(name: string, node: string, listing?: Listing): string | undefined {
if (!listing) return undefined;
const dot = name.indexOf(".");
const module = name.slice(0, dot);
const tool = name.slice(dot + 1);
const found = listing.tools.find((t) => t.module === module && t.name === tool && !t.seat);
const subjects = found?.subjects ?? [];
if (subjects.length === 0) return undefined;
if (node) return subjects.find((s) => s.endsWith(`.${node}`));
return subjects[0];
} }
/** /**
* 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."; (key.startsWith("seat:") ? ", or nothing holds that seat" : "") +
" — `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
View File
@@ -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);
}
+5 -1
View File
@@ -29,6 +29,9 @@ import { readFileSync } from "node:fs";
import { pathToFileURL } from "node:url"; import { pathToFileURL } from "node:url";
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js"; import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js";
import { runTools } from "./runtime.js"; import { runTools } from "./runtime.js";
/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */
let lastCredential: Credential | undefined;
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";
import type { Broker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging";
@@ -51,6 +54,7 @@ async function connectBroker(): Promise<Broker> {
let credential: Credential; let credential: Credential;
try { try {
credential = JSON.parse(readFileSync(file, "utf8")) as Credential; credential = JSON.parse(readFileSync(file, "utf8")) as Credential;
lastCredential = credential;
} catch (err) { } catch (err) {
console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`); console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`);
process.exit(1); process.exit(1);
@@ -121,7 +125,7 @@ async function serve(): Promise<void> {
.filter(Boolean); .filter(Boolean);
const broker = await connectBrokerPatiently(); const broker = await connectBrokerPatiently();
const stop = await runTools({ broker, moduleEntrypoints }); const stop = await runTools({ broker, moduleEntrypoints, credential: lastCredential });
const shutdown = async (): Promise<void> => { const shutdown = async (): Promise<void> => {
stop(); stop();
+191 -89
View File
@@ -1,47 +1,203 @@
/** /**
* 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, seatsIn, toolKey, toolsOn, whyItFailed, type Listing, type Seats } 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;
};
// The roles the last listing knew, so `<seat>.<verb>` resolves to the seat. Fetched once if a
// call arrives before any list did; a listing that failed leaves no roles, and the name is then
// a module's, which is the right fallback for a mesh whose control plane is away.
const roles = async (): Promise<Seats | undefined> => {
try {
return seatsIn(await listing());
} catch {
return undefined;
}
};
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, plus every role's tools from the mesh's records — the mesh's own verbs ` +
`(mesh-controller.status, .push, .assign …) among them; a module that did not answer ` +
`is named in the list's _meta and can still be called by <module>.<tool>.`,
});
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 — with `node`, the machine to ask when
// the module runs on several (novox/hq ADR 0159); a seat's verb takes none, the
// seat's scope decides. An empty object is a tool that takes nothing, which is a
// real answer and not a missing one.
inputSchema: t.seat ? asSchema(t.input) : withNode(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 given = String(request.params?.name ?? "");
const args = { ...((request.params?.arguments as Record<string, unknown> | undefined) ?? {}) };
// The machine, when the caller names one, travels in the subject and never reaches the
// module's arguments (novox/hq ADR 0159) — for a module's tool. A seat's verb takes no
// machine from the console (the seat's scope decides), so a `node` among its arguments
// is the verb's own, as `push` and `assign` take one, and is handed through untouched.
const have = await listing().catch(() => undefined);
const roles = have && seatsIn(have);
const isSeatVerb = roles ? toolKey(given.split("@", 1)[0], roles).startsWith("seat:") : false;
const node = !isSeatVerb && typeof args.node === "string" && args.node !== "" ? args.node : "";
if (!isSeatVerb) delete args.node;
const name = node && !given.includes("@") ? `${given}@${node}` : given;
try {
const { result, node: answeredBy } = await callTool(bus, name, args, roles, have);
// 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. Which machine answered follows it as its own line.
const content: { type: string; text: string }[] = [
{ type: "text", text: JSON.stringify(result, null, 2) },
];
if (answeredBy) content.push({ type: "text", text: `answered by ${answeredBy}` });
return answer(request.id, { content });
} 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 +207,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}`);
}
}
} }
} }
@@ -137,3 +226,16 @@ async function* lines(): AsyncGenerator<string> {
} }
if (buffered.trim() !== "") yield buffered.trim(); if (buffered.trim() !== "") yield buffered.trim();
} }
/** Every module tool takes an optional `node`: the machine to ask when the module runs on several
* (novox/hq ADR 0159). Added to the listing, stripped before the call, never seen by the module. */
function withNode(schema: Record<string, unknown>): Record<string, unknown> {
const properties = { ...((schema.properties as Record<string, unknown> | undefined) ?? {}) };
if (!("node" in properties)) {
properties.node = {
type: "string",
description: "the machine to ask, when this module runs on several; else whichever answers, and the answer says which",
};
}
return { ...schema, type: "object", properties };
}
+187 -54
View File
@@ -1,63 +1,101 @@
#!/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,
seatsIn,
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,53 +105,78 @@ 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); const line = t.description ? `${name.padEnd(36)} ${t.description}` : name;
console.log(t.seat ? `${line} (a role's tool: answered by whoever holds the ${t.module} seat)` : line);
}
if (have.notAnswering.length > 0) {
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) { if (who) {
console.log(`\nthis is what the mesh has. What ${who} may call was fixed when the credential was issued.`); 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 { try {
parsed = JSON.parse(raw); // `<seat>.<verb>` reaches the role when the mesh lists that verb for the seat; `seat:` says so
} catch (e) { // outright and asks nothing first.
console.error(`the arguments are not JSON: ${(e as Error).message}`); const have = key.startsWith("seat:") ? undefined : await toolsOn(bus).catch(() => undefined);
return 1; const { result, node } = await callTool(bus, key, parsed, have && seatsIn(have), have);
} console.log(JSON.stringify(result, null, 2));
} if (node) console.error(`answered by ${node}`);
try {
const answer = await callTool(bus, key, parsed);
console.log(JSON.stringify(answer, null, 2));
return 0; return 0;
} catch (e) { } catch (e) {
console.error(whyItFailed(key, e)); console.error(whyItFailed(key, e));
@@ -121,6 +184,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. */
+133 -4
View File
@@ -6,14 +6,38 @@
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, toolKey } from "@novox/mesh-sdk/tools";
import type { Broker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging";
import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
/**
* 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>>;
/** Where this tool is answered, as the mesh issued it (ADR 0160): the module's plain subject
* first when there is one, then this machine's. A caller composes nothing. */
subjects?: string[];
}[];
}
export interface RuntimeOptions { export interface RuntimeOptions {
/** The mesh broker to serve over. */ /** The mesh broker to serve over. */
broker: Broker; broker: Broker;
/** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */ /** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */
moduleEntrypoints: string[]; moduleEntrypoints: string[];
/** The credential the mesh delivered, for what it says about the seats this module claims
* (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat. */
credential?: Credential;
} }
/** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */ /** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */
@@ -28,8 +52,113 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the // Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
// audit logger) registers none, and its scoped account may not declare the serve queue — so a // audit logger) registers none, and its scoped account may not declare the serve queue — so a
// 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(); // A registration under a seat's name is the module's implementation of that seat's verbs
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {}; // (ADR 0159, 0160): served on the seat's subjects by serveClaimedSeats, never as a module's
// tools and never listed among them. Everything else is the module's own.
// A module named like its seat (the catalogue is the mesh-catalog seat) registers once and is
// both: its tools are the module's and the seat's verbs alike.
const self = opts.credential?.module;
const seatNames = new Set((opts.credential?.claims ?? []).map((c) => c.seat));
// A registration under a name that is neither this module nor a seat it claims is not served:
// said, and left out, rather than fatal — on 2026-10-01 the credential of a module that had just
// learned to implement a seat did not yet name the claim, and the whole runtime restarted for it.
const ownRegistrations = collectTools().filter(({ module }) => {
if (module === self || !self || seatNames.has(module)) return module === self || !self;
console.log(`[mesh-tools] ${self} registers tools under "${module}", which is neither this module nor a seat its credential claims; not served until the mesh issues the claim`);
return false;
});
const tools = ownRegistrations.flatMap(({ module, tools: own }) => own.map((t) => ({ module, name: t.name })));
const stops: Array<() => void> = [];
const stop = (): void => stops.splice(0).forEach((s) => s());
// Each tool on its own key, namespaced by its module (ADR 0047); where that key is answered is
// the broker's to know from the membership (ADR 0160).
for (const { module, tools: own } of ownRegistrations) {
const seen = new Set<string>();
for (const t of own) {
if (seen.has(t.name)) {
stop();
throw new Error(`${module} exposes two tools named ${t.name} — refused`);
}
seen.add(t.name);
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) => t.run(args ?? {})));
}
}
// 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 runtime = opts.broker as RuntimeBroker;
for (const { module, tools: own } of ownRegistrations) {
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 subjectsOf = (tool: string): string[] | undefined => {
const issued = typeof runtime.membership === "function" ? runtime.membership() : undefined;
if (!issued) return undefined;
const plain = issued.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
const mine = issued.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
return [...plain, ...mine];
};
const answer: ToolsAnswer = {
module,
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
};
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; stops.push(await serveClaimedSeats(opts.broker as RuntimeBroker, opts.credential));
return () => {
for (const s of stops) s();
};
}
/**
* Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). The credential names the
* seats this module claims and the verbs each promises; each verb is served on the seat's own
* subject by the module's tool of the same name. Whether this instance *holds* the seat is the bus's
* to decide: only the holder's account may subscribe the seat's subjects, so a claimant that does not
* hold it here is refused the subscription and serves nothing — never a failure of its own tools.
*/
async function serveClaimedSeats(broker: RuntimeBroker, credential?: Credential): Promise<() => void> {
const claims = credential?.claims ?? [];
if (claims.length === 0 || typeof broker.handleSubject !== "function") return () => {};
// A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
for (const { module, tools } of collectTools()) {
if (!claims.some((c) => c.seat === module)) continue;
const verbs = new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
for (const t of tools) verbs.set(t.name, (args) => t.run(args));
implementations.set(module, verbs);
}
let stops: (() => void)[] = [];
const serve = async (): Promise<void> => {
stops.forEach((s) => s());
stops = [];
const issued = typeof broker.membership === "function" ? broker.membership() : undefined;
for (const claim of claims) {
const verbs = implementations.get(claim.seat);
for (const verb of claim.serves ?? []) {
const run = verbs?.get(verb);
if (!run) {
console.log(`[mesh-tools] claims ${claim.seat} and implements no ${verb}, which that seat promises; not served`);
continue;
}
// Where the mesh issued the verb when it has; the derived shape until then.
const subject = issued?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
stops.push(await broker.handleSubject(subject, run));
console.log(`[mesh-tools] serving ${claim.seat}'s ${verb} on ${subject}, admitted where this module holds the seat`);
}
}
};
await serve();
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
return () => stops.forEach((s) => s());
} }
+93 -13
View File
@@ -18,37 +18,62 @@ 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 }));
// And the mesh's own records, served by the holder of the mesh-controller seat (ADR 0154): one
// role's tool, and the verb that lists every role's.
const controller = await connectNats({ url: url!, module: "mesh-controller" });
await controller.handle("seat:mesh-controller.tools", async () => ({
seats: [
{ seat: "mesh-controller", scope: "mesh", tools: [{ name: "status", description: "what is wrong", input: {} }] },
{ seat: "node-dns-resolver", scope: "node", tools: [{ name: "lookup", description: "one machine's", input: {} }] },
],
}));
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
return { return {
async close() { async close() {
await catalogue.close(); await catalogue.close();
await shop.close(); await shop.close();
await controller.close();
}, },
}; };
} }
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", "mesh-controller.status", "shop.price"],
"the list is what the catalogue answered, in a stable order", "the list is what the modules answered plus every role's tools, in a stable order",
); );
// A role's tool is marked as one; a node-scoped seat's waits for a caller naming the node.
assert.ok(have.tools.find((x) => x.module === "mesh-controller")!.seat);
assert.ok(!have.tools.some((x) => x.module === "node-dns-resolver"));
// 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();
@@ -61,7 +86,7 @@ test("a person calls a tool and gets the module's own answer, unshaped", async (
const person = await connectNats({ url, module: "person.ada" }); const person = await connectNats({ url, module: "person.ada" });
try { try {
const answer = await callTool(person, "shop.price", { of: "a hat" }); const answer = await callTool(person, "shop.price", { of: "a hat" });
assert.deepEqual(answer, { of: "a hat", cost: 12 }); assert.deepEqual(answer.result, { of: "a hat", cost: 12 });
} finally { } finally {
await person.close(); await person.close();
await mesh.close(); await mesh.close();
@@ -103,3 +128,58 @@ test("each way a call fails says what to do about it", () => {
assert.match(whyItFailed("shop.price", new Error("timeout")), /did not answer in time/); assert.match(whyItFailed("shop.price", new Error("timeout")), /did not answer in time/);
assert.match(whyItFailed("shop.price", new Error("something else")), /something else/); assert.match(whyItFailed("shop.price", new Error("something else")), /something else/);
}); });
test("a role's tool is reached through the seat, and a module's own name is never shadowed", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const { seatsIn, toolKey } = await import("../dist/client.js");
const mesh = await aMeshWithTools();
const person = await connectNats({ url, module: "person.ada" });
try {
const roles = seatsIn(await toolsOn(person));
assert.equal(toolKey("mesh-controller.status", roles), "seat:mesh-controller.status");
assert.equal(toolKey("mesh-controller.other", roles), "mesh-controller.other", "a verb the seat does not declare is a module's");
assert.equal(toolKey("shop.price", roles), "shop.price");
const answer = await callTool(person, "mesh-controller.status", {}, roles);
assert.deepEqual(answer.result, { output: "all quiet", ok: true });
const direct = await callTool(person, "seat:mesh-controller.status", {});
assert.deepEqual(direct.result, { output: "all quiet", ok: true });
} finally {
await person.close();
await mesh.close();
}
});
// A module on two machines (novox/hq ADR 0159): a call names the machine and reaches that instance
// and no other; a call that names none reaches one of them and says which; a seat's verb is served
// by the claimant's tool of the same name on the seat's own subject.
test("a tool call names the machine it is for, and every answer says which machine answered", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const onAnchor = await connectNats({ url: url!, module: "store", node: "anchor" });
const onHome = await connectNats({ url: url!, module: "store", node: "home-server" });
const asker = await connectNats({ url: url!, module: "console", node: "workstation" });
try {
await onAnchor.handle("databases", async () => ({ at: "anchor" }));
await onHome.handle("databases", async () => ({ at: "home-server" }));
const home = await callTool(asker, "store.databases@home-server", {});
assert.deepEqual(home.result, { at: "home-server" });
assert.equal(home.node, "home-server");
const anchor = await callTool(asker, "store.databases@anchor", {});
assert.deepEqual(anchor.result, { at: "anchor" });
assert.equal(anchor.node, "anchor");
const whichever = await callTool(asker, "store.databases", {});
assert.ok(["anchor", "home-server"].includes(whichever.node ?? ""), `an unnamed call still says who answered: ${whichever.node}`);
assert.deepEqual(whichever.result, { at: whichever.node });
// A seat's verb, served by the claimant on the seat's subject; a mesh seat's is flat.
await onAnchor.handleSubject("mesh.seat.mesh-store.tool.databases", async () => ({ seat: "mesh-store", at: "anchor" }));
const viaSeat = await callTool(asker, "seat:mesh-store.databases", {});
assert.deepEqual(viaSeat.result, { seat: "mesh-store", at: "anchor" });
assert.equal(viaSeat.node, "anchor");
} finally {
await onAnchor.close();
await onHome.close();
await asker.close();
}
});
+6
View File
@@ -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 () => ({}) },
]);
+12
View File
@@ -0,0 +1,12 @@
// The shop again, in a file of its own: a module imported once stays imported, so a second test needs a second entrypoint.
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 }) },
]);
+12
View File
@@ -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 }) },
]);
+12
View File
@@ -0,0 +1,12 @@
// A module that is both software and a role: postgres's own tools under its name, and its
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("postgres", () => [
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
]);
registerModuleTools("mesh-store", () => [
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
]);
+12
View File
@@ -0,0 +1,12 @@
// A module that is both software and a role: postgres's own tools under its name, and its
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("postgres", () => [
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
]);
registerModuleTools("mesh-store", () => [
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
]);
+113
View File
@@ -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();
}
});
+53 -6
View File
@@ -20,13 +20,27 @@ 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 }));
const controller = await connectNats({ url: url!, module: "mesh-controller" });
await controller.handle("seat:mesh-controller.tools", async () => ({
seats: [{ seat: "mesh-controller", scope: "mesh", tools: [
{ name: "status", description: "what is wrong", input: {} },
{ name: "push", description: "tell a machine", input: { node: { type: "string" } } },
] }],
}));
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
await controller.handle("seat:mesh-controller.push", async (body: { node?: string }) => ({ told: body.node ?? "nobody" }));
t.after(async () => { t.after(async () => {
await catalogue.close(); await catalogue.close();
await shop.close(); await shop.close();
await controller.close();
}); });
const { mkdtemp, writeFile } = await import("node:fs/promises"); const { mkdtemp, writeFile } = await import("node:fs/promises");
@@ -80,14 +94,23 @@ 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");
const listed = byId.get(2)!.result.tools; const listed = byId.get(2)!.result.tools;
assert.equal(listed.length, 1); assert.deepEqual(listed.map((x: { name: string }) => x.name), ["mesh-controller.push", "mesh-controller.status", "shop.price"],
assert.equal(listed[0].name, "shop.price", "a tool is named the way a person names it"); "the modules' tools and the roles', named the way a person names them");
assert.ok(listed[0].inputSchema, "a tool with no schema is one an agent cannot call"); const price = listed[2];
assert.ok(price.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 — and
// `node`, the machine to ask when the module runs on several (novox/hq ADR 0159), beside them.
assert.deepEqual(price.inputSchema.properties.of, { type: "string" });
assert.equal(price.inputSchema.properties.node.type, "string", "a module's tool takes the machine to ask");
assert.ok(!listed[1].inputSchema.properties?.node, "a seat's verb takes no machine; the seat's scope decides");
assert.equal(listed[0].inputSchema.properties?.node?.type, "string", "a seat's verb that takes a node of its own keeps it");
// 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)}`);
@@ -122,3 +145,27 @@ test("a method this surface does not have is refused, and a notification is not"
assert.equal(replies[0].error.code, -32601); assert.equal(replies[0].error.code, -32601);
assert.match(replies[0].error.message, /resources\/list/); assert.match(replies[0].error.message, /resources\/list/);
}); });
// A seat's verb that takes a machine as its own argument — `push <node>` — keeps it: the console
// moves `node` into the subject for a module's tool only (ADR 0159), never for a role's verb.
test("a seat's verb keeps a node of its own; only a module's tool gives it to the subject", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const credential = await aMeshAndACredential(t);
const replies = await driving(credential, [
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "mesh-controller.push", arguments: { node: "anchor" } } },
]);
const result = replies[0].result;
assert.ok(!result.isError, JSON.stringify(replies[0]));
assert.deepEqual(JSON.parse(result.content[0].text), { told: "anchor" });
});
test("a host calls the mesh's own verb through the seat", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const credential = await aMeshAndACredential(t);
const replies = await driving(credential, [
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "mesh-controller.status", arguments: {} } },
]);
const result = replies[0].result;
assert.ok(!result.isError, JSON.stringify(replies[0]));
assert.deepEqual(JSON.parse(result.content[0].text), { output: "all quiet", ok: true });
});
+225
View File
@@ -0,0 +1,225 @@
/**
* The mesh issues an assignment's subjects, and a runtime serves what it is issued (novox/hq ADR
* 0160). A runtime derives one address for itself — `mesh.assignment.<node>.<module>` — reads the
* membership there, serves exactly what it says, and re-serves when a new one arrives. Against a
* real bus with JetStream, because the membership is a direct get on a stream.
*
* 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/membership.test.ts
*/
import assert from "node:assert/strict";
import { test } from "node:test";
import { fileURLToPath } from "node:url";
import { connect, StringCodec } from "nats";
import { resetTools } from "@novox/mesh-sdk/tools";
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
import { callTool, subjectListed } from "../dist/client.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));
const sc = StringCodec();
/** The ASSIGNMENTS stream as the controller asserts it: last-per-subject, readable by direct get. */
async function anAssignmentsStream() {
const nc = await connect({ servers: url! });
const jsm = await nc.jetstreamManager();
try {
await jsm.streams.delete("ASSIGNMENTS");
} catch {
// none yet
}
await jsm.streams.add({
name: "ASSIGNMENTS",
subjects: ["mesh.assignment.>"],
max_msgs_per_subject: 1,
allow_direct: true,
} as never);
return {
async issue(m: object & { node: string; module: string }) {
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
},
async close() {
await nc.close();
},
};
}
test("a runtime serves exactly the subjects it is issued, and the listing carries them", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const stream = await anAssignmentsStream();
// The mesh placed shop on two machines that are not interchangeable: each is served by name only.
await stream.issue({
node: "anchor",
module: "shop",
serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }],
emits: "mesh.mod.shop.event.{event}",
tools: "mesh.mod.shop.tool.tools",
});
const shop = await connectNats({ url, node: "anchor", module: "shop" });
const asker = await connectNats({ url, module: "console", node: "workstation" });
let stop = () => {};
try {
stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
assert.equal(shop.membership()?.module, "shop", "the runtime read what the mesh issued");
// By name it answers; the plain subject was not issued, so nothing serves it.
const named = await callTool(asker, "shop.price@anchor", { of: "a hat" });
assert.deepEqual(named.result, { of: "a hat", cost: 12 });
assert.equal(named.node, "anchor");
await assert.rejects(callTool(asker, "shop.price", {}), /no responders|503/i, "a subject not issued is not served");
// The tools answer says where each is served, and the client composes nothing.
const tools = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>(
"shop.tools",
{},
);
assert.deepEqual(tools.tools.find((x) => x.name === "price")?.subjects, ["mesh.mod.shop.tool.price.anchor"]);
const listing = { tools: [{ module: "shop", name: "price", subjects: ["mesh.mod.shop.tool.price.anchor"] }], notAnswering: [] };
assert.equal(subjectListed("shop.price", "anchor", listing as never), "mesh.mod.shop.tool.price.anchor");
assert.equal(subjectListed("shop.price", "elsewhere", listing as never), undefined);
// The mesh re-issues the membership with the plain subject too (the modules became
// interchangeable); the runtime re-serves on it without a restart.
await stream.issue({
node: "anchor",
module: "shop",
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
emits: "mesh.mod.shop.event.{event}",
tools: "mesh.mod.shop.tool.tools",
});
await new Promise((r) => setTimeout(r, 300));
const plain = await callTool(asker, "shop.price", { of: "a coat" });
assert.deepEqual(plain.result, { of: "a coat", cost: 12 });
assert.equal(plain.node, "anchor");
} finally {
stop();
await asker.close();
await shop.close();
await stream.close();
}
});
test("a seat's verbs are implemented under the seat's name, served where issued, never listed as the module's", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const stream = await anAssignmentsStream();
await stream.issue({
node: "anchor",
module: "postgres",
serves: [{ subject: "mesh.mod.postgres.tool.{tool}", queue: "serve.postgres" }, { subject: "mesh.mod.postgres.tool.{tool}.anchor" }],
seats: [{ seat: "mesh-store", verb: "databases", subject: "mesh.seat.mesh-store.tool.databases" }],
emits: "mesh.mod.postgres.event.{event}",
tools: "mesh.mod.postgres.tool.tools",
});
const credential = { url, node: "anchor", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] };
const pg = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
let stop = () => {};
try {
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat.mjs")] });
// The module's own `databases` and the seat's are two tools: postgres's on its subject, the
// store's on the seat's, each answering as itself.
const own = await callTool(asker, "postgres.databases", {});
assert.deepEqual(own.result, { software: "postgres" });
const seat = await callTool(asker, "seat:mesh-store.databases", {});
assert.deepEqual(seat.result, { seat: "mesh-store" });
assert.equal(seat.node, "anchor");
// What the seat does not promise — creating a database — is postgres's tool and not the store's.
await assert.rejects(callTool(asker, "seat:mesh-store.postgres_create_database", {}), /no responders|503/i);
// And the module's `tools` lists only postgres's own, never the seat's implementation.
const tools = await asker.request<Record<string, never>, { module: string; tools: { name: string }[] }>("postgres.tools", {});
assert.deepEqual(tools.tools.map((x) => x.name).sort(), ["databases", "postgres_create_database"]);
await assert.rejects(asker.request("mesh-store.tools", {}), /no responders|503/i, "a seat is not a module with a tools verb");
} finally {
stop();
await asker.close();
await pg.close();
await stream.close();
}
});
test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const stream = await anAssignmentsStream();
await stream.issue({
node: "anchor",
module: "shop",
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
seats: [{ seat: "shop", verb: "price", subject: "mesh.seat.shop.tool.price" }],
emits: "mesh.mod.shop.event.{event}",
tools: "mesh.mod.shop.tool.tools",
});
const credential = { url, node: "anchor", module: "shop", claims: [{ seat: "shop", scope: "mesh", serves: ["price"] }] };
const shop = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
let stop = () => {};
try {
stop = await runTools({ broker: shop, credential, moduleEntrypoints: [fixture("shop-seat.mjs")] });
assert.deepEqual((await callTool(asker, "shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
assert.deepEqual((await callTool(asker, "seat:shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
const tools = await asker.request<Record<string, never>, { tools: { name: string }[] }>("shop.tools", {});
assert.deepEqual(tools.tools.map((x) => x.name), ["price", "refund"]);
} finally {
stop();
await asker.close();
await shop.close();
await stream.close();
}
});
test("a mesh that has issued nothing yet gets the derived shape, and says so", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const stream = await anAssignmentsStream();
const said: string[] = [];
const log = console.log;
console.log = (...a: unknown[]) => said.push(a.join(" "));
let lone;
try {
lone = await connectNats({ url, node: "home-server", module: "lone" });
} finally {
console.log = log;
}
const asker = await connectNats({ url, module: "console", node: "workstation" });
try {
assert.equal(lone.membership(), undefined);
assert.ok(said.some((s) => /no membership issued for lone on home-server/.test(s)), said.join("\n"));
await lone.handle("ping", async () => ({ pong: true }));
assert.deepEqual((await callTool(asker, "lone.ping", {})).result, { pong: true });
assert.deepEqual((await callTool(asker, "lone.ping@home-server", {})).result, { pong: true });
} finally {
await asker.close();
await lone.close();
await stream.close();
}
});
// A module that implements a seat its credential does not (yet) claim is not served for it and does
// not fall over either: the runtime says so and serves the module's own tools.
test("a registration under a seat the credential does not claim is said and skipped, not fatal", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const stream = await anAssignmentsStream();
const credential = { url, node: "anchor", module: "postgres" };
const pg = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
const said: string[] = [];
const log = console.log;
console.log = (...a: unknown[]) => said.push(a.join(" "));
let stop = () => {};
try {
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat-unclaimed.mjs")] });
console.log = log;
assert.ok(said.some((s) => /registers tools under "mesh-store".*not served/.test(s)), said.join("\n"));
assert.deepEqual((await callTool(asker, "postgres.databases", {})).result, { software: "postgres" });
await assert.rejects(callTool(asker, "seat:mesh-store.databases", {}), /no responders|503/i);
} finally {
console.log = log;
stop();
await asker.close();
await pg.close();
await stream.close();
}
});
+60
View File
@@ -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();
}
});