runtime: connect with a sealed, scoped credential over pinned amqps (ADR 0048) #2
@@ -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
@@ -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
@@ -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
@@ -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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user