diff --git a/modules/amqp-email-forwarder/module.json b/modules/amqp-email-forwarder/module.json deleted file mode 100644 index 94f84e6..0000000 --- a/modules/amqp-email-forwarder/module.json +++ /dev/null @@ -1,56 +0,0 @@ -{ - "module": "amqp-email-forwarder", - "version": "1", - "slug": "emailfwd", - "capabilities": [ - "container-runtime" - ], - "requires": [ - "amqp" - ], - "contributes": {}, - "binds": { - "amqp": "/var/lib/amqp-email-forwarder/amqp.json" - }, - "secrets": { - "amqp": "/var/lib/amqp-email-forwarder/amqp.secret" - }, - "own-secrets": { - "smtp-user": "/var/lib/amqp-email-forwarder/smtp-user.secret", - "smtp-password": "/var/lib/amqp-email-forwarder/smtp-password.secret" - }, - "resources": [ - { - "id": "state", - "type": "directory", - "path": "/var/lib/amqp-email-forwarder", - "mode": "0700" - }, - { - "id": "app-env", - "type": "file", - "path": "/var/lib/amqp-email-forwarder/app.env", - "mode": "0600", - "content": "AMQP_HOST=${bound:amqp:at}\nAMQP_PORT=${bound:amqp:port}\nAMQP_USER=${bound:amqp:as}\nAMQP_VHOST=EMAILDELIVERY_T\nAMQP_EXCHANGE=News.TransactionalEmailing.Command\nAMQP_QUEUE=email-forwarder\nAMQP_URL=amqp://${bound:amqp:as}:${secret:amqp}@${bound:amqp:at}:${bound:amqp:port}/EMAILDELIVERY_T\nSMTP_HOST=mail.novox.be\nSMTP_PORT=587\nSMTP_USER=${secret:smtp-user}\nSMTP_PASSWORD=${secret:smtp-password}\n" - }, - { - "id": "net", - "type": "network", - "name": "amqp-email-forwarder" - }, - { - "id": "app", - "type": "container", - "name": "amqp-email-forwarder", - "image": "registry-api.novox.be/novox/amqp-email-forwarder@sha256:f76d34646d9d3b2098c72688a63f6ae656f1888ffcb2f90a9c8dd2a44ad7f8af", - "network": "amqp-email-forwarder", - "env-file": [ - "/var/lib/amqp-email-forwarder/app.env" - ], - "restart-on": [ - "app-env" - ], - "secrets-in-environment": "the application's own code reads AMQP_URL, SMTP_USER and SMTP_PASSWORD from the environment (amqp-email-forwarder app.js); converting is that repository's change" - } - ] -} diff --git a/modules/amqp-ping/Dockerfile b/modules/amqp-ping/Dockerfile deleted file mode 100644 index 7814b0e..0000000 --- a/modules/amqp-ping/Dockerfile +++ /dev/null @@ -1,31 +0,0 @@ -# amqp-ping's runtime: the tool runtime, carrying this module's compiled code. -# -# **Built from this module's own directory and nothing else.** The sdk and the tool runtime are not -# copied out of neighbouring checkouts — they are in the base image, which is published like any -# other artifact. That is what makes this buildable by the mesh from a repository and a path -# (novox/hq ADR 0069) rather than only on a workstation that happens to have the siblings. -# -# Two bases, named rather than pinned: the image this is COMPILED in, and the image it RUNS in. -# They are different images on purpose — the first carries a compiler and the second must not, or -# every running container would carry one it never invokes. The mesh answers both with the copies it -# holds, because a fingerprint written here would name one particular copy and no other mesh has it -# (novox/hq issue 044). Declared in module.json's `build.on`; deliberately no defaults, so a build -# nobody told stops here and says which module to build first. -ARG BUILD_BASE -ARG RUNTIME_BASE - -FROM ${BUILD_BASE} AS build -# Compiled under /app/modules, so resolving `@novox/mesh-sdk` walks up to the base's own -# node_modules — the module is compiled against exactly the sdk it will run against. -WORKDIR /app/modules/amqp-ping -COPY . . -# The compiler is invoked by its real path, not through node_modules/.bin. Those are symlinks to -# a launcher that requires its library relatively, and the base image resolves them when copying — -# leaving a launcher whose relative require no longer points at anything. -RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts \ - --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist - -FROM ${RUNTIME_BASE} -COPY --from=build /app/modules/amqp-ping/dist /app/modules/amqp-ping/dist -# Declared rather than derived from which files happen to exist: the module knows what it serves. -ENV MESH_TOOL_MODULES=/app/modules/amqp-ping/dist/index.js diff --git a/modules/amqp-ping/client.ts b/modules/amqp-ping/client.ts deleted file mode 100644 index ee5a693..0000000 --- a/modules/amqp-ping/client.ts +++ /dev/null @@ -1,191 +0,0 @@ -// 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 deleted file mode 100644 index f3338c3..0000000 --- a/modules/amqp-ping/index.ts +++ /dev/null @@ -1,53 +0,0 @@ -// 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(); -} -// changed by the one-node test at build 66e54af151df -// changed by the one-node test at build 4fb41636cffd -// changed by the one-node test at build a69f083bf6a6 diff --git a/modules/amqp-ping/module.json b/modules/amqp-ping/module.json deleted file mode 100644 index 68e947c..0000000 --- a/modules/amqp-ping/module.json +++ /dev/null @@ -1,89 +0,0 @@ -{ - "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}\n" - }, - { - "id": "net", - "type": "network", - "name": "amqp-ping" - }, - { - "id": "runtime", - "type": "container", - "name": "amqp-ping", - "network": "amqp-ping", - "volumes": [ - "/var/lib/mesh/amqp-ping/broker:/run/secrets/broker:ro", - "/var/lib/amqp-ping/amqp.secret:/run/secrets/amqp:ro" - ], - "env": { - "MESH_BROKER_FILE": "/run/secrets/broker", - "MESH_AMQP_PASSWORD_FILE": "/run/secrets/amqp" - }, - "env-file": [ - "/var/lib/amqp-ping/amqp.env" - ], - "restart-on": [ - "amqp-env" - ], - "artifact": "runtime" - } - ], - "build": { - "on": [ - { - "arg": "BUILD_BASE", - "module": "mesh-tools", - "artifact": "build" - }, - { - "arg": "RUNTIME_BASE", - "module": "mesh-tools", - "artifact": "runtime" - } - ], - "artifacts": [ - { - "name": "runtime", - "kind": "image", - "from": "Dockerfile" - } - ] - } -} diff --git a/modules/amqp-ping/package.json b/modules/amqp-ping/package.json deleted file mode 100644 index d8883aa..0000000 --- a/modules/amqp-ping/package.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "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 deleted file mode 100644 index 63b6d63..0000000 --- a/modules/amqp-ping/tsconfig.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "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/Dockerfile b/modules/lavinmq/Dockerfile deleted file mode 100644 index a08580e..0000000 --- a/modules/lavinmq/Dockerfile +++ /dev/null @@ -1,42 +0,0 @@ -# lavinmq's runtime: the tool runtime, carrying this module's compiled bootstrap, provisioner, -# tools and event consumer. -# -# **Built from this module's own directory and nothing else.** The sdk is in the base image, so -# nothing is copied out of a neighbouring checkout — which is what lets the mesh build this from a -# repository and a path (novox/hq ADR 0069) rather than only on a workstation that happens to have -# the siblings. -# -# Two bases, named rather than pinned: the image this is COMPILED in, and the image it RUNS in. -# They are different images on purpose — the first carries a compiler and the second must not, or -# every running container would carry one it never invokes. The mesh answers both with the copies it -# holds, because a fingerprint written here would name one particular copy and no other mesh has it -# (novox/hq issue 044). Declared in module.json's `build.on`; deliberately no defaults, so a build -# nobody told stops here and says which module to build first. -ARG BUILD_BASE -ARG RUNTIME_BASE - -FROM ${BUILD_BASE} AS build -# Compiled under /app/modules so `@novox/mesh-sdk` resolves upward into the base's own -# node_modules — the module is compiled against exactly the sdk it will run against. -WORKDIR /app/modules/lavinmq -COPY . . -# The compiler is invoked by its real path rather than through node_modules/.bin, whose entries are -# symlinks to a launcher that requires its library relatively — resolved away when the base image -# was assembled. -# -# Four entrypoints and a client, because this module is four things: a run-once bootstrap that -# writes the broker's configuration before it first starts, a provisioner that grants consumers -# their own vhost and user, a set of tools, and an event consumer. -RUN node /app/node_modules/typescript/bin/tsc \ - client.ts index.ts bootstrap/index.ts provisioner/index.ts tools/index.ts \ - --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist - -FROM ${RUNTIME_BASE} -COPY --from=build /app/modules/lavinmq/dist /app/modules/lavinmq/dist -# What the runtime loads from this module in serve mode: its event consumer, its tools, and its -# provisioner — all three in one process, so the provisioner's reconcile loop runs with the broker -# connected (novox/hq issues 060/061; the provisioner used to be run as a separate container `args` -# command, which meant it served no tools and, once, ran nowhere at all). The run-once bootstrap is -# NOT listed here — it is named in its own container's `args`, because it runs to completion before -# the broker starts rather than serving. One image, because they are one module and share a client. -ENV MESH_TOOL_MODULES=/app/modules/lavinmq/dist/index.js,/app/modules/lavinmq/dist/tools/index.js,/app/modules/lavinmq/dist/provisioner/index.js diff --git a/modules/lavinmq/bootstrap/index.ts b/modules/lavinmq/bootstrap/index.ts deleted file mode 100644 index 3ba7369..0000000 --- a/modules/lavinmq/bootstrap/index.ts +++ /dev/null @@ -1,45 +0,0 @@ -// 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 deleted file mode 100644 index 04a30c4..0000000 --- a/modules/lavinmq/client.ts +++ /dev/null @@ -1,183 +0,0 @@ -// 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: ".*" }); - } - - /** - * Whether a consumer's user exists with exactly this password and full permissions on its own - * vhost. Read-only: the stored hash is salted SHA-256, the scheme `rabbitHash` writes, so the - * password is checked by hashing it with the stored salt rather than by logging in. `false` when - * the user or its permission is gone or the password differs; an unreachable API rejects - * (novox/hq issue 120). - */ - async holdsConsumer(login: string, password: string): Promise { - const v = encodeURIComponent(login); - const u = encodeURIComponent(login); - const user = await this.getOrNull<{ password_hash?: string; hashing_algorithm?: string }>(`/users/${u}`); - if (!user?.password_hash) return false; - if (user.hashing_algorithm && !/sha256/i.test(user.hashing_algorithm)) { - throw new Error(`lavinmq user ${login} is hashed with ${user.hashing_algorithm}, which this check cannot verify`); - } - const stored = Buffer.from(user.password_hash, "base64"); - if (stored.length < 5 || rabbitHash(password, stored.subarray(0, 4)) !== user.password_hash) return false; - const perm = await this.getOrNull<{ configure?: string; write?: string; read?: string }>(`/permissions/${v}/${u}`); - return perm?.configure === ".*" && perm?.write === ".*" && perm?.read === ".*"; - } - - /** A GET that answers null for a 404 and rejects on anything else that is not 2xx. */ - private async getOrNull(path: string): Promise { - const resp = await fetch(`${this.conn.base}/api${path}`, { - headers: { - Authorization: "Basic " + Buffer.from(`${this.conn.adminUser}:${this.conn.adminPassword}`).toString("base64"), - }, - }); - if (resp.status === 404) return null; - if (!resp.ok) throw new Error(`lavinmq management API GET ${path} -> ${resp.status}: ${await resp.text()}`); - return (await resp.json()) as T; - } - - /** 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 deleted file mode 100644 index 76f5128..0000000 --- a/modules/lavinmq/index.ts +++ /dev/null @@ -1,25 +0,0 @@ -// 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("amqp.provisioned", async (e) => { - console.log(`[lavinmq] broker provisioned for ${e.body.consumer ?? "?"} (user ${e.body.user}, vhost ${e.body.vhost})`); -}); - -await on("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 deleted file mode 100644 index f78368a..0000000 --- a/modules/lavinmq/module.json +++ /dev/null @@ -1,151 +0,0 @@ -{ - "module": "lavinmq", - "version": "1", - "provides": [ - { - "name": "amqp", - "scope": "mesh" - } - ], - "claims": [ - { - "name": "mesh-broker", - "scope": "mesh" - } - ], - "capabilities": [ - "container-runtime" - ], - "emits": [ - "amqp.provisioned", - "amqp.deprovisioned" - ], - "consumes": [ - "lavinmq.amqp.provisioned", - "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": { - "admin": "/var/lib/lavinmq-module/admin.secret", - "broker": "/var/lib/mesh/lavinmq/broker" - }, - "listens": [ - { - "port": 5671, - "protocol": "tcp", - "from": "mesh", - "why": "the mesh bus \u2014 amqps, every module's events and the control plane, reached over the overlay" - }, - { - "port": 5672, - "protocol": "tcp", - "from": "mesh", - "why": "modules on any machine that were granted a queue" - } - ], - "guards": [ - 15672 - ], - "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": "broker-data", - "type": "directory", - "path": "/var/lib/mesh-broker", - "mode": "0700" - }, - { - "id": "server", - "type": "container", - "name": "mesh-broker", - "image": "cloudamqp/lavinmq@sha256:3eb54c12916d700a978c2ea86e6362cd4974b0e3189508718006d4e6d341246b", - "ports": [ - "5671:5671", - "5672:5672", - "127.0.0.1:15672:15672" - ], - "volumes": [ - "/var/lib/mesh-broker:/var/lib/lavinmq", - "/var/lib/mesh-broker-tls:/tls:ro" - ], - "args": [ - "--amqps-port=5671", - "--cert=/tls/tls.crt", - "--key=/tls/tls.key" - ] - }, - { - "id": "runtime", - "type": "container", - "name": "mesh-lavinmq", - "artifact": "runtime", - "network": "host", - "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/admin.secret:/run/secrets/admin:ro" - ], - "env": { - "MESH_BROKER_FILE": "/run/secrets/broker", - "MESH_RECEIVES": "/var/lib/lavinmq-module/grants/mesh.json", - "MESH_PROVISION_LAVINMQ": "http://127.0.0.1:15672", - "MESH_PROVISION_ADMIN_USER": "guest", - "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/admin" - } - } - ], - "build": { - "on": [ - { - "arg": "BUILD_BASE", - "module": "mesh-tools", - "artifact": "build" - }, - { - "arg": "RUNTIME_BASE", - "module": "mesh-tools", - "artifact": "runtime" - } - ], - "artifacts": [ - { - "name": "runtime", - "kind": "image", - "from": "Dockerfile" - } - ] - }, - "accesses": [ - { - "path": "/var/lib/mesh-broker-tls", - "mode": "read" - } - ] -} diff --git a/modules/lavinmq/package.json b/modules/lavinmq/package.json deleted file mode 100644 index 9e21bb1..0000000 --- a/modules/lavinmq/package.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "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.1" - }, - "devDependencies": { - "@types/node": "^22.0.0", - "typescript": "^5.6.0" - } -} diff --git a/modules/lavinmq/provisioner/index.ts b/modules/lavinmq/provisioner/index.ts deleted file mode 100644 index acd1d97..0000000 --- a/modules/lavinmq/provisioner/index.ts +++ /dev/null @@ -1,56 +0,0 @@ -// 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("amqp.provisioned", { - consumer: p.consumer ?? "", - user: p.as, - vhost: p.as, - }); - }, - - async remove(p: { as: string }): Promise { - await lavinmq.removeConsumer(p.as); - await announce("amqp.deprovisioned", { user: p.as, vhost: p.as }); - }, - // Asked every minute by the harness: whether the backend still holds this consumer exactly as - // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). - async holds(p: Provision): Promise { - return lavinmq.holdsConsumer(p.as, p.password); - }, -}); diff --git a/modules/lavinmq/tools/index.ts b/modules/lavinmq/tools/index.ts deleted file mode 100644 index 55489c6..0000000 --- a/modules/lavinmq/tools/index.ts +++ /dev/null @@ -1,32 +0,0 @@ -// 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 deleted file mode 100644 index 95f347f..0000000 --- a/modules/lavinmq/tsconfig.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "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"] -}