diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index 930b8ac..b31a5eb 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -6,7 +6,11 @@ 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 +// Two topic exchanges, kept apart on purpose: tool invocations are request/reply and are not +// events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node +// events — never the RPC traffic. +const RPC_EXCHANGE = "mesh.rpc"; +const EVENTS_EXCHANGE = "mesh.events"; interface Reply { result?: unknown; @@ -17,7 +21,8 @@ interface Reply { export async function connectAmqp(url: string): Promise { const conn = await amqp.connect(url); const ch = await conn.createChannel(); - await ch.assertExchange(EXCHANGE, "topic", { durable: true }); + await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true }); + await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true }); // Request/reply: one exclusive reply queue, correlationId → resolver. const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true }); @@ -48,7 +53,7 @@ export async function connectAmqp(url: string): Promise { else resolve(r.result as Res); }); }); - ch.publish(EXCHANGE, key, Buffer.from(JSON.stringify(body)), { + ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), { correlationId: id, replyTo: replyQueue, }); @@ -57,7 +62,7 @@ export async function connectAmqp(url: string): Promise { async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true }); - await ch.bindQueue(queue, EXCHANGE, key); + await ch.bindQueue(queue, RPC_EXCHANGE, key); const consumer = await ch.consume(queue, (msg) => { if (!msg) return; void (async () => { @@ -79,12 +84,12 @@ export async function connectAmqp(url: string): Promise { }, async publish(env: Envelope): Promise { - ch.publish(EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body))); + ch.publish(EVENTS_EXCHANGE, env.key, Buffer.from(JSON.stringify(env.body))); }, async subscribe(pattern: string, handler: (env: Envelope) => Promise): Promise<() => void> { const { queue } = await ch.assertQueue("", { exclusive: true }); - await ch.bindQueue(queue, EXCHANGE, pattern); + await ch.bindQueue(queue, EVENTS_EXCHANGE, pattern); const consumer = await ch.consume(queue, (msg) => { if (!msg) return; void (async () => {