The runtime serves a list of modules on one credential, naming a bundle that fails to load (hq ADR 0175, to-be 38 WP1)
One tool runtime per node, host-side, is what the runtime was written to be; the catalogue built a container per module around it instead. This lets `serve` take a list — MESH_TOOL_MODULES as <module>=<entrypoint> entries — and do for every assigned module what it did for one: read that module's membership and follow it, serve its tools where the membership says, serve each held seat's verbs on the seat's subjects. The seats come from the memberships now, so the node's credential carries no claims; a module's own runtime still reads its credential's, so nothing built today changes behaviour. A bare path in MESH_TOOL_MODULES stays the one-module form. A bundle that throws on import is said in the log and in what `tools` answers for its module (`failed`), which discovery lists with the reason instead of as "not answering"; the other bundles serve. The filter that dropped every registration under a name but the one module goes; what stays is that a registration under a seat's name is served only where some served module claims the seat. A tool runs attributed to its module, so an event it emits lands on the module's subject and not the runtime's. MESH_OPERATOR_ACCOUNT and MESH_OPERATOR_HOME are read and said; tools take them from their environment. Proven against a real bus: three bundles, one broken; five tools and two seat verbs answer on their subjects; `tools` names the failed bundle; a membership re-issued mid-run re-serves.
This commit is contained in:
@@ -1,29 +1,42 @@
|
||||
# mesh-tools
|
||||
|
||||
The Novox Mesh **tool runtime** — the per-node process that makes a module's tools actually serve.
|
||||
The Novox Mesh **tool runtime** — the one process per node that makes every assigned module's tools
|
||||
actually serve (novox/hq ADR 0175).
|
||||
|
||||
A module ships its tools (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk)); this
|
||||
runtime is what loads them and puts them on the mesh. It:
|
||||
A module ships its tools as a bundle (built on [`@novox/mesh-sdk`](https://git.novox.be/novox/mesh-sdk));
|
||||
this runtime is what loads them and puts them on the mesh. It:
|
||||
|
||||
1. connects the mesh broker (novox/hq ADR 0001) — a concrete AMQP implementation of the sdk's
|
||||
`Broker` contract;
|
||||
2. imports the assigned modules' compiled tool entrypoints, each of which registers its tools as it
|
||||
loads;
|
||||
3. serves them through the sdk's `serveTools` harness, answering `tools.invoke` over the broker.
|
||||
1. connects the mesh bus on the node's credential — a concrete implementation of the sdk's `Broker`
|
||||
contract;
|
||||
2. reads one membership per module it serves — what the mesh issued that module on this machine
|
||||
(ADR 0160): where its tools are answered, which seats it holds — and follows each live;
|
||||
3. imports each module's compiled tool entrypoints, each of which registers its tools as it loads,
|
||||
guarded: a bundle that throws is named, in the log and in what `tools` answers for its module,
|
||||
and the others serve;
|
||||
4. serves every module's tools on that module's subjects and every held seat's verbs on the seat's.
|
||||
|
||||
Everything hard — dispatch, collection, duplicate-name safety — is the sdk's. This is the thin
|
||||
wrapper that binds the broker and loads the modules. Keeping the AMQP client here, out of the sdk,
|
||||
is deliberate: a broker-client change never rebuilds a module (ADR 0039).
|
||||
wrapper that binds the bus and loads the modules. Keeping the bus client here, out of the sdk, is
|
||||
deliberate: a bus-client change never rebuilds a module (ADR 0039). The runtime is module-agnostic:
|
||||
it knows bundles and subjects, nothing of what any module does. A tool that needs root escalates
|
||||
itself — root is the module's concern, not the runtime's.
|
||||
|
||||
## Running it
|
||||
|
||||
```
|
||||
MESH_BROKER_URL amqp://… the mesh broker
|
||||
MESH_TOOL_MODULES /a/tools/index.js,… the assigned modules' compiled tool entrypoints
|
||||
MESH_BROKER_FILE the node's sealed credential, as the mesh delivered it
|
||||
MESH_TOOL_MODULES alpha=/…/alpha/tools/index.js,beta=/…/beta/dist/index.js,…
|
||||
the modules to serve and their compiled entrypoints; several entries may
|
||||
name one module. A bare path is an entrypoint of the credential's own module
|
||||
— the one-module form a per-module container still sets.
|
||||
MESH_OPERATOR_ACCOUNT whose machine this is, and MESH_OPERATOR_HOME where their home is; set by
|
||||
the mesh when the node has an account, read by tools from their environment
|
||||
MESH_BROKER_URL a plain URL instead of the credential, for the bootstrap case
|
||||
```
|
||||
|
||||
`node dist/main.js`, or the container (`Dockerfile`). On a node the host resolves both variables
|
||||
and starts it like any other supervised workload.
|
||||
`node dist/main.js`. On a node the controller composes the variables and the host supervises the
|
||||
process like any other host-side workload (novox/hq to-be 38). The container (`Dockerfile`) is how
|
||||
a module's own *service* may still be built; it is no longer how tools reach a node.
|
||||
|
||||
## `mesh` — the tools for whoever is on a machine
|
||||
|
||||
@@ -41,6 +54,9 @@ dropped. A module may not name a tool of its own `tools`; the runtime refuses it
|
||||
|
||||
## Verified
|
||||
|
||||
`npm test` stands up LavinMQ (the mesh's broker) and proves the whole path over real AMQP: the
|
||||
runtime serves a registered tool, a separate connection invokes it by name and gets the result, and
|
||||
an unknown tool is refused over the wire.
|
||||
`npm test` runs against a real NATS server with JetStream (`MESH_TEST_NATS`, see any test's header
|
||||
for the one-line `docker run`) and proves the whole path over the wire: the runtime serves a
|
||||
registered tool, a separate connection invokes it by name and gets the result, an unknown tool is
|
||||
refused, a runtime serves exactly the subjects it is issued and re-serves on a new membership, and
|
||||
the node's runtime serves three modules' bundles on one credential — one of them broken, named and
|
||||
not fatal — with every seat verb answering where the membership put it.
|
||||
|
||||
+86
-42
@@ -16,6 +16,7 @@
|
||||
// derives the subject (design 29 §1), so reorganising the subject space leaves every module
|
||||
// correct.
|
||||
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import { createHash } from "node:crypto";
|
||||
import net from "node:net";
|
||||
import tls from "node:tls";
|
||||
@@ -71,19 +72,41 @@ export interface Answered<Res> {
|
||||
node?: string;
|
||||
}
|
||||
|
||||
/** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs —
|
||||
* an answer that says which machine gave it, and serving a subject that is not a module's own tool
|
||||
* (a seat's verb). */
|
||||
/** The bus as the runtime sees it: the sdk's contract, and the few things only the runtime needs —
|
||||
* an answer that says which machine gave it, serving a subject that is not a module's own tool
|
||||
* (a seat's verb), and following the memberships of the modules it serves beside its own. */
|
||||
export interface RuntimeBroker extends Broker {
|
||||
/** Call a tool by key, or — when `on` names a subject the mesh listed for it (ADR 0160) — there. */
|
||||
ask<Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>>;
|
||||
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void>;
|
||||
/** What the mesh issued this assignment, or undefined when nothing has been issued yet. */
|
||||
membership(): Membership | undefined;
|
||||
/** Called when the mesh issues a new membership; the runtime re-serves on it. */
|
||||
/**
|
||||
* Serve another module's tools from this connection: read its membership on this machine and
|
||||
* follow it, so `handle("<module>.<tool>")` is served where the mesh issued that module (novox/hq
|
||||
* ADR 0175: one runtime per node, every assigned module's tools). The account must be allowed
|
||||
* to read that membership and to subscribe its subjects — the node's is; a module's own is not,
|
||||
* and is refused by the bus, not here.
|
||||
*/
|
||||
follow(module: string): Promise<void>;
|
||||
/** What the mesh issued this assignment — this module's when none is named — or undefined when
|
||||
* nothing has been issued yet. */
|
||||
membership(module?: string): Membership | undefined;
|
||||
/** Called when the mesh issues a new membership to any module this connection follows; the
|
||||
* runtime re-serves on it. The membership says which module it is for. */
|
||||
onMembership(handler: (m: Membership) => void): void;
|
||||
/** The modules whose memberships this connection follows: its own and every one `follow` added. */
|
||||
serving(): string[];
|
||||
/** The module this connection is: what its credential named, and what a bare key serves as. */
|
||||
readonly module: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Which module's tool is at work on this connection, when one is (novox/hq ADR 0175). One runtime
|
||||
* carries many modules' tools, and an event a tool emits must land on the emitting *module's*
|
||||
* subject, not the runtime's — so the runtime runs each tool inside this store, and `publish`
|
||||
* reads the module from it. Outside a tool the connection's own module stands.
|
||||
*/
|
||||
export const atWork = new AsyncLocalStorage<{ module: string }>();
|
||||
|
||||
const ASSIGNMENTS_STREAM = "ASSIGNMENTS";
|
||||
|
||||
/** The one address a runtime derives for itself (ADR 0160). */
|
||||
@@ -153,38 +176,45 @@ export async function connectNats(
|
||||
|
||||
const node = cred.node;
|
||||
|
||||
// The membership, read once at connect and followed. A direct get is one request on the
|
||||
// stream's API, which is the whole of what this account may ask JetStream for its own subject;
|
||||
// a 404 is a mesh that has not issued one, which is a fact to say and not an error to retry.
|
||||
let issued: Membership | undefined;
|
||||
// The memberships, one per module this connection serves, each read once and followed. A
|
||||
// module's own runtime follows one, its own; the node's runtime follows one per assigned module
|
||||
// (novox/hq ADR 0175) — the same subject shape, the same stream, read as many times as there are
|
||||
// modules, and nothing on the bus learns a new shape for it. A direct get is one request on the
|
||||
// stream's API, which is the whole of what an account may ask JetStream for a subject it is
|
||||
// granted; a 404 is a mesh that has not issued one, which is a fact to say and not an error to
|
||||
// retry.
|
||||
const issued = new Map<string, Membership | undefined>();
|
||||
const issuedHandlers: ((m: Membership) => void)[] = [];
|
||||
const subjectOfMine = node ? membershipSubject(node, self) : "";
|
||||
if (subjectOfMine) {
|
||||
const follow = async (module: string): Promise<void> => {
|
||||
if (!node || issued.has(module)) return;
|
||||
issued.set(module, undefined);
|
||||
const subjectOfTheirs = membershipSubject(node, module);
|
||||
try {
|
||||
// The subject-addressed form of a direct get — the stream, then the subject, nothing in
|
||||
// the body — because that is the one address the mesh grants this account on the
|
||||
// stream's API; the body form asks the stream's root, which it may not.
|
||||
const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}.${subjectOfMine}`,
|
||||
const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}.${subjectOfTheirs}`,
|
||||
new Uint8Array(0), { timeout: 5_000 });
|
||||
const status = got.headers?.code ?? 0;
|
||||
if (status === 0 && got.data.length > 0) {
|
||||
issued = JSON.parse(sc.decode(got.data)) as Membership;
|
||||
issued.set(module, JSON.parse(sc.decode(got.data)) as Membership);
|
||||
}
|
||||
} catch {
|
||||
// Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below.
|
||||
}
|
||||
if (!issued) {
|
||||
console.log(`[mesh-tools] no membership issued for ${self} on ${node} yet; serving the derived shape until one arrives`);
|
||||
if (!issued.get(module)) {
|
||||
console.log(`[mesh-tools] no membership issued for ${module} on ${node} yet; serving the derived shape until one arrives`);
|
||||
}
|
||||
try {
|
||||
const live = conn.subscribe(subjectOfMine);
|
||||
const live = conn.subscribe(subjectOfTheirs);
|
||||
subs.push(live);
|
||||
void (async () => {
|
||||
for await (const msg of live) {
|
||||
try {
|
||||
issued = JSON.parse(sc.decode(msg.data)) as Membership;
|
||||
console.log(`[mesh-tools] ${self} on ${node} was issued a new membership; re-serving on it`);
|
||||
for (const h of issuedHandlers) h(issued);
|
||||
const m = JSON.parse(sc.decode(msg.data)) as Membership;
|
||||
issued.set(module, m);
|
||||
console.log(`[mesh-tools] ${module} on ${node} was issued a new membership; re-serving on it`);
|
||||
for (const h of issuedHandlers) h(m);
|
||||
} catch (err) {
|
||||
console.log(`[mesh-tools] a membership arrived that is not one: ${err}`);
|
||||
}
|
||||
@@ -193,32 +223,33 @@ export async function connectNats(
|
||||
} catch {
|
||||
// A subscription this account may not make is a mesh older than the membership.
|
||||
}
|
||||
}
|
||||
};
|
||||
await follow(self);
|
||||
|
||||
/** The subjects a tool of this module is served on: from the membership when issued, derived
|
||||
* otherwise (the shape the mesh issues on day one, so the two agree). */
|
||||
const servedOn = (tool: string): { subject: string; queue?: string }[] => {
|
||||
if (issued) {
|
||||
const m = issued;
|
||||
/** The subjects a tool of a served module is served on: from that module's membership when
|
||||
* issued, derived otherwise (the shape the mesh issues on day one, so the two agree). */
|
||||
const servedOn = (module: string, tool: string): { subject: string; queue?: string }[] => {
|
||||
const m = issued.get(module);
|
||||
if (m) {
|
||||
const out = m.serves.map((s) => ({ subject: s.subject.replace("{tool}", tool), queue: s.queue }));
|
||||
// The verb that lists what this module serves is answered on the mesh's plain address for it
|
||||
// whatever the placement — one answer suffices, so a queue — and on this machine's beside it.
|
||||
if (tool === "tools" && m.tools && !out.some((s) => s.subject === m.tools)) {
|
||||
out.unshift({ subject: m.tools, queue: `serve.${self}` });
|
||||
out.unshift({ subject: m.tools, queue: `serve.${module}` });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
const base = `mesh.mod.${self}.tool.${tool}`;
|
||||
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }];
|
||||
const base = `mesh.mod.${module}.tool.${tool}`;
|
||||
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${module}` }];
|
||||
if (node) out.push({ subject: `${base}.${node}` });
|
||||
return out;
|
||||
};
|
||||
|
||||
/** Where a call by key goes: a subject the membership says this module reaches, when it says
|
||||
* one — the machine's when named — else the derived shape. */
|
||||
/** Where a call by key goes: a subject this connection's own membership says it reaches, when it
|
||||
* says one — the machine's when named — else the derived shape. */
|
||||
const reachedAt = (key: string): string => {
|
||||
const [name, wanted] = key.split("@", 2);
|
||||
const reach = issued?.reaches?.[name];
|
||||
const reach = issued.get(self)?.reaches?.[name];
|
||||
if (reach && reach.length > 0) {
|
||||
if (wanted) {
|
||||
const at = reach.find((s) => s.endsWith(`.${wanted}`));
|
||||
@@ -288,27 +319,38 @@ export async function connectNats(
|
||||
*/
|
||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||
// A seat's verb named outright is served on the seat's subject as given, for a holder that
|
||||
// knows its role without a membership; everything else is this module's own tool, served
|
||||
// where the mesh issued it (ADR 0160). A key naming another module is not served here at all.
|
||||
// knows its role without a membership; everything else is a served module's own tool, served
|
||||
// where the mesh issued that module (ADR 0160). A key naming a module this connection does
|
||||
// not follow is not served here at all: a module serves its own tools, and the node's
|
||||
// runtime those of the modules it was given (ADR 0175) — never a stranger's.
|
||||
if (key.startsWith("seat:")) {
|
||||
const stop = answerOn(toolSubject(key, self), undefined, handler);
|
||||
return () => stop();
|
||||
}
|
||||
const dot = key.indexOf(".");
|
||||
if (dot >= 0 && key.slice(0, dot) !== self) {
|
||||
throw new Error(`${self} cannot serve ${key}: a module serves its own tools`);
|
||||
const module = dot < 0 ? self : key.slice(0, dot);
|
||||
if (!issued.has(module) && module !== self) {
|
||||
throw new Error(
|
||||
`${self} cannot serve ${key}: a module serves its own tools, and a runtime those of the ` +
|
||||
"modules it follows",
|
||||
);
|
||||
}
|
||||
const tool = dot < 0 ? key : key.slice(dot + 1);
|
||||
let stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
// When a new membership arrives, serve where it now says and stop serving where it no longer does.
|
||||
issuedHandlers.push(() => {
|
||||
let stops = servedOn(module, tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
// When that module's new membership arrives, serve where it now says and stop serving where
|
||||
// it no longer does.
|
||||
issuedHandlers.push((m) => {
|
||||
if (m.module !== module) return;
|
||||
stops.forEach((stop) => stop());
|
||||
stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
stops = servedOn(module, tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
});
|
||||
return () => stops.forEach((stop) => stop());
|
||||
},
|
||||
|
||||
membership: () => issued,
|
||||
module: self,
|
||||
follow,
|
||||
serving: () => [...issued.keys()],
|
||||
membership: (module?: string) => issued.get(module ?? self),
|
||||
onMembership: (handler: (m: Membership) => void) => {
|
||||
issuedHandlers.push(handler);
|
||||
},
|
||||
@@ -339,7 +381,9 @@ export async function connectNats(
|
||||
if (!meta["content-type"]) h.set("content-type", "application/json");
|
||||
if (env.node) h.set("x-node", env.node);
|
||||
|
||||
await js.publish(eventSubject(env.key, self), sc.encode(JSON.stringify(env.body)), {
|
||||
// The emitting module's subject: the tool at work's when a served module's tool emits from
|
||||
// the node's runtime (ADR 0175), this connection's own otherwise.
|
||||
await js.publish(eventSubject(env.key, atWork.getStore()?.module ?? self), sc.encode(JSON.stringify(env.body)), {
|
||||
headers: h,
|
||||
// De-duplicated by the server inside its window, so a redelivery after a crash between
|
||||
// publishing and acknowledging is not seen twice. Only the emitter can make this id.
|
||||
|
||||
+6
-1
@@ -192,7 +192,12 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
|
||||
}
|
||||
asked.forEach((outcome, i) => {
|
||||
const module = names[i]!;
|
||||
if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
|
||||
if (outcome.status === "fulfilled" && typeof outcome.value?.failed === "string") {
|
||||
// The runtime answered for it and serves nothing: the bundle failed to load (ADR 0175). Said
|
||||
// with the reason, because "not answering" would send somebody to check an assignment that
|
||||
// is fine.
|
||||
notAnswering.push(`${module} (its tools bundle failed to load: ${outcome.value.failed})`);
|
||||
} else if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
|
||||
for (const t of outcome.value.tools) {
|
||||
tools.push({ module, name: t.name, description: t.description, input: t.input, subjects: t.subjects });
|
||||
}
|
||||
|
||||
+40
-8
@@ -21,14 +21,15 @@
|
||||
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0043) — an
|
||||
// amqps account scoped to this module. Preferred: a module holds its own.
|
||||
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
|
||||
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
|
||||
// MESH_TOOL_MODULES <module>=/path/a,… the modules to serve and their compiled entrypoints (serve
|
||||
// mode; ADR 0175); a bare path is an entrypoint of the credential's own module
|
||||
// MESH_PREPARE /path/a,/path/b,… compiled entrypoints that prepare this module's state
|
||||
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
|
||||
|
||||
import { readFileSync } from "node:fs";
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js";
|
||||
import { runTools } from "./runtime.js";
|
||||
import { runTools, type ServedModule } from "./runtime.js";
|
||||
|
||||
/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */
|
||||
let lastCredential: Credential | undefined;
|
||||
@@ -118,14 +119,42 @@ async function connectBrokerPatiently(): Promise<Broker> {
|
||||
}
|
||||
}
|
||||
|
||||
async function serve(): Promise<void> {
|
||||
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
||||
.split(",")
|
||||
.map((s) => s.trim())
|
||||
.filter(Boolean);
|
||||
/**
|
||||
* What MESH_TOOL_MODULES names (novox/hq ADR 0175, to-be 38 WP1): `<module>=<entrypoint>` entries,
|
||||
* comma-separated, several per module allowed — the node's runtime serving every assigned module's
|
||||
* bundle. A bare path is the one-module form the per-module containers still set: an entrypoint of
|
||||
* the credential's own module. Both may appear; the result is one list of modules.
|
||||
*/
|
||||
export function servedModulesFrom(spec: string, own: string | undefined): { serves: ServedModule[]; moduleEntrypoints: string[] } {
|
||||
const serves = new Map<string, string[]>();
|
||||
const moduleEntrypoints: string[] = [];
|
||||
for (const raw of spec.split(",")) {
|
||||
const entry = raw.trim();
|
||||
if (!entry) continue;
|
||||
const eq = entry.indexOf("=");
|
||||
if (eq < 0) {
|
||||
moduleEntrypoints.push(entry);
|
||||
continue;
|
||||
}
|
||||
const module = entry.slice(0, eq).trim();
|
||||
const path = entry.slice(eq + 1).trim();
|
||||
if (!module || !path) {
|
||||
throw new Error(`MESH_TOOL_MODULES: "${entry}" is neither <module>=<entrypoint> nor an entrypoint of this module`);
|
||||
}
|
||||
if (module === own) {
|
||||
moduleEntrypoints.push(path);
|
||||
continue;
|
||||
}
|
||||
serves.set(module, [...(serves.get(module) ?? []), path]);
|
||||
}
|
||||
return { serves: [...serves].map(([module, entrypoints]) => ({ module, entrypoints })), moduleEntrypoints };
|
||||
}
|
||||
|
||||
async function serve(): Promise<void> {
|
||||
const broker = await connectBrokerPatiently();
|
||||
const stop = await runTools({ broker, moduleEntrypoints, credential: lastCredential });
|
||||
// Parsed after connecting: a bare entrypoint belongs to the module the credential names.
|
||||
const { serves, moduleEntrypoints } = servedModulesFrom(process.env.MESH_TOOL_MODULES ?? "", lastCredential?.module);
|
||||
const stop = await runTools({ broker, serves, moduleEntrypoints, credential: lastCredential });
|
||||
|
||||
const shutdown = async (): Promise<void> => {
|
||||
stop();
|
||||
@@ -250,4 +279,7 @@ async function main(): Promise<void> {
|
||||
await serve();
|
||||
}
|
||||
|
||||
// Only when run, so a test can import the pieces.
|
||||
if (process.argv[1] && import.meta.url === new URL(`file://${process.argv[1]}`).href) {
|
||||
void main();
|
||||
}
|
||||
|
||||
+224
-88
@@ -1,14 +1,20 @@
|
||||
// The tool runtime — the thin per-node process that makes a module's tools actually serve. It
|
||||
// binds the mesh broker, loads the assigned modules' tool entrypoints (each of which calls
|
||||
// registerModuleTools as it imports), and hands them to the sdk's serving harness. Everything hard
|
||||
// — dispatch, collection, duplicate-name safety — is the sdk's; this is the wrapper.
|
||||
// The tool runtime — the per-node process that makes the mesh's tools actually serve (novox/hq
|
||||
// ADR 0175). It binds the mesh broker, loads the served modules' tool bundles (each of which calls
|
||||
// registerModuleTools as it imports), and serves every module's tools on that module's subjects and
|
||||
// every held seat's verbs on the seat's. Everything hard — dispatch, collection, duplicate-name
|
||||
// safety — is the sdk's; this is the wrapper.
|
||||
//
|
||||
// One runtime, many modules. It was written for one module per process and ran that way in a
|
||||
// container per module; it now serves a list, as the one process per node the host supervises,
|
||||
// and the per-module shape is the list with one entry. A bundle that fails to import is named —
|
||||
// in the log and in what `tools` answers for it — and the others serve.
|
||||
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { resolve } from "node:path";
|
||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||
import { collectTools, toolKey } from "@novox/mesh-sdk/tools";
|
||||
import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||
import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
||||
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
||||
|
||||
/**
|
||||
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
|
||||
@@ -28,137 +34,267 @@ export interface ToolsAnswer {
|
||||
* first when there is one, then this machine's. A caller composes nothing. */
|
||||
subjects?: string[];
|
||||
}[];
|
||||
/** Why this module serves nothing here, when its bundle failed to load (ADR 0175): said where
|
||||
* discovery looks, so a module that is silent and one that is broken are told apart. */
|
||||
failed?: string;
|
||||
}
|
||||
|
||||
/** One module this runtime serves: its name and its compiled tool entrypoints. */
|
||||
export interface ServedModule {
|
||||
module: string;
|
||||
/** Absolute paths to the module's compiled tool entrypoints (e.g. .../umami/tools/index.js). */
|
||||
entrypoints: string[];
|
||||
}
|
||||
|
||||
export interface RuntimeOptions {
|
||||
/** The mesh broker to serve over. */
|
||||
broker: Broker;
|
||||
/** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */
|
||||
moduleEntrypoints: string[];
|
||||
/** The modules to serve, each with its entrypoints. */
|
||||
serves?: ServedModule[];
|
||||
/** The credential's own module's entrypoints — the one-module form, which the per-module
|
||||
* containers still use; the same as naming the credential's module in `serves`. */
|
||||
moduleEntrypoints?: string[];
|
||||
/** The credential the mesh delivered, for what it says about the seats this module claims
|
||||
* (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat. */
|
||||
* (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat its
|
||||
* memberships do not name. */
|
||||
credential?: Credential;
|
||||
}
|
||||
|
||||
/** Two environment words the mesh sets for the node's runtime and every tool reads from its
|
||||
* environment: whose machine this is (novox/hq to-be 37 §3, ADR 0175). */
|
||||
export const OPERATOR_ACCOUNT = "MESH_OPERATOR_ACCOUNT";
|
||||
export const OPERATOR_HOME = "MESH_OPERATOR_HOME";
|
||||
|
||||
/** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */
|
||||
export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
useBroker(() => opts.broker);
|
||||
const runtime = opts.broker as RuntimeBroker;
|
||||
// Whose runtime this is: the credential's module, or the connection's own when a runtime is
|
||||
// started by hand without one — the broker was told its module when it connected.
|
||||
const self = opts.credential?.module ?? (typeof runtime.module === "string" ? runtime.module : undefined);
|
||||
|
||||
for (const entry of opts.moduleEntrypoints) {
|
||||
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake.
|
||||
await import(pathToFileURL(resolve(entry)).href);
|
||||
// What to serve: the list, with the one-module form folded in as the credential's own entry.
|
||||
const served = new Map<string, string[]>();
|
||||
for (const s of opts.serves ?? []) {
|
||||
served.set(s.module, [...(served.get(s.module) ?? []), ...s.entrypoints]);
|
||||
}
|
||||
if (opts.moduleEntrypoints?.length) {
|
||||
if (!self) {
|
||||
throw new Error(
|
||||
"entrypoints were given with no module to serve them as: name the module (MESH_TOOL_MODULES " +
|
||||
"as <module>=<entrypoint>) or connect on a credential that names one",
|
||||
);
|
||||
}
|
||||
served.set(self, [...(served.get(self) ?? []), ...opts.moduleEntrypoints]);
|
||||
}
|
||||
|
||||
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
|
||||
// audit logger) registers none, and its scoped account may not declare the serve queue — so a
|
||||
// runtime that always served would fail for exactly the modules that never needed it.
|
||||
// A registration under a seat's name is the module's implementation of that seat's verbs
|
||||
// (ADR 0159, 0160): served on the seat's subjects by serveClaimedSeats, never as a module's
|
||||
// tools and never listed among them. Everything else is the module's own.
|
||||
// A module named like its seat (the catalogue is the mesh-catalog seat) registers once and is
|
||||
// both: its tools are the module's and the seat's verbs alike.
|
||||
const self = opts.credential?.module;
|
||||
const seatNames = new Set((opts.credential?.claims ?? []).map((c) => c.seat));
|
||||
// A registration under a name that is neither this module nor a seat it claims is not served:
|
||||
// said, and left out, rather than fatal — on 2026-10-01 the credential of a module that had just
|
||||
// learned to implement a seat did not yet name the claim, and the whole runtime restarted for it.
|
||||
const ownRegistrations = collectTools().filter(({ module }) => {
|
||||
if (module === self || !self || seatNames.has(module)) return module === self || !self;
|
||||
console.log(`[mesh-tools] ${self} registers tools under "${module}", which is neither this module nor a seat its credential claims; not served until the mesh issues the claim`);
|
||||
// The operator's machine, said once so a tool's behaviour under it can be read back from the
|
||||
// log. Tools read the two words from their own environment, which is this process's.
|
||||
const account = process.env[OPERATOR_ACCOUNT];
|
||||
if (account) {
|
||||
console.log(`[mesh-tools] the operator's account here is ${account}` +
|
||||
(process.env[OPERATOR_HOME] ? ` (home ${process.env[OPERATOR_HOME]})` : ""));
|
||||
}
|
||||
|
||||
// Follow every served module's membership before loading anything, so what each is issued is
|
||||
// known when its tools are bound. A module's own runtime already follows its own.
|
||||
if (typeof runtime.follow === "function") {
|
||||
for (const module of served.keys()) await runtime.follow(module);
|
||||
}
|
||||
|
||||
// Import each bundle, guarded (ADR 0175: one faulty bundle must not take the node's tools down).
|
||||
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake — and
|
||||
// the registrations it adds are the ones that appear after it, which is how each is attributed
|
||||
// to the module whose bundle made it.
|
||||
const failed = new Map<string, string>();
|
||||
const owner: string[] = []; // registration index → the module whose bundle registered it
|
||||
for (const [module, entrypoints] of served) {
|
||||
for (const entry of entrypoints) {
|
||||
const before = collectTools().length;
|
||||
try {
|
||||
await import(pathToFileURL(resolve(entry)).href);
|
||||
} catch (err) {
|
||||
const why = err instanceof Error ? err.message : String(err);
|
||||
failed.set(module, why);
|
||||
console.log(`[mesh-tools] ${module}'s bundle ${entry} failed to load: ${why}; its tools are not served here`);
|
||||
}
|
||||
const after = collectTools().length;
|
||||
for (let i = before; i < after; i++) owner[i] = module;
|
||||
}
|
||||
}
|
||||
|
||||
// A registration under a served module's name is that module's tools, served on its subjects.
|
||||
// One under a seat's name is the module's implementation of that seat's verbs (ADR 0159, 0160):
|
||||
// served on the seat's subjects by serveClaimedSeats where some served module claims the seat,
|
||||
// never as a module's tools and never listed among them. A module named like its seat (the
|
||||
// catalogue is the mesh-catalog seat) registers once and is both. Anything else is said and left
|
||||
// out rather than fatal — on 2026-10-01 the credential of a module that had just learned to
|
||||
// implement a seat did not yet name the claim, and the whole runtime restarted for it.
|
||||
const claimed = seatsClaimed(served.keys(), self, opts.credential, runtime);
|
||||
const registrations = collectTools().map((r, i) => ({ ...r, owner: owner[i] ?? self ?? r.module }));
|
||||
const ownRegistrations = registrations.filter(({ module, owner: by }) => {
|
||||
if (served.has(module)) return true;
|
||||
if (claimed.has(module)) return false;
|
||||
console.log(`[mesh-tools] ${by} registers tools under "${module}", which is neither a module served here nor a seat one of them claims; not served until the mesh issues the claim`);
|
||||
return false;
|
||||
});
|
||||
const tools = ownRegistrations.flatMap(({ module, tools: own }) => own.map((t) => ({ module, name: t.name })));
|
||||
const stops: Array<() => void> = [];
|
||||
const stop = (): void => stops.splice(0).forEach((s) => s());
|
||||
// Each tool on its own key, namespaced by its module (ADR 0047); where that key is answered is
|
||||
// the broker's to know from the membership (ADR 0160).
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
const seen = new Set<string>();
|
||||
for (const t of own) {
|
||||
if (seen.has(t.name)) {
|
||||
stop();
|
||||
throw new Error(`${module} exposes two tools named ${t.name} — refused`);
|
||||
}
|
||||
seen.add(t.name);
|
||||
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) => t.run(args ?? {})));
|
||||
}
|
||||
}
|
||||
|
||||
// And, for every module that serves any, the verb that says what it serves. Refused before
|
||||
// anything is bound if a module named a tool of its own `tools`: one name answering two things
|
||||
// is the fault nobody can diagnose afterwards, and the runtime is the only place that sees both.
|
||||
const runtime = opts.broker as RuntimeBroker;
|
||||
// Refused before anything is bound if a module named a tool of its own `tools`: one name
|
||||
// answering two things is the fault nobody can diagnose afterwards, and the runtime is the only
|
||||
// place that sees both. Likewise two tools of one module under one name.
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
if (own.length === 0) continue;
|
||||
if (own.some((t) => t.name === TOOLS_VERB)) {
|
||||
stop();
|
||||
throw new Error(
|
||||
`${module} names a tool "${TOOLS_VERB}", which is the verb the runtime answers for every ` +
|
||||
"module with what it serves (novox/hq ADR 0152) — refused, rename it",
|
||||
);
|
||||
}
|
||||
const subjectsOf = (tool: string): string[] | undefined => {
|
||||
const issued = typeof runtime.membership === "function" ? runtime.membership() : undefined;
|
||||
if (!issued) return undefined;
|
||||
const plain = issued.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
const mine = issued.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
return [...plain, ...mine];
|
||||
};
|
||||
const answer: ToolsAnswer = {
|
||||
module,
|
||||
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
|
||||
};
|
||||
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async () => answer));
|
||||
const seen = new Set<string>();
|
||||
for (const t of own) {
|
||||
if (seen.has(t.name)) throw new Error(`${module} exposes two tools named ${t.name} — refused`);
|
||||
seen.add(t.name);
|
||||
}
|
||||
}
|
||||
|
||||
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
||||
stops.push(await serveClaimedSeats(opts.broker as RuntimeBroker, opts.credential));
|
||||
return () => {
|
||||
for (const s of stops) s();
|
||||
// Each tool on its own key, namespaced by its module (ADR 0047); where that key is answered is
|
||||
// the broker's to know from the module's membership (ADR 0160). A tool runs attributed to its
|
||||
// module, so what it emits lands on the module's subject and not the runtime's.
|
||||
const names: string[] = [];
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
for (const t of own) {
|
||||
names.push(toolKey(module, t.name));
|
||||
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) =>
|
||||
atWork.run({ module }, () => t.run(args ?? {}))));
|
||||
}
|
||||
}
|
||||
|
||||
// And, for every served module, the verb that says what it serves — nothing, and why, for a
|
||||
// module whose bundle failed. A module that registered nothing and did not fail is a pure-events
|
||||
// module (the audit logger), whose scoped account may not declare the serve queue; it is left
|
||||
// silent as it always was.
|
||||
const byModule = new Map<string, ToolDefinition[]>();
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
byModule.set(module, [...(byModule.get(module) ?? []), ...own]);
|
||||
}
|
||||
for (const module of served.keys()) {
|
||||
const own = byModule.get(module) ?? [];
|
||||
const why = failed.get(module);
|
||||
if (own.length === 0 && !why) continue;
|
||||
const subjectsOf = (tool: string): string[] | undefined => {
|
||||
const m = typeof runtime.membership === "function" ? runtime.membership(module) : undefined;
|
||||
if (!m) return undefined;
|
||||
const plain = m.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
const mine = m.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
return [...plain, ...mine];
|
||||
};
|
||||
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async (): Promise<ToolsAnswer> => ({
|
||||
module,
|
||||
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
|
||||
...(why ? { failed: why } : {}),
|
||||
})));
|
||||
}
|
||||
|
||||
console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` +
|
||||
(failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : ""));
|
||||
stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations));
|
||||
return () => stop();
|
||||
}
|
||||
|
||||
/** The seats some served module claims: from the credential for its own module, and from every
|
||||
* served module's membership (ADR 0160) — the node's runtime holds no claims of its own. */
|
||||
function seatsClaimed(
|
||||
modules: Iterable<string>,
|
||||
self: string | undefined,
|
||||
credential: Credential | undefined,
|
||||
runtime: RuntimeBroker,
|
||||
): Set<string> {
|
||||
const out = new Set<string>();
|
||||
for (const c of credential?.claims ?? []) if (credential?.module === self) out.add(c.seat);
|
||||
for (const module of modules) {
|
||||
const m = typeof runtime.membership === "function" ? runtime.membership(module) : undefined;
|
||||
for (const s of m?.seats ?? []) out.add(s.seat);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** One seat's verb, where its callers ask, and which served module holds the seat. */
|
||||
interface SeatVerb {
|
||||
seat: string;
|
||||
verb: string;
|
||||
subject: string;
|
||||
holder: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). The credential names the
|
||||
* seats this module claims and the verbs each promises; each verb is served on the seat's own
|
||||
* subject by the module's tool of the same name. Whether this instance *holds* the seat is the bus's
|
||||
* to decide: only the holder's account may subscribe the seat's subjects, so a claimant that does not
|
||||
* hold it here is refused the subscription and serves nothing — never a failure of its own tools.
|
||||
* Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). What a served module
|
||||
* claims and promises comes from its membership (ADR 0160) — and, for a module's own runtime, from
|
||||
* its credential, which named the claims before memberships did. Each verb is served on the seat's
|
||||
* own subject by the tool of the same name registered under the seat's name. Whether this instance
|
||||
* *holds* the seat is the bus's to decide: only the holder's account may subscribe the seat's
|
||||
* subjects, so a claimant that does not hold it here is refused the subscription and serves nothing
|
||||
* — never a failure of its own tools.
|
||||
*/
|
||||
async function serveClaimedSeats(broker: RuntimeBroker, credential?: Credential): Promise<() => void> {
|
||||
const claims = credential?.claims ?? [];
|
||||
if (claims.length === 0 || typeof broker.handleSubject !== "function") return () => {};
|
||||
async function serveClaimedSeats(
|
||||
broker: RuntimeBroker,
|
||||
served: string[],
|
||||
self: string | undefined,
|
||||
credential: Credential | undefined,
|
||||
registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
|
||||
): Promise<() => void> {
|
||||
if (typeof broker.handleSubject !== "function") return () => {};
|
||||
// A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's
|
||||
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
|
||||
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
|
||||
for (const { module, tools } of collectTools()) {
|
||||
if (!claims.some((c) => c.seat === module)) continue;
|
||||
const verbs = new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
for (const t of tools) verbs.set(t.name, (args) => t.run(args));
|
||||
for (const { module, owner, tools } of registrations) {
|
||||
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args)));
|
||||
implementations.set(module, verbs);
|
||||
}
|
||||
|
||||
/** Every verb of every seat a served module claims, where the mesh issued it. */
|
||||
const wanted = (): SeatVerb[] => {
|
||||
const out: SeatVerb[] = [];
|
||||
const have = new Set<string>();
|
||||
const add = (v: SeatVerb): void => {
|
||||
if (have.has(v.subject)) return;
|
||||
have.add(v.subject);
|
||||
out.push(v);
|
||||
};
|
||||
for (const module of served) {
|
||||
const m = typeof broker.membership === "function" ? broker.membership(module) : undefined;
|
||||
for (const s of m?.seats ?? []) add({ seat: s.seat, verb: s.verb, subject: s.subject, holder: module });
|
||||
// The credential's claims, for the module's own runtime: where the mesh issued the verb when
|
||||
// it has; the derived shape until then.
|
||||
if (module !== self) continue;
|
||||
for (const claim of credential?.claims ?? []) {
|
||||
for (const verb of claim.serves ?? []) {
|
||||
const subject = m?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject
|
||||
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
|
||||
add({ seat: claim.seat, verb, subject, holder: module });
|
||||
}
|
||||
}
|
||||
}
|
||||
return out;
|
||||
};
|
||||
|
||||
let stops: (() => void)[] = [];
|
||||
const serve = async (): Promise<void> => {
|
||||
stops.forEach((s) => s());
|
||||
stops = [];
|
||||
const issued = typeof broker.membership === "function" ? broker.membership() : undefined;
|
||||
for (const claim of claims) {
|
||||
const verbs = implementations.get(claim.seat);
|
||||
for (const verb of claim.serves ?? []) {
|
||||
const run = verbs?.get(verb);
|
||||
for (const v of wanted()) {
|
||||
const run = implementations.get(v.seat)?.get(v.verb);
|
||||
if (!run) {
|
||||
console.log(`[mesh-tools] claims ${claim.seat} and implements no ${verb}, which that seat promises; not served`);
|
||||
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
|
||||
continue;
|
||||
}
|
||||
// Where the mesh issued the verb when it has; the derived shape until then.
|
||||
const subject = issued?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject
|
||||
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
|
||||
stops.push(await broker.handleSubject(subject, run));
|
||||
console.log(`[mesh-tools] serving ${claim.seat}'s ${verb} on ${subject}, admitted where this module holds the seat`);
|
||||
}
|
||||
stops.push(await broker.handleSubject(v.subject, run));
|
||||
console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`);
|
||||
}
|
||||
};
|
||||
await serve();
|
||||
// A membership issued to any served module may add, move or withdraw a seat's verbs.
|
||||
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
|
||||
return () => stops.forEach((s) => s());
|
||||
}
|
||||
|
||||
Vendored
+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 }; } },
|
||||
]);
|
||||
Vendored
+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 }) },
|
||||
]);
|
||||
Vendored
+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");
|
||||
@@ -0,0 +1,179 @@
|
||||
/**
|
||||
* 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 three modules' bundles on one credential, names the one that fails, 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"));
|
||||
// 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")] },
|
||||
],
|
||||
});
|
||||
console.log = log;
|
||||
assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "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 5 tool\(s\) for 3 module\(s\): alpha\.one, alpha\.two, beta\.three, beta\.four, beta\.five; 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 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" }] }));
|
||||
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"]);
|
||||
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();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user