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'.
310 lines
14 KiB
TypeScript
310 lines
14 KiB
TypeScript
// The runnable entrypoint. Three modes:
|
|
//
|
|
// mesh-tools serve — bind the broker and serve the assigned modules until
|
|
// stopped; as the node-tools module, also the console on loopback. A module entrypoint that subscribes to events (on("#"))
|
|
// starts consuming as it is imported, so this also runs consumers.
|
|
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
|
|
// and what an events test uses to put a message on the wire.
|
|
// mesh-tools prepare bring this module's state to the shape this version needs and exit
|
|
// — the runtime's answer to the word the mesh asks every module
|
|
// (novox/hq ADR 0135). The entrypoints come from MESH_PREPARE, which
|
|
// the module's own image names beside MESH_TOOL_MODULES.
|
|
// mesh-tools run ENTRYPOINT run one compiled module entrypoint to completion and exit — the
|
|
// runtime side of a run-once step (novox/hq ADR 0052). It imports
|
|
// the given entrypoint, whose top-level code does its work — seed a
|
|
// store, migrate, health-gate — and awaits it. It does NOT connect
|
|
// to the broker: a first-boot step runs offline, before the module
|
|
// has anything to talk to, and the host gates the container that
|
|
// depends on it on this process exiting 0.
|
|
//
|
|
// The broker, in order of preference:
|
|
// 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 <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 { takeToolEnvs, runTools, type ServedModule } from "./runtime.js";
|
|
import { serveMcpHttp, type Listening } from "./http.js";
|
|
|
|
/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */
|
|
let lastCredential: Credential | undefined;
|
|
import { invokeTool } from "@novox/mesh-sdk/tools";
|
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
|
import { emit } from "@novox/mesh-sdk/events";
|
|
|
|
/**
|
|
* Connect the way this process is meant to: with its sealed credential if the mesh gave it one, and
|
|
* over the plain bootstrap URL otherwise. A scoped module assumes the foundation's exchanges exist —
|
|
* its account may not declare them (ADR 0043).
|
|
*/
|
|
// fatalBrokerReasonFor is the reason a connection failure is final rather than "not yet", for
|
|
// whichever bus this runtime is on — each transport knows its own refusals.
|
|
function fatalBrokerReasonFor(err: unknown): string | null {
|
|
return fatalNatsReason(err);
|
|
}
|
|
|
|
async function connectBroker(): Promise<Broker> {
|
|
const file = process.env.MESH_BROKER_FILE;
|
|
if (file) {
|
|
let credential: Credential;
|
|
try {
|
|
credential = JSON.parse(readFileSync(file, "utf8")) as Credential;
|
|
lastCredential = credential;
|
|
} catch (err) {
|
|
console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`);
|
|
process.exit(1);
|
|
}
|
|
if (!credential.url) {
|
|
console.error(`mesh-tools: ${file} carries no url — it is not a broker credential`);
|
|
process.exit(1);
|
|
}
|
|
// The mesh scoped this account to a node and module; take the runtime's identity from the
|
|
// credential so its queue and the events it emits match what the mesh authorised, no matter
|
|
// what the environment says.
|
|
if (credential.node) process.env.MESH_NODE = credential.node;
|
|
if (credential.module) process.env.MESH_MODULE = credential.module;
|
|
// **The credential names the bus.** A module moved to the bus being built was handed a
|
|
// credential for it — `nats://…` with user, password and fingerprint beside the address — and
|
|
// nothing else in its environment changed (design 25; novox/hq design 28 task 5.2). The scheme
|
|
// is enough to know which bus to speak; a runtime that always dialled the old one would keep
|
|
// serving and answer nobody.
|
|
// One bus (novox/hq ADR 0131, design 28 task 5.5): the credential names it, and it is this.
|
|
return connectNats(credential);
|
|
}
|
|
const url = process.env.MESH_BROKER_URL;
|
|
if (!url) {
|
|
console.error(
|
|
"mesh-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach",
|
|
);
|
|
process.exit(1);
|
|
}
|
|
return connectNats({ url });
|
|
}
|
|
|
|
/**
|
|
* Connect for serve mode, retrying while the broker is merely not reachable yet. At startup that
|
|
* is the NORMAL case, not a failure: a container comes up in seconds and the overlay tunnel a
|
|
* moment later (novox/hq issue 058). Exiting instead delegated the retry to the container
|
|
* runtime, which read as a crash-loop to every restart-counting health check and every person
|
|
* watching. Retried indefinitely, aloud: the dependency appears or somebody reads why not.
|
|
*
|
|
* A failure that waiting cannot fix (see fatalBrokerReason) is thrown at once rather than retried —
|
|
* a permanent fault masquerading as "not reachable yet" is the silent non-progress this whole
|
|
* change exists to remove. Missing-file and empty-URL configuration errors exit inside
|
|
* connectBroker before they reach here; a malformed URL and a refused login are caught here.
|
|
*/
|
|
async function connectBrokerPatiently(): Promise<Broker> {
|
|
for (let delay = 2_000; ; delay = Math.min(delay * 2, 30_000)) {
|
|
try {
|
|
return await connectBroker();
|
|
} catch (err) {
|
|
const fatal = fatalBrokerReasonFor(err);
|
|
if (fatal !== null) {
|
|
console.error(`mesh-tools: ${fatal} — waiting will not fix this; giving up`);
|
|
throw err;
|
|
}
|
|
const why = err instanceof Error ? err.message : String(err);
|
|
// A little jitter so every module that was up when the broker bounced does not retry in
|
|
// lockstep and stampede it as it recovers.
|
|
const wait = delay + Math.floor(Math.random() * 1_000);
|
|
console.error(`mesh-tools: the broker is not reachable yet (${why}); retrying in ${Math.round(wait / 1000)}s`);
|
|
await new Promise((r) => setTimeout(r, wait));
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* 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 };
|
|
}
|
|
|
|
/** The module that is the node's tool runtime (novox/hq ADR 0175, to-be 38 WP3): on its credential,
|
|
* `serve` is also the console — MCP on the machine's loopback (design 34). */
|
|
export const RUNTIME_MODULE = "node-tools";
|
|
|
|
/** Where the console listens when the runtime is node-tools and nothing says otherwise: the port
|
|
* the module's manifest declares `from: machine`. MESH_CONSOLE_LISTEN overrides it either way. */
|
|
const CONSOLE_LISTEN = "127.0.0.1:4270";
|
|
|
|
async function serve(): Promise<void> {
|
|
const broker = await connectBrokerPatiently();
|
|
// Parsed after connecting: a bare entrypoint belongs to the module the credential names.
|
|
const { serves, moduleEntrypoints } = servedModulesFrom(process.env.MESH_TOOL_MODULES ?? "", lastCredential?.module);
|
|
// Each module's environment, composed by the mesh (ADR 0192): taken before any bundle is imported.
|
|
const envs = takeToolEnvs();
|
|
const stop = await runTools({ broker, serves, moduleEntrypoints, credential: lastCredential, envs });
|
|
|
|
// The console is this runtime's serving mode (ADR 0175 §6): as node-tools, or wherever the
|
|
// listen address is given, the same process answers MCP on loopback for whoever is on the
|
|
// machine. A module's own runtime in a container on the machine's network does not — two of
|
|
// them on one port would be the fault, and the console is one per machine.
|
|
const listen = process.env.MESH_CONSOLE_LISTEN ?? (lastCredential?.module === RUNTIME_MODULE ? CONSOLE_LISTEN : "");
|
|
let consoleUp: Listening | undefined;
|
|
if (listen) {
|
|
const who = `${lastCredential?.node ?? "?"}.${lastCredential?.module ?? RUNTIME_MODULE}`;
|
|
consoleUp = await serveMcpHttp(broker, who, listen);
|
|
console.log(`mesh console listening on http://${consoleUp.address}/mcp as ${who}`);
|
|
}
|
|
|
|
const shutdown = async (): Promise<void> => {
|
|
stop();
|
|
await consoleUp?.close();
|
|
await broker.close();
|
|
process.exit(0);
|
|
};
|
|
process.on("SIGTERM", () => void shutdown());
|
|
process.on("SIGINT", () => void shutdown());
|
|
}
|
|
|
|
async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
|
let body: unknown = {};
|
|
if (bodyJson) {
|
|
try {
|
|
body = JSON.parse(bodyJson);
|
|
} catch {
|
|
console.error(`mesh-tools emit: body is not JSON: ${bodyJson}`);
|
|
process.exit(1);
|
|
}
|
|
}
|
|
const broker = await connectBroker();
|
|
useBroker(() => broker);
|
|
// emit awaits the broker's publish confirm (ADR 0042), so the event is accepted before we close.
|
|
await emit(type, body);
|
|
await broker.close();
|
|
}
|
|
|
|
async function invokeOnce(module: string, tool: string, argsJson: string): Promise<void> {
|
|
let args: Record<string, unknown> = {};
|
|
if (argsJson) {
|
|
try {
|
|
args = JSON.parse(argsJson) as Record<string, unknown>;
|
|
} catch {
|
|
console.error(`mesh-tools invoke: args are not JSON: ${argsJson}`);
|
|
process.exit(1);
|
|
}
|
|
}
|
|
const broker = await connectBroker();
|
|
const result = await invokeTool(broker, module, tool, args);
|
|
process.stdout.write(JSON.stringify(result) + "\n");
|
|
await broker.close();
|
|
}
|
|
|
|
/**
|
|
* Run one compiled module entrypoint to completion — the runtime side of a run-once step
|
|
* (novox/hq ADR 0052). Importing it runs its top-level code and awaits any top-level await, so this
|
|
* returns only once the step's own code has finished; a step that throws rejects here and the
|
|
* process exits non-zero, which is how the host knows the step did not complete and must not start
|
|
* the container it gates. No broker is connected — a first-boot seed or migration runs offline.
|
|
*/
|
|
async function runEntry(entrypoint: string): Promise<void> {
|
|
if (!entrypoint) {
|
|
console.error("mesh-tools run <entrypoint> — a compiled module entrypoint path is required");
|
|
process.exit(1);
|
|
}
|
|
// A file URL, not a bare path: dynamic import of an absolute path is not portable, and the
|
|
// entrypoint the manifest names is an absolute path inside the image.
|
|
await import(pathToFileURL(entrypoint).href);
|
|
}
|
|
|
|
/**
|
|
* Bring this module's state to the shape this version needs, and exit — the runtime's answer to the
|
|
* one word the mesh asks every module (novox/hq ADR 0135).
|
|
*
|
|
* The entrypoints come from `MESH_PREPARE`, which a module's own image sets beside the entrypoints it
|
|
* already lists there: the module knows which of its files prepares its state, and nothing else could.
|
|
* Each is imported in the order given, to completion, with no broker — preparation runs before the
|
|
* version that would use it, so there is nothing yet to talk to.
|
|
*
|
|
* **An empty list is a failure, not a no-op.** The mesh only asks this of a module whose manifest says
|
|
* it prepares something; a module that says so and names nothing has been built wrong, and exiting 0
|
|
* would let that version serve against a state nobody shaped.
|
|
*/
|
|
async function prepareState(): Promise<void> {
|
|
const named = (process.env.MESH_PREPARE ?? "")
|
|
.split(",")
|
|
.map((entry) => entry.trim())
|
|
.filter((entry) => entry !== "");
|
|
if (named.length === 0) {
|
|
console.error(
|
|
"mesh-tools prepare: this module was asked to prepare its state and its image names nothing " +
|
|
"to do it with — set MESH_PREPARE to the compiled entrypoint(s) that prepare it, the way " +
|
|
"MESH_TOOL_MODULES names the ones it serves",
|
|
);
|
|
process.exit(1);
|
|
}
|
|
for (const entrypoint of named) {
|
|
console.log(`[mesh-tools] preparing with ${entrypoint}`);
|
|
await runEntry(entrypoint);
|
|
}
|
|
console.log(`[mesh-tools] prepared: ${named.length} entrypoint(s) ran to completion`);
|
|
}
|
|
|
|
async function main(): Promise<void> {
|
|
const [command, ...rest] = process.argv.slice(2);
|
|
if (command === "run") {
|
|
await runEntry(rest[0] ?? "");
|
|
return;
|
|
}
|
|
if (command === "prepare") {
|
|
await prepareState();
|
|
return;
|
|
}
|
|
if (command === "invoke") {
|
|
const [module, tool] = rest;
|
|
if (!module || !tool) {
|
|
console.error("mesh-tools invoke <module> <tool> [json-args] — a module and tool are required");
|
|
process.exit(1);
|
|
}
|
|
await invokeOnce(module, tool, rest[2] ?? "");
|
|
return;
|
|
}
|
|
if (command === "emit") {
|
|
const type = rest[0];
|
|
if (!type) {
|
|
console.error("mesh-tools emit <type> [json-body] — a routing key is required");
|
|
process.exit(1);
|
|
}
|
|
await emitOnce(type, rest[1] ?? "");
|
|
return;
|
|
}
|
|
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();
|
|
}
|