A tool call names the machine it is for, every answer says which machine answered, and a holder's

runtime serves its seat's verbs (novox/hq ADR 0159)

A module on several machines served one subject in one queue group, so a call reached whichever
instance answered first and nobody could ask one machine's instance. Now each instance also serves
its subject with its machine as the last token, `<module>.<tool>@<node>` addresses it, and every
answer carries the machine that gave it: the console lists `node` on every module tool, strips it
into the subject, and appends "answered by <node>" to the answer; `mesh call` prints it.

And a seat's verbs are served by whoever claims the seat: the credential names the claimed seats
and the verbs each promises, the runtime serves each verb with the module's tool of the same name
on the seat's own subject, and the bus — which admits only the holder's subscription — decides where
that serving is real. Design 33 §3 and §4, built.
This commit is contained in:
2026-10-01 13:58:22 +02:00
parent 6a91d144b3
commit c65f1993ba
8 changed files with 227 additions and 63 deletions
+94 -40
View File
@@ -40,6 +40,25 @@ export interface Credential {
module?: string;
user?: string;
password?: string;
/** The seats this module claims, with the verbs each promises (novox/hq ADR 0159). The runtime
* serves each claimed seat's verbs with its tools of the same name; the bus admits only the
* holder's subscription, so claiming and not holding costs a refused subscription and nothing else. */
claims?: { seat: string; scope?: string; serves?: 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> {
result: Res;
node?: string;
}
/** The bus as the runtime sees it: the sdk's contract, and the two things only the runtime needs —
* 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>>;
handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void>;
}
/** Whether a connection failure is worth retrying, or is a fact about this configuration that
@@ -69,7 +88,7 @@ export function fatalBrokerReason(err: unknown): string | null {
export async function connectNats(
target: string | Credential,
opts: { module?: string } = {},
): Promise<Broker> {
): Promise<RuntimeBroker> {
const cred: Credential = typeof target === "string" ? { url: target } : target;
const self = cred.module ?? opts.module;
if (!self) {
@@ -102,49 +121,72 @@ export async function connectNats(
let reading: Awaited<ReturnType<Awaited<ReturnType<typeof js.consumers.get>>["consume"]>> | undefined;
let closed = false;
const node = cred.node;
/** Answer one subject with one handler, and say which machine answered (novox/hq ADR 0159). */
const answerOn = <Req, Res>(
subject: string,
queue: string | undefined,
handler: (body: Req) => Promise<Res>,
): (() => void) => {
const sub = conn.subscribe(subject, queue ? { queue } : {});
subs.push(sub);
void (async () => {
for await (const msg of sub) {
let reply: { result?: Res; error?: string; node?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
if (node) reply.node = node;
msg.respond(sc.encode(JSON.stringify(reply)));
}
})();
return () => sub.unsubscribe();
};
/**
* Ask one question and await one answer, with the machine that gave it.
*
* Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3),
* and a lost one is a timeout the caller already handles. The reply travels on the inbox the
* 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)), {
timeout: REQUEST_TIMEOUT_MS,
});
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string; node?: string };
if (reply.error) throw new Error(reply.error);
return { result: reply.result as Res, node: reply.node };
};
return {
/**
* Ask one question and await one answer.
*
* Core NATS request/reply, not JetStream: a tool call must never be persisted (design 25 §3),
* and a lost one is a timeout the caller already handles. The reply travels on the inbox the
* request carries, which the responder may answer because its account has `allow_responses`
* — one reply to a message it actually received, and nothing wider.
*/
async request<Req, Res>(key: string, body: Req): Promise<Res> {
const msg = await conn.request(toolSubject(key, self), sc.encode(JSON.stringify(body)), {
timeout: REQUEST_TIMEOUT_MS,
});
const reply = JSON.parse(sc.decode(msg.data)) as { result?: Res; error?: string };
if (reply.error) throw new Error(reply.error);
return reply.result as Res;
return (await ask<Req, Res>(key, body)).result;
},
ask,
/**
* Answer a question.
*
* A queue group, so several nodes may serve one tool and exactly one of them answers each
* call.
* Answer a question, two ways (novox/hq ADR 0159): on the module's subject in a queue group,
* so several machines may serve one tool and exactly one of them answers each call; and on the
* same subject with this machine as its last token, so a caller that names the machine reaches
* this instance and no other. A runtime that does not know its machine serves only the first,
* which is how it always behaved.
*/
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
const sub = conn.subscribe(toolSubject(key, self), { queue: `serve.${self}` });
subs.push(sub);
void (async () => {
for await (const msg of sub) {
let reply: { result?: Res; error?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
msg.respond(sc.encode(JSON.stringify(reply)));
}
})();
return () => {
sub.unsubscribe();
};
const stops = [answerOn(toolSubject(key, self), `serve.${self}`, handler)];
if (node) stops.push(answerOn(`${toolSubject(key, self)}.${node}`, undefined, handler));
return () => stops.forEach((stop) => stop());
},
async handleSubject<Req, Res>(subject: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
return answerOn(subject, undefined, handler);
},
/**
@@ -325,9 +367,21 @@ function toolSubject(key: string, self: string): string {
const [verb, node] = rest.slice(dot + 1).split("@", 2);
return node ? `mesh.seat.${seat}.tool.${verb}.${node}` : `mesh.seat.${seat}.tool.${verb}`;
}
const dot = key.indexOf(".");
if (dot < 0) return `mesh.mod.${self}.tool.${key}`;
return `mesh.mod.${key.slice(0, dot)}.tool.${key.slice(dot + 1)}`;
// `<module>.<tool>@<node>` names the machine (novox/hq ADR 0159): the same subject with the
// machine as its last token, which is what that instance serves beside the queue.
const [name, node] = key.split("@", 2);
const dot = name.indexOf(".");
const base = dot < 0
? `mesh.mod.${self}.tool.${name}`
: `mesh.mod.${name.slice(0, dot)}.tool.${name.slice(dot + 1)}`;
return node ? `${base}.${node}` : base;
}
/** A seat's verb, as its holder serves it: flat for a mesh seat, carrying the machine for a
* node-scoped one (design 33 §4) — the same shape the controller grants. */
export function seatToolSubject(seat: string, verb: string, scope: string | undefined, node: string | undefined): string {
const base = `mesh.seat.${seat}.tool.${verb}`;
return scope === "node" && node ? `${base}.${node}` : base;
}
function normalizeFingerprint(fingerprint: string): string {