/** * The mesh issues an assignment's subjects, and a runtime serves what it is issued (novox/hq ADR * 0160). A runtime derives one address for itself — `mesh.assignment..` — reads the * membership there, serves exactly what it says, and re-serves when a new one arrives. Against a * real bus with JetStream, because the membership is a direct get on a stream. * * docker run -d --rm --name t -p 14232:4222 nats:2.10-alpine -js * MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/membership.test.ts */ import assert from "node:assert/strict"; import { test } from "node:test"; import { fileURLToPath } from "node:url"; import { connect, StringCodec } from "nats"; import { resetTools } from "@novox/mesh-sdk/tools"; import { connectNats, membershipSubject } from "../dist/broker-nats.js"; import { callTool, subjectListed } from "../dist/client.js"; import { runTools } from "../dist/runtime.js"; const url = process.env.MESH_TEST_NATS; const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url)); const sc = StringCodec(); /** The ASSIGNMENTS stream as the controller asserts it: last-per-subject, readable by direct get. */ async function anAssignmentsStream() { const nc = await connect({ servers: url! }); const jsm = await nc.jetstreamManager(); try { await jsm.streams.delete("ASSIGNMENTS"); } catch { // none yet } await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true, } as never); return { async issue(m: object & { node: string; module: string }) { await nc.jetstream().publish(membershipSubject(m.node, m.module), sc.encode(JSON.stringify(m))); }, async close() { await nc.close(); }, }; } test("a runtime serves exactly the subjects it is issued, and the listing carries them", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); resetTools(); const stream = await anAssignmentsStream(); // The mesh placed shop on two machines that are not interchangeable: each is served by name only. await stream.issue({ node: "anchor", module: "shop", serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }], emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools", }); const shop = await connectNats({ url, node: "anchor", module: "shop" }); const asker = await connectNats({ url, module: "console", node: "workstation" }); let stop = () => {}; try { stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] }); assert.equal(shop.membership()?.module, "shop", "the runtime read what the mesh issued"); // By name it answers; the plain subject was not issued, so nothing serves it. const named = await callTool(asker, "shop.price@anchor", { of: "a hat" }); assert.deepEqual(named.result, { of: "a hat", cost: 12 }); assert.equal(named.node, "anchor"); await assert.rejects(callTool(asker, "shop.price", {}), /no responders|503/i, "a subject not issued is not served"); // The tools answer says where each is served, and the client composes nothing. const tools = await asker.request, { tools: { name: string; subjects?: string[] }[] }>( "shop.tools", {}, ); assert.deepEqual(tools.tools.find((x) => x.name === "price")?.subjects, ["mesh.mod.shop.tool.price.anchor"]); const listing = { tools: [{ module: "shop", name: "price", subjects: ["mesh.mod.shop.tool.price.anchor"] }], notAnswering: [] }; assert.equal(subjectListed("shop.price", "anchor", listing as never), "mesh.mod.shop.tool.price.anchor"); assert.equal(subjectListed("shop.price", "elsewhere", listing as never), undefined); // The mesh re-issues the membership with the plain subject too (the modules became // interchangeable); the runtime re-serves on it without a restart. await stream.issue({ node: "anchor", module: "shop", serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }], emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools", }); await new Promise((r) => setTimeout(r, 300)); const plain = await callTool(asker, "shop.price", { of: "a coat" }); assert.deepEqual(plain.result, { of: "a coat", cost: 12 }); assert.equal(plain.node, "anchor"); } finally { stop(); await asker.close(); await shop.close(); await stream.close(); } }); test("a seat's verbs are implemented under the seat's name, served where issued, never listed as the module's", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); resetTools(); const stream = await anAssignmentsStream(); await stream.issue({ node: "anchor", module: "postgres", serves: [{ subject: "mesh.mod.postgres.tool.{tool}", queue: "serve.postgres" }, { subject: "mesh.mod.postgres.tool.{tool}.anchor" }], seats: [{ seat: "mesh-store", verb: "databases", subject: "mesh.seat.mesh-store.tool.databases" }], emits: "mesh.mod.postgres.event.{event}", tools: "mesh.mod.postgres.tool.tools", }); const credential = { url, node: "anchor", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] }; const pg = await connectNats(credential); const asker = await connectNats({ url, module: "console", node: "workstation" }); let stop = () => {}; try { stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat.mjs")] }); // The module's own `databases` and the seat's are two tools: postgres's on its subject, the // store's on the seat's, each answering as itself. const own = await callTool(asker, "postgres.databases", {}); assert.deepEqual(own.result, { software: "postgres" }); const seat = await callTool(asker, "seat:mesh-store.databases", {}); assert.deepEqual(seat.result, { seat: "mesh-store" }); assert.equal(seat.node, "anchor"); // What the seat does not promise — creating a database — is postgres's tool and not the store's. await assert.rejects(callTool(asker, "seat:mesh-store.postgres_create_database", {}), /no responders|503/i); // And the module's `tools` lists only postgres's own, never the seat's implementation. const tools = await asker.request, { module: string; tools: { name: string }[] }>("postgres.tools", {}); assert.deepEqual(tools.tools.map((x) => x.name).sort(), ["databases", "postgres_create_database"]); await assert.rejects(asker.request("mesh-store.tools", {}), /no responders|503/i, "a seat is not a module with a tools verb"); } finally { stop(); await asker.close(); await pg.close(); await stream.close(); } }); // novox/hq issue 218: the store seat is claimed by every machine running postgres and held by one. // Where the mesh issued a membership without the seat, the claim in the credential serves nothing: // the module's own tools answer, the seat's verbs do not, and the runtime does not announce them. test("a claimant the membership does not make the holder serves none of the seat's verbs", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); resetTools(); const stream = await anAssignmentsStream(); await stream.issue({ node: "elsewhere", module: "postgres", serves: [{ subject: "mesh.mod.postgres.tool.{tool}.elsewhere" }], emits: "mesh.mod.postgres.event.{event}", tools: "mesh.mod.postgres.tool.tools", }); const credential = { url, node: "elsewhere", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] }; const pg = await connectNats(credential); const asker = await connectNats({ url, module: "console", node: "workstation" }); let stop = () => {}; try { stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-claimant.mjs")] }); assert.deepEqual((await callTool(asker, "postgres.databases@elsewhere", {})).result, { software: "postgres" }); await assert.rejects(callTool(asker, "seat:mesh-store.databases", {}), /no responders|503/i, "the store is not held here"); } finally { stop(); await asker.close(); await pg.close(); await stream.close(); } }); test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); resetTools(); const stream = await anAssignmentsStream(); await stream.issue({ node: "anchor", module: "shop", serves: [{ subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }, { subject: "mesh.mod.shop.tool.{tool}.anchor" }], seats: [{ seat: "shop", verb: "price", subject: "mesh.seat.shop.tool.price" }], emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools", }); const credential = { url, node: "anchor", module: "shop", claims: [{ seat: "shop", scope: "mesh", serves: ["price"] }] }; const shop = await connectNats(credential); const asker = await connectNats({ url, module: "console", node: "workstation" }); let stop = () => {}; try { stop = await runTools({ broker: shop, credential, moduleEntrypoints: [fixture("shop-seat.mjs")] }); assert.deepEqual((await callTool(asker, "shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 }); assert.deepEqual((await callTool(asker, "seat:shop.price", { of: "a hat" })).result, { of: "a hat", cost: 12 }); const tools = await asker.request, { tools: { name: string }[] }>("shop.tools", {}); assert.deepEqual(tools.tools.map((x) => x.name), ["price", "refund"]); } finally { stop(); await asker.close(); await shop.close(); await stream.close(); } }); test("a mesh that has issued nothing yet gets the derived shape, and says so", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); const stream = await anAssignmentsStream(); const said: string[] = []; const log = console.log; console.log = (...a: unknown[]) => said.push(a.join(" ")); let lone; try { lone = await connectNats({ url, node: "home-server", module: "lone" }); } finally { console.log = log; } const asker = await connectNats({ url, module: "console", node: "workstation" }); try { assert.equal(lone.membership(), undefined); assert.ok(said.some((s) => /no membership issued for lone on home-server/.test(s)), said.join("\n")); await lone.handle("ping", async () => ({ pong: true })); assert.deepEqual((await callTool(asker, "lone.ping", {})).result, { pong: true }); assert.deepEqual((await callTool(asker, "lone.ping@home-server", {})).result, { pong: true }); } finally { await asker.close(); await lone.close(); await stream.close(); } }); // A module that implements a seat its credential does not (yet) claim is not served for it and does // not fall over either: the runtime says so and serves the module's own tools. test("a registration under a seat the credential does not claim is said and skipped, not fatal", async (t) => { if (!url) return t.skip("MESH_TEST_NATS unset"); resetTools(); const stream = await anAssignmentsStream(); const credential = { url, node: "anchor", module: "postgres" }; const pg = await connectNats(credential); const asker = await connectNats({ url, module: "console", node: "workstation" }); const said: string[] = []; const log = console.log; console.log = (...a: unknown[]) => said.push(a.join(" ")); let stop = () => {}; try { stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-seat-unclaimed.mjs")] }); console.log = log; assert.ok(said.some((s) => /registers tools under "mesh-store".*not served/.test(s)), said.join("\n")); assert.deepEqual((await callTool(asker, "postgres.databases", {})).result, { software: "postgres" }); await assert.rejects(callTool(asker, "seat:mesh-store.databases", {}), /no responders|503/i); } finally { console.log = log; stop(); await asker.close(); await pg.close(); await stream.close(); } });