Merge pull request 'A runtime serves what the mesh issued it, and a seat's verbs are implemented under the seat's name (hq ADR 0160)' (#23) from feat/a-runtime-serves-what-it-is-issued into main
This commit was merged in pull request #23.
This commit is contained in:
+130
-5
@@ -46,6 +46,24 @@ export interface Credential {
|
||||
claims?: { seat: string; scope?: string; serves?: string[] }[];
|
||||
}
|
||||
|
||||
/**
|
||||
* What the mesh issued this assignment (novox/hq ADR 0160): where its tools are served, in which
|
||||
* queue, the verbs of the seats it holds, where its events land, what it may reach. Read from the
|
||||
* ASSIGNMENTS stream at `mesh.assignment.<node>.<module>` — the one subject a runtime derives for
|
||||
* itself — and followed live. Absent for a mesh older than the membership, and then the runtime
|
||||
* serves the shape it always derived, and says so.
|
||||
*/
|
||||
export interface Membership {
|
||||
node: string;
|
||||
module: string;
|
||||
/** Addresses a tool is answered on; `{tool}` stands for the tool's name. */
|
||||
serves: { subject: string; queue?: string }[];
|
||||
seats?: { seat: string; verb: string; subject: string }[];
|
||||
emits: string;
|
||||
reaches?: Record<string, string[]>;
|
||||
tools: 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<Res> {
|
||||
@@ -57,8 +75,20 @@ export interface Answered<Res> {
|
||||
* 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<Req, Res>(key: string, body: Req): Promise<Answered<Res>>;
|
||||
/** Call a tool by key, or — when `on` names a subject the mesh listed for it (ADR 0160) — there. */
|
||||
ask<Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>>;
|
||||
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void>;
|
||||
/** What the mesh issued this assignment, or undefined when nothing has been issued yet. */
|
||||
membership(): Membership | undefined;
|
||||
/** Called when the mesh issues a new membership; the runtime re-serves on it. */
|
||||
onMembership(handler: (m: Membership) => void): void;
|
||||
}
|
||||
|
||||
const ASSIGNMENTS_STREAM = "ASSIGNMENTS";
|
||||
|
||||
/** The one address a runtime derives for itself (ADR 0160). */
|
||||
export function membershipSubject(node: string, module: string): string {
|
||||
return `mesh.assignment.${node}.${module}`;
|
||||
}
|
||||
|
||||
/** Whether a connection failure is worth retrying, or is a fact about this configuration that
|
||||
@@ -123,6 +153,80 @@ export async function connectNats(
|
||||
|
||||
const node = cred.node;
|
||||
|
||||
// The membership, read once at connect and followed. A direct get is one request on the
|
||||
// stream's API, which is the whole of what this account may ask JetStream for its own subject;
|
||||
// a 404 is a mesh that has not issued one, which is a fact to say and not an error to retry.
|
||||
let issued: Membership | undefined;
|
||||
const issuedHandlers: ((m: Membership) => void)[] = [];
|
||||
const subjectOfMine = node ? membershipSubject(node, self) : "";
|
||||
if (subjectOfMine) {
|
||||
try {
|
||||
const got = await conn.request(`$JS.API.DIRECT.GET.${ASSIGNMENTS_STREAM}`,
|
||||
sc.encode(JSON.stringify({ last_by_subj: subjectOfMine })), { timeout: 5_000 });
|
||||
const status = got.headers?.code ?? 0;
|
||||
if (status === 0 && got.data.length > 0) {
|
||||
issued = JSON.parse(sc.decode(got.data)) as Membership;
|
||||
}
|
||||
} catch {
|
||||
// Not readable here: an older mesh, a stream not yet asserted, or no grant. Said below.
|
||||
}
|
||||
if (!issued) {
|
||||
console.log(`[mesh-tools] no membership issued for ${self} on ${node} yet; serving the derived shape until one arrives`);
|
||||
}
|
||||
try {
|
||||
const live = conn.subscribe(subjectOfMine);
|
||||
subs.push(live);
|
||||
void (async () => {
|
||||
for await (const msg of live) {
|
||||
try {
|
||||
issued = JSON.parse(sc.decode(msg.data)) as Membership;
|
||||
console.log(`[mesh-tools] ${self} on ${node} was issued a new membership; re-serving on it`);
|
||||
for (const h of issuedHandlers) h(issued);
|
||||
} catch (err) {
|
||||
console.log(`[mesh-tools] a membership arrived that is not one: ${err}`);
|
||||
}
|
||||
}
|
||||
})();
|
||||
} catch {
|
||||
// A subscription this account may not make is a mesh older than the membership.
|
||||
}
|
||||
}
|
||||
|
||||
/** The subjects a tool of this module is served on: from the membership when issued, derived
|
||||
* otherwise (the shape the mesh issues on day one, so the two agree). */
|
||||
const servedOn = (tool: string): { subject: string; queue?: string }[] => {
|
||||
if (issued) {
|
||||
const m = issued;
|
||||
const out = m.serves.map((s) => ({ subject: s.subject.replace("{tool}", tool), queue: s.queue }));
|
||||
// The verb that lists what this module serves is answered on the mesh's plain address for it
|
||||
// whatever the placement — one answer suffices, so a queue — and on this machine's beside it.
|
||||
if (tool === "tools" && m.tools && !out.some((s) => s.subject === m.tools)) {
|
||||
out.unshift({ subject: m.tools, queue: `serve.${self}` });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
const base = `mesh.mod.${self}.tool.${tool}`;
|
||||
const out: { subject: string; queue?: string }[] = [{ subject: base, queue: `serve.${self}` }];
|
||||
if (node) out.push({ subject: `${base}.${node}` });
|
||||
return out;
|
||||
};
|
||||
|
||||
/** Where a call by key goes: a subject the membership says this module reaches, when it says
|
||||
* one — the machine's when named — else the derived shape. */
|
||||
const reachedAt = (key: string): string => {
|
||||
const [name, wanted] = key.split("@", 2);
|
||||
const reach = issued?.reaches?.[name];
|
||||
if (reach && reach.length > 0) {
|
||||
if (wanted) {
|
||||
const at = reach.find((s) => s.endsWith(`.${wanted}`));
|
||||
if (at) return at;
|
||||
} else {
|
||||
return reach[0];
|
||||
}
|
||||
}
|
||||
return toolSubject(key, self);
|
||||
};
|
||||
|
||||
/** Answer one subject with one handler, and say which machine answered (novox/hq ADR 0159). */
|
||||
const answerOn = <Req, Res>(
|
||||
subject: string,
|
||||
@@ -156,8 +260,8 @@ export async function connectNats(
|
||||
* 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 <Req, Res>(key: string, body: Req): Promise<Answered<Res>> => {
|
||||
const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), {
|
||||
const ask = async <Req, Res>(key: string, body: Req, on?: string): Promise<Answered<Res>> => {
|
||||
const msg = await conn.request(on ?? reachedAt(key), sc.encode(JSON.stringify(body)), {
|
||||
timeout: REQUEST_TIMEOUT_MS,
|
||||
});
|
||||
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string; node?: string };
|
||||
@@ -180,11 +284,32 @@ export async function connectNats(
|
||||
* which is how it always behaved.
|
||||
*/
|
||||
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||
const stops = [answerOn(toolSubject(key, self), `serve.${self}`, handler)];
|
||||
if (node) stops.push(answerOn(`${toolSubject(key, self)}.${node}`, undefined, handler));
|
||||
// A seat's verb named outright is served on the seat's subject as given, for a holder that
|
||||
// knows its role without a membership; everything else is this module's own tool, served
|
||||
// where the mesh issued it (ADR 0160). A key naming another module is not served here at all.
|
||||
if (key.startsWith("seat:")) {
|
||||
const stop = answerOn(toolSubject(key, self), undefined, handler);
|
||||
return () => stop();
|
||||
}
|
||||
const dot = key.indexOf(".");
|
||||
if (dot >= 0 && key.slice(0, dot) !== self) {
|
||||
throw new Error(`${self} cannot serve ${key}: a module serves its own tools`);
|
||||
}
|
||||
const tool = dot < 0 ? key : key.slice(dot + 1);
|
||||
let stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
// When a new membership arrives, serve where it now says and stop serving where it no longer does.
|
||||
issuedHandlers.push(() => {
|
||||
stops.forEach((stop) => stop());
|
||||
stops = servedOn(tool).map((s) => answerOn(s.subject, s.queue, handler));
|
||||
});
|
||||
return () => stops.forEach((stop) => stop());
|
||||
},
|
||||
|
||||
membership: () => issued,
|
||||
onMembership: (handler: (m: Membership) => void) => {
|
||||
issuedHandlers.push(handler);
|
||||
},
|
||||
|
||||
async handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
|
||||
return answerOn(subject, undefined, handler);
|
||||
},
|
||||
|
||||
+30
-4
@@ -40,6 +40,9 @@ export interface Tool {
|
||||
input?: unknown;
|
||||
/** True for a role's tool: addressed to the seat, answered by whoever holds it (ADR 0132). */
|
||||
seat?: boolean;
|
||||
/** Where the tool is answered, as the mesh issued it (ADR 0160): the plain subject first when the
|
||||
* module answers for itself anywhere, then one per machine. Absent for a runtime older than this. */
|
||||
subjects?: string[];
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -188,7 +191,7 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
|
||||
const module = names[i]!;
|
||||
if (outcome.status === "fulfilled" && Array.isArray(outcome.value?.tools)) {
|
||||
for (const t of outcome.value.tools) {
|
||||
tools.push({ module, name: t.name, description: t.description, input: t.input });
|
||||
tools.push({ module, name: t.name, description: t.description, input: t.input, subjects: t.subjects });
|
||||
}
|
||||
} else {
|
||||
notAnswering.push(module);
|
||||
@@ -204,7 +207,13 @@ export async function toolsOn(bus: Broker): Promise<Listing> {
|
||||
/** Call a tool and learn which machine answered (novox/hq ADR 0159). `<module>.<tool>@<node>` 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<Answered<unknown>> {
|
||||
export async function callTool(
|
||||
bus: Broker,
|
||||
key: string,
|
||||
args: unknown,
|
||||
seats?: Seats,
|
||||
listing?: Listing,
|
||||
): Promise<Answered<unknown>> {
|
||||
const at = key.indexOf("@");
|
||||
const name = at < 0 ? key : key.slice(0, at);
|
||||
const node = at < 0 ? "" : key.slice(at + 1);
|
||||
@@ -215,11 +224,28 @@ export async function callTool(bus: Broker, key: string, args: unknown, seats?:
|
||||
);
|
||||
}
|
||||
const resolved = toolKey(name, seats) + (node ? `@${node}` : "");
|
||||
const asking = bus as Broker & { ask?: <Req, Res>(k: string, b: Req) => Promise<Answered<Res>> };
|
||||
if (typeof asking.ask === "function") return asking.ask<unknown, unknown>(resolved, args ?? {});
|
||||
// Where the tool is answered is the module's to say and the mesh's to issue (ADR 0160): when the
|
||||
// listing carried subjects for it, the call goes to one of those and composes nothing.
|
||||
const on = subjectListed(name, node, listing);
|
||||
const asking = bus as Broker & { ask?: <Req, Res>(k: string, b: Req, on?: string) => Promise<Answered<Res>> };
|
||||
if (typeof asking.ask === "function") return asking.ask<unknown, unknown>(resolved, args ?? {}, on);
|
||||
return { result: await bus.request<unknown, unknown>(resolved, args ?? {}) };
|
||||
}
|
||||
|
||||
/** The subject the listing says answers `<module>.<tool>` — the machine's when one is named, else
|
||||
* the plain one — or undefined when the listing said none, and the key is composed as before. */
|
||||
export function subjectListed(name: string, node: string, listing?: Listing): string | undefined {
|
||||
if (!listing) return undefined;
|
||||
const dot = name.indexOf(".");
|
||||
const module = name.slice(0, dot);
|
||||
const tool = name.slice(dot + 1);
|
||||
const found = listing.tools.find((t) => t.module === module && t.name === tool && !t.seat);
|
||||
const subjects = found?.subjects ?? [];
|
||||
if (subjects.length === 0) return undefined;
|
||||
if (node) return subjects.find((s) => s.endsWith(`.${node}`));
|
||||
return subjects[0];
|
||||
}
|
||||
|
||||
/**
|
||||
* Why a call failed, said so that the remedy is in the words.
|
||||
*
|
||||
|
||||
+2
-1
@@ -139,7 +139,8 @@ export function mcpSurface(bus: Broker, who: string): Surface {
|
||||
delete args.node;
|
||||
const name = node && !given.includes("@") ? `${given}@${node}` : given;
|
||||
try {
|
||||
const { result, node: answeredBy } = await callTool(bus, name, args, await roles());
|
||||
const have = await listing().catch(() => undefined);
|
||||
const { result, node: answeredBy } = await callTool(bus, name, args, have && seatsIn(have), have);
|
||||
// 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. Which machine answered follows it as its own line.
|
||||
|
||||
+2
-2
@@ -173,8 +173,8 @@ async function calling(bus: Broker, args: string[]): Promise<number> {
|
||||
try {
|
||||
// `<seat>.<verb>` 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 { result, node } = await callTool(bus, key, parsed, roles);
|
||||
const have = key.startsWith("seat:") ? undefined : await toolsOn(bus).catch(() => undefined);
|
||||
const { result, node } = await callTool(bus, key, parsed, have && seatsIn(have), have);
|
||||
console.log(JSON.stringify(result, null, 2));
|
||||
if (node) console.error(`answered by ${node}`);
|
||||
return 0;
|
||||
|
||||
+72
-21
@@ -6,7 +6,7 @@
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { resolve } from "node:path";
|
||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||
import { collectTools, serveTools, listTools, toolKey } from "@novox/mesh-sdk/tools";
|
||||
import { collectTools, toolKey } from "@novox/mesh-sdk/tools";
|
||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||
import { seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
||||
|
||||
@@ -20,7 +20,14 @@ export const TOOLS_VERB = "tools";
|
||||
/** What `tools` answers for one module. */
|
||||
export interface ToolsAnswer {
|
||||
module: string;
|
||||
tools: { name: string; description: string; input: Readonly<Record<string, unknown>> }[];
|
||||
tools: {
|
||||
name: string;
|
||||
description: string;
|
||||
input: Readonly<Record<string, unknown>>;
|
||||
/** Where this tool is answered, as the mesh issued it (ADR 0160): the module's plain subject
|
||||
* first when there is one, then this machine's. A caller composes nothing. */
|
||||
subjects?: string[];
|
||||
}[];
|
||||
}
|
||||
|
||||
export interface RuntimeOptions {
|
||||
@@ -45,14 +52,36 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
|
||||
// audit logger) registers none, and its scoped account may not declare the serve queue — so a
|
||||
// runtime that always served would fail for exactly the modules that never needed it.
|
||||
const tools = listTools();
|
||||
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
|
||||
// A registration under a seat's name is the module's implementation of that seat's verbs
|
||||
// (ADR 0159, 0160): served on the seat's subjects by serveClaimedSeats, never as a module's
|
||||
// tools and never listed among them. Everything else is the module's own.
|
||||
// A module named like its seat (the catalogue is the mesh-catalog seat) registers once and is
|
||||
// both: its tools are the module's and the seat's verbs alike.
|
||||
const self = opts.credential?.module;
|
||||
const seatNames = new Set((opts.credential?.claims ?? []).map((c) => c.seat));
|
||||
const ownRegistrations = collectTools().filter(({ module }) => module === self || !seatNames.has(module));
|
||||
const tools = ownRegistrations.flatMap(({ module, tools: own }) => own.map((t) => ({ module, name: t.name })));
|
||||
const stops: Array<() => void> = [];
|
||||
const stop = (): void => stops.splice(0).forEach((s) => s());
|
||||
// Each tool on its own key, namespaced by its module (ADR 0047); where that key is answered is
|
||||
// the broker's to know from the membership (ADR 0160).
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
const seen = new Set<string>();
|
||||
for (const t of own) {
|
||||
if (seen.has(t.name)) {
|
||||
stop();
|
||||
throw new Error(`${module} exposes two tools named ${t.name} — refused`);
|
||||
}
|
||||
seen.add(t.name);
|
||||
stops.push(await opts.broker.handle(toolKey(module, t.name), (args: Record<string, unknown> | undefined) => t.run(args ?? {})));
|
||||
}
|
||||
}
|
||||
|
||||
// And, for every module that serves any, the verb that says what it serves. Refused before
|
||||
// anything is bound if a module named a tool of its own `tools`: one name answering two things
|
||||
// is the fault nobody can diagnose afterwards, and the runtime is the only place that sees both.
|
||||
const stops: Array<() => void> = [stop];
|
||||
for (const { module, tools: own } of collectTools()) {
|
||||
const runtime = opts.broker as RuntimeBroker;
|
||||
for (const { module, tools: own } of ownRegistrations) {
|
||||
if (own.length === 0) continue;
|
||||
if (own.some((t) => t.name === TOOLS_VERB)) {
|
||||
stop();
|
||||
@@ -61,9 +90,16 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
"module with what it serves (novox/hq ADR 0152) — refused, rename it",
|
||||
);
|
||||
}
|
||||
const subjectsOf = (tool: string): string[] | undefined => {
|
||||
const issued = typeof runtime.membership === "function" ? runtime.membership() : undefined;
|
||||
if (!issued) return undefined;
|
||||
const plain = issued.serves.filter((s) => s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
const mine = issued.serves.filter((s) => !s.queue).map((s) => s.subject.replace("{tool}", tool));
|
||||
return [...plain, ...mine];
|
||||
};
|
||||
const answer: ToolsAnswer = {
|
||||
module,
|
||||
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input })),
|
||||
tools: own.map((t) => ({ name: t.name, description: t.description, input: t.input, subjects: subjectsOf(t.name) })),
|
||||
};
|
||||
stops.push(await opts.broker.handle(toolKey(module, TOOLS_VERB), async () => answer));
|
||||
}
|
||||
@@ -85,22 +121,37 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
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<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
for (const { tools } of collectTools()) {
|
||||
for (const t of tools) byName.set(t.name, (args) => t.run(args));
|
||||
// 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>>>();
|
||||
for (const { module, tools } of collectTools()) {
|
||||
if (!claims.some((c) => c.seat === module)) continue;
|
||||
const verbs = new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
for (const t of tools) verbs.set(t.name, (args) => t.run(args));
|
||||
implementations.set(module, verbs);
|
||||
}
|
||||
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;
|
||||
let stops: (() => void)[] = [];
|
||||
const serve = async (): Promise<void> => {
|
||||
stops.forEach((s) => s());
|
||||
stops = [];
|
||||
const issued = typeof broker.membership === "function" ? broker.membership() : undefined;
|
||||
for (const claim of claims) {
|
||||
const verbs = implementations.get(claim.seat);
|
||||
for (const verb of claim.serves ?? []) {
|
||||
const run = verbs?.get(verb);
|
||||
if (!run) {
|
||||
console.log(`[mesh-tools] claims ${claim.seat} and implements no ${verb}, which that seat promises; not served`);
|
||||
continue;
|
||||
}
|
||||
// Where the mesh issued the verb when it has; the derived shape until then.
|
||||
const subject = issued?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.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`);
|
||||
}
|
||||
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`);
|
||||
}
|
||||
}
|
||||
};
|
||||
await serve();
|
||||
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
|
||||
return () => stops.forEach((s) => s());
|
||||
}
|
||||
|
||||
Vendored
+12
@@ -0,0 +1,12 @@
|
||||
// The shop again, in a file of its own: a module imported once stays imported, so a second test needs a second entrypoint.
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("shop", () => [
|
||||
{
|
||||
name: "price",
|
||||
description: "what something costs",
|
||||
input: { of: { type: "string", description: "the thing" } },
|
||||
run: async (args) => ({ of: args.of ?? "nothing", cost: 12 }),
|
||||
},
|
||||
{ name: "refund", description: "give it back", input: {}, run: async () => ({ done: true }) },
|
||||
]);
|
||||
Vendored
+12
@@ -0,0 +1,12 @@
|
||||
// A module that is both software and a role: postgres's own tools under its name, and its
|
||||
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("postgres", () => [
|
||||
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
|
||||
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
|
||||
]);
|
||||
|
||||
registerModuleTools("mesh-store", () => [
|
||||
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
|
||||
]);
|
||||
@@ -0,0 +1,197 @@
|
||||
/**
|
||||
* The mesh issues an assignment's subjects, and a runtime serves what it is issued (novox/hq ADR
|
||||
* 0160). A runtime derives one address for itself — `mesh.assignment.<node>.<module>` — reads the
|
||||
* membership there, serves exactly what it says, and re-serves when a new one arrives. Against a
|
||||
* real bus with JetStream, because the membership is a direct get on a stream.
|
||||
*
|
||||
* docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js
|
||||
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/membership.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 { callTool, subjectListed } from "../dist/client.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();
|
||||
|
||||
/** The ASSIGNMENTS stream as the controller asserts it: last-per-subject, readable by direct get. */
|
||||
async function anAssignmentsStream() {
|
||||
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);
|
||||
return {
|
||||
async issue(m: object & { node: string; module: string }) {
|
||||
await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m)));
|
||||
},
|
||||
async close() {
|
||||
await nc.close();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
test("a runtime serves exactly the subjects it is issued, and the listing carries them", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
// The mesh placed shop on two machines that are not interchangeable: each is served by name only.
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
const shop = await connectNats({ url, node: "anchor", module: "shop" });
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
|
||||
assert.equal(shop.membership()?.module, "shop", "the runtime read what the mesh issued");
|
||||
// By name it answers; the plain subject was not issued, so nothing serves it.
|
||||
const named = await callTool(asker, "shop.price@anchor", { of: "a hat" });
|
||||
assert.deepEqual(named.result, { of: "a hat", cost: 12 });
|
||||
assert.equal(named.node, "anchor");
|
||||
await assert.rejects(callTool(asker, "shop.price", {}), /no responders|503/i, "a subject not issued is not served");
|
||||
// The tools answer says where each is served, and the client composes nothing.
|
||||
const tools = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>(
|
||||
"shop.tools",
|
||||
{},
|
||||
);
|
||||
assert.deepEqual(tools.tools.find((x) => x.name === "price")?.subjects, ["mesh.mod.shop.tool.price.anchor"]);
|
||||
const listing = { tools: [{ module: "shop", name: "price", subjects: ["mesh.mod.shop.tool.price.anchor"] }], notAnswering: [] };
|
||||
assert.equal(subjectListed("shop.price", "anchor", listing as never), "mesh.mod.shop.tool.price.anchor");
|
||||
assert.equal(subjectListed("shop.price", "elsewhere", listing as never), undefined);
|
||||
|
||||
// The mesh re-issues the membership with the plain subject too (the modules became
|
||||
// interchangeable); the runtime re-serves on it without a restart.
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
await new Promise((r) => setTimeout(r, 300));
|
||||
const plain = await callTool(asker, "shop.price", { of: "a coat" });
|
||||
assert.deepEqual(plain.result, { of: "a coat", cost: 12 });
|
||||
assert.equal(plain.node, "anchor");
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await shop.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a seat's verbs are implemented under the seat's name, served where issued, never listed as the module's", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "postgres",
|
||||
serves: [{ subject: "mesh.mod.postgres.tool.{tool}", queue: "serve.postgres" }, { subject: "mesh.mod.postgres.tool.{tool}.anchor" }],
|
||||
seats: [{ seat: "mesh-store", verb: "databases", subject: "mesh.seat.mesh-store.tool.databases" }],
|
||||
emits: "mesh.mod.postgres.event.{event}",
|
||||
tools: "mesh.mod.postgres.tool.tools",
|
||||
});
|
||||
const credential = { url, node: "anchor", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] };
|
||||
const pg = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat.mjs")] });
|
||||
// The module's own `databases` and the seat's are two tools: postgres's on its subject, the
|
||||
// store's on the seat's, each answering as itself.
|
||||
const own = await callTool(asker, "postgres.databases", {});
|
||||
assert.deepEqual(own.result, { software: "postgres" });
|
||||
const seat = await callTool(asker, "seat:mesh-store.databases", {});
|
||||
assert.deepEqual(seat.result, { seat: "mesh-store" });
|
||||
assert.equal(seat.node, "anchor");
|
||||
// What the seat does not promise — creating a database — is postgres's tool and not the store's.
|
||||
await assert.rejects(callTool(asker, "seat:mesh-store.postgres_create_database", {}), /no responders|503/i);
|
||||
// And the module's `tools` lists only postgres's own, never the seat's implementation.
|
||||
const tools = await asker.request<Record<string, never>, { module: string; tools: { name: string }[] }>("postgres.tools", {});
|
||||
assert.deepEqual(tools.tools.map((x) => x.name).sort(), ["databases", "postgres_create_database"]);
|
||||
await assert.rejects(asker.request("mesh-store.tools", {}), /no responders|503/i, "a seat is not a module with a tools verb");
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await pg.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const stream = await anAssignmentsStream();
|
||||
await stream.issue({
|
||||
node: "anchor",
|
||||
module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }],
|
||||
seats: [{ seat: "shop", verb: "price", subject: "mesh.seat.shop.tool.price" }],
|
||||
emits: "mesh.mod.shop.event.{event}",
|
||||
tools: "mesh.mod.shop.tool.tools",
|
||||
});
|
||||
const credential = { url, node: "anchor", module: "shop", claims: [{ seat: "shop", scope: "mesh", serves: ["price"] }] };
|
||||
const shop = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: shop, credential, moduleEntrypoints: [fixture("shop-seat.mjs")] });
|
||||
assert.deepEqual((await callTool(asker, "shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
|
||||
assert.deepEqual((await callTool(asker, "seat:shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 });
|
||||
const tools = await asker.request<Record<string, never>, { tools: { name: string }[] }>("shop.tools", {});
|
||||
assert.deepEqual(tools.tools.map((x) => x.name), ["price", "refund"]);
|
||||
} finally {
|
||||
stop();
|
||||
await asker.close();
|
||||
await shop.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
|
||||
test("a mesh that has issued nothing yet gets the derived shape, and says so", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
const stream = await anAssignmentsStream();
|
||||
const said: string[] = [];
|
||||
const log = console.log;
|
||||
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||
let lone;
|
||||
try {
|
||||
lone = await connectNats({ url, node: "home-server", module: "lone" });
|
||||
} finally {
|
||||
console.log = log;
|
||||
}
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
try {
|
||||
assert.equal(lone.membership(), undefined);
|
||||
assert.ok(said.some((s) => /no membership issued for lone on home-server/.test(s)), said.join("\n"));
|
||||
await lone.handle("ping", async () => ({ pong: true }));
|
||||
assert.deepEqual((await callTool(asker, "lone.ping", {})).result, { pong: true });
|
||||
assert.deepEqual((await callTool(asker, "lone.ping@home-server", {})).result, { pong: true });
|
||||
} finally {
|
||||
await asker.close();
|
||||
await lone.close();
|
||||
await stream.close();
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user