1 Commits
Author SHA1 Message Date
jschoubben f04519e3e1 One consumer, one reader, however many patterns a module registers
A module has exactly one durable consumer, and each subscribe() started its own reader of it. Two
readers split the stream between them, and a reader that receives a message its own pattern does not
match acknowledges it — which is the right answer for a filter wider than anything registered, and
silent loss when the message was another handler's. The first module to subscribe twice would have
dropped roughly half of each kind of event with nothing reporting it.

Every registration is now dispatched from one reader, and a message is acknowledged once every handler
it is for has taken it.
2026-09-28 16:15:03 +02:00
19 changed files with 245 additions and 1771 deletions
-14
View File
@@ -25,20 +25,6 @@ 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
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
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the
+41 -233
View File
@@ -40,55 +40,6 @@ export interface Credential {
module?: string;
user?: 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
@@ -118,7 +69,7 @@ export function fatalBrokerReason(err: unknown): string | null {
export async function connectNats(
target: string | Credential,
opts: { module?: string } = {},
): Promise<RuntimeBroker> {
): Promise<Broker> {
const cred: Credential = typeof target === "string" ? { url: target } : target;
const self = cred.module ?? opts.module;
if (!self) {
@@ -151,170 +102,49 @@ export async function connectNats(
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
let closed = false;
const node = cred.node;
// The membership, read once at connect and followed. A direct get is one request on the
// 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.
let issued: Membership | undefined;
const issuedHandlers: ((m: Membership) => void)[] = [];
const subjectOfMine = node ? membershipSubject(node, self) : "";
if (subjectOfMine) {
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);
void (async () => {
for await (const msg of sub) {
let reply: { result?: Res; error?: string; node?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
if (node) reply.node = node;
msg.respond(sc.encode(JSON.stringify(reply)));
}
})();
return () => 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 {
/**
* 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> {
return (await ask<Req, Res>(key, body)).result;
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;
},
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.
* Answer a question.
*
* A queue group, so several nodes may serve one tool and exactly one of them answers each
* call.
*/
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
// 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);
const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` });
subs.push(sub);
void (async () => {
for await (const msg of sub) {
let reply: { result?: Res; error?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
msg.respond(sc.encode(JSON.stringify(reply)));
}
})();
return () => {
sub.unsubscribe();
};
},
/**
@@ -483,33 +313,11 @@ function eventSubject(type: string, self: string): string {
}
/** 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; `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). */
* another's, which is how a request reaches a module that is not this one. */
function toolSubject(key: string, self: string): string {
if (key.startsWith("seat:")) {
const rest = key.slice("seat:".length);
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;
const dot = key.indexOf(".");
if (dot < 0) return `mesh.mod.${self}.tool.${key}`;
return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`;
}
function normalizeFingerprint(fingerprint: string): string {
+35 -182
View File
@@ -1,91 +1,36 @@
/**
* The mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
* A person's client: the mesh's tools from a workstation (novox/hq design 25 §7).
*
* Two surfaces over one thing. A command line, for somebody at a terminal; an MCP server, for an
* agent. Both are adapters over the same three calls — what tools are there, what does this one take,
* call it — because a second way of reaching a tool is a second thing to keep correct.
*
* **It uses the same client a module's runtime uses.** Not a second protocol and not a bridge: the
* caller connects as its own bus user — a person's, or the console's — publishes on the tool subjects
* that account permits, and the server refuses anything else. So "what may this ask" is answered by the
* same permission list that answers it for a module, and there is nothing here for an audit to read
* separately.
* **It uses the same client a module's runtime uses.** Not a second protocol and not a bridge: a
* person connects as their own bus user, publishes on the tool subjects their account permits, and the
* server refuses anything else. So "what may this person do" is answered by the same permission list
* that answers it for a module, and there is nothing here for an audit to read separately.
*
* What the caller may NOT do is the more interesting half, and none of it is enforced here — it is the
* account (design 25 §4): it cannot publish an event, so it cannot claim a module said something; it
* has no consumer, so there is no delivery to acknowledge; and it cannot answer a request, so it cannot
* impersonate a module on a bus where anyone may serve a tool.
* What a person may NOT do is the more interesting half, and none of it is enforced here — it is the
* account (design 25 §4): they cannot publish an event, so they cannot claim a module said something;
* they have no consumer, so there is no delivery to acknowledge; and they cannot answer a request, so
* they cannot impersonate a module on a bus where anyone may serve a tool.
*/
import { readFile } from "node:fs/promises";
import type { Broker } from "@novox/mesh-sdk/messaging";
import { connectNats, type Answered, type Credential } from "./broker-nats.js";
import { TOOLS_VERB, type ToolsAnswer } from "./runtime.js";
import { connectNats, type Credential } from "./broker-nats.js";
/** Where the catalogue answers which modules the mesh holds. */
const CATALOGUE_MODULES = "mesh-catalog.catalog_modules";
/** Where the catalogue answers what tools the mesh has. */
const CATALOGUE_TOOLS = "mesh-catalog.catalog_tools";
/** 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. */
/** A tool as the catalogue describes one. */
export interface Tool {
/** The module that serves it — or, for a role's tool, the seat. */
module: string;
name: string;
description?: string;
/** The JSON schema of what it takes, as the module declared it. */
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;
}
/**
@@ -127,142 +72,50 @@ export async function connectAs(held: PersonCredential): Promise<Broker> {
}
/**
* 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
* and subjects derive from what the mesh authorised and from nothing in this process's environment.
*/
export async function connectAsTheConsole(path: string): Promise<{ bus: Broker; who: string }> {
const raw = await readFile(path, "utf8");
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",
);
}
return { bus: await connectNats(held), who: `${held.node ?? "?"}.${held.module}` };
}
/**
* What tools the mesh has, asked of the modules (design 34 §3).
* What tools the mesh has, asked of the catalogue.
*
* 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.
* **Asked, not configured.** The catalogue is the only thing that knows what is installed, and a
* client carrying its own list would be a list that goes stale the first time a module is assigned —
* silently, because a tool that is not offered looks exactly like a tool that does not exist.
*/
export async function toolsOn(bus: Broker): Promise<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}`, {})),
export async function toolsOn(bus: Broker): Promise<Tool[]> {
const answered = await bus.request<Record<string, never>, { tools?: Tool[] } | Tool[]>(
CATALOGUE_TOOLS,
{},
);
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 };
const tools = Array.isArray(answered) ? answered : (answered.tools ?? []);
return tools
.slice()
.sort((a: Tool, b: Tool) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
}
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what the account
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what their account
* permits — one vocabulary, so a refusal names the thing they asked for. */
/** Call a tool and learn which machine answered (novox/hq ADR 0159). `<module>.<tool>@<node>` asks
* 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(".")) {
export async function callTool(bus: Broker, key: string, args: unknown): Promise<unknown> {
if (!key.includes(".")) {
throw new Error(
`"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them, ` +
"or <module>.<tool>@<node> for the instance on one machine",
`"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them`,
);
}
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];
return bus.request<unknown, unknown>(key, args ?? {});
}
/**
* Why a call failed, said so that the remedy is in the words.
*
* Three answers a person actually gets, and they need different things done: nobody serves that tool,
* the mesh refused this account, or the tool itself failed. Without this they are one timeout and a
* the mesh refused this person, or the tool itself failed. Without this they are one timeout and a
* stack trace.
*/
export function whyItFailed(key: string, err: unknown): string {
const message = err instanceof Error ? err.message : String(err);
if (/no responders|503/i.test(message)) {
return `nothing serves ${key}. The module may not be assigned to any machine, or it is down` +
(key.startsWith("seat:") ? ", or nothing holds that seat" : "") +
" — `mesh tools` lists what answered.";
return `nothing serves ${key}. The module may not be assigned to any machine, or it is down — ` +
"`mesh tools` lists what the catalogue says is there.";
}
if (/permissions violation|authorization/i.test(message)) {
return `this account may not call ${key}. What it may call was fixed when it was issued — a ` +
"person's by `operator issue`, the console's by its manifest.";
return `this credential may not call ${key}. What it may call was fixed when it was issued; ` +
"`operator issue` again with the tool named, or ask somebody who can.";
}
if (/timeout/i.test(message)) {
return `${key} did not answer in time. Something is serving it, so this is the tool being slow ` +
-165
View File
@@ -1,165 +0,0 @@
/**
* 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);
}
+1 -47
View File
@@ -5,10 +5,6 @@
// starts consuming as it is imported, so this also runs consumers.
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
// and what an events test uses to put a message on the wire.
// mesh-tools prepare bring this module's state to the shape this version needs and exit
// — the runtime's answer to the word the mesh asks every module
// (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which
// the module's own image names beside MESH_TOOL_MODULES.
// mesh-tools run ENTRYPOINT run one compiled module entrypoint to completion and exit — the
// runtime side of a run-once step (novox/hq ADR 0052). It imports
// the given entrypoint, whose top-level code does its work — seed a
@@ -22,16 +18,12 @@
// amqps account scoped to this module. Preferred: a module holds its own.
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
// MESH_PREPARE /path/a,/path/b,… compiled entrypoints that prepare this module's state
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
import { readFileSync } from "node:fs";
import { pathToFileURL } from "node:url";
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.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 { useBroker } from "@novox/mesh-sdk/messaging";
import type { Broker } from "@novox/mesh-sdk/messaging";
@@ -54,7 +46,6 @@ async function connectBroker(): Promise<Broker> {
let credential: Credential;
try {
credential = JSON.parse(readFileSync(file, "utf8")) as Credential;
lastCredential = credential;
} catch (err) {
console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`);
process.exit(1);
@@ -125,7 +116,7 @@ async function serve(): Promise<void> {
.filter(Boolean);
const broker = await connectBrokerPatiently();
const stop = await runTools({ broker, moduleEntrypoints, credential: lastCredential });
const stop = await runTools({ broker, moduleEntrypoints });
const shutdown = async (): Promise<void> => {
stop();
@@ -186,49 +177,12 @@ async function runEntry(entrypoint: string): Promise<void> {
await import(pathToFileURL(entrypoint).href);
}
/**
* Bring this module's state to the shape this version needs, and exit — the runtime's answer to the
* one word the mesh asks every module (novox/hq ADR 0135).
*
* The entrypoints come from `MESH_PREPARE`, which a module's own image sets beside the entrypoints it
* already lists there: the module knows which of its files prepares its state, and nothing else could.
* Each is imported in the order given, to completion, with no broker — preparation runs before the
* version that would use it, so there is nothing yet to talk to.
*
* **An empty list is a failure, not a no-op.** The mesh only asks this of a module whose manifest says
* it prepares something; a module that says so and names nothing has been built wrong, and exiting 0
* would let that version serve against a state nobody shaped.
*/
async function prepareState(): Promise<void> {
const named = (process.env.MESH_PREPARE ?? "")
.split(",")
.map((entry) => entry.trim())
.filter((entry) => entry !== "");
if (named.length === 0) {
console.error(
"mesh-tools prepare: this module was asked to prepare its state and its image names nothing " +
"to do it with — set MESH_PREPARE to the compiled entrypoint(s) that prepare it, the way " +
"MESH_TOOL_MODULES names the ones it serves",
);
process.exit(1);
}
for (const entrypoint of named) {
console.log(`[mesh-tools] preparing with ${entrypoint}`);
await runEntry(entrypoint);
}
console.log(`[mesh-tools] prepared: ${named.length} entrypoint(s) ran to completion`);
}
async function main(): Promise<void> {
const [command, ...rest] = process.argv.slice(2);
if (command === "run") {
await runEntry(rest[0] ?? "");
return;
}
if (command === "prepare") {
await prepareState();
return;
}
if (command === "invoke") {
const [module, tool] = rest;
if (!module || !tool) {
+89 -191
View File
@@ -1,203 +1,47 @@
/**
* The mesh's tools as an MCP server (novox/hq design 25 §7, design 34).
* The mesh's tools as an MCP server, over stdio (novox/hq design 25 §7).
*
* **A thin adapter and nothing more.** Every tool an agent sees is one a module answered for and one
* this account may call; the schema is the module's own; the answer is the module's own. Nothing here
* decides anything, which is why it is short — an MCP surface that reshaped arguments or summarised
* answers would be a second definition of what a tool is, and the module's code is the first.
* **A thin adapter and nothing more.** Every tool an agent sees is one the catalogue listed and one
* this credential may call; the schema is the module's own; the answer is the module's own. Nothing
* here decides anything, which is why it is short — an MCP surface that reshaped arguments or
* summarised answers would be a second definition of what a tool is, and the module's manifest is the
* first.
*
* Implemented against the protocol directly rather than through a library: the surface is three
* methods and one framing, and a dependency here would be a dependency on every 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.
* methods and one framing, and a dependency here would be a dependency on every workstation.
*/
import type { Broker } from "@novox/mesh-sdk/messaging";
import { callTool, seatsIn, toolKey, toolsOn, whyItFailed, type Listing, type Seats } from "./client.js";
import { callTool, toolsOn, whyItFailed, type Tool } from "./client.js";
/** The protocol version this speaks. Stated, because a host that wants another should be told so
* rather than discovering it through a shape it did not expect. */
export const PROTOCOL = "2025-03-26";
const PROTOCOL = "2024-11-05";
export interface Request {
interface Request {
jsonrpc: string;
id?: number | string | null;
method: string;
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>;
}
/**
* The surface over one bus connection, as one account.
* Serve until stdin closes, which is how a host ends a session.
*
* The tool list is fetched when first asked and kept for a short while (design 34 §3): asking every
* module on every `tools/list` would fan out on every agent turn for something nobody changed, and
* 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.
* The tool list is fetched once, on the first `tools/list`, and kept. An agent asks for it repeatedly
* and the catalogue's answer does not change mid-session; refetching would make every turn cost a
* round trip to a module for something nobody changed.
*/
export async function serveMcp(bus: Broker, who: string): Promise<void> {
const surface = mcpSurface(bus, who);
let known: Tool[] | undefined;
const say = (message: unknown) => {
process.stdout.write(`${JSON.stringify(message)}\n`);
};
const answer = (id: Request["id"], result: unknown) => say({ jsonrpc: "2.0", id, result });
const refuse = (id: Request["id"], code: number, message: string) =>
say({ jsonrpc: "2.0", id, error: { code, message } });
for await (const line of lines()) {
let request: Request;
try {
@@ -207,8 +51,75 @@ export async function serveMcp(bus: Broker, who: string): Promise<void> {
// reply to a request that was never framed is noise on the same channel.
continue;
}
const reply = await surface.handle(request);
if (reply) say(reply);
// A notification has no id and expects no answer; `initialized` is the one every host sends.
const notification = request.id === undefined || request.id === null;
switch (request.method) {
case "initialize":
answer(request.id, {
protocolVersion: PROTOCOL,
capabilities: { tools: {} },
serverInfo: { name: "mesh", version: "1" },
// Said in the handshake, because an agent that knows whose authority it is acting under can
// say so when a call is refused — and a refusal is the one thing here that is not the
// mesh's fault or the tool's.
instructions:
`These are the tools of a Novox mesh, reached as ${who}. Every call goes to the module ` +
`that serves it; what may be called was fixed when this credential was issued, so a ` +
`refusal means the credential, not the tool.`,
});
break;
case "notifications/initialized":
break;
case "tools/list": {
try {
known ??= await toolsOn(bus);
} catch (e) {
refuse(request.id, -32603, whyItFailed("mesh-catalog.catalog_tools", e));
break;
}
answer(request.id, {
tools: known.map((t) => ({
name: `${t.module}.${t.name}`,
description: t.description ?? `${t.name}, served by ${t.module}`,
// The module's own schema, passed through. An empty object is a tool that takes nothing,
// which is a real answer and not a missing one.
inputSchema: t.input ?? { type: "object", properties: {} },
})),
});
break;
}
case "tools/call": {
const name = String(request.params?.name ?? "");
const args = request.params?.arguments ?? {};
try {
const result = await callTool(bus, name, args);
// Text, because that is what every host renders. The content is the module's answer as
// JSON, unshaped: an adapter that flattened it would be deciding what matters in somebody
// else's answer.
answer(request.id, {
content: [{ type: "text", text: JSON.stringify(result, null, 2) }],
});
} catch (e) {
// **An error the agent can act on, not a stack.** isError rather than a protocol failure,
// because the call was well-formed and the mesh answered it — with a refusal, an absence or
// a fault, and the words say which.
answer(request.id, {
content: [{ type: "text", text: whyItFailed(name, e) }],
isError: true,
});
}
break;
}
default:
if (!notification) {
refuse(request.id, -32601, `mesh's MCP surface has no ${request.method}`);
}
}
}
}
@@ -226,16 +137,3 @@ async function* lines(): AsyncGenerator<string> {
}
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 };
}
+56 -189
View File
@@ -1,101 +1,63 @@
#!/usr/bin/env node
/**
* `mesh` — the mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
* `mesh` — the mesh's tools from a workstation, for a person (novox/hq design 25 §7).
*
* Four verbs and nothing else. What tools are there, call one, serve the same two to an agent over
* stdio, and serve them on a machine's loopback as the console the mesh assigns. Deliberately thin:
* everything that could be a decision is one the mesh already made, and a client that grew opinions
* would be a second place the mesh's behaviour is defined.
* Three verbs and nothing else. What tools are there, call one, and serve the same two to an agent
* over MCP. Deliberately thin: everything that could be a decision is one the mesh already made, and a
* client that grew opinions would be a second place the mesh's behaviour is defined.
*
* mesh tools what the running modules answer
* mesh tools what this credential may call
* mesh call <module>.<tool> [json] call one, arguments as JSON on the command line or on stdin
* mesh mcp the same two, as an MCP server over stdio
* mesh serve [--listen host:port] the console: MCP over HTTP on this machine's loopback
* mesh mcp the same, as an MCP server over stdio
*
* 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
* The credential comes from MESH_CREDENTIAL, or --credential. It is the JSON `operator issue` printed.
*/
import { readFile } from "node:fs/promises";
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 { callTool, connectAs, credentialFrom, toolsOn, whyItFailed, type Tool } from "./client.js";
import { serveMcp } from "./mcp.js";
const usage = `mesh tools
mesh call <module>.<tool> [json]
mesh mcp
mesh serve [--listen host:port]
--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";
--credential <file> the JSON \`operator issue\` printed; default $MESH_CREDENTIAL`;
async function main(argv: string[]): Promise<number> {
const args = [...argv];
let credentialPath = process.env.MESH_CREDENTIAL ?? "";
let consoleUrl = process.env.MESH_CONSOLE ?? "";
let listen = process.env.MESH_CONSOLE_LISTEN ?? DEFAULT_LISTEN;
for (let i = 0; i < args.length; i++) {
const take = () => {
const v = args[i + 1] ?? "";
if (args[i] === "--credential") {
credentialPath = args[i + 1] ?? "";
args.splice(i, 2);
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();
if (!verb || verb === "help" || verb === "--help") {
console.log(usage);
return verb ? 0 : 1;
}
// Through a console already on this machine: no credential to hold, which is the point of one.
if (consoleUrl && (verb === "tools" || verb === "call")) {
return verb === "tools" ? listingVia(consoleUrl) : callingVia(consoleUrl, args);
if (!credentialPath) {
console.error(
"no credential: set MESH_CREDENTIAL or pass --credential <file>. It is the JSON " +
"`operator issue` printed, saved verbatim.",
);
return 1;
}
const { bus, who } = await connecting(credentialPath, verb);
const held = await credentialFrom(credentialPath);
const bus = await connectAs(held);
try {
switch (verb) {
case "tools":
return await listing(bus, who);
return await listing(bus, held.person);
case "call":
return await calling(bus, args);
case "mcp":
// Serves until stdin closes, which is how an MCP host ends a session.
await serveMcp(bus, who);
await serveMcp(bus, held.person ?? held.user ?? "somebody");
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:
console.error(`mesh has no "${verb}".\n\n${usage}`);
return 1;
@@ -105,78 +67,53 @@ async function main(argv: string[]): Promise<number> {
}
}
/**
* Who this process is on the bus. The module credential first: a console is started by the mesh with
* MESH_BROKER_FILE and nothing else, and must not fall back to a person's file lying around.
*/
async function connecting(credentialPath: string, verb: string): Promise<{ bus: Broker; who: string }> {
const delivered = process.env.MESH_BROKER_FILE;
if (delivered) {
return connectAsTheConsole(delivered);
}
if (!credentialPath) {
throw new Error(
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",
);
}
const held = await credentialFrom(credentialPath);
return { bus: await connectAs(held), who: held.person ?? held.user ?? "somebody" };
}
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 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) {
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;
async function listing(bus: Awaited<ReturnType<typeof connectAs>>, who?: string): Promise<number> {
let tools: Tool[];
try {
have = await toolsOn(bus);
tools = await toolsOn(bus);
} catch (e) {
console.error(whyItFailed("mesh-catalog.catalog_modules", e));
console.error(whyItFailed("mesh-catalog.catalog_tools", e));
return 1;
}
printListing(have, who);
if (tools.length === 0) {
console.log("the catalogue lists no tools; nothing on this mesh serves any");
return 0;
}
// **What the catalogue has, not what this credential may call.** The two differ and the difference
// is the point: a person seeing only their own tools cannot tell "not installed" from "not yours",
// and those need different people to fix them.
for (const t of tools) {
const name = `${t.module}.${t.name}`;
console.log(t.description ? `${name.padEnd(36)} ${t.description}` : name);
}
if (who) {
console.log(`\nthis is what the mesh has. What ${who} may call was fixed when the credential was issued.`);
}
return 0;
}
async function calling(bus: Broker, args: string[]): Promise<number> {
async function calling(
bus: Awaited<ReturnType<typeof connectAs>>,
args: string[],
): Promise<number> {
const key = args.shift();
if (!key) {
console.error("mesh call <module>.<tool> [json]");
return 1;
}
const parsed = await argumentsFrom(args);
if (parsed === undefined) return 1;
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
let parsed: unknown = {};
if (raw.trim() !== "") {
try {
parsed = JSON.parse(raw);
} catch (e) {
console.error(`the arguments are not JSON: ${(e as Error).message}`);
return 1;
}
}
try {
// `<seat>.<verb>` reaches the role when the mesh lists that verb for the seat; `seat:` says so
// outright and asks nothing first.
const have = key.startsWith("seat:") ? undefined : await toolsOn(bus).catch(() => undefined);
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}`);
const answer = await callTool(bus, key, parsed);
console.log(JSON.stringify(answer, null, 2));
return 0;
} catch (e) {
console.error(whyItFailed(key, e));
@@ -184,76 +121,6 @@ async function calling(bus: Broker, args: string[]): Promise<number> {
}
}
/** 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
* terminal, so `mesh call x.y` with no arguments does not hang waiting for something nobody is
* typing. */
+4 -133
View File
@@ -6,38 +6,14 @@
import { pathToFileURL } from "node:url";
import { resolve } from "node:path";
import { useBroker } from "@novox/mesh-sdk/messaging";
import { collectTools, toolKey } from "@novox/mesh-sdk/tools";
import { serveTools, listTools } from "@novox/mesh-sdk/tools";
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 {
/** The mesh broker to serve over. */
broker: Broker;
/** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */
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. */
@@ -52,113 +28,8 @@ 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
// 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.
// A registration under a seat's name is the module's implementation of that seat's verbs
// (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));
}
const tools = listTools();
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
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());
return stop;
}
+13 -93
View File
@@ -18,62 +18,37 @@ import { callTool, toolsOn, whyItFailed } from "../dist/client.js";
const url = process.env.MESH_TEST_NATS;
/** A mesh as discovery sees it (design 34 §3): the catalogue holding three modules, two of them up
* and answering `tools`, one held and not running. Two connections, because a person and a module
* are different users even in a test. */
/** A module serving the catalogue's tool list and one tool of its own, so the client has a mesh to
* talk to. Two connections, because a person and a module are different users even in a test. */
async function aMeshWithTools() {
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
const shop = await connectNats({ url: url!, module: "shop" });
await catalogue.handle("catalog_modules", async () => ({
modules: [{ module: "shop" }, { module: "mesh-catalog" }, { module: "ghost" }],
}));
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 }));
// 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 catalogue.handle("catalog_tools", async () => ({
tools: [
{ module: "shop", name: "price", description: "what something costs", input: { type: "object" } },
{ module: "mesh-catalog", name: "catalog_tools", description: "what tools the mesh has" },
],
}));
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
return {
async close() {
await catalogue.close();
await shop.close();
await controller.close();
},
};
}
test("a person sees what the running modules answer, sorted, and who did not answer", async (t) => {
test("a person sees what the catalogue says the mesh has, sorted", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const mesh = await aMeshWithTools();
const person = await connectNats({ url, module: "person.ada" });
try {
const began = Date.now();
const have = await toolsOn(person);
const tools = await toolsOn(person);
assert.deepEqual(
have.tools.map((x) => `${x.module}.${x.name}`),
["mesh-catalog.catalog_modules", "mesh-controller.status", "shop.price"],
"the list is what the modules answered plus every role's tools, in a stable order",
tools.map((x) => `${x.module}.${x.name}`),
["mesh-catalog.catalog_tools", "shop.price"],
"the list is what the catalogue answered, 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 {
await person.close();
await mesh.close();
@@ -86,7 +61,7 @@ test("a person calls a tool and gets the module's own answer, unshaped", async (
const person = await connectNats({ url, module: "person.ada" });
try {
const answer = await callTool(person, "shop.price", { of: "a hat" });
assert.deepEqual(answer.result, { of: "a hat", cost: 12 });
assert.deepEqual(answer, { of: "a hat", cost: 12 });
} finally {
await person.close();
await mesh.close();
@@ -128,58 +103,3 @@ 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("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
@@ -1,6 +0,0 @@
// 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
@@ -1,12 +0,0 @@
// 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
@@ -1,12 +0,0 @@
// 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
@@ -1,12 +0,0 @@
// 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
@@ -1,12 +0,0 @@
// 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
@@ -1,113 +0,0 @@
/**
* 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();
}
});
+6 -53
View File
@@ -20,27 +20,13 @@ const url = process.env.MESH_TEST_NATS;
async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void) => void }) {
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
const shop = await connectNats({ url: url!, module: "shop" });
await catalogue.handle("catalog_modules", async () => ({
modules: [{ module: "shop" }, { module: "ghost" }],
}));
await shop.handle("tools", async () => ({
module: "shop",
tools: [{ name: "price", description: "what something costs", input: { of: { type: "string" } } }],
await catalogue.handle("catalog_tools", async () => ({
tools: [{ module: "shop", name: "price", description: "what something costs" }],
}));
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
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 () => {
await catalogue.close();
await shop.close();
await controller.close();
});
const { mkdtemp, writeFile } = await import("node:fs/promises");
@@ -94,23 +80,14 @@ 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)}`);
const hello = byId.get(1)!.result;
assert.equal(hello.protocolVersion, "2025-03-26");
assert.equal(hello.protocolVersion, "2024-11-05");
assert.ok(hello.capabilities.tools, "a server offering no tools is not this one");
assert.match(hello.instructions, /ada/, "the handshake says whose authority a call is made under");
const listed = byId.get(2)!.result.tools;
assert.deepEqual(listed.map((x: { name: string }) => x.name), ["mesh-controller.push", "mesh-controller.status", "shop.price"],
"the modules' tools and the roles', named the way a person names them");
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"]);
assert.equal(listed.length, 1);
assert.equal(listed[0].name, "shop.price", "a tool is named the way a person names it");
assert.ok(listed[0].inputSchema, "a tool with no schema is one an agent cannot call");
const called = byId.get(3)!.result;
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
@@ -145,27 +122,3 @@ test("a method this surface does not have is refused, and a notification is not"
assert.equal(replies[0].error.code, -32601);
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
@@ -1,225 +0,0 @@
/**
* 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();
}
});
-19
View File
@@ -1,4 +1,3 @@
import { spawn } from "node:child_process";
import { test } from "node:test";
import assert from "node:assert/strict";
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
@@ -43,24 +42,6 @@ test("a non-Error value does not crash the classifier", () => {
assert.equal(fatalBrokerReason(undefined), null);
});
// **A module asked to prepare its state and naming nothing is a failure, not a no-op** (novox/hq
// ADR 0135). The mesh asks this only of a module whose manifest says it prepares something, so an
// image that names nothing was built wrong, and exiting 0 would let that version serve against a
// state nobody shaped.
test("preparing with nothing named fails rather than passing quietly", async () => {
const runtime = new URL("../dist/main.js", import.meta.url).pathname;
const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
const child = spawn(process.execPath, [runtime, "prepare"], {
env: { ...process.env, MESH_PREPARE: "" },
});
let said = "";
child.stderr.on("data", (chunk) => (said += String(chunk)));
child.on("close", (code) => resolve({ code, said }));
});
assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
assert.match(ran.said, /MESH_PREPARE/);
});
// **One durable consumer feeds one reader, however many patterns a module registers.**
//
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
-60
View File
@@ -1,60 +0,0 @@
/**
* 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();
}
});