Merge pull request 'The runtime serves a list of modules on one credential, naming a bundle that fails to load (hq ADR 0175, to-be 38 WP1)' (#29) from feat/the-operators-machine into main

This commit was merged in pull request #29.
This commit is contained in:
2026-10-02 19:27:49 +00:00
12 changed files with 843 additions and 161 deletions
+41 -17
View File
@@ -1,29 +1,50 @@
# mesh-tools
The Novox Mesh **tool runtime** — the per-node process that makes a module's tools actually serve.
The Novox Mesh **tool runtime** — the one process per node that makes every assigned module's tools
actually serve (novox/hq ADR 0175).
A module ships its tools (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk)); this
runtime is what loads them and puts them on the mesh. It:
A module ships its tools as a bundle (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk));
this runtime is what loads them and puts them on the mesh. It:
1. connects the mesh broker (novox/hq ADR 0001) — a concrete AMQP implementation of the sdk's
`Broker` contract;
2. imports the assigned modules' compiled tool entrypoints, each of which registers its tools as it
loads;
3. serves them through the sdk's `serveTools` harness, answering `tools.invoke` over the broker.
1. connects the mesh bus on the node's credential — a concrete implementation of the sdk's `Broker`
contract;
2. reads one membership per module it serves — what the mesh issued that module on this machine
(ADR 0160): where its tools are answered, which seats it holds — and follows each live;
3. imports each module's compiled tool entrypoints, each of which registers its tools as it loads,
guarded: a bundle that throws is named, in the log and in what `tools` answers for its module,
and the others serve;
4. serves every module's tools on that module's subjects and every held seat's verbs on the seat's.
A bundle that is not plain JavaScript — a Go or Rust binary, a Python script, or a JavaScript file
marked executable — is **launched** rather than imported (novox/hq ADR 0187): the runtime starts it
as a child with its own environment and speaks MCP over stdio to it, `tools/list` once and
`tools/call` per call. A tool it lists as `<seat>.<verb>` is the seat's implementation. A child that
exits is named in the log and started again on its next call. So a tools bundle may be written in
any language; the mesh's SDK for each is the stdio loop and nothing more (`src/launch.ts` is the
runtime's side of it).
Everything hard — dispatch, collection, duplicate-name safety — is the sdk's. This is the thin
wrapper that binds the broker and loads the modules. Keeping the AMQP client here, out of the sdk,
is deliberate: a broker-client change never rebuilds a module (ADR 0039).
wrapper that binds the bus and loads the modules. Keeping the bus client here, out of the sdk, is
deliberate: a bus-client change never rebuilds a module (ADR 0039). The runtime is module-agnostic:
it knows bundles and subjects, nothing of what any module does. A tool that needs root escalates
itself — root is the module's concern, not the runtime's.
## Running it
```
MESH_BROKER_URL amqp://… the mesh broker
MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool entrypoints
MESH_BROKER_FILE the node's sealed credential, as the mesh delivered it
MESH_TOOL_MODULES alpha=/…/alpha/tools/index.js,beta=/…/beta/dist/index.js,…
the modules to serve and their compiled entrypoints; several entries may
name one module. A bare path is an entrypoint of the credential's own module
— the one-module form a per-module container still sets.
MESH_OPERATOR_ACCOUNT whose machine this is, and MESH_OPERATOR_HOME where their home is; set by
the mesh when the node has an account, read by tools from their environment
MESH_BROKER_URL a plain URL instead of the credential, for the bootstrap case
```
`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables
and starts it like any other supervised workload.
`node dist/main.js`. On a node the controller composes the variables and the host supervises the
process like any other host-side workload (novox/hq to-be 38). The container (`Dockerfile`) is how
a module's own *service* may still be built; it is no longer how tools reach a node.
## `mesh` — the tools for whoever is on a machine
@@ -41,6 +62,9 @@ dropped. A module may not name a tool of its own `tools`; the runtime refuses it
## Verified
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the
runtime serves a registered tool, a separate connection invokes it by name and gets the result, and
an unknown tool is refused over the wire.
`npm test` runs against a real NATS server with JetStream (`MESH_TEST_NATS`, see any test's header
for the one-line `docker run`) and proves the whole path over the wire: the runtime serves a
registered tool, a separate connection invokes it by name and gets the result, an unknown tool is
refused, a runtime serves exactly the subjects it is issued and re-serves on a new membership, and
the node's runtime serves three modules' bundles on one credential — one of them broken, named and
not fatal — with every seat verb answering where the membership put it.
+86 -42
View File
@@ -16,6 +16,7 @@
// derives the subject (design 29 §1), so reorganising the subject space leaves every module
// correct.
import { AsyncLocalStorage } from "node:async_hooks";
import { createHash } from "node:crypto";
import net from "node:net";
import tls from "node:tls";
@@ -71,19 +72,41 @@ export interface Answered<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). */
/** The bus as the runtime sees it: the sdk's contract, and the few things only the runtime needs —
* an answer that says which machine gave it, serving a subject that is not a module's own tool
* (a seat's verb), and following the memberships of the modules it serves beside its own. */
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. */
/**
* Serve another module's tools from this connection: read its membership on this machine and
* follow it, so `handle("<module>.<tool>")` is served where the mesh issued that module (novox/hq
* ADR 0175: one runtime per node, every assigned module's tools). The account must be allowed
* to read that membership and to subscribe its subjects — the node's is; a module's own is not,
* and is refused by the bus, not here.
*/
follow(module: string): Promise<void>;
/** What the mesh issued this assignment — this module's when none is named — or undefined when
* nothing has been issued yet. */
membership(module?: string): Membership | undefined;
/** Called when the mesh issues a new membership to any module this connection follows; the
* runtime re-serves on it. The membership says which module it is for. */
onMembership(handler: (m: Membership) => void): void;
/** The modules whose memberships this connection follows: its own and every one `follow` added. */
serving(): string[];
/** The module this connection is: what its credential named, and what a bare key serves as. */
readonly module: string;
}
/**
* Which module's tool is at work on this connection, when one is (novox/hq ADR 0175). One runtime
* carries many modules' tools, and an event a tool emits must land on the emitting *module's*
* subject, not the runtime's — so the runtime runs each tool inside this store, and `publish`
* reads the module from it. Outside a tool the connection's own module stands.
*/
export const atWork = new AsyncLocalStorage<{ module: string }>();
const ASSIGNMENTS_STREAM = "ASSIGNMENTS";
/** The one address a runtime derives for itself (ADR 0160). */
@@ -153,38 +176,45 @@ export async function connectNats(
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;
// The memberships, one per module this connection serves, each read once and followed. A
// module's own runtime follows one, its own; the node's runtime follows one per assigned module
// (novox/hq ADR 0175) — the same subject shape, the same stream, read as many times as there are
// modules, and nothing on the bus learns a new shape for it. A direct get is one request on the
// stream's API, which is the whole of what an account may ask JetStream for a subject it is
// granted; a 404 is a mesh that has not issued one, which is a fact to say and not an error to
// retry.
const issued = new Map<string, Membership | undefined>();
const issuedHandlers: ((m: Membership) => void)[] = [];
const subjectOfMine = node ? membershipSubject(node, self) : "";
if (subjectOfMine) {
const follow = async (module: string): Promise<void> => {
if (!node || issued.has(module)) return;
issued.set(module, undefined);
const subjectOfTheirs = membershipSubject(node, module);
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}`,
const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}.${subjectOfTheirs}`,
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;
issued.set(module, 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`);
if (!issued.get(module)) {
console.log(`[mesh-tools] no membership issued for ${module} on ${node} yet; serving the derived shape until one arrives`);
}
try {
const live = conn.subscribe(subjectOfMine);
const live = conn.subscribe(subjectOfTheirs);
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);
const m = JSON.parse(sc.decode(msg.data)) as Membership;
issued.set(module, m);
console.log(`[mesh-tools] ${module} on ${node} was issued a new membership; re-serving on it`);
for (const h of issuedHandlers) h(m);
} catch (err) {
console.log(`[mesh-tools] a membership arrived that is not one: ${err}`);
}
@@ -193,32 +223,33 @@ export async function connectNats(
} catch {
// A subscription this account may not make is a mesh older than the membership.
}
}
};
await follow(self);
/** 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;
/** The subjects a tool of a served module is served on: from that module's membership when
* issued, derived otherwise (the shape the mesh issues on day one, so the two agree). */
const servedOn = (module: string, tool: string): { subject: string; queue?: string }[] => {
const m = issued.get(module);
if (m) {
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}` });
out.unshift({ subject: m.tools, queue: `serve.${module}` });
}
return out;
}
const base = `mesh.mod.${self}.tool.${tool}`;
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }];
const base = `mesh.mod.${module}.tool.${tool}`;
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${module}` }];
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. */
/** Where a call by key goes: a subject this connection's own membership says it 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];
const reach = issued.get(self)?.reaches?.[name];
if (reach && reach.length > 0) {
if (wanted) {
const at = reach.find((s) => s.endsWith(`.${wanted}`));
@@ -288,27 +319,38 @@ export async function connectNats(
*/
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.
// knows its role without a membership; everything else is a served module's own tool, served
// where the mesh issued that module (ADR 0160). A key naming a module this connection does
// not follow is not served here at all: a module serves its own tools, and the node's
// runtime those of the modules it was given (ADR 0175) — never a stranger's.
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 module = dot < 0 ? self : key.slice(0, dot);
if (!issued.has(module) && module !== self) {
throw new Error(
`${self} cannot serve ${key}: a module serves its own tools, and a runtime those of the ` +
"modules it follows",
);
}
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(() => {
let stops = servedOn(module, tool).map((s) => answerOn(s.subject, s.queue, handler));
// When that module's new membership arrives, serve where it now says and stop serving where
// it no longer does.
issuedHandlers.push((m) => {
if (m.module !== module) return;
stops.forEach((stop) => stop());
stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
stops = servedOn(module, tool).map((s) => answerOn(s.subject, s.queue, handler));
});
return () => stops.forEach((stop) => stop());
},
membership: () => issued,
module: self,
follow,
serving: () => [...issued.keys()],
membership: (module?: string) => issued.get(module ?? self),
onMembership: (handler: (m: Membership) => void) => {
issuedHandlers.push(handler);
},
@@ -339,7 +381,9 @@ export async function connectNats(
if (!meta["content-type"]) h.set("content-type", "application/json");
if (env.node) h.set("x-node", env.node);
await js.publish(eventSubject(env.key, self), sc.encode(JSON.stringify(env.body)), {
// The emitting module's subject: the tool at work's when a served module's tool emits from
// the node's runtime (ADR 0175), this connection's own otherwise.
await js.publish(eventSubject(env.key, atWork.getStore()?.module ?? self), sc.encode(JSON.stringify(env.body)), {
headers: h,
// De-duplicated by the server inside its window, so a redelivery after a crash between
// publishing and acknowledging is not seen twice. Only the emitter can make this id.
+6 -1
View File
@@ -192,7 +192,12 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
}
asked.forEach((outcome, i) => {
const module = names[i]!;
if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
if (outcome.status === "fulfilled" && typeof outcome.value?.failed === "string") {
// The runtime answered for it and serves nothing: the bundle failed to load (ADR 0175). Said
// with the reason, because "not answering" would send somebody to check an assignment that
// is fine.
notAnswering.push(`${module} (its tools bundle failed to load: ${outcome.value.failed})`);
} else 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 });
}
+169
View File
@@ -0,0 +1,169 @@
// A tools bundle as a process the runtime launches (novox/hq ADR 0187).
//
// The runtime does not run a tool's code itself when the bundle is not JavaScript: it starts the
// bundle's executable as a child with the runtime's environment and speaks MCP over stdio to it —
// `initialize`, `tools/list` once, `tools/call` per call. A Rust binary, a Go binary, a Python
// script and a Node script are the same thing from here: a process that answers those. Everything
// the mesh adds — the subjects from the membership, the held seats, the `tools` answer, a bundle
// that failed named and the others serving — is the runtime's, outside this file.
//
// A tool the child lists as `<seat>.<verb>` is the module's implementation of that seat's verb;
// any other name is the module's own tool. The same rule the in-process registration follows.
import { spawn, type ChildProcess } from "node:child_process";
import { accessSync, constants } from "node:fs";
import type { ToolDefinition } from "@novox/mesh-sdk/tools";
/** The protocol version this speaks; a bundle says the same. */
export const PROTOCOL = "2025-03-26";
/** How long a child has to answer `initialize` and `tools/list` before it is a failed bundle, and
* how long a call may take before the caller is told the tool is slow rather than absent. */
const HANDSHAKE_MS = 10_000;
const CALL_MS = 30_000;
/** Whether an entrypoint is launched as a process rather than imported: anything that is not a
* plain JavaScript file, and a JavaScript file marked executable — a bundle written against the
* protocol in TypeScript, served the same way as any other language. */
export function launches(entry: string): boolean {
const javascript = /\.(m|c)?js$/.test(entry);
let executable = false;
try {
accessSync(entry, constants.X_OK);
executable = true;
} catch {
// not executable, or not there — importing will say which
}
return !javascript || executable;
}
/** What a launched bundle registers: the groups the in-process path would have, by name. */
export interface Launched {
registrations: { module: string; tools: ToolDefinition[] }[];
stop(): void;
}
interface Pending {
resolve(v: any): void;
reject(e: Error): void;
timer: NodeJS.Timeout;
}
/**
* Launch a bundle and learn its tools. Rejects when the child cannot be started or does not complete
* the handshake, which the runtime records as the bundle having failed. A child that exits later is
* started again on the next call, once; a call in flight when it died is told so.
*/
export async function launch(module: string, entry: string, env: NodeJS.ProcessEnv = process.env): Promise<Launched> {
let child: ChildProcess | undefined;
let nextId = 1;
const pending = new Map<number, Pending>();
let stopped = false;
const start = async (): Promise<void> => {
const proc = spawn(entry, [], { stdio: ["pipe", "pipe", "pipe"], env });
child = proc;
let buffered = "";
proc.stdout!.on("data", (chunk: Buffer) => {
buffered += chunk.toString("utf8");
let at: number;
while ((at = buffered.indexOf("\n")) >= 0) {
const line = buffered.slice(0, at).trim();
buffered = buffered.slice(at + 1);
if (!line) continue;
let reply: { id?: number; result?: unknown; error?: { message?: string } };
try {
reply = JSON.parse(line);
} catch {
console.log(`[mesh-tools] ${module}'s bundle said something that is not a reply: ${line.slice(0, 120)}`);
continue;
}
const waiting = typeof reply.id === "number" ? pending.get(reply.id) : undefined;
if (!waiting) continue;
pending.delete(reply.id!);
clearTimeout(waiting.timer);
if (reply.error) waiting.reject(new Error(reply.error.message ?? "the bundle refused the request"));
else waiting.resolve(reply.result);
}
});
// stderr is the bundle's log; kept under the module's name so a fault reads where it belongs.
proc.stderr!.on("data", (chunk: Buffer) => {
for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) console.log(`[${module}] ${line}`);
});
const exited = new Promise<never>((_, reject) => {
proc.once("error", (err) => reject(err));
proc.once("exit", (code, signal) => {
const why = `${module}'s bundle exited (${signal ?? code})`;
for (const [id, p] of pending) {
pending.delete(id);
clearTimeout(p.timer);
p.reject(new Error(why));
}
if (child === proc) child = undefined;
if (!stopped) console.log(`[mesh-tools] ${why}; started again on its next call`);
reject(new Error(why));
});
});
const ask = (method: string, params: unknown, ms: number): Promise<any> =>
Promise.race([
new Promise<any>((resolve, reject) => {
const id = nextId++;
const timer = setTimeout(() => {
pending.delete(id);
reject(new Error(`${module}'s bundle did not answer ${method} in ${ms / 1000}s`));
}, ms);
pending.set(id, { resolve, reject, timer });
proc.stdin!.write(JSON.stringify({ jsonrpc: "2.0", id, method, params }) + "\n");
}),
exited,
]);
exited.catch(() => {}); // observed through the race; never unhandled
(proc as ChildProcess & { ask?: typeof ask }).ask = ask;
await ask("initialize", { protocolVersion: PROTOCOL, capabilities: {}, clientInfo: { name: "node-tools", version: "1" } }, HANDSHAKE_MS);
proc.stdin!.write(JSON.stringify({ jsonrpc: "2.0", method: "notifications/initialized" }) + "\n");
};
const asking = async (method: string, params: unknown, ms: number): Promise<any> => {
if (!child) await start();
return (child as ChildProcess & { ask: (m: string, p: unknown, ms: number) => Promise<any> }).ask(method, params, ms);
};
await start();
const listed = (await asking("tools/list", {}, HANDSHAKE_MS)) as { tools?: { name: string; description?: string; inputSchema?: unknown }[] };
const groups = new Map<string, ToolDefinition[]>();
for (const t of listed.tools ?? []) {
const dot = t.name.indexOf(".");
const under = dot < 0 ? module : t.name.slice(0, dot);
const name = dot < 0 ? t.name : t.name.slice(dot + 1);
const tools = groups.get(under) ?? [];
tools.push({
name,
description: t.description ?? "",
input: (t.inputSchema as Record<string, unknown> | undefined) ?? {},
run: async (args) => {
const result = (await asking("tools/call", { name: t.name, arguments: args ?? {} }, CALL_MS)) as {
content?: { type: string; text?: string }[];
isError?: boolean;
};
const text = result?.content?.find((c) => c.type === "text")?.text ?? "";
if (result?.isError) throw new Error(text || `${module}.${t.name} failed`);
// The bundle's answer is JSON as text (that is what every MCP host renders); handed back as
// the value it encodes so a caller on the bus sees what an in-process tool would return.
try {
return JSON.parse(text);
} catch {
return text;
}
},
});
groups.set(under, tools);
}
return {
registrations: [...groups].map(([under, tools]) => ({ module: under, tools })),
stop: () => {
stopped = true;
child?.kill("SIGTERM");
child = undefined;
},
};
}
+41 -9
View File
@@ -21,14 +21,15 @@
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0043) — an
// 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_TOOL_MODULES <module>=/path/a,… the modules to serve and their compiled entrypoints (serve
// mode; ADR 0175); a bare path is an entrypoint of the credential's own module
// 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";
import { runTools, type ServedModule } from "./runtime.js";
/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */
let lastCredential: Credential | undefined;
@@ -118,14 +119,42 @@ async function connectBrokerPatiently(): Promise<Broker> {
}
}
async function serve(): Promise<void> {
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
.split(",")
.map((s) => s.trim())
.filter(Boolean);
/**
* What MESH_TOOL_MODULES names (novox/hq ADR 0175, to-be 38 WP1): `<module>=<entrypoint>` entries,
* comma-separated, several per module allowed — the node's runtime serving every assigned module's
* bundle. A bare path is the one-module form the per-module containers still set: an entrypoint of
* the credential's own module. Both may appear; the result is one list of modules.
*/
export function servedModulesFrom(spec: string, own: string | undefined): { serves: ServedModule[]; moduleEntrypoints: string[] } {
const serves = new Map<string, string[]>();
const moduleEntrypoints: string[] = [];
for (const raw of spec.split(",")) {
const entry = raw.trim();
if (!entry) continue;
const eq = entry.indexOf("=");
if (eq < 0) {
moduleEntrypoints.push(entry);
continue;
}
const module = entry.slice(0, eq).trim();
const path = entry.slice(eq + 1).trim();
if (!module || !path) {
throw new Error(`MESH_TOOL_MODULES: "${entry}" is neither <module>=<entrypoint> nor an entrypoint of this module`);
}
if (module === own) {
moduleEntrypoints.push(path);
continue;
}
serves.set(module, [...(serves.get(module) ?? []), path]);
}
return { serves: [...serves].map(([module, entrypoints]) => ({ module, entrypoints })), moduleEntrypoints };
}
async function serve(): Promise<void> {
const broker = await connectBrokerPatiently();
const stop = await runTools({ broker, moduleEntrypoints, credential: lastCredential });
// Parsed after connecting: a bare entrypoint belongs to the module the credential names.
const { serves, moduleEntrypoints } = servedModulesFrom(process.env.MESH_TOOL_MODULES ?? "", lastCredential?.module);
const stop = await runTools({ broker, serves, moduleEntrypoints, credential: lastCredential });
const shutdown = async (): Promise<void> => {
stop();
@@ -250,4 +279,7 @@ async function main(): Promise<void> {
await serve();
}
void main();
// Only when run, so a test can import the pieces.
if (process.argv[1] && import.meta.url === new URL(`file://${process.argv[1]}`).href) {
void main();
}
+243 -92
View File
@@ -1,14 +1,21 @@
// The tool runtime — the thin per-node process that makes a module's tools actually serve. It
// binds the mesh broker, loads the assigned modules' tool entrypoints (each of which calls
// registerModuleTools as it imports), and hands them to the sdk's serving harness. Everything hard
// — dispatch, collection, duplicate-name safety — is the sdk's; this is the wrapper.
// The tool runtime — the per-node process that makes the mesh's tools actually serve (novox/hq
// ADR 0175). It binds the mesh broker, loads the served modules' tool bundles (each of which calls
// registerModuleTools as it imports), and serves every module's tools on that module's subjects and
// every held seat's verbs on the seat's. Everything hard — dispatch, collection, duplicate-name
// safety — is the sdk's; this is the wrapper.
//
// One runtime, many modules. It was written for one module per process and ran that way in a
// container per module; it now serves a list, as the one process per node the host supervises,
// and the per-module shape is the list with one entry. A bundle that fails to import is named —
// in the log and in what `tools` answers for it — and the others serve.
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 { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools";
import type { Broker } from "@novox/mesh-sdk/messaging";
import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
import { launch, launches } from "./launch.js";
/**
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
@@ -28,137 +35,281 @@ export interface ToolsAnswer {
* first when there is one, then this machine's. A caller composes nothing. */
subjects?: string[];
}[];
/** Why this module serves nothing here, when its bundle failed to load (ADR 0175): said where
* discovery looks, so a module that is silent and one that is broken are told apart. */
failed?: string;
}
/** One module this runtime serves: its name and its compiled tool entrypoints. */
export interface ServedModule {
module: string;
/** Absolute paths to the module's compiled tool entrypoints (e.g. .../umami/tools/index.js). */
entrypoints: 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 modules to serve, each with its entrypoints. */
serves?: ServedModule[];
/** The credential's own module's entrypoints — the one-module form, which the per-module
* containers still use; the same as naming the credential's module in `serves`. */
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. */
* (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat its
* memberships do not name. */
credential?: Credential;
}
/** Two environment words the mesh sets for the node's runtime and every tool reads from its
* environment: whose machine this is (novox/hq to-be 37 §3, ADR 0175). */
export const OPERATOR_ACCOUNT = "MESH_OPERATOR_ACCOUNT";
export const OPERATOR_HOME = "MESH_OPERATOR_HOME";
/** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */
export async function runTools(opts: RuntimeOptions): Promise<() => void> {
useBroker(() => opts.broker);
const runtime = opts.broker as RuntimeBroker;
// Whose runtime this is: the credential's module, or the connection's own when a runtime is
// started by hand without one — the broker was told its module when it connected.
const self = opts.credential?.module ?? (typeof runtime.module === "string" ? runtime.module : undefined);
for (const entry of opts.moduleEntrypoints) {
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake.
await import(pathToFileURL(resolve(entry)).href);
// What to serve: the list, with the one-module form folded in as the credential's own entry.
const served = new Map<string, string[]>();
for (const s of opts.serves ?? []) {
served.set(s.module, [...(served.get(s.module) ?? []), ...s.entrypoints]);
}
if (opts.moduleEntrypoints?.length) {
if (!self) {
throw new Error(
"entrypoints were given with no module to serve them as: name the module (MESH_TOOL_MODULES " +
"as <module>=<entrypoint>) or connect on a credential that names one",
);
}
served.set(self, [...(served.get(self) ?? []), ...opts.moduleEntrypoints]);
}
// 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`);
// The operator's machine, said once so a tool's behaviour under it can be read back from the
// log. Tools read the two words from their own environment, which is this process's.
const account = process.env[OPERATOR_ACCOUNT];
if (account) {
console.log(`[mesh-tools] the operator's account here is ${account}` +
(process.env[OPERATOR_HOME] ? ` (home ${process.env[OPERATOR_HOME]})` : ""));
}
// Follow every served module's membership before loading anything, so what each is issued is
// known when its tools are bound. A module's own runtime already follows its own.
if (typeof runtime.follow === "function") {
for (const module of served.keys()) await runtime.follow(module);
}
// Import each bundle, guarded (ADR 0175: one faulty bundle must not take the node's tools down).
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake — and
// the registrations it adds are the ones that appear after it, which is how each is attributed
// to the module whose bundle made it.
// A bundle that is not plain JavaScript — or is marked executable — is launched as a process
// and spoken to over MCP on stdio instead (ADR 0187); what it lists is registered the same way.
const failed = new Map<string, string>();
const owner: string[] = []; // registration index → the module whose bundle registered it
const launched: { module: string; owner: string; tools: ToolDefinition[] }[] = [];
const children: Array<() => void> = [];
for (const [module, entrypoints] of served) {
for (const entry of entrypoints) {
const path = resolve(entry);
try {
if (launches(path)) {
const child = await launch(module, path);
children.push(child.stop);
for (const r of child.registrations) launched.push({ ...r, owner: module });
continue;
}
const before = collectTools().length;
await import(pathToFileURL(path).href);
const after = collectTools().length;
for (let i = before; i < after; i++) owner[i] = module;
} catch (err) {
const why = err instanceof Error ? err.message : String(err);
failed.set(module, why);
console.log(`[mesh-tools] ${module}'s bundle ${entry} failed to load: ${why}; its tools are not served here`);
}
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;
// A registration under a served module's name is that module's tools, served on its subjects.
// One 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 where some served module claims the seat,
// never as a module's tools and never listed among them. A module named like its seat (the
// catalogue is the mesh-catalog seat) registers once and is both. Anything else is 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 claimed = seatsClaimed(served.keys(), self, opts.credential, runtime);
const registrations = [
...collectTools().map((r, i) => ({ ...r, owner: owner[i] ?? self ?? r.module })),
...launched,
];
const ownRegistrations = registrations.filter(({ module, owner: by }) => {
if (served.has(module)) return true;
if (claimed.has(module)) return false;
console.log(`[mesh-tools] ${by} registers tools under "${module}", which is neither a module served here nor a seat one of them claims; not served until the mesh issues the claim`);
return false;
});
const stops: Array<() => void> = [...children];
const stop = (): void => stops.splice(0).forEach((s) => s());
// 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. Likewise two tools of one module under one name.
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 seen = new Set<string>();
for (const t of own) {
if (seen.has(t.name)) throw new Error(`${module} exposes two tools named ${t.name} — refused`);
seen.add(t.name);
}
}
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();
};
// 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 module's membership (ADR 0160). A tool runs attributed to its
// module, so what it emits lands on the module's subject and not the runtime's.
const names: string[] = [];
for (const { module, tools: own } of ownRegistrations) {
for (const t of own) {
names.push(toolKey(module, t.name));
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) =>
atWork.run({ module }, () => t.run(args ?? {}))));
}
}
// And, for every served module, the verb that says what it serves — nothing, and why, for a
// module whose bundle failed. A module that registered nothing and did not fail is a pure-events
// module (the audit logger), whose scoped account may not declare the serve queue; it is left
// silent as it always was.
const byModule = new Map<string, ToolDefinition[]>();
for (const { module, tools: own } of ownRegistrations) {
byModule.set(module, [...(byModule.get(module) ?? []), ...own]);
}
for (const module of served.keys()) {
const own = byModule.get(module) ?? [];
const why = failed.get(module);
if (own.length === 0 && !why) continue;
const subjectsOf = (tool: string): string[] | undefined => {
const m = typeof runtime.membership === "function" ? runtime.membership(module) : undefined;
if (!m) return undefined;
const plain = m.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
const mine = m.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
return [...plain, ...mine];
};
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async (): Promise<ToolsAnswer> => ({
module,
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
...(why ? { failed: why } : {}),
})));
}
console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` +
(failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : ""));
stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations));
return () => stop();
}
/** The seats some served module claims: from the credential for its own module, and from every
* served module's membership (ADR 0160) — the node's runtime holds no claims of its own. */
function seatsClaimed(
modules: Iterable<string>,
self: string | undefined,
credential: Credential | undefined,
runtime: RuntimeBroker,
): Set<string> {
const out = new Set<string>();
for (const c of credential?.claims ?? []) if (credential?.module === self) out.add(c.seat);
for (const module of modules) {
const m = typeof runtime.membership === "function" ? runtime.membership(module) : undefined;
for (const s of m?.seats ?? []) out.add(s.seat);
}
return out;
}
/** One seat's verb, where its callers ask, and which served module holds the seat. */
interface SeatVerb {
seat: string;
verb: string;
subject: string;
holder: string;
}
/**
* 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.
* Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). What a served module
* claims and promises comes from its membership (ADR 0160) — and, for a module's own runtime, from
* its credential, which named the claims before memberships did. Each verb is served on the seat's
* own subject by the tool of the same name registered under the seat's 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 () => {};
async function serveClaimedSeats(
broker: RuntimeBroker,
served: string[],
self: string | undefined,
credential: Credential | undefined,
registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
): Promise<() => void> {
if (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));
for (const { module, owner, tools } of registrations) {
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args)));
implementations.set(module, verbs);
}
/** Every verb of every seat a served module claims, where the mesh issued it. */
const wanted = (): SeatVerb[] => {
const out: SeatVerb[] = [];
const have = new Set<string>();
const add = (v: SeatVerb): void => {
if (have.has(v.subject)) return;
have.add(v.subject);
out.push(v);
};
for (const module of served) {
const m = typeof broker.membership === "function" ? broker.membership(module) : undefined;
for (const s of m?.seats ?? []) add({ seat: s.seat, verb: s.verb, subject: s.subject, holder: module });
// The credential's claims, for the module's own runtime: where the mesh issued the verb when
// it has; the derived shape until then.
if (module !== self) continue;
for (const claim of credential?.claims ?? []) {
for (const verb of claim.serves ?? []) {
const subject = m?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
add({ seat: claim.seat, verb, subject, holder: module });
}
}
}
return out;
};
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`);
for (const v of wanted()) {
const run = implementations.get(v.seat)?.get(v.verb);
if (!run) {
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
continue;
}
stops.push(await broker.handleSubject(v.subject, run));
console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`);
}
};
await serve();
// A membership issued to any served module may add, move or withdraw a seat's verbs.
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
return () => stops.forEach((s) => s());
}
+9
View File
@@ -0,0 +1,9 @@
// One of several bundles the node's runtime loads (novox/hq ADR 0175): a module with two tools.
import { registerModuleTools } from "@novox/mesh-sdk/tools";
import { emit } from "@novox/mesh-sdk/events";
registerModuleTools("alpha", () => [
{ name: "one", description: "alpha's first", input: {}, run: async () => ({ alpha: 1 }) },
// A tool that emits: the event must land on alpha's subject, not the runtime's.
{ name: "two", description: "alpha's second, which emits", input: {}, run: async () => { await emit("happened", { by: "alpha" }); return { alpha: 2 }; } },
]);
+14
View File
@@ -0,0 +1,14 @@
// One of several bundles the node's runtime loads (novox/hq ADR 0175): a module with three tools of
// its own and the implementation of a node seat's two verbs under the seat's name (ADR 0159).
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("beta", () => [
{ name: "three", description: "beta's", input: {}, run: async () => ({ beta: 3 }) },
{ name: "four", description: "beta's", input: {}, run: async () => ({ beta: 4 }) },
{ name: "five", description: "beta's", input: {}, run: async () => ({ beta: 5 }) },
]);
registerModuleTools("node-shelf", () => [
{ name: "list", description: "what is on the shelf", input: {}, run: async () => ({ shelf: ["a", "b"] }) },
{ name: "clear", description: "take it all off", input: {}, run: async () => ({ cleared: true }) },
]);
+3
View File
@@ -0,0 +1,3 @@
// A bundle that throws on import — the fault ADR 0175 names as what got harder: one module's bundle
// must not take the node's other tools down.
throw new Error("gamma's bundle cannot find its client");
Vendored Executable
+29
View File
@@ -0,0 +1,29 @@
#!/usr/bin/env python3
# A tools bundle in a second language (novox/hq ADR 0187): MCP over stdio, no SDK, no dependencies.
# One tool of its own, one seat verb, and one that exits the process mid-call.
import json, sys
def say(m):
m["jsonrpc"] = "2.0"; sys.stdout.write(json.dumps(m) + "\n"); sys.stdout.flush()
TOOLS = [
{"name": "greet", "description": "say hello", "inputSchema": {"type": "object", "properties": {"who": {"type": "string"}}}},
{"name": "node-lamp.on", "description": "the seat's verb", "inputSchema": {"type": "object", "properties": {}}},
{"name": "die", "description": "exit without answering", "inputSchema": {"type": "object", "properties": {}}},
]
for line in sys.stdin:
req = json.loads(line); rid = req.get("id"); m = req.get("method"); p = req.get("params") or {}
if m == "initialize":
say({"id": rid, "result": {"protocolVersion": "2025-03-26", "capabilities": {"tools": {}}, "serverInfo": {"name": "delta", "version": "1"}}})
elif m == "tools/list":
say({"id": rid, "result": {"tools": TOOLS}})
elif m == "tools/call":
name = p.get("name"); args = p.get("arguments") or {}
if name == "greet":
say({"id": rid, "result": {"content": [{"type": "text", "text": json.dumps({"greeting": "hello " + args.get("who", "world"), "language": "python"})}]}})
elif name == "node-lamp.on":
say({"id": rid, "result": {"content": [{"type": "text", "text": json.dumps({"on": True, "language": "python"})}]}})
elif name == "die":
print("delta: told to die", file=sys.stderr); sys.exit(3)
else:
say({"id": rid, "error": {"code": -32602, "message": "no such tool"}})
+8
View File
@@ -0,0 +1,8 @@
#!/usr/bin/env node
// A TypeScript bundle written against the protocol and marked executable: served through the
// launcher like any other language, with the in-process shortcut off (novox/hq ADR 0187).
import { serveStdio } from "@novox/mesh-sdk/stdio";
await serveStdio("epsilon", [
{ name: "seven", description: "epsilon's", input: {}, run: async () => ({ epsilon: 7, via: "stdio" }) },
]);
+194
View File
@@ -0,0 +1,194 @@
/**
* One runtime per node serves every assigned module's tools (novox/hq ADR 0175, to-be 38 WP1). The
* node's runtime is handed a list of modules and their bundles on the node's credential; it reads
* one membership per module, serves each module's tools on that module's subjects and each held
* seat's verbs on the seat's, names a bundle that fails to load without dropping the others, and
* re-serves a module whose membership is re-issued mid-run. Against a real bus with JetStream.
*
* 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/node-runtime.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, toolsOn } from "../dist/client.js";
import { servedModulesFrom } from "../dist/main.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 controller's job, done by hand: the ASSIGNMENTS stream (last-per-subject, direct get) and an
* EVENTS stream for what a tool emits. */
async function aMesh() {
const nc = await connect({ servers: url! });
const jsm = await nc.jetstreamManager();
for (const name of ["ASSIGNMENTS", "EVENTS"]) {
try {
await jsm.streams.delete(name);
} catch {
// none yet
}
}
await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never);
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>"] });
return {
async issue(m: object & { node: string; module: string }) {
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
},
/** The next subject an event lands on, under a pattern. */
nextEvent(pattern: string): Promise<string> {
const sub = nc.subscribe(pattern, { max: 1 });
return (async () => {
for await (const m of sub) return m.subject;
throw new Error("no event");
})();
},
async close() {
await nc.close();
},
};
}
/** A membership as the controller issues one on a machine, with the module's own subject when it
* answers for the module anywhere, and the verbs of the node seats it holds. */
function membershipOf(module: string, node: string, opts: { plain?: boolean; seats?: Record<string, string[]> } = {}) {
const own = `mesh.mod.${module}`;
const serves: { subject: string; queue?: string }[] = [{ subject: `${own}.tool.{tool}.${node}` }];
if (opts.plain) serves.push({ subject: `${own}.tool.{tool}`, queue: `serve.${module}` });
const seats = Object.entries(opts.seats ?? {}).flatMap(([seat, verbs]) =>
verbs.map((verb) => ({ seat, verb, subject: `mesh.seat.${seat}.tool.${verb}.${node}` })));
return { node, module, serves, seats, emits: `${own}.event.{event}`, tools: `${own}.tool.tools` };
}
/** Wait for something to be served: the bus answers "no responders" at once until it is. */
async function until<T>(attempt: () => Promise<T>, tries = 50): Promise<T> {
for (let i = 0; ; i++) {
try {
return await attempt();
} catch (e) {
if (i >= tries) throw e;
await new Promise((r) => setTimeout(r, 100));
}
}
}
test("MESH_TOOL_MODULES names modules and their entrypoints; a bare path is the credential's own module's", () => {
const have = servedModulesFrom(" alpha=/a/tools/index.js, beta=/b/one.js ,beta=/b/two.js, /mine/index.js ,node-tools=/own/x.js", "node-tools");
assert.deepEqual(have.serves, [
{ module: "alpha", entrypoints: ["/a/tools/index.js"] },
{ module: "beta", entrypoints: ["/b/one.js", "/b/two.js"] },
]);
// The credential's own module, named or bare, is the one-module form either way.
assert.deepEqual(have.moduleEntrypoints, ["/mine/index.js", "/own/x.js"]);
assert.throws(() => servedModulesFrom("=/nothing.js", undefined), /neither <module>=<entrypoint>/);
assert.deepEqual(servedModulesFrom("", "x"), { serves: [], moduleEntrypoints: [] });
});
test("the node's runtime serves five modules' bundles on one credential — two of them launched, one broken — and follows a re-issued membership", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const mesh = await aMesh();
// Three modules assigned to the machine: alpha answers for itself anywhere, beta only here and
// holds the node-shelf seat, gamma's bundle is broken.
await mesh.issue(membershipOf("alpha", "anchor", { plain: true }));
await mesh.issue(membershipOf("beta", "anchor", { seats: { "node-shelf": ["list", "clear"] } }));
await mesh.issue(membershipOf("gamma", "anchor"));
// Two more, launched rather than loaded (ADR 0187): delta is Python and holds the node-lamp seat;
// epsilon is TypeScript written against the protocol and marked executable.
await mesh.issue(membershipOf("delta", "anchor", { seats: { "node-lamp": ["on"] } }));
await mesh.issue(membershipOf("epsilon", "anchor"));
// The node's credential: the runtime module's name, no claims (seats come from the memberships).
const credential = { url, node: "anchor", module: "node-tools" };
const nodeTools = 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 {
process.env.MESH_OPERATOR_ACCOUNT = "somebody";
process.env.MESH_OPERATOR_HOME = "/home/somebody";
stop = await runTools({
broker: nodeTools,
credential,
serves: [
{ module: "alpha", entrypoints: [fixture("many-alpha.mjs")] },
{ module: "beta", entrypoints: [fixture("many-beta.mjs")] },
{ module: "gamma", entrypoints: [fixture("many-broken.mjs")] },
{ module: "delta", entrypoints: [fixture("many-delta.py")] },
{ module: "epsilon", entrypoints: [fixture("many-epsilon.mjs")] },
],
});
console.log = log;
assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "delta", "epsilon", "gamma", "node-tools"]);
assert.ok(said.some((s) => /the operator's account here is somebody \(home \/home\/somebody\)/.test(s)), said.join("\n"));
assert.ok(said.some((s) => /gamma's bundle .*many-broken\.mjs failed to load: gamma's bundle cannot find its client; its tools are not served here/.test(s)), said.join("\n"));
assert.ok(said.some((s) => /serving 8 tool\(s\) for 5 module\(s\): alpha\.one, alpha\.two, beta\.three, beta\.four, beta\.five, delta\.greet, delta\.die, epsilon\.seven; not serving gamma/.test(s)), said.join("\n"));
// Five tools answer, each where its module's membership says: alpha anywhere and here, beta here only.
assert.deepEqual((await callTool(asker, "alpha.one", {})).result, { alpha: 1 });
assert.deepEqual((await callTool(asker, "alpha.one@anchor", {})).result, { alpha: 1 });
assert.deepEqual((await callTool(asker, "beta.three@anchor", {})).result, { beta: 3 });
assert.deepEqual((await callTool(asker, "beta.four@anchor", {})).result, { beta: 4 });
assert.deepEqual((await callTool(asker, "beta.five@anchor", {})).result, { beta: 5 });
await assert.rejects(callTool(asker, "beta.three", {}), /no responders|503/i, "beta was not issued the module's plain subject");
// Two seat verbs answer on the seat's subjects, held by beta.
assert.deepEqual((await callTool(asker, "seat:node-shelf.list@anchor", {})).result, { shelf: ["a", "b"] });
assert.deepEqual((await callTool(asker, "seat:node-shelf.clear@anchor", {})).result, { cleared: true });
// A bundle in another language answers the same way, its seat verb among them; so does a
// TypeScript bundle served through the protocol rather than imported.
assert.deepEqual((await callTool(asker, "delta.greet@anchor", { who: "mesh" })).result, { greeting: "hello mesh", language: "python" });
assert.deepEqual((await callTool(asker, "seat:node-lamp.on@anchor", {})).result, { on: true, language: "python" });
assert.deepEqual((await callTool(asker, "epsilon.seven@anchor", {})).result, { epsilon: 7, via: "stdio" });
// A child that exits mid-call tells the caller so and is started again on the next call.
await assert.rejects(callTool(asker, "delta.die@anchor", {}), /delta's bundle exited \(3\)/);
assert.deepEqual((await callTool(asker, "delta.greet@anchor", {})).result, { greeting: "hello world", language: "python" });
// A tool that emits does so as its module, not as the runtime.
const landed = mesh.nextEvent("mesh.mod.*.event.>");
assert.deepEqual((await callTool(asker, "alpha.two", {})).result, { alpha: 2 });
assert.equal(await landed, "mesh.mod.alpha.event.happened");
// `tools` answers for each: what alpha and beta serve, and why gamma serves nothing.
const gamma = await asker.request<Record<string, never>, { module: string; tools: unknown[]; failed?: string }>("gamma.tools@anchor", {});
assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle cannot find its client" });
const beta = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>("beta.tools@anchor", {});
assert.deepEqual(beta.tools.map((x) => x.name), ["three", "four", "five"]);
assert.deepEqual(beta.tools[0]!.subjects, ["mesh.mod.beta.tool.three.anchor"]);
// And discovery says so, with the reason, beside the modules that answered.
const catalogue = await connectNats({ url, module: "mesh-catalog" });
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "alpha" }, { module: "beta" }, { module: "gamma" }, { module: "delta" }, { module: "epsilon" }] }));
try {
const have = await toolsOn(asker);
assert.deepEqual(have.tools.map((x) => `${x.module}.${x.name}`), ["alpha.one", "alpha.two", "beta.five", "beta.four", "beta.three", "delta.die", "delta.greet", "epsilon.seven"]);
assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle cannot find its client)", "mesh-controller (seat)"]);
} finally {
await catalogue.close();
}
// The mesh re-issues beta's membership mid-run — now answering for the module anywhere — and
// the runtime serves the new subject without a restart.
await mesh.issue(membershipOf("beta", "anchor", { plain: true, seats: { "node-shelf": ["list", "clear"] } }));
assert.deepEqual((await until(() => callTool(asker, "beta.three", {}))).result, { beta: 3 });
assert.deepEqual((await callTool(asker, "seat:node-shelf.list@anchor", {})).result, { shelf: ["a", "b"] });
} finally {
console.log = log;
delete process.env.MESH_OPERATOR_ACCOUNT;
delete process.env.MESH_OPERATOR_HOME;
stop();
await asker.close();
await nodeTools.close();
await mesh.close();
resetTools();
}
});