mesh-tools — the tool runtime and AMQP broker binding
The per-node process that makes a module's tools serve on the mesh: - broker-amqp.ts: a concrete AMQP implementation of the sdk's Broker contract (request/reply over a reply queue + correlation id, publish/ subscribe over a topic exchange). Kept here, not in the sdk, so a broker-client change never rebuilds a module (ADR 0044). - runtime.ts: bind the broker, import the assigned modules' tool entrypoints (each registers as it loads), serveTools. The thin wrapper. - main.ts: the entrypoint, configured by MESH_BROKER_URL + MESH_TOOL_MODULES. - Dockerfile: the container a node runs it as. Verified over a REAL broker: the test spins LavinMQ (the mesh's broker), the runtime serves a registered tool, a separate connection invokes it by name over AMQP and gets the result, and an unknown tool is refused over the wire. So 'modules can serve' is now running-on-the-mesh, not just proven in a mock. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
@@ -0,0 +1,103 @@
|
||||
// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker
|
||||
// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime
|
||||
// provides the binding — so a broker-client change never rebuilds the modules.
|
||||
|
||||
import amqp from "amqplib";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import type { Broker, Envelope } from "@novox/mesh-sdk/messaging";
|
||||
|
||||
const EXCHANGE = "mesh.tools"; // one topic exchange carries tool invocations and events
|
||||
|
||||
interface Reply {
|
||||
result?: unknown;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
/** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */
|
||||
export async function connectAmqp(url: string): Promise<Broker> {
|
||||
const conn = await amqp.connect(url);
|
||||
const ch = await conn.createChannel();
|
||||
await ch.assertExchange(EXCHANGE, "topic", { durable: true });
|
||||
|
||||
// Request/reply: one exclusive reply queue, correlationId → resolver.
|
||||
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
|
||||
const pending = new Map<string, (r: Reply) => void>();
|
||||
await ch.consume(
|
||||
replyQueue,
|
||||
(msg) => {
|
||||
if (!msg) return;
|
||||
const resolve = pending.get(msg.properties.correlationId);
|
||||
if (resolve) {
|
||||
pending.delete(msg.properties.correlationId);
|
||||
resolve(JSON.parse(msg.content.toString()) as Reply);
|
||||
}
|
||||
},
|
||||
{ noAck: true },
|
||||
);
|
||||
|
||||
return {
|
||||
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||
const id = randomUUID();
|
||||
const answered = new Promise<Res>((resolve, reject) => {
|
||||
const timer = setTimeout(() => {
|
||||
if (pending.delete(id)) reject(new Error(`request ${key} timed out`));
|
||||
}, 30_000);
|
||||
pending.set(id, (r) => {
|
||||
clearTimeout(timer);
|
||||
if (r.error) reject(new Error(r.error));
|
||||
else resolve(r.result as Res);
|
||||
});
|
||||
});
|
||||
ch.publish(EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
||||
correlationId: id,
|
||||
replyTo: replyQueue,
|
||||
});
|
||||
return answered;
|
||||
},
|
||||
|
||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
|
||||
await ch.bindQueue(queue, EXCHANGE, key);
|
||||
const consumer = await ch.consume(queue, (msg) => {
|
||||
if (!msg) return;
|
||||
void (async () => {
|
||||
let reply: Reply;
|
||||
try {
|
||||
reply = { result: await handler(JSON.parse(msg.content.toString()) as Req) };
|
||||
} catch (err) {
|
||||
reply = { error: err instanceof Error ? err.message : String(err) };
|
||||
}
|
||||
if (msg.properties.replyTo) {
|
||||
ch.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), {
|
||||
correlationId: msg.properties.correlationId,
|
||||
});
|
||||
}
|
||||
ch.ack(msg);
|
||||
})();
|
||||
});
|
||||
return () => void ch.cancel(consumer.consumerTag);
|
||||
},
|
||||
|
||||
async publish<T>(env: Envelope<T>): Promise<void> {
|
||||
ch.publish(EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body)));
|
||||
},
|
||||
|
||||
async subscribe<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
||||
await ch.bindQueue(queue, EXCHANGE, pattern);
|
||||
const consumer = await ch.consume(queue, (msg) => {
|
||||
if (!msg) return;
|
||||
void (async () => {
|
||||
await handler({ key: msg.fields.routingKey, node: "", body: JSON.parse(msg.content.toString()) as T });
|
||||
ch.ack(msg);
|
||||
})();
|
||||
});
|
||||
return () => void ch.cancel(consumer.consumerTag);
|
||||
},
|
||||
|
||||
async close(): Promise<void> {
|
||||
await ch.close();
|
||||
await conn.close();
|
||||
},
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user