Author SHA1 Message Date
jschoubben 6a91d144b3 Merge pull request 'The console lists and calls a role's tools' (#21) from feat/the-mesh-answers-for-itself into main
Reviewed-on: #21
2026-09-30 15:54:31 +00:00
jschoubben 71965ef958 The console lists and calls a role's tools
seat:<seat>.<verb> addresses a role's tool (with @<node> for a node-scoped seat); the listing asks the
mesh-controller seat's tools verb beside the modules and marks a role's tools; <seat>.<verb> resolves
to the seat when the seat declares that verb, a module's own name otherwise (novox/hq ADR 0154).
2026-09-30 17:41:56 +02:00
jschoubben dea98e509a Merge pull request 'The console: mesh serve on loopback, and every runtime answers tools' (#20) from feat/the-console into main
Reviewed-on: #20
2026-09-30 14:46:54 +00:00
jschoubben 80b02740ab The console: mesh serve on loopback, and every runtime answers tools
The runtime serves a tools verb per module with names, descriptions and schemas (design 34 §3), and
refuses a module naming its own tool tools. Discovery asks catalog_modules then each module, naming
what did not answer. One MCP handler over two transports: stdio (mesh mcp) and loopback HTTP (mesh
serve, the mesh-console module, novox/hq ADR 0152); serve refuses any bind but loopback. tools/call
may go through a running console with --console and no credential.
2026-09-30 16:19:22 +02:00
mesh-admin 621d033d53 Merge pull request 'One consumer, one reader, however many patterns a module registers' (#19) from fix/one-consumer-one-loop into main 2026-09-28 14:25:47 +00:00
jschoubben 38831c5c56 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:25:42 +02:00
mesh-admin fdad2f3268 Merge pull request 'A module answers the word the mesh asks: prepare' (#18) from feat/a-module-answers-prepare into main 2026-09-28 10:45:02 +00:00
jschoubben 81972a4995 A module answers the word the mesh asks: prepare
The runtime gains `prepare`, which brings this module's state to the shape this version needs and
exits (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which a module's own image names
beside the entrypoints it already lists there — the module knows which of its files prepares its
state and nothing else could. No broker is connected: preparation runs before the version that would
use it. An empty list fails rather than passing quietly, because the mesh asks this only of a module
whose manifest says it prepares something, and exiting 0 would let that version serve against a
state nobody shaped.
2026-09-28 12:45:00 +02:00
mesh-admin 10e8191717 Merge pull request 'The runtime hears on its own inbox' (#17) from fix/the-runtime-hears-on-its-own-inbox into main 2026-09-28 02:30:24 +00:00
jschoubben d703cebff4 The runtime hears on its own inbox
Every user's inbox is private to it and the grant names it; a reply space the client invented was
refused, and with it every pull for the next message and every answer to a tool call.
2026-09-28 04:30:22 +02:00
mesh-admin 46b56d53a6 Merge pull request 'The pin is the only check: the runtime stops verifying the bus's name' (#16) from fix/the-pin-is-the-only-check into main 2026-09-28 02:03:46 +00:00
15 changed files with 1132 additions and 217 deletions
+14
View File
@@ -25,6 +25,20 @@ MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool
`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables
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
+61 -20
View File
@@ -85,6 +85,10 @@ export async function connectNats(
pass: cred.password,
name: `${cred.node ?? "?"}.${self}`,
tls: cred.fingerprint ? await pinnedTls(cred.url, cred.fingerprint) : undefined,
// Its own inbox, not a random one: every user's inbox is private to it (design 25 §4), and the
// grant names `_INBOX.<user>.>` — a reply space the client invented would be refused, and with
// it every pull for the next message and every answer to a tool call.
inboxPrefix: cred.user ? `_INBOX.${cred.user}` : undefined,
// Reconnect forever: the bus being restarted is an upgrade, not a reason for every module on
// the mesh to exit. `close()` stays the only thing that ends the connection.
maxReconnectAttempts: -1,
@@ -92,6 +96,10 @@ export async function connectNats(
const js = conn.jetstream();
const subs: Subscription[] = [];
// Every registration, and the one reader that dispatches to them. A module has one durable
// consumer; the loop belongs to the connection rather than to a subscription.
const listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[] = [];
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
let closed = false;
return {
@@ -176,21 +184,38 @@ export async function connectNats(
* consumes (design 29 §3) — this binds to it and never creates one. A runtime that created
* its own would be a module deciding its own delivery semantics, and its account cannot
* reach the JetStream API to do it anyway.
*
* **One consumer, one loop, however many patterns a module registers.** A module has exactly one
* durable consumer, so two loops reading it would each take half the messages — and a loop that
* received one its own pattern does not match acknowledges it, which is the right answer for a
* filter wider than anything registered and silent loss when it is another handler's. Every
* registration is therefore dispatched from one reader, and a message is acknowledged once every
* handler it is for has taken it.
*/
async subscribe<T>(
pattern: string,
handler: (env: Envelope<T>) => Promise<void>,
): Promise<() => void> {
const durable = `${cred.node ?? "?"}_${self}`;
const consumer = await js.consumers.get("EVENTS", durable);
const messages = await consumer.consume();
void (async () => {
for await (const msg of messages) {
await deliver(msg, pattern, handler);
}
})();
const listener = { pattern, handler: handler as (env: Envelope<unknown>) => Promise<void> };
listeners.push(listener);
if (!reading) {
const durable = `${cred.node ?? "?"}_${self}`;
const consumer = await js.consumers.get("EVENTS", durable);
const messages = await consumer.consume();
reading = messages;
void (async () => {
for await (const msg of messages) {
await deliver(msg, listeners);
}
})();
}
return () => {
void messages.close();
const at = listeners.indexOf(listener);
if (at >= 0) listeners.splice(at, 1);
if (listeners.length === 0 && reading) {
void reading.close();
reading = undefined;
}
};
},
@@ -205,15 +230,20 @@ export async function connectNats(
};
}
/** Deliver one event, acknowledging only once a handler has taken it. */
async function deliver<T>(
/**
* Deliver one event to every handler it is for, acknowledging only once each has taken it.
*
* Several registrations share one durable consumer, so matching happens here rather than by having
* each registration read the stream: two readers of one consumer would split it between them, and a
* message that reached the wrong one would be acknowledged as not-for-me and lost.
*/
async function deliver(
msg: JsMsg,
pattern: string,
handler: (env: Envelope<T>) => Promise<void>,
listeners: { pattern: string; handler: (env: Envelope<unknown>) => Promise<void> }[],
): Promise<void> {
let env: Envelope<T>;
let env: Envelope<unknown>;
try {
env = toEnvelope<T>(msg);
env = toEnvelope<unknown>(msg);
} catch {
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
@@ -221,15 +251,16 @@ async function deliver<T>(
msg.term();
return;
}
if (!topicMatches(pattern, env.key)) {
// The consumer's filters are the controller's, and may be wider than one subscription's
// pattern when a module subscribes twice. Acknowledge what this handler is not for, or it
const forThis = listeners.filter((l) => topicMatches(l.pattern, env.key));
if (forThis.length === 0) {
// The consumer's filters are the controller's, derived from what the module declared it
// consumes, and may be wider than anything it registered a handler for. Acknowledge it, or it
// would be redelivered until it expired.
msg.ack();
return;
}
try {
await handler(env);
for (const l of forThis) await l.handler(env);
msg.ack();
} catch {
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
@@ -282,8 +313,18 @@ 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. */
* another's, which is how a request reaches a module that is not this one; `seat:<seat>.<verb>`
* addresses a role's tool, answered by whoever holds the seat (novox/hq ADR 0132) — with
* `seat:<seat>.<verb>@<node>` for a node-scoped seat, whose tool carries the machine (design 33 §4). */
function toolSubject(key: string, self: string): string {
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}`;
}
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)}`;
+145 -34
View File
@@ -1,36 +1,88 @@
/**
* A person's client: the mesh's tools from a workstation (novox/hq design 25 §7).
* The mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
*
* Two surfaces over one thing. A command line, for somebody at a terminal; an MCP server, for an
* agent. Both are adapters over the same three calls — what tools are there, what does this one take,
* call it — because a second way of reaching a tool is a second thing to keep correct.
*
* **It uses the same client a module's runtime uses.** Not a second protocol and not a bridge: a
* person connects as their own bus user, publishes on the tool subjects their account permits, and the
* server refuses anything else. So "what may this person do" is answered by the same permission list
* that answers it for a module, and there is nothing here for an audit to read separately.
* **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.
*
* 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.
* 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.
*/
import { readFile } from "node:fs/promises";
import type { Broker } from "@novox/mesh-sdk/messaging";
import { connectNats, type Credential } from "./broker-nats.js";
import { TOOLS_VERB, type ToolsAnswer } from "./runtime.js";
/** Where the catalogue answers what tools the mesh has. */
const CATALOGUE_TOOLS = "mesh-catalog.catalog_tools";
/** Where the catalogue answers which modules the mesh holds. */
const CATALOGUE_MODULES = "mesh-catalog.catalog_modules";
/** A tool as the catalogue describes one. */
/** Where the mesh answers every role's tools, from its records: the mesh-controller seat's own
* `tools` verb (novox/hq ADR 0154, design 33 §5). */
const SEAT_TOOLS = "seat:mesh-controller.tools";
/** A tool as its module describes it. */
export interface Tool {
/** 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;
}
/**
* What the mesh could say about its tools when asked (design 34 §3).
*
* **Silence is named, never dropped.** A module the catalogue holds and nothing answered for is in
* `notAnswering`, because a tool that is not offered looks exactly like a tool that does not exist,
* and those need different people to fix them.
*/
export interface Listing {
tools: Tool[];
/** Modules the catalogue holds whose runtime did not answer `tools`: not assigned, not up, or built
* before the runtime answered it. Each may still be called by name. The mesh's own records are
* listed here as `mesh-controller (seat)` when the control plane did not answer. */
notAnswering: string[];
}
/** The seats and the verbs each declares, from the last listing, so a call can tell a role's tool
* from a module's when the two share a prefix (a module and a seat may share a name). */
export type Seats = Map<string, Set<string>>;
/** The roles' tools, keyed the way `toolKey` names them. */
export function seatsIn(have: Listing): Seats {
const seats: Seats = new Map();
for (const t of have.tools) {
if (!t.seat) continue;
if (!seats.has(t.module)) seats.set(t.module, new Set());
seats.get(t.module)!.add(t.name);
}
return seats;
}
/** The key a call uses for `<prefix>.<name>`: a role's when the prefix is a seat declaring that
* verb, a module's otherwise. Both names for one capability are deliberate and bounded (ADR 0132);
* the seat wins only for a verb it actually declares, so a module's own tool is never shadowed. */
export function toolKey(name: string, seats?: Seats): string {
if (name.startsWith("seat:")) return name;
const dot = name.indexOf(".");
if (dot < 0) return name;
const prefix = name.slice(0, dot);
const verb = name.slice(dot + 1);
if (seats?.get(prefix)?.has(verb)) return `seat:${prefix}.${verb}`;
return name;
}
/**
@@ -72,50 +124,109 @@ export async function connectAs(held: PersonCredential): Promise<Broker> {
}
/**
* What tools the mesh has, asked of the catalogue.
*
* **Asked, not configured.** The catalogue is the only thing that knows what is installed, and a
* client carrying its own list would be a list that goes stale the first time a module is assigned —
* silently, because a tool that is not offered looks exactly like a tool that does not exist.
* 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 toolsOn(bus: Broker): Promise<Tool[]> {
const answered = await bus.request<Record<string, never>, { tools?: Tool[] } | Tool[]>(
CATALOGUE_TOOLS,
{},
);
const tools = Array.isArray(answered) ? answered : (answered.tools ?? []);
return tools
.slice()
.sort((a: Tool, b: Tool) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
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}` };
}
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what their account
/**
* What tools the mesh has, asked of the modules (design 34 §3).
*
* The catalogue says which modules the mesh holds; each module says what it serves, through the one
* verb its runtime answers for it. **Asked, not configured**: a client carrying its own list would be a
* list that goes stale the first time a module is assigned. Every module is asked at once, and the bus
* refuses at once a request nothing serves, so the cost is bounded by the modules that are up.
*/
export async function toolsOn(bus: Broker): Promise<Listing> {
const [answered, roles] = await Promise.all([
bus.request<Record<string, never>, { modules?: { module: string }[] }>(CATALOGUE_MODULES, {}),
// The roles' tools, from the mesh's records (design 33 §5). Asked beside the modules rather
// than first: a control plane that is restarting must not hide every module's tools with it.
bus
.request<Record<string, never>, { seats?: { seat: string; scope?: string; tools?: ToolsAnswer["tools"] }[] }>(
SEAT_TOOLS,
{},
)
.catch(() => undefined),
]);
const names = (answered.modules ?? []).map((m) => m.module).filter((m) => typeof m === "string");
const asked = await Promise.allSettled(
names.map((module) => bus.request<Record<string, never>, ToolsAnswer>(`${module}.${TOOLS_VERB}`, {})),
);
const tools: Tool[] = [];
const notAnswering: string[] = [];
if (roles) {
for (const s of roles.seats ?? []) {
// A node-scoped seat's tool is asked of one machine, and the listing does not know which;
// those wait for a caller naming the node (`seat:<seat>.<verb>@<node>`).
if (s.scope === "node") continue;
for (const t of s.tools ?? []) {
tools.push({ module: s.seat, name: t.name, description: t.description, input: t.input, seat: true });
}
}
} else {
notAnswering.push("mesh-controller (seat)");
}
asked.forEach((outcome, i) => {
const module = names[i]!;
if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
for (const t of outcome.value.tools) {
tools.push({ module, name: t.name, description: t.description, input: t.input });
}
} else {
notAnswering.push(module);
}
});
tools.sort((a, b) => `${a.module}.${a.name}`.localeCompare(`${b.module}.${b.name}`));
notAnswering.sort();
return { tools, notAnswering };
}
/** Call one tool. The key is `<module>.<tool>`, which is what a person types and what the account
* permits — one vocabulary, so a refusal names the thing they asked for. */
export async function callTool(bus: Broker, key: string, args: unknown): Promise<unknown> {
export async function callTool(bus: Broker, key: string, args: unknown, seats?: Seats): Promise<unknown> {
if (!key.includes(".")) {
throw new Error(
`"${key}" does not name a tool: write <module>.<tool>, as \`mesh tools\` lists them`,
);
}
return bus.request<unknown, unknown>(key, args ?? {});
return bus.request<unknown, unknown>(toolKey(key, seats), args ?? {});
}
/**
* Why a call failed, said so that the remedy is in the words.
*
* Three answers a person actually gets, and they need different things done: nobody serves that tool,
* the mesh refused this person, or the tool itself failed. Without this they are one timeout and a
* the mesh refused this account, or the tool itself failed. Without this they are one timeout and a
* stack trace.
*/
export function whyItFailed(key: string, err: unknown): string {
const message = err instanceof Error ? err.message : String(err);
if (/no responders|503/i.test(message)) {
return `nothing serves ${key}. The module may not be assigned to any machine, or it is down — ` +
"`mesh tools` lists what the catalogue says is there.";
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.";
}
if (/permissions violation|authorization/i.test(message)) {
return `this credential may not call ${key}. What it may call was fixed when it was issued; ` +
"`operator issue` again with the tool named, or ask somebody who can.";
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.";
}
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
@@ -0,0 +1,165 @@
/**
* The console's endpoint: MCP over HTTP, on a machine's loopback (novox/hq ADR 0152, design 34 §2).
*
* **Loopback is the authority boundary.** Whoever can connect is on the machine, and whoever is on the
* machine is the account that owns the mesh there (ADR 0034, ADR 0144). So there is no token and no
* login here, and the one thing this file enforces is that it binds nothing else: a console reachable
* from another machine would be authority over the mesh handed to whoever finds the port.
*
* The transport is the streamable-HTTP shape an agent host speaks: `POST /mcp` with one JSON-RPC
* message, answered with one JSON body. No session, because the surface holds nothing per caller; no
* event stream, because nothing here has anything to say unasked.
*/
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
import type { Broker } from "@novox/mesh-sdk/messaging";
import { mcpSurface, type Reply, type Request } from "./mcp.js";
/** The most a request body may be. A tool's arguments are small; a megabyte is somebody else's file. */
const BODY_LIMIT = 1 << 20;
export interface Listening {
/** Where it listens, as `host:port`, with the port the machine actually gave. */
address: string;
close(): Promise<void>;
}
/** Hosts that are this machine and no other. */
const loopback = new Set(["127.0.0.1", "::1", "localhost", "[::1]"]);
/**
* Listen on `host:port`. Refused unless the host is loopback — said before binding, so a manifest or
* a flag that would open the console to a network is a startup failure rather than something
* discovered by whoever finds it.
*/
export async function serveMcpHttp(bus: Broker, who: string, listen: string): Promise<Listening> {
const at = listen.lastIndexOf(":");
if (at < 0) {
throw new Error(`"${listen}" is not host:port`);
}
const host = listen.slice(0, at);
const port = Number(listen.slice(at + 1));
if (!loopback.has(host)) {
throw new Error(
`the console listens on loopback and nowhere else (novox/hq ADR 0152): "${host}" is not this ` +
"machine's own address — whoever is on the machine owns the mesh there, and nobody else may reach this",
);
}
if (!Number.isInteger(port) || port < 0 || port > 65535) {
throw new Error(`"${listen.slice(at + 1)}" is not a port`);
}
const surface = mcpSurface(bus, who);
const server = createServer((req, res) => {
void route(req, res, surface.handle).catch((e) => {
json(res, 500, { jsonrpc: "2.0", id: null, error: { code: -32603, message: String(e) } });
});
});
await new Promise<void>((resolve, reject) => {
server.once("error", reject);
server.listen(port, host.replace(/^\[|\]$/g, ""), () => resolve());
});
const bound = server.address();
const address = typeof bound === "object" && bound ? `${host}:${bound.port}` : listen;
return {
address,
close: () =>
new Promise<void>((resolve) => {
server.close(() => resolve());
}),
};
}
async function route(
req: IncomingMessage,
res: ServerResponse,
handle: (r: Request) => Promise<Reply | undefined>,
): Promise<void> {
const path = (req.url ?? "/").split("?")[0];
if (path === "/") {
res.writeHead(200, { "content-type": "text/plain; charset=utf-8" });
res.end("the mesh's console: MCP over HTTP at POST /mcp (novox/hq design 34)\n");
return;
}
if (path !== "/mcp") {
json(res, 404, { error: "the console serves /mcp and nothing else" });
return;
}
switch (req.method) {
case "POST":
break;
case "DELETE":
// A host ending a session. There is no session to end; saying so is the truthful answer.
res.writeHead(204).end();
return;
case "GET":
// A host opening an event stream. The console has nothing to say unasked.
res.writeHead(405, { allow: "POST, DELETE" }).end();
return;
default:
res.writeHead(405, { allow: "POST, DELETE" }).end();
return;
}
let body: string;
try {
body = await read(req);
} catch (e) {
json(res, 413, { jsonrpc: "2.0", id: null, error: { code: -32600, message: String(e) } });
return;
}
let parsed: unknown;
try {
parsed = JSON.parse(body);
} catch {
json(res, 400, { jsonrpc: "2.0", id: null, error: { code: -32700, message: "the body is not JSON" } });
return;
}
// One message, or a batch of them; a batch is answered as a batch. A notification gets no reply
// and, alone, no body: 202 is how the transport says "heard".
if (Array.isArray(parsed)) {
const replies = (await Promise.all(parsed.map((r) => handle(r as Request)))).filter(Boolean);
if (replies.length === 0) {
res.writeHead(202).end();
} else {
json(res, 200, replies);
}
return;
}
const reply = await handle(parsed as Request);
if (!reply) {
res.writeHead(202).end();
return;
}
json(res, 200, reply);
}
function read(req: IncomingMessage): Promise<string> {
return new Promise((resolve, reject) => {
let size = 0;
const chunks: Buffer[] = [];
req.on("data", (chunk: Buffer) => {
size += chunk.length;
if (size > BODY_LIMIT) {
reject(new Error(`the request is larger than ${BODY_LIMIT} bytes`));
req.destroy();
return;
}
chunks.push(chunk);
});
req.on("end", () => resolve(Buffer.concat(chunks).toString("utf8")));
req.on("error", reject);
});
}
function json(res: ServerResponse, status: number, body: unknown): void {
const text = JSON.stringify(body);
res.writeHead(status, {
"content-type": "application/json; charset=utf-8",
"content-length": Buffer.byteLength(text),
});
res.end(text);
}
+42
View File
@@ -5,6 +5,10 @@
// 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
@@ -18,6 +22,7 @@
// 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";
@@ -177,12 +182,49 @@ 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) {
+164 -89
View File
@@ -1,47 +1,189 @@
/**
* The mesh's tools as an MCP server, over stdio (novox/hq design 25 §7).
* The mesh's tools as an MCP server (novox/hq design 25 §7, design 34).
*
* **A thin adapter and nothing more.** Every tool an agent sees is one the catalogue listed and one
* 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.
* **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.
*
* Implemented against the protocol directly rather than through a library: the surface is three
* methods and one framing, and a dependency here would be a dependency on every workstation.
* methods and one framing, and a dependency here would be a dependency on every machine.
*
* One handler, two transports. Over stdio for a program a person starts (`mesh mcp`), over HTTP on a
* machine's loopback for the console the mesh assigns there (`mesh serve`, http.ts). The handler does
* not know which asked.
*/
import type { Broker } from "@novox/mesh-sdk/messaging";
import { callTool, toolsOn, whyItFailed, type Tool } from "./client.js";
import { callTool, seatsIn, toolsOn, whyItFailed, type Listing, type Seats } from "./client.js";
/** The protocol version this speaks. Stated, because a host that wants another should be told so
* rather than discovering it through a shape it did not expect. */
const PROTOCOL = "2024-11-05";
export const PROTOCOL = "2025-03-26";
interface Request {
export interface Request {
jsonrpc: string;
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>;
}
/**
* Serve until stdin closes, which is how a host ends a session.
* The surface over one bus connection, as one account.
*
* The tool list is fetched once, on the first `tools/list`, and kept. An agent asks for it repeatedly
* 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.
* 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. An empty object is a tool that takes
// nothing, which is a real answer and not a missing one.
inputSchema: asSchema(t.input),
})),
// Silence, named (design 34 §3): the modules the catalogue holds and nothing answered
// for. Not a tool, so not in `tools`; not dropped either.
_meta: { notAnswering: have.notAnswering },
});
}
case "tools/call": {
const name = String(request.params?.name ?? "");
const args = request.params?.arguments ?? {};
try {
const result = await callTool(bus, name, args, await roles());
// Text, because that is what every host renders. The content is the module's answer
// as JSON, unshaped: an adapter that flattened it would be deciding what matters in
// somebody else's answer.
return answer(request.id, {
content: [{ type: "text", text: JSON.stringify(result, null, 2) }],
});
} catch (e) {
// **An error the agent can act on, not a stack.** isError rather than a protocol
// failure, because the call was well-formed and the mesh answered it — with a refusal,
// an absence or a fault, and the words say which.
return answer(request.id, {
content: [{ type: "text", text: whyItFailed(name, e) }],
isError: true,
});
}
}
default:
return notification ? undefined : refuse(request.id, -32601, `mesh's MCP surface has no ${request.method}`);
}
},
};
}
/**
* A module's declared input as a JSON schema an agent can read.
*
* The sdk keeps a tool's input opaque, and the catalogue's modules write it as a bare map of
* property to description — `{ module: { type, description } }` — which is the `properties` of a
* schema rather than a schema. Wrapped here when that is what arrived; passed through when a module
* already wrote a schema; an empty object when it declared nothing. The module's words are kept
* either way.
*/
export function asSchema(input: unknown): Record<string, unknown> {
if (!input || typeof input !== "object" || Array.isArray(input)) {
return { type: "object", properties: {} };
}
const given = input as Record<string, unknown>;
if (given.type === "object" || "properties" in given) return given;
if (Object.keys(given).length === 0) return { type: "object", properties: {} };
return { type: "object", properties: given };
}
/**
* Serve over stdio until stdin closes, which is how a host ends a session.
*/
export async function serveMcp(bus: Broker, who: string): Promise<void> {
let known: Tool[] | undefined;
const surface = mcpSurface(bus, who);
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 {
@@ -51,75 +193,8 @@ export async function serveMcp(bus: Broker, who: string): Promise<void> {
// reply to a request that was never framed is noise on the same channel.
continue;
}
// A notification has no id and expects no answer; `initialized` is the one every host sends.
const notification = request.id === undefined || request.id === null;
switch (request.method) {
case "initialize":
answer(request.id, {
protocolVersion: PROTOCOL,
capabilities: { tools: {} },
serverInfo: { name: "mesh", version: "1" },
// Said in the handshake, because an agent that knows whose authority it is acting under can
// say so when a call is refused — and a refusal is the one thing here that is not the
// mesh's fault or the tool's.
instructions:
`These are the tools of a Novox mesh, reached as ${who}. Every call goes to the module ` +
`that serves it; what may be called was fixed when this credential was issued, so a ` +
`refusal means the credential, not the tool.`,
});
break;
case "notifications/initialized":
break;
case "tools/list": {
try {
known ??= await toolsOn(bus);
} catch (e) {
refuse(request.id, -32603, whyItFailed("mesh-catalog.catalog_tools", e));
break;
}
answer(request.id, {
tools: known.map((t) => ({
name: `${t.module}.${t.name}`,
description: t.description ?? `${t.name}, served by ${t.module}`,
// The module's own schema, passed through. An empty object is a tool that takes nothing,
// which is a real answer and not a missing one.
inputSchema: t.input ?? { type: "object", properties: {} },
})),
});
break;
}
case "tools/call": {
const name = String(request.params?.name ?? "");
const args = request.params?.arguments ?? {};
try {
const result = await callTool(bus, name, args);
// Text, because that is what every host renders. The content is the module's answer as
// JSON, unshaped: an adapter that flattened it would be deciding what matters in somebody
// else's answer.
answer(request.id, {
content: [{ type: "text", text: JSON.stringify(result, null, 2) }],
});
} catch (e) {
// **An error the agent can act on, not a stack.** isError rather than a protocol failure,
// because the call was well-formed and the mesh answered it — with a refusal, an absence or
// a fault, and the words say which.
answer(request.id, {
content: [{ type: "text", text: whyItFailed(name, e) }],
isError: true,
});
}
break;
}
default:
if (!notification) {
refuse(request.id, -32601, `mesh's MCP surface has no ${request.method}`);
}
}
const reply = await surface.handle(request);
if (reply) say(reply);
}
}
+185 -53
View File
@@ -1,63 +1,101 @@
#!/usr/bin/env node
/**
* `mesh` — the mesh's tools from a workstation, for a person (novox/hq design 25 §7).
* `mesh` — the mesh's tools, for whoever is on a machine (novox/hq design 25 §7, design 34).
*
* Three verbs and nothing else. What tools are there, call one, and serve the same two to an agent
* 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.
* 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.
*
* mesh tools what this credential may call
* mesh tools what the running modules answer
* mesh call <module>.<tool> [json] call one, arguments as JSON on the command line or on stdin
* mesh mcp the same, as an MCP server over stdio
* mesh mcp the same two, as an MCP server over stdio
* mesh serve [--listen host:port] the console: MCP over HTTP on this machine's loopback
*
* The credential comes from MESH_CREDENTIAL, or --credential. It is the JSON `operator issue` printed.
* Who it speaks as, in order of preference:
* MESH_BROKER_FILE the module credential the mesh delivered — the console's (ADR 0152)
* MESH_CREDENTIAL / --credential <file> a person's, as `operator issue` printed it (design 25 §7)
* MESH_CONSOLE / --console <url> no credential: `tools` and `call` go through a console
* already running on this machine, over loopback HTTP
*/
import { readFile } from "node:fs/promises";
import { callTool, connectAs, credentialFrom, toolsOn, whyItFailed, type Tool } from "./client.js";
import type { Broker } from "@novox/mesh-sdk/messaging";
import {
callTool,
connectAs,
connectAsTheConsole,
credentialFrom,
seatsIn,
toolsOn,
whyItFailed,
type Listing,
} from "./client.js";
import { serveMcpHttp } from "./http.js";
import { serveMcp } from "./mcp.js";
const usage = `mesh tools
mesh call <module>.<tool> [json]
mesh mcp
mesh serve [--listen host:port]
--credential <file> the JSON \`operator issue\` printed; default $MESH_CREDENTIAL`;
--credential <file> a person's credential, the JSON \`operator issue\` printed; default $MESH_CREDENTIAL
--console <url> a console on this machine to ask through instead; default $MESH_CONSOLE
--listen <host:port> where \`serve\` listens; loopback only; default $MESH_CONSOLE_LISTEN or 127.0.0.1:4270
MESH_BROKER_FILE the module credential the mesh delivered, which \`serve\` runs on`;
/** What the console listens on when nothing says otherwise. */
const DEFAULT_LISTEN = "127.0.0.1:4270";
async function main(argv: string[]): Promise<number> {
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++) {
if (args[i] === "--credential") {
credentialPath = args[i + 1] ?? "";
const take = () => {
const v = 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;
}
if (!credentialPath) {
console.error(
"no credential: set MESH_CREDENTIAL or pass --credential <file>. It is the JSON " +
"`operator issue` printed, saved verbatim.",
);
return 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);
}
const held = await credentialFrom(credentialPath);
const bus = await connectAs(held);
const { bus, who } = await connecting(credentialPath, verb);
try {
switch (verb) {
case "tools":
return await listing(bus, held.person);
return await listing(bus, who);
case "call":
return await calling(bus, args);
case "mcp":
// Serves until stdin closes, which is how an MCP host ends a session.
await serveMcp(bus, held.person ?? held.user ?? "somebody");
await serveMcp(bus, who);
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;
@@ -67,52 +105,76 @@ async function main(argv: string[]): Promise<number> {
}
}
async function listing(bus: Awaited<ReturnType<typeof connectAs>>, who?: string): Promise<number> {
let tools: Tool[];
try {
tools = await toolsOn(bus);
} catch (e) {
console.error(whyItFailed("mesh-catalog.catalog_tools", e));
return 1;
/**
* 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 (tools.length === 0) {
console.log("the catalogue lists no tools; nothing on this mesh serves any");
return 0;
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",
);
}
// **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 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}`;
console.log(t.description ? `${name.padEnd(36)} ${t.description}` : name);
const line = t.description ? `${name.padEnd(36)} ${t.description}` : name;
console.log(t.seat ? `${line} (a role's tool: answered by whoever holds the ${t.module} seat)` : line);
}
if (have.notAnswering.length > 0) {
console.log(
`\nheld by the mesh and not answering: ${have.notAnswering.join(", ")} — not assigned, not up, ` +
"or built before the runtime answered `tools`; each can still be called by name",
);
}
if (who) {
console.log(`\nthis is what the mesh has. What ${who} may call was fixed when the credential was issued.`);
console.log(`\nasked as ${who}; what ${who} may call was fixed when the account was issued.`);
}
}
async function listing(bus: Broker, who?: string): Promise<number> {
let have: Listing;
try {
have = await toolsOn(bus);
} catch (e) {
console.error(whyItFailed("mesh-catalog.catalog_modules", e));
return 1;
}
printListing(have, who);
return 0;
}
async function calling(
bus: Awaited<ReturnType<typeof connectAs>>,
args: string[],
): Promise<number> {
async function calling(bus: Broker, args: string[]): Promise<number> {
const key = args.shift();
if (!key) {
console.error("mesh call <module>.<tool> [json]");
return 1;
}
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
let parsed: unknown = {};
if (raw.trim() !== "") {
try {
parsed = JSON.parse(raw);
} catch (e) {
console.error(`the arguments are not JSON: ${(e as Error).message}`);
return 1;
}
}
const parsed = await argumentsFrom(args);
if (parsed === undefined) return 1;
try {
const answer = await callTool(bus, key, parsed);
// `<seat>.<verb>` reaches the role when the mesh lists that verb for the seat; `seat:` says so
// outright and asks nothing first.
const roles = key.startsWith("seat:") ? undefined : seatsIn(await toolsOn(bus).catch(() => ({ tools: [], notAnswering: [] })));
const answer = await callTool(bus, key, parsed, roles);
console.log(JSON.stringify(answer, null, 2));
return 0;
} catch (e) {
@@ -121,6 +183,76 @@ async function calling(
}
}
/** The console's answer to one MCP request, over loopback HTTP. */
async function viaConsole(consoleUrl: string, method: string, params?: unknown): Promise<any> {
const endpoint = consoleUrl.endsWith("/mcp") ? consoleUrl : `${consoleUrl.replace(/\/$/, "")}/mcp`;
const res = await fetch(endpoint, {
method: "POST",
headers: { "content-type": "application/json", accept: "application/json" },
body: JSON.stringify({ jsonrpc: "2.0", id: 1, method, params }),
});
if (!res.ok) {
throw new Error(`the console at ${endpoint} answered ${res.status}`);
}
const reply = (await res.json()) as { result?: any; error?: { message: string } };
if (reply.error) throw new Error(reply.error.message);
return reply.result;
}
async function listingVia(consoleUrl: string): Promise<number> {
try {
const result = await viaConsole(consoleUrl, "tools/list");
const have: Listing = {
tools: (result.tools ?? []).map((t: { name: string; description?: string; inputSchema?: unknown }) => {
const at = t.name.indexOf(".");
return { module: t.name.slice(0, at), name: t.name.slice(at + 1), description: t.description, input: t.inputSchema };
}),
notAnswering: result._meta?.notAnswering ?? [],
};
printListing(have);
return 0;
} catch (e) {
console.error(e instanceof Error ? e.message : String(e));
return 1;
}
}
async function callingVia(consoleUrl: string, args: string[]): Promise<number> {
const key = args.shift();
if (!key) {
console.error("mesh call <module>.<tool> [json]");
return 1;
}
const parsed = await argumentsFrom(args);
if (parsed === undefined) return 1;
try {
const result = await viaConsole(consoleUrl, "tools/call", { name: key, arguments: parsed });
const text = result?.content?.[0]?.text ?? JSON.stringify(result);
if (result?.isError) {
console.error(text);
return 1;
}
console.log(text);
return 0;
} catch (e) {
console.error(e instanceof Error ? e.message : String(e));
return 1;
}
}
/** A call's arguments: JSON on the command line, else on stdin, else nothing. Undefined when what
* was given is not JSON, after saying so. */
async function argumentsFrom(args: string[]): Promise<unknown> {
const raw = args.length > 0 ? args.join(" ") : await maybeStdin();
if (raw.trim() === "") return {};
try {
return JSON.parse(raw);
} catch (e) {
console.error(`the arguments are not JSON: ${(e as Error).message}`);
return undefined;
}
}
/** Arguments on stdin, for a call whose JSON is too long or too quoted to type. Empty when stdin is a
* terminal, so `mesh call x.y` with no arguments does not hang waiting for something nobody is
* typing. */
+38 -2
View File
@@ -6,9 +6,22 @@
import { pathToFileURL } from "node:url";
import { resolve } from "node:path";
import { useBroker } from "@novox/mesh-sdk/messaging";
import { serveTools, listTools } from "@novox/mesh-sdk/tools";
import { collectTools, serveTools, listTools, toolKey } from "@novox/mesh-sdk/tools";
import type { Broker } from "@novox/mesh-sdk/messaging";
/**
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
* module's tool names, descriptions and argument schemas, from the code that answers them and
* from nowhere else. Discovery asks the module, because a copy kept anywhere else drifts.
*/
export const TOOLS_VERB = "tools";
/** What `tools` answers for one module. */
export interface ToolsAnswer {
module: string;
tools: { name: string; description: string; input: Readonly<Record<string, unknown>> }[];
}
export interface RuntimeOptions {
/** The mesh broker to serve over. */
broker: Broker;
@@ -30,6 +43,29 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
// runtime that always served would fail for exactly the modules that never needed it.
const tools = listTools();
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
// And, for every module that serves any, the verb that says what it serves. Refused before
// anything is bound if a module named a tool of its own `tools`: one name answering two things
// is the fault nobody can diagnose afterwards, and the runtime is the only place that sees both.
const stops: Array<() => void> = [stop];
for (const { module, tools: own } of collectTools()) {
if (own.length === 0) continue;
if (own.some((t) => t.name === TOOLS_VERB)) {
stop();
throw new Error(
`${module} names a tool "${TOOLS_VERB}", which is the verb the runtime answers for every ` +
"module with what it serves (novox/hq ADR 0152) — refused, rename it",
);
}
const answer: ToolsAnswer = {
module,
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input })),
};
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async () => answer));
}
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
return stop;
return () => {
for (const s of stops) s();
};
}
+57 -12
View File
@@ -18,37 +18,62 @@ import { callTool, toolsOn, whyItFailed } from "../dist/client.js";
const url = process.env.MESH_TEST_NATS;
/** A module serving the catalogue's tool list and one tool of its own, so the client has a mesh to
* talk to. Two connections, because a person and a module are different users even in a test. */
/** 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. */
async function aMeshWithTools() {
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
const shop = await connectNats({ url: url!, module: "shop" });
await catalogue.handle("catalog_tools", async () => ({
tools: [
{ module: "shop", name: "price", description: "what something costs", input: { type: "object" } },
{ module: "mesh-catalog", name: "catalog_tools", description: "what tools the mesh has" },
],
await 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 controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
return {
async close() {
await catalogue.close();
await shop.close();
await controller.close();
},
};
}
test("a person sees what the catalogue says the mesh has, sorted", async (t) => {
test("a person sees what the running modules answer, sorted, and who did not answer", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const mesh = await aMeshWithTools();
const person = await connectNats({ url, module: "person.ada" });
try {
const tools = await toolsOn(person);
const began = Date.now();
const have = await toolsOn(person);
assert.deepEqual(
tools.map((x) => `${x.module}.${x.name}`),
["mesh-catalog.catalog_tools", "shop.price"],
"the list is what the catalogue answered, in a stable order",
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",
);
// 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();
@@ -103,3 +128,23 @@ 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, { output: "all quiet", ok: true });
const direct = await callTool(person, "seat:mesh-controller.status", {});
assert.deepEqual(direct, { output: "all quiet", ok: true });
} finally {
await person.close();
await mesh.close();
}
});
+6
View File
@@ -0,0 +1,6 @@
// A module naming a tool after the verb the runtime answers for every module — refused at load.
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("clash", () => [
{ name: "tools", description: "mine, not the runtime's", input: {}, run: async () => ({}) },
]);
+12
View File
@@ -0,0 +1,12 @@
// 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 }) },
]);
+113
View File
@@ -0,0 +1,113 @@
/**
* The console: the same surface over HTTP on loopback, started the way the mesh starts it — on the
* module credential in MESH_BROKER_FILE (novox/hq ADR 0152, design 34 §2).
*
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/http.test.ts
*/
import assert from "node:assert/strict";
import { test } from "node:test";
import { spawn, type ChildProcess } from "node:child_process";
import { mkdtemp, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { connectNats } from "../dist/broker-nats.js";
import { serveMcpHttp } from "../dist/http.js";
const url = process.env.MESH_TEST_NATS;
async function aMesh(t: { after: (fn: () => Promise<void> | void) => void }) {
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
const shop = await connectNats({ url: url!, module: "shop" });
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "shop" }] }));
await shop.handle("tools", async () => ({
module: "shop",
tools: [{ name: "price", description: "what something costs", input: {} }],
}));
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
t.after(async () => {
await catalogue.close();
await shop.close();
});
}
/** `mesh serve` as the mesh runs it: MESH_BROKER_FILE, a listen address, nothing else. */
async function aConsole(t: { after: (fn: () => Promise<void> | void) => void }): Promise<string> {
const dir = await mkdtemp("/tmp/mesh-console-");
const credential = join(dir, "broker");
await writeFile(credential, JSON.stringify({ url, node: "desk", module: "mesh-console", user: "desk.mesh-console", password: "x" }));
const child: ChildProcess = spawn(process.execPath, ["dist/mesh.js", "serve", "--listen", "127.0.0.1:0"], {
env: { ...process.env, MESH_BROKER_FILE: credential, MESH_CREDENTIAL: "" },
stdio: ["ignore", "pipe", "pipe"],
});
t.after(() => {
child.kill("SIGTERM");
});
return new Promise((resolve, reject) => {
let out = "";
let err = "";
child.stdout!.on("data", (d) => {
out += d.toString();
const m = /listening on (http:\/\/[^/]+\/mcp)/.exec(out);
if (m) resolve(m[1]!);
});
child.stderr!.on("data", (d) => (err += d.toString()));
child.on("exit", (code) => reject(new Error(`serve exited ${code}: ${err}`)));
});
}
async function post(endpoint: string, body: unknown): Promise<{ status: number; json?: any }> {
const res = await fetch(endpoint, {
method: "POST",
headers: { "content-type": "application/json", accept: "application/json" },
body: JSON.stringify(body),
});
const text = await res.text();
return { status: res.status, json: text ? JSON.parse(text) : undefined };
}
test("the console answers a host on loopback, as the account the mesh gave it", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
await aMesh(t);
const endpoint = await aConsole(t);
const hello = await post(endpoint, { jsonrpc: "2.0", id: 1, method: "initialize", params: {} });
assert.equal(hello.status, 200);
assert.match(hello.json.result.instructions, /desk\.mesh-console/, "the handshake names the console's account");
const heard = await post(endpoint, { jsonrpc: "2.0", method: "notifications/initialized" });
assert.equal(heard.status, 202, "a notification is heard and not answered");
const listed = await post(endpoint, { jsonrpc: "2.0", id: 2, method: "tools/list" });
assert.deepEqual(listed.json.result.tools.map((x: { name: string }) => x.name), ["shop.price"]);
const called = await post(endpoint, {
jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "shop.price", arguments: { of: "a hat" } },
});
assert.deepEqual(JSON.parse(called.json.result.content[0].text), { of: "a hat", cost: 12 });
// A person's client through the same endpoint, with no credential of its own.
const { main } = await import("../dist/mesh.js");
const logged: string[] = [];
const was = console.log;
console.log = (line: string) => logged.push(String(line));
try {
assert.equal(await main(["tools", "--console", endpoint]), 0);
} finally {
console.log = was;
}
assert.ok(logged.some((l) => l.startsWith("shop.price")), `the client did not list through the console: ${logged}`);
});
test("the console binds loopback and nowhere else", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
const bus = await connectNats({ url, module: "mesh-console", node: "desk" });
try {
await assert.rejects(() => serveMcpHttp(bus, "desk.mesh-console", "0.0.0.0:0"), /loopback and nowhere else/);
const up = await serveMcpHttp(bus, "desk.mesh-console", "127.0.0.1:0");
assert.match(up.address, /^127\.0\.0\.1:\d+$/);
await up.close();
} finally {
await bus.close();
}
});
+32 -6
View File
@@ -20,13 +20,23 @@ 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_tools", async () => ({
tools: [{ module: "shop", name: "price", description: "what something costs" }],
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 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: {} }] }],
}));
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
t.after(async () => {
await catalogue.close();
await shop.close();
await controller.close();
});
const { mkdtemp, writeFile } = await import("node:fs/promises");
@@ -80,14 +90,19 @@ 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, "2024-11-05");
assert.equal(hello.protocolVersion, "2025-03-26");
assert.ok(hello.capabilities.tools, "a server offering no tools is not this one");
assert.match(hello.instructions, /ada/, "the handshake says whose authority a call is made under");
const listed = byId.get(2)!.result.tools;
assert.equal(listed.length, 1);
assert.equal(listed[0].name, "shop.price", "a tool is named the way a person names it");
assert.ok(listed[0].inputSchema, "a tool with no schema is one an agent cannot call");
assert.deepEqual(listed.map((x: { name: string }) => x.name), ["mesh-controller.status", "shop.price"],
"the modules' tools and the roles', named the way a person names them");
const price = listed[1];
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.
assert.deepEqual(price.inputSchema, { type: "object", properties: { of: { type: "string" } } });
// Silence is named: the module the catalogue holds and nothing answered for.
assert.deepEqual(byId.get(2)!.result._meta.notAnswering, ["ghost"]);
const called = byId.get(3)!.result;
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
@@ -122,3 +137,14 @@ 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/);
});
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 });
});
+38 -1
View File
@@ -1,6 +1,7 @@
import { spawn } from "node:child_process";
import { test } from "node:test";
import assert from "node:assert/strict";
import { fatalBrokerReason, PinMismatchError } from "../src/broker-nats.ts";
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
@@ -41,3 +42,39 @@ test("a non-Error value does not crash the classifier", () => {
assert.equal(fatalBrokerReason("just a string"), null);
assert.equal(fatalBrokerReason(undefined), null);
});
// **A module asked to prepare its state and naming nothing is a failure, not a no-op** (novox/hq
// ADR 0135). The mesh asks this only of a module whose manifest says it prepares something, so an
// image that names nothing was built wrong, and exiting 0 would let that version serve against a
// state nobody shaped.
test("preparing with nothing named fails rather than passing quietly", async () => {
const runtime = new URL("../dist/main.js", import.meta.url).pathname;
const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
const child = spawn(process.execPath, [runtime, "prepare"], {
env: { ...process.env, MESH_PREPARE: "" },
});
let said = "";
child.stderr.on("data", (chunk) => (said += String(chunk)));
child.on("close", (code) => resolve({ code, said }));
});
assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
assert.match(ran.said, /MESH_PREPARE/);
});
// **One durable consumer feeds one reader, however many patterns a module registers.**
//
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
// reader that received one its own pattern does not match acknowledges it, which is right for a
// filter wider than anything registered and silent loss when it is another handler's. The matching is
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not
// by which reader happened to fetch it.
test("a message is for every pattern that matches it, and nothing else", () => {
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"];
const matched = (key: string) => registered.filter((p) => topicMatches(p, key));
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]);
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]);
// Nothing registered for it: the consumer's filter is the controller's and may be wider.
assert.deepEqual(matched("mesh-catalog.upgraded"), []);
// And a handler that asked for everything gets both, which is what the audit logger does.
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]);
});
+60
View File
@@ -0,0 +1,60 @@
/**
* Every module's runtime answers `tools` for it (novox/hq ADR 0152, design 34 §3): the names,
* descriptions and schemas from the code that answers them. Against a real bus, because the claim is
* what a second connection gets back.
*
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/runtime-tools.test.ts
*/
import assert from "node:assert/strict";
import { test } from "node:test";
import { fileURLToPath } from "node:url";
import { resetTools } from "@novox/mesh-sdk/tools";
import { connectNats } from "../dist/broker-nats.js";
import { runTools } from "../dist/runtime.js";
const url = process.env.MESH_TEST_NATS;
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
test("a module registering two tools answers three names, the third being what it serves", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const shop = await connectNats({ url, node: "one", module: "shop" });
const asker = await connectNats({ url, module: "person.ada" });
const stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
try {
const answer = await asker.request<Record<string, never>, { module: string; tools: { name: string; input: unknown }[] }>(
"shop.tools",
{},
);
assert.equal(answer.module, "shop");
assert.deepEqual(answer.tools.map((x) => x.name), ["price", "refund"]);
// The schema travels with the name: a name alone is not callable by something that has never
// seen the mesh before.
assert.deepEqual(answer.tools[0].input, { of: { type: "string", description: "the thing" } });
// And the tools themselves still answer beside it.
const priced = await asker.request<{ of: string }, { cost: number }>("shop.price", { of: "a hat" });
assert.equal(priced.cost, 12);
} finally {
stop();
await asker.close();
await shop.close();
}
});
test("a module naming a tool of its own `tools` is refused at load", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const clash = await connectNats({ url, node: "one", module: "clash" });
try {
await assert.rejects(
() => runTools({ broker: clash, moduleEntrypoints: [fixture("clash-tools.mjs")] }),
/names a tool "tools"/,
);
} finally {
resetTools();
await clash.close();
}
});