diff --git a/node-tools/src/announce.ts b/node-tools/src/announce.ts new file mode 100644 index 0000000..a152a51 --- /dev/null +++ b/node-tools/src/announce.ts @@ -0,0 +1,134 @@ +// What a runtime serves, it announces (novox/hq ADR 0197): the NATS services protocol's discovery — +// `$SRV.PING`, `$SRV.INFO`, `$SRV.STATS`, and the same followed by the service's name and its id — +// answered in the io.nats.micro.v1 format with what is served at the moment of the request. Serving is +// unchanged; this only says what is served. One service per runtime process, because the bus admits +// one reply per request from each responder: one endpoint per tool per subject, its metadata saying +// which module, seat, scope and machine it is. The same shape the Go runtime answers. + +import { StringCodec } from "nats"; +import { asSchema } from "@novox/mesh-sdk/stdio"; +import type { ToolDefinition } from "@novox/mesh-sdk/tools"; +import type { RuntimeBroker } from "./broker-nats.js"; + +const sc = StringCodec(); + +export const VERSION = "0.1.0"; +export const INFO_RESPONSE = "io.nats.micro.v1.info_response"; +export const PING_RESPONSE = "io.nats.micro.v1.ping_response"; +export const STATS_RESPONSE = "io.nats.micro.v1.stats_response"; + +/** One tool served on one subject, as announced. */ +export interface Endpoint { + kind: "tool" | "seat"; + module: string; + tool: string; + seat?: string; + scope?: "mesh" | "node"; + node: string; + description: string; + schema: unknown; + interchangeable: boolean; + subject: string; + queue?: string; +} + +/** A seat's verb as the runtime serves it, with the definition that answers it. */ +export interface ServedSeatVerb { + seat: string; + verb: string; + subject: string; + holder: string; + tool: ToolDefinition; +} + +export interface Service { + name: string; + id: string; + description: string; + metadata: Record; +} + +/** The endpoint's name as the protocol allows it; the metadata, not the name, identifies it. */ +function nameOf(e: Endpoint): string { + const clean = (s: string) => s.replace(/[^A-Za-z0-9_-]/g, "_"); + return `${clean(e.kind === "seat" ? e.seat ?? "" : e.module)}__${clean(e.tool)}`; +} + +/** The info_response for these endpoints. */ +export function info(s: Service, endpoints: Endpoint[]): Record { + return { + name: s.name, id: s.id, version: VERSION, metadata: s.metadata, type: INFO_RESPONSE, description: s.description, + endpoints: endpoints.map((e) => { + const metadata: Record = { + kind: e.kind, module: e.module, tool: e.tool, node: e.node, description: e.description, + schema: JSON.stringify(e.schema ?? {}), interchangeable: e.interchangeable ? "true" : "false", + }; + if (e.kind === "seat") { + metadata.seat = e.seat ?? ""; + metadata.scope = e.scope ?? "mesh"; + } + return { name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", metadata }; + }), + }; +} + +/** Everything served now: each module's tools on every subject issued for them, and each held + * seat's verbs on the seat's subject. */ +export function endpointsOf( + broker: RuntimeBroker, + own: { module: string; tools: ToolDefinition[] }[], + seats: ServedSeatVerb[], +): Endpoint[] { + const node = broker.node ?? ""; + const out: Endpoint[] = []; + for (const { module, tools } of own) { + for (const t of tools) { + const served = broker.servedOn ? broker.servedOn(module, t.name) : []; + const interchangeable = served.some((s) => s.subject === `mesh.mod.${module}.tool.${t.name}`); + for (const s of served) { + out.push({ kind: "tool", module, tool: t.name, node, description: t.description, schema: asSchema(t.input), + interchangeable, subject: s.subject, queue: s.queue }); + } + } + } + for (const v of seats) { + out.push({ kind: "seat", module: v.holder, tool: v.verb, seat: v.seat, + scope: node && v.subject.endsWith(`.${node}`) ? "node" : "mesh", node, description: v.tool.description, + schema: asSchema(v.tool.input), interchangeable: false, subject: v.subject }); + } + return out; +} + +/** Answer discovery for one service until stopped. */ +export function announce(broker: RuntimeBroker, s: Service, current: () => Endpoint[]): () => void { + const started = new Date().toISOString(); + const identity = { name: s.name, id: s.id, version: VERSION, metadata: s.metadata }; + const answer = (subject: string): Uint8Array | undefined => { + const parts = subject.split("."); + if (parts[0] !== "$SRV" || parts.length < 2) return undefined; + if (parts.length >= 3 && parts[2] !== s.name) return undefined; // another service's + if (parts.length >= 4 && parts[3] !== s.id) return undefined; // another instance's + let v: unknown; + switch (parts[1]) { + case "PING": + v = { ...identity, type: PING_RESPONSE }; + break; + case "INFO": + v = info(s, current()); + break; + case "STATS": + v = { ...identity, type: STATS_RESPONSE, started, endpoints: current().map((e) => ({ + name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", num_requests: 0, num_errors: 0, + last_error: "", processing_time: 0, average_processing_time: 0 })) }; + break; + default: + return undefined; + } + return sc.encode(JSON.stringify(v)); + }; + const stops: (() => void)[] = []; + for (const verb of ["PING", "INFO", "STATS"]) { + for (const subject of [`$SRV.${verb}`, `$SRV.${verb}.>`]) stops.push(broker.raw!(subject, (subj) => answer(subj))); + } + return () => stops.forEach((stop) => stop()); +} diff --git a/node-tools/src/broker-nats.ts b/node-tools/src/broker-nats.ts index cbf53c9..e7546eb 100644 --- a/node-tools/src/broker-nats.ts +++ b/node-tools/src/broker-nats.ts @@ -97,6 +97,13 @@ export interface RuntimeBroker extends Broker { serving(): string[]; /** The module this connection is: what its credential named, and what a bare key serves as. */ readonly module: string; + /** Where a served module's tool is answered right now (ADR 0197: what it serves, it announces). */ + servedOn?(module: string, tool: string): { subject: string; queue?: string }[]; + /** Answer a subject in a format of its own, not a tool's reply envelope — the NATS services + * protocol's discovery (novox/hq ADR 0197). An undefined answer is no reply. */ + raw?(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void; + /** The machine this connection serves on, when its credential names one. */ + readonly node?: string; } /** @@ -348,6 +355,19 @@ export async function connectNats( }, module: self, + node, + servedOn, + raw(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void { + const sub = conn.subscribe(subject); + subs.push(sub); + void (async () => { + for await (const msg of sub) { + const body = answer(msg.subject, msg.data); + if (body) msg.respond(body); + } + })(); + return () => sub.unsubscribe(); + }, follow, serving: () => [...issued.keys()], membership: (module?: string) => issued.get(module ?? self), diff --git a/node-tools/src/runtime.ts b/node-tools/src/runtime.ts index e18b1c9..578d32b 100644 --- a/node-tools/src/runtime.ts +++ b/node-tools/src/runtime.ts @@ -16,6 +16,7 @@ import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tool import type { Broker } from "@novox/mesh-sdk/messaging"; import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js"; import { launch, launches } from "./launch.js"; +import { announce, endpointsOf, type ServedSeatVerb } from "./announce.js"; /** * The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the @@ -259,7 +260,16 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> { console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` + (failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : "")); - stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations)); + const seats = await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations); + stops.push(seats.stop); + // **What it serves, it announces** (novox/hq ADR 0197), asked at the moment of the request. + if (typeof runtime.raw === "function") { + stops.push(announce(runtime, { + name: self ?? "runtime", id: runtime.node ?? self ?? "runtime", + description: `the tool runtime of ${self ?? "a module"}${runtime.node ? ` on ${runtime.node}` : ""}`, + metadata: runtime.node ? { node: runtime.node } : {}, + }, () => endpointsOf(runtime, ownRegistrations, seats.serving()))); + } return () => stop(); } @@ -303,11 +313,17 @@ async function serveClaimedSeats( self: string | undefined, credential: Credential | undefined, registrations: { module: string; owner: string; tools: ToolDefinition[] }[], -): Promise<() => void> { - if (typeof broker.handleSubject !== "function") return () => {}; +): Promise<{ stop: () => void; serving: () => ServedSeatVerb[] }> { + if (typeof broker.handleSubject !== "function") return { stop: () => {}, serving: () => [] }; // A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's // name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools. const implementations = new Map) => Promise>>(); + const definitions = new Map>(); + for (const { module, tools } of registrations) { + const defs = definitions.get(module) ?? new Map(); + for (const t of tools) defs.set(t.name, t); + definitions.set(module, defs); + } for (const { module, owner, tools } of registrations) { const verbs = implementations.get(module) ?? new Map) => Promise>(); for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args))); @@ -341,21 +357,25 @@ async function serveClaimedSeats( }; let stops: (() => void)[] = []; + let servingNow: ServedSeatVerb[] = []; const serve = async (): Promise => { stops.forEach((s) => s()); stops = []; + servingNow = []; for (const v of wanted()) { const run = implementations.get(v.seat)?.get(v.verb); + const tool = definitions.get(v.seat)?.get(v.verb); if (!run) { console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`); continue; } stops.push(await broker.handleSubject(v.subject, run)); + if (tool) servingNow.push({ ...v, tool }); console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`); } }; await serve(); // A membership issued to any served module may add, move or withdraw a seat's verbs. if (typeof broker.onMembership === "function") broker.onMembership(() => void serve()); - return () => stops.forEach((s) => s()); + return { stop: () => stops.forEach((s) => s()), serving: () => servingNow }; } diff --git a/node-tools/test/announce.test.ts b/node-tools/test/announce.test.ts new file mode 100644 index 0000000..7df625c --- /dev/null +++ b/node-tools/test/announce.test.ts @@ -0,0 +1,69 @@ +/** + * What a runtime serves, it announces (novox/hq ADR 0197): the TypeScript runtime the per-module + * containers still run answers the NATS services protocol's discovery in the same shape as the Go + * tool runtime — its module's tools on every subject issued, and the seat verbs it serves. + * + * MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/announce.test.ts + */ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { fileURLToPath } from "node:url"; +import { connect, StringCodec } from "nats"; + +import { resetTools } from "@novox/mesh-sdk/tools"; + +import { connectNats, membershipSubject } from "../dist/broker-nats.js"; +import { runTools } from "../dist/runtime.js"; + +const url = process.env.MESH_TEST_NATS; +const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url)); +const sc = StringCodec(); + +test("the runtime answers $SRV.INFO with what it serves, in the services protocol's format", async (t) => { + if (!url) return t.skip("MESH_TEST_NATS unset"); + resetTools(); + const nc = await connect({ servers: url }); + const jsm = await nc.jetstreamManager(); + try { + await jsm.streams.delete("ASSIGNMENTS"); + } catch { + // none yet + } + await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never); + await nc.jetstream().publish(membershipSubject("anchor", "shop"), sc.encode(JSON.stringify({ + node: "anchor", module: "shop", + serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }, { subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }], + emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools", + }))); + const shop = await connectNats({ url, node: "anchor", module: "shop" }); + let stop = () => {}; + try { + stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] }); + const msg = await nc.request("$SRV.INFO.shop.anchor", sc.encode(""), { timeout: 2000 }); + const info = JSON.parse(sc.decode(msg.data)) as { + type: string; name: string; id: string; version: string; + endpoints: { name: string; subject: string; queue_group: string; metadata: Record }[]; + }; + assert.equal(info.type, "io.nats.micro.v1.info_response"); + assert.equal(info.name, "shop"); + assert.equal(info.id, "anchor"); + assert.ok(info.version); + const price = info.endpoints.filter((e) => e.metadata.tool === "price").map((e) => `${e.subject}|${e.queue_group}`).sort(); + assert.deepEqual(price, ["mesh.mod.shop.tool.price.anchor|", "mesh.mod.shop.tool.price|serve.shop"]); + const one = info.endpoints.find((e) => e.metadata.tool === "price")!; + assert.equal(one.metadata.kind, "tool"); + assert.equal(one.metadata.module, "shop"); + assert.equal(one.metadata.node, "anchor"); + assert.equal(one.metadata.interchangeable, "true"); + assert.ok(JSON.parse(one.metadata.schema).type === "object"); + // Ping answers with the same identity; another service's request is not answered. + const ping = JSON.parse(sc.decode((await nc.request("$SRV.PING", sc.encode(""), { timeout: 2000 })).data)); + assert.equal(ping.type, "io.nats.micro.v1.ping_response"); + await assert.rejects(nc.request("$SRV.INFO.somebody-else", sc.encode(""), { timeout: 300 })); + } finally { + stop(); + await shop.close(); + await nc.close(); + resetTools(); + } +});