runtime: connect with a sealed credential, scoped, over pinned amqps (ADR 0048)
A module reads its broker credential from MESH_BROKER_FILE — the sealed
{url,fingerprint} the mesh delivered — and connects over amqps pinned to
exactly that certificate. The pin is two-phase (fetch cert, verify, then
trust only it), because Node's checkServerIdentity does not run under
rejectUnauthorized:false, so a naive connect-then-check would already have
sent the password to whoever answered.
A scoped module (assumeExchanges) never declares the exchanges (its account
may not) nor its own queue with a dead-letter (the broker refuses that to a
non-administrator) — the mesh pre-declared the queue, so it passively checks
it, binds and consumes. The RPC reply queue is lazy, and a module that
registered no tools serves none: a pure-events consumer touches only what its
account allows.
Verified end-to-end against a real broker as the scoped account: the audit
logger consumes # and records events, over an account that is not the
broker's own.
Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
This commit is contained in:
@@ -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)}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
+100
-12
@@ -5,7 +5,8 @@
|
|||||||
// 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";
|
||||||
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";
|
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 0047): tool invocations are request/reply and are
|
||||||
@@ -23,30 +24,57 @@ interface Reply {
|
|||||||
error?: string;
|
error?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Connect to the mesh broker and return a Broker. `close()` tears both channel and connection down. */
|
/** A broker credential as the mesh delivers it (novox/hq ADR 0048): an amqps URL and the
|
||||||
export async function connectAmqp(url: string): Promise<Broker> {
|
* fingerprint of the certificate the broker must present. A plain string is a bootstrap URL. */
|
||||||
const conn = await amqp.connect(url);
|
export interface Credential {
|
||||||
|
url: string;
|
||||||
|
fingerprint?: 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
|
// 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 0047).
|
||||||
const ch = await conn.createConfirmChannel();
|
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(RPC_EXCHANGE, "topic", { durable: true });
|
||||||
await ch.assertExchange(EVENTS_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.assertExchange(DEAD_EXCHANGE, "topic", { durable: true });
|
||||||
await ch.assertQueue(DEAD_EXCHANGE, { durable: true });
|
await ch.assertQueue(DEAD_EXCHANGE, { durable: true });
|
||||||
await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#");
|
await ch.bindQueue(DEAD_EXCHANGE, DEAD_EXCHANGE, "#");
|
||||||
|
}
|
||||||
|
|
||||||
await ch.prefetch(EVENT_PREFETCH);
|
await ch.prefetch(EVENT_PREFETCH);
|
||||||
|
|
||||||
// Request/reply: one exclusive reply queue, correlationId → resolver.
|
// Request/reply is set up lazily: a consumer-only module (the audit logger) never calls a tool,
|
||||||
const { queue: replyQueue } = await ch.assertQueue("", { exclusive: true });
|
// and its scoped account may not declare the exclusive reply queue this would otherwise need.
|
||||||
const pending = new Map<string, (r: Reply) => void>();
|
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(
|
await ch.consume(
|
||||||
replyQueue,
|
queue,
|
||||||
(msg) => {
|
(msg) => {
|
||||||
if (!msg) return;
|
if (!msg) return;
|
||||||
const resolve = pending.get(msg.properties.correlationId);
|
const resolve = pending.get(msg.properties.correlationId);
|
||||||
@@ -57,6 +85,8 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
},
|
},
|
||||||
{ noAck: true },
|
{ noAck: true },
|
||||||
);
|
);
|
||||||
|
return queue;
|
||||||
|
}
|
||||||
|
|
||||||
// One durable event queue per consumer (ADR 0047: <node>.<module>.events), with many bindings and
|
// 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
|
// a single consumer that fans out to the handlers whose pattern matches. AMQP delivers a message
|
||||||
@@ -94,6 +124,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
|
|
||||||
return {
|
return {
|
||||||
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
async request<Req, Res>(key: string, body: Req): Promise<Res> {
|
||||||
|
const reply = await ensureReply();
|
||||||
const id = randomUUID();
|
const id = randomUUID();
|
||||||
const answered = new Promise<Res>((resolve, reject) => {
|
const answered = new Promise<Res>((resolve, reject) => {
|
||||||
const timer = setTimeout(() => {
|
const timer = setTimeout(() => {
|
||||||
@@ -107,7 +138,7 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
});
|
});
|
||||||
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
ch.publish(RPC_EXCHANGE, key, Buffer.from(JSON.stringify(body)), {
|
||||||
correlationId: id,
|
correlationId: id,
|
||||||
replyTo: replyQueue,
|
replyTo: reply,
|
||||||
});
|
});
|
||||||
return answered;
|
return answered;
|
||||||
},
|
},
|
||||||
@@ -168,7 +199,14 @@ export async function connectAmqp(url: string): Promise<Broker> {
|
|||||||
if (node && mod) {
|
if (node && mod) {
|
||||||
if (!eventQueue) {
|
if (!eventQueue) {
|
||||||
const name = `${node}.${mod}.events`;
|
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 });
|
await ch.assertQueue(name, { durable: true, deadLetterExchange: DEAD_EXCHANGE });
|
||||||
|
}
|
||||||
eventQueue = name;
|
eventQueue = name;
|
||||||
}
|
}
|
||||||
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern);
|
await ch.bindQueue(eventQueue, EVENTS_EXCHANGE, pattern);
|
||||||
@@ -216,6 +254,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. */
|
/** Read a broker message back into an Envelope: string headers, contentType folded in, body parsed. */
|
||||||
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
function toEnvelope<T>(msg: amqp.ConsumeMessage): Envelope<T> {
|
||||||
const raw = msg.properties.headers ?? {};
|
const raw = msg.properties.headers ?? {};
|
||||||
|
|||||||
+40
-14
@@ -6,23 +6,59 @@
|
|||||||
// mesh-tools emit TYPE [JSON] emit one event onto the mesh and exit — an operable primitive,
|
// 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.
|
// 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_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 0047)
|
||||||
|
|
||||||
|
import { readFileSync } from "node:fs";
|
||||||
import { connectAmqp } from "./broker-amqp.js";
|
import { connectAmqp } from "./broker-amqp.js";
|
||||||
|
import type { Credential } from "./broker-amqp.js";
|
||||||
import { runTools } from "./runtime.js";
|
import { runTools } from "./runtime.js";
|
||||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||||
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
import { emit } from "@novox/mesh-sdk/events";
|
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);
|
||||||
|
}
|
||||||
|
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> {
|
async function serve(): Promise<void> {
|
||||||
const url = requireEnv("MESH_BROKER_URL");
|
|
||||||
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
const moduleEntrypoints = (process.env.MESH_TOOL_MODULES ?? "")
|
||||||
.split(",")
|
.split(",")
|
||||||
.map((s) => s.trim())
|
.map((s) => s.trim())
|
||||||
.filter(Boolean);
|
.filter(Boolean);
|
||||||
|
|
||||||
const broker = await connectAmqp(url);
|
const broker = await connectBroker();
|
||||||
const stop = await runTools({ broker, moduleEntrypoints });
|
const stop = await runTools({ broker, moduleEntrypoints });
|
||||||
|
|
||||||
const shutdown = async (): Promise<void> => {
|
const shutdown = async (): Promise<void> => {
|
||||||
@@ -35,7 +71,6 @@ async function serve(): Promise<void> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
||||||
const url = requireEnv("MESH_BROKER_URL");
|
|
||||||
let body: unknown = {};
|
let body: unknown = {};
|
||||||
if (bodyJson) {
|
if (bodyJson) {
|
||||||
try {
|
try {
|
||||||
@@ -45,7 +80,7 @@ async function emitOnce(type: string, bodyJson: string): Promise<void> {
|
|||||||
process.exit(1);
|
process.exit(1);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
const broker = await connectAmqp(url);
|
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 0047), so the event is accepted before we close.
|
||||||
await emit(type, body);
|
await emit(type, body);
|
||||||
@@ -66,13 +101,4 @@ async function main(): Promise<void> {
|
|||||||
await serve();
|
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();
|
void main();
|
||||||
|
|||||||
+4
-1
@@ -25,8 +25,11 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
|||||||
await import(pathToFileURL(resolve(entry)).href);
|
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 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)"}`);
|
console.log(`[mesh-tools] serving ${tools.length} tool(s): ${tools.map((t) => t.name).join(", ") || "(none)"}`);
|
||||||
return stop;
|
return stop;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user