Files
mesh-sdk/test/sdk.test.ts
jochen 3192491df6 Brake counts only successful re-applies; only a timeout ends a pass
A lost consumer whose re-apply fails is retried at the next check with
its count unchanged, instead of waiting out a backoff meant for adapters
whose create and holds disagree. A check that fails for one consumer no
longer stops checking the consumers after it; only a timeout does.
2026-09-26 01:30:13 +02:00

485 lines
18 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("provisioner backs off when create does not satisfy holds, and says so", async () => {
const dir = await mkdtemp(join(tmpdir(), "prov-brake-"));
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 logged: string[] = [];
const original = console.error;
console.error = (...a: unknown[]) => logged.push(a.join(" "));
try {
const stop = runProvisioner(
"cache",
{
async create() {
creates++;
},
async remove() {},
async holds() {
asked++;
return false; // an adapter that disagrees with itself: nothing create does is ever held
},
},
{ receives, everyMs: 5, verifyEveryMs: 20 },
);
await new Promise((r) => setTimeout(r, 400));
stop();
} finally {
console.error = original;
}
// Without the brake this would be ~20 creates (one per 20ms check). With doubling waits
// (20, 40, 80, 160ms…) it is a handful.
assert.ok(creates >= 3 && creates <= 7, `creates: ${creates}`);
assert.equal(asked, creates - 1);
assert.ok(logged.some((l) => l.includes("create does not produce what holds checks")));
});
test("provisioner treats a failed or hung holds as could-not-ask, and never logs the password", async () => {
const dir = await mkdtemp(join(tmpdir(), "prov-timeout-"));
await writeFile(join(dir, "a.secret"), "s3cret-pw");
await writeFile(join(dir, "b.secret"), "other-pw");
const receives = join(dir, "cache.json");
await writeFile(
receives,
JSON.stringify({ requirement: "cache", given: [{ as: "a", secret: join(dir, "a.secret") }, { as: "b", secret: join(dir, "b.secret") }] }),
);
let creates = 0;
const askedFor: string[] = [];
const logged: string[] = [];
const original = console.error;
console.error = (...a: unknown[]) => logged.push(a.join(" "));
try {
const stop = runProvisioner(
"cache",
{
async create() {
creates++;
},
async remove() {},
holds(p) {
askedFor.push(p.as);
// "a" fails the way a shelled-out tool does, with its arguments in the message.
if (p.as === "a") return Promise.reject(new Error(`Command failed: tool --password ${p.password}`));
return Promise.resolve(true);
},
},
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 20 },
);
await new Promise((r) => setTimeout(r, 200));
stop();
} finally {
console.error = original;
}
assert.equal(creates, 2); // the first pass only; nothing is applied again on "could not ask"
// A failure that is not a timeout may be one consumer's alone: the consumers after it are still
// asked on every pass.
assert.ok(askedFor.filter((as) => as === "b").length >= 2, askedFor.join(","));
assert.ok(logged.some((l) => l.includes("Command failed: tool --password ***")));
assert.ok(!logged.some((l) => l.includes("s3cret-pw")));
// A check that never answers counts as could-not-ask once its time is up.
let hungCreates = 0;
const hungLogged: string[] = [];
console.error = (...a: unknown[]) => hungLogged.push(a.join(" "));
try {
const stopHung = runProvisioner(
"cache",
{
async create() {
hungCreates++;
},
async remove() {},
holds: () => new Promise<boolean>(() => {}),
},
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 20 },
);
await new Promise((r) => setTimeout(r, 200));
stopHung();
} finally {
console.error = original;
}
assert.equal(hungCreates, 2);
assert.ok(hungLogged.some((l) => l.includes("no answer within 20ms")));
});
test("provisioner stops a verify pass at a timeout, and retries a failed re-apply at the next check", async () => {
const dir = await mkdtemp(join(tmpdir(), "prov-retry-"));
await writeFile(join(dir, "a.secret"), "pa");
await writeFile(join(dir, "b.secret"), "pb");
const receives = join(dir, "cache.json");
await writeFile(
receives,
JSON.stringify({ requirement: "cache", given: [{ as: "a", secret: join(dir, "a.secret") }, { as: "b", secret: join(dir, "b.secret") }] }),
);
const original = console.error;
console.error = () => {};
try {
// A timeout on "a" means "b" is not asked in that pass.
const askedFor: string[] = [];
const stopHung = runProvisioner(
"cache",
{
async create() {},
async remove() {},
holds(p) {
askedFor.push(p.as);
return p.as === "a" ? new Promise<boolean>(() => {}) : Promise.resolve(true);
},
},
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 10 },
);
await new Promise((r) => setTimeout(r, 150));
stopHung();
assert.ok(askedFor.length >= 2 && askedFor.every((as) => as === "a"), askedFor.join(","));
// A lost consumer whose re-apply fails is retried at the next check, not held back by the brake,
// and the retry that succeeds counts as the first re-apply, not the second.
let held = true;
let createCalls = 0;
const logged: string[] = [];
console.error = (...a: unknown[]) => logged.push(a.join(" "));
const stop = runProvisioner(
"cache",
{
async create(p) {
if (p.as !== "a") return;
createCalls++;
if (createCalls === 2) throw new Error("write refused"); // the first re-apply fails
if (createCalls === 3) held = true; // the retry works
},
async remove() {},
async holds(p) {
return p.as === "a" ? held : true;
},
},
{ receives, everyMs: 5, verifyEveryMs: 20 },
);
await waitFor(() => createCalls === 1, 2000);
held = false; // the backend loses "a"
await waitFor(() => createCalls === 3, 2000);
await new Promise((r) => setTimeout(r, 100));
stop();
assert.equal(createCalls, 3);
assert.ok(!logged.some((l) => l.includes("create does not produce what holds checks")), logged.join("\n"));
} finally {
console.error = original;
}
});
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));
}
}