Both runtimes subscribed $SRV.<verb>.> as a wildcard; the grants allow the bare question and the service's own name and instance. The bus refused the wildcard, and the TypeScript runtime treats a refused subscription as fatal, so every per-module container crash-looped after the image rolled. They now subscribe exactly $SRV.<verb>, $SRV.<verb>.<name> and $SRV.<verb>.<name>.<id>.
140 lines
5.6 KiB
TypeScript
140 lines
5.6 KiB
TypeScript
// 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)[] = [];
|
|
// Exactly the questions asked of every service and of this one by name and instance — what the
|
|
// grants allow (novox/hq ADR 0197). A wildcard is refused by the bus, and a refused subscription
|
|
// ends a runtime: on 2026-10-03 it crash-looped every container that announced itself.
|
|
for (const verb of ["PING", "INFO", "STATS"]) {
|
|
for (const subject of [`$SRV.${verb}`, `$SRV.${verb}.${s.name}`, `$SRV.${verb}.${s.name}.${s.id}`]) {
|
|
stops.push(broker.raw!(subject, (subj) => answer(subj)));
|
|
}
|
|
}
|
|
return () => stops.forEach((stop) => stop());
|
|
}
|