/** * 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(); } }); 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(); } });