diff --git a/src/broker-amqp.ts b/src/broker-amqp.ts index 78c57dd..3bd8e76 100644 --- a/src/broker-amqp.ts +++ b/src/broker-amqp.ts @@ -75,6 +75,10 @@ export async function connectAmqp( async function ensureReply(): Promise { if (replyQueue) return replyQueue; const { queue } = await ch.assertQueue("", { exclusive: true }); + // Replies come back through the RPC exchange keyed by this queue's own name, not the default + // exchange (novox/hq ADR 0052): a serving module's scoped account may write to mesh.rpc but not + // the default exchange, which would let it publish into any queue on the broker. + await ch.bindQueue(queue, RPC_EXCHANGE, queue); replyQueue = queue; await ch.consume( queue, @@ -161,7 +165,9 @@ export async function connectAmqp( reply = { error: err instanceof Error ? err.message : String(err) }; } if (msg.properties.replyTo) { - ch.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { + // Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped + // account answers with write on mesh.rpc alone — never the default exchange (ADR 0052). + ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { correlationId: msg.properties.correlationId, }); } diff --git a/src/main.ts b/src/main.ts index bea038a..769b6a5 100644 --- a/src/main.ts +++ b/src/main.ts @@ -17,6 +17,7 @@ import { readFileSync } from "node:fs"; import { connectAmqp } from "./broker-amqp.js"; import type { Credential } from "./broker-amqp.js"; import { runTools } from "./runtime.js"; +import { invokeTool } from "@novox/mesh-sdk/tools"; import { useBroker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging"; import { emit } from "@novox/mesh-sdk/events"; @@ -92,8 +93,33 @@ async function emitOnce(type: string, bodyJson: string): Promise { await broker.close(); } +async function invokeOnce(module: string, tool: string, argsJson: string): Promise { + let args: Record = {}; + if (argsJson) { + try { + args = JSON.parse(argsJson) as Record; + } catch { + console.error(`mesh-tools invoke: args are not JSON: ${argsJson}`); + process.exit(1); + } + } + const broker = await connectBroker(); + const result = await invokeTool(broker, module, tool, args); + process.stdout.write(JSON.stringify(result) + "\n"); + await broker.close(); +} + async function main(): Promise { const [command, ...rest] = process.argv.slice(2); + if (command === "invoke") { + const [module, tool] = rest; + if (!module || !tool) { + console.error("mesh-tools invoke [json-args] — a module and tool are required"); + process.exit(1); + } + await invokeOnce(module, tool, rest[2] ?? ""); + return; + } if (command === "emit") { const type = rest[0]; if (!type) { diff --git a/test/amqp.test.ts b/test/amqp.test.ts index 205eafb..8a406e4 100644 --- a/test/amqp.test.ts +++ b/test/amqp.test.ts @@ -1,7 +1,7 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { registerModuleTools, resetTools } from "@novox/mesh-sdk/tools"; +import { registerModuleTools, resetTools, invokeTool } from "@novox/mesh-sdk/tools"; import { connectAmqp } from "../dist/broker-amqp.js"; import { runTools } from "../dist/runtime.js"; @@ -27,17 +27,12 @@ test("a module's tool serves and is invoked over a real AMQP broker", { skip: !u const serverBroker = await connectAmqp(url!); const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] }); - // A separate connection — a caller, like mesh-control's command API — invokes over the broker. + // A separate connection — a caller, like mesh-control's command API — invokes over the broker, + // by module and tool (novox/hq ADR 0052: served on serve.demo.greet, invoked as demo.greet). const caller = await connectAmqp(url!); - const result = await caller.request<{ tool: string; args: Record }, { hello: string }>( - "tools.invoke", - { tool: "greet", args: { who: "mesh" } }, - ); + const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string }; assert.equal(result.hello, "mesh"); - // An unknown tool is refused over the wire, not silently dropped. - await assert.rejects(caller.request("tools.invoke", { tool: "nope", args: {} }), /no such tool/); - stop(); await caller.close(); await serverBroker.close();