runtime: an emit primitive, so an events test can put a message on the wire
'mesh-tools emit <type> [json]' connects, emits one ADR 0047 event (awaiting
the publish confirm), and exits. The serve path already runs a module's
on('#') subscription as an import side effect, so the runtime hosts both an
emitter and the audit-logger consumer.
Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
+45
-5
@@ -1,13 +1,21 @@
|
||||
// The runnable entrypoint. Reads its configuration from the environment the host resolved for it,
|
||||
// connects the mesh broker, and serves the assigned modules' tools until stopped.
|
||||
// The runnable entrypoint. Two modes:
|
||||
//
|
||||
// MESH_BROKER_URL amqp://… the mesh broker
|
||||
// MESH_TOOL_MODULES /path/a,/path/b,… compiled tool entrypoints of the assigned modules
|
||||
// mesh-tools serve — bind the broker and serve the assigned modules until
|
||||
// stopped. 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_BROKER_URL amqp://… the mesh broker (both modes)
|
||||
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
|
||||
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047)
|
||||
|
||||
import { connectAmqp } from "./broker-amqp.js";
|
||||
import { runTools } from "./runtime.js";
|
||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||
import { emit } from "@novox/mesh-sdk/events";
|
||||
|
||||
async function main(): Promise<void> {
|
||||
async function serve(): Promise<void> {
|
||||
const url = requireEnv("MESH_BROKER_URL");
|
||||
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
||||
.split(",")
|
||||
@@ -26,6 +34,38 @@ async function main(): Promise<void> {
|
||||
process.on("SIGINT", () => void shutdown());
|
||||
}
|
||||
|
||||
async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
||||
const url = requireEnv("MESH_BROKER_URL");
|
||||
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 connectAmqp(url);
|
||||
useBroker(() => broker);
|
||||
// emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close.
|
||||
await emit(type, body);
|
||||
await broker.close();
|
||||
}
|
||||
|
||||
async function main(): Promise<void> {
|
||||
const [command, ...rest] = process.argv.slice(2);
|
||||
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();
|
||||
}
|
||||
|
||||
function requireEnv(name: string): string {
|
||||
const v = process.env[name];
|
||||
if (!v) {
|
||||
|
||||
Reference in New Issue
Block a user