The mesh composes every served module's words into MESH_TOOL_ENV; the runtime takes it at start and removes it from its own environment, then gives each registration's contributor and each launched child the runtime's words plus its own module's, never another's. Against an SDK without collectToolsEach it says so and serves with the runtime's words only. The test serves two imported bundles and one launched, each answering with its own words and none of the others'.
292 lines
16 KiB
TypeScript
292 lines
16 KiB
TypeScript
/**
|
|
* One runtime per node serves every assigned module's tools (novox/hq ADR 0175, to-be 38 WP1). The
|
|
* node's runtime is handed a list of modules and their bundles on the node's credential; it reads
|
|
* one membership per module, serves each module's tools on that module's subjects and each held
|
|
* seat's verbs on the seat's, names a bundle that fails to load without dropping the others, and
|
|
* re-serves a module whose membership is re-issued mid-run. Against a real bus with JetStream.
|
|
*
|
|
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
|
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/node-runtime.test.ts
|
|
*/
|
|
import assert from "node:assert/strict";
|
|
import { test } from "node:test";
|
|
import { fileURLToPath } from "node:url";
|
|
import { cpSync, mkdirSync, mkdtempSync, realpathSync, rmSync, writeFileSync } from "node:fs";
|
|
import { join } from "node:path";
|
|
import { tmpdir } from "node:os";
|
|
import { connect, StringCodec } from "nats";
|
|
|
|
import { resetTools } from "@novox/mesh-sdk/tools";
|
|
|
|
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
|
|
import { callTool, toolsOn } from "../dist/client.js";
|
|
import { servedModulesFrom } from "../dist/main.js";
|
|
import { runTools, takeToolEnvs } from "../dist/runtime.js";
|
|
|
|
const url = process.env.MESH_TEST_NATS;
|
|
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
|
const sc = StringCodec();
|
|
|
|
/** The controller's job, done by hand: the ASSIGNMENTS stream (last-per-subject, direct get) and an
|
|
* EVENTS stream for what a tool emits. */
|
|
async function aMesh() {
|
|
const nc = await connect({ servers: url! });
|
|
const jsm = await nc.jetstreamManager();
|
|
for (const name of ["ASSIGNMENTS", "EVENTS"]) {
|
|
try {
|
|
await jsm.streams.delete(name);
|
|
} catch {
|
|
// none yet
|
|
}
|
|
}
|
|
await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never);
|
|
await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>"] });
|
|
return {
|
|
async issue(m: object & { node: string; module: string }) {
|
|
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
|
|
},
|
|
/** The next subject an event lands on, under a pattern. */
|
|
nextEvent(pattern: string): Promise<string> {
|
|
const sub = nc.subscribe(pattern, { max: 1 });
|
|
return (async () => {
|
|
for await (const m of sub) return m.subject;
|
|
throw new Error("no event");
|
|
})();
|
|
},
|
|
async close() {
|
|
await nc.close();
|
|
},
|
|
};
|
|
}
|
|
|
|
/** A membership as the controller issues one on a machine, with the module's own subject when it
|
|
* answers for the module anywhere, and the verbs of the node seats it holds. */
|
|
function membershipOf(module: string, node: string, opts: { plain?: boolean; seats?: Record<string, string[]> } = {}) {
|
|
const own = `mesh.mod.${module}`;
|
|
const serves: { subject: string; queue?: string }[] = [{ subject: `${own}.tool.{tool}.${node}` }];
|
|
if (opts.plain) serves.push({ subject: `${own}.tool.{tool}`, queue: `serve.${module}` });
|
|
const seats = Object.entries(opts.seats ?? {}).flatMap(([seat, verbs]) =>
|
|
verbs.map((verb) => ({ seat, verb, subject: `mesh.seat.${seat}.tool.${verb}.${node}` })));
|
|
return { node, module, serves, seats, emits: `${own}.event.{event}`, tools: `${own}.tool.tools` };
|
|
}
|
|
|
|
/** Wait for something to be served: the bus answers "no responders" at once until it is. */
|
|
async function until<T>(attempt: () => Promise<T>, tries = 50): Promise<T> {
|
|
for (let i = 0; ; i++) {
|
|
try {
|
|
return await attempt();
|
|
} catch (e) {
|
|
if (i >= tries) throw e;
|
|
await new Promise((r) => setTimeout(r, 100));
|
|
}
|
|
}
|
|
}
|
|
|
|
test("MESH_TOOL_MODULES names modules and their entrypoints; a bare path is the credential's own module's", () => {
|
|
const have = servedModulesFrom(" alpha=/a/tools/index.js, beta=/b/one.js ,beta=/b/two.js, /mine/index.js ,node-tools=/own/x.js", "node-tools");
|
|
assert.deepEqual(have.serves, [
|
|
{ module: "alpha", entrypoints: ["/a/tools/index.js"] },
|
|
{ module: "beta", entrypoints: ["/b/one.js", "/b/two.js"] },
|
|
]);
|
|
// The credential's own module, named or bare, is the one-module form either way.
|
|
assert.deepEqual(have.moduleEntrypoints, ["/mine/index.js", "/own/x.js"]);
|
|
assert.throws(() => servedModulesFrom("=/nothing.js", undefined), /neither <module>=<entrypoint>/);
|
|
assert.deepEqual(servedModulesFrom("", "x"), { serves: [], moduleEntrypoints: [] });
|
|
});
|
|
|
|
test("the node's runtime serves five modules' bundles on one credential — two of them launched, one broken — and follows a re-issued membership", async (t) => {
|
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
|
resetTools();
|
|
const mesh = await aMesh();
|
|
// Three modules assigned to the machine: alpha answers for itself anywhere, beta only here and
|
|
// holds the node-shelf seat, gamma's bundle is broken.
|
|
await mesh.issue(membershipOf("alpha", "anchor", { plain: true }));
|
|
await mesh.issue(membershipOf("beta", "anchor", { seats: { "node-shelf": ["list", "clear"] } }));
|
|
await mesh.issue(membershipOf("gamma", "anchor"));
|
|
// Two more, launched rather than loaded (ADR 0188): delta is Python and holds the node-lamp seat;
|
|
// epsilon is TypeScript written against the protocol and marked executable.
|
|
await mesh.issue(membershipOf("delta", "anchor", { seats: { "node-lamp": ["on"] } }));
|
|
await mesh.issue(membershipOf("epsilon", "anchor"));
|
|
// The node's credential: the runtime module's name, no claims (seats come from the memberships).
|
|
const credential = { url, node: "anchor", module: "node-tools" };
|
|
const nodeTools = await connectNats(credential);
|
|
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
|
const said: string[] = [];
|
|
const log = console.log;
|
|
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
|
let stop = () => {};
|
|
try {
|
|
process.env.MESH_OPERATOR_ACCOUNT = "somebody";
|
|
process.env.MESH_OPERATOR_HOME = "/home/somebody";
|
|
stop = await runTools({
|
|
broker: nodeTools,
|
|
credential,
|
|
serves: [
|
|
{ module: "alpha", entrypoints: [fixture("many-alpha.mjs")] },
|
|
{ module: "beta", entrypoints: [fixture("many-beta.mjs")] },
|
|
{ module: "gamma", entrypoints: [fixture("many-broken.mjs")] },
|
|
{ module: "delta", entrypoints: [fixture("many-delta.py")] },
|
|
{ module: "epsilon", entrypoints: [fixture("many-epsilon.mjs")] },
|
|
],
|
|
});
|
|
console.log = log;
|
|
assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "delta", "epsilon", "gamma", "node-tools"]);
|
|
assert.ok(said.some((s) => /the operator's account here is somebody \(home \/home\/somebody\)/.test(s)), said.join("\n"));
|
|
assert.ok(said.some((s) => /gamma's bundle .*many-broken\.mjs failed to load: gamma's bundle cannot find its client; its tools are not served here/.test(s)), said.join("\n"));
|
|
assert.ok(said.some((s) => /serving 8 tool\(s\) for 5 module\(s\): alpha\.one, alpha\.two, beta\.three, beta\.four, beta\.five, delta\.greet, delta\.die, epsilon\.seven; not serving gamma/.test(s)), said.join("\n"));
|
|
|
|
// Five tools answer, each where its module's membership says: alpha anywhere and here, beta here only.
|
|
assert.deepEqual((await callTool(asker, "alpha.one", {})).result, { alpha: 1 });
|
|
assert.deepEqual((await callTool(asker, "alpha.one@anchor", {})).result, { alpha: 1 });
|
|
assert.deepEqual((await callTool(asker, "beta.three@anchor", {})).result, { beta: 3 });
|
|
assert.deepEqual((await callTool(asker, "beta.four@anchor", {})).result, { beta: 4 });
|
|
assert.deepEqual((await callTool(asker, "beta.five@anchor", {})).result, { beta: 5 });
|
|
await assert.rejects(callTool(asker, "beta.three", {}), /no responders|503/i, "beta was not issued the module's plain subject");
|
|
|
|
// Two seat verbs answer on the seat's subjects, held by beta.
|
|
assert.deepEqual((await callTool(asker, "seat:node-shelf.list@anchor", {})).result, { shelf: ["a", "b"] });
|
|
assert.deepEqual((await callTool(asker, "seat:node-shelf.clear@anchor", {})).result, { cleared: true });
|
|
|
|
// A bundle in another language answers the same way, its seat verb among them; so does a
|
|
// TypeScript bundle served through the protocol rather than imported.
|
|
assert.deepEqual((await callTool(asker, "delta.greet@anchor", { who: "mesh" })).result, { greeting: "hello mesh", language: "python" });
|
|
assert.deepEqual((await callTool(asker, "seat:node-lamp.on@anchor", {})).result, { on: true, language: "python" });
|
|
assert.deepEqual((await callTool(asker, "epsilon.seven@anchor", {})).result, { epsilon: 7, via: "stdio" });
|
|
// A child that exits mid-call tells the caller so and is started again on the next call.
|
|
await assert.rejects(callTool(asker, "delta.die@anchor", {}), /delta's bundle exited \(3\)/);
|
|
assert.deepEqual((await callTool(asker, "delta.greet@anchor", {})).result, { greeting: "hello world", language: "python" });
|
|
|
|
// A tool that emits does so as its module, not as the runtime.
|
|
const landed = mesh.nextEvent("mesh.mod.*.event.>");
|
|
assert.deepEqual((await callTool(asker, "alpha.two", {})).result, { alpha: 2 });
|
|
assert.equal(await landed, "mesh.mod.alpha.event.happened");
|
|
|
|
// `tools` answers for each: what alpha and beta serve, and why gamma serves nothing.
|
|
const gamma = await asker.request<Record<string, never>, { module: string; tools: unknown[]; failed?: string }>("gamma.tools@anchor", {});
|
|
assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle cannot find its client" });
|
|
const beta = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>("beta.tools@anchor", {});
|
|
assert.deepEqual(beta.tools.map((x) => x.name), ["three", "four", "five"]);
|
|
assert.deepEqual(beta.tools[0]!.subjects, ["mesh.mod.beta.tool.three.anchor"]);
|
|
|
|
// And discovery says so, with the reason, beside the modules that answered.
|
|
const catalogue = await connectNats({ url, module: "mesh-catalog" });
|
|
await catalogue.handle("catalog_modules", async () => ({ modules: [{ module: "alpha" }, { module: "beta" }, { module: "gamma" }, { module: "delta" }, { module: "epsilon" }] }));
|
|
try {
|
|
const have = await toolsOn(asker);
|
|
assert.deepEqual(have.tools.map((x) => `${x.module}.${x.name}`), ["alpha.one", "alpha.two", "beta.five", "beta.four", "beta.three", "delta.die", "delta.greet", "epsilon.seven"]);
|
|
assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle cannot find its client)", "mesh-controller (seat)"]);
|
|
} finally {
|
|
await catalogue.close();
|
|
}
|
|
|
|
// The mesh re-issues beta's membership mid-run — now answering for the module anywhere — and
|
|
// the runtime serves the new subject without a restart.
|
|
await mesh.issue(membershipOf("beta", "anchor", { plain: true, seats: { "node-shelf": ["list", "clear"] } }));
|
|
assert.deepEqual((await until(() => callTool(asker, "beta.three", {}))).result, { beta: 3 });
|
|
assert.deepEqual((await callTool(asker, "seat:node-shelf.list@anchor", {})).result, { shelf: ["a", "b"] });
|
|
} finally {
|
|
console.log = log;
|
|
delete process.env.MESH_OPERATOR_ACCOUNT;
|
|
delete process.env.MESH_OPERATOR_HOME;
|
|
stop();
|
|
await asker.close();
|
|
await nodeTools.close();
|
|
await mesh.close();
|
|
resetTools();
|
|
}
|
|
});
|
|
|
|
test("a bundle carrying its own copy of the SDK registers into the runtime's registry, and its tools are served (issue 209)", async (t) => {
|
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
|
resetTools();
|
|
let dir = "";
|
|
let stop = () => {};
|
|
const closing: Array<() => Promise<void>> = [];
|
|
const said: string[] = [];
|
|
const log = console.log;
|
|
try {
|
|
// A bundle as the toolchain packs one: its compiled entrypoint, a package.json saying ES modules,
|
|
// and its dependencies copied in — the SDK among them, a second copy beside the runtime's own,
|
|
// and a dependency of the bundle's own that the runtime does not carry.
|
|
dir = mkdtempSync(join(tmpdir(), "mesh-bundle-"));
|
|
const sdk = realpathSync(fileURLToPath(new URL("../node_modules/@novox/mesh-sdk/", import.meta.url)));
|
|
cpSync(sdk, join(dir, "node_modules", "@novox", "mesh-sdk"), { recursive: true });
|
|
mkdirSync(join(dir, "node_modules", "zeta-flavour"), { recursive: true });
|
|
writeFileSync(join(dir, "node_modules", "zeta-flavour", "package.json"), '{"name":"zeta-flavour","type":"module","main":"index.js"}\n');
|
|
writeFileSync(join(dir, "node_modules", "zeta-flavour", "index.js"), 'export const flavour = "the bundle\'s own";\n');
|
|
writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n');
|
|
writeFileSync(join(dir, "index.js"),
|
|
'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' +
|
|
'import { flavour } from "zeta-flavour";\n' +
|
|
'registerModuleTools("zeta", () => [{ name: "probe", description: "answers", input: {}, run: async () => ({ zeta: true, flavour }) }]);\n');
|
|
const mesh = await aMesh();
|
|
closing.push(() => mesh.close());
|
|
await mesh.issue(membershipOf("zeta", "anchor"));
|
|
const credential = { url, node: "anchor", module: "node-tools" };
|
|
const nodeTools = await connectNats(credential);
|
|
closing.push(() => nodeTools.close());
|
|
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
|
closing.push(() => asker.close());
|
|
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
|
stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.js")] }] });
|
|
console.log = log;
|
|
assert.ok(said.some((s) => /serving 1 tool\(s\) for 1 module\(s\): zeta\.probe/.test(s)), said.join("\n"));
|
|
// The SDK is the runtime's (the registration arrived); the bundle's other dependency is its own.
|
|
assert.deepEqual((await callTool(asker, "zeta.probe@anchor", {})).result, { zeta: true, flavour: "the bundle's own" });
|
|
} finally {
|
|
console.log = log;
|
|
stop();
|
|
for (const close of closing.reverse()) await close();
|
|
if (dir) rmSync(dir, { recursive: true, force: true });
|
|
resetTools();
|
|
}
|
|
});
|
|
|
|
test("each bundle is given its own environment and none of another's, imported or launched (ADR 0192)", async (t) => {
|
|
if (!url) return t.skip("MESH_TEST_NATS unset");
|
|
resetTools();
|
|
const mesh = await aMesh();
|
|
for (const m of ["gamma", "delta", "zeta"]) await mesh.issue(membershipOf(m, "anchor"));
|
|
const credential = { url, node: "anchor", module: "node-tools" };
|
|
const nodeTools = await connectNats(credential);
|
|
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
|
const log = console.log;
|
|
let stop = () => {};
|
|
const before = process.env.MESH_TOOL_ENV;
|
|
try {
|
|
process.env.MESH_OPERATOR_ACCOUNT = "somebody";
|
|
process.env.MESH_TOOL_ENV = JSON.stringify({
|
|
gamma: { GAMMA_CONFIG_FILE: "/var/lib/mesh/gamma/config.json" },
|
|
delta: { DELTA_TOKEN_FILE: "/var/lib/mesh/delta/token" },
|
|
zeta: { ZETA_URL: "http://127.0.0.1:3000" },
|
|
});
|
|
const envs = takeToolEnvs();
|
|
assert.equal(process.env.MESH_TOOL_ENV, undefined, "the composed environments were left in the process's");
|
|
console.log = () => {};
|
|
stop = await runTools({
|
|
broker: nodeTools, credential, envs,
|
|
serves: [
|
|
{ module: "gamma", entrypoints: [fixture("env-gamma.mjs")] },
|
|
{ module: "delta", entrypoints: [fixture("env-delta.mjs")] },
|
|
{ module: "zeta", entrypoints: [fixture("env-zeta.mjs")] },
|
|
],
|
|
});
|
|
console.log = log;
|
|
assert.deepEqual((await callTool(asker, "gamma.given@anchor", {})).result,
|
|
{ mine: "/var/lib/mesh/gamma/config.json", theirs: null, runtime: "somebody", composed: null });
|
|
assert.deepEqual((await callTool(asker, "delta.given@anchor", {})).result,
|
|
{ mine: "/var/lib/mesh/delta/token", theirs: null });
|
|
assert.deepEqual((await callTool(asker, "zeta.given@anchor", {})).result,
|
|
{ mine: "http://127.0.0.1:3000", theirs: null, composed: null });
|
|
} finally {
|
|
console.log = log;
|
|
if (before === undefined) delete process.env.MESH_TOOL_ENV; else process.env.MESH_TOOL_ENV = before;
|
|
delete process.env.MESH_OPERATOR_ACCOUNT;
|
|
stop();
|
|
await asker.close();
|
|
await nodeTools.close();
|
|
await mesh.close();
|
|
resetTools();
|
|
}
|
|
});
|