Author SHA1 Message Date
jochen a269a76c87 The console announces five tools and reaches everything by address (hq ADR 0195)
mesh_overview, mesh_machine, mesh_search, mesh_describe and mesh_call walk the mesh's structure;
every tool has one address per layer: <seat>.<verb>, <node>/<seat>.<verb>, <node>/<module>.<tool>,
and <module>.<tool> for a module the mesh issued a plain subject. A stateful module called without
its machine, a node seat without one, a mesh seat with one, or a module on the wrong machine is
refused naming what would work. Answers come from the mesh when asked, kept five seconds, so a tool
that arrives mid-session is found. The flat catalogue stays behind MESH_CONSOLE_FLAT=1 and old
<module>.<tool> names still answer.
2026-10-03 21:54:19 +02:00
jochen e6a33dd969 node-tools is built from Go (hq ADR 0193)
The node's runtime is the Go binary: one static executable the host runs as ./node-tools from its own
bundle. Node.js stays on the machine for the TypeScript bundles the runtime launches. The TypeScript
source stays in the repository: it is the runtime image the per-module containers run until WP4c.
2026-10-03 21:26:06 +02:00
jochen e45389f75a node-tools in Go: the node's runtime, launch-only, wire-compatible with the TypeScript (hq ADR 0193)
The runtime knows no language, so nothing ties it to Node.js. This ports its serve mode — the
pinned bus connection and patient connect, following memberships, launching every served bundle
over MCP on stdio with its own environment, a child's emit published as its module, each tool,
the tools verb and seat verbs served where the mesh issued them, and the console on loopback —
to one static binary. Same subjects, request and reply bodies, event headers and MCP answers.
The TypeScript stays: it is still the runtime inside the per-module containers until WP4c.
Tests run against a real bus and share the TypeScript fixtures.
2026-10-03 21:21:52 +02:00
mesh-admin 26e7fcb4cb Merge pull request 'A launched bundle is told the module it serves (hq ADR 0193)' (#35) from feat/0193-a-launched-bundle-is-told-its-module into main 2026-10-03 19:08:06 +00:00
jochen 302645dd01 A launched bundle is told the module it serves (hq ADR 0193)
node-tools sets MESH_SERVED_MODULE for each child it launches, so the SDK's stdio loop lists that
module's tools by their names and a seat's verbs as the seat's, whichever was registered first.
2026-10-03 21:07:52 +02:00
mesh-admin 8e30ea9ff5 Merge pull request 'node-tools hands each bundle its own environment (hq ADR 0192)' (#33) from feat/0192-each-bundle-its-own-env into main 2026-10-03 13:34:05 +00:00
jochen 193ed0ac63 node-tools hands each bundle its own environment (hq ADR 0192)
The mesh composes every served module's words into MESH_TOOL_ENV; the runtime takes it at start
and removes it from its own environment, then gives each registration's contributor and each
launched child the runtime's words plus its own module's, never another's. Against an SDK
without collectToolsEach it says so and serves with the runtime's words only. The test serves
two imported bundles and one launched, each answering with its own words and none of the others'.
2026-10-03 15:28:49 +02:00
mesh-admin 4537494150 Merge pull request 'One SDK per runtime: a bundle's SDK import resolves to the runtime's copy (hq issue 209)' (#32) from fix/issue-209-one-sdk-per-runtime into main 2026-10-03 11:07:08 +00:00
jochen 7152148410 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 (module.registerHooks, in
thread; module.register is deprecated from Node 26) 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 13:00:19 +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
69 changed files with 5131 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.
+130
View File
@@ -0,0 +1,130 @@
// node-tools — the node's tool runtime, in Go (novox/hq ADR 0175, ADR 0193; to-be 38 WP4d).
//
// One process per machine: it connects to the bus on the node's credential, launches every assigned
// module's tools bundle and serves its tools and its seats' verbs; as the node-tools module — or
// wherever MESH_CONSOLE_LISTEN says — it is also the console, MCP over HTTP on loopback.
//
// MESH_BROKER_FILE the credential the mesh sealed to this machine for the runtime
// MESH_TOOL_MODULES <module>=<entrypoint>,… the bundles to launch
// MESH_TOOL_ENV {"<module>": {"<word>": "<value>"}} what each is given (ADR 0192)
// MESH_OPERATOR_ACCOUNT whose machine this is, and MESH_OPERATOR_HOME their home
// MESH_CONSOLE_LISTEN where the console listens, overriding 127.0.0.1:4270
package main
import (
"encoding/json"
"fmt"
"log"
"math/rand/v2"
"os"
"os/signal"
"syscall"
"time"
"github.com/novox/mesh-tools/node-tools/internal/bus"
"github.com/novox/mesh-tools/node-tools/internal/console"
"github.com/novox/mesh-tools/node-tools/internal/runtime"
)
// runtimeModule is the module that is the node's tool runtime; on its credential, serving is also
// the console.
const runtimeModule = "node-tools"
// consoleListen is where the console listens when nothing says otherwise.
const consoleListen = "127.0.0.1:4270"
func main() {
log.SetFlags(0)
if len(os.Args) > 1 && os.Args[1] != "serve" {
fmt.Fprintf(os.Stderr, "node-tools: %q is not a mode; this runtime serves (novox/hq ADR 0193)\n", os.Args[1])
os.Exit(2)
}
cred := credential()
conn := connectPatiently(cred)
served, err := runtime.ServedModulesFrom(os.Getenv(runtime.ToolModules), cred.Module)
if err != nil {
log.Fatalf("node-tools: %v", err)
}
envs, err := runtime.TakeToolEnvs()
if err != nil {
log.Fatalf("node-tools: %v", err)
}
stop, err := runtime.Run(conn, served, envs, log.Printf)
if err != nil {
log.Fatalf("node-tools: %v", err)
}
listen := os.Getenv("MESH_CONSOLE_LISTEN")
if listen == "" && cred.Module == runtimeModule {
listen = consoleListen
}
var up *console.Listening
if listen != "" {
node := cred.Node
if node == "" {
node = "?"
}
who := node + "." + cred.Module
up, err = console.Serve(console.NewSurface(conn, who), listen)
if err != nil {
log.Fatalf("node-tools: %v", err)
}
log.Printf("mesh console listening on http://%s/mcp as %s", up.Address, who)
}
signals := make(chan os.Signal, 1)
signal.Notify(signals, syscall.SIGTERM, syscall.SIGINT)
<-signals
stop()
if up != nil {
_ = up.Close()
}
conn.Close()
}
// credential reads the credential the mesh delivered; a missing or unreadable one is a fault of
// configuration, said and final.
func credential() bus.Credential {
file := os.Getenv("MESH_BROKER_FILE")
if file == "" {
if url := os.Getenv("MESH_BROKER_URL"); url != "" {
return bus.Credential{URL: url, Module: runtimeModule}
}
fmt.Fprintln(os.Stderr, "node-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach")
os.Exit(1)
}
raw, err := os.ReadFile(file)
if err != nil {
fmt.Fprintf(os.Stderr, "node-tools: cannot read the broker credential at %s: %v\n", file, err)
os.Exit(1)
}
var cred bus.Credential
if err := json.Unmarshal(raw, &cred); err != nil {
fmt.Fprintf(os.Stderr, "node-tools: cannot read the broker credential at %s: %v\n", file, err)
os.Exit(1)
}
if cred.URL == "" {
fmt.Fprintf(os.Stderr, "node-tools: %s carries no url — it is not a broker credential\n", file)
os.Exit(1)
}
return cred
}
// connectPatiently retries while the bus is merely not reachable yet — the normal case at startup —
// and gives up at once on what waiting cannot fix (issue 058).
func connectPatiently(cred bus.Credential) *bus.Conn {
for delay := 2 * time.Second; ; delay = min(delay*2, 30*time.Second) {
conn, err := bus.Connect(cred)
if err == nil {
return conn
}
if fatal := bus.Fatal(err); fatal != "" {
fmt.Fprintf(os.Stderr, "node-tools: %s — waiting will not fix this; giving up\n", fatal)
os.Exit(1)
}
wait := delay + time.Duration(rand.IntN(1000))*time.Millisecond
fmt.Fprintf(os.Stderr, "node-tools: the broker is not reachable yet (%v); retrying in %ds\n", err, int(wait.Round(time.Second)/time.Second))
time.Sleep(wait)
}
}
+15
View File
@@ -0,0 +1,15 @@
module github.com/novox/mesh-tools/node-tools
go 1.26.0
require (
github.com/nats-io/nats.go v1.54.0
golang.org/x/sys v0.48.0
)
require (
github.com/klauspost/compress v1.20.0 // indirect
github.com/nats-io/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // indirect
golang.org/x/crypto v0.57.0 // indirect
)
+12
View File
@@ -0,0 +1,12 @@
github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA=
github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI=
github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA=
github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k=
github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg=
github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs=
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
+568
View File
@@ -0,0 +1,568 @@
// Package bus is the node runtime's connection to the mesh bus, on NATS (novox/hq design 25, design
// 29, ADR 0160, ADR 0175). It is the Go port of node-tools' broker-nats.ts, and speaks the same
// wire: the same subjects, the same JSON request and reply bodies, the same event headers.
//
// mesh.mod.<module>.event.<type> an event a module emits
// mesh.mod.<module>.tool.<tool> a tool a module serves
// mesh.seat.<seat>.tool.<verb> a role's verb, answered by whoever holds the seat
package bus
import (
"crypto/sha256"
"crypto/tls"
"crypto/x509"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"log"
"regexp"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-tools/node-tools/internal/wire"
)
// RequestTimeout is how long a call waits for its answer — what modules already expect.
const RequestTimeout = 30 * time.Second
const assignmentsStream = "ASSIGNMENTS"
// Credential is a broker credential as the mesh delivers it (novox/hq ADR 0120).
type Credential struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
Node string `json:"node,omitempty"`
Module string `json:"module,omitempty"`
User string `json:"user,omitempty"`
Password string `json:"password,omitempty"`
Claims []Claim `json:"claims,omitempty"`
}
// Claim is a seat a module claims, with the verbs it promises (novox/hq ADR 0159).
type Claim struct {
Seat string `json:"seat"`
Scope string `json:"scope,omitempty"`
Serves []string `json:"serves,omitempty"`
}
// Membership is what the mesh issued one assignment (novox/hq ADR 0160).
type Membership struct {
Node string `json:"node"`
Module string `json:"module"`
Serves []Served `json:"serves"`
Seats []SeatVerb `json:"seats,omitempty"`
Emits string `json:"emits"`
Reaches map[string][]string `json:"reaches,omitempty"`
Tools string `json:"tools"`
}
// Served is an address a tool is answered on; `{tool}` stands for the tool's name.
type Served struct {
Subject string `json:"subject"`
Queue string `json:"queue,omitempty"`
}
// SeatVerb is one verb of a seat a module holds, where it is answered.
type SeatVerb struct {
Seat string `json:"seat"`
Verb string `json:"verb"`
Subject string `json:"subject"`
}
// Envelope is an event as a module emits it: the body is the payload, the metadata rides as headers
// (novox/hq ADR 0042).
type Envelope struct {
Key string `json:"key"`
Node string `json:"node,omitempty"`
Body json.RawMessage `json:"body"`
Headers map[string]string `json:"headers,omitempty"`
}
// Answered is a call's result and the machine that gave it (novox/hq ADR 0159).
type Answered struct {
Result json.RawMessage
Node string
}
// Handler answers one request body.
type Handler func(body json.RawMessage) (any, error)
// ErrPin is a bus whose certificate is not the one the mesh pinned: final, never retried.
var ErrPin = errors.New("the bus's certificate does not match the pin")
// Fatal is why a connection failure is final rather than "not yet", or "" when waiting may fix it —
// the same classification the TypeScript runtime makes.
func Fatal(err error) string {
if err == nil {
return ""
}
msg := err.Error()
if errors.Is(err, ErrPin) || strings.Contains(msg, ErrPin.Error()) {
return "the bus's certificate does not match the pin"
}
if regexp.MustCompile(`(?i)invalid url|no servers available for connection: .*url`).MatchString(msg) ||
strings.Contains(msg, "nats: invalid url") {
return "the bus address is not a usable URL"
}
if regexp.MustCompile(`(?i)authorization violation|user authentication expired|permissions violation`).MatchString(msg) {
return "the bus refused this account"
}
return ""
}
func normalizeFingerprint(f string) string {
f = strings.TrimSpace(f)
if len(f) > 7 && strings.EqualFold(f[:7], "sha256:") {
f = f[7:]
}
return strings.ToLower(strings.ReplaceAll(f, ":", ""))
}
// pinned accepts exactly the certificate with this SHA-256 and no other. The pin is the only check:
// the bus's certificate names the seat, not the address a machine dials it by.
func pinned(want string) *tls.Config {
want = normalizeFingerprint(want)
return &tls.Config{
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin, which is stricter
MinVersion: tls.VersionTLS12,
VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error {
if len(raw) == 0 {
return fmt.Errorf("%w: it presented none", ErrPin)
}
sum := sha256.Sum256(raw[0])
if got := hex.EncodeToString(sum[:]); got != want {
return fmt.Errorf("%w: it presented %s, not the pinned %s", ErrPin, got, want)
}
return nil
},
}
}
// Conn is the runtime's connection: the module it is, the memberships it follows, the subjects it
// answers.
type Conn struct {
nc *nats.Conn
js nats.JetStreamContext
self string
node string
cred Credential
mu sync.Mutex
issued map[string]*Membership // module → membership; present with nil = followed, none issued
onNew []func(Membership)
subs []*nats.Subscription
Logf func(format string, args ...any)
}
// Connect dials the bus as the credential's module. A module's subjects come from its credential,
// never from its calls (ADR 0074).
func Connect(cred Credential) (*Conn, error) {
if cred.Module == "" {
return nil, errors.New("a broker credential with no module: the runtime derives its subjects " +
"from the account the mesh issued, and cannot guess which module it is")
}
node := cred.Node
if node == "" {
node = "?"
}
opts := []nats.Option{
nats.Name(node + "." + cred.Module),
// Reconnect forever: the bus restarting is an upgrade, not a reason to exit.
nats.MaxReconnects(-1),
}
if cred.User != "" {
opts = append(opts, nats.UserInfo(cred.User, cred.Password),
// Its own inbox: every user's inbox is private to it (design 25 §4).
nats.CustomInboxPrefix("_INBOX."+cred.User))
}
if strings.TrimSpace(cred.Fingerprint) != "" {
opts = append(opts, nats.Secure(pinned(cred.Fingerprint)))
}
nc, err := nats.Connect(cred.URL, opts...)
if err != nil {
return nil, err
}
js, err := nc.JetStream()
if err != nil {
nc.Close()
return nil, err
}
c := &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred,
issued: map[string]*Membership{}, Logf: log.Printf}
c.Follow(cred.Module)
return c, nil
}
// Module is what this connection is.
func (c *Conn) Module() string { return c.self }
// Node is the machine this connection's account is scoped to.
func (c *Conn) Node() string { return c.node }
// Credential is what this connection was opened with.
func (c *Conn) Credential() Credential { return c.cred }
// MembershipSubject is the one address a runtime derives for an assignment (ADR 0160).
func MembershipSubject(node, module string) string {
return "mesh.assignment." + node + "." + module
}
// Follow reads a module's membership on this machine once and follows it live, so its tools are
// served where the mesh issued them (ADR 0175).
func (c *Conn) Follow(module string) {
c.mu.Lock()
if c.node == "" {
c.mu.Unlock()
return
}
if _, has := c.issued[module]; has {
c.mu.Unlock()
return
}
c.issued[module] = nil
c.mu.Unlock()
subject := MembershipSubject(c.node, module)
// The subject-addressed direct get: the one address the mesh grants this account on the
// stream's API.
if got, err := c.nc.Request("$JS.API.DIRECT.GET."+assignmentsStream+"."+subject, nil, 5*time.Second); err == nil {
if got.Header.Get("Status") == "" && len(got.Data) > 0 {
var m Membership
if json.Unmarshal(got.Data, &m) == nil {
c.mu.Lock()
c.issued[module] = &m
c.mu.Unlock()
}
}
}
if c.Membership(module) == nil {
c.Logf("[mesh-tools] no membership issued for %s on %s yet; serving the derived shape until one arrives", module, c.node)
}
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
var m Membership
if err := json.Unmarshal(msg.Data, &m); err != nil {
c.Logf("[mesh-tools] a membership arrived that is not one: %v", err)
return
}
c.mu.Lock()
c.issued[module] = &m
handlers := append([]func(Membership){}, c.onNew...)
c.mu.Unlock()
c.Logf("[mesh-tools] %s on %s was issued a new membership; re-serving on it", module, c.node)
for _, h := range handlers {
h(m)
}
_ = c.nc.Flush()
})
if err == nil {
c.track(sub)
}
}
// Membership is what the mesh issued a module here, or nil when nothing has been issued.
func (c *Conn) Membership(module string) *Membership {
c.mu.Lock()
defer c.mu.Unlock()
return c.issued[module]
}
// Following says whether this connection follows a module's membership.
func (c *Conn) Following(module string) bool {
c.mu.Lock()
defer c.mu.Unlock()
_, has := c.issued[module]
return has
}
// Serving is every module whose membership this connection follows.
func (c *Conn) Serving() []string {
c.mu.Lock()
defer c.mu.Unlock()
out := make([]string, 0, len(c.issued))
for m := range c.issued {
out = append(out, m)
}
return out
}
// OnMembership is called with every new membership any followed module is issued.
func (c *Conn) OnMembership(h func(Membership)) {
c.mu.Lock()
c.onNew = append(c.onNew, h)
c.mu.Unlock()
}
func (c *Conn) track(s *nats.Subscription) {
c.mu.Lock()
c.subs = append(c.subs, s)
c.mu.Unlock()
}
// servedOn is where a served module's tool is answered: its membership's subjects when issued, the
// derived shape otherwise (the shape the mesh issues on day one).
func (c *Conn) servedOn(module, tool string) []Served {
if m := c.Membership(module); m != nil {
out := make([]Served, 0, len(m.Serves)+1)
for _, s := range m.Serves {
out = append(out, Served{Subject: strings.ReplaceAll(s.Subject, "{tool}", tool), Queue: s.Queue})
}
if tool == "tools" && m.Tools != "" {
found := false
for _, s := range out {
found = found || s.Subject == m.Tools
}
if !found {
out = append([]Served{{Subject: m.Tools, Queue: "serve." + module}}, out...)
}
}
return out
}
base := "mesh.mod." + module + ".tool." + tool
out := []Served{{Subject: base, Queue: "serve." + module}}
if c.node != "" {
out = append(out, Served{Subject: base + "." + c.node})
}
return out
}
type reply struct {
Result any `json:"result,omitempty"`
Error string `json:"error,omitempty"`
Node string `json:"node,omitempty"`
}
// answerOn answers one subject with one handler, and says which machine answered (ADR 0159).
func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) {
cb := func(msg *nats.Msg) {
go func() {
var r reply
result, err := h(json.RawMessage(msg.Data))
if err != nil {
r.Error = err.Error()
} else {
r.Result = nullable(result)
}
r.Node = c.node
body, _ := wire.Marshal(r)
_ = msg.Respond(body)
}()
}
var sub *nats.Subscription
var err error
if queue != "" {
sub, err = c.nc.QueueSubscribe(subject, queue, cb)
} else {
sub, err = c.nc.Subscribe(subject, cb)
}
if err != nil {
return func() {}, err
}
c.track(sub)
return func() { _ = sub.Unsubscribe() }, nil
}
// nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always
// carries `result` when the handler did not throw.
func nullable(v any) any {
if v == nil {
return json.RawMessage("null")
}
return v
}
// HandleSubject answers a subject outright: a seat's verb where the mesh issued it.
func (c *Conn) HandleSubject(subject string, h Handler) (func(), error) {
return c.answerOn(subject, "", h)
}
// Handle serves `<module>.<tool>` where the mesh issued that module, and follows its membership:
// when a new one arrives, it serves where it now says and stops where it no longer does.
func (c *Conn) Handle(key string, h Handler) (func(), error) {
if strings.HasPrefix(key, "seat:") {
subject, err := ToolSubject(key, c.self)
if err != nil {
return func() {}, err
}
return c.answerOn(subject, "", h)
}
module, tool := c.self, key
if dot := strings.Index(key, "."); dot >= 0 {
module, tool = key[:dot], key[dot+1:]
}
if module != c.self && !c.Following(module) {
return func() {}, fmt.Errorf("%s cannot serve %s: a module serves its own tools, and a runtime "+
"those of the modules it follows", c.self, key)
}
var mu sync.Mutex
var stops []func()
serve := func() {
mu.Lock()
defer mu.Unlock()
for _, s := range stops {
s()
}
stops = nil
for _, s := range c.servedOn(module, tool) {
if stop, err := c.answerOn(s.Subject, s.Queue, h); err == nil {
stops = append(stops, stop)
} else {
c.Logf("[mesh-tools] cannot serve %s on %s: %v", key, s.Subject, err)
}
}
}
serve()
c.OnMembership(func(m Membership) {
if m.Module == module {
serve()
}
})
return func() {
mu.Lock()
defer mu.Unlock()
for _, s := range stops {
s()
}
stops = nil
}, nil
}
// reachedAt is where a call by key goes: a subject this connection's own membership says it
// reaches — the machine's when named — else the derived shape.
func (c *Conn) reachedAt(key string) (string, error) {
name, wanted, _ := strings.Cut(key, "@")
if m := c.Membership(c.self); m != nil {
if reach := m.Reaches[name]; len(reach) > 0 {
if wanted == "" {
return reach[0], nil
}
for _, s := range reach {
if strings.HasSuffix(s, "."+wanted) {
return s, nil
}
}
}
}
return ToolSubject(key, c.self)
}
// Ask calls a tool by key — or on a subject the mesh listed for it — and learns which machine
// answered. Core request/reply: a tool call is never persisted (design 25 §3).
func (c *Conn) Ask(key string, body any, on string) (Answered, error) {
subject := on
if subject == "" {
s, err := c.reachedAt(key)
if err != nil {
return Answered{}, err
}
subject = s
}
data, err := wire.Marshal(body)
if err != nil {
return Answered{}, err
}
msg, err := c.nc.Request(subject, data, RequestTimeout)
if err != nil {
if errors.Is(err, nats.ErrNoResponders) {
return Answered{}, errors.New("503 no responders")
}
if errors.Is(err, nats.ErrTimeout) {
return Answered{}, errors.New("timeout")
}
return Answered{}, err
}
var r struct {
Result json.RawMessage `json:"result"`
Error string `json:"error"`
Node string `json:"node"`
}
if err := json.Unmarshal(msg.Data, &r); err != nil {
return Answered{}, err
}
if r.Error != "" {
return Answered{}, errors.New(r.Error)
}
return Answered{Result: r.Result, Node: r.Node}, nil
}
// PublishAs emits an event as a module: published into JetStream and awaited, de-duplicated by its
// own id (ADR 0042). The body is the payload; the metadata rides as headers.
func (c *Conn) PublishAs(module string, env Envelope) error {
msg := nats.NewMsg("mesh.mod." + module + ".event." + env.Key)
for k, v := range env.Headers {
msg.Header.Set(k, v)
}
if env.Headers["content-type"] == "" {
msg.Header.Set("content-type", "application/json")
}
if env.Node != "" {
msg.Header.Set("x-node", env.Node)
}
body := env.Body
if len(body) == 0 {
body = json.RawMessage("null")
}
msg.Data = body
var opts []nats.PubOpt
if id := env.Headers["x-event-id"]; id != "" {
opts = append(opts, nats.MsgId(id))
}
_, err := c.js.PublishMsg(msg, opts...)
return err
}
// Flush waits until the bus has every subscription made so far, so what is served is answerable
// when this returns.
func (c *Conn) Flush() { _ = c.nc.Flush() }
// Close unsubscribes everything and drains, so an in-flight reply is finished rather than dropped.
func (c *Conn) Close() {
c.mu.Lock()
subs := c.subs
c.subs = nil
c.mu.Unlock()
for _, s := range subs {
_ = s.Unsubscribe()
}
_ = c.nc.Drain()
}
// ToolSubject is a tool's subject. A bare name is this module's own; `<module>.<tool>` another's;
// `seat:<seat>.<verb>` a role's, with `@<node>` for a node-scoped seat (design 33 §4).
func ToolSubject(key, self string) (string, error) {
if strings.HasPrefix(key, "seat:") {
rest := strings.TrimPrefix(key, "seat:")
dot := strings.Index(rest, ".")
if dot < 0 {
return "", fmt.Errorf("%q names a seat and no verb: seat:<seat>.<verb>", key)
}
seat := rest[:dot]
verb, node, _ := strings.Cut(rest[dot+1:], "@")
if node != "" {
return "mesh.seat." + seat + ".tool." + verb + "." + node, nil
}
return "mesh.seat." + seat + ".tool." + verb, nil
}
name, node, _ := strings.Cut(key, "@")
var base string
if dot := strings.Index(name, "."); dot < 0 {
base = "mesh.mod." + self + ".tool." + name
} else {
base = "mesh.mod." + name[:dot] + ".tool." + name[dot+1:]
}
if node != "" {
return base + "." + node, nil
}
return base, nil
}
// SeatToolSubject is a seat's verb as its holder serves it: flat for a mesh seat, carrying the
// machine for a node-scoped one.
func SeatToolSubject(seat, verb, scope, node string) string {
base := "mesh.seat." + seat + ".tool." + verb
if scope == "node" && node != "" {
return base + "." + node
}
return base
}
+48
View File
@@ -0,0 +1,48 @@
package bus
import (
"errors"
"fmt"
"testing"
)
func TestSubjectsAreTheOnesTheTypeScriptRuntimeUses(t *testing.T) {
for key, want := range map[string]string{
"status": "mesh.mod.self.tool.status",
"alpha.one": "mesh.mod.alpha.tool.one",
"alpha.one@anchor": "mesh.mod.alpha.tool.one.anchor",
"seat:node-shelf.list": "mesh.seat.node-shelf.tool.list",
"seat:node-shelf.list@anchor": "mesh.seat.node-shelf.tool.list.anchor",
} {
if got, err := ToolSubject(key, "self"); err != nil || got != want {
t.Errorf("%s: %s %v, want %s", key, got, err, want)
}
}
if _, err := ToolSubject("seat:nothing", "self"); err == nil {
t.Error("a seat with no verb was accepted")
}
if SeatToolSubject("s", "v", "node", "n") != "mesh.seat.s.tool.v.n" || SeatToolSubject("s", "v", "mesh", "n") != "mesh.seat.s.tool.v" {
t.Error("seat subjects")
}
}
func TestAFingerprintIsReadHoweverItIsWritten(t *testing.T) {
for _, f := range []string{"sha256:AB:CD:ef", "abcdef", "ABCDEF", "SHA256:abcdef"} {
if got := normalizeFingerprint(f); got != "abcdef" {
t.Errorf("%s → %s", f, got)
}
}
}
func TestWhatWaitingCannotFixIsFinal(t *testing.T) {
for err, want := range map[error]string{
fmt.Errorf("x509: %w: it presented aa", ErrPin): "the bus's certificate does not match the pin",
errors.New("nats: Authorization Violation"): "the bus refused this account",
errors.New("nats: invalid url"): "the bus address is not a usable URL",
errors.New("dial tcp 10.0.0.1:4222: connect: connection refused"): "",
} {
if got := Fatal(err); got != want {
t.Errorf("%v → %q, want %q", err, got, want)
}
}
}
+735
View File
@@ -0,0 +1,735 @@
package console
// The mesh's tools found by address, not announced whole (novox/hq ADR 0195, to-be 34 §3a).
//
// The console announces five tools. Everything the mesh answers is reached through them by one
// address per layer:
//
// <seat>.<verb> a seat held once for the mesh — its holder answers
// <node>/<seat>.<verb> a seat held once per machine — that machine's holder answers
// <node>/<module>.<tool> a module assigned to a machine — that assignment answers
// <module>.<tool> also, for a module whose instances are interchangeable (ADR 0160)
//
// A module that is not interchangeable is called with its machine or refused, naming the machines
// it runs on: "whichever answers" is no answer for state a machine holds.
//
// Every discovery verb asks the mesh when it is called — kept a few seconds at most, never for a
// session — so a tool that arrived a minute ago is found without the client reconnecting.
import (
"encoding/json"
"fmt"
"regexp"
"sort"
"strings"
"sync"
"time"
"github.com/novox/mesh-tools/node-tools/internal/bus"
)
// IndexKept is how long what the mesh answered is kept before it is asked again: long enough that
// one agent turn's search, describe and call ask once, short enough that nothing goes stale.
var IndexKept = 5 * time.Second
// The five tools the console announces. Names of the API's kind — letters, digits, `_`, `-`.
const (
verbOverview = "mesh_overview"
verbMachine = "mesh_machine"
verbSearch = "mesh_search"
verbDescribe = "mesh_describe"
verbCall = "mesh_call"
)
// searchCap is how many matches a search answers before it says how many more there were.
const searchCap = 25
const grammar = "Addresses: `<seat>.<verb>` for a seat held once for the mesh (e.g. `mesh-controller.nodes`); " +
"`<node>/<seat>.<verb>` for a seat every machine holds (e.g. `ace/node-packet-filter.rules`); " +
"`<node>/<module>.<tool>` for a module on one machine (e.g. `novox/postgres.postgres_list_databases`); " +
"and `<module>.<tool>` also for a module whose instances are interchangeable."
// discovery is the five tools as tools/list announces them.
func discovery() []map[string]any {
str := func(desc string) map[string]any { return map[string]any{"type": "string", "description": desc} }
obj := func(props map[string]any, required ...string) map[string]any {
s := map[string]any{"type": "object", "properties": props}
if len(required) > 0 {
s["required"] = required
}
return s
}
return []map[string]any{
{"name": verbOverview, "inputSchema": obj(map[string]any{}),
"description": "The mesh at a glance: the seats it holds once for the whole mesh with their verbs, the seats " +
"every machine holds, and its machines. Start here, then `mesh_machine` for one machine. " + grammar},
{"name": verbMachine, "inputSchema": obj(map[string]any{"node": str("the machine, as mesh_overview names it")}, "node"),
"description": "One machine: the seats it holds with their verbs, and the modules assigned to it with their tools — " +
"each with the address to describe or call it by. " + grammar},
{"name": verbSearch, "inputSchema": obj(map[string]any{"query": str("words to find in tool names and descriptions, e.g. `postgres databases`")}, "query"),
"description": "Find tools anywhere in the mesh by words: every match's address and a line of what it does, across the " +
"mesh's seats, the machines' seats and every module on every machine. " + grammar},
{"name": verbDescribe, "inputSchema": obj(map[string]any{"address": str("the tool's address")}, "address"),
"description": "What one tool does and the arguments it takes, as a JSON schema. The machine is in the address, " +
"never an argument. " + grammar},
{"name": verbCall, "inputSchema": obj(map[string]any{
"address": str("the tool's address"),
"arguments": map[string]any{"type": "object", "description": "the tool's arguments, as mesh_describe gives its schema"},
}, "address"),
"description": "Call one tool by its address with its arguments; the answer says which machine gave it. " + grammar},
}
}
func isDiscovery(name string) bool {
switch name {
case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall:
return true
}
return false
}
// seatInfo is a seat as the mesh's records define it, and who holds it where.
type seatInfo struct {
Seat string
Scope string // "mesh" or "node"
Verbs []Tool
Holders []holder
}
type holder struct {
Module string `json:"module"`
Node string `json:"node"`
}
// moduleInfo is a module that answers tools: where it runs, whether any instance will do, and its
// tools as one of its instances described them.
type moduleInfo struct {
Module string
On []string
Interchangeable bool
Tools []Tool
}
// index is what the mesh answered about itself, at one moment.
type index struct {
Seats []seatInfo
Machines []string
Modules map[string]*moduleInfo
NotAnswering []string
Listing *Listing // the flat catalogue the call path resolves subjects with
}
func (x *index) seat(name string) *seatInfo {
for i := range x.Seats {
if x.Seats[i].Seat == name {
return &x.Seats[i]
}
}
return nil
}
func findTool(tools []Tool, name string) *Tool {
for i := range tools {
if tools[i].Name == name {
return &tools[i]
}
}
return nil
}
// controllerOutput is a controller seat verb's answer: the command's printed output.
func controllerOutput(conn *bus.Conn, verb string) (string, error) {
got, err := conn.Ask("seat:mesh-controller."+verb, map[string]any{}, "")
if err != nil {
return "", err
}
var r struct {
Output string `json:"output"`
OK *bool `json:"ok"`
}
if json.Unmarshal(got.Result, &r) != nil {
return "", fmt.Errorf("mesh-controller.%s answered something that is not its output", verb)
}
if r.OK != nil && !*r.OK {
return "", fmt.Errorf("mesh-controller.%s: %s", verb, strings.TrimSpace(r.Output))
}
return r.Output, nil
}
// jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs
// the controller's command, and a command may warn before it answers.
func jsonIn(output string) string {
if strings.HasPrefix(strings.TrimSpace(output), "{") {
return strings.TrimSpace(output)
}
if i := strings.Index(output, "\n{"); i >= 0 {
return strings.TrimSpace(output[i+1:])
}
return ""
}
// A machine as `node list` prints it: its name, when it was last heard from, its mode — converged
// or adopted, the only two the command prints — and its id. Lines the command says around them
// (a warning about the bus's users, "no node records yet") are not machines and are skipped.
var nodeLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+.*\s(converged|adopted)\s+\S+\s*$`)
// machinesIn reads the machines from `node list`'s output.
func machinesIn(output string) []string {
var out []string
for _, line := range strings.Split(output, "\n") {
if m := nodeLine.FindStringSubmatch(line); m != nil {
out = append(out, m[1])
}
}
sort.Strings(out)
return out
}
// A module as `module list` prints it: name, version, how it was built, and where it runs.
var moduleLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+\S+\s+.*?\s+on (.+)$`)
// assignmentsIn reads, from `module list`'s output, the machines each module runs on.
func assignmentsIn(output string) map[string][]string {
out := map[string][]string{}
for _, line := range strings.Split(output, "\n") {
if line == "" || line[0] == ' ' || line[0] == '\t' {
continue
}
m := moduleLine.FindStringSubmatch(strings.TrimRight(line, " "))
if m == nil {
continue
}
on := strings.TrimSpace(m[2])
if on == "nothing" {
out[m[1]] = []string{}
continue
}
var nodes []string
for _, n := range strings.Split(on, ",") {
if n = strings.TrimSpace(n); n != "" {
nodes = append(nodes, n)
}
}
sort.Strings(nodes)
out[m[1]] = nodes
}
return out
}
// interchangeable is whether the mesh issued the module a plain subject for this tool — one any of
// its instances answers (ADR 0160): the module's own subject with no machine after it.
func interchangeable(module string, t Tool) bool {
plain := "mesh.mod." + module + ".tool." + t.Name
for _, s := range t.Subjects {
if s == plain {
return true
}
}
return false
}
// indexOn asks the mesh what it holds: the flat listing (catalogue, modules, the seats' tools),
// and from the controller the seats' holders, the machines and the assignments — at once.
func indexOn(conn *bus.Conn) (*index, error) {
var wg sync.WaitGroup
var seatsOut, nodesOut, modulesOut string
wg.Add(3)
go func() { defer wg.Done(); seatsOut, _ = controllerOutput(conn, "seats") }()
go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }()
go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }()
l, err := toolsOn(conn)
wg.Wait()
if err != nil {
return nil, err
}
x := &index{Modules: map[string]*moduleInfo{}, NotAnswering: l.NotAnswering, Listing: l}
// The seats: their verbs from the mesh's records, their holders from the controller.
var held struct {
Seats []struct {
Seat string `json:"seat"`
Scope string `json:"scope"`
Holders []holder `json:"holders"`
} `json:"seats"`
}
_ = json.Unmarshal([]byte(jsonIn(seatsOut)), &held)
holders := map[string][]holder{}
scopes := map[string]string{}
for _, s := range held.Seats {
holders[s.Seat] = s.Holders
scopes[s.Seat] = s.Scope
}
bySeat := map[string]*seatInfo{}
for _, t := range l.Tools {
if !t.Seat {
continue
}
s := bySeat[t.Module]
if s == nil {
scope := t.Scope
if scope == "" {
scope = scopes[t.Module]
}
if scope != "node" {
scope = "mesh"
}
s = &seatInfo{Seat: t.Module, Scope: scope, Holders: holders[t.Module]}
bySeat[t.Module] = s
}
s.Verbs = append(s.Verbs, t)
}
for _, s := range bySeat {
x.Seats = append(x.Seats, *s)
}
sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat })
// The machines: what the controller knows, and any that hold a seat.
seen := map[string]bool{}
for _, n := range machinesIn(nodesOut) {
seen[n] = true
}
for _, s := range x.Seats {
for _, h := range s.Holders {
if h.Node != "" {
seen[h.Node] = true
}
}
}
// The modules that answer tools, where they run, and whether any instance will do.
on := assignmentsIn(modulesOut)
for _, t := range l.Tools {
if t.Seat {
continue
}
m := x.Modules[t.Module]
if m == nil {
m = &moduleInfo{Module: t.Module, On: on[t.Module]}
x.Modules[t.Module] = m
}
m.Tools = append(m.Tools, t)
if interchangeable(t.Module, t) {
m.Interchangeable = true
}
}
for _, m := range x.Modules {
// A module the controller could not place is placed where its own answer says it runs.
if len(m.On) == 0 {
nodes := map[string]bool{}
for _, t := range m.Tools {
for _, s := range t.Subjects {
base := "mesh.mod." + m.Module + ".tool." + t.Name + "."
if strings.HasPrefix(s, base) {
nodes[strings.TrimPrefix(s, base)] = true
}
}
}
for n := range nodes {
m.On = append(m.On, n)
}
sort.Strings(m.On)
}
for _, n := range m.On {
seen[n] = true
}
}
for n := range seen {
x.Machines = append(x.Machines, n)
}
sort.Strings(x.Machines)
return x, nil
}
func (s *Surface) index() (*index, error) {
s.mu.Lock()
if s.idx != nil && time.Since(s.idxAt) <= IndexKept {
x := s.idx
s.mu.Unlock()
return x, nil
}
s.mu.Unlock()
x, err := indexOn(s.conn)
if err != nil {
return nil, err
}
s.mu.Lock()
s.idx, s.idxAt = x, time.Now()
s.mu.Unlock()
return x, nil
}
// target is what an address resolves to.
type target struct {
Address string
Key string // the call key the existing path takes: seat:<s>.<v>[@node] or <m>.<t>[@node]
Name string // <seat>.<verb> or <module>.<tool>, for the listing's subject lookup
Node string
Tool Tool
Seat bool
}
// resolve turns an address into exactly one target, or says why it cannot.
func resolve(x *index, address string) (target, error) {
address = strings.TrimSpace(address)
node, rest, hasNode := strings.Cut(address, "/")
if !hasNode {
rest, node = address, ""
}
dot := strings.Index(rest, ".")
if dot <= 0 || dot == len(rest)-1 || strings.Contains(rest, "/") {
return target{}, fmt.Errorf("%q is not an address. %s", address, grammar)
}
prefix, name := rest[:dot], rest[dot+1:]
if s := x.seat(prefix); s != nil {
verb := findTool(s.Verbs, name)
if verb == nil {
return target{}, fmt.Errorf("the seat %s has no verb %s; it has %s", prefix, name, toolNames(s.Verbs))
}
if s.Scope == "node" {
if node == "" {
return target{}, fmt.Errorf("%s is held once per machine: write <node>/%s — it is held on %s",
prefix, rest, orNobody(nodesOf(s.Holders)))
}
return target{Address: node + "/" + rest, Key: "seat:" + rest + "@" + node, Name: rest, Node: node, Tool: *verb, Seat: true}, nil
}
if node != "" {
return target{}, fmt.Errorf("%s is held once for the whole mesh: write %s, without a machine", prefix, rest)
}
return target{Address: rest, Key: "seat:" + rest, Name: rest, Tool: *verb, Seat: true}, nil
}
m := x.Modules[prefix]
if m == nil {
return target{}, fmt.Errorf("nothing in the mesh is called %s: no seat, and no module that answers tools. "+
"mesh_search finds a tool by words", prefix)
}
tool := findTool(m.Tools, name)
if tool == nil {
return target{}, fmt.Errorf("%s has no tool %s; it has %s", prefix, name, toolNames(m.Tools))
}
if node == "" {
if !m.Interchangeable {
return target{}, fmt.Errorf("%s keeps state on each machine it runs on, so a call names the machine: "+
"write <node>/%s — it runs on %s", prefix, rest, orNobody(m.On))
}
return target{Address: rest, Key: rest, Name: rest, Tool: *tool}, nil
}
if len(m.On) > 0 && !contains(m.On, node) {
return target{}, fmt.Errorf("%s does not run on %s; it runs on %s", prefix, node, orNobody(m.On))
}
return target{Address: node + "/" + rest, Key: rest + "@" + node, Name: rest, Node: node, Tool: *tool}, nil
}
func contains(xs []string, s string) bool {
for _, x := range xs {
if x == s {
return true
}
}
return false
}
func toolNames(ts []Tool) string {
names := make([]string, 0, len(ts))
for _, t := range ts {
names = append(names, t.Name)
}
sort.Strings(names)
return strings.Join(names, ", ")
}
func nodesOf(hs []holder) []string {
var out []string
for _, h := range hs {
if h.Node != "" && !contains(out, h.Node) {
out = append(out, h.Node)
}
}
sort.Strings(out)
return out
}
func orNobody(nodes []string) string {
if len(nodes) == 0 {
return "no machine the mesh knows of"
}
return strings.Join(nodes, ", ")
}
// firstLine is a description's first line, for a list.
func firstLine(s string) string {
s = strings.TrimSpace(s)
if i := strings.IndexAny(s, "\n"); i >= 0 {
s = s[:i]
}
if len(s) > 160 {
s = s[:157] + "…"
}
return s
}
// schemaWithoutNode is a tool's schema as an agent passes it: the machine is in the address.
func schemaWithoutNode(raw json.RawMessage) map[string]any {
schema := asSchema(raw)
out := map[string]any{}
for k, v := range schema {
out[k] = v
}
if p, ok := schema["properties"].(map[string]any); ok {
props := map[string]any{}
for k, v := range p {
if k != "node" {
props[k] = v
}
}
out["properties"] = props
}
if r, ok := schema["required"].([]any); ok {
var keep []any
for _, k := range r {
if k != "node" {
keep = append(keep, k)
}
}
if len(keep) == 0 {
delete(out, "required")
} else {
out["required"] = keep
}
}
return out
}
// answerText is a discovery verb's answer as MCP content: JSON, indented.
func answerText(v any) map[string]any {
var b strings.Builder
enc := json.NewEncoder(&b)
enc.SetEscapeHTML(false) // `<node>/…` is read by an agent, not embedded in a page
enc.SetIndent("", " ")
_ = enc.Encode(v)
return map[string]any{"content": []map[string]any{{"type": "text", "text": strings.TrimRight(b.String(), "\n")}}}
}
func failure(text string) map[string]any {
return map[string]any{"content": []map[string]any{{"type": "text", "text": text}}, "isError": true}
}
// discover answers one of the five.
func (s *Surface) discover(name string, args map[string]any) map[string]any {
x, err := s.index()
if err != nil {
return failure(whyItFailed(catalogueModules, err))
}
str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) }
switch name {
case verbOverview:
type verbLine struct {
Address string `json:"address"`
Description string `json:"description"`
}
var mesh, node []map[string]any
for _, st := range x.Seats {
var verbs []verbLine
for _, v := range st.Verbs {
addr := st.Seat + "." + v.Name
if st.Scope == "node" {
addr = "<node>/" + addr
}
verbs = append(verbs, verbLine{addr, firstLine(v.Description)})
}
entry := map[string]any{"seat": st.Seat, "verbs": verbs}
if st.Scope == "node" {
entry["held on"] = nodesOf(st.Holders)
node = append(node, entry)
} else {
if h := nodesOf(st.Holders); len(h) > 0 {
entry["held on"] = h
}
mesh = append(mesh, entry)
}
}
modules := 0
tools := 0
for _, m := range x.Modules {
modules++
tools += len(m.Tools)
}
return answerText(map[string]any{
"seats of the mesh": mesh,
"seats every machine": node,
"machines": x.Machines,
"modules with tools": fmt.Sprintf("%d modules, %d tools — mesh_machine lists a machine's, mesh_search finds one", modules, tools),
"not answering": x.NotAnswering,
})
case verbMachine:
node := str("node")
if node == "" {
return failure("mesh_machine needs `node`: one of " + orNobody(x.Machines))
}
if !contains(x.Machines, node) {
return failure(fmt.Sprintf("the mesh knows no machine %q; it has %s", node, orNobody(x.Machines)))
}
var seats []map[string]any
for _, st := range x.Seats {
if st.Scope != "node" || !contains(nodesOf(st.Holders), node) {
continue
}
var verbs []string
for _, v := range st.Verbs {
verbs = append(verbs, node+"/"+st.Seat+"."+v.Name)
}
var by string
for _, h := range st.Holders {
if h.Node == node {
by = h.Module
}
}
seats = append(seats, map[string]any{"seat": st.Seat, "held by": by, "verbs": verbs})
}
var modules []map[string]any
names := make([]string, 0, len(x.Modules))
for n := range x.Modules {
names = append(names, n)
}
sort.Strings(names)
for _, n := range names {
m := x.Modules[n]
if !contains(m.On, node) {
continue
}
var tools []map[string]string
for _, t := range m.Tools {
tools = append(tools, map[string]string{"address": node + "/" + m.Module + "." + t.Name, "does": firstLine(t.Description)})
}
entry := map[string]any{"module": m.Module, "tools": tools}
if m.Interchangeable {
entry["interchangeable"] = "any instance answers <module>.<tool> as well"
}
modules = append(modules, entry)
}
return answerText(map[string]any{"machine": node, "seats": seats, "modules": modules})
case verbSearch:
query := strings.ToLower(str("query"))
if query == "" {
return failure("mesh_search needs `query`: words to find, e.g. `postgres databases`")
}
words := strings.Fields(query)
type hit struct {
Address string `json:"address"`
Does string `json:"does"`
Also string `json:"also,omitempty"`
}
var hits []hit
matches := func(parts ...string) bool {
hay := strings.ToLower(strings.Join(parts, " "))
for _, w := range words {
if !strings.Contains(hay, w) {
return false
}
}
return true
}
for _, st := range x.Seats {
for _, v := range st.Verbs {
if !matches(st.Seat, v.Name, v.Description) {
continue
}
if st.Scope == "node" {
on := nodesOf(st.Holders)
first := "<node>"
also := ""
if len(on) > 0 {
first = on[0]
if len(on) > 1 {
also = "also on " + strings.Join(on[1:], ", ")
}
}
hits = append(hits, hit{first + "/" + st.Seat + "." + v.Name, firstLine(v.Description), also})
} else {
hits = append(hits, hit{st.Seat + "." + v.Name, firstLine(v.Description), ""})
}
}
}
names := make([]string, 0, len(x.Modules))
for n := range x.Modules {
names = append(names, n)
}
sort.Strings(names)
for _, n := range names {
m := x.Modules[n]
for _, t := range m.Tools {
if !matches(m.Module, t.Name, t.Description) {
continue
}
switch {
case m.Interchangeable:
hits = append(hits, hit{m.Module + "." + t.Name, firstLine(t.Description), "any instance; on " + orNobody(m.On)})
case len(m.On) == 0:
hits = append(hits, hit{"<node>/" + m.Module + "." + t.Name, firstLine(t.Description), "the mesh places it on no machine"})
default:
also := ""
if len(m.On) > 1 {
also = "also on " + strings.Join(m.On[1:], ", ")
}
hits = append(hits, hit{m.On[0] + "/" + m.Module + "." + t.Name, firstLine(t.Description), also})
}
}
}
out := map[string]any{"matches": hits}
if len(hits) > searchCap {
out["matches"] = hits[:searchCap]
out["more"] = fmt.Sprintf("%d more; narrow the words", len(hits)-searchCap)
}
if len(hits) == 0 {
out["matches"] = []hit{}
out["hint"] = "nothing matched every word; try fewer words, or mesh_overview and mesh_machine to browse"
}
return answerText(out)
case verbDescribe:
t, err := resolve(x, str("address"))
if err != nil {
return failure(err.Error())
}
description := t.Tool.Description
if description == "" {
description = t.Name
}
return answerText(map[string]any{"address": t.Address, "description": description,
"arguments": schemaWithoutNode(t.Tool.Input)})
case verbCall:
t, err := resolve(x, str("address"))
if err != nil {
return failure(err.Error())
}
callArgs := map[string]any{}
if a, ok := args["arguments"].(map[string]any); ok {
for k, v := range a {
if k != "node" {
callArgs[k] = v
}
}
}
var got bus.Answered
if t.Seat {
got, err = s.conn.Ask(t.Key, callArgs, "")
} else {
got, err = callTool(s.conn, t.Key, callArgs, nil, x.Listing)
}
if err != nil {
return failure(whyItFailed(t.Address, err))
}
content := []map[string]any{{"type": "text", "text": pretty(got.Result)}}
if got.Node != "" {
content = append(content, map[string]any{"type": "text", "text": "answered by " + got.Node})
}
return map[string]any{"content": content}
}
return failure("no discovery verb " + name)
}
+222
View File
@@ -0,0 +1,222 @@
package console
import (
"encoding/json"
"fmt"
"strings"
"testing"
"time"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
"github.com/novox/mesh-tools/node-tools/internal/runtime"
)
// text is a tool result's first text, and whether it was an error.
func text(t *testing.T, reply map[string]any) (string, bool) {
t.Helper()
result, ok := reply["result"].(map[string]any)
if !ok {
t.Fatalf("no result: %v", reply)
}
content := result["content"].([]any)
isErr, _ := result["isError"].(bool)
var parts []string
for _, c := range content {
parts = append(parts, c.(map[string]any)["text"].(string))
}
return strings.Join(parts, "\n"), isErr
}
func call(t *testing.T, endpoint, tool string, args map[string]any) (string, bool) {
t.Helper()
body, _ := json.Marshal(map[string]any{"jsonrpc": "2.0", "id": 9, "method": "tools/call",
"params": map[string]any{"name": tool, "arguments": args}})
return text(t, post(t, endpoint, string(body)))
}
// novox/hq ADR 0195: the console announces five tools, and everything the mesh answers is reached
// through them by one address per layer.
func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
was := IndexKept
IndexKept = 0 // every discovery asks the mesh, so a module arriving mid-test is found
t.Cleanup(func() { IndexKept = was })
mesh := mt.New(t)
// alpha: interchangeable (the mesh issued it a plain subject); beta: state on its machine, holds
// the node-shelf seat there.
mesh.Issue(t, mt.MembershipOf("alpha", "desk", true, nil))
mesh.Issue(t, mt.MembershipOf("beta", "desk", false, map[string][]string{"node-shelf": {"list", "clear"}}))
nodeTools := connect(t, "node-tools", "desk")
stop, err := runtime.Run(nodeTools, []runtime.Served{
{Module: "alpha", Entrypoints: []string{mt.Fixture("many-alpha.serve.mjs")}},
{Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}},
}, nil, (&mt.Logs{}).Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
catalogue := connect(t, "mesh-catalog", "")
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) {
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "gamma"}}}, nil
})
defer stopCat()
// The controller, answering as its seat verbs do: the command's printed output.
controller := connect(t, "mesh-controller", "")
out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} }
serve := func(verb string, answer func() any) {
stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) {
return answer(), nil
})
if err != nil {
t.Fatal(err)
}
t.Cleanup(stop)
}
serve("tools", func() any {
return map[string]any{"seats": []map[string]any{
{"seat": "mesh-controller", "scope": "mesh", "tools": []map[string]any{
{"name": "nodes", "description": "Every machine the mesh knows.", "input": map[string]any{}}}},
{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
{"name": "clear", "description": "take it all off", "input": map[string]any{}}}},
}}
})
serve("nodes", func() any {
return out("bench 3m ago converged 1f2e\ndesk here converged 9a8b\n")
})
serve("seats", func() any {
return out("a warning the command printed first\n" + `{
"seats": [
{"seat": "mesh-controller", "scope": "mesh", "decision": "x", "holders": [{"module": "mesh-controller", "node": "bench"}]},
{"seat": "node-shelf", "scope": "node", "decision": "y", "holders": [{"module": "beta", "node": "desk"}]}
]
}`)
})
gammaOn := "nothing"
serve("modules", func() any {
return out(fmt.Sprintf("alpha 1 built 1a2b3c4d on desk\n needs container-runtime\n"+
"beta 1 built 1a2b3c4d on desk\n"+
"gamma 1 built 1a2b3c4d on %s\n", gammaOn))
})
catalogue.Flush()
controller.Flush()
up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer up.Close()
endpoint := "http://" + up.Address + "/mcp"
// Five tools, nothing else.
listed := post(t, endpoint, `{"jsonrpc":"2.0","id":2,"method":"tools/list"}`)["result"].(map[string]any)
var names []string
for _, x := range listed["tools"].([]any) {
name := x.(map[string]any)["name"].(string)
names = append(names, name)
for _, r := range name {
if !(r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '_' || r == '-') || len(name) > 64 {
t.Errorf("%q is not a name the API takes", name)
}
}
}
if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call" {
t.Errorf("announced %s", got)
}
// The overview names the mesh's seats, the machines' seats and the machines.
overview, isErr := call(t, endpoint, "mesh_overview", nil)
if isErr || !strings.Contains(overview, "mesh-controller.nodes") || !strings.Contains(overview, "<node>/node-shelf.list") ||
!strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) {
t.Errorf("overview: %s", overview)
}
machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"})
if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") ||
!strings.Contains(machine, "desk/alpha.one") {
t.Errorf("machine: %s", machine)
}
// A mesh seat, a node seat, an assignment and an interchangeable module, each by address.
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "mesh-controller.nodes"}); isErr || !strings.Contains(got, "converged") {
t.Errorf("mesh seat: %s", got)
}
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "desk/node-shelf.list"}); isErr || !strings.Contains(got, `"a"`) {
t.Errorf("node seat: %s", got)
}
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "desk/beta.three"}); isErr ||
!strings.Contains(got, `"beta": 3`) || !strings.Contains(got, "answered by desk") {
t.Errorf("assignment: %s", got)
}
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "alpha.one"}); isErr || !strings.Contains(got, `"alpha": 1`) {
t.Errorf("interchangeable module: %s", got)
}
// Refused, by name, where the address does not say enough or says the wrong thing.
for address, want := range map[string]string{
"beta.three": "keeps state on each machine it runs on, so a call names the machine: write <node>/beta.three — it runs on desk",
"node-shelf.list": "held once per machine: write <node>/node-shelf.list — it is held on desk",
"desk/mesh-controller.nodes": "held once for the whole mesh",
"bench/beta.three": "beta does not run on bench; it runs on desk",
"nonsense": "is not an address",
} {
got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": address})
if !isErr || !strings.Contains(got, want) {
t.Errorf("%s: %s (want %q)", address, got, want)
}
}
// Described without `node`: the address carries the machine.
described, isErr := call(t, endpoint, "mesh_describe", map[string]any{"address": "desk/beta.three"})
if isErr || strings.Contains(described, `"node"`) || !strings.Contains(described, `"address": "desk/beta.three"`) {
t.Errorf("describe: %s", described)
}
// Search finds across the layers; a module that starts serving after the first answer is found.
if got, _ := call(t, endpoint, "mesh_search", map[string]any{"query": "shelf"}); !strings.Contains(got, "desk/node-shelf.list") {
t.Errorf("search a node seat: %s", got)
}
if got, _ := call(t, endpoint, "mesh_search", map[string]any{"query": "gamma"}); strings.Contains(got, "gamma.given") {
t.Fatalf("gamma was found before it served: %s", got)
}
mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil))
gammaOn = "desk"
late := connect(t, "node-tools", "desk")
stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}},
nil, (&mt.Logs{}).Logf)
if err != nil {
t.Fatal(err)
}
defer stopLate()
var found string
for i := 0; i < 30; i++ {
found, _ = call(t, endpoint, "mesh_search", map[string]any{"query": "gamma given"})
if strings.Contains(found, "desk/gamma.given") {
break
}
time.Sleep(100 * time.Millisecond)
}
if !strings.Contains(found, "desk/gamma.given") {
t.Errorf("a module that arrived later was not found: %s", found)
}
// The old names still answer, unannounced.
if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) {
t.Errorf("an old name: %s", got)
}
}
func TestTheControllersPrintedListsAreRead(t *testing.T) {
machines := machinesIn("the bus's user list leaves out 2 user(s)\nace 2m ago converged 0c1d\n" +
"g14 here converged 77aa\nnovox 5s ago adopted 3e4f\n")
if got := strings.Join(machines, ","); got != "ace,g14,novox" {
t.Errorf("machines: %s", got)
}
on := assignmentsIn("baserow 1 built 7c800705 on ace\n requires postgres-database\n" +
"confluence 1 built 7c800705 on nothing\n" +
"mesh-wireguard 1 with the control plane on ace, g14, novox, shanks\n")
if strings.Join(on["baserow"], ",") != "ace" || len(on["confluence"]) != 0 || strings.Join(on["mesh-wireguard"], ",") != "ace,g14,novox,shanks" {
t.Errorf("assignments: %v", on)
}
}
+240
View File
@@ -0,0 +1,240 @@
package console
import (
"encoding/json"
"errors"
"fmt"
"regexp"
"sort"
"strings"
"sync"
"github.com/novox/mesh-tools/node-tools/internal/bus"
)
const (
catalogueModules = "mesh-catalog.catalog_modules"
seatTools = "seat:mesh-controller.tools"
toolsVerb = "tools"
)
// Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it.
type Tool struct {
Module string
Name string
Description string
Input json.RawMessage
Seat bool
Scope string
Subjects []string
}
// Listing is what the mesh could say about its tools; silence is named, never dropped (design 34 §3).
type Listing struct {
Tools []Tool
NotAnswering []string
}
// Seats is the seats and their verbs, so `<seat>.<verb>` resolves to the role.
type Seats map[string]map[string]bool
func seatsIn(l *Listing) Seats {
out := Seats{}
for _, t := range l.Tools {
if !t.Seat {
continue
}
if out[t.Module] == nil {
out[t.Module] = map[string]bool{}
}
out[t.Module][t.Name] = true
}
return out
}
// toolKey is the key a call uses: a role's when the prefix is a seat declaring that verb.
func toolKey(name string, seats Seats) string {
if strings.HasPrefix(name, "seat:") {
return name
}
dot := strings.Index(name, ".")
if dot < 0 {
return name
}
if seats[name[:dot]][name[dot+1:]] {
return "seat:" + name
}
return name
}
// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it
// serves, the controller's seat every role's tools — at once, so a restarting control plane hides
// nothing else.
func toolsOn(conn *bus.Conn) (*Listing, error) {
type rolesAnswer struct {
Seats []struct {
Seat string `json:"seat"`
Scope string `json:"scope"`
Tools []struct {
Name string `json:"name"`
Description string `json:"description"`
Input json.RawMessage `json:"input"`
} `json:"tools"`
} `json:"seats"`
}
var wg sync.WaitGroup
var roles *rolesAnswer
wg.Add(1)
go func() {
defer wg.Done()
if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil {
var r rolesAnswer
if json.Unmarshal(got.Result, &r) == nil {
roles = &r
}
}
}()
answered, err := conn.Ask(catalogueModules, map[string]any{}, "")
if err != nil {
wg.Wait()
return nil, err
}
var held struct {
Modules []struct {
Module string `json:"module"`
} `json:"modules"`
}
_ = json.Unmarshal(answered.Result, &held)
names := make([]string, 0, len(held.Modules))
for _, m := range held.Modules {
if m.Module != "" {
names = append(names, m.Module)
}
}
type outcome struct {
ok bool
answer struct {
Tools *[]struct {
Name string `json:"name"`
Description string `json:"description"`
Input json.RawMessage `json:"input"`
Subjects []string `json:"subjects"`
} `json:"tools"`
Failed *string `json:"failed"`
}
}
outcomes := make([]outcome, len(names))
for i, module := range names {
wg.Add(1)
go func(i int, module string) {
defer wg.Done()
got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "")
if err != nil {
return
}
if json.Unmarshal(got.Result, &outcomes[i].answer) == nil {
outcomes[i].ok = true
}
}(i, module)
}
wg.Wait()
l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
if roles != nil {
for _, s := range roles.Seats {
for _, t := range s.Tools {
l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description,
Input: t.Input, Seat: true, Scope: s.Scope})
}
}
} else {
l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)")
}
for i, module := range names {
o := outcomes[i]
switch {
case o.ok && o.answer.Failed != nil:
l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed))
case o.ok && o.answer.Tools != nil:
for _, t := range *o.answer.Tools {
l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description,
Input: t.Input, Subjects: t.Subjects})
}
default:
l.NotAnswering = append(l.NotAnswering, module)
}
}
sort.SliceStable(l.Tools, func(i, j int) bool {
return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name
})
sort.Strings(l.NotAnswering)
return l, nil
}
// callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one.
func callTool(conn *bus.Conn, key string, args any, seats Seats, l *Listing) (bus.Answered, error) {
name, node, _ := strings.Cut(key, "@")
if !strings.Contains(name, ".") {
return bus.Answered{}, fmt.Errorf("%q does not name a tool: write <module>.<tool>, as `mesh tools` lists "+
"them, or <module>.<tool>@<node> for the instance on one machine", key)
}
resolved := toolKey(name, seats)
if node != "" {
resolved += "@" + node
}
return conn.Ask(resolved, args, subjectListed(name, node, l))
}
func subjectListed(name, node string, l *Listing) string {
if l == nil {
return ""
}
dot := strings.Index(name, ".")
module, tool := name[:dot], name[dot+1:]
for _, t := range l.Tools {
if t.Module == module && t.Name == tool && !t.Seat {
if len(t.Subjects) == 0 {
return ""
}
if node == "" {
return t.Subjects[0]
}
for _, s := range t.Subjects {
if strings.HasSuffix(s, "."+node) {
return s
}
}
return ""
}
}
return ""
}
var (
noResponders = regexp.MustCompile(`(?i)no responders|503`)
refused = regexp.MustCompile(`(?i)permissions violation|authorization`)
timedOut = regexp.MustCompile(`(?i)timeout`)
)
// whyItFailed says why a call failed, so the remedy is in the words.
func whyItFailed(key string, err error) string {
if err == nil {
err = errors.New("failed")
}
msg := err.Error()
switch {
case noResponders.MatchString(msg):
extra := ""
if strings.HasPrefix(key, "seat:") {
extra = ", or nothing holds that seat"
}
return "nothing serves " + key + ". The module may not be assigned to any machine, or it is down" +
extra + " — `mesh tools` lists what answered."
case refused.MatchString(msg):
return "this account may not call " + key + ". What it may call was fixed when it was issued — a " +
"person's by `operator issue`, the console's by its manifest."
case timedOut.MatchString(msg):
return key + " did not answer in time. Something is serving it, so this is the tool being slow " +
"rather than absent."
}
return key + " failed: " + msg
}
+138
View File
@@ -0,0 +1,138 @@
package console
import (
"bytes"
"encoding/json"
"io"
"net/http"
"strings"
"testing"
"github.com/novox/mesh-tools/node-tools/internal/bus"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
"github.com/novox/mesh-tools/node-tools/internal/runtime"
)
func connect(t *testing.T, module, node string) *bus.Conn {
t.Helper()
c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node})
if err != nil {
t.Fatal(err)
}
c.Logf = func(string, ...any) {}
t.Cleanup(c.Close)
return c
}
func post(t *testing.T, endpoint string, body string) map[string]any {
t.Helper()
res, err := http.Post(endpoint, "application/json", bytes.NewBufferString(body))
if err != nil {
t.Fatal(err)
}
defer res.Body.Close()
raw, _ := io.ReadAll(res.Body)
var out map[string]any
if err := json.Unmarshal(raw, &out); err != nil {
t.Fatalf("%d %s", res.StatusCode, raw)
}
return out
}
// As node-tools, the runtime serves the bundles and is the console on loopback: the listing is what
// the modules and the mesh's records answered, and a call reaches the module on the machine named.
func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
mesh := mt.New(t)
mesh.Issue(t, mt.MembershipOf("alpha", "desk", true, nil))
mesh.Issue(t, mt.MembershipOf("beta", "desk", false, map[string][]string{"node-shelf": {"list", "clear"}}))
nodeTools := connect(t, "node-tools", "desk")
stop, err := runtime.Run(nodeTools, []runtime.Served{
{Module: "alpha", Entrypoints: []string{mt.Fixture("many-alpha.serve.mjs")}},
{Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}},
}, nil, (&mt.Logs{}).Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
catalogue := connect(t, "mesh-catalog", "")
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) {
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "ghost"}}}, nil
})
defer stopCat()
controller := connect(t, "mesh-controller", "")
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) {
return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
{"name": "clear", "description": "take it all off", "input": map[string]any{}},
}}}}, nil
})
defer stopSeat()
catalogue.Flush()
controller.Flush()
flat := NewSurface(nodeTools, "desk.node-tools")
flat.Flat = true // the whole catalogue, as before ADR 0195: still reachable, no longer announced
up, err := Serve(flat, "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer up.Close()
endpoint := "http://" + up.Address + "/mcp"
init := post(t, endpoint, `{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-03-26","capabilities":{}}}`)
if !strings.Contains(init["result"].(map[string]any)["instructions"].(string), "reached as desk.node-tools") {
t.Errorf("initialize: %v", init)
}
listed := post(t, endpoint, `{"jsonrpc":"2.0","id":2,"method":"tools/list"}`)["result"].(map[string]any)
var names []string
for _, x := range listed["tools"].([]any) {
names = append(names, x.(map[string]any)["name"].(string))
}
if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" {
t.Errorf("listed %s", got)
}
if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost" {
t.Errorf("not answering: %v", got)
}
for _, x := range listed["tools"].([]any) {
tool := x.(map[string]any)
if tool["name"] == "node-shelf.list" {
schema := tool["inputSchema"].(map[string]any)
if req, _ := schema["required"].([]any); len(req) != 1 || req[0] != "node" {
t.Errorf("a node seat's verb does not require node: %v", schema)
}
}
}
called := post(t, endpoint, `{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"alpha.one","arguments":{"node":"desk"}}}`)
content := called["result"].(map[string]any)["content"].([]any)
var got map[string]any
_ = json.Unmarshal([]byte(content[0].(map[string]any)["text"].(string)), &got)
if got["alpha"] != float64(1) || content[1].(map[string]any)["text"] != "answered by desk" {
t.Errorf("called: %v", called)
}
seat := post(t, endpoint, `{"jsonrpc":"2.0","id":4,"method":"tools/call","params":{"name":"node-shelf.list","arguments":{"node":"desk"}}}`)
if text := seat["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string); !strings.Contains(text, `"a"`) {
t.Errorf("seat verb: %v", seat)
}
refused := post(t, endpoint, `{"jsonrpc":"2.0","id":5,"method":"tools/call","params":{"name":"node-shelf.list","arguments":{}}}`)
if refused["error"] == nil {
t.Errorf("a node seat's verb was called without its machine: %v", refused)
}
absent := post(t, endpoint, `{"jsonrpc":"2.0","id":6,"method":"tools/call","params":{"name":"ghost.boo","arguments":{}}}`)
result := absent["result"].(map[string]any)
if result["isError"] != true || !strings.Contains(result["content"].([]any)[0].(map[string]any)["text"].(string), "nothing serves ghost.boo") {
t.Errorf("an absent tool: %v", absent)
}
res, err := http.Post(endpoint, "application/json", strings.NewReader(`{"jsonrpc":"2.0","method":"notifications/initialized"}`))
if err != nil || res.StatusCode != 202 {
t.Errorf("a notification: %v %v", res, err)
}
}
func TestTheConsoleListensOnLoopbackAndNowhereElse(t *testing.T) {
if _, err := Serve(NewSurface(nil, "x"), "0.0.0.0:0"); err == nil || !strings.Contains(err.Error(), "loopback and nowhere else") {
t.Errorf("a non-loopback console was not refused: %v", err)
}
}
+125
View File
@@ -0,0 +1,125 @@
package console
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"strconv"
"strings"
"github.com/novox/mesh-tools/node-tools/internal/wire"
)
// bodyLimit is the most a request body may be.
const bodyLimit = 1 << 20
var loopback = map[string]bool{"127.0.0.1": true, "::1": true, "localhost": true, "[::1]": true}
// Listening is a console that listens: where, with the port the machine gave, and how to stop it.
type Listening struct {
Address string
Close func() error
}
// Serve listens on host:port, refused unless the host is loopback — said before binding, so a
// console that would open to a network is a startup failure (ADR 0152).
func Serve(s *Surface, listen string) (*Listening, error) {
at := strings.LastIndex(listen, ":")
if at < 0 {
return nil, fmt.Errorf("%q is not host:port", listen)
}
host, portText := listen[:at], listen[at+1:]
if !loopback[host] {
return nil, fmt.Errorf(`the console listens on loopback and nowhere else (novox/hq ADR 0152): %q is not this `+
"machine's own address — whoever is on the machine owns the mesh there, and nobody else may reach this", host)
}
port, err := strconv.Atoi(portText)
if err != nil || port < 0 || port > 65535 {
return nil, fmt.Errorf("%q is not a port", portText)
}
ln, err := net.Listen("tcp", net.JoinHostPort(strings.Trim(host, "[]"), portText))
if err != nil {
return nil, err
}
server := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { route(w, r, s) })}
go func() { _ = server.Serve(ln) }()
bound := ln.Addr().(*net.TCPAddr).Port
return &Listening{Address: host + ":" + strconv.Itoa(bound), Close: func() error { return server.Shutdown(context.Background()) }}, nil
}
func route(w http.ResponseWriter, r *http.Request, s *Surface) {
switch r.URL.Path {
case "/":
w.Header().Set("content-type", "text/plain; charset=utf-8")
_, _ = io.WriteString(w, "the mesh's console: MCP over HTTP at POST /mcp (novox/hq design 34)\n")
return
case "/mcp":
default:
writeJSON(w, 404, map[string]any{"error": "the console serves /mcp and nothing else"})
return
}
switch r.Method {
case http.MethodPost:
case http.MethodDelete:
w.WriteHeader(204) // no session to end
return
default:
w.Header().Set("allow", "POST, DELETE")
w.WriteHeader(405)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, bodyLimit+1))
if err != nil || len(body) > bodyLimit {
if err == nil {
err = fmt.Errorf("the request is larger than %d bytes", bodyLimit)
}
writeJSON(w, 413, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32600, "message": err.Error()}})
return
}
trimmed := strings.TrimSpace(string(body))
if strings.HasPrefix(trimmed, "[") {
var batch []Request
if err := json.Unmarshal(body, &batch); err != nil {
writeJSON(w, 400, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32700, "message": "the body is not JSON"}})
return
}
replies := []*Reply{}
for _, req := range batch {
if reply := s.Handle(req); reply != nil {
replies = append(replies, reply)
}
}
if len(replies) == 0 {
w.WriteHeader(202)
return
}
writeJSON(w, 200, replies)
return
}
var req Request
if err := json.Unmarshal(body, &req); err != nil {
writeJSON(w, 400, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32700, "message": "the body is not JSON"}})
return
}
reply := s.Handle(req)
if reply == nil {
w.WriteHeader(202)
return
}
writeJSON(w, 200, reply)
}
func writeJSON(w http.ResponseWriter, status int, body any) {
text, err := wire.Marshal(body)
if err != nil {
text, _ = json.Marshal(map[string]any{"error": errors.New("unencodable answer").Error()})
}
w.Header().Set("content-type", "application/json; charset=utf-8")
w.Header().Set("content-length", strconv.Itoa(len(text)))
w.WriteHeader(status)
_, _ = w.Write(text)
}
+297
View File
@@ -0,0 +1,297 @@
// Package console is the mesh's tools as an MCP server on a machine's loopback (novox/hq design 34,
// ADR 0152, ADR 0175 §6): the Go port of node-tools' http.ts, mcp.ts and client.ts. A thin adapter:
// every tool listed is one a module answered for, the schema is the module's, the answer the module's.
package console
import (
"encoding/json"
"os"
"strings"
"sync"
"time"
"github.com/novox/mesh-tools/node-tools/internal/bus"
)
// Protocol is the MCP version spoken to an agent host.
const Protocol = "2025-03-26"
// ListingKept is how long a fetched tool list is kept before the modules are asked again.
var ListingKept = 30 * time.Second
// Request is one JSON-RPC message from a host.
type Request struct {
JSONRPC string `json:"jsonrpc"`
ID json.RawMessage `json:"id,omitempty"`
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
}
// Reply is one JSON-RPC answer.
type Reply struct {
JSONRPC string `json:"jsonrpc"`
ID json.RawMessage `json:"id"`
Result any `json:"result,omitempty"`
Error *RPCError `json:"error,omitempty"`
}
// RPCError is a protocol-level refusal.
type RPCError struct {
Code int `json:"code"`
Message string `json:"message"`
}
// Surface answers MCP requests over one bus connection, as one account.
type Surface struct {
conn *bus.Conn
who string
mu sync.Mutex
known *Listing
at time.Time
idx *index
idxAt time.Time
// Flat announces the whole catalogue, as the console did before ADR 0195: for a person reading
// it or a client that wants it. Off by default; MESH_CONSOLE_FLAT=1 turns it on.
Flat bool
}
// NewSurface is the surface over a connection, as `who`.
func NewSurface(conn *bus.Conn, who string) *Surface {
return &Surface{conn: conn, who: who, Flat: os.Getenv("MESH_CONSOLE_FLAT") == "1"}
}
func (s *Surface) listing() (*Listing, error) {
s.mu.Lock()
if s.known != nil && time.Since(s.at) <= ListingKept {
l := s.known
s.mu.Unlock()
return l, nil
}
s.mu.Unlock()
l, err := toolsOn(s.conn)
if err != nil {
return nil, err
}
s.mu.Lock()
s.known, s.at = l, time.Now()
s.mu.Unlock()
return l, nil
}
func isNotification(id json.RawMessage) bool {
t := strings.TrimSpace(string(id))
return t == "" || t == "null"
}
func answer(id json.RawMessage, result any) *Reply {
return &Reply{JSONRPC: "2.0", ID: idOrNull(id), Result: result}
}
func refuse(id json.RawMessage, code int, message string) *Reply {
return &Reply{JSONRPC: "2.0", ID: idOrNull(id), Error: &RPCError{Code: code, Message: message}}
}
func idOrNull(id json.RawMessage) json.RawMessage {
if isNotification(id) {
return json.RawMessage("null")
}
return id
}
// Handle answers one request; nil for a notification, which expects none.
func (s *Surface) Handle(r Request) *Reply {
notification := isNotification(r.ID)
switch r.Method {
case "initialize":
return answer(r.ID, map[string]any{
"protocolVersion": Protocol,
"capabilities": map[string]any{"tools": map[string]any{}},
"serverInfo": map[string]any{"name": "mesh", "version": "1"},
"instructions": s.instructions(),
})
case "notifications/initialized":
return nil
case "ping":
if notification {
return nil
}
return answer(r.ID, map[string]any{})
case "tools/list":
if !s.Flat {
return answer(r.ID, map[string]any{"tools": discovery()})
}
l, err := s.listing()
if err != nil {
return refuse(r.ID, -32603, whyItFailed(catalogueModules, err))
}
tools := make([]map[string]any, 0, len(l.Tools))
for _, t := range l.Tools {
var schema map[string]any
switch {
case t.Seat && t.Scope != "node":
schema = asSchema(t.Input)
case t.Seat:
schema = withNode(asSchema(t.Input), "the machine whose seat answers; required, the seat is held once per machine", true)
default:
schema = withNode(asSchema(t.Input), "", false)
}
description := t.Description
if description == "" {
description = t.Name + ", served by " + t.Module
}
tools = append(tools, map[string]any{"name": t.Module + "." + t.Name, "description": description, "inputSchema": schema})
}
return answer(r.ID, map[string]any{"tools": tools, "_meta": map[string]any{"notAnswering": l.NotAnswering}})
case "tools/call":
var p struct {
Name string `json:"name"`
Arguments map[string]any `json:"arguments"`
}
_ = json.Unmarshal(r.Params, &p)
if isDiscovery(p.Name) {
if p.Arguments == nil {
p.Arguments = map[string]any{}
}
return answer(r.ID, s.discover(p.Name, p.Arguments))
}
args := map[string]any{}
for k, v := range p.Arguments {
args[k] = v
}
l, _ := s.listing()
var roles Seats
if l != nil {
roles = seatsIn(l)
}
bare, _, _ := strings.Cut(p.Name, "@")
isSeatVerb := roles != nil && strings.HasPrefix(toolKey(bare, roles), "seat:")
nodeScoped := false
if isSeatVerb && l != nil {
for _, t := range l.Tools {
if t.Seat && t.Scope == "node" && t.Module+"."+t.Name == bare {
nodeScoped = true
}
}
}
takesNode := !isSeatVerb || nodeScoped
node := ""
if takesNode {
if n, ok := args["node"].(string); ok {
node = n
}
delete(args, "node")
}
if nodeScoped && node == "" && !strings.Contains(p.Name, "@") {
return refuse(r.ID, -32602, p.Name+" is a machine's seat's verb: name the machine with `node`")
}
name := p.Name
if node != "" && !strings.Contains(p.Name, "@") {
name = p.Name + "@" + node
}
got, err := callTool(s.conn, name, args, roles, l)
if err != nil {
return answer(r.ID, map[string]any{
"content": []map[string]any{{"type": "text", "text": whyItFailed(name, err)}},
"isError": true,
})
}
content := []map[string]any{{"type": "text", "text": pretty(got.Result)}}
if got.Node != "" {
content = append(content, map[string]any{"type": "text", "text": "answered by " + got.Node})
}
return answer(r.ID, map[string]any{"content": content})
}
if notification {
return nil
}
return refuse(r.ID, -32601, "mesh's MCP surface has no "+r.Method)
}
// pretty is a module's answer as JSON text, indented as JSON.stringify(result, null, 2) writes it.
func pretty(raw json.RawMessage) string {
if len(raw) == 0 {
return "null"
}
var v any
if json.Unmarshal(raw, &v) != nil {
return string(raw)
}
b, err := json.MarshalIndent(v, "", " ")
if err != nil {
return string(raw)
}
return strings.NewReplacer(`<`, "<", `>`, ">", `&`, "&").Replace(string(b))
}
// asSchema is a module's declared input as a JSON schema: wrapped when it is a bare map of
// properties, passed through when it is a schema, empty when nothing was declared.
func asSchema(raw json.RawMessage) map[string]any {
var given map[string]any
if json.Unmarshal(raw, &given) != nil || given == nil {
return map[string]any{"type": "object", "properties": map[string]any{}}
}
if given["type"] == "object" {
return given
}
if _, has := given["properties"]; has {
return given
}
if len(given) == 0 {
return map[string]any{"type": "object", "properties": map[string]any{}}
}
return map[string]any{"type": "object", "properties": given}
}
// withNode adds the optional — or, for a node seat, required — `node` argument (ADR 0159).
func withNode(schema map[string]any, description string, required bool) map[string]any {
properties := map[string]any{}
if p, ok := schema["properties"].(map[string]any); ok {
for k, v := range p {
properties[k] = v
}
}
if _, has := properties["node"]; !has {
if description == "" {
description = "the machine to ask, when this module runs on several; else whichever answers, and the answer says which"
}
properties["node"] = map[string]any{"type": "string", "description": description}
}
out := map[string]any{}
for k, v := range schema {
out[k] = v
}
out["type"] = "object"
out["properties"] = properties
if required {
have := []any{}
if r, ok := schema["required"].([]any); ok {
have = r
}
hasNode := false
for _, x := range have {
hasNode = hasNode || x == "node"
}
if !hasNode {
have = append(have, "node")
}
out["required"] = have
}
return out
}
// instructions is what an agent host is told about this surface when it connects.
func (s *Surface) instructions() string {
if s.Flat {
return "These are the tools of a Novox mesh, reached as " + s.who + ". Every call goes to the module " +
"that serves it; what may be called was fixed when this account was issued, so a " +
"refusal means the account, not the tool. The list is what the running modules " +
"answered, plus every role's tools from the mesh's records — the mesh's own verbs " +
"(mesh-controller.status, .push, .assign …) among them; a module that did not answer " +
"is named in the list's _meta and can still be called by <module>.<tool>."
}
return "The tools of a Novox mesh, reached as " + s.who + ", found by address rather than listed " +
"whole (novox/hq ADR 0195). mesh_overview shows the mesh's seats and machines; mesh_machine one " +
"machine's seats and modules; mesh_search finds a tool by words; mesh_describe gives one tool's " +
"arguments; mesh_call calls it. " + grammar + " What may be called was fixed when this account " +
"was issued, so a refusal means the account, not the tool."
}
+401
View File
@@ -0,0 +1,401 @@
// Package launch starts a bundle the runtime serves and speaks MCP over stdio to it (novox/hq ADR
// 0188, ADR 0193): `initialize`, `tools/list` once, `tools/call` per call. A Go binary, a Python
// script and a Node launcher are the same thing here: an executable that answers those. The runtime
// knows no language; it starts the path it is given.
package launch
import (
"bufio"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"os/exec"
"regexp"
"strings"
"sync"
"syscall"
"time"
"golang.org/x/sys/unix"
"github.com/novox/mesh-tools/node-tools/internal/wire"
)
// Protocol is the MCP version spoken to a bundle.
const Protocol = "2025-03-26"
// How long a child has for its handshake, and a call before the caller is told it is slow.
var (
HandshakeTimeout = 10 * time.Second
CallTimeout = 30 * time.Second
)
// Tool is one tool a launched bundle listed, and how to call it.
type Tool struct {
Name string
Description string
Input json.RawMessage
Run func(args json.RawMessage) (json.RawMessage, error)
}
// Registration is the tools a bundle listed under one name: its module's, or a seat's.
type Registration struct {
Module string
Tools []Tool
}
// Publisher publishes an event a bundle asked the runtime to emit, as the bundle's module.
type Publisher func(params json.RawMessage) error
// Launched is a running bundle: what it registered, and how to stop it.
type Launched struct {
Registrations []Registration
Stop func()
}
// Executable says whether an entrypoint can be started: a bundle the runtime serves is executable,
// and one that is not was not built to be served (ADR 0193).
func Executable(path string) bool {
info, err := os.Stat(path)
if err != nil || info.IsDir() {
return false
}
return info.Mode().Perm()&0o111 != 0
}
type message struct {
JSONRPC string `json:"jsonrpc,omitempty"`
ID json.RawMessage `json:"id,omitempty"`
Method string `json:"method,omitempty"`
Params json.RawMessage `json:"params,omitempty"`
Result json.RawMessage `json:"result,omitempty"`
Error *struct {
Code int `json:"code"`
Message string `json:"message"`
} `json:"error,omitempty"`
}
type child struct {
cmd *exec.Cmd
stdin io.WriteCloser
writeMu sync.Mutex
mu sync.Mutex
pending map[int64]chan message
next int64
dead chan struct{}
why error
}
func (c *child) write(m any) error {
b, err := wire.Marshal(m)
if err != nil {
return err
}
c.writeMu.Lock()
defer c.writeMu.Unlock()
_, err = c.stdin.Write(append(b, '\n'))
return err
}
var stackLine = regexp.MustCompile(`^\s+at\s`)
// Start launches one bundle and learns its tools. It fails when the child cannot be started or does
// not complete the handshake. A child that exits later is started again on its next call.
func Start(module, entry string, env []string, publish Publisher, logf func(string, ...any)) (*Launched, error) {
var mu sync.Mutex
var current *child
stopped := false
start := func() (*child, error) {
cmd := exec.Command(entry)
cmd.Env = env
stdin, err := cmd.StdinPipe()
if err != nil {
return nil, err
}
stdout, err := cmd.StdoutPipe()
if err != nil {
return nil, err
}
stderr, err := cmd.StderrPipe()
if err != nil {
return nil, err
}
if err := cmd.Start(); err != nil {
return nil, err
}
c := &child{cmd: cmd, stdin: stdin, pending: map[int64]chan message{}, next: 1, dead: make(chan struct{})}
var lastSaid string
var saidMu sync.Mutex
stderrDone := make(chan struct{})
go func() {
defer close(stderrDone)
scan := bufio.NewScanner(stderr)
scan.Buffer(make([]byte, 64*1024), 1<<20)
for scan.Scan() {
line := scan.Text()
if strings.TrimSpace(line) == "" {
continue
}
logf("[%s] %s", module, line)
if !stackLine.MatchString(line) && !strings.HasPrefix(line, "Node.js v") {
saidMu.Lock()
lastSaid = strings.TrimSpace(line)
saidMu.Unlock()
}
}
}()
go func() {
scan := bufio.NewScanner(stdout)
scan.Buffer(make([]byte, 64*1024), 16<<20)
for scan.Scan() {
line := strings.TrimSpace(scan.Text())
if line == "" {
continue
}
var m message
if err := json.Unmarshal([]byte(line), &m); err != nil {
if len(line) > 120 {
line = line[:120]
}
logf("[mesh-tools] %s's bundle said something that is not a reply: %s", module, line)
continue
}
// The bundle asks the runtime to emit (ADR 0193): published as this module, answered
// once the bus has accepted it. Nothing else a bundle may ask.
if m.Method != "" {
go func(m message) {
if m.Method != "mesh/publish" {
if len(m.ID) > 0 {
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID,
"error": map[string]any{"code": -32601, "message": "the runtime answers no " + m.Method + " from a bundle"}})
}
return
}
err := publish(m.Params)
if len(m.ID) == 0 {
return
}
if err != nil {
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID,
"error": map[string]any{"code": -32000, "message": err.Error()}})
return
}
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": map[string]any{}})
}(m)
continue
}
var id int64
if json.Unmarshal(m.ID, &id) != nil {
continue
}
c.mu.Lock()
ch := c.pending[id]
delete(c.pending, id)
c.mu.Unlock()
if ch != nil {
ch <- m
}
}
<-stderrDone
err := cmd.Wait()
code := "0"
if exit := (*exec.ExitError)(nil); errors.As(err, &exit) {
if status, ok := exit.Sys().(syscall.WaitStatus); ok && status.Signaled() {
code = unix.SignalName(status.Signal())
} else {
code = fmt.Sprint(exit.ExitCode())
}
}
saidMu.Lock()
why := fmt.Sprintf("%s's bundle exited (%s)", module, code)
if lastSaid != "" {
why += ": " + lastSaid
}
saidMu.Unlock()
c.mu.Lock()
c.why = errors.New(why)
c.mu.Unlock()
close(c.dead)
mu.Lock()
if current == c {
current = nil
}
wasStopped := stopped
mu.Unlock()
if !wasStopped {
logf("[mesh-tools] %s; started again on its next call", why)
}
}()
if _, err := c.ask(module, "initialize", map[string]any{"protocolVersion": Protocol, "capabilities": map[string]any{},
"clientInfo": map[string]any{"name": "node-tools", "version": "1"}}, HandshakeTimeout); err != nil {
_ = cmd.Process.Kill()
return nil, err
}
_ = c.write(map[string]any{"jsonrpc": "2.0", "method": "notifications/initialized"})
return c, nil
}
asking := func(method string, params any, timeout time.Duration) (json.RawMessage, error) {
mu.Lock()
c := current
mu.Unlock()
if c == nil {
fresh, err := start()
if err != nil {
return nil, err
}
mu.Lock()
current = fresh
c = fresh
mu.Unlock()
}
return c.ask(module, method, params, timeout)
}
first, err := start()
if err != nil {
return nil, err
}
mu.Lock()
current = first
mu.Unlock()
raw, err := asking("tools/list", map[string]any{}, HandshakeTimeout)
if err != nil {
return nil, err
}
var listed struct {
Tools []struct {
Name string `json:"name"`
Description string `json:"description"`
InputSchema json.RawMessage `json:"inputSchema"`
} `json:"tools"`
}
if err := json.Unmarshal(raw, &listed); err != nil {
return nil, fmt.Errorf("%s's bundle listed its tools in a shape that is not MCP's: %w", module, err)
}
order := []string{}
groups := map[string][]Tool{}
for _, t := range listed.Tools {
under, name := module, t.Name
if dot := strings.Index(t.Name, "."); dot >= 0 {
under, name = t.Name[:dot], t.Name[dot+1:]
}
full := t.Name
input := t.InputSchema
if len(input) == 0 || string(input) == "null" {
input = json.RawMessage("{}")
}
tool := Tool{Name: name, Description: t.Description, Input: input,
Run: func(args json.RawMessage) (json.RawMessage, error) {
if len(args) == 0 || string(args) == "null" {
args = json.RawMessage("{}")
}
res, err := asking("tools/call", map[string]any{"name": full, "arguments": args}, CallTimeout)
if err != nil {
return nil, err
}
var called struct {
Content []struct {
Type string `json:"type"`
Text string `json:"text"`
} `json:"content"`
IsError bool `json:"isError"`
}
_ = json.Unmarshal(res, &called)
text := ""
for _, c := range called.Content {
if c.Type == "text" {
text = c.Text
break
}
}
if called.IsError {
if text == "" {
text = module + "." + full + " failed"
}
return nil, errors.New(text)
}
// The bundle's answer is JSON as text; handed back as the value it encodes.
if json.Valid([]byte(text)) && text != "" {
return json.RawMessage(text), nil
}
b, _ := json.Marshal(text)
return b, nil
}}
if _, seen := groups[under]; !seen {
order = append(order, under)
}
groups[under] = append(groups[under], tool)
}
out := &Launched{Stop: func() {
mu.Lock()
stopped = true
c := current
current = nil
mu.Unlock()
if c != nil && c.cmd.Process != nil {
_ = c.cmd.Process.Signal(syscall.SIGTERM)
}
}}
for _, under := range order {
out.Registrations = append(out.Registrations, Registration{Module: under, Tools: groups[under]})
}
return out, nil
}
func (c *child) ask(module, method string, params any, timeout time.Duration) (json.RawMessage, error) {
c.mu.Lock()
if c.why != nil {
c.mu.Unlock()
return nil, c.why
}
id := c.next
c.next++
ch := make(chan message, 1)
c.pending[id] = ch
c.mu.Unlock()
if err := c.write(map[string]any{"jsonrpc": "2.0", "id": id, "method": method, "params": params}); err != nil {
c.mu.Lock()
delete(c.pending, id)
c.mu.Unlock()
select {
case <-c.dead:
return nil, c.deathReason()
case <-time.After(100 * time.Millisecond):
return nil, err
}
}
timer := time.NewTimer(timeout)
defer timer.Stop()
select {
case m := <-ch:
if m.Error != nil {
msg := m.Error.Message
if msg == "" {
msg = "the bundle refused the request"
}
return nil, errors.New(msg)
}
return m.Result, nil
case <-c.dead:
return nil, c.deathReason()
case <-timer.C:
c.mu.Lock()
delete(c.pending, id)
c.mu.Unlock()
return nil, fmt.Errorf("%s's bundle did not answer %s in %ds", module, method, int(timeout/time.Second))
}
}
func (c *child) deathReason() error {
c.mu.Lock()
defer c.mu.Unlock()
if c.why != nil {
return c.why
}
return errors.New("the bundle exited")
}
+143
View File
@@ -0,0 +1,143 @@
// Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS
// and EVENTS streams, memberships issued by hand, and a fixture's path.
package meshtest
import (
"encoding/json"
"os"
"path/filepath"
"runtime"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-tools/node-tools/internal/bus"
)
// URL is the test bus, or the test is skipped.
func URL(t *testing.T) string {
t.Helper()
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
return url
}
// Mesh is the controller's job, done by hand.
type Mesh struct {
nc *nats.Conn
js nats.JetStreamContext
}
// New raises the streams afresh.
func New(t *testing.T) *Mesh {
t.Helper()
nc, err := nats.Connect(URL(t))
if err != nil {
t.Fatal(err)
}
js, _ := nc.JetStream()
for _, s := range []string{"ASSIGNMENTS", "EVENTS"} {
_ = js.DeleteStream(s)
}
if _, err := js.AddStream(&nats.StreamConfig{Name: "ASSIGNMENTS", Subjects: []string{"mesh.assignment.>"},
MaxMsgsPerSubject: 1, AllowDirect: true}); err != nil {
t.Fatal(err)
}
if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil {
t.Fatal(err)
}
t.Cleanup(nc.Close)
return &Mesh{nc: nc, js: js}
}
// Issue publishes a membership.
func (m *Mesh) Issue(t *testing.T, mem bus.Membership) {
t.Helper()
body, _ := json.Marshal(mem)
if _, err := m.js.Publish(bus.MembershipSubject(mem.Node, mem.Module), body); err != nil {
t.Fatal(err)
}
}
// NextEvent is the subject the next event under a pattern lands on.
func (m *Mesh) NextEvent(t *testing.T, pattern string) <-chan string {
t.Helper()
ch := make(chan string, 1)
sub, err := m.nc.Subscribe(pattern, func(msg *nats.Msg) {
select {
case ch <- msg.Subject:
default:
}
})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
_ = m.nc.Flush()
return ch
}
// MembershipOf is a membership as the controller issues one on a machine.
func MembershipOf(module, node string, plain bool, seats map[string][]string) bus.Membership {
own := "mesh.mod." + module
serves := []bus.Served{{Subject: own + ".tool.{tool}." + node}}
if plain {
serves = append(serves, bus.Served{Subject: own + ".tool.{tool}", Queue: "serve." + module})
}
var verbs []bus.SeatVerb
for seat, vs := range seats {
for _, v := range vs {
verbs = append(verbs, bus.SeatVerb{Seat: seat, Verb: v, Subject: "mesh.seat." + seat + ".tool." + v + "." + node})
}
}
return bus.Membership{Node: node, Module: module, Serves: serves, Seats: verbs,
Emits: own + ".event.{event}", Tools: own + ".tool.tools"}
}
// Fixture is a path in node-tools/test/fixtures, which these tests share with the TypeScript ones.
func Fixture(name string) string {
_, here, _, _ := runtime.Caller(0)
return filepath.Join(filepath.Dir(here), "..", "..", "test", "fixtures", name)
}
// Until retries while the bus answers "no responders" — something not yet served.
func Until(t *testing.T, try func() error) {
t.Helper()
var err error
for i := 0; i < 50; i++ {
if err = try(); err == nil {
return
}
time.Sleep(100 * time.Millisecond)
}
t.Fatal(err)
}
// Logs collects what the runtime says.
type Logs struct{ lines []string }
// Logf is a logger that keeps the lines.
func (l *Logs) Logf(format string, args ...any) {
l.lines = append(l.lines, sprintf(format, args...))
}
// Has says whether a line contains every fragment.
func (l *Logs) Has(fragments ...string) bool {
for _, line := range l.lines {
all := true
for _, f := range fragments {
all = all && strings.Contains(line, f)
}
if all {
return true
}
}
return false
}
// All is every line.
func (l *Logs) All() string { return strings.Join(l.lines, "\n") }
+5
View File
@@ -0,0 +1,5 @@
package meshtest
import "fmt"
func sprintf(format string, args ...any) string { return fmt.Sprintf(format, args...) }
+378
View File
@@ -0,0 +1,378 @@
// Package runtime is the node's tool runtime (novox/hq ADR 0175, ADR 0193), the Go port of
// node-tools' runtime.ts in its launch-only form: it launches every assigned module's bundle,
// serves each module's tools on that module's subjects and each held seat's verbs on the seat's,
// and answers for every module the verb that says what it serves. It imports nothing and knows no
// language.
package runtime
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"github.com/novox/mesh-tools/node-tools/internal/bus"
"github.com/novox/mesh-tools/node-tools/internal/launch"
)
// ToolsVerb is the verb every module's runtime answers for it (ADR 0152): its tools, from the code
// that answers them.
const ToolsVerb = "tools"
// Words the mesh sets for the runtime (ADR 0175, ADR 0192).
const (
ToolModules = "MESH_TOOL_MODULES"
ToolEnv = "MESH_TOOL_ENV"
OperatorAccount = "MESH_OPERATOR_ACCOUNT"
OperatorHome = "MESH_OPERATOR_HOME"
)
// Served is one module this runtime serves and its entrypoints.
type Served struct {
Module string
Entrypoints []string
}
// ToolsAnswer is what `tools` answers for one module.
type ToolsAnswer struct {
Module string `json:"module"`
Tools []ListedTool `json:"tools"`
Failed string `json:"failed,omitempty"`
}
// ListedTool is one tool as a module's `tools` answer lists it.
type ListedTool struct {
Name string `json:"name"`
Description string `json:"description"`
Input json.RawMessage `json:"input"`
Subjects []string `json:"subjects,omitempty"`
}
// ServedModulesFrom reads MESH_TOOL_MODULES: `<module>=<entrypoint>` entries, comma-separated, several
// per module. The one-module form — a bare path, or the runtime's own module — is the per-module
// containers' (to-be 38 WP4c) and refused here: the node's runtime imports nothing.
func ServedModulesFrom(spec, own string) ([]Served, error) {
order := []string{}
by := map[string][]string{}
for _, raw := range strings.Split(spec, ",") {
entry := strings.TrimSpace(raw)
if entry == "" {
continue
}
module, path, ok := strings.Cut(entry, "=")
module, path = strings.TrimSpace(module), strings.TrimSpace(path)
if !ok || module == "" || path == "" || module == own {
return nil, fmt.Errorf("%s: %q is not <module>=<entrypoint> of another module; the node's "+
"runtime launches the bundles it is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry)
}
if _, seen := by[module]; !seen {
order = append(order, module)
}
by[module] = append(by[module], path)
}
out := make([]Served, 0, len(order))
for _, m := range order {
out = append(out, Served{Module: m, Entrypoints: by[m]})
}
return out, nil
}
// TakeToolEnvs reads the composed environments (ADR 0192) and removes them from the process's, so no
// bundle finds another's there.
func TakeToolEnvs() (map[string]map[string]string, error) {
raw := os.Getenv(ToolEnv)
os.Unsetenv(ToolEnv)
out := map[string]map[string]string{}
if raw == "" {
return out, nil
}
var parsed map[string]map[string]any
if err := json.Unmarshal([]byte(raw), &parsed); err != nil {
return nil, fmt.Errorf(`%s is not JSON of the shape {"<module>": {"<word>": "<value>"}}: %w`, ToolEnv, err)
}
for module, words := range parsed {
own := map[string]string{}
for k, v := range words {
if s, ok := v.(string); ok {
own[k] = s
} else {
b, _ := json.Marshal(v)
own[k] = string(b)
}
}
out[module] = own
}
return out, nil
}
type registration struct {
module string // the module's own name, or a seat's
owner string // the module whose bundle made it
tools []launch.Tool
}
// Run launches, binds and serves. It answers a stop function.
func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, logf func(string, ...any)) (func(), error) {
if account := os.Getenv(OperatorAccount); account != "" {
home := ""
if h := os.Getenv(OperatorHome); h != "" {
home = " (home " + h + ")"
}
logf("[mesh-tools] the operator's account here is %s%s", account, home)
}
modules := make([]string, 0, len(served))
isServed := map[string]bool{}
for _, s := range served {
modules = append(modules, s.Module)
isServed[s.Module] = true
conn.Follow(s.Module)
}
node := conn.Node()
base := os.Environ()
envFor := func(module string) []string {
words := map[string]string{}
for _, kv := range base {
if k, v, ok := strings.Cut(kv, "="); ok && k != ToolEnv {
words[k] = v
}
}
for k, v := range envs[module] {
words[k] = v
}
words["MESH_SERVED_MODULE"] = module
words["MESH_MODULE"] = module
if node != "" {
words["MESH_NODE"] = node
}
out := make([]string, 0, len(words))
for k, v := range words {
out = append(out, k+"="+v)
}
sort.Strings(out)
return out
}
failed := map[string]string{}
var registrations []registration
var stops []func()
for _, s := range served {
for _, entry := range s.Entrypoints {
path, _ := filepath.Abs(entry)
module := s.Module
fail := func(why string) {
failed[module] = why
logf("[mesh-tools] %s's bundle %s failed to load: %s; its tools are not served here", module, entry, why)
}
if !launch.Executable(path) {
fail(path + " is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)")
continue
}
child, err := launch.Start(module, path, envFor(module), func(params json.RawMessage) error {
var env bus.Envelope
if err := json.Unmarshal(params, &env); err != nil {
return fmt.Errorf("not an event envelope: %w", err)
}
return conn.PublishAs(module, env)
}, logf)
if err != nil {
fail(err.Error())
continue
}
stops = append(stops, child.Stop)
for _, r := range child.Registrations {
registrations = append(registrations, registration{module: r.Module, owner: module, tools: r.Tools})
}
}
}
claimed := map[string]bool{}
for _, m := range modules {
if mem := conn.Membership(m); mem != nil {
for _, s := range mem.Seats {
claimed[s.Seat] = true
}
}
}
var own []registration
for _, r := range registrations {
switch {
case isServed[r.module]:
own = append(own, r)
case claimed[r.module]:
default:
logf(`[mesh-tools] %s registers tools under "%s", which is neither a module served here nor a seat one of them claims; not served until the mesh issues the claim`, r.owner, r.module)
}
}
stopAll := func() {
for i := len(stops) - 1; i >= 0; i-- {
stops[i]()
}
}
for _, r := range own {
seen := map[string]bool{}
for _, t := range r.tools {
if t.Name == ToolsVerb {
stopAll()
return nil, fmt.Errorf(`%s names a tool "%s", which is the verb the runtime answers for every module with what it serves (novox/hq ADR 0152) — refused, rename it`, r.module, ToolsVerb)
}
if seen[t.Name] {
stopAll()
return nil, fmt.Errorf("%s exposes two tools named %s — refused", r.module, t.Name)
}
seen[t.Name] = true
}
}
var names []string
byModule := map[string][]launch.Tool{}
for _, r := range own {
for _, t := range r.tools {
t := t
stop, err := conn.Handle(r.module+"."+t.Name, func(body json.RawMessage) (any, error) {
return t.Run(argsOf(body))
})
if err != nil {
logf("[mesh-tools] cannot serve %s.%s: %v", r.module, t.Name, err)
continue
}
names = append(names, r.module+"."+t.Name)
stops = append(stops, stop)
}
byModule[r.module] = append(byModule[r.module], r.tools...)
}
for _, module := range modules {
module := module
tools := byModule[module]
why := failed[module]
if len(tools) == 0 && why == "" {
continue // a pure-events module: silent, as it always was
}
stop, err := conn.Handle(module+"."+ToolsVerb, func(json.RawMessage) (any, error) {
answer := ToolsAnswer{Module: module, Tools: []ListedTool{}, Failed: why}
for _, t := range tools {
answer.Tools = append(answer.Tools, ListedTool{Name: t.Name, Description: t.Description,
Input: t.Input, Subjects: subjectsOf(conn, module, t.Name)})
}
return answer, nil
})
if err == nil {
stops = append(stops, stop)
}
}
failedNames := make([]string, 0, len(failed))
for _, m := range modules {
if _, f := failed[m]; f {
failedNames = append(failedNames, m)
}
}
line := fmt.Sprintf("[mesh-tools] serving %d tool(s) for %d module(s): %s", len(names), len(modules), orNone(names))
if len(failedNames) > 0 {
line += fmt.Sprintf("; not serving %s, whose bundle(s) failed to load", strings.Join(failedNames, ", "))
}
logf("%s", line)
stops = append(stops, serveSeats(conn, modules, registrations, logf))
conn.Flush()
return stopAll, nil
}
func orNone(names []string) string {
if len(names) == 0 {
return "(none)"
}
return strings.Join(names, ", ")
}
// argsOf is a call's arguments as the tool receives them: an object, `{}` for none.
func argsOf(body json.RawMessage) json.RawMessage {
trimmed := strings.TrimSpace(string(body))
if trimmed == "" || trimmed == "null" {
return json.RawMessage("{}")
}
return body
}
// subjectsOf is where a tool is answered as the mesh issued it: the plain subject first, then this
// machine's; nothing before a membership is issued.
func subjectsOf(conn *bus.Conn, module, tool string) []string {
m := conn.Membership(module)
if m == nil {
return nil
}
var plain, mine []string
for _, s := range m.Serves {
subject := strings.ReplaceAll(s.Subject, "{tool}", tool)
if s.Queue != "" {
plain = append(plain, subject)
} else {
mine = append(mine, subject)
}
}
return append(plain, mine...)
}
// serveSeats serves every verb of every seat a served module holds, where the mesh issued it, by
// the tool of the same name registered under the seat's name (ADR 0159, 0160) — and serves again
// whenever a membership changes. Whether this machine holds the seat is the bus's to decide.
func serveSeats(conn *bus.Conn, modules []string, registrations []registration, logf func(string, ...any)) func() {
impl := map[string]map[string]launch.Tool{}
for _, r := range registrations {
if impl[r.module] == nil {
impl[r.module] = map[string]launch.Tool{}
}
for _, t := range r.tools {
impl[r.module][t.Name] = t
}
}
var mu sync.Mutex
var stops []func()
serve := func() {
mu.Lock()
defer mu.Unlock()
for _, s := range stops {
s()
}
stops = nil
have := map[string]bool{}
for _, module := range modules {
m := conn.Membership(module)
if m == nil {
continue
}
for _, v := range m.Seats {
if have[v.Subject] {
continue
}
have[v.Subject] = true
t, ok := impl[v.Seat][v.Verb]
if !ok {
logf("[mesh-tools] %s claims %s and implements no %s, which that seat promises; not served", module, v.Seat, v.Verb)
continue
}
stop, err := conn.HandleSubject(v.Subject, func(body json.RawMessage) (any, error) {
return t.Run(argsOf(body))
})
if err != nil {
logf("[mesh-tools] cannot serve %s's %s on %s: %v", v.Seat, v.Verb, v.Subject, err)
continue
}
stops = append(stops, stop)
logf("[mesh-tools] serving %s's %s on %s, admitted where %s holds the seat", v.Seat, v.Verb, v.Subject, module)
}
}
}
serve()
conn.OnMembership(func(bus.Membership) { go func() { serve(); conn.Flush() }() })
return func() {
mu.Lock()
defer mu.Unlock()
for _, s := range stops {
s()
}
stops = nil
}
}
+251
View File
@@ -0,0 +1,251 @@
package runtime
import (
"encoding/json"
"os"
"strings"
"testing"
"time"
"github.com/novox/mesh-tools/node-tools/internal/bus"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
)
func connect(t *testing.T, module, node string) *bus.Conn {
t.Helper()
c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node})
if err != nil {
t.Fatal(err)
}
c.Logf = func(string, ...any) {}
t.Cleanup(c.Close)
return c
}
func call(t *testing.T, asker *bus.Conn, key string, args any) (string, error) {
t.Helper()
got, err := asker.Ask(key, args, "")
return string(got.Result), err
}
func same(t *testing.T, got, want string) {
t.Helper()
var a, b any
if json.Unmarshal([]byte(got), &a) != nil || json.Unmarshal([]byte(want), &b) != nil {
t.Fatalf("not JSON: got %s want %s", got, want)
}
ga, _ := json.Marshal(a)
gb, _ := json.Marshal(b)
if string(ga) != string(gb) {
t.Errorf("got %s, want %s", got, want)
}
}
// The node's runtime serves five modules' bundles on one credential — two TypeScript, one broken,
// one Python, one written against the protocol — and follows a re-issued membership (ADR 0175, 0193).
func TestTheNodesRuntimeServesFiveModulesAndFollowsAReissuedMembership(t *testing.T) {
mesh := mt.New(t)
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil))
mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}}))
mesh.Issue(t, mt.MembershipOf("gamma", "anchor", false, nil))
mesh.Issue(t, mt.MembershipOf("delta", "anchor", false, map[string][]string{"node-lamp": {"on"}}))
mesh.Issue(t, mt.MembershipOf("epsilon", "anchor", false, nil))
nodeTools := connect(t, "node-tools", "anchor")
asker := connect(t, "console", "workstation")
logs := &mt.Logs{}
t.Setenv("MESH_OPERATOR_ACCOUNT", "somebody")
t.Setenv("MESH_OPERATOR_HOME", "/home/somebody")
stop, err := Run(nodeTools, []Served{
{"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}},
{"beta", []string{mt.Fixture("many-beta.serve.mjs")}},
{"gamma", []string{mt.Fixture("many-broken.serve.mjs")}},
{"delta", []string{mt.Fixture("many-delta.py")}},
{"epsilon", []string{mt.Fixture("many-epsilon.mjs")}},
}, nil, logs.Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
if !logs.Has("the operator's account here is somebody (home /home/somebody)") {
t.Errorf("the operator was not said:\n%s", logs.All())
}
if !logs.Has("gamma's bundle", "many-broken.serve.mjs failed to load: gamma's bundle exited (1): Error: gamma's bundle cannot find its client; its tools are not served here") {
t.Errorf("the broken bundle was not named with its own words:\n%s", logs.All())
}
if !logs.Has("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") {
t.Errorf("not serving what it should:\n%s", logs.All())
}
for key, want := range map[string]string{
"alpha.one": `{"alpha":1}`, "alpha.one@anchor": `{"alpha":1}`, "beta.three@anchor": `{"beta":3}`,
"beta.four@anchor": `{"beta":4}`, "beta.five@anchor": `{"beta":5}`,
"seat:node-shelf.list@anchor": `{"shelf":["a","b"]}`, "seat:node-shelf.clear@anchor": `{"cleared":true}`,
"seat:node-lamp.on@anchor": `{"on":true,"language":"python"}`, "epsilon.seven@anchor": `{"epsilon":7,"via":"stdio"}`,
} {
got, err := call(t, asker, key, map[string]any{})
if err != nil {
t.Errorf("%s: %v", key, err)
continue
}
same(t, got, want)
}
if _, err := call(t, asker, "beta.three", map[string]any{}); err == nil || !strings.Contains(err.Error(), "no responders") {
t.Errorf("beta was not issued the plain subject, and answered on it: %v", err)
}
got, err := call(t, asker, "delta.greet@anchor", map[string]any{"who": "mesh"})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"greeting":"hello mesh","language":"python"}`)
// A launched tool that emits does so as its module, through the runtime.
landed := mesh.NextEvent(t, "mesh.mod.*.event.>")
got, err = call(t, asker, "alpha.two", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"alpha":2}`)
select {
case subject := <-landed:
if subject != "mesh.mod.alpha.event.happened" {
t.Errorf("the event landed on %s", subject)
}
case <-time.After(5 * time.Second):
t.Error("the event never landed")
}
// A bundle that dies mid-call is said, and started again on its next call.
if _, err := call(t, asker, "delta.die@anchor", map[string]any{}); err == nil {
t.Error("a bundle that died answered")
}
got, err = call(t, asker, "delta.greet@anchor", map[string]any{"who": "again"})
if err != nil {
t.Fatalf("not started again: %v", err)
}
same(t, got, `{"greeting":"hello again","language":"python"}`)
// `tools` answers for each, and why gamma serves nothing.
gamma, err := call(t, asker, "gamma.tools@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, gamma, `{"module":"gamma","tools":[],"failed":"gamma's bundle exited (1): Error: gamma's bundle cannot find its client"}`)
beta, _ := call(t, asker, "beta.tools@anchor", map[string]any{})
var answer ToolsAnswer
_ = json.Unmarshal([]byte(beta), &answer)
if len(answer.Tools) != 3 || answer.Tools[0].Name != "three" || strings.Join(answer.Tools[0].Subjects, ",") != "mesh.mod.beta.tool.three.anchor" {
t.Errorf("beta's tools answer: %s", beta)
}
// Re-issued mid-run, now answering for the module anywhere: served without a restart.
mesh.Issue(t, mt.MembershipOf("beta", "anchor", true, map[string][]string{"node-shelf": {"list", "clear"}}))
mt.Until(t, func() error { _, err := call(t, asker, "beta.three", map[string]any{}); return err })
got, err = call(t, asker, "seat:node-shelf.list@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"shelf":["a","b"]}`)
}
// Each bundle is given its own environment and none of another's (ADR 0192), and the composed
// environments are not left in the runtime's.
func TestEachBundleIsGivenItsOwnEnvironment(t *testing.T) {
mesh := mt.New(t)
for _, m := range []string{"gamma", "delta", "zeta"} {
mesh.Issue(t, mt.MembershipOf(m, "anchor", false, nil))
}
nodeTools := connect(t, "node-tools", "anchor")
asker := connect(t, "console", "workstation")
t.Setenv("MESH_OPERATOR_ACCOUNT", "somebody")
t.Setenv(ToolEnv, `{"gamma":{"GAMMA_CONFIG_FILE":"/var/lib/mesh/gamma/config.json"},"delta":{"DELTA_TOKEN_FILE":"/var/lib/mesh/delta/token"},"zeta":{"ZETA_URL":"http://127.0.0.1:3000"}}`)
envs, err := TakeToolEnvs()
if err != nil {
t.Fatal(err)
}
if _, left := os.LookupEnv(ToolEnv); left {
t.Error("the composed environments were left in the runtime's")
}
logs := &mt.Logs{}
stop, err := Run(nodeTools, []Served{
{"gamma", []string{mt.Fixture("env-gamma.serve.mjs")}},
{"delta", []string{mt.Fixture("env-delta.serve.mjs")}},
{"zeta", []string{mt.Fixture("env-zeta.mjs")}},
}, envs, logs.Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
for key, want := range map[string]string{
"gamma.given@anchor": `{"mine":"/var/lib/mesh/gamma/config.json","theirs":null,"runtime":"somebody","composed":null}`,
"delta.given@anchor": `{"mine":"/var/lib/mesh/delta/token","theirs":null}`,
"zeta.given@anchor": `{"mine":"http://127.0.0.1:3000","theirs":null,"composed":null}`,
} {
got, err := call(t, asker, key, map[string]any{})
if err != nil {
t.Fatalf("%s: %v\n%s", key, err, logs.All())
}
same(t, got, want)
}
}
// An entrypoint that is not executable is refused by name, and the others serve (ADR 0193).
func TestAnEntrypointThatIsNotExecutableIsRefused(t *testing.T) {
mesh := mt.New(t)
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil))
mesh.Issue(t, mt.MembershipOf("plain", "anchor", false, nil))
nodeTools := connect(t, "node-tools", "anchor")
asker := connect(t, "console", "workstation")
logs := &mt.Logs{}
stop, err := Run(nodeTools, []Served{
{"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}},
{"plain", []string{mt.Fixture("many-alpha.mjs")}},
}, nil, logs.Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
if !logs.Has("plain's bundle", "many-alpha.mjs failed to load:", "is not executable; a bundle the runtime serves is started, never imported") {
t.Errorf("the non-executable entrypoint was not refused by name:\n%s", logs.All())
}
got, err := call(t, asker, "alpha.one@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"alpha":1}`)
}
// A launched bundle is told the module it serves, so its seat's verbs stay the seat's (ADR 0193).
func TestALaunchedBundleRegisteringItsSeatFirstServesTheSeat(t *testing.T) {
mesh := mt.New(t)
mesh.Issue(t, mt.MembershipOf("theta", "anchor", false, map[string][]string{"node-shelf": {"list"}}))
nodeTools := connect(t, "node-tools", "anchor")
asker := connect(t, "console", "workstation")
stop, err := Run(nodeTools, []Served{{"theta", []string{mt.Fixture("served-seat-first.mjs")}}},
map[string]map[string]string{"theta": {"THETA_WORD": "given"}}, (&mt.Logs{}).Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
got, err := call(t, asker, "theta.own@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"theta":"given"}`)
got, err = call(t, asker, "seat:node-shelf.list@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"shelf":["x"]}`)
}
func TestToolModulesNamesOtherModulesOnly(t *testing.T) {
got, err := ServedModulesFrom(" alpha=/a/tools/index.serve.mjs, beta=/b/one, beta=/b/two ", "node-tools")
if err != nil || len(got) != 2 || got[1].Module != "beta" || len(got[1].Entrypoints) != 2 {
t.Fatalf("%v %v", got, err)
}
for _, bad := range []string{"/mine/index.js", "node-tools=/own.js", "=/x"} {
if _, err := ServedModulesFrom(bad, "node-tools"); err == nil {
t.Errorf("%q was accepted", bad)
}
}
}
+18
View File
@@ -0,0 +1,18 @@
// Package wire is JSON as the TypeScript runtime writes it: no HTML escaping of <, > and &.
package wire
import (
"bytes"
"encoding/json"
)
// Marshal encodes v the way JSON.stringify does, without a trailing newline.
func Marshal(v any) ([]byte, error) {
var b bytes.Buffer
enc := json.NewEncoder(&b)
enc.SetEscapeHTML(false)
if err := enc.Encode(v); err != nil {
return nil, err
}
return bytes.TrimRight(b.Bytes(), "\n"), nil
}
+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": "go",
"system": "arch",
"from": "cmd/node-tools",
"binary": "node-tools"
}
]
}
}
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;
},
};
}
+66 -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 { takeToolEnvs, 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,68 @@ 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);
// Each module's environment, composed by the mesh (ADR 0192): taken before any bundle is imported.
const envs = takeToolEnvs();
const stop = await runTools({ broker, serves, moduleEntrypoints, credential: lastCredential, envs });
// 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 +303,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);
+427
View File
@@ -0,0 +1,427 @@
// 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 { fileURLToPath, pathToFileURL } from "node:url";
import { dirname, join, resolve } from "node:path";
import { existsSync, readFileSync } from "node:fs";
import { registerHooks } from "node:module";
import { useBroker } from "@novox/mesh-sdk/messaging";
import * as sdkTools from "@novox/mesh-sdk/tools";
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;
/** What each served module's bundles are given (novox/hq ADR 0192): module → words, composed by
* the mesh per machine. A module absent here is given the runtime's own words and nothing more. */
envs?: ReadonlyMap<string, Readonly<Record<string, string>>>;
}
/** The variable the mesh composes every served module's environment into, as JSON (ADR 0192). Read
* once at start and removed from the process's environment, so no bundle finds another's there. */
export const TOOL_ENV = "MESH_TOOL_ENV";
/** Read and remove the composed environments from an environment (the process's, by default). */
export function takeToolEnvs(env: NodeJS.ProcessEnv = process.env): Map<string, Record<string, string>> {
const raw = env[TOOL_ENV];
delete env[TOOL_ENV];
const out = new Map<string, Record<string, string>>();
if (!raw) return out;
let parsed: unknown;
try {
parsed = JSON.parse(raw);
} catch {
throw new Error(`${TOOL_ENV} is not JSON; the mesh composes it as {"<module>": {"<word>": "<value>"}}`);
}
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) {
throw new Error(`${TOOL_ENV} is not an object of modules`);
}
for (const [module, words] of Object.entries(parsed as Record<string, unknown>)) {
if (!words || typeof words !== "object" || Array.isArray(words)) {
throw new Error(`${TOOL_ENV}: ${module}'s environment is not an object of words`);
}
const own: Record<string, string> = {};
for (const [k, v] of Object.entries(words as Record<string, unknown>)) own[k] = String(v);
out.set(module, own);
}
return out;
}
/** 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);
}
// What each module's bundles are given: the runtime's own words, and over them the module's own.
const envs = opts.envs ?? new Map<string, Record<string, string>>();
const envFor = (module: string): NodeJS.ProcessEnv => ({ ...process.env, ...(envs.get(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)) {
// Told the module it serves it as, so a seat's verbs are the seat's (ADR 0193).
const child = await launch(module, path, { ...envFor(module), MESH_SERVED_MODULE: module });
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);
// Each registration's contributor is given its own module's environment and no other's (ADR 0192):
// the runtime's own words, and over them what the mesh composed for the module whose bundle made
// the registration. An SDK too old to ask per registration cannot do that; said, not hidden.
const each = (sdkTools as { collectToolsEach?: (f: (module: string, i: number) => NodeJS.ProcessEnv) => { module: string; tools: ToolDefinition[] }[] }).collectToolsEach;
if (!each && envs.size > 0) {
console.log(`[mesh-tools] this runtime's SDK cannot give each bundle its own environment; ${[...envs.keys()].join(", ")} serve with the runtime's words only (novox/hq ADR 0192)`);
}
const collected = each ? each((module, i) => envFor(owner[i] ?? self ?? module)) : collectTools();
const registrations = [
...collected.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;
const SDK = "@novox/mesh-sdk";
/**
* 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 the 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. A launched bundle (ADR 0188) is another process and is untouched.
*
* Installed once, in-thread, before the first bundle is imported; the hook sees every import after,
* `require` included. A bundle whose own copy is another version than the runtime's is said once,
* so a tool failing against the runtime's SDK points at the bundle rather than at the runtime.
*/
function oneSdk(): void {
if (sdkHooked) return;
registerHooks({
resolve(specifier, context, next) {
if (specifier === SDK || specifier.startsWith(SDK + "/")) {
if (context.parentURL) sayOtherSdk(context.parentURL);
return next(specifier, { ...context, parentURL: import.meta.url });
}
return next(specifier, context);
},
});
sdkHooked = true;
}
const sdkSaid = new Set<string>();
/** The version of the SDK copy nearest a file, by its package.json, or nothing when the file has none above it. */
function sdkVersionNear(fileURL: string): { dir: string; version: string } | undefined {
let dir = dirname(fileURLToPath(fileURL));
for (;;) {
const pkg = join(dir, "node_modules", SDK, "package.json");
if (existsSync(pkg)) {
try {
return { dir, version: String((JSON.parse(readFileSync(pkg, "utf8")) as { version?: string }).version ?? "?") };
} catch {
return { dir, version: "?" };
}
}
const up = dirname(dir);
if (up === dir) return undefined;
dir = up;
}
}
function sayOtherSdk(parentURL: string): void {
if (!parentURL.startsWith("file:")) return;
const own = sdkVersionNear(import.meta.url);
const theirs = sdkVersionNear(parentURL);
if (!theirs || !own || theirs.dir === own.dir || sdkSaid.has(theirs.dir)) return;
sdkSaid.add(theirs.dir);
if (theirs.version !== own.version) {
console.log(`[mesh-tools] ${theirs.dir} carries ${SDK} ${theirs.version}; this runtime's is ${own.version}, and the bundle speaks to the runtime's`);
}
}
@@ -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");
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for clash-tools.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./clash-tools.mjs");
await serveRegisteredOverStdio();
+4
View File
@@ -0,0 +1,4 @@
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("delta", (env) => [
{ name: "given", description: "what delta was given", input: {}, run: async () => ({ mine: env.DELTA_TOKEN_FILE ?? null, theirs: env.GAMMA_CONFIG_FILE ?? null }) },
]);
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for env-delta.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./env-delta.mjs");
await serveRegisteredOverStdio();
+5
View File
@@ -0,0 +1,5 @@
// A bundle that reads what it was given (novox/hq ADR 0192): its contributor's environment.
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("gamma", (env) => [
{ name: "given", description: "what gamma was given", input: {}, run: async () => ({ mine: env.GAMMA_CONFIG_FILE ?? null, theirs: env.DELTA_TOKEN_FILE ?? null, runtime: env.MESH_OPERATOR_ACCOUNT ?? null, composed: env.MESH_TOOL_ENV ?? null }) },
]);
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for env-gamma.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./env-gamma.mjs");
await serveRegisteredOverStdio();
+6
View File
@@ -0,0 +1,6 @@
#!/usr/bin/env node
// Launched (ADR 0188): its environment is the child's own.
import { serveStdio } from "@novox/mesh-sdk/stdio";
await serveStdio("zeta", [
{ name: "given", description: "what zeta was given", input: {}, run: async () => ({ mine: process.env.ZETA_URL ?? null, theirs: process.env.GAMMA_CONFIG_FILE ?? null, composed: process.env.MESH_TOOL_ENV ?? null }) },
]);
+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 }; } },
]);
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for many-alpha.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./many-alpha.mjs");
await serveRegisteredOverStdio();
+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 }) },
]);
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for many-beta.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./many-beta.mjs");
await serveRegisteredOverStdio();
+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");
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for many-broken.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./many-broken.mjs");
await serveRegisteredOverStdio();
+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" }) },
]);
+8
View File
@@ -0,0 +1,8 @@
#!/usr/bin/env node
// A launched bundle registering its seat before its own tools, served through the SDK's loop
// (novox/hq ADR 0193): which is the seat's and which the module's comes from the name it is given.
import { registerModuleTools } from "@novox/mesh-sdk/tools";
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
registerModuleTools("node-shelf", () => [{ name: "list", description: "the shelf", input: {}, run: async () => ({ shelf: ["x"] }) }]);
registerModuleTools("theta", (env) => [{ name: "own", description: "theta's", input: {}, run: async () => ({ theta: env.THETA_WORD ?? null }) }]);
await serveRegisteredOverStdio();
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for shop-seat.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./shop-seat.mjs");
await serveRegisteredOverStdio();
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for shop-tools.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./shop-tools.mjs");
await serveRegisteredOverStdio();
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for store-seat-unclaimed.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./store-seat-unclaimed.mjs");
await serveRegisteredOverStdio();
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for store-seat.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./store-seat.mjs");
await serveRegisteredOverStdio();
@@ -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"]);
+318
View File
@@ -0,0 +1,318 @@
/**
* 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, mkdirSync, 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, takeToolEnvs } 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();
let dir = "";
let stop = () => {};
const closing: Array<() => Promise<void>> = [];
const said: string[] = [];
const log = console.log;
try {
// 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,
// and a dependency of the bundle's own that the runtime does not carry.
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 });
mkdirSync(join(dir, "node_modules", "zeta-flavour"), { recursive: true });
writeFileSync(join(dir, "node_modules", "zeta-flavour", "package.json"), '{"name":"zeta-flavour","type":"module","main":"index.js"}\n');
writeFileSync(join(dir, "node_modules", "zeta-flavour", "index.js"), 'export const flavour = "the bundle\'s own";\n');
writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n');
writeFileSync(join(dir, "index.js"),
'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' +
'import { flavour } from "zeta-flavour";\n' +
'registerModuleTools("zeta", () => [{ name: "probe", description: "answers", input: {}, run: async () => ({ zeta: true, flavour }) }]);\n');
const mesh = await aMesh();
closing.push(() => mesh.close());
await mesh.issue(membershipOf("zeta", "anchor"));
const credential = { url, node: "anchor", module: "node-tools" };
const nodeTools = await connectNats(credential);
closing.push(() => nodeTools.close());
const asker = await connectNats({ url, module: "console", node: "workstation" });
closing.push(() => asker.close());
console.log = (...a: unknown[]) => said.push(a.join(" "));
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"));
// The SDK is the runtime's (the registration arrived); the bundle's other dependency is its own.
assert.deepEqual((await callTool(asker, "zeta.probe@anchor", {})).result, { zeta: true, flavour: "the bundle's own" });
} finally {
console.log = log;
stop();
for (const close of closing.reverse()) await close();
if (dir) rmSync(dir, { recursive: true, force: true });
resetTools();
}
});
test("each bundle is given its own environment and none of another's, imported or launched (ADR 0192)", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const mesh = await aMesh();
for (const m of ["gamma", "delta", "zeta"]) await mesh.issue(membershipOf(m, "anchor"));
const credential = { url, node: "anchor", module: "node-tools" };
const nodeTools = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
const log = console.log;
let stop = () => {};
const before = process.env.MESH_TOOL_ENV;
try {
process.env.MESH_OPERATOR_ACCOUNT = "somebody";
process.env.MESH_TOOL_ENV = JSON.stringify({
gamma: { GAMMA_CONFIG_FILE: "/var/lib/mesh/gamma/config.json" },
delta: { DELTA_TOKEN_FILE: "/var/lib/mesh/delta/token" },
zeta: { ZETA_URL: "http://127.0.0.1:3000" },
});
const envs = takeToolEnvs();
assert.equal(process.env.MESH_TOOL_ENV, undefined, "the composed environments were left in the process's");
console.log = () => {};
stop = await runTools({
broker: nodeTools, credential, envs,
serves: [
{ module: "gamma", entrypoints: [fixture("env-gamma.mjs")] },
{ module: "delta", entrypoints: [fixture("env-delta.mjs")] },
{ module: "zeta", entrypoints: [fixture("env-zeta.mjs")] },
],
});
console.log = log;
assert.deepEqual((await callTool(asker, "gamma.given@anchor", {})).result,
{ mine: "/var/lib/mesh/gamma/config.json", theirs: null, runtime: "somebody", composed: null });
assert.deepEqual((await callTool(asker, "delta.given@anchor", {})).result,
{ mine: "/var/lib/mesh/delta/token", theirs: null });
assert.deepEqual((await callTool(asker, "zeta.given@anchor", {})).result,
{ mine: "http://127.0.0.1:3000", theirs: null, composed: null });
} finally {
console.log = log;
if (before === undefined) delete process.env.MESH_TOOL_ENV; else process.env.MESH_TOOL_ENV = before;
delete process.env.MESH_OPERATOR_ACCOUNT;
stop();
await asker.close();
await nodeTools.close();
await mesh.close();
resetTools();
}
});
test("a launched bundle is told the module it serves, so its seat's verbs stay the seat's (ADR 0193)", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const mesh = await aMesh();
await mesh.issue(membershipOf("theta", "anchor", { seats: { "node-shelf": ["list"] } }));
const credential = { url, node: "anchor", module: "node-tools" };
const nodeTools = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
const log = console.log;
let stop = () => {};
try {
console.log = () => {};
stop = await runTools({ broker: nodeTools, credential, envs: new Map([["theta", { THETA_WORD: "given" }]]),
serves: [{ module: "theta", entrypoints: [fixture("served-seat-first.mjs")] }] });
console.log = log;
assert.deepEqual((await callTool(asker, "theta.own@anchor", {})).result, { theta: "given" });
assert.deepEqual((await callTool(asker, "seat:node-shelf.list@anchor", {})).result, { shelf: ["x"] });
} finally {
console.log = log;
stop();
await asker.close();
await nodeTools.close();
await mesh.close();
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());
}