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.
This commit is contained in:
@@ -0,0 +1,188 @@
|
||||
/**
|
||||
* A person's client, against a real bus.
|
||||
*
|
||||
* What is worth checking is not that a request/reply works — the runtime's own tests cover that — but
|
||||
* that the two surfaces are the same thing. An agent and a person must see the same tools and get the
|
||||
* same answers, or the MCP surface becomes a second definition of what a tool is.
|
||||
*
|
||||
* 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/client.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
|
||||
// The built output, not the source: the client imports its siblings as `.js`, which is what ships and
|
||||
// what every other file here does, and cannot be loaded as TypeScript directly. `pretest` builds.
|
||||
import { connectNats } from "../dist/broker-nats.js";
|
||||
import { callTool, seatsIn, toolKey, toolsOn, whyItFailed } from "../dist/client.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
|
||||
/** A mesh as discovery sees it (design 34 §3): the catalogue holding three modules, two of them up
|
||||
* and answering `tools`, one held and not running. Two connections, because a person and a module
|
||||
* are different users even in a test. */
|
||||
async function aMeshWithTools() {
|
||||
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||
const shop = await connectNats({ url: url!, module: "shop" });
|
||||
await catalogue.handle("catalog_modules", async () => ({
|
||||
modules: [{ module: "shop" }, { module: "mesh-catalog" }, { module: "ghost" }],
|
||||
}));
|
||||
await catalogue.handle("tools", async () => ({
|
||||
module: "mesh-catalog",
|
||||
tools: [{ name: "catalog_modules", description: "what modules the mesh has", input: {} }],
|
||||
}));
|
||||
await shop.handle("tools", async () => ({
|
||||
module: "shop",
|
||||
tools: [{ name: "price", description: "what something costs", input: { type: "object" } }],
|
||||
}));
|
||||
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||
// And the mesh's own records, served by the holder of the mesh-controller seat (ADR 0154): one
|
||||
// role's tool, and the verb that lists every role's.
|
||||
const controller = await connectNats({ url: url!, module: "mesh-controller" });
|
||||
await controller.handle("seat:mesh-controller.tools", async () => ({
|
||||
seats: [
|
||||
{ seat: "mesh-controller", scope: "mesh", tools: [{ name: "status", description: "what is wrong", input: {} }] },
|
||||
{ seat: "node-dns-resolver", scope: "node", tools: [{ name: "lookup", description: "one machine's", input: {} }] },
|
||||
],
|
||||
}));
|
||||
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
|
||||
return {
|
||||
async close() {
|
||||
await catalogue.close();
|
||||
await shop.close();
|
||||
await controller.close();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
test("a person sees what the running modules answer, sorted, and who did not answer", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const mesh = await aMeshWithTools();
|
||||
const person = await connectNats({ url, module: "person.ada" });
|
||||
try {
|
||||
const began = Date.now();
|
||||
const have = await toolsOn(person);
|
||||
assert.deepEqual(
|
||||
have.tools.map((x) => `${x.module}.${x.name}`),
|
||||
["mesh-catalog.catalog_modules", "mesh-controller.status", "node-dns-resolver.lookup", "shop.price"],
|
||||
"the list is what the modules answered plus every role's tools, in a stable order",
|
||||
);
|
||||
// A role's tool is marked as one; a node-scoped seat's carries its scope, so a caller names
|
||||
// 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);
|
||||
const lookup = have.tools.find((x) => x.module === "node-dns-resolver")!;
|
||||
assert.ok(lookup.seat && lookup.scope === "node");
|
||||
assert.equal(toolKey("node-dns-resolver.lookup", seatsIn(have)), "seat:node-dns-resolver.lookup");
|
||||
// Silence is named, never dropped: a module the catalogue holds and nothing answered for.
|
||||
assert.deepEqual(have.notAnswering, ["ghost"]);
|
||||
// And at once: a module that is not running costs nothing, or the list is unusable.
|
||||
assert.ok(Date.now() - began < 5_000, "an absent module waited out the timeout");
|
||||
} finally {
|
||||
await person.close();
|
||||
await mesh.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a person calls a tool and gets the module's own answer, unshaped", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const mesh = await aMeshWithTools();
|
||||
const person = await connectNats({ url, module: "person.ada" });
|
||||
try {
|
||||
const answer = await callTool(person, "shop.price", { of: "a hat" });
|
||||
assert.deepEqual(answer.result, { of: "a hat", cost: 12 });
|
||||
} finally {
|
||||
await person.close();
|
||||
await mesh.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a tool nobody serves says so at once, and says what to do about it", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const person = await connectNats({ url: url!, module: "person.ada" });
|
||||
try {
|
||||
const began = Date.now();
|
||||
await assert.rejects(() => callTool(person, "ghost.missing", {}));
|
||||
// At once, not after the whole wait: "that module is down" and "that tool is slow" need
|
||||
// different things done, and a timeout cannot tell them apart.
|
||||
assert.ok(Date.now() - began < 5_000, "a tool nobody serves waited out the timeout");
|
||||
} finally {
|
||||
await person.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a name that is not <module>.<tool> is refused before anything is sent", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const person = await connectNats({ url: url!, module: "person.ada" });
|
||||
try {
|
||||
await assert.rejects(() => callTool(person, "price", {}), /does not name a tool/);
|
||||
} finally {
|
||||
await person.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("each way a call fails says what to do about it", () => {
|
||||
// The three answers a person actually gets. Without this they are one timeout and a stack trace,
|
||||
// and the remedies are in three different places.
|
||||
assert.match(whyItFailed("shop.price", new Error("no responders")), /nothing serves shop\.price/);
|
||||
assert.match(
|
||||
whyItFailed("shop.price", new Error("Permissions Violation for Publish")),
|
||||
/may not call shop\.price/,
|
||||
);
|
||||
assert.match(whyItFailed("shop.price", new Error("timeout")), /did not answer in time/);
|
||||
assert.match(whyItFailed("shop.price", new Error("something else")), /something else/);
|
||||
});
|
||||
|
||||
test("a role's tool is reached through the seat, and a module's own name is never shadowed", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const { seatsIn, toolKey } = await import("../dist/client.js");
|
||||
const mesh = await aMeshWithTools();
|
||||
const person = await connectNats({ url, module: "person.ada" });
|
||||
try {
|
||||
const roles = seatsIn(await toolsOn(person));
|
||||
assert.equal(toolKey("mesh-controller.status", roles), "seat:mesh-controller.status");
|
||||
assert.equal(toolKey("mesh-controller.other", roles), "mesh-controller.other", "a verb the seat does not declare is a module's");
|
||||
assert.equal(toolKey("shop.price", roles), "shop.price");
|
||||
const answer = await callTool(person, "mesh-controller.status", {}, roles);
|
||||
assert.deepEqual(answer.result, { output: "all quiet", ok: true });
|
||||
const direct = await callTool(person, "seat:mesh-controller.status", {});
|
||||
assert.deepEqual(direct.result, { output: "all quiet", ok: true });
|
||||
} finally {
|
||||
await person.close();
|
||||
await mesh.close();
|
||||
}
|
||||
});
|
||||
|
||||
// A module on two machines (novox/hq ADR 0159): a call names the machine and reaches that instance
|
||||
// and no other; a call that names none reaches one of them and says which; a seat's verb is served
|
||||
// by the claimant's tool of the same name on the seat's own subject.
|
||||
test("a tool call names the machine it is for, and every answer says which machine answered", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const onAnchor = await connectNats({ url: url!, module: "store", node: "anchor" });
|
||||
const onHome = await connectNats({ url: url!, module: "store", node: "home-server" });
|
||||
const asker = await connectNats({ url: url!, module: "console", node: "workstation" });
|
||||
try {
|
||||
await onAnchor.handle("databases", async () => ({ at: "anchor" }));
|
||||
await onHome.handle("databases", async () => ({ at: "home-server" }));
|
||||
|
||||
const home = await callTool(asker, "store.databases@home-server", {});
|
||||
assert.deepEqual(home.result, { at: "home-server" });
|
||||
assert.equal(home.node, "home-server");
|
||||
const anchor = await callTool(asker, "store.databases@anchor", {});
|
||||
assert.deepEqual(anchor.result, { at: "anchor" });
|
||||
assert.equal(anchor.node, "anchor");
|
||||
|
||||
const whichever = await callTool(asker, "store.databases", {});
|
||||
assert.ok(["anchor", "home-server"].includes(whichever.node ?? ""), `an unnamed call still says who answered: ${whichever.node}`);
|
||||
assert.deepEqual(whichever.result, { at: whichever.node });
|
||||
|
||||
// A seat's verb, served by the claimant on the seat's subject; a mesh seat's is flat.
|
||||
await onAnchor.handleSubject("mesh.seat.mesh-store.tool.databases", async () => ({ seat: "mesh-store", at: "anchor" }));
|
||||
const viaSeat = await callTool(asker, "seat:mesh-store.databases", {});
|
||||
assert.deepEqual(viaSeat.result, { seat: "mesh-store", at: "anchor" });
|
||||
assert.equal(viaSeat.node, "anchor");
|
||||
} finally {
|
||||
await onAnchor.close();
|
||||
await onHome.close();
|
||||
await asker.close();
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,54 @@
|
||||
// The TypeScript implementation, held to the shared fixtures (novox/hq ADR 0074, design 19).
|
||||
//
|
||||
// Run against a NATS server, because the question is what actually reaches the wire:
|
||||
//
|
||||
// docker run -d --rm --name c -p 14222:4222 nats:2.10-alpine -js
|
||||
// node test/conformance.mjs
|
||||
//
|
||||
// **The runner lives with the implementation it exercises; the fixture does not.** It is read
|
||||
// from the sdk's conformance directory by sibling path — the same file the Go suite reads. A
|
||||
// fixture copied into each implementation is two fixtures, and two fixtures drift, which is the
|
||||
// failure the suite exists to prevent.
|
||||
import { readFileSync } from "node:fs";
|
||||
import { connect } from "nats";
|
||||
|
||||
const clientPath = process.argv[2] ?? "../dist/broker-nats.js";
|
||||
const { connectNats } = await import(clientPath);
|
||||
const f = JSON.parse(readFileSync(new URL("../../mesh-sdk/conformance/events/module-event.json", import.meta.url)));
|
||||
const URL_ = process.env.MESH_TEST_NATS ?? "nats://127.0.0.1:14222";
|
||||
|
||||
let failed = 0;
|
||||
const check = (ok, what) => { console.log(` ${ok ? "ok " : "FAIL"} ${what}`); if (!ok) failed++; };
|
||||
|
||||
const admin = await connect({ servers: URL_ });
|
||||
const jsm = await admin.jetstreamManager();
|
||||
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>", "mesh.seat.*.event.>"] });
|
||||
|
||||
const shop = await connectNats({ url: URL_, node: f.given.node, module: f.given.module });
|
||||
await shop.publish({ key: f.given.key, node: f.given.node, body: f.given.body, headers: f.given.headers });
|
||||
|
||||
// What actually landed, read back from the stream rather than from the client that wrote it.
|
||||
const msg = await jsm.streams.getMessage("EVENTS", { last_by_subj: f.wire.subject });
|
||||
check(!!msg, `it lands on ${f.wire.subject}`);
|
||||
|
||||
if (msg) {
|
||||
const got = {};
|
||||
if (msg.header) for (const k of msg.header.keys()) got[k] = msg.header.get(k);
|
||||
for (const h of f.wire.requiredHeaders) {
|
||||
check(got[h] !== undefined && got[h] !== "", `${h} is set`);
|
||||
}
|
||||
check(got["content-type"] === f.wire.headerFormats["content-type"], "content-type is as pinned");
|
||||
check(new RegExp(f.wire.headerFormats["x-event-id"]).test(got["x-event-id"]), "x-event-id is as pinned");
|
||||
check(!Number.isNaN(Date.parse(got["x-time"])), "x-time parses as a date");
|
||||
check(got["x-source"] === f.given.module, "x-source agrees with the subject's module");
|
||||
|
||||
const payload = JSON.parse(new TextDecoder().decode(msg.data));
|
||||
check(payload.envelope === undefined && payload.key === undefined,
|
||||
"the payload is the body alone, not the envelope (the fixture refuses nesting)");
|
||||
check(JSON.stringify(payload) === JSON.stringify(f.given.body),
|
||||
"the body round-trips — semantic, not byte-exact, per the README");
|
||||
}
|
||||
|
||||
await shop.close(); await admin.close();
|
||||
console.log(failed ? `\n${failed} failed` : "\nall passed");
|
||||
process.exit(failed ? 1 : 0);
|
||||
+6
@@ -0,0 +1,6 @@
|
||||
// A module naming a tool after the verb the runtime answers for every module — refused at load.
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("clash", () => [
|
||||
{ name: "tools", description: "mine, not the runtime's", input: {}, run: async () => ({}) },
|
||||
]);
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
// One of several bundles the node's runtime loads (novox/hq ADR 0175): a module with two tools.
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
import { emit } from "@novox/mesh-sdk/events";
|
||||
|
||||
registerModuleTools("alpha", () => [
|
||||
{ name: "one", description: "alpha's first", input: {}, run: async () => ({ alpha: 1 }) },
|
||||
// A tool that emits: the event must land on alpha's subject, not the runtime's.
|
||||
{ name: "two", description: "alpha's second, which emits", input: {}, run: async () => { await emit("happened", { by: "alpha" }); return { alpha: 2 }; } },
|
||||
]);
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
// One of several bundles the node's runtime loads (novox/hq ADR 0175): a module with three tools of
|
||||
// its own and the implementation of a node seat's two verbs under the seat's name (ADR 0159).
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("beta", () => [
|
||||
{ name: "three", description: "beta's", input: {}, run: async () => ({ beta: 3 }) },
|
||||
{ name: "four", description: "beta's", input: {}, run: async () => ({ beta: 4 }) },
|
||||
{ name: "five", description: "beta's", input: {}, run: async () => ({ beta: 5 }) },
|
||||
]);
|
||||
|
||||
registerModuleTools("node-shelf", () => [
|
||||
{ name: "list", description: "what is on the shelf", input: {}, run: async () => ({ shelf: ["a", "b"] }) },
|
||||
{ name: "clear", description: "take it all off", input: {}, run: async () => ({ cleared: true }) },
|
||||
]);
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
// A bundle that throws on import — the fault ADR 0175 names as what got harder: one module's bundle
|
||||
// must not take the node's other tools down.
|
||||
throw new Error("gamma's bundle cannot find its client");
|
||||
+29
@@ -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
@@ -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" }) },
|
||||
]);
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
// The shop again, in a file of its own: a module imported once stays imported, so a second test needs a second entrypoint.
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("shop", () => [
|
||||
{
|
||||
name: "price",
|
||||
description: "what something costs",
|
||||
input: { of: { type: "string", description: "the thing" } },
|
||||
run: async (args) => ({ of: args.of ?? "nothing", cost: 12 }),
|
||||
},
|
||||
{ name: "refund", description: "give it back", input: {}, run: async () => ({ done: true }) },
|
||||
]);
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
// A module's tool entrypoint, as the runtime imports one: registers and returns.
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("shop", () => [
|
||||
{
|
||||
name: "price",
|
||||
description: "what something costs",
|
||||
input: { of: { type: "string", description: "the thing" } },
|
||||
run: async (args) => ({ of: args.of ?? "nothing", cost: 12 }),
|
||||
},
|
||||
{ name: "refund", description: "give it back", input: {}, run: async () => ({ done: true }) },
|
||||
]);
|
||||
@@ -0,0 +1,12 @@
|
||||
// A module that is both software and a role: postgres's own tools under its name, and its
|
||||
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("postgres", () => [
|
||||
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
|
||||
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
|
||||
]);
|
||||
|
||||
registerModuleTools("mesh-store", () => [
|
||||
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
|
||||
]);
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
// A module that is both software and a role: postgres's own tools under its name, and its
|
||||
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("postgres", () => [
|
||||
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
|
||||
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
|
||||
]);
|
||||
|
||||
registerModuleTools("mesh-store", () => [
|
||||
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
|
||||
]);
|
||||
@@ -0,0 +1,113 @@
|
||||
/**
|
||||
* The console: the same surface over HTTP on loopback, started the way the mesh starts it — on the
|
||||
* module credential in MESH_BROKER_FILE (novox/hq ADR 0152, design 34 §2).
|
||||
*
|
||||
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/http.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { mkdtemp, writeFile } from "node:fs/promises";
|
||||
import { join } from "node:path";
|
||||
|
||||
import { connectNats } from "../dist/broker-nats.js";
|
||||
import { serveMcpHttp } from "../dist/http.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
|
||||
async function aMesh(t: { after: (fn: () => Promise<void> | void) => void }) {
|
||||
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||
const shop = await connectNats({ url: url!, module: "shop" });
|
||||
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "shop" }] }));
|
||||
await shop.handle("tools", async () => ({
|
||||
module: "shop",
|
||||
tools: [{ name: "price", description: "what something costs", input: {} }],
|
||||
}));
|
||||
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||
t.after(async () => {
|
||||
await catalogue.close();
|
||||
await shop.close();
|
||||
});
|
||||
}
|
||||
|
||||
/** `mesh serve` as the mesh runs it: MESH_BROKER_FILE, a listen address, nothing else. */
|
||||
async function aConsole(t: { after: (fn: () => Promise<void> | void) => void }): Promise<string> {
|
||||
const dir = await mkdtemp("/tmp/mesh-console-");
|
||||
const credential = join(dir, "broker");
|
||||
await writeFile(credential, JSON.stringify({ url, node: "desk", module: "mesh-console", user: "desk.mesh-console", password: "x" }));
|
||||
const child: ChildProcess = spawn(process.execPath, ["dist/mesh.js", "serve", "--listen", "127.0.0.1:0"], {
|
||||
env: { ...process.env, MESH_BROKER_FILE: credential, MESH_CREDENTIAL: "" },
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
});
|
||||
t.after(() => {
|
||||
child.kill("SIGTERM");
|
||||
});
|
||||
return new Promise((resolve, reject) => {
|
||||
let out = "";
|
||||
let err = "";
|
||||
child.stdout!.on("data", (d) => {
|
||||
out += d.toString();
|
||||
const m = /listening on (http:\/\/[^/]+\/mcp)/.exec(out);
|
||||
if (m) resolve(m[1]!);
|
||||
});
|
||||
child.stderr!.on("data", (d) => (err += d.toString()));
|
||||
child.on("exit", (code) => reject(new Error(`serve exited ${code}: ${err}`)));
|
||||
});
|
||||
}
|
||||
|
||||
async function post(endpoint: string, body: unknown): Promise<{ status: number; json?: any }> {
|
||||
const res = await fetch(endpoint, {
|
||||
method: "POST",
|
||||
headers: { "content-type": "application/json", accept: "application/json" },
|
||||
body: JSON.stringify(body),
|
||||
});
|
||||
const text = await res.text();
|
||||
return { status: res.status, json: text ? JSON.parse(text) : undefined };
|
||||
}
|
||||
|
||||
test("the console answers a host on loopback, as the account the mesh gave it", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
await aMesh(t);
|
||||
const endpoint = await aConsole(t);
|
||||
|
||||
const hello = await post(endpoint, { jsonrpc: "2.0", id: 1, method: "initialize", params: {} });
|
||||
assert.equal(hello.status, 200);
|
||||
assert.match(hello.json.result.instructions, /desk\.mesh-console/, "the handshake names the console's account");
|
||||
|
||||
const heard = await post(endpoint, { jsonrpc: "2.0", method: "notifications/initialized" });
|
||||
assert.equal(heard.status, 202, "a notification is heard and not answered");
|
||||
|
||||
const listed = await post(endpoint, { jsonrpc: "2.0", id: 2, method: "tools/list" });
|
||||
assert.deepEqual(listed.json.result.tools.map((x: { name: string }) => x.name), ["shop.price"]);
|
||||
|
||||
const called = await post(endpoint, {
|
||||
jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "shop.price", arguments: { of: "a hat" } },
|
||||
});
|
||||
assert.deepEqual(JSON.parse(called.json.result.content[0].text), { of: "a hat", cost: 12 });
|
||||
|
||||
// A person's client through the same endpoint, with no credential of its own.
|
||||
const { main } = await import("../dist/mesh.js");
|
||||
const logged: string[] = [];
|
||||
const was = console.log;
|
||||
console.log = (line: string) => logged.push(String(line));
|
||||
try {
|
||||
assert.equal(await main(["tools", "--console", endpoint]), 0);
|
||||
} finally {
|
||||
console.log = was;
|
||||
}
|
||||
assert.ok(logged.some((l) => l.startsWith("shop.price")), `the client did not list through the console: ${logged}`);
|
||||
});
|
||||
|
||||
test("the console binds loopback and nowhere else", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const bus = await connectNats({ url, module: "mesh-console", node: "desk" });
|
||||
try {
|
||||
await assert.rejects(() => serveMcpHttp(bus, "desk.mesh-console", "0.0.0.0:0"), /loopback and nowhere else/);
|
||||
const up = await serveMcpHttp(bus, "desk.mesh-console", "127.0.0.1:0");
|
||||
assert.match(up.address, /^127\.0\.0\.1:\d+$/);
|
||||
await up.close();
|
||||
} finally {
|
||||
await bus.close();
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,196 @@
|
||||
/**
|
||||
* The MCP surface, driven the way a host drives it.
|
||||
*
|
||||
* **The claim worth checking is that it is the same thing the command line is.** An agent and a
|
||||
* person must see the same tools and get the same answers, or this becomes a second definition of what
|
||||
* a tool is — which is exactly what a thin adapter is supposed to avoid.
|
||||
*
|
||||
* 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/mcp.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { spawn } from "node:child_process";
|
||||
|
||||
import { connectNats } from "../dist/broker-nats.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
|
||||
/** A module answering the catalogue's list and one tool, plus a credential file the client reads. */
|
||||
async function aMeshAndACredential(t: { after: (fn: () => Promise<void> | void) => void }) {
|
||||
const catalogue = await connectNats({ url: url!, module: "mesh-catalog" });
|
||||
const shop = await connectNats({ url: url!, module: "shop" });
|
||||
await catalogue.handle("catalog_modules", async () => ({
|
||||
modules: [{ module: "shop" }, { module: "ghost" }],
|
||||
}));
|
||||
await shop.handle("tools", async () => ({
|
||||
module: "shop",
|
||||
tools: [{ name: "price", description: "what something costs", input: { of: { type: "string" } } }],
|
||||
}));
|
||||
await shop.handle("price", async (body: { of?: string }) => ({ of: body.of ?? "nothing", cost: 12 }));
|
||||
const controller = await connectNats({ url: url!, module: "mesh-controller" });
|
||||
await controller.handle("seat:mesh-controller.tools", async () => ({
|
||||
seats: [
|
||||
{ seat: "mesh-controller", scope: "mesh", tools: [
|
||||
{ name: "status", description: "what is wrong", input: {} },
|
||||
{ name: "push", description: "tell a machine", input: { node: { type: "string" } } },
|
||||
] },
|
||||
// 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: {} }] },
|
||||
],
|
||||
}));
|
||||
await controller.handle("seat:node-dns-resolver.lookup@anchor", async () => ({ machine: "anchor", answered: true }));
|
||||
await controller.handle("seat:mesh-controller.status", async () => ({ output: "all quiet", ok: true }));
|
||||
await controller.handle("seat:mesh-controller.push", async (body: { node?: string }) => ({ told: body.node ?? "nobody" }));
|
||||
t.after(async () => {
|
||||
await catalogue.close();
|
||||
await shop.close();
|
||||
await controller.close();
|
||||
});
|
||||
|
||||
const { mkdtemp, writeFile } = await import("node:fs/promises");
|
||||
const { join } = await import("node:path");
|
||||
const dir = await mkdtemp("/tmp/mesh-client-");
|
||||
const path = join(dir, "credential.json");
|
||||
await writeFile(
|
||||
path,
|
||||
JSON.stringify({ url, user: "person.ada", password: "x", person: "ada", invokes: ["shop.price"] }),
|
||||
);
|
||||
return path;
|
||||
}
|
||||
|
||||
/** Drive `mesh mcp` over stdio and collect the replies, as a host would. */
|
||||
function driving(credential: string, requests: unknown[]): Promise<Record<string, any>[]> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const child = spawn(process.execPath, ["dist/mesh.js", "mcp", "--credential", credential], {
|
||||
stdio: ["pipe", "pipe", "pipe"],
|
||||
});
|
||||
let out = "";
|
||||
let err = "";
|
||||
child.stdout.on("data", (d) => (out += d.toString()));
|
||||
child.stderr.on("data", (d) => (err += d.toString()));
|
||||
child.on("error", reject);
|
||||
child.on("close", () => {
|
||||
const replies = out
|
||||
.split("\n")
|
||||
.filter((l) => l.trim() !== "")
|
||||
.map((l) => JSON.parse(l) as Record<string, any>);
|
||||
if (replies.length === 0 && err !== "") reject(new Error(err));
|
||||
else resolve(replies);
|
||||
});
|
||||
for (const r of requests) child.stdin.write(`${JSON.stringify(r)}\n`);
|
||||
child.stdin.end();
|
||||
});
|
||||
}
|
||||
|
||||
test("a host initialises, lists the mesh's tools and calls one", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "initialize", params: {} },
|
||||
{ jsonrpc: "2.0", method: "notifications/initialized" },
|
||||
{ jsonrpc: "2.0", id: 2, method: "tools/list" },
|
||||
{ jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "shop.price", arguments: { of: "a hat" } } },
|
||||
]);
|
||||
|
||||
const byId = new Map(replies.map((r) => [r.id, r]));
|
||||
// A notification is answered with nothing, or a host waiting on ids sees a reply it cannot match.
|
||||
assert.equal(replies.length, 3, `expected three replies, got ${JSON.stringify(replies)}`);
|
||||
|
||||
const hello = byId.get(1)!.result;
|
||||
assert.equal(hello.protocolVersion, "2025-03-26");
|
||||
assert.ok(hello.capabilities.tools, "a server offering no tools is not this one");
|
||||
assert.match(hello.instructions, /ada/, "the handshake says whose authority a call is made under");
|
||||
|
||||
const listed = byId.get(2)!.result.tools;
|
||||
assert.deepEqual(listed.map((x: { name: string }) => x.name),
|
||||
["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");
|
||||
// A node-scoped seat's verb takes the machine, and requires it (ADR 0170).
|
||||
const lookup = listed[2];
|
||||
assert.equal(lookup.inputSchema.properties.node.type, "string");
|
||||
assert.deepEqual(lookup.inputSchema.required, ["node"]);
|
||||
const price = listed[3];
|
||||
assert.ok(price.inputSchema, "a tool with no schema is one an agent cannot call");
|
||||
// A module's bare property map arrives as a schema an agent can read, its words kept — and
|
||||
// `node`, the machine to ask when the module runs on several (novox/hq ADR 0159), beside them.
|
||||
assert.deepEqual(price.inputSchema.properties.of, { type: "string" });
|
||||
assert.equal(price.inputSchema.properties.node.type, "string", "a module's tool takes the machine to ask");
|
||||
assert.ok(!listed[1].inputSchema.properties?.node, "a seat's verb takes no machine; the seat's scope decides");
|
||||
assert.equal(listed[0].inputSchema.properties?.node?.type, "string", "a seat's verb that takes a node of its own keeps it");
|
||||
// Silence is named: the module the catalogue holds and nothing answered for.
|
||||
assert.deepEqual(byId.get(2)!.result._meta.notAnswering, ["ghost"]);
|
||||
|
||||
const called = byId.get(3)!.result;
|
||||
assert.ok(!called.isError, `the call failed: ${JSON.stringify(called)}`);
|
||||
// The module's own answer, unshaped. An adapter that summarised it would be deciding what matters
|
||||
// in somebody else's answer.
|
||||
assert.deepEqual(JSON.parse(called.content[0].text), { of: "a hat", cost: 12 });
|
||||
});
|
||||
|
||||
test("a tool nobody serves comes back as an error the agent can act on", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "ghost.missing", arguments: {} } },
|
||||
]);
|
||||
const result = replies[0].result;
|
||||
// isError, not a protocol failure: the call was well-formed and the mesh answered it — with an
|
||||
// absence. A JSON-RPC error would tell the agent its request was malformed, which it was not.
|
||||
assert.ok(result?.isError, `expected a tool error, got ${JSON.stringify(replies[0])}`);
|
||||
assert.match(result.content[0].text, /nothing serves ghost\.missing/);
|
||||
});
|
||||
|
||||
test("a method this surface does not have is refused, and a notification is not", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "resources/list" },
|
||||
{ jsonrpc: "2.0", method: "notifications/cancelled" },
|
||||
]);
|
||||
assert.equal(replies.length, 1, "a notification was answered");
|
||||
assert.equal(replies[0].error.code, -32601);
|
||||
assert.match(replies[0].error.message, /resources\/list/);
|
||||
});
|
||||
|
||||
// A seat's verb that takes a machine as its own argument — `push <node>` — keeps it: the console
|
||||
// moves `node` into the subject for a module's tool only (ADR 0159), never for a role's verb.
|
||||
test("a seat's verb keeps a node of its own; only a module's tool gives it to the subject", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "mesh-controller.push", arguments: { node: "anchor" } } },
|
||||
]);
|
||||
const result = replies[0].result;
|
||||
assert.ok(!result.isError, JSON.stringify(replies[0]));
|
||||
assert.deepEqual(JSON.parse(result.content[0].text), { told: "anchor" });
|
||||
});
|
||||
|
||||
test("a host calls the mesh's own verb through the seat", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "tools/call", params: { name: "mesh-controller.status", arguments: {} } },
|
||||
]);
|
||||
const result = replies[0].result;
|
||||
assert.ok(!result.isError, JSON.stringify(replies[0]));
|
||||
assert.deepEqual(JSON.parse(result.content[0].text), { output: "all quiet", ok: true });
|
||||
});
|
||||
|
||||
test("a node-scoped seat's verb is asked of the machine named, and refused without one", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const credential = await aMeshAndACredential(t);
|
||||
const replies = await driving(credential, [
|
||||
{ jsonrpc: "2.0", id: 1, method: "initialize", params: {} },
|
||||
{ jsonrpc: "2.0", id: 2, method: "tools/call", params: { name: "node-dns-resolver.lookup", arguments: { node: "anchor" } } },
|
||||
{ jsonrpc: "2.0", id: 3, method: "tools/call", params: { name: "node-dns-resolver.lookup", arguments: {} } },
|
||||
]);
|
||||
const byId = new Map(replies.map((r) => [r.id, r]));
|
||||
const answered = byId.get(2)!.result;
|
||||
assert.ok(!answered.isError, JSON.stringify(answered));
|
||||
assert.match(answered.content[0].text, /"machine": "anchor"/, "the machine's holder answered");
|
||||
assert.match(byId.get(3)!.error?.message ?? JSON.stringify(byId.get(3)), /name the machine/);
|
||||
});
|
||||
@@ -0,0 +1,225 @@
|
||||
/**
|
||||
* The mesh issues an assignment's subjects, and a runtime serves what it is issued (novox/hq ADR
|
||||
* 0160). A runtime derives one address for itself — `mesh.assignment.<node>.<module>` — reads the
|
||||
* membership there, serves exactly what it says, and re-serves when a new one arrives. Against a
|
||||
* real bus with JetStream, because the membership is a direct get on a stream.
|
||||
*
|
||||
* 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/membership.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { connect, StringCodec } from "nats";
|
||||
|
||||
import { resetTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
|
||||
import { callTool, subjectListed } from "../dist/client.js";
|
||||
import { runTools } from "../dist/runtime.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
||||
const sc = StringCodec();
|
||||
|
||||
/** The ASSIGNMENTS stream as the controller asserts it: last-per-subject, readable by direct get. */
|
||||
async function anAssignmentsStream() {
|
||||
const nc = await connect({ servers: url! });
|
||||
const jsm = await nc.jetstreamManager();
|
||||
try {
|
||||
await jsm.streams.delete("ASSIGNMENTS");
|
||||
} catch {
|
||||
// none yet
|
||||
}
|
||||
await jsm.streams.add({
|
||||
name: "ASSIGNMENTS",
|
||||
subjects: ["mesh.assignment.>"],
|
||||
max_msgs_per_subject: 1,
|
||||
allow_direct: true,
|
||||
} as never);
|
||||
return {
|
||||
async issue(m: object & { node: string; module: string }) {
|
||||
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
|
||||
},
|
||||
async close() {
|
||||
await nc.close();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
test("a runtime serves exactly the subjects it is issued, and the listing carries them", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
// The mesh placed shop on two machines that are not interchangeable: each is served by name only.
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
const shop = await connectNats({ url, node: "anchor", module: "shop" });
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
|
||||
assert.equal(shop.membership()?.module, "shop", "the runtime read what the mesh issued");
|
||||
// By name it answers; the plain subject was not issued, so nothing serves it.
|
||||
const named = await callTool(asker, "shop.price@anchor", { of: "a hat" });
|
||||
assert.deepEqual(named.result, { of: "a hat", cost: 12 });
|
||||
assert.equal(named.node, "anchor");
|
||||
await assert.rejects(callTool(asker, "shop.price", {}), /no responders|503/i, "a subject not issued is not served");
|
||||
// The tools answer says where each is served, and the client composes nothing.
|
||||
const tools = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>(
|
||||
"shop.tools",
|
||||
{},
|
||||
);
|
||||
assert.deepEqual(tools.tools.find((x) => x.name === "price")?.subjects, ["mesh.mod.shop.tool.price.anchor"]);
|
||||
const listing = { tools: [{ module: "shop", name: "price", subjects: ["mesh.mod.shop.tool.price.anchor"] }], notAnswering: [] };
|
||||
assert.equal(subjectListed("shop.price", "anchor", listing as never), "mesh.mod.shop.tool.price.anchor");
|
||||
assert.equal(subjectListed("shop.price", "elsewhere", listing as never), undefined);
|
||||
|
||||
// The mesh re-issues the membership with the plain subject too (the modules became
|
||||
// interchangeable); the runtime re-serves on it without a restart.
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
await new Promise((r) => setTimeout(r, 300));
|
||||
const plain = await callTool(asker, "shop.price", { of: "a coat" });
|
||||
assert.deepEqual(plain.result, { of: "a coat", cost: 12 });
|
||||
assert.equal(plain.node, "anchor");
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await shop.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a seat's verbs are implemented under the seat's name, served where issued, never listed as the module's", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "postgres",
|
||||
serves: [{ subject: "mesh.mod.postgres.tool.{tool}", queue: "serve.postgres" }, { subject: "mesh.mod.postgres.tool.{tool}.anchor" }],
|
||||
seats: [{ seat: "mesh-store", verb: "databases", subject: "mesh.seat.mesh-store.tool.databases" }],
|
||||
emits: "mesh.mod.postgres.event.{event}",
|
||||
tools: "mesh.mod.postgres.tool.tools",
|
||||
});
|
||||
const credential = { url, node: "anchor", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] };
|
||||
const pg = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat.mjs")] });
|
||||
// The module's own `databases` and the seat's are two tools: postgres's on its subject, the
|
||||
// store's on the seat's, each answering as itself.
|
||||
const own = await callTool(asker, "postgres.databases", {});
|
||||
assert.deepEqual(own.result, { software: "postgres" });
|
||||
const seat = await callTool(asker, "seat:mesh-store.databases", {});
|
||||
assert.deepEqual(seat.result, { seat: "mesh-store" });
|
||||
assert.equal(seat.node, "anchor");
|
||||
// What the seat does not promise — creating a database — is postgres's tool and not the store's.
|
||||
await assert.rejects(callTool(asker, "seat:mesh-store.postgres_create_database", {}), /no responders|503/i);
|
||||
// And the module's `tools` lists only postgres's own, never the seat's implementation.
|
||||
const tools = await asker.request<Record<string, never>, { module: string; tools: { name: string }[] }>("postgres.tools", {});
|
||||
assert.deepEqual(tools.tools.map((x) => x.name).sort(), ["databases", "postgres_create_database"]);
|
||||
await assert.rejects(asker.request("mesh-store.tools", {}), /no responders|503/i, "a seat is not a module with a tools verb");
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await pg.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
seats: [{ seat: "shop", verb: "price", subject: "mesh.seat.shop.tool.price" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
const credential = { url, node: "anchor", module: "shop", claims: [{ seat: "shop", scope: "mesh", serves: ["price"] }] };
|
||||
const shop = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: shop, credential, moduleEntrypoints: [fixture("shop-seat.mjs")] });
|
||||
assert.deepEqual((await callTool(asker, "shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
|
||||
assert.deepEqual((await callTool(asker, "seat:shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
|
||||
const tools = await asker.request<Record<string, never>, { tools: { name: string }[] }>("shop.tools", {});
|
||||
assert.deepEqual(tools.tools.map((x) => x.name), ["price", "refund"]);
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await shop.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a mesh that has issued nothing yet gets the derived shape, and says so", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const stream = await anAssignmentsStream();
|
||||
const said: string[] = [];
|
||||
const log = console.log;
|
||||
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||
let lone;
|
||||
try {
|
||||
lone = await connectNats({ url, node: "home-server", module: "lone" });
|
||||
} finally {
|
||||
console.log = log;
|
||||
}
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
try {
|
||||
assert.equal(lone.membership(), undefined);
|
||||
assert.ok(said.some((s) => /no membership issued for lone on home-server/.test(s)), said.join("\n"));
|
||||
await lone.handle("ping", async () => ({ pong: true }));
|
||||
assert.deepEqual((await callTool(asker, "lone.ping", {})).result, { pong: true });
|
||||
assert.deepEqual((await callTool(asker, "lone.ping@home-server", {})).result, { pong: true });
|
||||
} finally {
|
||||
await asker.close();
|
||||
await lone.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
// A module that implements a seat its credential does not (yet) claim is not served for it and does
|
||||
// not fall over either: the runtime says so and serves the module's own tools.
|
||||
test("a registration under a seat the credential does not claim is said and skipped, not fatal", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
const credential = { url, node: "anchor", module: "postgres" };
|
||||
const pg = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
const said: string[] = [];
|
||||
const log = console.log;
|
||||
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat-unclaimed.mjs")] });
|
||||
console.log = log;
|
||||
assert.ok(said.some((s) => /registers tools under "mesh-store".*not served/.test(s)), said.join("\n"));
|
||||
assert.deepEqual((await callTool(asker, "postgres.databases", {})).result, { software: "postgres" });
|
||||
await assert.rejects(callTool(asker, "seat:mesh-store.databases", {}), /no responders|503/i);
|
||||
} finally {
|
||||
console.log = log;
|
||||
stop();
|
||||
await asker.close();
|
||||
await pg.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,194 @@
|
||||
/**
|
||||
* One runtime per node serves every assigned module's tools (novox/hq ADR 0175, to-be 38 WP1). The
|
||||
* node's runtime is handed a list of modules and their bundles on the node's credential; it reads
|
||||
* one membership per module, serves each module's tools on that module's subjects and each held
|
||||
* seat's verbs on the seat's, names a bundle that fails to load without dropping the others, and
|
||||
* re-serves a module whose membership is re-issued mid-run. Against a real bus with JetStream.
|
||||
*
|
||||
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/node-runtime.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { connect, StringCodec } from "nats";
|
||||
|
||||
import { resetTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
|
||||
import { callTool, toolsOn } from "../dist/client.js";
|
||||
import { servedModulesFrom } from "../dist/main.js";
|
||||
import { runTools } from "../dist/runtime.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
||||
const sc = StringCodec();
|
||||
|
||||
/** The controller's job, done by hand: the ASSIGNMENTS stream (last-per-subject, direct get) and an
|
||||
* EVENTS stream for what a tool emits. */
|
||||
async function aMesh() {
|
||||
const nc = await connect({ servers: url! });
|
||||
const jsm = await nc.jetstreamManager();
|
||||
for (const name of ["ASSIGNMENTS", "EVENTS"]) {
|
||||
try {
|
||||
await jsm.streams.delete(name);
|
||||
} catch {
|
||||
// none yet
|
||||
}
|
||||
}
|
||||
await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never);
|
||||
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>"] });
|
||||
return {
|
||||
async issue(m: object & { node: string; module: string }) {
|
||||
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
|
||||
},
|
||||
/** The next subject an event lands on, under a pattern. */
|
||||
nextEvent(pattern: string): Promise<string> {
|
||||
const sub = nc.subscribe(pattern, { max: 1 });
|
||||
return (async () => {
|
||||
for await (const m of sub) return m.subject;
|
||||
throw new Error("no event");
|
||||
})();
|
||||
},
|
||||
async close() {
|
||||
await nc.close();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** A membership as the controller issues one on a machine, with the module's own subject when it
|
||||
* answers for the module anywhere, and the verbs of the node seats it holds. */
|
||||
function membershipOf(module: string, node: string, opts: { plain?: boolean; seats?: Record<string, string[]> } = {}) {
|
||||
const own = `mesh.mod.${module}`;
|
||||
const serves: { subject: string; queue?: string }[] = [{ subject: `${own}.tool.{tool}.${node}` }];
|
||||
if (opts.plain) serves.push({ subject: `${own}.tool.{tool}`, queue: `serve.${module}` });
|
||||
const seats = Object.entries(opts.seats ?? {}).flatMap(([seat, verbs]) =>
|
||||
verbs.map((verb) => ({ seat, verb, subject: `mesh.seat.${seat}.tool.${verb}.${node}` })));
|
||||
return { node, module, serves, seats, emits: `${own}.event.{event}`, tools: `${own}.tool.tools` };
|
||||
}
|
||||
|
||||
/** Wait for something to be served: the bus answers "no responders" at once until it is. */
|
||||
async function until<T>(attempt: () => Promise<T>, tries = 50): Promise<T> {
|
||||
for (let i = 0; ; i++) {
|
||||
try {
|
||||
return await attempt();
|
||||
} catch (e) {
|
||||
if (i >= tries) throw e;
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
test("MESH_TOOL_MODULES names modules and their entrypoints; a bare path is the credential's own module's", () => {
|
||||
const have = servedModulesFrom(" alpha=/a/tools/index.js, beta=/b/one.js ,beta=/b/two.js, /mine/index.js ,node-tools=/own/x.js", "node-tools");
|
||||
assert.deepEqual(have.serves, [
|
||||
{ module: "alpha", entrypoints: ["/a/tools/index.js"] },
|
||||
{ module: "beta", entrypoints: ["/b/one.js", "/b/two.js"] },
|
||||
]);
|
||||
// The credential's own module, named or bare, is the one-module form either way.
|
||||
assert.deepEqual(have.moduleEntrypoints, ["/mine/index.js", "/own/x.js"]);
|
||||
assert.throws(() => servedModulesFrom("=/nothing.js", undefined), /neither <module>=<entrypoint>/);
|
||||
assert.deepEqual(servedModulesFrom("", "x"), { serves: [], moduleEntrypoints: [] });
|
||||
});
|
||||
|
||||
test("the node's runtime serves five modules' bundles on one credential — two of them launched, one broken — and follows a re-issued membership", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const mesh = await aMesh();
|
||||
// Three modules assigned to the machine: alpha answers for itself anywhere, beta only here and
|
||||
// holds the node-shelf seat, gamma's bundle is broken.
|
||||
await mesh.issue(membershipOf("alpha", "anchor", { plain: true }));
|
||||
await mesh.issue(membershipOf("beta", "anchor", { seats: { "node-shelf": ["list", "clear"] } }));
|
||||
await mesh.issue(membershipOf("gamma", "anchor"));
|
||||
// Two more, launched rather than loaded (ADR 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();
|
||||
}
|
||||
});
|
||||
@@ -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 });
|
||||
});
|
||||
@@ -0,0 +1,80 @@
|
||||
import { spawn } from "node:child_process";
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { fatalBrokerReason, PinMismatchError, topicMatches } from "../src/broker-nats.ts";
|
||||
|
||||
// novox/hq issue 058 (and its review): serve mode retries a broker that is not up yet, but must
|
||||
// give up at once on a failure waiting cannot fix — otherwise a permanent fault loops for ever
|
||||
// disguised as "not reachable". This is the classifier that draws the line; the bed cannot test it
|
||||
// (it starts the consumer only after the broker is up), so it is proven here.
|
||||
|
||||
test("a broker that is not up yet is retryable, not fatal", () => {
|
||||
for (const err of [
|
||||
Object.assign(new Error("connect ECONNREFUSED 10.42.0.1:5671"), { code: "ECONNREFUSED" }),
|
||||
Object.assign(new Error("connect ETIMEDOUT"), { code: "ETIMEDOUT" }),
|
||||
Object.assign(new Error("getaddrinfo EAI_AGAIN anchor.internal"), { code: "EAI_AGAIN" }),
|
||||
new Error("timed out fetching the broker's certificate"),
|
||||
]) {
|
||||
assert.equal(fatalBrokerReason(err), null, `should retry: ${(err as Error).message}`);
|
||||
}
|
||||
});
|
||||
|
||||
test("a certificate that does not match the pin is fatal, by type not by message", () => {
|
||||
// Typed, so rewording the message cannot turn an impostor back into an infinite retry.
|
||||
assert.notEqual(fatalBrokerReason(new PinMismatchError("anything at all")), null);
|
||||
// A plain Error with pin-ish words is NOT treated as the pin case — only the type is.
|
||||
assert.equal(fatalBrokerReason(new Error("the pinned value was fine")), null);
|
||||
});
|
||||
|
||||
test("a malformed broker URL is fatal — it never parses on the next try", () => {
|
||||
assert.notEqual(fatalBrokerReason(Object.assign(new Error("Invalid URL"), { code: "ERR_INVALID_URL" })), null);
|
||||
assert.notEqual(fatalBrokerReason(new Error("Invalid URL: not-a-url")), null);
|
||||
});
|
||||
|
||||
test("a refused login is fatal — a wrong or revoked credential, not an absent broker", () => {
|
||||
// The bus refuses a login in its own words; each is final, because the next try says the same.
|
||||
for (const msg of ["Authorization Violation", "nats: user authentication expired", "Permissions Violation for Subscription to \"x\""]) {
|
||||
assert.notEqual(fatalBrokerReason(new Error(msg)), null, `should be fatal: ${msg}`);
|
||||
}
|
||||
});
|
||||
|
||||
test("a non-Error value does not crash the classifier", () => {
|
||||
assert.equal(fatalBrokerReason("just a string"), null);
|
||||
assert.equal(fatalBrokerReason(undefined), null);
|
||||
});
|
||||
|
||||
// **A module asked to prepare its state and naming nothing is a failure, not a no-op** (novox/hq
|
||||
// ADR 0135). The mesh asks this only of a module whose manifest says it prepares something, so an
|
||||
// image that names nothing was built wrong, and exiting 0 would let that version serve against a
|
||||
// state nobody shaped.
|
||||
test("preparing with nothing named fails rather than passing quietly", async () => {
|
||||
const runtime = new URL("../dist/main.js", import.meta.url).pathname;
|
||||
const ran = await new Promise<{ code: number | null; said: string }>((resolve) => {
|
||||
const child = spawn(process.execPath, [runtime, "prepare"], {
|
||||
env: { ...process.env, MESH_PREPARE: "" },
|
||||
});
|
||||
let said = "";
|
||||
child.stderr.on("data", (chunk) => (said += String(chunk)));
|
||||
child.on("close", (code) => resolve({ code, said }));
|
||||
});
|
||||
assert.notEqual(ran.code, 0, "a module that prepares nothing exited 0, so its version would serve");
|
||||
assert.match(ran.said, /MESH_PREPARE/);
|
||||
});
|
||||
|
||||
// **One durable consumer feeds one reader, however many patterns a module registers.**
|
||||
//
|
||||
// A module has exactly one consumer, so two readers of it would each take half the messages — and a
|
||||
// reader that received one its own pattern does not match acknowledges it, which is right for a
|
||||
// filter wider than anything registered and silent loss when it is another handler's. The matching is
|
||||
// therefore pure and tested as such: what a message is for is decided by the patterns registered, not
|
||||
// by which reader happened to fetch it.
|
||||
test("a message is for every pattern that matches it, and nothing else", () => {
|
||||
const registered = ["mesh-build-machine.built", "mesh-controller.built-before"];
|
||||
const matched = (key: string) => registered.filter((p) => topicMatches(p, key));
|
||||
assert.deepEqual(matched("mesh-build-machine.built"), ["mesh-build-machine.built"]);
|
||||
assert.deepEqual(matched("mesh-controller.built-before"), ["mesh-controller.built-before"]);
|
||||
// Nothing registered for it: the consumer's filter is the controller's and may be wider.
|
||||
assert.deepEqual(matched("mesh-catalog.upgraded"), []);
|
||||
// And a handler that asked for everything gets both, which is what the audit logger does.
|
||||
assert.deepEqual(["#"].filter((p) => topicMatches(p, "mesh-controller.built-before")), ["#"]);
|
||||
});
|
||||
@@ -0,0 +1,57 @@
|
||||
// Round-trip the runtime's NATS client against a real server: a tool call answered, and an
|
||||
// event emitted and received with its envelope intact.
|
||||
import { connect } from "nats";
|
||||
import { connectNats } from "../dist/broker-nats.js";
|
||||
|
||||
const URL = "nats://127.0.0.1:14222";
|
||||
|
||||
// The controller's job, done by hand here: the stream and the module's durable consumer.
|
||||
const admin = await connect({ servers: URL });
|
||||
const jsm = await admin.jetstreamManager();
|
||||
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>", "mesh.seat.*.event.>"] });
|
||||
await jsm.consumers.add("EVENTS", {
|
||||
durable_name: "one_audit", ack_policy: "explicit",
|
||||
filter_subjects: ["mesh.mod.shop.event.order.placed"],
|
||||
});
|
||||
|
||||
const shop = await connectNats({ url: URL, node: "one", module: "shop" });
|
||||
const audit = await connectNats({ url: URL, node: "one", module: "audit" });
|
||||
|
||||
let failures = 0;
|
||||
const check = (ok, what) => { console.log(` ${ok ? "ok " : "FAIL"} ${what}`); if (!ok) failures++; };
|
||||
|
||||
// A tool, served and called.
|
||||
await shop.handle("price", async (body) => ({ total: body.qty * 3 }));
|
||||
const answer = await audit.request("shop.price", { qty: 4 });
|
||||
check(answer.total === 12, "a tool call is answered across two connections");
|
||||
|
||||
// A handler that throws reaches the caller as an error, not a timeout.
|
||||
await shop.handle("boom", async () => { throw new Error("no"); });
|
||||
let threw = null;
|
||||
try { await audit.request("shop.boom", {}); } catch (e) { threw = e.message; }
|
||||
check(threw === "no", "a handler that throws answers the caller instead of timing out");
|
||||
|
||||
// An event, emitted and received with its envelope intact.
|
||||
const seen = [];
|
||||
await audit.subscribe("order.placed", async (env) => { seen.push(env); });
|
||||
await shop.publish({
|
||||
key: "order.placed", node: "one", body: { id: "a1" },
|
||||
headers: { "x-event-id": "e1", "x-node": "one", "content-type": "application/json" },
|
||||
});
|
||||
await new Promise((r) => setTimeout(r, 800));
|
||||
check(seen.length === 1, `exactly one delivery (saw ${seen.length})`);
|
||||
if (seen[0]) {
|
||||
check(seen[0].key === "order.placed", "the key survives the subject round trip");
|
||||
check(seen[0].body?.id === "a1", "the body is the payload, not the whole envelope");
|
||||
check(seen[0].node === "one", "the node comes back from the headers");
|
||||
check(seen[0].headers?.["x-event-id"] === "e1", "the event id survives as a header");
|
||||
}
|
||||
|
||||
// A module cannot reach into another's namespace by naming its own event oddly.
|
||||
await shop.publish({ key: "other", node: "one", body: {}, headers: { "x-event-id": "e2" } });
|
||||
const msg = await jsm.streams.getMessage("EVENTS", { last_by_subj: "mesh.mod.shop.event.other" });
|
||||
check(!!msg, "an event lands under the emitting module's own namespace");
|
||||
|
||||
await shop.close(); await audit.close(); await admin.close();
|
||||
console.log(failures ? `\n${failures} failed` : "\nall passed");
|
||||
process.exit(failures ? 1 : 0);
|
||||
@@ -0,0 +1,60 @@
|
||||
/**
|
||||
* Every module's runtime answers `tools` for it (novox/hq ADR 0152, design 34 §3): the names,
|
||||
* descriptions and schemas from the code that answers them. Against a real bus, because the claim is
|
||||
* what a second connection gets back.
|
||||
*
|
||||
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/runtime-tools.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
|
||||
import { resetTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
import { connectNats } from "../dist/broker-nats.js";
|
||||
import { runTools } from "../dist/runtime.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
||||
|
||||
test("a module registering two tools answers three names, the third being what it serves", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const shop = await connectNats({ url, node: "one", module: "shop" });
|
||||
const asker = await connectNats({ url, module: "person.ada" });
|
||||
const stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
|
||||
try {
|
||||
const answer = await asker.request<Record<string, never>, { module: string; tools: { name: string; input: unknown }[] }>(
|
||||
"shop.tools",
|
||||
{},
|
||||
);
|
||||
assert.equal(answer.module, "shop");
|
||||
assert.deepEqual(answer.tools.map((x) => x.name), ["price", "refund"]);
|
||||
// The schema travels with the name: a name alone is not callable by something that has never
|
||||
// seen the mesh before.
|
||||
assert.deepEqual(answer.tools[0].input, { of: { type: "string", description: "the thing" } });
|
||||
// And the tools themselves still answer beside it.
|
||||
const priced = await asker.request<{ of: string }, { cost: number }>("shop.price", { of: "a hat" });
|
||||
assert.equal(priced.cost, 12);
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await shop.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a module naming a tool of its own `tools` is refused at load", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const clash = await connectNats({ url, node: "one", module: "clash" });
|
||||
try {
|
||||
await assert.rejects(
|
||||
() => runTools({ broker: clash, moduleEntrypoints: [fixture("clash-tools.mjs")] }),
|
||||
/names a tool "tools"/,
|
||||
);
|
||||
} finally {
|
||||
resetTools();
|
||||
await clash.close();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user