node-tools launches every bundle it serves; a child's emit is published as its module (hq ADR 0193)
Every served entrypoint is started as a process speaking MCP over stdio, told its module and node; one that is not executable is refused by name. The import path, the SDK resolve hook (issue 209) and the per-registration hand-off go. The one-module form the per-module containers use is still imported until they move (to-be 38 WP4c). A child's mesh/publish is published as its module and answered once accepted; a child that dies says why in its own last words. Fixtures are served through launchers exactly as the builder writes them.
This commit is contained in:
@@ -13,6 +13,8 @@
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { accessSync, constants } from "node:fs";
|
||||
import type { ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { broker, type Envelope } from "@novox/mesh-sdk/messaging";
|
||||
import { atWork } from "./broker-nats.js";
|
||||
|
||||
/** The protocol version this speaks; a bundle says the same. */
|
||||
export const PROTOCOL = "2025-03-26";
|
||||
@@ -71,29 +73,50 @@ export async function launch(module: string, entry: string, env: NodeJS.ProcessE
|
||||
const line = buffered.slice(0, at).trim();
|
||||
buffered = buffered.slice(at + 1);
|
||||
if (!line) continue;
|
||||
let reply: { id?: number; result?: unknown; error?: { message?: string } };
|
||||
let reply: { id?: number | string; method?: string; params?: unknown; result?: unknown; error?: { message?: string } };
|
||||
try {
|
||||
reply = JSON.parse(line);
|
||||
} catch {
|
||||
console.log(`[mesh-tools] ${module}'s bundle said something that is not a reply: ${line.slice(0, 120)}`);
|
||||
continue;
|
||||
}
|
||||
// **The bundle asks the runtime to emit** (novox/hq ADR 0193): published on the bus as this
|
||||
// module, and answered once the bus has accepted it, so the tool's emit means what it means
|
||||
// in-process. Nothing else a bundle may ask.
|
||||
if (typeof reply.method === "string") {
|
||||
const id = reply.id;
|
||||
const answer = (m: Record<string, unknown>) => proc.stdin!.write(JSON.stringify({ jsonrpc: "2.0", id, ...m }) + "\n");
|
||||
if (reply.method !== "mesh/publish") {
|
||||
if (id !== undefined) answer({ error: { code: -32601, message: `the runtime answers no ${reply.method} from a bundle` } });
|
||||
continue;
|
||||
}
|
||||
atWork.run({ module }, () => broker().publish(reply.params as Envelope<unknown>))
|
||||
.then(() => { if (id !== undefined) answer({ result: {} }); })
|
||||
.catch((err: unknown) => { if (id !== undefined) answer({ error: { code: -32000, message: err instanceof Error ? err.message : String(err) } }); });
|
||||
continue;
|
||||
}
|
||||
const waiting = typeof reply.id === "number" ? pending.get(reply.id) : undefined;
|
||||
if (!waiting) continue;
|
||||
pending.delete(reply.id!);
|
||||
pending.delete(reply.id as number);
|
||||
clearTimeout(waiting.timer);
|
||||
if (reply.error) waiting.reject(new Error(reply.error.message ?? "the bundle refused the request"));
|
||||
else waiting.resolve(reply.result);
|
||||
}
|
||||
});
|
||||
// stderr is the bundle's log; kept under the module's name so a fault reads where it belongs.
|
||||
// The last thing it said is kept, so a bundle that dies says why in its own words, not by code.
|
||||
let lastSaid = "";
|
||||
proc.stderr!.on("data", (chunk: Buffer) => {
|
||||
for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) console.log(`[${module}] ${line}`);
|
||||
for (const line of chunk.toString("utf8").split("\n")) {
|
||||
if (!line.trim()) continue;
|
||||
console.log(`[${module}] ${line}`);
|
||||
if (/\S/.test(line) && !/^\s+at\s/.test(line) && !/^Node\.js v/.test(line)) lastSaid = line.trim();
|
||||
}
|
||||
});
|
||||
const exited = new Promise<never>((_, reject) => {
|
||||
proc.once("error", (err) => reject(err));
|
||||
proc.once("exit", (code, signal) => {
|
||||
const why = `${module}'s bundle exited (${signal ?? code})`;
|
||||
const why = `${module}'s bundle exited (${signal ?? code})` + (lastSaid ? `: ${lastSaid}` : "");
|
||||
for (const [id, p] of pending) {
|
||||
pending.delete(id);
|
||||
clearTimeout(p.timer);
|
||||
|
||||
Reference in New Issue
Block a user