The mesh's tools are found by address, from what announces itself on the bus (hq ADR 0195, 0197) #39
@@ -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<string, string>;
|
||||
}
|
||||
|
||||
/** 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<string, unknown> {
|
||||
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<string, string> = {
|
||||
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());
|
||||
}
|
||||
@@ -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),
|
||||
|
||||
@@ -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<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
|
||||
const definitions = new Map<string, Map<string, ToolDefinition>>();
|
||||
for (const { module, tools } of registrations) {
|
||||
const defs = definitions.get(module) ?? new Map<string, ToolDefinition>();
|
||||
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<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
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<void> => {
|
||||
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 };
|
||||
}
|
||||
|
||||
@@ -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<string, string> }[];
|
||||
};
|
||||
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();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user