SDK: per-key tool serving (ADR 0052) and the provider contract (ADR 0053) #2
+29
-1
@@ -44,9 +44,37 @@ export interface ToolDefinition {
|
||||
readonly run: (args: Readonly<Record<string, unknown>>) => Promise<unknown>;
|
||||
}
|
||||
|
||||
/** A message crossing the broker: a routing key and a JSON body, per-node addressed. */
|
||||
/**
|
||||
* The metadata that rides an event as AMQP headers (novox/hq ADR 0047). An event's identity and
|
||||
* provenance live here, not in the body, so a consumer — or the broker, or an audit tool — reads
|
||||
* who/when/what without parsing the payload. An unknown `x-` header is ignored, not refused: an
|
||||
* event is observed by parties that need not all understand every header.
|
||||
*/
|
||||
export interface EventHeaders {
|
||||
/** A unique id — for dedup and audit (delivery is at-least-once). */
|
||||
readonly "x-event-id": string;
|
||||
/** The emitter: the module, context or node name. */
|
||||
readonly "x-source": string;
|
||||
/** The node it was emitted from. */
|
||||
readonly "x-node": string;
|
||||
/** Emit time, RFC-3339. */
|
||||
readonly "x-time": string;
|
||||
/** Always `application/json`. */
|
||||
readonly "content-type": string;
|
||||
/** The event or command that caused this one — tracing. */
|
||||
readonly "x-causation-id"?: string;
|
||||
/** A version of the body's shape, so a body evolves without silent misreads. */
|
||||
readonly "x-schema"?: string;
|
||||
readonly [header: string]: string | undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* A message crossing the broker: a routing key and a JSON body, per-node addressed. For an event,
|
||||
* `headers` carries the ADR 0047 metadata; plain request/reply transport leaves it absent.
|
||||
*/
|
||||
export interface Envelope<T = unknown> {
|
||||
readonly key: string;
|
||||
readonly node: string;
|
||||
readonly body: T;
|
||||
readonly headers?: EventHeaders;
|
||||
}
|
||||
|
||||
+58
-18
@@ -4,45 +4,85 @@
|
||||
// credential — just the broker's topic routing). Declared as emits/consumes on the manifest so the
|
||||
// mesh knows the event graph.
|
||||
//
|
||||
// This is a thin, audit-ready surface over the broker's publish/subscribe: every event carries who
|
||||
// emitted it, on which node, and when — so a logger consuming `#` can write a real audit trail.
|
||||
// Identity and provenance ride as AMQP headers, not in the body (novox/hq ADR 0047): a consumer —
|
||||
// or the broker, or an audit tool — reads who/when/what without parsing the payload, and the body
|
||||
// is only the domain payload. This surface hides the header/exchange/queue mechanics; a module
|
||||
// names a type and a body and never sees the wire.
|
||||
|
||||
import { randomUUID } from "node:crypto";
|
||||
|
||||
import { broker } from "../messaging/index.js";
|
||||
import type { Envelope, EventHeaders } from "../contracts/index.js";
|
||||
|
||||
/** An event on the mesh: a topic key, its source, and a body — with the metadata audit needs. */
|
||||
/** An event on the mesh: a topic key, its provenance, and a body — the shape a handler receives. */
|
||||
export interface Event<T = unknown> {
|
||||
/** The routing key, e.g. "module.umami.site.created". Dotted, so listeners can match by prefix. */
|
||||
readonly type: string;
|
||||
/** The emitting module. */
|
||||
/** The unique event id (x-event-id) — at-least-once delivery means a handler must dedup on it. */
|
||||
readonly id: string;
|
||||
/** The emitting module, context or node (x-source). */
|
||||
readonly source: string;
|
||||
/** The node it was emitted from. */
|
||||
/** The node it was emitted from (x-node). */
|
||||
readonly node: string;
|
||||
/** ISO-8601 emit time. */
|
||||
/** RFC-3339 emit time (x-time). */
|
||||
readonly at: string;
|
||||
/** The event or command that caused this one, if any (x-causation-id) — tracing. */
|
||||
readonly causationId?: string;
|
||||
/** A version tag for the body's shape, if the emitter set one (x-schema). */
|
||||
readonly schema?: string;
|
||||
readonly body: T;
|
||||
}
|
||||
|
||||
/** Extra provenance a caller may attach when emitting. */
|
||||
export interface EmitOptions {
|
||||
/** The event or command that caused this one (x-causation-id). */
|
||||
readonly causationId?: string;
|
||||
/** A version tag for the body's shape (x-schema). */
|
||||
readonly schema?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Emit an event. Source and node come from the environment the runtime set for the module
|
||||
* (MESH_MODULE, MESH_NODE), so a module names only the type and the body.
|
||||
* (MESH_MODULE, MESH_NODE), so a module names only the type and the body; the sdk stamps the
|
||||
* ADR 0047 headers (id, source, node, time) and the runtime rides them on the broker.
|
||||
*/
|
||||
export async function emit<T>(type: string, body: T): Promise<void> {
|
||||
const event: Event<T> = {
|
||||
type,
|
||||
source: process.env.MESH_MODULE ?? "unknown",
|
||||
node: process.env.MESH_NODE ?? "unknown",
|
||||
at: new Date().toISOString(),
|
||||
body,
|
||||
export async function emit<T>(type: string, body: T, opts: EmitOptions = {}): Promise<void> {
|
||||
const source = process.env.MESH_MODULE ?? "unknown";
|
||||
const node = process.env.MESH_NODE ?? "unknown";
|
||||
const headers: EventHeaders = {
|
||||
"x-event-id": randomUUID(),
|
||||
"x-source": source,
|
||||
"x-node": node,
|
||||
"x-time": new Date().toISOString(),
|
||||
"content-type": "application/json",
|
||||
...(opts.causationId ? { "x-causation-id": opts.causationId } : {}),
|
||||
...(opts.schema ? { "x-schema": opts.schema } : {}),
|
||||
};
|
||||
await broker().publish<Event<T>>({ key: type, node: event.node, body: event });
|
||||
await broker().publish<T>({ key: type, node, body, headers });
|
||||
}
|
||||
|
||||
/**
|
||||
* React to events whose type matches a topic pattern (`*` one segment, `#` any). The audit logger
|
||||
* is just `on("#", …)`. The handler receives the whole event, metadata included.
|
||||
* is just `on("#", …)`. The handler receives the reconstructed event — its metadata read back from
|
||||
* the headers, its body the domain payload.
|
||||
*/
|
||||
export async function on<T>(pattern: string, handler: (event: Event<T>) => Promise<void>): Promise<() => void> {
|
||||
return broker().subscribe<Event<T>>(pattern, async (envelope) => {
|
||||
await handler(envelope.body);
|
||||
return broker().subscribe<T>(pattern, async (envelope) => {
|
||||
await handler(fromEnvelope<T>(envelope));
|
||||
});
|
||||
}
|
||||
|
||||
/** Rebuild the Event a handler sees from a broker envelope's headers (ADR 0047) and body. */
|
||||
function fromEnvelope<T>(env: Envelope<T>): Event<T> {
|
||||
const h = env.headers ?? ({} as EventHeaders);
|
||||
return {
|
||||
type: env.key,
|
||||
id: h["x-event-id"] ?? "",
|
||||
source: h["x-source"] ?? "unknown",
|
||||
node: h["x-node"] ?? env.node ?? "unknown",
|
||||
at: h["x-time"] ?? "",
|
||||
causationId: h["x-causation-id"],
|
||||
schema: h["x-schema"],
|
||||
body: env.body,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -3,9 +3,9 @@
|
||||
// binding is provided by the runtime that hosts a module's code — the sdk defines the contract so
|
||||
// module code, the tool runtime and provisioners all speak it the same way.
|
||||
|
||||
import type { Envelope } from "../contracts/index.js";
|
||||
import type { Envelope, EventHeaders } from "../contracts/index.js";
|
||||
|
||||
export type { Envelope };
|
||||
export type { Envelope, EventHeaders };
|
||||
|
||||
/** A request/reply call and a publish/subscribe surface over the mesh broker. */
|
||||
export interface Broker {
|
||||
|
||||
+4
-36
@@ -1,42 +1,10 @@
|
||||
// Small, stable primitives every module's code may need. No behaviour here changes when a module
|
||||
// changes; that is the whole point of it living in the sdk.
|
||||
|
||||
import { createCipheriv, createDecipheriv, randomBytes, scryptSync } from "node:crypto";
|
||||
|
||||
// --- sealing ---
|
||||
//
|
||||
// A secret is sealed to a key so a copy of it at rest is not a working credential. AES-256-GCM;
|
||||
// the key is derived from a per-node passphrase the host holds. The host unseals on the machine;
|
||||
// nothing else does (novox/hq ADR 0043's link is the boundary — the sdk only carries the mechanism).
|
||||
|
||||
const MAGIC = "msk1"; // versions the sealed format, so it can change without silent misreads
|
||||
|
||||
/** Seal plaintext to a passphrase. Returns `msk1:<salt>:<iv>:<tag>:<ciphertext>`, base64 parts. */
|
||||
export function seal(plaintext: string, passphrase: string): string {
|
||||
const salt = randomBytes(16);
|
||||
const iv = randomBytes(12);
|
||||
const key = scryptSync(passphrase, salt, 32);
|
||||
const cipher = createCipheriv("aes-256-gcm", key, iv);
|
||||
const enc = Buffer.concat([cipher.update(plaintext, "utf8"), cipher.final()]);
|
||||
const tag = cipher.getAuthTag();
|
||||
return [MAGIC, b64(salt), b64(iv), b64(tag), b64(enc)].join(":");
|
||||
}
|
||||
|
||||
/** Unseal what seal produced. Throws — loudly — on any tamper or wrong key. */
|
||||
export function unseal(sealed: string, passphrase: string): string {
|
||||
const parts = sealed.split(":");
|
||||
if (parts.length !== 5 || parts[0] !== MAGIC) {
|
||||
throw new Error("not a sealed value this version understands");
|
||||
}
|
||||
const [, salt, iv, tag, enc] = parts.map((p, i) => (i === 0 ? Buffer.alloc(0) : ub64(p)));
|
||||
const key = scryptSync(passphrase, salt, 32);
|
||||
const decipher = createDecipheriv("aes-256-gcm", key, iv);
|
||||
decipher.setAuthTag(tag);
|
||||
return Buffer.concat([decipher.update(enc), decipher.final()]).toString("utf8");
|
||||
}
|
||||
|
||||
const b64 = (b: Buffer): string => b.toString("base64url");
|
||||
const ub64 = (s: string): Buffer => Buffer.from(s, "base64url");
|
||||
// There was a symmetric seal()/unseal() here, for a provider to seal a credential to a key before
|
||||
// writing it. It is gone: a provider is handed the credential the mesh minted and seals nothing
|
||||
// (novox/hq ADR 0053), the mesh's own secret delivery is asymmetric and belongs to the host, and
|
||||
// nothing else called it. The primitive left with the provisioner that was its only caller.
|
||||
|
||||
// --- semver ---
|
||||
//
|
||||
|
||||
+99
-76
@@ -1,74 +1,109 @@
|
||||
// The reconcile harness every provider shares. The loop below is identical for postgres, redis,
|
||||
// minio and umami: read the grants the control plane has written, bring each requested resource to
|
||||
// existence through the provider's adapter, seal and hand back the credential, and remove what is
|
||||
// no longer requested. A module writes ONLY the adapter — the per-service half — which is why this
|
||||
// lives in the sdk and the adapter lives in the module (novox/hq ADR 0044).
|
||||
// minio and umami: read the contributions the mesh delivered, bring each consumer's resource into
|
||||
// being through the provider's adapter under the login and password the mesh minted, and withdraw
|
||||
// what the mesh no longer asks for. A module writes ONLY the adapter — the per-service half — which
|
||||
// is why this lives in the sdk and the adapter lives in the module (novox/hq ADR 0044).
|
||||
//
|
||||
// **A provider creates the credential the mesh minted, and seals nothing (novox/hq ADR 0053).** The
|
||||
// control plane mints one password per consumer and seals it to this node; the host unseals it into
|
||||
// the file the contribution names. The provider does not generate a password, does not seal one, and
|
||||
// does not hand one back — the consumer already receives its copy through the mesh's own channel.
|
||||
// This is why the old $MESH_SEAL_KEY, the symmetric seal, and the *.grant.json / *.credential files
|
||||
// are gone: they described a second, contradictory credential path the mesh does not have.
|
||||
|
||||
import { readdir, readFile, writeFile, rm } from "node:fs/promises";
|
||||
import { join } from "node:path";
|
||||
import type { Grant, Credential } from "../contracts/index.js";
|
||||
import { seal } from "../primitives/index.js";
|
||||
import { readFile } from "node:fs/promises";
|
||||
|
||||
/** What a provider implements — the only per-service code. It is handed a Grant and creates (or
|
||||
* removes) the resource in its own software, returning the credential a consumer receives. */
|
||||
/** One consumer's resource to bring into being — everything the mesh derived and delivered. */
|
||||
export interface Provision {
|
||||
/** The login the mesh derived and gave the consumer to present. The provider creates exactly
|
||||
* this name — a name the consumer cannot learn is a name it cannot authenticate with. */
|
||||
readonly as: string;
|
||||
/** The password the mesh minted for this consumer, read from the file the mesh sealed to this
|
||||
* node and the host unsealed. The provider sets it; it never invents one. */
|
||||
readonly password: string;
|
||||
/** What the consumer contributed, per the interface's spec keys (e.g. `{ name: "umami" }`). */
|
||||
readonly values: Readonly<Record<string, unknown>>;
|
||||
/** Where the consumer is, for a provider that must reach back to it. Usually unused. */
|
||||
readonly at?: string;
|
||||
/** The consumer node, for logging and lifecycle events. */
|
||||
readonly consumer?: string;
|
||||
}
|
||||
|
||||
/** What a provider implements — the only per-service code. It is handed the login and password the
|
||||
* mesh made and brings the resource into being under them, or withdraws it. It returns nothing:
|
||||
* the credential is the mesh's, already delivered to the consumer. */
|
||||
export interface Adapter {
|
||||
create(grant: Grant): Promise<Credential>;
|
||||
remove(grant: Grant): Promise<void>;
|
||||
create(p: Provision): Promise<void>;
|
||||
remove(p: { readonly as: string }): Promise<void>;
|
||||
}
|
||||
|
||||
export interface ProvisionerOptions {
|
||||
/** Directory the control plane writes grant requests into and the harness writes credentials to.
|
||||
* Defaults to $GRANTS. */
|
||||
grants?: string;
|
||||
/** Passphrase a credential is sealed to before it is written. Defaults to $MESH_SEAL_KEY. */
|
||||
sealKey?: string;
|
||||
/** The contributions file the mesh writes for this resource (the module's `receives` path).
|
||||
* Defaults to $MESH_RECEIVES. */
|
||||
receives?: string;
|
||||
/** Reconcile interval in ms. Defaults to 5000. */
|
||||
everyMs?: number;
|
||||
}
|
||||
|
||||
export type { Grant, Credential };
|
||||
/** One entry in the mesh's contributions file: a consumer the provider must serve. */
|
||||
interface Contribution {
|
||||
readonly as: string;
|
||||
readonly secret: string;
|
||||
readonly node?: string;
|
||||
readonly at?: string;
|
||||
readonly values?: Readonly<Record<string, unknown>>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Run the reconcile loop for one provided resource. Returns a stop function. Never throws for a
|
||||
* single bad grant — it logs and keeps converging, because one consumer's failure must not stop
|
||||
* the others' provisioning.
|
||||
* single bad contribution — it logs and keeps converging, because one consumer's failure must not
|
||||
* stop the others' provisioning.
|
||||
*/
|
||||
export function runProvisioner(resource: string, adapter: Adapter, opts: ProvisionerOptions = {}): () => void {
|
||||
const dir = opts.grants ?? envOrThrow("GRANTS");
|
||||
const sealKey = opts.sealKey ?? envOrThrow("MESH_SEAL_KEY");
|
||||
const receives = opts.receives ?? envOrThrow("MESH_RECEIVES");
|
||||
const everyMs = opts.everyMs ?? 5000;
|
||||
|
||||
const applied = new Map<string, string>(); // consumer -> hash of the grant last applied
|
||||
const applied = new Map<string, string>(); // login (`as`) -> hash of what was last applied
|
||||
let stopped = false;
|
||||
|
||||
async function reconcile(): Promise<void> {
|
||||
const requested = await readGrants(dir, resource);
|
||||
const wantById = new Map(requested.map((g) => [g.consumer, g]));
|
||||
const given = await readContributions(receives, resource);
|
||||
const wantByAs = new Map(given.map((g) => [g.as, g]));
|
||||
|
||||
// Create or update anything requested whose grant has changed.
|
||||
for (const grant of requested) {
|
||||
const h = hash(grant);
|
||||
if (applied.get(grant.consumer) === h) continue;
|
||||
// Create or update every consumer whose login, password or values changed.
|
||||
for (const g of given) {
|
||||
let password: string;
|
||||
try {
|
||||
const cred = await adapter.create(grant);
|
||||
await writeSealedCredential(dir, grant, cred, sealKey);
|
||||
applied.set(grant.consumer, h);
|
||||
// The file the mesh sealed to this node, unsealed by the host into plaintext. Trailing
|
||||
// newline trimmed: a sealed value is exactly the secret, and a host writing a file may add
|
||||
// one.
|
||||
password = (await readFile(g.secret, "utf8")).replace(/\n$/, "");
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:${resource}] ${grant.consumer}: create failed, will retry: ${err}`);
|
||||
// The contribution names a secret the host has not written yet — a normal race on the
|
||||
// first pass. Skipped, retried on the next tick, never fatal.
|
||||
console.error(`[provisioner:${resource}] ${g.as}: secret not readable yet (${g.secret}): ${err}`);
|
||||
continue;
|
||||
}
|
||||
const h = hash(g.as, password, g.values ?? {});
|
||||
if (applied.get(g.as) === h) continue;
|
||||
try {
|
||||
await adapter.create({ as: g.as, password, values: g.values ?? {}, at: g.at, consumer: g.node });
|
||||
applied.set(g.as, h);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:${resource}] ${g.as}: create failed, will retry: ${err}`);
|
||||
}
|
||||
}
|
||||
|
||||
// Remove anything applied that is no longer requested. Only the harness's own outputs are
|
||||
// touched — the credential file it wrote — never anything it did not create.
|
||||
for (const consumer of [...applied.keys()]) {
|
||||
if (wantById.has(consumer)) continue;
|
||||
const grant = lastGrant(dir, resource, consumer);
|
||||
// Withdraw every login the provider made that the mesh no longer asks for. The mesh drops a
|
||||
// consumer from the file when it goes away; that is how a login is withdrawn rather than kept
|
||||
// working for ever. Only logins this harness created are touched.
|
||||
for (const as of [...applied.keys()]) {
|
||||
if (wantByAs.has(as)) continue;
|
||||
try {
|
||||
await adapter.remove(grant);
|
||||
await rm(credentialPath(dir, resource, consumer), { force: true });
|
||||
applied.delete(consumer);
|
||||
await adapter.remove({ as });
|
||||
applied.delete(as);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:${resource}] ${consumer}: remove failed, will retry: ${err}`);
|
||||
console.error(`[provisioner:${resource}] ${as}: remove failed, will retry: ${err}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -89,47 +124,35 @@ export function runProvisioner(resource: string, adapter: Adapter, opts: Provisi
|
||||
};
|
||||
}
|
||||
|
||||
// --- grant/credential files ---
|
||||
|
||||
async function readGrants(dir: string, resource: string): Promise<Grant[]> {
|
||||
let names: string[];
|
||||
/**
|
||||
* Read the mesh's contributions file for one resource. Absent or empty means "no consumer asks for
|
||||
* this" — the file is always written, so a provider can tell that from "the mesh never wrote it".
|
||||
*/
|
||||
async function readContributions(path: string, resource: string): Promise<Contribution[]> {
|
||||
let raw: string;
|
||||
try {
|
||||
names = await readdir(dir);
|
||||
raw = await readFile(path, "utf8");
|
||||
} catch {
|
||||
return []; // not written yet, or nobody serves anything here — nothing to converge to
|
||||
}
|
||||
let doc: { requirement?: string; given?: Contribution[] };
|
||||
try {
|
||||
doc = JSON.parse(raw) as { requirement?: string; given?: Contribution[] };
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:${resource}] contributions file is not JSON (${path}): ${err}`);
|
||||
return [];
|
||||
}
|
||||
const out: Grant[] = [];
|
||||
for (const name of names) {
|
||||
if (!name.endsWith(".grant.json")) continue;
|
||||
try {
|
||||
const raw = await readFile(join(dir, name), "utf8");
|
||||
const g = JSON.parse(raw) as Grant;
|
||||
if (g.resource === resource) out.push(g);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:${resource}] unreadable grant ${name}: ${err}`);
|
||||
if (doc.requirement && doc.requirement !== resource) {
|
||||
console.error(`[provisioner:${resource}] ${path} is for ${doc.requirement}, not ${resource}`);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
return out;
|
||||
// A contribution with no `as` is not a credential grant (a module offering something on its own
|
||||
// machine) — the provisioner has nothing to create for it.
|
||||
return (doc.given ?? []).filter((g) => g.as && g.secret);
|
||||
}
|
||||
|
||||
function credentialPath(dir: string, resource: string, consumer: string): string {
|
||||
return join(dir, `${consumer}.${resource}.credential`);
|
||||
}
|
||||
|
||||
async function writeSealedCredential(dir: string, grant: Grant, cred: Credential, key: string): Promise<void> {
|
||||
const body = JSON.stringify({ resource: grant.resource, consumer: grant.consumer, fields: cred.fields });
|
||||
await writeFile(credentialPath(dir, grant.resource, grant.consumer), seal(body, key), { mode: 0o600 });
|
||||
}
|
||||
|
||||
function lastGrant(dir: string, resource: string, consumer: string): Grant {
|
||||
// For removal the harness needs a Grant to hand the adapter; the identity is enough to act on.
|
||||
return { resource, consumer, node: "", values: {} };
|
||||
}
|
||||
|
||||
// --- helpers ---
|
||||
|
||||
function hash(g: Grant): string {
|
||||
return JSON.stringify([g.resource, g.consumer, g.node, g.values]);
|
||||
function hash(as: string, password: string, values: Readonly<Record<string, unknown>>): string {
|
||||
return JSON.stringify([as, password, values]);
|
||||
}
|
||||
|
||||
function envOrThrow(name: string): string {
|
||||
|
||||
+32
-9
@@ -66,20 +66,43 @@ export function listTools(env: NodeJS.ProcessEnv = process.env): { module: strin
|
||||
* other — two tools answering one name is a fault, not a race to resolve.
|
||||
*/
|
||||
export async function serveTools(broker: Broker, env: NodeJS.ProcessEnv = process.env): Promise<() => void> {
|
||||
const byName = new Map<string, ToolDefinition>();
|
||||
// Each tool is served on its own key, namespaced by its module (novox/hq ADR 0052): a caller
|
||||
// invokes `<module>.<tool>`, only the module that serves it answers, and the module's account is
|
||||
// scoped to serve.<module>.* — so one module cannot answer another's calls. The module is the
|
||||
// namespace, so a tool name need only be unique within its module, not across the whole mesh.
|
||||
const stops: Array<() => void> = [];
|
||||
for (const { module, tools } of collectTools(env)) {
|
||||
const seen = new Set<string>();
|
||||
for (const t of tools) {
|
||||
if (byName.has(t.name)) {
|
||||
throw new Error(`tool name ${t.name} is exposed by two modules (one is ${module}) — refused`);
|
||||
if (seen.has(t.name)) {
|
||||
throw new Error(`${module} exposes two tools named ${t.name} — refused`);
|
||||
}
|
||||
byName.set(t.name, t);
|
||||
seen.add(t.name);
|
||||
const stop = await broker.handle<Readonly<Record<string, unknown>>, unknown>(
|
||||
toolKey(module, t.name),
|
||||
(args) => t.run(args ?? {}),
|
||||
);
|
||||
stops.push(stop);
|
||||
}
|
||||
}
|
||||
return broker.handle<Invocation, unknown>("tools.invoke", async (inv) => {
|
||||
const tool = byName.get(inv.tool);
|
||||
if (!tool) throw new Error(`no such tool: ${inv.tool}`);
|
||||
return tool.run(inv.args ?? {});
|
||||
});
|
||||
return () => {
|
||||
for (const stop of stops) stop();
|
||||
};
|
||||
}
|
||||
|
||||
/** The broker key a tool is served on and invoked by — the module namespaces the tool (ADR 0052). */
|
||||
export function toolKey(module: string, tool: string): string {
|
||||
return `${module}.${tool}`;
|
||||
}
|
||||
|
||||
/** Invoke a module's tool over the broker — the caller's side of serving. */
|
||||
export async function invokeTool(
|
||||
broker: Broker,
|
||||
module: string,
|
||||
tool: string,
|
||||
args: Readonly<Record<string, unknown>> = {},
|
||||
): Promise<unknown> {
|
||||
return broker.request(toolKey(module, tool), args);
|
||||
}
|
||||
|
||||
/** Testing/inspection: drop all registrations. */
|
||||
|
||||
+46
-39
@@ -4,19 +4,12 @@ import { mkdtemp, writeFile, readFile, readdir } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
|
||||
import { seal, unseal, compareVersions } from "../dist/primitives/index.js";
|
||||
import { registerModuleTools, collectTools, resetTools, serveTools, listTools } from "../dist/tools/index.js";
|
||||
import { runProvisioner, type Grant } from "../dist/provisioner/index.js";
|
||||
import { compareVersions } from "../dist/primitives/index.js";
|
||||
import { registerModuleTools, collectTools, resetTools, serveTools, listTools, invokeTool } from "../dist/tools/index.js";
|
||||
import { runProvisioner } from "../dist/provisioner/index.js";
|
||||
import { emit, on, type Event } from "../dist/events/index.js";
|
||||
import { useBroker } from "../dist/messaging/index.js";
|
||||
|
||||
test("seal round-trips and rejects the wrong key", () => {
|
||||
const sealed = seal("hunter2", "node-key");
|
||||
assert.equal(unseal(sealed, "node-key"), "hunter2");
|
||||
assert.notEqual(sealed, "hunter2");
|
||||
assert.throws(() => unseal(sealed, "wrong-key"));
|
||||
});
|
||||
|
||||
test("semver orders releases", () => {
|
||||
assert.equal(compareVersions("1.2.3", "1.2.10"), -1);
|
||||
assert.equal(compareVersions("2.0.0", "1.9.9"), 1);
|
||||
@@ -38,39 +31,46 @@ test("module tools register and collect, a thrower is skipped not fatal", () =>
|
||||
assert.equal(broken?.tools.length, 0);
|
||||
});
|
||||
|
||||
test("provisioner creates a sealed credential for a grant, then removes on withdrawal", async () => {
|
||||
test("provisioner creates each consumer with the mesh's login and password, removes on withdrawal", async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), "prov-"));
|
||||
const created: string[] = [];
|
||||
const created: { as: string; password: string; values: unknown }[] = [];
|
||||
const removed: string[] = [];
|
||||
|
||||
const grant: Grant = { resource: "analytics", consumer: "webapp", node: "anchor", values: { name: "webapp" } };
|
||||
await writeFile(join(dir, "webapp.grant.json"), JSON.stringify(grant));
|
||||
// What the mesh delivers: the password it minted (as the host leaves it after unsealing) and a
|
||||
// contributions file naming the consumer's login and where that password is.
|
||||
await writeFile(join(dir, "webapp.secret"), "minted-pw\n");
|
||||
const receives = join(dir, "analytics.json");
|
||||
const doc = (given: unknown[]): string =>
|
||||
JSON.stringify({ contributions: 1, requirement: "analytics", given });
|
||||
await writeFile(
|
||||
receives,
|
||||
doc([{ from: "webapp", node: "anchor", as: "webapp-anchor", secret: join(dir, "webapp.secret"), values: { name: "webapp" } }]),
|
||||
);
|
||||
|
||||
const stop = runProvisioner(
|
||||
"analytics",
|
||||
{
|
||||
async create(g) {
|
||||
created.push(g.consumer);
|
||||
return { fields: { siteId: "abc", snippet: "<script>", dashboard: "https://x/webapp" } };
|
||||
async create(p) {
|
||||
created.push({ as: p.as, password: p.password, values: p.values });
|
||||
},
|
||||
async remove(g) {
|
||||
removed.push(g.consumer);
|
||||
async remove(p) {
|
||||
removed.push(p.as);
|
||||
},
|
||||
},
|
||||
{ grants: dir, sealKey: "k", everyMs: 20 },
|
||||
{ receives, everyMs: 20 },
|
||||
);
|
||||
|
||||
await waitFor(() => created.length === 1, 2000);
|
||||
const files = await readdir(dir);
|
||||
const credFile = files.find((f) => f.endsWith(".credential"));
|
||||
assert.ok(credFile, "a credential file was written");
|
||||
const sealed = await readFile(join(dir, credFile!), "utf8");
|
||||
const body = JSON.parse(unseal(sealed, "k")) as { fields: Record<string, string> };
|
||||
assert.equal(body.fields.siteId, "abc");
|
||||
// The adapter is handed the mesh's login and the mesh's password — not one it generated, and the
|
||||
// trailing newline of the unsealed file is stripped.
|
||||
assert.equal(created[0].as, "webapp-anchor");
|
||||
assert.equal(created[0].password, "minted-pw");
|
||||
assert.deepEqual(created[0].values, { name: "webapp" });
|
||||
|
||||
// Withdraw the grant → the harness removes via the adapter and deletes the credential.
|
||||
await (await import("node:fs/promises")).rm(join(dir, "webapp.grant.json"));
|
||||
// Withdraw: the mesh drops the consumer from the file → the harness removes it via the adapter.
|
||||
await writeFile(receives, doc([]));
|
||||
await waitFor(() => removed.length === 1, 2000);
|
||||
assert.equal(removed[0], "webapp-anchor");
|
||||
stop();
|
||||
});
|
||||
|
||||
@@ -106,26 +106,28 @@ test("modules can SERVE: a real async tool, loaded and invoked over the broker",
|
||||
const broker = memBroker();
|
||||
const stop = await serveTools(broker, {});
|
||||
|
||||
// Invoke it the way a caller (mesh-control's command API) would — over the broker, by name.
|
||||
const result = await broker.request<{ tool: string; args: Record<string, unknown> }, { snippet: string }>(
|
||||
"tools.invoke",
|
||||
{ tool: "create_site", args: { domain: "my-app" } },
|
||||
);
|
||||
// Invoke it the way a caller (mesh-control's command API) would — over the broker, by module and
|
||||
// tool, each served on its own key `demo.create_site` (novox/hq ADR 0052).
|
||||
const result = (await invokeTool(broker, "demo", "create_site", { domain: "my-app" })) as { snippet: string };
|
||||
assert.match(result.snippet, /data-website-id="site-123"/);
|
||||
|
||||
// Discovery works, and an unknown tool is refused rather than silently dropped.
|
||||
assert.deepEqual(listTools({}).map((t) => t.name), ["create_site"]);
|
||||
await assert.rejects(broker.request("tools.invoke", { tool: "nope", args: {} }));
|
||||
await assert.rejects(invokeTool(broker, "demo", "nope", {}));
|
||||
|
||||
stop();
|
||||
server.close();
|
||||
});
|
||||
|
||||
test("serving refuses two modules exposing one tool name", async () => {
|
||||
test("serving refuses one module exposing two tools of the same name", async () => {
|
||||
// Two *modules* may share a tool name — each is served on its own `module.tool` key (ADR 0052).
|
||||
// What is refused is one module exposing the same name twice, where the key would collide.
|
||||
resetTools();
|
||||
registerModuleTools("a", () => [{ name: "dup", description: "", input: {}, run: async () => 1 }]);
|
||||
registerModuleTools("b", () => [{ name: "dup", description: "", input: {}, run: async () => 2 }]);
|
||||
await assert.rejects(serveTools(memBroker(), {}), /exposed by two modules/);
|
||||
registerModuleTools("a", () => [
|
||||
{ name: "dup", description: "", input: {}, run: async () => 1 },
|
||||
{ name: "dup", description: "", input: {}, run: async () => 2 },
|
||||
]);
|
||||
await assert.rejects(serveTools(memBroker(), {}), /exposes two tools named dup/);
|
||||
});
|
||||
|
||||
// A minimal in-memory broker: request routes to a registered handle, and publish routes to every
|
||||
@@ -195,12 +197,17 @@ test("events: a module emits, a listener and the audit sink (#) both receive it,
|
||||
assert.deepEqual(heard.map((e) => e.type), ["module.umami.site.created"]);
|
||||
assert.deepEqual(audited.map((e) => e.type), ["module.umami.site.created", "module.plex.play.started"]);
|
||||
|
||||
// The metadata an audit trail needs is present.
|
||||
// The metadata an audit trail needs is present — read back from the ADR 0047 headers, not the body.
|
||||
const e = heard[0];
|
||||
assert.equal(e.source, "umami");
|
||||
assert.equal(e.node, "anchor");
|
||||
assert.equal((e.body as { domain: string }).domain, "my-app");
|
||||
assert.match(e.at, /^\d{4}-\d{2}-\d{2}T/);
|
||||
// Every event carries a unique id (x-event-id) — the handle a consumer dedups on.
|
||||
assert.ok(e.id, "event has an x-event-id");
|
||||
assert.notEqual(audited[0].id, audited[1].id);
|
||||
// The body is exactly the domain payload — provenance never leaks into it.
|
||||
assert.deepEqual(Object.keys(e.body as object), ["domain"]);
|
||||
|
||||
delete process.env.MESH_MODULE;
|
||||
delete process.env.MESH_NODE;
|
||||
|
||||
Reference in New Issue
Block a user