tools serve: the dispatch harness, proven end to end
Add serveTools(broker) — the tool runtime's core: collect every module's registered tools, index by name (refusing a duplicate name across two modules rather than silently shadowing), and answer 'tools.invoke' requests by running the named tool and returning its result. Plus listTools() for discovery and Broker.handle() (the server side of request/reply). Proven by test: a real async, network-calling tool is registered (as a module does), served over an in-memory broker, and invoked by name — it reaches its upstream and returns the computed result. So a module's tools genuinely serve: register -> collect -> serve -> invoke -> real work -> result. The per-node runtime process that binds the mesh's real broker and imports the assigned modules is the thin wrapper over this. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
@@ -9,8 +9,12 @@ export type { Envelope };
|
|||||||
|
|
||||||
/** A request/reply call and a publish/subscribe surface over the mesh broker. */
|
/** A request/reply call and a publish/subscribe surface over the mesh broker. */
|
||||||
export interface Broker {
|
export interface Broker {
|
||||||
/** Ask one question and await one answer — the shape mesh-control's command API is reached by. */
|
/** Ask one question and await one answer — the client side; the shape mesh-control's command
|
||||||
|
* API is reached by. */
|
||||||
request<Req, Res>(key: string, body: Req): Promise<Res>;
|
request<Req, Res>(key: string, body: Req): Promise<Res>;
|
||||||
|
/** Answer a question — the server side of request/reply. A tool runtime serves invocations this
|
||||||
|
* way. Returns an unregister. */
|
||||||
|
handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void>;
|
||||||
/** Emit an event onto the mesh. */
|
/** Emit an event onto the mesh. */
|
||||||
publish<T>(env: Envelope<T>): Promise<void>;
|
publish<T>(env: Envelope<T>): Promise<void>;
|
||||||
/** React to events matching a routing-key pattern. Returns an unsubscribe. */
|
/** React to events matching a routing-key pattern. Returns an unsubscribe. */
|
||||||
|
|||||||
@@ -4,9 +4,16 @@
|
|||||||
// the API client they call, live in the module (novox/hq ADR 0044).
|
// the API client they call, live in the module (novox/hq ADR 0044).
|
||||||
|
|
||||||
import type { ToolDefinition } from "../contracts/index.js";
|
import type { ToolDefinition } from "../contracts/index.js";
|
||||||
|
import type { Broker } from "../messaging/index.js";
|
||||||
|
|
||||||
export type { ToolDefinition };
|
export type { ToolDefinition };
|
||||||
|
|
||||||
|
/** One tool invocation crossing the broker: which tool, and its arguments. */
|
||||||
|
export interface Invocation {
|
||||||
|
readonly tool: string;
|
||||||
|
readonly args: Readonly<Record<string, unknown>>;
|
||||||
|
}
|
||||||
|
|
||||||
/** A module contributes its tools as a function of its resolved environment. Returning [] (e.g.
|
/** A module contributes its tools as a function of its resolved environment. Returning [] (e.g.
|
||||||
* when a token is absent) is normal — the module simply exposes nothing until it can. */
|
* when a token is absent) is normal — the module simply exposes nothing until it can. */
|
||||||
export type ToolContributor = (env: NodeJS.ProcessEnv) => ToolDefinition[];
|
export type ToolContributor = (env: NodeJS.ProcessEnv) => ToolDefinition[];
|
||||||
@@ -42,6 +49,39 @@ export function collectTools(env: NodeJS.ProcessEnv = process.env): { module: st
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** List the tools every registered module exposes — for discovery, without invoking anything. */
|
||||||
|
export function listTools(env: NodeJS.ProcessEnv = process.env): { module: string; name: string; description: string }[] {
|
||||||
|
return collectTools(env).flatMap(({ module, tools }) =>
|
||||||
|
tools.map((t) => ({ module, name: t.name, description: t.description })),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Serve the registered modules' tools over the mesh broker — the tool runtime's core. It collects
|
||||||
|
* every module's tools, indexes them by name, and answers `tools.invoke` requests by running the
|
||||||
|
* named tool and returning its result. This is the stable half of serving; a specific tool's work
|
||||||
|
* (and its client) lives in its module. Returns a stop function.
|
||||||
|
*
|
||||||
|
* A duplicate tool name across modules is refused loudly rather than one silently shadowing the
|
||||||
|
* other — two tools answering one name is a fault, not a race to resolve.
|
||||||
|
*/
|
||||||
|
export async function serveTools(broker: Broker, env: NodeJS.ProcessEnv = process.env): Promise<() => void> {
|
||||||
|
const byName = new Map<string, ToolDefinition>();
|
||||||
|
for (const { module, tools } of collectTools(env)) {
|
||||||
|
for (const t of tools) {
|
||||||
|
if (byName.has(t.name)) {
|
||||||
|
throw new Error(`tool name ${t.name} is exposed by two modules (one is ${module}) — refused`);
|
||||||
|
}
|
||||||
|
byName.set(t.name, t);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return broker.handle<Invocation, unknown>("tools.invoke", async (inv) => {
|
||||||
|
const tool = byName.get(inv.tool);
|
||||||
|
if (!tool) throw new Error(`no such tool: ${inv.tool}`);
|
||||||
|
return tool.run(inv.args ?? {});
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
/** Testing/inspection: drop all registrations. */
|
/** Testing/inspection: drop all registrations. */
|
||||||
export function resetTools(): void {
|
export function resetTools(): void {
|
||||||
registrations.length = 0;
|
registrations.length = 0;
|
||||||
|
|||||||
+77
-1
@@ -5,7 +5,7 @@ import { tmpdir } from "node:os";
|
|||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
|
|
||||||
import { seal, unseal, compareVersions } from "../dist/primitives/index.js";
|
import { seal, unseal, compareVersions } from "../dist/primitives/index.js";
|
||||||
import { registerModuleTools, collectTools, resetTools } from "../dist/tools/index.js";
|
import { registerModuleTools, collectTools, resetTools, serveTools, listTools } from "../dist/tools/index.js";
|
||||||
import { runProvisioner, type Grant } from "../dist/provisioner/index.js";
|
import { runProvisioner, type Grant } from "../dist/provisioner/index.js";
|
||||||
|
|
||||||
test("seal round-trips and rejects the wrong key", () => {
|
test("seal round-trips and rejects the wrong key", () => {
|
||||||
@@ -72,6 +72,82 @@ test("provisioner creates a sealed credential for a grant, then removes on withd
|
|||||||
stop();
|
stop();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("modules can SERVE: a real async tool, loaded and invoked over the broker", async () => {
|
||||||
|
resetTools();
|
||||||
|
|
||||||
|
// Stand up a fake upstream the tool actually calls over the network — proving a tool that does
|
||||||
|
// real work (not a pure function) serves end to end.
|
||||||
|
const { createServer } = await import("node:http");
|
||||||
|
const server = createServer((_req, res) => {
|
||||||
|
res.setHeader("content-type", "application/json");
|
||||||
|
res.end(JSON.stringify({ id: "site-123" }));
|
||||||
|
});
|
||||||
|
await new Promise<void>((r) => server.listen(0, r));
|
||||||
|
const addr = server.address() as { port: number };
|
||||||
|
const base = `http://127.0.0.1:${addr.port}`;
|
||||||
|
|
||||||
|
// A module registers its tool exactly as umami does — the tool calls the upstream and returns a
|
||||||
|
// result computed from it.
|
||||||
|
registerModuleTools("demo", () => [
|
||||||
|
{
|
||||||
|
name: "create_site",
|
||||||
|
description: "register a site and return its snippet",
|
||||||
|
input: { domain: { type: "string" } },
|
||||||
|
run: async (args) => {
|
||||||
|
const res = await fetch(`${base}/api/websites`, { method: "POST" });
|
||||||
|
const { id } = (await res.json()) as { id: string };
|
||||||
|
return { domain: String(args.domain), snippet: `<script data-website-id="${id}"></script>` };
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
|
||||||
|
const broker = memBroker();
|
||||||
|
const stop = await serveTools(broker, {});
|
||||||
|
|
||||||
|
// Invoke it the way a caller (mesh-control's command API) would — over the broker, by name.
|
||||||
|
const result = await broker.request<{ tool: string; args: Record<string, unknown> }, { snippet: string }>(
|
||||||
|
"tools.invoke",
|
||||||
|
{ tool: "create_site", args: { domain: "my-app" } },
|
||||||
|
);
|
||||||
|
assert.match(result.snippet, /data-website-id="site-123"/);
|
||||||
|
|
||||||
|
// Discovery works, and an unknown tool is refused rather than silently dropped.
|
||||||
|
assert.deepEqual(listTools({}).map((t) => t.name), ["create_site"]);
|
||||||
|
await assert.rejects(broker.request("tools.invoke", { tool: "nope", args: {} }));
|
||||||
|
|
||||||
|
stop();
|
||||||
|
server.close();
|
||||||
|
});
|
||||||
|
|
||||||
|
test("serving refuses two modules exposing one tool name", async () => {
|
||||||
|
resetTools();
|
||||||
|
registerModuleTools("a", () => [{ name: "dup", description: "", input: {}, run: async () => 1 }]);
|
||||||
|
registerModuleTools("b", () => [{ name: "dup", description: "", input: {}, run: async () => 2 }]);
|
||||||
|
await assert.rejects(serveTools(memBroker(), {}), /exposed by two modules/);
|
||||||
|
});
|
||||||
|
|
||||||
|
// A minimal in-memory broker: request routes to a registered handle. Enough to serve tools; the
|
||||||
|
// real broker binding is the mesh's, provided by the hosting runtime.
|
||||||
|
function memBroker() {
|
||||||
|
const handlers = new Map<string, (b: unknown) => Promise<unknown>>();
|
||||||
|
return {
|
||||||
|
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||||
|
const h = handlers.get(key);
|
||||||
|
if (!h) throw new Error(`no handler for ${key}`);
|
||||||
|
return (await h(body)) as Res;
|
||||||
|
},
|
||||||
|
async handle<Req, Res>(key: string, handler: (b: Req) => Promise<Res>): Promise<() => void> {
|
||||||
|
handlers.set(key, handler as (b: unknown) => Promise<unknown>);
|
||||||
|
return () => handlers.delete(key);
|
||||||
|
},
|
||||||
|
async publish(): Promise<void> {},
|
||||||
|
async subscribe(): Promise<() => void> {
|
||||||
|
return () => {};
|
||||||
|
},
|
||||||
|
async close(): Promise<void> {},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
||||||
const start = Date.now();
|
const start = Date.now();
|
||||||
while (!cond()) {
|
while (!cond()) {
|
||||||
|
|||||||
Reference in New Issue
Block a user