commit 9eddcdb0a0976462b5f418ebec52be9b0e5e5a9c Author: jochen Date: Thu Sep 3 23:21:01 2026 +0200 mesh-tools — the tool runtime and AMQP broker binding The per-node process that makes a module's tools serve on the mesh: - broker-amqp.ts: a concrete AMQP implementation of the sdk's Broker contract (request/reply over a reply queue + correlation id, publish/ subscribe over a topic exchange). Kept here, not in the sdk, so a broker-client change never rebuilds a module (ADR 0044). - runtime.ts: bind the broker, import the assigned modules' tool entrypoints (each registers as it loads), serveTools. The thin wrapper. - main.ts: the entrypoint, configured by MESH_BROKER_URL + MESH_TOOL_MODULES. - Dockerfile: the container a node runs it as. Verified over a REAL broker: the test spins LavinMQ (the mesh's broker), the runtime serves a registered tool, a separate connection invokes it by name over AMQP and gets the result, and an unknown tool is refused over the wire. So 'modules can serve' is now running-on-the-mesh, not just proven in a mock. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..b947077 --- /dev/null +++ b/.gitignore @@ -0,0 +1,2 @@ +node_modules/ +dist/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..d867136 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,15 @@ +# The tool runtime, as the container a node runs. It is handed the broker URL and the assigned +# modules' tool entrypoints at deploy time (MESH_BROKER_URL, MESH_TOOL_MODULES) and serves them. +FROM node:22-alpine AS build +WORKDIR /app +COPY package.json ./ +RUN npm install --omit=dev --no-audit --no-fund +COPY dist ./dist + +FROM node:22-alpine +WORKDIR /app +COPY --from=build /app/node_modules ./node_modules +COPY --from=build /app/dist ./dist +COPY package.json ./ +USER node +ENTRYPOINT ["node", "dist/main.js"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..5f5406c --- /dev/null +++ b/README.md @@ -0,0 +1,32 @@ +# mesh-tools + +The Novox Mesh **tool runtime** — the per-node process that makes a module's tools actually serve. + +A module ships its tools (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk)); this +runtime is what loads them and puts them on the mesh. It: + +1. connects the mesh broker (novox/hq ADR 0001) — a concrete AMQP implementation of the sdk's + `Broker` contract; +2. imports the assigned modules' compiled tool entrypoints, each of which registers its tools as it + loads; +3. serves them through the sdk's `serveTools` harness, answering `tools.invoke` over the broker. + +Everything hard — dispatch, collection, duplicate-name safety — is the sdk's. This is the thin +wrapper that binds the broker and loads the modules. Keeping the AMQP client here, out of the sdk, +is deliberate: a broker-client change never rebuilds a module (ADR 0044). + +## Running it + +``` +MESH_BROKER_URL amqp://… the mesh broker +MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool entrypoints +``` + +`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables +and starts it like any other supervised workload. + +## Verified + +`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the +runtime serves a registered tool, a separate connection invokes it by name and gets the result, and +an unknown tool is refused over the wire. diff --git a/package-lock.json b/package-lock.json new file mode 100644 index 0000000..2d08849 --- /dev/null +++ b/package-lock.json @@ -0,0 +1,118 @@ +{ + "name": "@novox/mesh-tools", + "version": "0.1.0", + "lockfileVersion": 3, + "requires": true, + "packages": { + "": { + "name": "@novox/mesh-tools", + "version": "0.1.0", + "dependencies": { + "@novox/mesh-sdk": "file:../mesh-sdk", + "amqplib": "^0.10.9" + }, + "bin": { + "mesh-tools": "dist/main.js" + }, + "devDependencies": { + "@types/amqplib": "^0.10.8", + "@types/node": "^22.20.1", + "typescript": "^5.9.3" + } + }, + "../mesh-sdk": { + "name": "@novox/mesh-sdk", + "version": "0.1.0", + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } + }, + "node_modules/@novox/mesh-sdk": { + "resolved": "../mesh-sdk", + "link": true + }, + "node_modules/@types/amqplib": { + "version": "0.10.8", + "resolved": "https://registry.npmjs.org/@types/amqplib/-/amqplib-0.10.8.tgz", + "integrity": "sha512-vtDp8Pk1wsE/AuQ8/Rgtm6KUZYqcnTgNvEHwzCkX8rL7AGsC6zqAfKAAJhUZXFhM/Pp++tbnUHiam/8vVpPztA==", + "dev": true, + "license": "MIT", + "dependencies": { + "@types/node": "*" + } + }, + "node_modules/@types/node": { + "version": "22.20.1", + "resolved": "https://registry.npmjs.org/@types/node/-/node-22.20.1.tgz", + "integrity": "sha512-EANqOCF9QFyra+4pfxUcX9STKJpCLjMbObVzljIJomAWSnuSIEAvyzEU53GaajbXJEgdh0iEcPL+DGvpUd4k1Q==", + "dev": true, + "license": "MIT", + "dependencies": { + "undici-types": "~6.21.0" + } + }, + "node_modules/amqplib": { + "version": "0.10.9", + "resolved": "https://registry.npmjs.org/amqplib/-/amqplib-0.10.9.tgz", + "integrity": "sha512-jwSftI4QjS3mizvnSnOrPGYiUnm1vI2OP1iXeOUz5pb74Ua0nbf6nPyyTzuiCLEE3fMpaJORXh2K/TQ08H5xGA==", + "license": "MIT", + "dependencies": { + "buffer-more-ints": "~1.0.0", + "url-parse": "~1.5.10" + }, + "engines": { + "node": ">=10" + } + }, + "node_modules/buffer-more-ints": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/buffer-more-ints/-/buffer-more-ints-1.0.0.tgz", + "integrity": "sha512-EMetuGFz5SLsT0QTnXzINh4Ksr+oo4i+UGTXEshiGCQWnsgSs7ZhJ8fzlwQ+OzEMs0MpDAMr1hxnblp5a4vcHg==", + "license": "MIT" + }, + "node_modules/querystringify": { + "version": "2.2.0", + "resolved": "https://registry.npmjs.org/querystringify/-/querystringify-2.2.0.tgz", + "integrity": "sha512-FIqgj2EUvTa7R50u0rGsyTftzjYmv/a3hO345bZNrqabNqjtgiDMgmo4mkUjd+nzU5oF3dClKqFIPUKybUyqoQ==", + "license": "MIT" + }, + "node_modules/requires-port": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/requires-port/-/requires-port-1.0.0.tgz", + "integrity": "sha512-KigOCHcocU3XODJxsu8i/j8T9tzT4adHiecwORRQ0ZZFcp7ahwXuRU1m+yuO90C5ZUyGeGfocHDI14M3L3yDAQ==", + "license": "MIT" + }, + "node_modules/typescript": { + "version": "5.9.3", + "resolved": "https://registry.npmjs.org/typescript/-/typescript-5.9.3.tgz", + "integrity": "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw==", + "dev": true, + "license": "Apache-2.0", + "bin": { + "tsc": "bin/tsc", + "tsserver": "bin/tsserver" + }, + "engines": { + "node": ">=14.17" + } + }, + "node_modules/undici-types": { + "version": "6.21.0", + "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", + "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", + "dev": true, + "license": "MIT" + }, + "node_modules/url-parse": { + "version": "1.5.10", + "resolved": "https://registry.npmjs.org/url-parse/-/url-parse-1.5.10.tgz", + "integrity": "sha512-WypcfiRhfeUP9vvF0j6rw0J3hrWrw6iZv3+22h6iRMJ/8z1Tj6XfLP4DsUix5MhMPnXpiHDoKyoZ/bdCkwBCiQ==", + "license": "MIT", + "dependencies": { + "querystringify": "^2.1.1", + "requires-port": "^1.0.0" + } + } + } +} diff --git a/package.json b/package.json new file mode 100644 index 0000000..fb4377b --- /dev/null +++ b/package.json @@ -0,0 +1,22 @@ +{ + "name": "@novox/mesh-tools", + "version": "0.1.0", + "description": "The Novox Mesh tool runtime — binds the mesh broker and serves the assigned modules' tools.", + "type": "module", + "bin": { + "mesh-tools": "./dist/main.js" + }, + "scripts": { + "build": "tsc", + "test": "node --test --experimental-strip-types 'test/*.test.ts'" + }, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0", + "amqplib": "^0.10.9" + }, + "devDependencies": { + "@types/amqplib": "^0.10.8", + "@types/node": "^22.20.1", + "typescript": "^5.9.3" + } +} diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts new file mode 100644 index 0000000..930b8ac --- /dev/null +++ b/src/broker-amqp.ts @@ -0,0 +1,103 @@ +// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker +// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime +// provides the binding — so a broker-client change never rebuilds the modules. + +import amqp from "amqplib"; +import { randomUUID } from "node:crypto"; +import type { Broker, Envelope } from "@novox/mesh-sdk/messaging"; + +const EXCHANGE = "mesh.tools"; // one topic exchange carries tool invocations and events + +interface Reply { + result?: unknown; + error?: string; +} + +/** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */ +export async function connectAmqp(url: string): Promise { + const conn = await amqp.connect(url); + const ch = await conn.createChannel(); + await ch.assertExchange(EXCHANGE, "topic", { durable: true }); + + // Request/reply: one exclusive reply queue, correlationId → resolver. + const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true }); + const pending = new Map void>(); + await ch.consume( + replyQueue, + (msg) => { + if (!msg) return; + const resolve = pending.get(msg.properties.correlationId); + if (resolve) { + pending.delete(msg.properties.correlationId); + resolve(JSON.parse(msg.content.toString()) as Reply); + } + }, + { noAck: true }, + ); + + return { + async request(key: string, body: Req): Promise { + const id = randomUUID(); + const answered = new Promise((resolve, reject) => { + const timer = setTimeout(() => { + if (pending.delete(id)) reject(new Error(`request ${key} timed out`)); + }, 30_000); + pending.set(id, (r) => { + clearTimeout(timer); + if (r.error) reject(new Error(r.error)); + else resolve(r.result as Res); + }); + }); + ch.publish(EXCHANGE, key, Buffer.from(JSON.stringify(body)), { + correlationId: id, + replyTo: replyQueue, + }); + return answered; + }, + + async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { + const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true }); + await ch.bindQueue(queue, EXCHANGE, key); + const consumer = await ch.consume(queue, (msg) => { + if (!msg) return; + void (async () => { + let reply: Reply; + try { + reply = { result: await handler(JSON.parse(msg.content.toString()) as Req) }; + } catch (err) { + reply = { error: err instanceof Error ? err.message : String(err) }; + } + if (msg.properties.replyTo) { + ch.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { + correlationId: msg.properties.correlationId, + }); + } + ch.ack(msg); + })(); + }); + return () => void ch.cancel(consumer.consumerTag); + }, + + async publish(env: Envelope): Promise { + ch.publish(EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body))); + }, + + async subscribe(pattern: string, handler: (env: Envelope) => Promise): Promise<() => void> { + const { queue } = await ch.assertQueue("", { exclusive: true }); + await ch.bindQueue(queue, EXCHANGE, pattern); + const consumer = await ch.consume(queue, (msg) => { + if (!msg) return; + void (async () => { + await handler({ key: msg.fields.routingKey, node: "", body: JSON.parse(msg.content.toString()) as T }); + ch.ack(msg); + })(); + }); + return () => void ch.cancel(consumer.consumerTag); + }, + + async close(): Promise { + await ch.close(); + await conn.close(); + }, + }; +} diff --git a/src/main.ts b/src/main.ts new file mode 100644 index 0000000..4c17be6 --- /dev/null +++ b/src/main.ts @@ -0,0 +1,38 @@ +// The runnable entrypoint. Reads its configuration from the environment the host resolved for it, +// connects the mesh broker, and serves the assigned modules' tools until stopped. +// +// MESH_BROKER_URL amqp://… the mesh broker +// MESH_TOOL_MODULES /path/a,/path/b,… compiled tool entrypoints of the assigned modules + +import { connectAmqp } from "./broker-amqp.js"; +import { runTools } from "./runtime.js"; + +async function main(): Promise { + const url = requireEnv("MESH_BROKER_URL"); + const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "") + .split(",") + .map((s) => s.trim()) + .filter(Boolean); + + const broker = await connectAmqp(url); + const stop = await runTools({ broker, moduleEntrypoints }); + + const shutdown = async (): Promise => { + stop(); + await broker.close(); + process.exit(0); + }; + process.on("SIGTERM", () => void shutdown()); + process.on("SIGINT", () => void shutdown()); +} + +function requireEnv(name: string): string { + const v = process.env[name]; + if (!v) { + console.error(`mesh-tools: ${name} is not set — the runtime cannot serve without it`); + process.exit(1); + } + return v; +} + +void main(); diff --git a/src/runtime.ts b/src/runtime.ts new file mode 100644 index 0000000..d43b33e --- /dev/null +++ b/src/runtime.ts @@ -0,0 +1,32 @@ +// 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 { serveTools, listTools } from "@novox/mesh-sdk/tools"; +import type { Broker } from "@novox/mesh-sdk/messaging"; + +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[]; +} + +/** 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); + } + + const stop = await serveTools(opts.broker); + const tools = listTools(); + console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`); + return stop; +} diff --git a/test/amqp.test.ts b/test/amqp.test.ts new file mode 100644 index 0000000..205eafb --- /dev/null +++ b/test/amqp.test.ts @@ -0,0 +1,44 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { registerModuleTools, resetTools } from "@novox/mesh-sdk/tools"; +import { connectAmqp } from "../dist/broker-amqp.js"; +import { runTools } from "../dist/runtime.js"; + +// Requires a real broker at $MESH_BROKER_URL. The test harness spins LavinMQ (the mesh's broker) +// and points this at it; skipped if it is not set, never failed for the environment. +const url = process.env.MESH_BROKER_URL; + +test("a module's tool serves and is invoked over a real AMQP broker", { skip: !url }, async () => { + resetTools(); + + // A module registers a real tool, exactly as umami does. + registerModuleTools("demo", () => [ + { + name: "greet", + description: "return a greeting", + input: { who: { type: "string" } }, + run: async (args) => ({ hello: String(args.who), from: "the tool runtime" }), + }, + ]); + + // The runtime binds the broker and serves the module (no module entrypoints to import here — the + // tool is registered in-process — but this is the exact runtime that serves on a node). + const serverBroker = await connectAmqp(url!); + const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] }); + + // A separate connection — a caller, like mesh-control's command API — invokes over the broker. + const caller = await connectAmqp(url!); + const result = await caller.request<{ tool: string; args: Record }, { hello: string }>( + "tools.invoke", + { tool: "greet", args: { who: "mesh" } }, + ); + assert.equal(result.hello, "mesh"); + + // An unknown tool is refused over the wire, not silently dropped. + await assert.rejects(caller.request("tools.invoke", { tool: "nope", args: {} }), /no such tool/); + + stop(); + await caller.close(); + await serverBroker.close(); +}); diff --git a/tsconfig.json b/tsconfig.json new file mode 100644 index 0000000..2606b43 --- /dev/null +++ b/tsconfig.json @@ -0,0 +1,15 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "declaration": true, + "outDir": "dist", + "rootDir": "src", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "forceConsistentCasingInFileNames": true + }, + "include": ["src/**/*.ts"] +}