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.
226 lines
10 KiB
TypeScript
226 lines
10 KiB
TypeScript
/**
|
|
* 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();
|
|
}
|
|
});
|