Resync hq ADR references (0044-0054 -> 0039-0049) #4

Merged
jschoubben merged 1 commits from feat/adr-ref-resync into main 2026-09-05 10:45:51 +00:00
5 changed files with 23 additions and 23 deletions
Showing only changes of commit ef540042bd - Show all commits
+1 -1
View File
@@ -13,7 +13,7 @@ runtime is what loads them and puts them on the mesh. It:
Everything hard — dispatch, collection, duplicate-name safety — is the sdk's. This is the thin Everything hard — dispatch, collection, duplicate-name safety — is the sdk's. This is the thin
wrapper that binds the broker and loads the modules. Keeping the AMQP client here, out of the sdk, wrapper that binds the broker and loads the modules. Keeping the AMQP client here, out of the sdk,
is deliberate: a broker-client change never rebuilds a module (ADR 0044). is deliberate: a broker-client change never rebuilds a module (ADR 0039).
## Running it ## Running it
+15 -15
View File
@@ -1,7 +1,7 @@
// A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker // A concrete AMQP implementation of the sdk's Broker contract, over the mesh broker
// (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime // (novox/hq ADR 0001). The sdk deliberately keeps this out — it defines the interface; the runtime
// provides the binding — so a broker-client change never rebuilds the modules. This is where the // provides the binding — so a broker-client change never rebuilds the modules. This is where the
// ADR 0047 wire shape lives: the two exchanges, persistent events, per-consumer durable queues, // ADR 0042 wire shape lives: the two exchanges, persistent events, per-consumer durable queues,
// prefetch, dead-letter — none of which a module ever sees. // prefetch, dead-letter — none of which a module ever sees.
import amqp from "amqplib"; import amqp from "amqplib";
@@ -9,14 +9,14 @@ import * as tls from "node:tls";
import { randomUUID, createHash } from "node:crypto"; import { randomUUID, createHash } from "node:crypto";
import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging"; import type { Broker, Envelope, EventHeaders } from "@novox/mesh-sdk/messaging";
// Two topic exchanges, kept apart on purpose (ADR 0047): tool invocations are request/reply and are // Two topic exchanges, kept apart on purpose (ADR 0042): tool invocations are request/reply and are
// not events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node // not events, so an audit sink subscribing to `#` on the events exchange sees module, mesh and node
// events — never the RPC traffic. // events — never the RPC traffic.
const RPC_EXCHANGE = "mesh.rpc"; const RPC_EXCHANGE = "mesh.rpc";
const EVENTS_EXCHANGE = "mesh.events"; const EVENTS_EXCHANGE = "mesh.events";
// Where an event rejected past its redelivery limit is set aside for inspection. // Where an event rejected past its redelivery limit is set aside for inspection.
const DEAD_EXCHANGE = "mesh.events.dead"; const DEAD_EXCHANGE = "mesh.events.dead";
// Bound in-flight events so one slow consumer can't pull the whole backlog into memory (ADR 0047). // Bound in-flight events so one slow consumer can't pull the whole backlog into memory (ADR 0042).
const EVENT_PREFETCH = 32; const EVENT_PREFETCH = 32;
interface Reply { interface Reply {
@@ -24,7 +24,7 @@ interface Reply {
error?: string; error?: string;
} }
/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL, the fingerprint /** A broker credential as the mesh delivers it (novox/hq ADR 0043): an amqps URL, the fingerprint
* of the certificate the broker must present, and the node and module the account is scoped to (so * of the certificate the broker must present, and the node and module the account is scoped to (so
* the runtime names its queue as the mesh did). A plain string is a bootstrap URL. */ * the runtime names its queue as the mesh did). A plain string is a bootstrap URL. */
export interface Credential { export interface Credential {
@@ -37,7 +37,7 @@ export interface Credential {
/** /**
* Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. * Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down.
* *
* A scoped module (novox/hq ADR 0048) passes `assumeExchanges: true`: its account may not declare * A scoped module (novox/hq ADR 0043) passes `assumeExchanges: true`: its account may not declare
* an exchange, and the substrate already owns them, so it declares only its own queue. A credential * an exchange, and the substrate already owns them, so it declares only its own queue. A credential
* carrying a fingerprint is dialled over amqps, pinned to exactly that certificate. * carrying a fingerprint is dialled over amqps, pinned to exactly that certificate.
*/ */
@@ -52,10 +52,10 @@ export async function connectAmqp(
// A confirm channel, so an event publish awaits the broker's ack: a publish the broker never // A confirm channel, so an event publish awaits the broker's ack: a publish the broker never
// accepted (it was mid-restart, the connection dropped) fails the emit rather than vanishing — // accepted (it was mid-restart, the connection dropped) fails the emit rather than vanishing —
// at-least-once starts at the emitter, not only the consumer (ADR 0047). // at-least-once starts at the emitter, not only the consumer (ADR 0042).
const ch = await conn.createConfirmChannel(); const ch = await conn.createConfirmChannel();
// The substrate owns the exchanges (ADR 0048). A bootstrap/admin connection declares them; a // The substrate owns the exchanges (ADR 0043). A bootstrap/admin connection declares them; a
// scoped module assumes they exist and never tries — its account could not, and the dead-letter // scoped module assumes they exist and never tries — its account could not, and the dead-letter
// queue behind the exchange is the substrate's to keep, not a module's. // queue behind the exchange is the substrate's to keep, not a module's.
if (!opts.assumeExchanges) { if (!opts.assumeExchanges) {
@@ -76,7 +76,7 @@ export async function connectAmqp(
if (replyQueue) return replyQueue; if (replyQueue) return replyQueue;
const { queue } = await ch.assertQueue("", { exclusive: true }); const { queue } = await ch.assertQueue("", { exclusive: true });
// Replies come back through the RPC exchange keyed by this queue's own name, not the default // Replies come back through the RPC exchange keyed by this queue's own name, not the default
// exchange (novox/hq ADR 0052): a serving module's scoped account may write to mesh.rpc but not // exchange (novox/hq ADR 0047): a serving module's scoped account may write to mesh.rpc but not
// the default exchange, which would let it publish into any queue on the broker. // the default exchange, which would let it publish into any queue on the broker.
await ch.bindQueue(queue, RPC_EXCHANGE, queue); await ch.bindQueue(queue, RPC_EXCHANGE, queue);
replyQueue = queue; replyQueue = queue;
@@ -95,7 +95,7 @@ export async function connectAmqp(
return queue; return queue;
} }
// One durable event queue per consumer (ADR 0047: <node>.<module>.events), with many bindings and // One durable event queue per consumer (ADR 0042: <node>.<module>.events), with many bindings and
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message // a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
// once however many bindings match, so the local match is what keeps a two-pattern module from // once however many bindings match, so the local match is what keeps a two-pattern module from
// running the wrong handler. // running the wrong handler.
@@ -123,7 +123,7 @@ export async function connectAmqp(
ch.ack(msg); ch.ack(msg);
} catch { } catch {
// First handler failure: requeue once. A second (already redelivered) dead-letters it, so a // First handler failure: requeue once. A second (already redelivered) dead-letters it, so a
// poison event is set aside rather than looping forever or vanishing (ADR 0047). Redelivery // poison event is set aside rather than looping forever or vanishing (ADR 0042). Redelivery
// re-runs every matching handler, so a consumer must be idempotent — which the ADR requires. // re-runs every matching handler, so a consumer must be idempotent — which the ADR requires.
ch.nack(msg, false, !msg.fields.redelivered); ch.nack(msg, false, !msg.fields.redelivered);
} }
@@ -151,7 +151,7 @@ export async function connectAmqp(
}, },
async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> { async handle<Req, Res>(key: string, handler: (body: Req) => Promise<Res>): Promise<() => void> {
// A durable, shared serve queue (ADR 0047): several runtimes serving one tool key compete for // A durable, shared serve queue (ADR 0042): several runtimes serving one tool key compete for
// invocations rather than each answering the same call. // invocations rather than each answering the same call.
const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true }); const { queue } = await ch.assertQueue(`serve.${key}`, { durable: true });
await ch.bindQueue(queue, RPC_EXCHANGE, key); await ch.bindQueue(queue, RPC_EXCHANGE, key);
@@ -166,7 +166,7 @@ export async function connectAmqp(
} }
if (msg.properties.replyTo) { if (msg.properties.replyTo) {
// Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped // Reply through the RPC exchange, keyed by the caller's reply-queue name, so a scoped
// account answers with write on mesh.rpc alone — never the default exchange (ADR 0052). // account answers with write on mesh.rpc alone — never the default exchange (ADR 0047).
ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), { ch.publish(RPC_EXCHANGE, msg.properties.replyTo, Buffer.from(JSON.stringify(reply)), {
correlationId: msg.properties.correlationId, correlationId: msg.properties.correlationId,
}); });
@@ -180,7 +180,7 @@ export async function connectAmqp(
async publish<T>(env: Envelope<T>): Promise<void> { async publish<T>(env: Envelope<T>): Promise<void> {
const headers = env.headers ?? ({} as EventHeaders); const headers = env.headers ?? ({} as EventHeaders);
// Events are persistent (delivery-mode 2): an audit trail that loses events on a broker // Events are persistent (delivery-mode 2): an audit trail that loses events on a broker
// restart is not one (ADR 0047). Metadata rides as headers; the body is only the payload. // restart is not one (ADR 0042). Metadata rides as headers; the body is only the payload.
// The publish is awaited to the broker's confirm — an unaccepted publish rejects here. // The publish is awaited to the broker's confirm — an unaccepted publish rejects here.
await new Promise<void>((resolve, reject) => { await new Promise<void>((resolve, reject) => {
ch.publish( ch.publish(
@@ -203,7 +203,7 @@ export async function connectAmqp(
const node = process.env.MESH_NODE; const node = process.env.MESH_NODE;
const mod = process.env.MESH_MODULE; const mod = process.env.MESH_MODULE;
// A module we can name gets its ADR 0047 durable queue; an anonymous subscriber (a test, an // A module we can name gets its ADR 0042 durable queue; an anonymous subscriber (a test, an
// ad-hoc listener) gets a transient exclusive one that dies with the connection. // ad-hoc listener) gets a transient exclusive one that dies with the connection.
if (node && mod) { if (node && mod) {
if (!eventQueue) { if (!eventQueue) {
@@ -271,7 +271,7 @@ function normalizeFingerprint(fingerprint: string): string {
/** /**
* Socket options that pin the broker to exactly the certificate whose fingerprint the mesh * Socket options that pin the broker to exactly the certificate whose fingerprint the mesh
* delivered (novox/hq ADR 0048, as the builder does). Done in two phases so a credential never * delivered (novox/hq ADR 0043, as the builder does). Done in two phases so a credential never
* reaches an impostor: first a bare TLS connection that sends nothing fetches the certificate and * reaches an impostor: first a bare TLS connection that sends nothing fetches the certificate and
* the fingerprint is checked; only then does the real connection trust *that* certificate as its * the fingerprint is checked; only then does the real connection trust *that* certificate as its
* own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key. * own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key.
+4 -4
View File
@@ -7,11 +7,11 @@
// and what an events test uses to put a message on the wire. // and what an events test uses to put a message on the wire.
// //
// The broker, in order of preference: // The broker, in order of preference:
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0048) — an // MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0043) — an
// amqps account scoped to this module. Preferred: a module holds its own. // amqps account scoped to this module. Preferred: a module holds its own.
// MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account. // MESH_BROKER_URL a plain URL, for the bootstrap/admin case before a module has an account.
// MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode) // MESH_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047) // MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0042)
import { readFileSync } from "node:fs"; import { readFileSync } from "node:fs";
import { connectAmqp } from "./broker-amqp.js"; import { connectAmqp } from "./broker-amqp.js";
@@ -25,7 +25,7 @@ import { emit } from "@novox/mesh-sdk/events";
/** /**
* Connect the way this process is meant to: with its sealed credential if the mesh gave it one, and * Connect the way this process is meant to: with its sealed credential if the mesh gave it one, and
* over the plain bootstrap URL otherwise. A scoped module assumes the substrate's exchanges exist — * over the plain bootstrap URL otherwise. A scoped module assumes the substrate's exchanges exist —
* its account may not declare them (ADR 0048). * its account may not declare them (ADR 0043).
*/ */
async function connectBroker(): Promise<Broker> { async function connectBroker(): Promise<Broker> {
const file = process.env.MESH_BROKER_FILE; const file = process.env.MESH_BROKER_FILE;
@@ -88,7 +88,7 @@ async function emitOnce(type: string, bodyJson: string): Promise<void> {
} }
const broker = await connectBroker(); const broker = await connectBroker();
useBroker(() => broker); useBroker(() => broker);
// emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close. // emit awaits the broker's publish confirm (ADR 0042), so the event is accepted before we close.
await emit(type, body); await emit(type, body);
await broker.close(); await broker.close();
} }
+1 -1
View File
@@ -28,7 +28,7 @@ test("a module's tool serves and is invoked over a real AMQP broker", { skip: !u
const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] }); const stop = await runTools({ broker: serverBroker, moduleEntrypoints: [] });
// A separate connection — a caller, like mesh-control's command API — invokes over the broker, // A separate connection — a caller, like mesh-control's command API — invokes over the broker,
// by module and tool (novox/hq ADR 0052: served on serve.demo.greet, invoked as demo.greet). // by module and tool (novox/hq ADR 0047: served on serve.demo.greet, invoked as demo.greet).
const caller = await connectAmqp(url!); const caller = await connectAmqp(url!);
const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string }; const result = (await invokeTool(caller, "demo", "greet", { who: "mesh" })) as { hello: string };
assert.equal(result.hello, "mesh"); assert.equal(result.hello, "mesh");
+2 -2
View File
@@ -6,7 +6,7 @@ import { connectAmqp } from "../dist/broker-amqp.js";
import { useBroker } from "@novox/mesh-sdk/messaging"; import { useBroker } from "@novox/mesh-sdk/messaging";
import { emit, on, type Event } from "@novox/mesh-sdk/events"; import { emit, on, type Event } from "@novox/mesh-sdk/events";
// Binding conformance: does the mesh-tools AMQP *adapter* honour the ADR 0047 wire contract — // Binding conformance: does the mesh-tools AMQP *adapter* honour the ADR 0042 wire contract —
// headers on the wire, a body that is only the payload, persistent messages, a durable per-consumer // headers on the wire, a body that is only the payload, persistent messages, a durable per-consumer
// queue, dead-letter, poison handling? This tests the adapter in isolation, against a disposable // queue, dead-letter, poison handling? This tests the adapter in isolation, against a disposable
// broker that stands in for the one the mesh hosts (ADR 0001). It is NOT an event test of the mesh: // broker that stands in for the one the mesh hosts (ADR 0001). It is NOT an event test of the mesh:
@@ -14,7 +14,7 @@ import { emit, on, type Event } from "@novox/mesh-sdk/events";
// $MESH_BROKER_URL (a throwaway broker); skipped, never failed, when it is not set. // $MESH_BROKER_URL (a throwaway broker); skipped, never failed, when it is not set.
const url = process.env.MESH_BROKER_URL; const url = process.env.MESH_BROKER_URL;
test("an event rides the wire with ADR 0047 headers and a body that is only the payload", { skip: !url }, async () => { test("an event rides the wire with ADR 0042 headers and a body that is only the payload", { skip: !url }, async () => {
const sub = await connectAmqp(url!); const sub = await connectAmqp(url!);
const pub = await connectAmqp(url!); const pub = await connectAmqp(url!);