Task 3.6 of novox/hq ADR 0116. A module is still written against request, handle, publish, subscribe, close; only what is underneath changes. main.ts still selects the AMQP client — steps 1 to 4 leave every node on AMQP, so this ships beside it and is selected at the rollout. Round-tripped against a real server (test/roundtrip.mjs): a tool answered across two connections, a throwing handler reaching the caller as an error rather than a timeout, an event delivered once with its key, body, node and event id intact, and an event landing under its emitter's own namespace. Three things the compiler and the server corrected: - the envelope's field is `key`, not `type`, and the payload is `env.body` with metadata in headers — not the whole envelope re-encoded. An implementation that nested the envelope would pass all its own tests and agree with nobody, which is what the conformance suite exists to stop. - the NATS client's TLS options are PEM strings with no verify hook, so the AMQP client's `checkServerIdentity: () => undefined` has no equivalent. The fingerprint check still happens and is still the guarantee, but the bus's certificate must now carry a SAN matching the address nodes dial. That is a constraint on the mesh's certificates, recorded where it bites. - a durable consumer is bound, never created: a module's account cannot reach the JetStream API, and a runtime creating its own would be a module choosing its own delivery semantics.
58 lines
2.8 KiB
JavaScript
58 lines
2.8 KiB
JavaScript
// 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);
|