diff --git a/src/broker-nats.ts b/src/broker-nats.ts index 3650310..198817e 100644 --- a/src/broker-nats.ts +++ b/src/broker-nats.ts @@ -40,6 +40,25 @@ export interface Credential { module?: string; user?: string; password?: string; + /** The seats this module claims, with the verbs each promises (novox/hq ADR 0159). The runtime + * serves each claimed seat's verbs with its tools of the same name; the bus admits only the + * holder's subscription, so claiming and not holding costs a refused subscription and nothing else. */ + claims?: { seat: string; scope?: string; serves?: string[] }[]; +} + +/** What a tool call answers: the module's own result, and which machine answered it + * (novox/hq ADR 0159) — a module on several machines is otherwise an answer from nowhere. */ +export interface Answered { + result: Res; + node?: string; +} + +/** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs — + * an answer that says which machine gave it, and serving a subject that is not a module's own tool + * (a seat's verb). */ +export interface RuntimeBroker extends Broker { + ask(key: string, body: Req): Promise>; + handleSubject(subject: string, handler: (body: Req) => Promise): Promise<() => void>; } /** Whether a connection failure is worth retrying, or is a fact about this configuration that @@ -69,7 +88,7 @@ export function fatalBrokerReason(err: unknown): string | null { export async function connectNats( target: string | Credential, opts: { module?: string } = {}, -): Promise { +): Promise { const cred: Credential = typeof target === "string" ? { url: target } : target; const self = cred.module ?? opts.module; if (!self) { @@ -102,49 +121,72 @@ export async function connectNats( let reading: Awaited>["consume"]>> | undefined; let closed = false; + const node = cred.node; + + /** Answer one subject with one handler, and say which machine answered (novox/hq ADR 0159). */ + const answerOn = ( + subject: string, + queue: string | undefined, + handler: (body: Req) => Promise, + ): (() => void) => { + const sub = conn.subscribe(subject, queue ? { queue } : {}); + subs.push(sub); + void (async () => { + for await (const msg of sub) { + let reply: { result?: Res; error?: string; node?: string }; + try { + reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) }; + } catch (err) { + // The caller is told, rather than left to time out: a handler that threw is a + // different failure from a tool nobody serves, and only one of them is worth retrying. + reply = { error: err instanceof Error ? err.message : String(err) }; + } + if (node) reply.node = node; + msg.respond(sc.encode(JSON.stringify(reply))); + } + })(); + return () => sub.unsubscribe(); + }; + + /** + * Ask one question and await one answer, with the machine that gave it. + * + * Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3), + * and a lost one is a timeout the caller already handles. The reply travels on the inbox the + * request carries, which the responder may answer because its account has `allow_responses` + * — one reply to a message it actually received, and nothing wider. + */ + const ask = async (key: string, body: Req): Promise> => { + const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), { + timeout: REQUEST_TIMEOUT_MS, + }); + const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string; node?: string }; + if (reply.error) throw new Error(reply.error); + return { result: reply.result as Res, node: reply.node }; + }; + return { - /** - * Ask one question and await one answer. - * - * Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3), - * and a lost one is a timeout the caller already handles. The reply travels on the inbox the - * request carries, which the responder may answer because its account has `allow_responses` - * — one reply to a message it actually received, and nothing wider. - */ async request(key: string, body: Req): Promise { - const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), { - timeout: REQUEST_TIMEOUT_MS, - }); - const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string }; - if (reply.error) throw new Error(reply.error); - return reply.result as Res; + return (await ask(key, body)).result; }, + ask, + /** - * Answer a question. - * - * A queue group, so several nodes may serve one tool and exactly one of them answers each - * call. + * Answer a question, two ways (novox/hq ADR 0159): on the module's subject in a queue group, + * so several machines may serve one tool and exactly one of them answers each call; and on the + * same subject with this machine as its last token, so a caller that names the machine reaches + * this instance and no other. A runtime that does not know its machine serves only the first, + * which is how it always behaved. */ async handle(key: string, handler: (body: Req) => Promise): Promise<() => void> { - const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` }); - subs.push(sub); - void (async () => { - for await (const msg of sub) { - let reply: { result?: Res; error?: string }; - try { - reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) }; - } catch (err) { - // The caller is told, rather than left to time out: a handler that threw is a - // different failure from a tool nobody serves, and only one of them is worth retrying. - reply = { error: err instanceof Error ? err.message : String(err) }; - } - msg.respond(sc.encode(JSON.stringify(reply))); - } - })(); - return () => { - sub.unsubscribe(); - }; + const stops = [answerOn(toolSubject(key, self), `serve.${self}`, handler)]; + if (node) stops.push(answerOn(`${toolSubject(key, self)}.${node}`, undefined, handler)); + return () => stops.forEach((stop) => stop()); + }, + + async handleSubject(subject: string, handler: (body: Req) => Promise): Promise<() => void> { + return answerOn(subject, undefined, handler); }, /** @@ -325,9 +367,21 @@ function toolSubject(key: string, self: string): string { const [verb, node] = rest.slice(dot + 1).split("@", 2); return node ? `mesh.seat.${seat}.tool.${verb}.${node}` : `mesh.seat.${seat}.tool.${verb}`; } - const dot = key.indexOf("."); - if (dot < 0) return `mesh.mod.${self}.tool.${key}`; - return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`; + // `.@` names the machine (novox/hq ADR 0159): the same subject with the + // machine as its last token, which is what that instance serves beside the queue. + const [name, node] = key.split("@", 2); + const dot = name.indexOf("."); + const base = dot < 0 + ? `mesh.mod.${self}.tool.${name}` + : `mesh.mod.${name.slice(0, dot)}.tool.${name.slice(dot + 1)}`; + return node ? `${base}.${node}` : base; +} + +/** A seat's verb, as its holder serves it: flat for a mesh seat, carrying the machine for a + * node-scoped one (design 33 §4) — the same shape the controller grants. */ +export function seatToolSubject(seat: string, verb: string, scope: string | undefined, node: string | undefined): string { + const base = `mesh.seat.${seat}.tool.${verb}`; + return scope === "node" && node ? `${base}.${node}` : base; } function normalizeFingerprint(fingerprint: string): string { diff --git a/src/client.ts b/src/client.ts index e46640c..af1ea90 100644 --- a/src/client.ts +++ b/src/client.ts @@ -20,7 +20,7 @@ import { readFile } from "node:fs/promises"; import type { Broker } from "@novox/mesh-sdk/messaging"; -import { connectNats, type Credential } from "./broker-nats.js"; +import { connectNats, type Answered, type Credential } from "./broker-nats.js"; import { TOOLS_VERB, type ToolsAnswer } from "./runtime.js"; /** Where the catalogue answers which modules the mesh holds. */ @@ -201,13 +201,23 @@ export async function toolsOn(bus: Broker): Promise { /** Call one tool. The key is `.`, which is what a person types and what the account * permits — one vocabulary, so a refusal names the thing they asked for. */ -export async function callTool(bus: Broker, key: string, args: unknown, seats?: Seats): Promise { - if (!key.includes(".")) { +/** Call a tool and learn which machine answered (novox/hq ADR 0159). `.@` asks + * the instance on one machine; without it, whichever instance answers first does, and the answer + * says which. */ +export async function callTool(bus: Broker, key: string, args: unknown, seats?: Seats): Promise> { + const at = key.indexOf("@"); + const name = at < 0 ? key : key.slice(0, at); + const node = at < 0 ? "" : key.slice(at + 1); + if (!name.includes(".")) { throw new Error( - `"${key}" does not name a tool: write ., as \`mesh tools\` lists them`, + `"${key}" does not name a tool: write ., as \`mesh tools\` lists them, ` + + "or .@ for the instance on one machine", ); } - return bus.request(toolKey(key, seats), args ?? {}); + const resolved = toolKey(name, seats) + (node ? `@${node}` : ""); + const asking = bus as Broker & { ask?: (k: string, b: Req) => Promise> }; + if (typeof asking.ask === "function") return asking.ask(resolved, args ?? {}); + return { result: await bus.request(resolved, args ?? {}) }; } /** diff --git a/src/main.ts b/src/main.ts index 22671f9..b81bba3 100644 --- a/src/main.ts +++ b/src/main.ts @@ -29,6 +29,9 @@ import { readFileSync } from "node:fs"; import { pathToFileURL } from "node:url"; import { connectNats, fatalBrokerReason as fatalNatsReason, type Credential } from "./broker-nats.js"; import { runTools } from "./runtime.js"; + +/** The credential this process connected with, for what it says beyond the connection (ADR 0159). */ +let lastCredential: Credential | undefined; import { invokeTool } from "@novox/mesh-sdk/tools"; import { useBroker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging"; @@ -51,6 +54,7 @@ async function connectBroker(): Promise { let credential: Credential; try { credential = JSON.parse(readFileSync(file, "utf8")) as Credential; + lastCredential = credential; } catch (err) { console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`); process.exit(1); @@ -121,7 +125,7 @@ async function serve(): Promise { .filter(Boolean); const broker = await connectBrokerPatiently(); - const stop = await runTools({ broker, moduleEntrypoints }); + const stop = await runTools({ broker, moduleEntrypoints, credential: lastCredential }); const shutdown = async (): Promise => { stop(); diff --git a/src/mcp.ts b/src/mcp.ts index 2456bbe..71ce922 100644 --- a/src/mcp.ts +++ b/src/mcp.ts @@ -118,9 +118,11 @@ export function mcpSurface(bus: Broker, who: string): Surface { tools: have.tools.map((t) => ({ name: `${t.module}.${t.name}`, description: t.description ?? `${t.name}, served by ${t.module}`, - // The module's own schema, passed through. An empty object is a tool that takes - // nothing, which is a real answer and not a missing one. - inputSchema: asSchema(t.input), + // The module's own schema, passed through — with `node`, the machine to ask when + // the module runs on several (novox/hq ADR 0159); a seat's verb takes none, the + // seat's scope decides. An empty object is a tool that takes nothing, which is a + // real answer and not a missing one. + inputSchema: t.seat ? asSchema(t.input) : withNode(asSchema(t.input)), })), // Silence, named (design 34 §3): the modules the catalogue holds and nothing answered // for. Not a tool, so not in `tools`; not dropped either. @@ -129,16 +131,23 @@ export function mcpSurface(bus: Broker, who: string): Surface { } case "tools/call": { - const name = String(request.params?.name ?? ""); - const args = request.params?.arguments ?? {}; + const given = String(request.params?.name ?? ""); + const args = { ...((request.params?.arguments as Record | undefined) ?? {}) }; + // The machine, when the caller names one, travels in the subject and never reaches the + // module's arguments (novox/hq ADR 0159). + const node = typeof args.node === "string" && args.node !== "" ? args.node : ""; + delete args.node; + const name = node && !given.includes("@") ? `${given}@${node}` : given; try { - const result = await callTool(bus, name, args, await roles()); + const { result, node: answeredBy } = await callTool(bus, name, args, await roles()); // Text, because that is what every host renders. The content is the module's answer // as JSON, unshaped: an adapter that flattened it would be deciding what matters in - // somebody else's answer. - return answer(request.id, { - content: [{ type: "text", text: JSON.stringify(result, null, 2) }], - }); + // somebody else's answer. Which machine answered follows it as its own line. + const content: { type: string; text: string }[] = [ + { type: "text", text: JSON.stringify(result, null, 2) }, + ]; + if (answeredBy) content.push({ type: "text", text: `answered by ${answeredBy}` }); + return answer(request.id, { content }); } catch (e) { // **An error the agent can act on, not a stack.** isError rather than a protocol // failure, because the call was well-formed and the mesh answered it — with a refusal, @@ -212,3 +221,16 @@ async function* lines(): AsyncGenerator { } if (buffered.trim() !== "") yield buffered.trim(); } + +/** Every module tool takes an optional `node`: the machine to ask when the module runs on several + * (novox/hq ADR 0159). Added to the listing, stripped before the call, never seen by the module. */ +function withNode(schema: Record): Record { + const properties = { ...((schema.properties as Record | undefined) ?? {}) }; + if (!("node" in properties)) { + properties.node = { + type: "string", + description: "the machine to ask, when this module runs on several; else whichever answers, and the answer says which", + }; + } + return { ...schema, type: "object", properties }; +} diff --git a/src/mesh.ts b/src/mesh.ts index 1361752..d6b9313 100644 --- a/src/mesh.ts +++ b/src/mesh.ts @@ -174,8 +174,9 @@ async function calling(bus: Broker, args: string[]): Promise { // `.` reaches the role when the mesh lists that verb for the seat; `seat:` says so // outright and asks nothing first. const roles = key.startsWith("seat:") ? undefined : seatsIn(await toolsOn(bus).catch(() => ({ tools: [], notAnswering: [] }))); - const answer = await callTool(bus, key, parsed, roles); - console.log(JSON.stringify(answer, null, 2)); + const { result, node } = await callTool(bus, key, parsed, roles); + console.log(JSON.stringify(result, null, 2)); + if (node) console.error(`answered by ${node}`); return 0; } catch (e) { console.error(whyItFailed(key, e)); diff --git a/src/runtime.ts b/src/runtime.ts index 70c0f92..de2846f 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -8,6 +8,7 @@ import { resolve } from "node:path"; import { useBroker } from "@novox/mesh-sdk/messaging"; import { collectTools, serveTools, listTools, toolKey } from "@novox/mesh-sdk/tools"; import type { Broker } from "@novox/mesh-sdk/messaging"; +import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js"; /** * The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the @@ -27,6 +28,9 @@ export interface RuntimeOptions { broker: Broker; /** Absolute paths to the assigned modules' compiled tool entrypoints (e.g. .../umami/tools/index.js). */ moduleEntrypoints: string[]; + /** The credential the mesh delivered, for what it says about the seats this module claims + * (novox/hq ADR 0159). Absent for a runtime started by hand, which then serves no seat. */ + credential?: Credential; } /** Load the modules, bind the broker, and serve. Returns a stop function that unhooks serving. */ @@ -65,7 +69,38 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> { } console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`); + stops.push(await serveClaimedSeats(opts.broker as RuntimeBroker, opts.credential)); return () => { for (const s of stops) s(); }; } + +/** + * Holding a seat means serving its tools (design 33 §3, novox/hq ADR 0159). The credential names the + * seats this module claims and the verbs each promises; each verb is served on the seat's own + * subject by the module's tool of the same name. Whether this instance *holds* the seat is the bus's + * to decide: only the holder's account may subscribe the seat's subjects, so a claimant that does not + * hold it here is refused the subscription and serves nothing — never a failure of its own tools. + */ +async function serveClaimedSeats(broker: RuntimeBroker, credential?: Credential): Promise<() => void> { + const claims = credential?.claims ?? []; + if (claims.length === 0 || typeof broker.handleSubject !== "function") return () => {}; + const byName = new Map) => Promise>(); + for (const { tools } of collectTools()) { + for (const t of tools) byName.set(t.name, (args) => t.run(args)); + } + const stops: (() => void)[] = []; + for (const claim of claims) { + for (const verb of claim.serves ?? []) { + const run = byName.get(verb); + if (!run) { + console.log(`[mesh-tools] claims ${claim.seat} and has no tool named ${verb}, which that seat promises; not served`); + continue; + } + const subject = seatToolSubject(claim.seat, verb, claim.scope, credential?.node); + stops.push(await broker.handleSubject(subject, run)); + console.log(`[mesh-tools] serving ${claim.seat}'s ${verb} on ${subject}, admitted where this module holds the seat`); + } + } + return () => stops.forEach((s) => s()); +} diff --git a/test/client.test.ts b/test/client.test.ts index 1e26e48..f7fe14f 100644 --- a/test/client.test.ts +++ b/test/client.test.ts @@ -86,7 +86,7 @@ test("a person calls a tool and gets the module's own answer, unshaped", async ( const person = await connectNats({ url, module: "person.ada" }); try { const answer = await callTool(person, "shop.price", { of: "a hat" }); - assert.deepEqual(answer, { of: "a hat", cost: 12 }); + assert.deepEqual(answer.result, { of: "a hat", cost: 12 }); } finally { await person.close(); await mesh.close(); @@ -140,11 +140,46 @@ test("a role's tool is reached through the seat, and a module's own name is neve assert.equal(toolKey("mesh-controller.other", roles), "mesh-controller.other", "a verb the seat does not declare is a module's"); assert.equal(toolKey("shop.price", roles), "shop.price"); const answer = await callTool(person, "mesh-controller.status", {}, roles); - assert.deepEqual(answer, { output: "all quiet", ok: true }); + assert.deepEqual(answer.result, { output: "all quiet", ok: true }); const direct = await callTool(person, "seat:mesh-controller.status", {}); - assert.deepEqual(direct, { output: "all quiet", ok: true }); + assert.deepEqual(direct.result, { output: "all quiet", ok: true }); } finally { await person.close(); await mesh.close(); } }); + +// A module on two machines (novox/hq ADR 0159): a call names the machine and reaches that instance +// and no other; a call that names none reaches one of them and says which; a seat's verb is served +// by the claimant's tool of the same name on the seat's own subject. +test("a tool call names the machine it is for, and every answer says which machine answered", async (t) => { + if (!url) return t.skip("MESH_TEST_NATS unset"); + const onAnchor = await connectNats({ url: url!, module: "store", node: "anchor" }); + const onHome = await connectNats({ url: url!, module: "store", node: "home-server" }); + const asker = await connectNats({ url: url!, module: "console", node: "workstation" }); + try { + await onAnchor.handle("databases", async () => ({ at: "anchor" })); + await onHome.handle("databases", async () => ({ at: "home-server" })); + + const home = await callTool(asker, "store.databases@home-server", {}); + assert.deepEqual(home.result, { at: "home-server" }); + assert.equal(home.node, "home-server"); + const anchor = await callTool(asker, "store.databases@anchor", {}); + assert.deepEqual(anchor.result, { at: "anchor" }); + assert.equal(anchor.node, "anchor"); + + const whichever = await callTool(asker, "store.databases", {}); + assert.ok(["anchor", "home-server"].includes(whichever.node ?? ""), `an unnamed call still says who answered: ${whichever.node}`); + assert.deepEqual(whichever.result, { at: whichever.node }); + + // A seat's verb, served by the claimant on the seat's subject; a mesh seat's is flat. + await onAnchor.handleSubject("mesh.seat.mesh-store.tool.databases", async () => ({ seat: "mesh-store", at: "anchor" })); + const viaSeat = await callTool(asker, "seat:mesh-store.databases", {}); + assert.deepEqual(viaSeat.result, { seat: "mesh-store", at: "anchor" }); + assert.equal(viaSeat.node, "anchor"); + } finally { + await onAnchor.close(); + await onHome.close(); + await asker.close(); + } +}); diff --git a/test/mcp.test.ts b/test/mcp.test.ts index 3f9c4db..81e5858 100644 --- a/test/mcp.test.ts +++ b/test/mcp.test.ts @@ -99,8 +99,11 @@ test("a host initialises, lists the mesh's tools and calls one", async (t) => { "the modules' tools and the roles', named the way a person names them"); const price = listed[1]; assert.ok(price.inputSchema, "a tool with no schema is one an agent cannot call"); - // A module's bare property map arrives as a schema an agent can read, its words kept. - assert.deepEqual(price.inputSchema, { type: "object", properties: { of: { type: "string" } } }); + // A module's bare property map arrives as a schema an agent can read, its words kept — and + // `node`, the machine to ask when the module runs on several (novox/hq ADR 0159), beside them. + assert.deepEqual(price.inputSchema.properties.of, { type: "string" }); + assert.equal(price.inputSchema.properties.node.type, "string", "a module's tool takes the machine to ask"); + assert.ok(!listed[0].inputSchema.properties?.node, "a seat's verb takes no machine; the seat's scope decides"); // Silence is named: the module the catalogue holds and nothing answered for. assert.deepEqual(byId.get(2)!.result._meta.notAnswering, ["ghost"]);