Author SHA1 Message Date
jochen d4a4f4d1ba One SDK per runtime: a bundle's SDK import resolves to the runtime's copy (hq issue 209)
A bundle carries its dependencies, the SDK among them; imported in-process that copy was a
second SDK with its own tool registry, so a bundle registered its tools into a list the
runtime never read and served nothing, silently. A resolve hook now sends every import of
@novox/mesh-sdk, from whichever bundle, to the runtime's own copy: one registry, one broker.
The test loads a bundle from a directory holding its own SDK copy and sees its tool served.
2026-10-03 12:44:21 +02:00
mesh-admin 5d488a45f7 Merge pull request 'node-tools is a module beside mesh-tools: the runtime as a bundle, and serve is the console (hq ADR 0175, to-be 38 WP3)' (#31) from feat/wp3-node-tools into main 2026-10-02 19:57:24 +00:00
jochen c46f9502ee node-tools is a module beside mesh-tools: the runtime as a bundle, and serve is the console (hq ADR 0175, to-be 38 WP3)
One repository, two modules (ADR 0069). `node-tools/` holds the runtime — its code, tests, package
and the manifest of the module the controller composes a process for on every machine it is
assigned to: a bundle of `src/main.js`, the interpreter as a package, a place for the node's
credential, the loopback port the console declared, and leave to call every tool. Nothing about
how it runs: which bundles to load, where the credential is and whose machine it is are the
controller's to compose (WP2). The root module `mesh-tools` keeps the two images TypeScript
bundles are compiled in and a module's own service may run in; it is no longer how tools reach a
node.

As node-tools, `serve` is also the console (ADR 0175 §6): the same process answers MCP on
loopback for whoever is on the machine, through which the tools it serves can be called. A
module's own runtime in a container keeps serving without a listener.

The toolchain image now carries /app/runtime — a package.json saying the compiled files are ES
modules and the production node_modules — for the builder to copy into every TypeScript bundle,
so a bundle unpacked on a machine starts (ADR 0188 §5; the builder's side is the controller's).
Proven here by compiling node-tools with the toolchain's exact flags and starting the result.
The AMQP probe script is gone with the bus it probed.
2026-10-02 21:43:37 +02:00
mesh-admin 2746fd31f0 Merge pull request 'Cite ADR 0188, not 0187: the record was renumbered before it merged' (#30) from fix/cite-adr-0188 into main 2026-10-02 19:29:30 +00:00
jochen 82d7306ee0 Cite ADR 0188, not 0187: the record was renumbered before it merged (0187 is the dead-tracker record) 2026-10-02 21:29:04 +02:00
mesh-admin 2818f99b17 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 2026-10-02 19:27:49 +00:00
jochen 1436b02755 A bundle that is not JavaScript is launched and spoken to over MCP on stdio (hq ADR 0187, to-be 38 WP1b)
The runtime imported a bundle into its own process, which only JavaScript can be. Now an entrypoint
that is not a plain JavaScript file — or is one marked executable — is started as a child with the
runtime's environment and asked `tools/list` once and `tools/call` per call; what it lists is
registered exactly as an imported bundle's registrations are, a `<seat>.<verb>` name as the seat's
implementation. So a tools bundle may be in any language, and the mesh's part — the subjects, the
seats, the `tools` answer, a failed bundle named — stays in the runtime and is shared by all of
them. A child that exits mid-call tells the caller so and is started again on its next call.

Proven against a real bus beside the three bundles already there: a Python bundle with no SDK at
all answers its tool and its seat verb; a TypeScript bundle written against the protocol and marked
executable is served through the launcher, shortcut off; a bundle told to exit is relaunched.
2026-10-02 21:24:20 +02:00
mesh-admin 4e0559c31a Merge pull request 'Cite hq ADR 0170, not 0169: the firewall seat's record was renumbered' (#28) from fix/adr-0170-cited into main 2026-10-02 16:44:07 +00:00
jochen 6390d1d7fb The runtime serves a list of modules on one credential, naming a bundle that fails to load (hq ADR 0175, to-be 38 WP1)
One tool runtime per node, host-side, is what the runtime was written to be; the catalogue built a
container per module around it instead. This lets `serve` take a list — MESH_TOOL_MODULES as
<module>=<entrypoint> entries — and do for every assigned module what it did for one: read that
module's membership and follow it, serve its tools where the membership says, serve each held
seat's verbs on the seat's subjects. The seats come from the memberships now, so the node's
credential carries no claims; a module's own runtime still reads its credential's, so nothing
built today changes behaviour. A bare path in MESH_TOOL_MODULES stays the one-module form.

A bundle that throws on import is said in the log and in what `tools` answers for its module
(`failed`), which discovery lists with the reason instead of as "not answering"; the other bundles
serve. The filter that dropped every registration under a name but the one module goes; what
stays is that a registration under a seat's name is served only where some served module claims
the seat. A tool runs attributed to its module, so an event it emits lands on the module's
subject and not the runtime's. MESH_OPERATOR_ACCOUNT and MESH_OPERATOR_HOME are read and said;
tools take them from their environment.

Proven against a real bus: three bundles, one broken; five tools and two seat verbs answer on
their subjects; `tools` names the failed bundle; a membership re-issued mid-run re-serves.
2026-10-02 18:22:30 +02:00
jschoubben 020003ea3f Cite hq ADR 0170, not 0169: the firewall seat's record was renumbered after a collision on hq main 2026-10-02 14:52:28 +02:00
mesh-admin 809f07d085 Merge pull request 'A node-scoped seat's verb is callable through the console, naming its machine (hq ADR 0169)' (#27) from fix/a-node-scoped-seats-verb-is-callable-through-the-console into main 2026-10-02 12:24:02 +00:00
39 changed files with 1184 additions and 303 deletions
+1
View File
@@ -1,2 +1,3 @@
node_modules/ node_modules/
dist/ dist/
.mesh-build/
+25 -16
View File
@@ -1,5 +1,8 @@
ARG NODE_BASE=node:22-bookworm-slim ARG NODE_BASE=node:22-bookworm-slim
# Three stages, two published images: the one modules are COMPILED in, and the one they RUN in. # The mesh-tools module: the two images every TypeScript module is COMPILED in and may RUN in. The
# runtime itself ships as the node-tools module's bundle (node-tools/, novox/hq ADR 0175, to-be 38
# WP3); these images are the toolchain for TypeScript bundles and the base a module's own service
# may still be built on. They are no longer how tools reach a node.
# #
# **They were the same image, and that was a mistake.** A module's recipe starts from this and # **They were the same image, and that was a mistake.** A module's recipe starts from this and
# invokes the compiler out of it, so the compiler had to be here — and because the same image was # invokes the compiler out of it, so the compiler had to be here — and because the same image was
@@ -19,36 +22,42 @@ RUN apt-get update \
&& apt-get install -y --no-install-recommends git ca-certificates \ && apt-get install -y --no-install-recommends git ca-certificates \
&& rm -rf /var/lib/apt/lists/* && rm -rf /var/lib/apt/lists/*
WORKDIR /app WORKDIR /app
COPY package.json ./ COPY node-tools/package.json ./
# The builder writes .npmrc into the build context; it authenticates to the mesh's package registry # The builder writes .npmrc into the build context; it authenticates to the mesh's package registry
# for the @novox scope, which is where @novox/mesh-sdk resolves. This stage is not published, so the # for the @novox scope, which is where @novox/mesh-sdk resolves. This stage is not published, so the
# credential travels no further than here. Development dependencies included: the compiler is one. # credential travels no further than here. Development dependencies included: the compiler is one.
COPY .npmrc ./.npmrc COPY .npmrc ./.npmrc
RUN npm install --no-audit --no-fund RUN npm install --no-audit --no-fund
# ---- toolchain: what a module is compiled in, WITHOUT the credential -------------------------- # ---- compiling: the runtime's own code built, WITHOUT the credential -------------------------
FROM ${NODE_BASE} AS toolchain FROM ${NODE_BASE} AS compiling
WORKDIR /app WORKDIR /app
COPY package.json ./ COPY node-tools/package.json ./
# The resolved libraries, but not the .npmrc that resolved them. # The resolved libraries, but not the .npmrc that resolved them.
COPY --from=deps /app/node_modules ./node_modules COPY --from=deps /app/node_modules ./node_modules
# The toolkit arrives compiled. It used to arrive as sources, and this compiled it by hand — the COPY node-tools/tsconfig.json ./
# hook that builds it on install was running all along, and the result was then packed out of the COPY node-tools/src ./src
# package, because with no explicit file list npm falls back to .gitignore and that ignores the
# build output. Fixed where it belonged, in the toolkit.
COPY tsconfig.json ./
COPY src ./src
RUN npm run build RUN npm run build
# ---- what the running image needs, and nothing else ------------------------------------------- # ---- what a running bundle needs, and nothing else -------------------------------------------
# Its own stage so the toolchain image keeps its build tools while the runtime image does not. # Its own stage so the toolchain image keeps its build tools while the runtime image does not.
FROM toolchain AS lean FROM compiling AS lean
RUN npm prune --omit=dev RUN npm prune --omit=dev
# ---- runtime: what a module runs in ----------------------------------------------------------- # ---- toolchain: what a TypeScript bundle is compiled in ---------------------------------------
# Beside the compiler, at /app/runtime, what every TypeScript bundle runs with: the production
# dependencies the SDK and the runtime need, and a package.json saying the compiled files are ES
# modules. The builder copies this directory whole into a compiled bundle (novox/hq ADR 0188 §5),
# so a bundle unpacked on a machine starts — a `.js` without that package.json is read as
# CommonJS, and an import of `nats` without node_modules beside it resolves to nothing.
FROM compiling AS toolchain
COPY --from=lean /app/node_modules /app/runtime/node_modules
RUN printf '{"type":"module","private":true}\n' > /app/runtime/package.json
# ---- runtime: what a module's own service may run in ------------------------------------------
FROM ${NODE_BASE} AS runtime FROM ${NODE_BASE} AS runtime
WORKDIR /app WORKDIR /app
COPY package.json ./ COPY node-tools/package.json ./
COPY --from=lean /app/node_modules ./node_modules COPY --from=lean /app/node_modules ./node_modules
COPY --from=toolchain /app/dist ./dist COPY --from=compiling /app/dist ./dist
ENTRYPOINT ["node", "dist/main.js"] ENTRYPOINT ["node", "dist/main.js"]
+47 -17
View File
@@ -1,29 +1,56 @@
# mesh-tools # mesh-tools
The Novox Mesh **tool runtime** — the per-node process that makes a module's tools actually serve. Two modules in one repository (novox/hq ADR 0069), one piece of software:
A module ships its tools (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk)); this - **`node-tools`** (`node-tools/`) — the node's **tool runtime** as a module (ADR 0175, to-be 38
runtime is what loads them and puts them on the mesh. It: WP3): one process per machine the host runs from this bundle, serving every assigned module's tools
and every held seat's verbs on the bus, and answering MCP on the machine's loopback — the console
(design 34). The code, its tests and the `mesh` client all live there.
- **`mesh-tools`** (this directory) — the two images TypeScript bundles are compiled in and a module's
own *service* may still run in. Built from the same code; no longer how tools reach a node.
1. connects the mesh broker (novox/hq ADR 0001) — a concrete AMQP implementation of the sdk's The runtime:
`Broker` contract;
2. imports the assigned modules' compiled tool entrypoints, each of which registers its tools as it 1. connects the mesh bus on the node's credential — a concrete implementation of the sdk's `Broker`
loads; contract;
3. serves them through the sdk's `serveTools` harness, answering `tools.invoke` over the broker. 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 0188): 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 (`node-tools/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 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, wrapper that binds the bus and loads the modules. Keeping the bus client here, out of the sdk, is
is deliberate: a broker-client change never rebuilds a module (ADR 0039). 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 ## Running it
``` ```
MESH_BROKER_URL amqp://… the mesh broker MESH_BROKER_FILE the node's sealed credential, as the mesh delivered it
MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool entrypoints 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 `node dist/main.js`. On a node the controller composes the variables and the host supervises the
and starts it like any other supervised workload. process like any other host-side workload (novox/hq to-be 38). As `node-tools` the same process is
the console: MCP on `127.0.0.1:4270` (or `MESH_CONSOLE_LISTEN`). 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 ## `mesh` — the tools for whoever is on a machine
@@ -41,6 +68,9 @@ dropped. A module may not name a tool of its own `tools`; the runtime refuses it
## Verified ## Verified
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the `npm test` runs against a real NATS server with JetStream (`MESH_TEST_NATS`, see any test's header
runtime serves a registered tool, a separate connection invokes it by name and gets the result, and for the one-line `docker run`) and proves the whole path over the wire: the runtime serves a
an unknown tool is refused over the wire. 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.
-48
View File
@@ -1,48 +0,0 @@
import amqp from "amqplib";
const PORT = process.argv[2];
const MPORT = process.argv[3];
const B = `http://127.0.0.1:${MPORT}`;
const AUTH = "Basic " + Buffer.from("guest:guest").toString("base64");
async function api(method, path, body) {
const r = await fetch(B + path, {
method,
headers: { "content-type": "application/json", authorization: AUTH },
body: body ? JSON.stringify(body) : undefined,
});
if (r.status >= 300 && r.status !== 404) throw new Error(`${method} ${path} -> ${r.status}`);
}
await api("PUT", "/api/exchanges/%2f/mesh.events.dead", { type: "topic", durable: true });
await api("PUT", "/api/users/al", { password: "s", tags: "" });
const Q = "anchor.al.events";
const D = "mesh.events.dead";
const q = Q.replace(/\./g, "\\.");
const d = D.replace(/\./g, "\\.");
// configure, write, read patterns per grant on the dead exchange
const combos = {
"none": { configure: `^${q}$`, write: `^${q}$`, read: `^${q}$` },
"read-dead": { configure: `^${q}$`, write: `^${q}$`, read: `^(${q}|${d})$` },
"write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^${q}$` },
"configure-dead": { configure: `^(${q}|${d})$`, write: `^${q}$`, read: `^${q}$` },
"read+write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^(${q}|${d})$` },
};
let i = 0;
for (const [label, perms] of Object.entries(combos)) {
await api("PUT", "/api/permissions/%2f/al", perms);
const queue = `${Q}.${i++}`; // fresh each time
try {
const c = await amqp.connect(`amqp://al:s@127.0.0.1:${PORT}/`);
const ch = await c.createChannel();
ch.on("error", () => {});
await ch.assertQueue(queue, { durable: true, deadLetterExchange: D });
console.log(`${label}: declare-with-DLX OK`);
await c.close();
} catch (e) {
console.log(`${label}: FAIL - ${String(e.message).slice(0, 70)}`);
}
}
+11
View File
@@ -0,0 +1,11 @@
# node-tools
The node's tool runtime as a module (novox/hq ADR 0175, to-be 38 WP3). Assigned to a machine, it is
one process the host runs from this bundle, as the operator's account: it serves every assigned
module's tools and every held seat's verbs on the bus, and answers MCP on the machine's loopback —
the console (design 34). The controller composes the process (which bundles to load, where the
credential is, whose machine it is); this manifest says only what the machine must have for it: the
interpreter, a place for the credential, the loopback port, and leave to call every tool.
The code is the `mesh-tools` package in this directory; the module at the repository root,
`mesh-tools`, builds the images TypeScript bundles are compiled in. See the repository README.
+45
View File
@@ -0,0 +1,45 @@
{
"module": "node-tools",
"version": "1",
"slug": "node-tools",
"invokes": [
"*"
],
"own-secrets": {
"broker": "${dir:mesh-state}/broker"
},
"listens": [
{
"name": "mcp",
"port": 4270,
"protocol": "tcp",
"from": "machine",
"why": "the mesh's tools for whoever is on this machine, over MCP on loopback; the machine's login is the authority (novox/hq ADR 0152, 0175)"
}
],
"resources": [
{
"id": "mesh-state",
"type": "directory",
"mode": "0755",
"place": "mesh"
},
{
"id": "interpreter",
"type": "package",
"package": "nodejs"
}
],
"build": {
"artifacts": [
{
"name": "runtime",
"kind": "bundle",
"language": "typescript",
"entrypoints": [
"src/main.js"
]
}
]
}
}
View File
@@ -16,6 +16,7 @@
// derives the subject (design 29 §1), so reorganising the subject space leaves every module // derives the subject (design 29 §1), so reorganising the subject space leaves every module
// correct. // correct.
import { AsyncLocalStorage } from "node:async_hooks";
import { createHash } from "node:crypto"; import { createHash } from "node:crypto";
import net from "node:net"; import net from "node:net";
import tls from "node:tls"; import tls from "node:tls";
@@ -71,19 +72,41 @@ export interface Answered<Res> {
node?: string; node?: string;
} }
/** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs — /** 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, and serving a subject that is not a module's own tool * an answer that says which machine gave it, serving a subject that is not a module's own tool
* (a seat's verb). */ * (a seat's verb), and following the memberships of the modules it serves beside its own. */
export interface RuntimeBroker extends Broker { export interface RuntimeBroker extends Broker {
/** Call a tool by key, or — when `on` names a subject the mesh listed for it (ADR 0160) — there. */ /** 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>>; ask<Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>>;
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void>; 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; * Serve another module's tools from this connection: read its membership on this machine and
/** Called when the mesh issues a new membership; the runtime re-serves on it. */ * 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; 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"; const ASSIGNMENTS_STREAM = "ASSIGNMENTS";
/** The one address a runtime derives for itself (ADR 0160). */ /** The one address a runtime derives for itself (ADR 0160). */
@@ -153,38 +176,45 @@ export async function connectNats(
const node = cred.node; const node = cred.node;
// The membership, read once at connect and followed. A direct get is one request on the // The memberships, one per module this connection serves, each read once and followed. A
// stream's API, which is the whole of what this account may ask JetStream for its own subject; // module's own runtime follows one, its own; the node's runtime follows one per assigned module
// a 404 is a mesh that has not issued one, which is a fact to say and not an error to retry. // (novox/hq ADR 0175) — the same subject shape, the same stream, read as many times as there are
let issued: Membership | undefined; // 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 issuedHandlers: ((m: Membership) => void)[] = [];
const subjectOfMine = node ? membershipSubject(node, self) : ""; const follow = async (module: string): Promise<void> => {
if (subjectOfMine) { if (!node || issued.has(module)) return;
issued.set(module, undefined);
const subjectOfTheirs = membershipSubject(node, module);
try { try {
// The subject-addressed form of a direct get — the stream, then the subject, nothing in // 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 // 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. // 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 }); new Uint8Array(0), { timeout: 5_000 });
const status = got.headers?.code ?? 0; const status = got.headers?.code ?? 0;
if (status === 0 && got.data.length > 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 { } catch {
// Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below. // Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below.
} }
if (!issued) { if (!issued.get(module)) {
console.log(`[mesh-tools] no membership issued for ${self} on ${node} yet; serving the derived shape until one arrives`); console.log(`[mesh-tools] no membership issued for ${module} on ${node} yet; serving the derived shape until one arrives`);
} }
try { try {
const live = conn.subscribe(subjectOfMine); const live = conn.subscribe(subjectOfTheirs);
subs.push(live); subs.push(live);
void (async () => { void (async () => {
for await (const msg of live) { for await (const msg of live) {
try { try {
issued = JSON.parse(sc.decode(msg.data)) as Membership; const m = JSON.parse(sc.decode(msg.data)) as Membership;
console.log(`[mesh-tools] ${self} on ${node} was issued a new membership; re-serving on it`); issued.set(module, m);
for (const h of issuedHandlers) h(issued); 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) { } catch (err) {
console.log(`[mesh-tools] a membership arrived that is not one: ${err}`); console.log(`[mesh-tools] a membership arrived that is not one: ${err}`);
} }
@@ -193,32 +223,33 @@ export async function connectNats(
} catch { } catch {
// A subscription this account may not make is a mesh older than the membership. // 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 /** The subjects a tool of a served module is served on: from that module's membership when
* otherwise (the shape the mesh issues on day one, so the two agree). */ * issued, derived otherwise (the shape the mesh issues on day one, so the two agree). */
const servedOn = (tool: string): { subject: string; queue?: string }[] => { const servedOn = (module: string, tool: string): { subject: string; queue?: string }[] => {
if (issued) { const m = issued.get(module);
const m = issued; if (m) {
const out = m.serves.map((s) => ({ subject: s.subject.replace("{tool}", tool), queue: s.queue })); 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 // 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. // 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)) { 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; return out;
} }
const base = `mesh.mod.${self}.tool.${tool}`; const base = `mesh.mod.${module}.tool.${tool}`;
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }]; const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${module}` }];
if (node) out.push({ subject: `${base}.${node}` }); if (node) out.push({ subject: `${base}.${node}` });
return out; return out;
}; };
/** Where a call by key goes: a subject the membership says this module reaches, when it says /** Where a call by key goes: a subject this connection's own membership says it reaches, when it
* one — the machine's when named — else the derived shape. */ * says one — the machine's when named — else the derived shape. */
const reachedAt = (key: string): string => { const reachedAt = (key: string): string => {
const [name, wanted] = key.split("@", 2); const [name, wanted] = key.split("@", 2);
const reach = issued?.reaches?.[name]; const reach = issued.get(self)?.reaches?.[name];
if (reach && reach.length > 0) { if (reach && reach.length > 0) {
if (wanted) { if (wanted) {
const at = reach.find((s) => s.endsWith(`.${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> { 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 // 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 // knows its role without a membership; everything else is a served module's own tool, served
// where the mesh issued it (ADR 0160). A key naming another module is not served here at all. // 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:")) { if (key.startsWith("seat:")) {
const stop = answerOn(toolSubject(key, self), undefined, handler); const stop = answerOn(toolSubject(key, self), undefined, handler);
return () => stop(); return () => stop();
} }
const dot = key.indexOf("."); const dot = key.indexOf(".");
if (dot >= 0 && key.slice(0, dot) !== self) { const module = dot < 0 ? self : key.slice(0, dot);
throw new Error(`${self} cannot serve ${key}: a module serves its own tools`); 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); const tool = dot < 0 ? key : key.slice(dot + 1);
let stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler)); let stops = servedOn(module, 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. // When that module's new membership arrives, serve where it now says and stop serving where
issuedHandlers.push(() => { // it no longer does.
issuedHandlers.push((m) => {
if (m.module !== module) return;
stops.forEach((stop) => stop()); 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()); 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) => { onMembership: (handler: (m: Membership) => void) => {
issuedHandlers.push(handler); issuedHandlers.push(handler);
}, },
@@ -339,7 +381,9 @@ export async function connectNats(
if (!meta["content-type"]) h.set("content-type", "application/json"); if (!meta["content-type"]) h.set("content-type", "application/json");
if (env.node) h.set("x-node", env.node); 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, headers: h,
// De-duplicated by the server inside its window, so a redelivery after a crash between // 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. // publishing and acknowledging is not seen twice. Only the emitter can make this id.
+7 -2
View File
@@ -182,7 +182,7 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
for (const s of roles.seats ?? []) { for (const s of roles.seats ?? []) {
// A node-scoped seat's tool is asked of one machine (design 33 §4): listed with its scope, so // A node-scoped seat's tool is asked of one machine (design 33 §4): listed with its scope, so
// a caller names the machine and the call carries it — `seat:<seat>.<verb>@<node>`. Left out // a caller names the machine and the call carries it — `seat:<seat>.<verb>@<node>`. Left out
// of the listing, the verb never resolved as a seat's and nothing served it (ADR 0169). // of the listing, the verb never resolved as a seat's and nothing served it (ADR 0170).
for (const t of s.tools ?? []) { for (const t of s.tools ?? []) {
tools.push({ module: s.seat, name: t.name, description: t.description, input: t.input, seat: true, scope: s.scope }); tools.push({ module: s.seat, name: t.name, description: t.description, input: t.input, seat: true, scope: s.scope });
} }
@@ -192,7 +192,12 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
} }
asked.forEach((outcome, i) => { asked.forEach((outcome, i) => {
const module = names[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) { for (const t of outcome.value.tools) {
tools.push({ module, name: t.name, description: t.description, input: t.input, subjects: t.subjects }); 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 0188).
//
// 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;
},
};
}
+64 -10
View File
@@ -1,7 +1,7 @@
// The runnable entrypoint. Three modes: // The runnable entrypoint. Three modes:
// //
// mesh-tools serve — bind the broker and serve the assigned modules until // mesh-tools serve — bind the broker and serve the assigned modules until
// stopped. A module entrypoint that subscribes to events (on("#")) // stopped; as the node-tools module, also the console on loopback. A module entrypoint that subscribes to events (on("#"))
// starts consuming as it is imported, so this also runs consumers. // starts consuming as it is imported, so this also runs consumers.
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive, // mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
// and what an events test uses to put a message on the wire. // and what an events test uses to put a message on the wire.
@@ -21,14 +21,16 @@
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0043) — an // 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. // amqps account scoped to this module. Preferred: a module holds its own.
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account. // MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode) // MESH_TOOL_MODULES <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_PREPARE /path/a,/path/b,… compiled entrypoints that prepare this module's state
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042) // MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
import { readFileSync } from "node:fs"; import { readFileSync } from "node:fs";
import { pathToFileURL } from "node:url"; import { pathToFileURL } from "node:url";
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js"; import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js";
import { runTools } from "./runtime.js"; import { runTools, type ServedModule } from "./runtime.js";
import { serveMcpHttp, type Listening } from "./http.js";
/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */ /** The credential this process connected with, for what it says beyond the connection (ADR 0159). */
let lastCredential: Credential | undefined; let lastCredential: Credential | undefined;
@@ -118,17 +120,66 @@ async function connectBrokerPatiently(): Promise<Broker> {
} }
} }
async function serve(): Promise<void> { /**
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "") * What MESH_TOOL_MODULES names (novox/hq ADR 0175, to-be 38 WP1): `<module>=<entrypoint>` entries,
.split(",") * comma-separated, several per module allowed — the node's runtime serving every assigned module's
.map((s) => s.trim()) * bundle. A bare path is the one-module form the per-module containers still set: an entrypoint of
.filter(Boolean); * 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 };
}
/** The module that is the node's tool runtime (novox/hq ADR 0175, to-be 38 WP3): on its credential,
* `serve` is also the console — MCP on the machine's loopback (design 34). */
export const RUNTIME_MODULE = "node-tools";
/** Where the console listens when the runtime is node-tools and nothing says otherwise: the port
* the module's manifest declares `from: machine`. MESH_CONSOLE_LISTEN overrides it either way. */
const CONSOLE_LISTEN = "127.0.0.1:4270";
async function serve(): Promise<void> {
const broker = await connectBrokerPatiently(); 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 });
// The console is this runtime's serving mode (ADR 0175 §6): as node-tools, or wherever the
// listen address is given, the same process answers MCP on loopback for whoever is on the
// machine. A module's own runtime in a container on the machine's network does not — two of
// them on one port would be the fault, and the console is one per machine.
const listen = process.env.MESH_CONSOLE_LISTEN ?? (lastCredential?.module === RUNTIME_MODULE ? CONSOLE_LISTEN : "");
let consoleUp: Listening | undefined;
if (listen) {
const who = `${lastCredential?.node ?? "?"}.${lastCredential?.module ?? RUNTIME_MODULE}`;
consoleUp = await serveMcpHttp(broker, who, listen);
console.log(`mesh console listening on http://${consoleUp.address}/mcp as ${who}`);
}
const shutdown = async (): Promise<void> => { const shutdown = async (): Promise<void> => {
stop(); stop();
await consoleUp?.close();
await broker.close(); await broker.close();
process.exit(0); process.exit(0);
}; };
@@ -250,4 +301,7 @@ async function main(): Promise<void> {
await serve(); 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();
}
+1 -1
View File
@@ -143,7 +143,7 @@ export function mcpSurface(bus: Broker, who: string): Surface {
const roles = have && seatsIn(have); const roles = have && seatsIn(have);
const bare = given.split("@", 1)[0]; const bare = given.split("@", 1)[0];
const isSeatVerb = roles ? toolKey(bare, roles).startsWith("seat:") : false; const isSeatVerb = roles ? toolKey(bare, roles).startsWith("seat:") : false;
// A node-scoped seat's verb is asked of one machine (design 33 §4, ADR 0169): `node` // A node-scoped seat's verb is asked of one machine (design 33 §4, ADR 0170): `node`
// names it and travels in the subject, as for a module's tool. // names it and travels in the subject, as for a module's tool.
const nodeScoped = isSeatVerb && (have?.tools.some((t) => t.seat && t.scope === "node" && const nodeScoped = isSeatVerb && (have?.tools.some((t) => t.seat && t.scope === "node" &&
`${t.module}.${t.name}` === bare) ?? false); `${t.module}.${t.name}` === bare) ?? false);
+327
View File
@@ -0,0 +1,327 @@
// 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 { register } from "node:module";
import { useBroker } from "@novox/mesh-sdk/messaging";
import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools";
import type { Broker } from "@novox/mesh-sdk/messaging";
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
* module's tool names, descriptions and argument schemas, from the code that answers them and
* from nowhere else. Discovery asks the module, because a copy kept anywhere else drifts.
*/
export const TOOLS_VERB = "tools";
/** What `tools` answers for one module. */
export interface ToolsAnswer {
module: string;
tools: {
name: string;
description: string;
input: Readonly<Record<string, unknown>>;
/** Where this tool is answered, as the mesh issued it (ADR 0160): the module's plain subject
* first when there is one, then this machine's. A caller composes nothing. */
subjects?: string[];
}[];
/** 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;
/** 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 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);
// 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]);
}
// 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 0188); what it lists is registered the same way.
// Every bundle's import of the SDK resolves to this runtime's copy (04-ISSUES/209): one registry
// of tools, one broker. Installed before the first bundle is imported.
oneSdk();
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`);
}
}
}
// 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.some((t) => t.name === TOOLS_VERB)) {
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 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);
}
}
// 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). 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,
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, 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 = [];
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());
}
let sdkHooked = false;
/** Resolve `@novox/mesh-sdk` for every bundle to the copy this runtime imported (sdk-hooks.ts). */
function oneSdk(): void {
if (sdkHooked) return;
sdkHooked = true;
register("./sdk-hooks.js", { parentURL: import.meta.url, data: { runtimeURL: import.meta.url } });
}
+38
View File
@@ -0,0 +1,38 @@
// The one SDK in a node's runtime (novox/hq ADR 0175, 04-ISSUES/209).
//
// A bundle carries its own dependencies — the toolchain copies them in so a bundle starts anywhere
// (to-be 38 WP3) — and among them is a copy of @novox/mesh-sdk. Imported in this process, that copy
// would be a second SDK: its own registry of tools, its own broker handle. A bundle calling
// registerModuleTools through it registers into a list this runtime never reads, and its tools are
// silently not served. So every import of the SDK, from whichever bundle, is resolved as if this
// runtime had written it: one registry, one broker — the runtime's. Everything else a bundle
// carries resolves from the bundle's own tree, as before.
//
// Installed through module.register(), whose resolve hook sees every import that follows.
interface ResolveContext {
parentURL?: string;
conditions: string[];
importAttributes: Record<string, string>;
}
interface Resolved {
url: string;
format?: string | null;
shortCircuit?: boolean;
}
type NextResolve = (specifier: string, context?: Partial<ResolveContext>) => Promise<Resolved>;
const SDK = "@novox/mesh-sdk";
let runtimeURL = "";
/** Told, once, which file's tree holds the runtime's SDK. */
export function initialize(data: { runtimeURL: string }): void {
runtimeURL = data.runtimeURL;
}
export async function resolve(specifier: string, context: ResolveContext, next: NextResolve): Promise<Resolved> {
if (runtimeURL && (specifier === SDK || specifier.startsWith(SDK + "/"))) {
return next(specifier, { ...context, parentURL: runtimeURL });
}
return next(specifier, context);
}
@@ -68,7 +68,7 @@ test("a person sees what the running modules answer, sorted, and who did not ans
"the list is what the modules answered plus every role's tools, in a stable order", "the list is what the modules answered plus every role's tools, in a stable order",
); );
// A role's tool is marked as one; a node-scoped seat's carries its scope, so a caller names // A role's tool is marked as one; a node-scoped seat's carries its scope, so a caller names
// the machine and the verb resolves as the seat's (design 33 §4, ADR 0169). // the machine and the verb resolves as the seat's (design 33 §4, ADR 0170).
assert.ok(have.tools.find((x) => x.module === "mesh-controller")!.seat); assert.ok(have.tools.find((x) => x.module === "mesh-controller")!.seat);
const lookup = have.tools.find((x) => x.module === "node-dns-resolver")!; const lookup = have.tools.find((x) => x.module === "node-dns-resolver")!;
assert.ok(lookup.seat && lookup.scope === "node"); assert.ok(lookup.seat && lookup.scope === "node");
+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");
+29
View File
@@ -0,0 +1,29 @@
#!/usr/bin/env python3
# A tools bundle in a second language (novox/hq ADR 0188): 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 0188).
import { serveStdio } from "@novox/mesh-sdk/stdio";
await serveStdio("epsilon", [
{ name: "seven", description: "epsilon's", input: {}, run: async () => ({ epsilon: 7, via: "stdio" }) },
]);
@@ -35,7 +35,7 @@ async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void)
{ name: "status", description: "what is wrong", input: {} }, { name: "status", description: "what is wrong", input: {} },
{ name: "push", description: "tell a machine", input: { node: { type: "string" } } }, { name: "push", description: "tell a machine", input: { node: { type: "string" } } },
] }, ] },
// A seat held once per machine (design 33 §4, ADR 0169): its verb is asked of one. // A seat held once per machine (design 33 §4, ADR 0170): its verb is asked of one.
{ seat: "node-dns-resolver", scope: "node", tools: [{ name: "lookup", description: "one machine's", input: {} }] }, { seat: "node-dns-resolver", scope: "node", tools: [{ name: "lookup", description: "one machine's", input: {} }] },
], ],
})); }));
@@ -107,7 +107,7 @@ test("a host initialises, lists the mesh's tools and calls one", async (t) => {
assert.deepEqual(listed.map((x: { name: string }) => x.name), assert.deepEqual(listed.map((x: { name: string }) => x.name),
["mesh-controller.push", "mesh-controller.status", "node-dns-resolver.lookup", "shop.price"], ["mesh-controller.push", "mesh-controller.status", "node-dns-resolver.lookup", "shop.price"],
"the modules' tools and the roles', named the way a person names them"); "the modules' tools and the roles', named the way a person names them");
// A node-scoped seat's verb takes the machine, and requires it (ADR 0169). // A node-scoped seat's verb takes the machine, and requires it (ADR 0170).
const lookup = listed[2]; const lookup = listed[2];
assert.equal(lookup.inputSchema.properties.node.type, "string"); assert.equal(lookup.inputSchema.properties.node.type, "string");
assert.deepEqual(lookup.inputSchema.required, ["node"]); assert.deepEqual(lookup.inputSchema.required, ["node"]);
+234
View File
@@ -0,0 +1,234 @@
/**
* 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 { cpSync, mkdtempSync, realpathSync, rmSync, writeFileSync } from "node:fs";
import { join } from "node:path";
import { tmpdir } from "node:os";
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 0188): 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();
}
});
test("a bundle carrying its own copy of the SDK registers into the runtime's registry, and its tools are served (issue 209)", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
// A bundle as the toolchain packs one: its compiled entrypoint, a package.json saying ES modules,
// and its dependencies copied in — the SDK among them, a second copy beside the runtime's own.
const dir = mkdtempSync(join(tmpdir(), "mesh-bundle-"));
const sdk = realpathSync(fileURLToPath(new URL("../node_modules/@novox/mesh-sdk/", import.meta.url)));
cpSync(sdk, join(dir, "node_modules", "@novox", "mesh-sdk"), { recursive: true });
writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n');
writeFileSync(join(dir, "index.js"),
'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' +
'registerModuleTools("zeta", () => [{ name: "probe", description: "answers", input: {}, run: async () => ({ zeta: true }) }]);\n');
const mesh = await aMesh();
await mesh.issue(membershipOf("zeta", "anchor"));
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 {
stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.js")] }] });
console.log = log;
assert.ok(said.some((s) => /serving 1 tool\(s\) for 1 module\(s\): zeta\.probe/.test(s)), said.join("\n"));
assert.deepEqual((await callTool(asker, "zeta.probe@anchor", {})).result, { zeta: true });
} finally {
console.log = log;
stop();
await asker.close();
await nodeTools.close();
await mesh.close();
rmSync(dir, { recursive: true, force: true });
resetTools();
}
});
+63
View File
@@ -0,0 +1,63 @@
/**
* The runtime as the node-tools module (novox/hq ADR 0175 §6, to-be 38 WP3): started the way the
* host starts it — `main.js` with the node's credential and MESH_TOOL_MODULES — it serves the bundles
* AND answers MCP on loopback as the console, through which a tool it serves can be called.
*
* 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-tools-serve.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 { fileURLToPath } from "node:url";
import { connectNats } from "../dist/broker-nats.js";
const url = process.env.MESH_TEST_NATS;
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
async function post(endpoint: string, body: unknown): Promise<any> {
const res = await fetch(endpoint, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(body) });
return res.json();
}
test("as node-tools, serve loads the bundles and is the console on loopback", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
// The catalogue, for discovery; the node credential names the runtime module and no claims.
const catalogue = await connectNats({ url, module: "mesh-catalog" });
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "alpha" }] }));
t.after(() => catalogue.close());
const dir = await mkdtemp("/tmp/node-tools-");
const credential = join(dir, "broker");
await writeFile(credential, JSON.stringify({ url, node: "desk", module: "node-tools", user: "desk.node-tools", password: "x" }));
const child: ChildProcess = spawn(process.execPath, ["dist/main.js"], {
env: {
...process.env,
MESH_BROKER_FILE: credential,
MESH_TOOL_MODULES: `alpha=${fixture("many-alpha.mjs")}`,
MESH_CONSOLE_LISTEN: "127.0.0.1:0",
},
stdio: ["ignore", "pipe", "pipe"],
});
t.after(() => {
child.kill("SIGTERM");
});
const endpoint = await new Promise<string>((resolve, reject) => {
let out = "";
let err = "";
child.stdout!.on("data", (d) => {
out += d.toString();
const m = /listening on (http:\/\/[^/]+\/mcp) as desk\.node-tools/.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}`)));
});
const listed = await post(endpoint, { jsonrpc: "2.0", id: 1, method: "tools/list" });
assert.deepEqual(listed.result.tools.map((x: any) => x.name).filter((n: string) => n.startsWith("alpha.")), ["alpha.one", "alpha.two"]);
// A tool the same process serves on the bus, called through the console it also is.
const called = await post(endpoint, { jsonrpc: "2.0", id: 2, method: "tools/call", params: { name: "alpha.one", arguments: { node: "desk" } } });
assert.deepEqual(JSON.parse(called.result.content[0].text), { alpha: 1 });
});
-164
View File
@@ -1,164 +0,0 @@
// 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.
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 type { Broker } from "@novox/mesh-sdk/messaging";
import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
/**
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
* module's tool names, descriptions and argument schemas, from the code that answers them and
* from nowhere else. Discovery asks the module, because a copy kept anywhere else drifts.
*/
export const TOOLS_VERB = "tools";
/** What `tools` answers for one module. */
export interface ToolsAnswer {
module: string;
tools: {
name: string;
description: string;
input: Readonly<Record<string, unknown>>;
/** Where this tool is answered, as the mesh issued it (ADR 0160): the module's plain subject
* first when there is one, then this machine's. A caller composes nothing. */
subjects?: string[];
}[];
}
export interface RuntimeOptions {
/** The mesh broker to serve over. */
broker: Broker;
/** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */
moduleEntrypoints: string[];
/** The credential the mesh delivered, for what it says about the seats this module claims
* (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat. */
credential?: Credential;
}
/** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */
export async function runTools(opts: RuntimeOptions): Promise<() => void> {
useBroker(() => opts.broker);
for (const entry of opts.moduleEntrypoints) {
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake.
await import(pathToFileURL(resolve(entry)).href);
}
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
// audit logger) registers none, and its scoped account may not declare the serve queue — so a
// runtime that always served would fail for exactly the modules that never needed it.
// A registration under a seat's name is the module's implementation of that seat's verbs
// (ADR 0159, 0160): served on the seat's subjects by serveClaimedSeats, never as a module's
// tools and never listed among them. Everything else is the module's own.
// A module named like its seat (the catalogue is the mesh-catalog seat) registers once and is
// both: its tools are the module's and the seat's verbs alike.
const self = opts.credential?.module;
const seatNames = new Set((opts.credential?.claims ?? []).map((c) => c.seat));
// A registration under a name that is neither this module nor a seat it claims is not served:
// said, and left out, rather than fatal — on 2026-10-01 the credential of a module that had just
// learned to implement a seat did not yet name the claim, and the whole runtime restarted for it.
const ownRegistrations = collectTools().filter(({ module }) => {
if (module === self || !self || seatNames.has(module)) return module === self || !self;
console.log(`[mesh-tools] ${self} registers tools under "${module}", which is neither this module nor a seat its credential claims; not served until the mesh issues the claim`);
return false;
});
const tools = ownRegistrations.flatMap(({ module, tools: own }) => own.map((t) => ({ module, name: t.name })));
const stops: Array<() => void> = [];
const stop = (): void => stops.splice(0).forEach((s) => s());
// Each tool on its own key, namespaced by its module (ADR 0047); where that key is answered is
// the broker's to know from the membership (ADR 0160).
for (const { module, tools: own } of ownRegistrations) {
const seen = new Set<string>();
for (const t of own) {
if (seen.has(t.name)) {
stop();
throw new Error(`${module} exposes two tools named ${t.name} — refused`);
}
seen.add(t.name);
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) => t.run(args ?? {})));
}
}
// And, for every module that serves any, the verb that says what it serves. Refused before
// anything is bound if a module named a tool of its own `tools`: one name answering two things
// is the fault nobody can diagnose afterwards, and the runtime is the only place that sees both.
const runtime = opts.broker as RuntimeBroker;
for (const { module, tools: own } of ownRegistrations) {
if (own.length === 0) continue;
if (own.some((t) => t.name === TOOLS_VERB)) {
stop();
throw new Error(
`${module} names a tool "${TOOLS_VERB}", which is the verb the runtime answers for every ` +
"module with what it serves (novox/hq ADR 0152) — refused, rename it",
);
}
const subjectsOf = (tool: string): string[] | undefined => {
const issued = typeof runtime.membership === "function" ? runtime.membership() : undefined;
if (!issued) return undefined;
const plain = issued.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
const mine = issued.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
return [...plain, ...mine];
};
const answer: ToolsAnswer = {
module,
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
};
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async () => answer));
}
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
stops.push(await serveClaimedSeats(opts.broker as RuntimeBroker, opts.credential));
return () => {
for (const s of stops) s();
};
}
/**
* Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). The credential names the
* seats this module claims and the verbs each promises; each verb is served on the seat's own
* subject by the module's tool of the same name. Whether this instance *holds* the seat is the bus's
* to decide: only the holder's account may subscribe the seat's subjects, so a claimant that does not
* hold it here is refused the subscription and serves nothing — never a failure of its own tools.
*/
async function serveClaimedSeats(broker: RuntimeBroker, credential?: Credential): Promise<() => void> {
const claims = credential?.claims ?? [];
if (claims.length === 0 || typeof broker.handleSubject !== "function") return () => {};
// A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
for (const { module, tools } of collectTools()) {
if (!claims.some((c) => c.seat === module)) continue;
const verbs = new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
for (const t of tools) verbs.set(t.name, (args) => t.run(args));
implementations.set(module, verbs);
}
let stops: (() => void)[] = [];
const serve = async (): Promise<void> => {
stops.forEach((s) => s());
stops = [];
const issued = typeof broker.membership === "function" ? broker.membership() : undefined;
for (const claim of claims) {
const verbs = implementations.get(claim.seat);
for (const verb of claim.serves ?? []) {
const run = verbs?.get(verb);
if (!run) {
console.log(`[mesh-tools] claims ${claim.seat} and implements no ${verb}, which that seat promises; not served`);
continue;
}
// Where the mesh issued the verb when it has; the derived shape until then.
const subject = issued?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
stops.push(await broker.handleSubject(subject, run));
console.log(`[mesh-tools] serving ${claim.seat}'s ${verb} on ${subject}, admitted where this module holds the seat`);
}
}
};
await serve();
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
return () => stops.forEach((s) => s());
}