// Round-trip the runtime's NATS client against a real server: a tool call answered, and an // event emitted and received with its envelope intact. import { connect } from "nats"; import { connectNats } from "../dist/broker-nats.js"; const URL = "nats://127.0.0.1:14222"; // The controller's job, done by hand here: the stream and the module's durable consumer. const admin = await connect({ servers: URL }); const jsm = await admin.jetstreamManager(); await jsm.streams.add({ name: "EVENTS", subjects: ["mesh.mod.*.event.>", "mesh.seat.*.event.>"] }); await jsm.consumers.add("EVENTS", { durable_name: "one_audit", ack_policy: "explicit", filter_subjects: ["mesh.mod.shop.event.order.placed"], }); const shop = await connectNats({ url: URL, node: "one", module: "shop" }); const audit = await connectNats({ url: URL, node: "one", module: "audit" }); let failures = 0; const check = (ok, what) => { console.log(` ${ok ? "ok " : "FAIL"} ${what}`); if (!ok) failures++; }; // A tool, served and called. await shop.handle("price", async (body) => ({ total: body.qty * 3 })); const answer = await audit.request("shop.price", { qty: 4 }); check(answer.total === 12, "a tool call is answered across two connections"); // A handler that throws reaches the caller as an error, not a timeout. await shop.handle("boom", async () => { throw new Error("no"); }); let threw = null; try { await audit.request("shop.boom", {}); } catch (e) { threw = e.message; } check(threw === "no", "a handler that throws answers the caller instead of timing out"); // An event, emitted and received with its envelope intact. const seen = []; await audit.subscribe("order.placed", async (env) => { seen.push(env); }); await shop.publish({ key: "order.placed", node: "one", body: { id: "a1" }, headers: { "x-event-id": "e1", "x-node": "one", "content-type": "application/json" }, }); await new Promise((r) => setTimeout(r, 800)); check(seen.length === 1, `exactly one delivery (saw ${seen.length})`); if (seen[0]) { check(seen[0].key === "order.placed", "the key survives the subject round trip"); check(seen[0].body?.id === "a1", "the body is the payload, not the whole envelope"); check(seen[0].node === "one", "the node comes back from the headers"); check(seen[0].headers?.["x-event-id"] === "e1", "the event id survives as a header"); } // A module cannot reach into another's namespace by naming its own event oddly. await shop.publish({ key: "other", node: "one", body: {}, headers: { "x-event-id": "e2" } }); const msg = await jsm.streams.getMessage("EVENTS", { last_by_subj: "mesh.mod.shop.event.other" }); check(!!msg, "an event lands under the emitting module's own namespace"); await shop.close(); await audit.close(); await admin.close(); console.log(failures ? `\n${failures} failed` : "\nall passed"); process.exit(failures ? 1 : 0);