runtime: connect with a sealed, scoped credential over pinned amqps (ADR 0048) #2

Closed
jschoubben wants to merge 2 commits from events/module-credential into events/adr-0047-alignment
4 changed files with 218 additions and 45 deletions
+48
View File
@@ -0,0 +1,48 @@
import amqp from "amqplib";
const PORT = process.argv[2];
const MPORT = process.argv[3];
const B = `http://127.0.0.1:${MPORT}`;
const AUTH = "Basic " + Buffer.from("guest:guest").toString("base64");
async function api(method, path, body) {
const r = await fetch(B + path, {
method,
headers: { "content-type": "application/json", authorization: AUTH },
body: body ? JSON.stringify(body) : undefined,
});
if (r.status >= 300 && r.status !== 404) throw new Error(`${method} ${path} -> ${r.status}`);
}
await api("PUT", "/api/exchanges/%2f/mesh.events.dead", { type: "topic", durable: true });
await api("PUT", "/api/users/al", { password: "s", tags: "" });
const Q = "anchor.al.events";
const D = "mesh.events.dead";
const q = Q.replace(/\./g, "\\.");
const d = D.replace(/\./g, "\\.");
// configure, write, read patterns per grant on the dead exchange
const combos = {
"none": { configure: `^${q}$`, write: `^${q}$`, read: `^${q}$` },
"read-dead": { configure: `^${q}$`, write: `^${q}$`, read: `^(${q}|${d})$` },
"write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^${q}$` },
"configure-dead": { configure: `^(${q}|${d})$`, write: `^${q}$`, read: `^${q}$` },
"read+write-dead": { configure: `^${q}$`, write: `^(${q}|${d})$`, read: `^(${q}|${d})$` },
};
let i = 0;
for (const [label, perms] of Object.entries(combos)) {
await api("PUT", "/api/permissions/%2f/al", perms);
const queue = `${Q}.${i++}`; // fresh each time
try {
const c = await amqp.connect(`amqp://al:s@127.0.0.1:${PORT}/`);
const ch = await c.createChannel();
ch.on("error", () => {});
await ch.assertQueue(queue, { durable: true, deadLetterExchange: D });
console.log(`${label}: declare-with-DLX OK`);
await c.close();
} catch (e) {
console.log(`${label}: FAIL - ${String(e.message).slice(0, 70)}`);
}
}
+103 -12
View File
@@ -5,7 +5,8 @@
// prefetch, dead-letter — none of which a module ever sees.
import amqp from "amqplib";
import { randomUUID } from "node:crypto";
import * as tls from "node:tls";
import { randomUUID, createHash } from "node:crypto";
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
@@ -23,30 +24,60 @@ interface Reply {
error?: string;
}
/** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */
export async function connectAmqp(url: string): Promise<Broker> {
const conn = await amqp.connect(url);
/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL, the fingerprint
* 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. */
export interface Credential {
url: string;
fingerprint?: string;
node?: string;
module?: string;
}
/**
* 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
* 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.
*/
export async function connectAmqp(
target: string | Credential,
opts: { assumeExchanges?: boolean } = {},
): Promise<Broker> {
const cred: Credential = typeof target === "string" ? { url: target } : target;
const conn = cred.fingerprint
? await amqp.connect(cred.url, await pinnedOptions(cred.url, cred.fingerprint))
: await amqp.connect(cred.url);
// 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 —
// at-least-once starts at the emitter, not only the consumer (ADR 0047).
const ch = await conn.createConfirmChannel();
// The substrate owns the exchanges (ADR 0048). A bootstrap/admin connection declares them; a
// 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.
if (!opts.assumeExchanges) {
await ch.assertExchange(RPC_EXCHANGE, "topic", { durable: true });
await ch.assertExchange(EVENTS_EXCHANGE, "topic", { durable: true });
// The dead-letter home for poison events. A durable queue bound to `#` retains them for
// inspection — a dead-letter exchange with no queue behind it would drop them silently, which is
// exactly the loss the audit trail exists to prevent.
await ch.assertExchange(DEAD_EXCHANGE, "topic", { durable: true });
await ch.assertQueue(DEAD_EXCHANGE, { durable: true });
await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#");
}
await ch.prefetch(EVENT_PREFETCH);
// Request/reply: one exclusive reply queue, correlationId → resolver.
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
// Request/reply is set up lazily: a consumer-only module (the audit logger) never calls a tool,
// and its scoped account may not declare the exclusive reply queue this would otherwise need.
const pending = new Map<string, (r: Reply) => void>();
let replyQueue: string | undefined;
async function ensureReply(): Promise<string> {
if (replyQueue) return replyQueue;
const { queue } = await ch.assertQueue("", { exclusive: true });
replyQueue = queue;
await ch.consume(
replyQueue,
queue,
(msg) => {
if (!msg) return;
const resolve = pending.get(msg.properties.correlationId);
@@ -57,6 +88,8 @@ export async function connectAmqp(url: string): Promise<Broker> {
},
{ noAck: true },
);
return queue;
}
// One durable event queue per consumer (ADR 0047: <node>.<module>.events), with many bindings and
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
@@ -94,6 +127,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
return {
async request<Req, Res>(key: string, body: Req): Promise<Res> {
const reply = await ensureReply();
const id = randomUUID();
const answered = new Promise<Res>((resolve, reject) => {
const timer = setTimeout(() => {
@@ -107,7 +141,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
});
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
correlationId: id,
replyTo: replyQueue,
replyTo: reply,
});
return answered;
},
@@ -168,7 +202,14 @@ export async function connectAmqp(url: string): Promise<Broker> {
if (node && mod) {
if (!eventQueue) {
const name = `${node}.${mod}.events`;
if (opts.assumeExchanges) {
// The mesh pre-declared this queue with its dead-letter when it issued the account: a
// scoped account may not declare a dead-lettered queue itself (the broker refuses that
// to a non-administrator). Passively check it is there, then bind and consume.
await ch.checkQueue(name);
} else {
await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE });
}
eventQueue = name;
}
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern);
@@ -216,6 +257,56 @@ export async function connectAmqp(url: string): Promise<Broker> {
};
}
/** Normalise a certificate fingerprint to bare lower-case hex, dropping an `sha256:` prefix and
* any colon grouping, so two spellings of the same fingerprint compare equal. */
function normalizeFingerprint(fingerprint: string): string {
return fingerprint.replace(/^sha256:/i, "").replace(/:/g, "").toLowerCase();
}
/**
* 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
* 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
* own authority, so the AMQP login flows solely to the broker that proved it holds the pinned key.
* Node's `checkServerIdentity` does not run under `rejectUnauthorized: false`, so a one-phase
* "connect then check" would have already sent the password to whoever answered.
*/
async function pinnedOptions(rawUrl: string, fingerprint: string): Promise<tls.ConnectionOptions> {
const url = new URL(rawUrl);
const host = url.hostname;
const port = url.port ? Number(url.port) : 5671;
const certificate = await new Promise<tls.DetailedPeerCertificate>((resolve, reject) => {
const socket = tls.connect({ host, port, servername: host, rejectUnauthorized: false }, () => {
const peer = socket.getPeerCertificate(true);
socket.destroy();
if (!peer || !peer.raw) reject(new Error("the broker presented no certificate to pin"));
else resolve(peer);
});
socket.setTimeout(15_000, () => {
socket.destroy();
reject(new Error("timed out fetching the broker's certificate"));
});
socket.on("error", reject);
});
const seen = createHash("sha256").update(certificate.raw).digest("hex");
if (seen !== normalizeFingerprint(fingerprint)) {
throw new Error(
`the broker's certificate (sha256:${seen}) does not match the pinned ${fingerprint} — refusing`,
);
}
const pem =
"-----BEGIN CERTIFICATE-----\n" +
(certificate.raw.toString("base64").match(/.{1,64}/g) ?? []).join("\n") +
"\n-----END CERTIFICATE-----\n";
// Trust that one certificate and nothing else; the mesh's own name is not in any public store,
// so identity is the pin, not the hostname — checkServerIdentity is satisfied deliberately.
return { ca: [pem], checkServerIdentity: () => undefined };
}
/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
const raw = msg.properties.headers ?? {};
+45 -14
View File
@@ -6,23 +6,64 @@
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
// and what an events test uses to put a message on the wire.
//
// MESH_BROKER_URL amqp://… the mesh broker (both modes)
// The broker, in order of preference:
// MESH_BROKER_FILE a sealed {url, fingerprint} the mesh delivered (novox/hq ADR 0048) — an
// 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_TOOL_MODULES /path/a,/path/b,… compiled module entrypoints (serve mode)
// MESH_MODULE / MESH_NODE the identity stamped onto emitted events (ADR 0047)
import { readFileSync } from "node:fs";
import { connectAmqp } from "./broker-amqp.js";
import type { Credential } from "./broker-amqp.js";
import { runTools } from "./runtime.js";
import { useBroker } from "@novox/mesh-sdk/messaging";
import type { Broker } from "@novox/mesh-sdk/messaging";
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
* over the plain bootstrap URL otherwise. A scoped module assumes the substrate's exchanges exist —
* its account may not declare them (ADR 0048).
*/
async function connectBroker(): Promise<Broker> {
const file = process.env.MESH_BROKER_FILE;
if (file) {
let credential: Credential;
try {
credential = JSON.parse(readFileSync(file, "utf8")) as Credential;
} catch (err) {
console.error(`mesh-tools: cannot read the broker credential at ${file}: ${err}`);
process.exit(1);
}
if (!credential.url) {
console.error(`mesh-tools: ${file} carries no url — it is not a broker credential`);
process.exit(1);
}
// The mesh scoped this account to a node and module; take the runtime's identity from the
// credential so its queue and the events it emits match what the mesh authorised, no matter
// what the environment says.
if (credential.node) process.env.MESH_NODE = credential.node;
if (credential.module) process.env.MESH_MODULE = credential.module;
return connectAmqp(credential, { assumeExchanges: true });
}
const url = process.env.MESH_BROKER_URL;
if (!url) {
console.error(
"mesh-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach",
);
process.exit(1);
}
return connectAmqp(url);
}
async function serve(): Promise<void> {
const url = requireEnv("MESH_BROKER_URL");
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
.split(",")
.map((s) => s.trim())
.filter(Boolean);
const broker = await connectAmqp(url);
const broker = await connectBroker();
const stop = await runTools({ broker, moduleEntrypoints });
const shutdown = async (): Promise<void> => {
@@ -35,7 +76,6 @@ async function serve(): Promise<void> {
}
async function emitOnce(type: string, bodyJson: string): Promise<void> {
const url = requireEnv("MESH_BROKER_URL");
let body: unknown = {};
if (bodyJson) {
try {
@@ -45,7 +85,7 @@ async function emitOnce(type: string, bodyJson: string): Promise<void> {
process.exit(1);
}
}
const broker = await connectAmqp(url);
const broker = await connectBroker();
useBroker(() => broker);
// emit awaits the broker's publish confirm (ADR 0047), so the event is accepted before we close.
await emit(type, body);
@@ -66,13 +106,4 @@ async function main(): Promise<void> {
await serve();
}
function requireEnv(name: string): string {
const v = process.env[name];
if (!v) {
console.error(`mesh-tools: ${name} is not set — the runtime cannot serve without it`);
process.exit(1);
}
return v;
}
void main();
+4 -1
View File
@@ -25,8 +25,11 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
await import(pathToFileURL(resolve(entry)).href);
}
const stop = await serveTools(opts.broker);
// Serve the RPC endpoint only if a module actually registered a tool. A pure-events module (the
// audit logger) registers none, and its scoped account may not declare the serve queue — so a
// runtime that always served would fail for exactly the modules that never needed it.
const tools = listTools();
const stop = tools.length > 0 ? await serveTools(opts.broker) : () => {};
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
return stop;
}