holds() for postgres, mssql, mongodb, minio, lavinmq, mosquitto, mailu and gitea, so the harness makes again a login the backend lost (hq issue 120). Each checks the mesh's password as the consumer presents it, or compares it read-only, and returns false only when the backend says the credential is absent or wrong; an unreachable backend throws.
184 lines
8.8 KiB
TypeScript
184 lines
8.8 KiB
TypeScript
// 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<unknown> {
|
|
const headers: Record<string, string> = {
|
|
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<boolean> {
|
|
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<void> {
|
|
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 `<login>` to vhost `<login>` and can reach nothing
|
|
* else (novox/hq ADR 0048).
|
|
*/
|
|
async createConsumer(login: string, password: string): Promise<void> {
|
|
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<boolean> {
|
|
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<T>(path: string): Promise<T | null> {
|
|
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<void> {
|
|
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;
|
|
}
|
|
}
|