From 97707480d6677b2e8f34cc6bbf57c386f0db2436 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 6 Sep 2026 15:50:19 +0200 Subject: [PATCH] Add lavinmq amqp provider and amqp-ping consumer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit lavinmq becomes a provider of a user-facing amqp interface: a consumer that requires a message queue is given its OWN broker — a scoped vhost and user on a lavinmq provider — not an account on the mesh's own control-plane broker (ADR 0048). Vhost-per-login is the isolation model, the exact analog of postgres's database-per-login: the provider names a vhost after the consumer's login and a user with full rights on that vhost and none elsewhere, so a login is a broker the consumer alone can reach. The provider drives lavinmq through its HTTP management API (client.ts, the module's one impure seam), with a run-once bootstrap that computes the RabbitMQ-compatible password hash lavinmq's config wants from the plain admin secret the mesh mints — the value no ${secret:...} placeholder can produce and the reason the bootstrap exists (ADR 0052). serves.amqp carries the port so consumers reference ${bound:amqp:port}. amqp-ping is a demo consumer: it contributes nothing (the vhost is the login), reads its grant from an env-file the mesh fills, and uses ${bound:amqp:as} for BOTH its username and its vhost — the db-name lesson applied to AMQP. It speaks AMQP 0-9-1 over a raw socket with no npm dependency (the way redis speaks RESP) and round-trips one message. It carries a slug so its identity fits the 20-char backend bound (ADR 0049). Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- modules/amqp-ping/client.ts | 191 +++++++++++++++++++++++++++ modules/amqp-ping/index.ts | 50 +++++++ modules/amqp-ping/module.json | 64 +++++++++ modules/amqp-ping/package.json | 14 ++ modules/amqp-ping/tsconfig.json | 12 ++ modules/lavinmq/bootstrap/index.ts | 45 +++++++ modules/lavinmq/client.ts | 150 +++++++++++++++++++++ modules/lavinmq/index.ts | 25 ++++ modules/lavinmq/module.json | 133 +++++++++++++++++++ modules/lavinmq/package.json | 14 ++ modules/lavinmq/provisioner/index.ts | 51 +++++++ modules/lavinmq/tools/index.ts | 32 +++++ modules/lavinmq/tsconfig.json | 12 ++ 13 files changed, 793 insertions(+) create mode 100644 modules/amqp-ping/client.ts create mode 100644 modules/amqp-ping/index.ts create mode 100644 modules/amqp-ping/module.json create mode 100644 modules/amqp-ping/package.json create mode 100644 modules/amqp-ping/tsconfig.json create mode 100644 modules/lavinmq/bootstrap/index.ts create mode 100644 modules/lavinmq/client.ts create mode 100644 modules/lavinmq/index.ts create mode 100644 modules/lavinmq/module.json create mode 100644 modules/lavinmq/package.json create mode 100644 modules/lavinmq/provisioner/index.ts create mode 100644 modules/lavinmq/tools/index.ts create mode 100644 modules/lavinmq/tsconfig.json diff --git a/modules/amqp-ping/client.ts b/modules/amqp-ping/client.ts new file mode 100644 index 0000000..ee5a693 --- /dev/null +++ b/modules/amqp-ping/client.ts @@ -0,0 +1,191 @@ +// amqp-ping's AMQP client — the demo consumer's own code (novox/hq ADR 0039). It speaks AMQP 0-9-1 +// directly over a raw TCP socket (node:net), the way redis's client speaks RESP: the module carries +// NO npm dependency beyond @novox/mesh-sdk — no amqplib, no CLI in the image. It does exactly one +// thing, the round-trip that proves the grant works: connect, authenticate with PLAIN to the vhost +// the mesh named, declare a queue, publish one message and get it back. +// +// This is the consumer half of the `amqp` interface. It connects as the login the mesh derived +// (`${bound:amqp:as}`) with the password the mesh minted (`${secret:amqp}`) to a vhost of that SAME +// name — the provider named the vhost after the login, so the consumer must too. Nothing here is +// hardcoded: user AND vhost are both the bound login, and a wrong vhost is refused by the broker. + +import { createConnection, type Socket } from "node:net"; +import { readFileSync } from "node:fs"; + +const FRAME_END = 0xce; +const PROTOCOL_HEADER = Buffer.from([0x41, 0x4d, 0x51, 0x50, 0x00, 0x00, 0x09, 0x01]); // "AMQP" 0-9-1 + +export interface AmqpConn { + readonly host: string; + readonly port: number; + readonly user: string; + readonly password: string; + readonly vhost: string; +} + +/** Build the connection facts from the environment the mesh's env-file set (see the module manifest). */ +export function connFromEnv(env: NodeJS.ProcessEnv = process.env): AmqpConn { + const host = env.MESH_AMQP_HOST ?? ""; + const port = Number(env.MESH_AMQP_PORT ?? "5672") || 5672; + const user = env.MESH_AMQP_USER ?? ""; + const vhost = env.MESH_AMQP_VHOST ?? user; // the provider names the vhost after the login + const password = env.MESH_AMQP_PASSWORD ?? readMaybe(env.MESH_AMQP_PASSWORD_FILE); + if (!host || !user || !password) { + throw new Error(`amqp-ping: connection is not fully set yet (host=${host} user=${user} password=${password ? "set" : "unset"})`); + } + return { host, port, user, password, vhost }; +} + +// --- wire helpers --------------------------------------------------------------------------------- + +function shortstr(s: string): Buffer { + const b = Buffer.from(s, "utf8"); + const o = Buffer.alloc(1 + b.length); + o.writeUInt8(b.length, 0); + b.copy(o, 1); + return o; +} +function longstr(s: Buffer | string): Buffer { + const b = Buffer.isBuffer(s) ? s : Buffer.from(s, "utf8"); + const o = Buffer.alloc(4 + b.length); + o.writeUInt32BE(b.length, 0); + b.copy(o, 4); + return o; +} +function u16(n: number): Buffer { + const o = Buffer.alloc(2); + o.writeUInt16BE(n, 0); + return o; +} +function u32(n: number): Buffer { + const o = Buffer.alloc(4); + o.writeUInt32BE(n, 0); + return o; +} +function frame(type: number, channel: number, payload: Buffer): Buffer { + const o = Buffer.alloc(7 + payload.length + 1); + o.writeUInt8(type, 0); + o.writeUInt16BE(channel, 1); + o.writeUInt32BE(payload.length, 3); + payload.copy(o, 7); + o.writeUInt8(FRAME_END, 7 + payload.length); + return o; +} +function method(channel: number, classId: number, methodId: number, ...parts: Buffer[]): Buffer { + return frame(1, channel, Buffer.concat([u16(classId), u16(methodId), ...parts])); +} + +interface MethodWaiter { + classId: number; + methodId: number; + resolve: (args: Buffer) => void; + reject: (e: Error) => void; +} + +/** + * Connect, authenticate to the vhost, declare a queue, publish one message and get it back. Returns + * the body that came back — the caller checks it equals what went out. Throws on any protocol error, + * including the broker's `NOT_ALLOWED` refusal of a vhost the login has no permission on (the + * isolation the provider builds, seen from the consumer's side). + */ +export function roundTrip(conn: AmqpConn, queue = "amqp-ping", payload?: string): Promise { + const body = Buffer.from(payload ?? `ping-${Date.now()}`); + return new Promise((resolve, reject) => { + const sock: Socket = createConnection({ host: conn.host, port: conn.port }); + let buf = Buffer.alloc(0); + const waiters: MethodWaiter[] = []; + let lastBody: Buffer | null = null; + let done = false; + + const fail = (e: Error): void => { + if (done) return; + done = true; + sock.destroy(); + reject(e); + }; + const expect = (classId: number, methodId: number): Promise => + new Promise((res, rej) => waiters.push({ classId, methodId, resolve: res, reject: rej })); + + sock.on("error", (e) => fail(e)); + sock.on("close", () => fail(new Error("amqp connection closed before the round-trip completed"))); + sock.on("data", (chunk: Buffer) => { + buf = Buffer.concat([buf, chunk]); + for (;;) { + if (buf.length < 7) return; + const type = buf.readUInt8(0); + const size = buf.readUInt32BE(3); + if (buf.length < 7 + size + 1) return; + const framePayload = buf.subarray(7, 7 + size); + buf = buf.subarray(7 + size + 1); + if (type === 1) { + const classId = framePayload.readUInt16BE(0); + const methodId = framePayload.readUInt16BE(2); + const args = framePayload.subarray(4); + const w = waiters.shift(); + if (!w) continue; + if (w.classId === classId && w.methodId === methodId) w.resolve(args); + else w.reject(new Error(`expected method ${w.classId}/${w.methodId}, got ${classId}/${methodId}: ${args.toString("utf8")}`)); + } else if (type === 3) { + lastBody = framePayload; // a content body frame + } + // type 2 (content header) and type 8 (heartbeat) need no handling for this round-trip. + } + }); + + sock.on("connect", () => { + void (async () => { + try { + sock.write(PROTOCOL_HEADER); + await expect(10, 10); // Connection.Start + const response = Buffer.concat([ + Buffer.from([0]), Buffer.from(conn.user, "utf8"), Buffer.from([0]), Buffer.from(conn.password, "utf8"), + ]); + // Connection.Start-Ok: empty client-properties table, PLAIN, the SASL response, locale. + sock.write(method(0, 10, 11, u32(0), shortstr("PLAIN"), longstr(response), shortstr("en_US"))); + const tune = await expect(10, 30); // Connection.Tune + const frameMax = tune.readUInt32BE(2) || 131072; + sock.write(method(0, 10, 31, u16(tune.readUInt16BE(0)), u32(frameMax), u16(0))); // Tune-Ok, no heartbeat + sock.write(method(0, 10, 40, shortstr(conn.vhost), shortstr(""), Buffer.from([0]))); // Connection.Open + await expect(10, 41); // Open-Ok — authenticated and into the vhost + + sock.write(method(1, 20, 10, shortstr(""))); // Channel.Open + await expect(20, 11); + // Queue.Declare: reserved, queue, bits(auto-delete=1), empty arguments table. + sock.write(method(1, 50, 10, u16(0), shortstr(queue), Buffer.from([0b00001000]), u32(0))); + await expect(50, 11); + + // Basic.Publish to the default exchange, routing-key = queue; then content header + body. + sock.write(method(1, 60, 40, u16(0), shortstr(""), shortstr(queue), Buffer.from([0]))); + const bodySize = Buffer.alloc(8); + bodySize.writeBigUInt64BE(BigInt(body.length), 0); + sock.write(frame(2, 1, Buffer.concat([u16(60), u16(0), bodySize, u16(0)]))); // content header, no properties + sock.write(frame(3, 1, body)); // content body + + await new Promise((r) => setTimeout(r, 200)); + lastBody = null; + sock.write(method(1, 60, 70, u16(0), shortstr(queue), Buffer.from([1]))); // Basic.Get, no-ack + await expect(60, 71); // Get-Ok (a Get-Empty would arrive as 60/72 and reject the expect) + await new Promise((r) => setTimeout(r, 200)); + const received = lastBody ? (lastBody as Buffer).toString("utf8") : ""; + + sock.write(method(0, 10, 50, u16(200), shortstr("bye"), u16(0), u16(0))); // Connection.Close + await expect(10, 51).catch(() => undefined); + done = true; + sock.end(); + resolve(received); + } catch (e) { + fail(e instanceof Error ? e : new Error(String(e))); + } + })(); + }); + }); +} + +function readMaybe(path: string | undefined): string { + if (!path) return ""; + try { + return readFileSync(path, "utf8").replace(/\n$/, ""); + } catch { + return ""; + } +} diff --git a/modules/amqp-ping/index.ts b/modules/amqp-ping/index.ts new file mode 100644 index 0000000..6f606f1 --- /dev/null +++ b/modules/amqp-ping/index.ts @@ -0,0 +1,50 @@ +// amqp-ping — a tiny demo consumer of the mesh `amqp` interface, run as a long-lived container by +// `mesh-tools run` (it never returns, so the container stays up). It exists to PROVE the grant end to +// end: the mesh gave it a scoped login and a vhost of that name on the lavinmq provider, and this +// connects with exactly those and round-trips a message. +// +// The connection facts arrive the way every consumer's do — the mesh writes them into an env-file the +// container reads (novox/hq ADR 0048): MESH_AMQP_HOST/PORT from the binding, MESH_AMQP_USER and +// MESH_AMQP_VHOST both from `${bound:amqp:as}` (the provider named the vhost after the login, so the +// consumer uses the login for both — the db-name lesson applied to AMQP), and MESH_AMQP_PASSWORD from +// `${secret:amqp}`. +// +// It retries: on first boot the provider may not have provisioned this consumer yet (the reconcile is +// asynchronous and cross-container), so a refused or unreachable connection is a "not yet", not a +// failure — it waits and tries again until the round-trip succeeds, then holds the connection idle +// and re-pings on a slow cadence so the container is a stable, running proof. + +import { connFromEnv, roundTrip } from "./client.js"; + +async function sleep(ms: number): Promise { + await new Promise((r) => setTimeout(r, ms)); +} + +async function pingOnce(): Promise { + try { + const conn = connFromEnv(); + const sent = `ping-${Date.now()}`; + const got = await roundTrip(conn, "amqp-ping", sent); + if (got === sent) { + console.log(`[amqp-ping] round-trip ok as ${conn.user} on vhost ${conn.vhost} (${conn.host}:${conn.port})`); + return true; + } + console.error(`[amqp-ping] round-trip mismatch: sent ${sent}, got ${got}`); + return false; + } catch (err) { + console.error(`[amqp-ping] not ready yet: ${err instanceof Error ? err.message : err}`); + return false; + } +} + +// Wait for the first successful round-trip — the proof this consumer's grant works — then stay up. +let first = false; +for (let i = 0; !first; i++) { + first = await pingOnce(); + if (!first) await sleep(3000); +} +console.log("[amqp-ping] connected and round-tripped; holding steady"); +for (;;) { + await sleep(30000); + await pingOnce(); +} diff --git a/modules/amqp-ping/module.json b/modules/amqp-ping/module.json new file mode 100644 index 0000000..ab4a0cb --- /dev/null +++ b/modules/amqp-ping/module.json @@ -0,0 +1,64 @@ +{ + "module": "amqp-ping", + "slug": "ping", + "version": "1", + "capabilities": [ + "container-runtime" + ], + "requires": [ + "amqp" + ], + "contributes": {}, + "binds": { + "amqp": "/var/lib/amqp-ping/amqp.json" + }, + "secrets": { + "amqp": "/var/lib/amqp-ping/amqp.secret" + }, + "own-secrets": { + "broker": "/var/lib/mesh/amqp-ping/broker" + }, + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/amqp-ping", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/amqp-ping", + "mode": "0700" + }, + { + "id": "amqp-env", + "type": "file", + "path": "/var/lib/amqp-ping/amqp.env", + "mode": "0600", + "content": "MESH_AMQP_HOST=${bound:amqp:at}\nMESH_AMQP_PORT=${bound:amqp:port}\nMESH_AMQP_USER=${bound:amqp:as}\nMESH_AMQP_VHOST=${bound:amqp:as}\nMESH_AMQP_PASSWORD=${secret:amqp}\n" + }, + { + "id": "net", + "type": "network", + "name": "amqp-ping" + }, + { + "id": "runtime", + "type": "container", + "name": "amqp-ping", + "image": "mesh-runtime-amqp-ping@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "amqp-ping", + "env-file": [ + "/var/lib/amqp-ping/amqp.env" + ], + "args": [ + "run", + "/app/modules/amqp-ping/dist/index.js" + ], + "restart-on": [ + "amqp-env" + ] + } + ] +} diff --git a/modules/amqp-ping/package.json b/modules/amqp-ping/package.json new file mode 100644 index 0000000..d8883aa --- /dev/null +++ b/modules/amqp-ping/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-amqp-ping", + "version": "0.1.0", + "description": "amqp-ping — a demo consumer of the mesh amqp interface. Connects with its scoped grant and round-trips one message to prove the broker the mesh gave it (novox/hq ADR 0039).", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/amqp-ping/tsconfig.json b/modules/amqp-ping/tsconfig.json new file mode 100644 index 0000000..63b6d63 --- /dev/null +++ b/modules/amqp-ping/tsconfig.json @@ -0,0 +1,12 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["client.ts", "index.ts"] +} diff --git a/modules/lavinmq/bootstrap/index.ts b/modules/lavinmq/bootstrap/index.ts new file mode 100644 index 0000000..3ba7369 --- /dev/null +++ b/modules/lavinmq/bootstrap/index.ts @@ -0,0 +1,45 @@ +// lavinmq's run-once bootstrap — lavinmq's own code (novox/hq ADR 0039), run once before the broker +// first starts (ADR 0052). lavinmq's default admin is set at first boot from a config file's +// `default_password_hash`, and that value is a HASH of the mesh-minted admin password, not the +// password itself — a form the mesh's plain-secret delivery cannot produce and no `${secret:...}` +// placeholder can compute. So this step computes it: it reads the admin password the mesh minted and +// the host unsealed, hashes it the way lavinmq expects (see client.rabbitHash), and writes the config +// file the broker container reads with `--config`. The host runs it to completion and only then +// starts the broker the manifest places after it — so the broker's first boot finds a config with an +// admin it can authenticate, and the provisioner (which reaches the management API as that admin) +// can do its work. +// +// It runs in the module's own runtime image, under the module's own account, as `mesh-tools run` +// imports it — no broker connection, because writing a config file is an offline operation and there +// is no broker to reach yet. +// +// lavinmq consults `default_user`/`default_password_hash` only on a first boot with an empty data +// dir; a later boot uses the persisted user database and ignores them. So this seeds the admin once, +// and a rotation of the admin secret does not re-key an already-initialised broker — the same +// first-boot-only shape the RabbitMQ-compatible default user has always had. + +import { writeFileSync } from "node:fs"; +import { readFileSync } from "node:fs"; +import { rabbitHash } from "../client.js"; + +const adminUser = process.env.MESH_PROVISION_ADMIN_USER ?? process.env.MESH_LAVINMQ_ADMIN_USER ?? "mesh-admin"; +const passwordFile = process.env.MESH_PROVISION_PASSWORD_FILE ?? process.env.MESH_LAVINMQ_ADMIN_PASSWORD_FILE ?? "/run/secrets/default"; +const configOut = process.env.MESH_LAVINMQ_CONFIG_OUT ?? "/var/lib/lavinmq-module/lavinmq.ini"; +const dataDir = process.env.MESH_LAVINMQ_DATA_DIR ?? "/var/lib/lavinmq"; + +const password = readFileSync(passwordFile, "utf8").replace(/\n$/, ""); +if (!password) { + throw new Error(`[lavinmq:bootstrap] admin password file ${passwordFile} is empty — cannot seed the admin`); +} + +// The broker reads only what it needs to authenticate its admin on first boot; bind/ports come from +// the container's entrypoint (`-b 0.0.0.0`), so this file names the admin and nothing else about the +// network. +const ini = + "[main]\n" + + `data_dir = ${dataDir}\n` + + `default_user = ${adminUser}\n` + + `default_password_hash = ${rabbitHash(password)}\n`; + +writeFileSync(configOut, ini, { mode: 0o600 }); +console.log(`[lavinmq:bootstrap] wrote ${configOut} with admin '${adminUser}' (password hashed for lavinmq)`); diff --git a/modules/lavinmq/client.ts b/modules/lavinmq/client.ts new file mode 100644 index 0000000..04716b2 --- /dev/null +++ b/modules/lavinmq/client.ts @@ -0,0 +1,150 @@ +// lavinmq's admin client — lavinmq's own code, living in the module (novox/hq ADR 0039). Both this +// module's tools and its provisioner import it, and nothing outside lavinmq does. +// +// It drives lavinmq through its HTTP management API (the RabbitMQ-compatible surface lavinmq serves +// on 15672), not a hand-rolled AMQP admin stack: the module may take NO npm dependency beyond +// @novox/mesh-sdk, and the management API is exactly the admin surface — create/remove a vhost, a +// user, and its permissions — reached with `fetch` (global on node 22) and HTTP Basic auth. One +// boundary, `api()`, and every method is built on it. This is the module's one impure seam, the way +// postgres's is `psql` and redis's is a RESP socket. +// +// **The login and password are the mesh's, not the provisioner's (novox/hq ADR 0048).** The mesh +// derives the login and hands it to both ends so they agree, and mints the password and delivers a +// copy to each. lavinmq creates exactly that user with exactly that password on a vhost of the same +// name — a name or password the provisioner invented is one the consumer could never present. + +import { createHash, randomBytes } from "node:crypto"; +import { readFileSync } from "node:fs"; + +export interface LavinmqConn { + /** Base URL of the management API, e.g. http://lavinmq:15672 (no trailing /api). */ + readonly base: string; + readonly adminUser: string; + readonly adminPassword: string; +} + +export class LavinmqClient { + constructor(private readonly conn: LavinmqConn) {} + + /** + * Build from the module's resolved environment. Reads MESH_LAVINMQ_* first (the documented + * names), falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner + * container so the module runs unchanged there. Throws if it cannot find a management endpoint and + * an admin password — the right failure, because without them nothing it does can work. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): LavinmqClient { + const base = (env.MESH_LAVINMQ_MANAGEMENT ?? env.MESH_PROVISION_LAVINMQ ?? "").replace(/\/+$/, ""); + const adminUser = env.MESH_LAVINMQ_ADMIN_USER ?? env.MESH_PROVISION_ADMIN_USER ?? "mesh-admin"; + const adminPassword = env.MESH_LAVINMQ_ADMIN_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); + if (!base || !adminPassword) { + throw new Error("lavinmq management endpoint or admin password is not set — lavinmq's own code cannot reach the server"); + } + return new LavinmqClient({ base, adminUser, adminPassword }); + } + + /** One request against the management API. A non-2xx reply rejects, carrying the body for the log. */ + async api(method: string, path: string, body?: unknown): Promise { + const headers: Record = { + Authorization: "Basic " + Buffer.from(`${this.conn.adminUser}:${this.conn.adminPassword}`).toString("base64"), + }; + if (body !== undefined) headers["Content-Type"] = "application/json"; + const resp = await fetch(`${this.conn.base}/api${path}`, { + method, + headers, + body: body !== undefined ? JSON.stringify(body) : undefined, + }); + if (!resp.ok) { + throw new Error(`lavinmq management API ${method} ${path} -> ${resp.status}: ${await resp.text()}`); + } + const text = await resp.text(); + return text ? JSON.parse(text) : null; + } + + /** True once the management API answers — the server has finished starting. */ + async ready(): Promise { + try { + await this.api("GET", "/overview"); + return true; + } catch { + return false; + } + } + + /** Block until the management API answers, or throw once the budget is spent. */ + async waitReady(retries = 30, delayMs = 1000): Promise { + for (let i = 0; i < retries; i++) { + if (await this.ready()) return; + await new Promise((r) => setTimeout(r, delayMs)); + } + throw new Error("lavinmq management API did not become ready"); + } + + /** + * Create (or reset to a known state) one consumer's broker: a vhost and a user both named for the + * consumer's login, with the login owning full permissions on exactly that vhost. Idempotent — a + * PUT of a vhost or user that exists is a no-op or a password reset, so a reconcile can call it + * again without harm. The consumer connects as `` to vhost `` and can reach nothing + * else (novox/hq ADR 0048). + */ + async createConsumer(login: string, password: string): Promise { + const v = encodeURIComponent(login); + const u = encodeURIComponent(login); + await this.api("PUT", `/vhosts/${v}`); + await this.api("PUT", `/users/${u}`, { password, tags: "" }); + await this.api("PUT", `/permissions/${v}/${u}`, { configure: ".*", write: ".*", read: ".*" }); + } + + /** Remove a consumer's vhost and user, idempotently. A DELETE of what is already gone is tolerated. */ + async removeConsumer(login: string): Promise { + const v = encodeURIComponent(login); + const u = encodeURIComponent(login); + try { + await this.api("DELETE", `/vhosts/${v}`); + } catch (err) { + console.error(`[lavinmq] delete vhost ${login} failed (continuing): ${err}`); + } + try { + await this.api("DELETE", `/users/${u}`); + } catch (err) { + console.error(`[lavinmq] delete user ${login} failed (continuing): ${err}`); + } + } + + /** The vhosts, for the amqp_list_vhosts tool. */ + async listVhosts(): Promise<{ name: string; messages: number }[]> { + const vhosts = (await this.api("GET", "/vhosts")) as { name: string; messages?: number }[]; + return vhosts.map((v) => ({ name: v.name, messages: v.messages ?? 0 })); + } + + /** The queues on one vhost (default the root vhost), for the amqp_list_queues tool. */ + async listQueues(vhost = "/"): Promise<{ name: string; messages: number; consumers: number }[]> { + const queues = (await this.api("GET", `/queues/${encodeURIComponent(vhost)}`)) as + { name: string; messages?: number; consumers?: number }[]; + return queues.map((q) => ({ name: q.name, messages: q.messages ?? 0, consumers: q.consumers ?? 0 })); + } +} + +/** + * The RabbitMQ-compatible SHA-256 password hash lavinmq's `default_password_hash` expects: + * base64( salt[4] || sha256( salt || utf8(password) ) ). The salt is any four bytes — random here, + * because a fixed salt buys nothing and a fresh one is free. Verified against `lavinmqctl + * hash_password`: a hash produced here is accepted by the server unchanged. + */ +export function rabbitHash(password: string, salt: Buffer = randomBytes(4)): string { + const digest = createHash("sha256").update(Buffer.concat([salt, Buffer.from(password, "utf8")])).digest(); + return Buffer.concat([salt, digest]).toString("base64"); +} + +/** Generate a URL-safe password. */ +export function generatePassword(): string { + return randomBytes(24).toString("base64url"); +} + +function readSecretFile(path: string | undefined): string | undefined { + if (!path) return undefined; + try { + return readFileSync(path, "utf8").trim(); + } catch { + return undefined; + } +} diff --git a/modules/lavinmq/index.ts b/modules/lavinmq/index.ts new file mode 100644 index 0000000..d440014 --- /dev/null +++ b/modules/lavinmq/index.ts @@ -0,0 +1,25 @@ +// lavinmq's events entrypoint, loaded by the per-node tool host (the provisioner container runs +// ./provisioner separately). The broker lifecycle events are EMITTED from the provisioner, where the +// lifecycle actually happens (novox/hq ADR 0041/0042): +// module.lavinmq.amqp.provisioned — a consumer's vhost + user was created +// module.lavinmq.amqp.deprovisioned — that vhost + user was removed +// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a +// broker and who lost one — observability the provider itself is best placed to log. + +import { on } from "@novox/mesh-sdk/events"; + +interface AmqpEvent { + consumer?: string; + user: string; + vhost?: string; +} + +await on("module.lavinmq.amqp.provisioned", async (e) => { + console.log(`[lavinmq] broker provisioned for ${e.body.consumer ?? "?"} (user ${e.body.user}, vhost ${e.body.vhost})`); +}); + +await on("module.lavinmq.amqp.deprovisioned", async (e) => { + console.log(`[lavinmq] broker deprovisioned (user ${e.body.user})`); +}); + +console.log("[lavinmq] auditing broker lifecycle events"); diff --git a/modules/lavinmq/module.json b/modules/lavinmq/module.json new file mode 100644 index 0000000..ad94b4d --- /dev/null +++ b/modules/lavinmq/module.json @@ -0,0 +1,133 @@ +{ + "module": "lavinmq", + "version": "1", + "provides": [ + { + "name": "amqp", + "scope": "mesh" + } + ], + "capabilities": [ + "container-runtime" + ], + "emits": [ + "module.lavinmq.amqp.provisioned", + "module.lavinmq.amqp.deprovisioned" + ], + "consumes": [ + "module.lavinmq.amqp.provisioned", + "module.lavinmq.amqp.deprovisioned" + ], + "serves": { + "amqp": { + "port": 5672 + } + }, + "receives": { + "amqp": "/var/lib/lavinmq-module/grants/mesh.json" + }, + "grants": { + "amqp": "/var/lib/lavinmq-module/grants" + }, + "own-secrets": { + "default": "/var/lib/lavinmq-module/default.secret", + "broker": "/var/lib/mesh/lavinmq/broker" + }, + "listens": [ + { + "port": 5672, + "protocol": "tcp", + "from": "mesh", + "why": "modules on any machine that were granted a queue" + } + ], + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/lavinmq", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/lavinmq-module", + "mode": "0700" + }, + { + "id": "grants-dir", + "type": "directory", + "path": "/var/lib/lavinmq-module/grants", + "mode": "0700" + }, + { + "id": "data", + "type": "directory", + "path": "/services/lavinmq/data", + "mode": "0700" + }, + { + "id": "net", + "type": "network", + "name": "lavinmq" + }, + { + "id": "bootstrap", + "type": "container", + "name": "lavinmq-bootstrap", + "image": "mesh-runtime-lavinmq@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "run-once": true, + "volumes": [ + "/var/lib/lavinmq-module:/var/lib/lavinmq-module", + "/var/lib/lavinmq-module/default.secret:/run/secrets/default:ro" + ], + "env": { + "MESH_PROVISION_ADMIN_USER": "mesh-admin", + "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/default", + "MESH_LAVINMQ_CONFIG_OUT": "/var/lib/lavinmq-module/lavinmq.ini", + "MESH_LAVINMQ_DATA_DIR": "/var/lib/lavinmq" + }, + "args": [ + "run", + "/app/modules/lavinmq/dist/bootstrap/index.js" + ] + }, + { + "id": "server", + "type": "container", + "name": "lavinmq", + "image": "cloudamqp/lavinmq@sha256:3eb54c12916d700a978c2ea86e6362cd4974b0e3189508718006d4e6d341246b", + "network": "lavinmq", + "ports": [ + "5672" + ], + "volumes": [ + "/services/lavinmq/data:/var/lib/lavinmq", + "/var/lib/lavinmq-module/lavinmq.ini:/etc/lavinmq/lavinmq.ini:ro" + ], + "args": [ + "--config", + "/etc/lavinmq/lavinmq.ini" + ] + }, + { + "id": "runtime", + "type": "container", + "name": "mesh-lavinmq", + "image": "mesh-runtime-lavinmq@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "lavinmq", + "volumes": [ + "/var/lib/mesh/lavinmq/broker:/run/secrets/broker:ro", + "/var/lib/lavinmq-module/grants:/var/lib/lavinmq-module/grants:ro", + "/var/lib/lavinmq-module/default.secret:/run/secrets/default:ro" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_RECEIVES": "/var/lib/lavinmq-module/grants/mesh.json", + "MESH_PROVISION_LAVINMQ": "http://lavinmq:15672", + "MESH_PROVISION_ADMIN_USER": "mesh-admin", + "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/default" + } + } + ] +} diff --git a/modules/lavinmq/package.json b/modules/lavinmq/package.json new file mode 100644 index 0000000..1a4b297 --- /dev/null +++ b/modules/lavinmq/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-lavinmq", + "version": "0.1.0", + "description": "lavinmq — provides the mesh amqp interface (a per-consumer AMQP message broker). Its management client, provisioner, run-once bootstrap, tools and events live here (novox/hq ADR 0039).", + "type": "module", + "private": true, + "dependencies": { + "@novox/mesh-sdk": "^0.1.0" + }, + "devDependencies": { + "@types/node": "^22.0.0", + "typescript": "^5.6.0" + } +} diff --git a/modules/lavinmq/provisioner/index.ts b/modules/lavinmq/provisioner/index.ts new file mode 100644 index 0000000..f5514bf --- /dev/null +++ b/modules/lavinmq/provisioner/index.ts @@ -0,0 +1,51 @@ +// lavinmq's provisioner — the adapter that makes lavinmq a provider of the mesh `amqp` interface. +// The reconcile loop, the contributions file, and reading the mesh's minted password are the sdk +// harness's; this writes only the per-service half: how lavinmq creates and removes a consumer's own +// broker (novox/hq ADR 0039/0040/0048). +// +// The `amqp` interface: a consumer connects as `as` with the password the mesh minted, to a vhost +// named for that same login — its own message broker, isolated from every other consumer's by the +// vhost boundary. It is a broker of its own, not a shared account on the mesh's control-plane broker. +// +// **The login and password are the mesh's, not the provisioner's (novox/hq ADR 0048).** The mesh +// derives the login and hands it to both ends so they agree, and mints the password and delivers a +// copy to each. lavinmq creates exactly that user with exactly that password — a name or password the +// provisioner invented is one the consumer could never present. +// +// Vhost-per-login is the isolation model, the exact analog of postgres's database-per-login: the +// consumer owns one vhost, named for its login, and a user with full rights on that vhost and no +// rights anywhere else. lavinmq enforces it — a user with no permission on `/` is refused the moment +// it opens that vhost (`NOT_ALLOWED`), so a login is a broker the consumer alone can reach. + +import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { LavinmqClient } from "../client.js"; + +const lavinmq = LavinmqClient.fromEnv(); + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:amqp] emit ${type} failed: ${err}`); + } +} + +runProvisioner("amqp", { + async create(p: Provision): Promise { + // The vhost and the user share the consumer's login, so one cannot reach another's broker. + await lavinmq.waitReady(); + await lavinmq.createConsumer(p.as, p.password); + await announce("module.lavinmq.amqp.provisioned", { + consumer: p.consumer ?? "", + user: p.as, + vhost: p.as, + }); + }, + + async remove(p: { as: string }): Promise { + await lavinmq.removeConsumer(p.as); + await announce("module.lavinmq.amqp.deprovisioned", { user: p.as, vhost: p.as }); + }, +}); diff --git a/modules/lavinmq/tools/index.ts b/modules/lavinmq/tools/index.ts new file mode 100644 index 0000000..55489c6 --- /dev/null +++ b/modules/lavinmq/tools/index.ts @@ -0,0 +1,32 @@ +// lavinmq's tools — lavinmq's own code (novox/hq ADR 0039), importing lavinmq's own management +// client. They return structured data; the mesh serves them through the sdk's tool harness. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { LavinmqClient } from "../client.js"; + +export function getLavinmqTools(lavinmq: LavinmqClient): ToolDefinition[] { + return [ + { + name: "amqp_list_vhosts", + description: "List the lavinmq virtual hosts — one per consumer that was granted a broker.", + input: {}, + run: async () => ({ vhosts: await lavinmq.listVhosts() }), + }, + { + name: "amqp_list_queues", + description: "List the queues on a vhost with message and consumer counts. Omit vhost for the root '/'.", + input: { vhost: { type: "string", description: "the vhost to list, e.g. a consumer's login; defaults to '/'" } }, + run: async (args) => ({ queues: await lavinmq.listQueues(args.vhost ? String(args.vhost) : undefined) }), + }, + ]; +} + +// The tools exist only when the server can be reached from the environment; without it, lavinmq +// contributes none rather than failing the whole tool runtime. +registerModuleTools("lavinmq", (env) => { + try { + return getLavinmqTools(LavinmqClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/lavinmq/tsconfig.json b/modules/lavinmq/tsconfig.json new file mode 100644 index 0000000..95f347f --- /dev/null +++ b/modules/lavinmq/tsconfig.json @@ -0,0 +1,12 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts", "bootstrap/index.ts"] +}