An optional holds() on the adapter is asked for every applied consumer every minute; false applies it again. A backend that forgets what it was given while the provisioner runs (hq issue 120) is healed within a minute instead of failing its consumers in silence. Unable to ask is not treated as loss. Adapters without holds() behave as before.
316 lines
12 KiB
TypeScript
316 lines
12 KiB
TypeScript
import { test } from "node:test";
|
|
import assert from "node:assert/strict";
|
|
import { mkdtemp, writeFile, readFile, readdir } from "node:fs/promises";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
|
|
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("semver orders releases", () => {
|
|
assert.equal(compareVersions("1.2.3", "1.2.10"), -1);
|
|
assert.equal(compareVersions("2.0.0", "1.9.9"), 1);
|
|
assert.equal(compareVersions("v1.0.0", "1.0.0"), 0);
|
|
});
|
|
|
|
test("module tools register and collect, a thrower is skipped not fatal", () => {
|
|
resetTools();
|
|
registerModuleTools("umami", () => [
|
|
{ name: "umami_stats", description: "d", input: {}, run: async () => 1 },
|
|
]);
|
|
registerModuleTools("broken", () => {
|
|
throw new Error("no token");
|
|
});
|
|
const collected = collectTools({});
|
|
const umami = collected.find((c) => c.module === "umami");
|
|
const broken = collected.find((c) => c.module === "broken");
|
|
assert.equal(umami?.tools.length, 1);
|
|
assert.equal(broken?.tools.length, 0);
|
|
});
|
|
|
|
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: { as: string; password: string; values: unknown }[] = [];
|
|
const removed: string[] = [];
|
|
|
|
// 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(p) {
|
|
created.push({ as: p.as, password: p.password, values: p.values });
|
|
},
|
|
async remove(p) {
|
|
removed.push(p.as);
|
|
},
|
|
},
|
|
{ receives, everyMs: 20 },
|
|
);
|
|
|
|
await waitFor(() => created.length === 1, 2000);
|
|
// 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 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();
|
|
});
|
|
|
|
test("provisioner applies again what the backend no longer holds, and trusts memory without holds", async () => {
|
|
const dir = await mkdtemp(join(tmpdir(), "prov-holds-"));
|
|
await writeFile(join(dir, "webapp.secret"), "minted-pw\n");
|
|
const receives = join(dir, "cache.json");
|
|
await writeFile(
|
|
receives,
|
|
JSON.stringify({ contributions: 1, requirement: "cache", given: [{ from: "webapp", node: "anchor", as: "webapp-anchor", secret: join(dir, "webapp.secret") }] }),
|
|
);
|
|
|
|
// A backend that forgets: what was created is held until it "restarts".
|
|
const backend = new Set<string>();
|
|
let creates = 0;
|
|
let asked = 0;
|
|
const stop = runProvisioner(
|
|
"cache",
|
|
{
|
|
async create(p) {
|
|
creates++;
|
|
backend.add(`${p.as}:${p.password}`);
|
|
},
|
|
async remove(p) {
|
|
for (const k of backend) if (k.startsWith(`${p.as}:`)) backend.delete(k);
|
|
},
|
|
async holds(p) {
|
|
asked++;
|
|
return backend.has(`${p.as}:${p.password}`);
|
|
},
|
|
},
|
|
{ receives, everyMs: 10, verifyEveryMs: 30 },
|
|
);
|
|
|
|
await waitFor(() => creates === 1, 2000);
|
|
// While the backend holds it, asking changes nothing: no second create.
|
|
await waitFor(() => asked >= 2, 2000);
|
|
assert.equal(creates, 1);
|
|
|
|
// The backend restarts and forgets. The contributions did not change; only asking can notice.
|
|
backend.clear();
|
|
await waitFor(() => creates === 2, 2000);
|
|
assert.ok(backend.has("webapp-anchor:minted-pw"));
|
|
stop();
|
|
|
|
// An adapter without holds is trusted from memory, as before: a forgotten backend stays forgotten.
|
|
const forgetful = new Set<string>();
|
|
let plainCreates = 0;
|
|
const stopPlain = runProvisioner(
|
|
"cache",
|
|
{
|
|
async create(p) {
|
|
plainCreates++;
|
|
forgetful.add(p.as);
|
|
},
|
|
async remove() {},
|
|
},
|
|
{ receives, everyMs: 10, verifyEveryMs: 10 },
|
|
);
|
|
await waitFor(() => plainCreates === 1, 2000);
|
|
forgetful.clear();
|
|
await new Promise((r) => setTimeout(r, 100));
|
|
assert.equal(plainCreates, 1);
|
|
stopPlain();
|
|
});
|
|
|
|
test("provisioner keeps a consumer applied when the backend cannot be asked", async () => {
|
|
const dir = await mkdtemp(join(tmpdir(), "prov-unreachable-"));
|
|
await writeFile(join(dir, "webapp.secret"), "minted-pw");
|
|
const receives = join(dir, "cache.json");
|
|
await writeFile(
|
|
receives,
|
|
JSON.stringify({ requirement: "cache", given: [{ as: "webapp-anchor", secret: join(dir, "webapp.secret") }] }),
|
|
);
|
|
let creates = 0;
|
|
let asked = 0;
|
|
const stop = runProvisioner(
|
|
"cache",
|
|
{
|
|
async create() {
|
|
creates++;
|
|
},
|
|
async remove() {},
|
|
async holds() {
|
|
asked++;
|
|
throw new Error("connection refused");
|
|
},
|
|
},
|
|
{ receives, everyMs: 10, verifyEveryMs: 20 },
|
|
);
|
|
await waitFor(() => asked >= 3, 2000);
|
|
// Unable to ask is not evidence of loss: nothing is applied again.
|
|
assert.equal(creates, 1);
|
|
stop();
|
|
});
|
|
|
|
test("modules can SERVE: a real async tool, loaded and invoked over the broker", async () => {
|
|
resetTools();
|
|
|
|
// Stand up a fake upstream the tool actually calls over the network — proving a tool that does
|
|
// real work (not a pure function) serves end to end.
|
|
const { createServer } = await import("node:http");
|
|
const server = createServer((_req, res) => {
|
|
res.setHeader("content-type", "application/json");
|
|
res.end(JSON.stringify({ id: "site-123" }));
|
|
});
|
|
await new Promise<void>((r) => server.listen(0, r));
|
|
const addr = server.address() as { port: number };
|
|
const base = `http://127.0.0.1:${addr.port}`;
|
|
|
|
// A module registers its tool exactly as umami does — the tool calls the upstream and returns a
|
|
// result computed from it.
|
|
registerModuleTools("demo", () => [
|
|
{
|
|
name: "create_site",
|
|
description: "register a site and return its snippet",
|
|
input: { domain: { type: "string" } },
|
|
run: async (args) => {
|
|
const res = await fetch(`${base}/api/websites`, { method: "POST" });
|
|
const { id } = (await res.json()) as { id: string };
|
|
return { domain: String(args.domain), snippet: `<script data-website-id="${id}"></script>` };
|
|
},
|
|
},
|
|
]);
|
|
|
|
const broker = memBroker();
|
|
const stop = await serveTools(broker, {});
|
|
|
|
// Invoke it the way a caller (mesh-controller's command API) would — over the broker, by module and
|
|
// tool, each served on its own key `demo.create_site` (novox/hq ADR 0047).
|
|
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(invokeTool(broker, "demo", "nope", {}));
|
|
|
|
stop();
|
|
server.close();
|
|
});
|
|
|
|
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 0047).
|
|
// 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 },
|
|
{ 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
|
|
// subscriber whose topic pattern matches. Enough to exercise tools and events; the real binding is
|
|
// the mesh's AMQP one, provided by the hosting runtime.
|
|
function memBroker() {
|
|
const handlers = new Map<string, (b: unknown) => Promise<unknown>>();
|
|
const subs: { pattern: string; handler: (env: { key: string; node: string; body: unknown }) => Promise<void> }[] = [];
|
|
return {
|
|
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
|
const h = handlers.get(key);
|
|
if (!h) throw new Error(`no handler for ${key}`);
|
|
return (await h(body)) as Res;
|
|
},
|
|
async handle<Req, Res>(key: string, handler: (b: Req) => Promise<Res>): Promise<() => void> {
|
|
handlers.set(key, handler as (b: unknown) => Promise<unknown>);
|
|
return () => handlers.delete(key);
|
|
},
|
|
async publish<T>(env: { key: string; node: string; body: T }): Promise<void> {
|
|
for (const s of subs) if (topicMatch(s.pattern, env.key)) await s.handler(env);
|
|
},
|
|
async subscribe<T>(pattern: string, handler: (env: { key: string; node: string; body: T }) => Promise<void>): Promise<() => void> {
|
|
const entry = { pattern, handler: handler as (env: { key: string; node: string; body: unknown }) => Promise<void> };
|
|
subs.push(entry);
|
|
return () => {
|
|
const i = subs.indexOf(entry);
|
|
if (i >= 0) subs.splice(i, 1);
|
|
};
|
|
},
|
|
async close(): Promise<void> {},
|
|
};
|
|
}
|
|
|
|
// AMQP topic matching: `*` one segment, `#` any run of segments.
|
|
function topicMatch(pattern: string, key: string): boolean {
|
|
if (pattern === "#") return true;
|
|
const p = pattern.split(".");
|
|
const k = key.split(".");
|
|
let pi = 0;
|
|
let ki = 0;
|
|
while (pi < p.length) {
|
|
if (p[pi] === "#") return true; // simplification: # only as a trailing wildcard
|
|
if (ki >= k.length) return false;
|
|
if (p[pi] !== "*" && p[pi] !== k[ki]) return false;
|
|
pi++;
|
|
ki++;
|
|
}
|
|
return ki === k.length;
|
|
}
|
|
|
|
test("events: a module emits, a listener and the audit sink (#) both receive it, with metadata", async () => {
|
|
const broker = memBroker();
|
|
useBroker(() => broker);
|
|
|
|
process.env.MESH_MODULE = "umami";
|
|
process.env.MESH_NODE = "anchor";
|
|
|
|
const heard: Event[] = [];
|
|
const audited: Event[] = [];
|
|
await on("module.umami.#", async (e) => void heard.push(e));
|
|
await on("#", async (e) => void audited.push(e)); // the audit logger is just this
|
|
|
|
await emit("module.umami.site.created", { domain: "my-app" });
|
|
await emit("module.plex.play.started", { title: "x" }); // a different module's event
|
|
|
|
// The umami listener heard only umami's event; the audit sink heard both.
|
|
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 — read back from the ADR 0042 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;
|
|
});
|
|
|
|
async function waitFor(cond: () => boolean, ms: number): Promise<void> {
|
|
const start = Date.now();
|
|
while (!cond()) {
|
|
if (Date.now() - start > ms) throw new Error("timed out waiting");
|
|
await new Promise((r) => setTimeout(r, 10));
|
|
}
|
|
}
|