broker: split events (mesh.events) from tool RPC (mesh.rpc)
Per ADR 0047: two topic exchanges kept apart, so a '#' subscription on mesh.events is a clean audit of module/mesh/node events without the request/reply tool traffic. Header-in-AMQP, persistent messages, per- consumer durable dead-lettered queues and prefetch are the remaining alignment to 0047. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
+11
-6
@@ -6,7 +6,11 @@ import amqp from "amqplib";
|
|||||||
import { randomUUID } from "node:crypto";
|
import { randomUUID } from "node:crypto";
|
||||||
import type { Broker, Envelope } from "@novox/mesh-sdk/messaging";
|
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 {
|
interface Reply {
|
||||||
result?: unknown;
|
result?: unknown;
|
||||||
@@ -17,7 +21,8 @@ interface Reply {
|
|||||||
export async function connectAmqp(url: string): Promise<Broker> {
|
export async function connectAmqp(url: string): Promise<Broker> {
|
||||||
const conn = await amqp.connect(url);
|
const conn = await amqp.connect(url);
|
||||||
const ch = await conn.createChannel();
|
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.
|
// Request/reply: one exclusive reply queue, correlationId → resolver.
|
||||||
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
|
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
|
||||||
@@ -48,7 +53,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
else resolve(r.result as Res);
|
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,
|
correlationId: id,
|
||||||
replyTo: replyQueue,
|
replyTo: replyQueue,
|
||||||
});
|
});
|
||||||
@@ -57,7 +62,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
|
|
||||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||||
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
|
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) => {
|
const consumer = await ch.consume(queue, (msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
void (async () => {
|
void (async () => {
|
||||||
@@ -79,12 +84,12 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
},
|
},
|
||||||
|
|
||||||
async publish<T>(env: Envelope<T>): Promise<void> {
|
async publish<T>(env: Envelope<T>): Promise<void> {
|
||||||
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<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
async subscribe<T>(pattern: string, handler: (env: Envelope<T>) => Promise<void>): Promise<() => void> {
|
||||||
const { queue } = await ch.assertQueue("", { exclusive: true });
|
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) => {
|
const consumer = await ch.consume(queue, (msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
void (async () => {
|
void (async () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user