A check that hangs no longer stalls every consumer: it times out after 30s and counts as could-not-ask, and the rest of that pass is not asked. A consumer still not held after being applied again is checked at doubling intervals up to an hour, and said loudly, so an adapter whose create and holds disagree costs one re-apply an hour, not one a minute. The consumer's password is scrubbed from every error the harness logs.
421 lines
16 KiB
TypeScript
421 lines
16 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"
|
|
// After one consumer's check fails, the rest of that pass is not asked: "b" follows "a" and is
|
|
// never reached, because every pass stops at "a".
|
|
assert.ok(askedFor.length >= 2 && askedFor.every((as) => as === "a"), 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("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));
|
|
}
|
|
}
|