// mosquitto's admin client — mosquitto's own code, living in the module (novox/hq ADR 0039). Both // this module's tools and its provisioner import it, and nothing outside mosquitto does. // // Client, role and ACL administration is driven through `mosquitto_ctrl dynsec`, not a hand-rolled // MQTT stack: the module may take NO npm dependency beyond @novox/mesh-sdk, and mosquitto ships the // exact admin client for its Dynamic Security plugin — so it shells out to it, the same way postgres // drives itself through `psql`, minio through `mc` and mailu through `doveadm`. One boundary, // `ctl()`, and every method is built on it. // // Why the Dynamic Security plugin and not `password_file`: dynsec creates and revokes clients while // the broker runs, over an admin connection, with no broker restart and no file the host rewrites — // the true analog of redis's runtime ACL users. A `password_file` would have to be re-read on a // SIGHUP the runtime cannot cleanly send across containers, and — declared as a managed file — would // be rewritten by the host on every reconcile, wiping every provisioned user (novox/nox issue 011). // The one cost dynsec carries is the bootstrap file; see initBootstrapFile() and the module README. import { randomBytes } from "node:crypto"; import { readFileSync } from "node:fs"; import { execFile } from "node:child_process"; import { promisify } from "node:util"; const run = promisify(execFile); export interface MqttConn { readonly host: string; readonly port: number; /** The Dynamic Security admin client the runtime authenticates as. */ readonly adminUser: string; readonly adminPassword: string; } export class MosquittoClient { constructor(private readonly conn: MqttConn) {} /** * Build from the module's resolved environment. Reads MESH_MQTT_* first (the documented names), * falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner * container. Throws if it cannot find a host and an admin password — the right failure, because * without them nothing it does can work. */ static fromEnv(env: NodeJS.ProcessEnv = process.env): MosquittoClient { const endpoint = env.MESH_PROVISION_MQTT ?? ""; // "host:port" const host = env.MESH_MQTT_HOST ?? (endpoint ? endpoint.split(":")[0] : undefined); const port = Number(env.MESH_MQTT_PORT ?? (endpoint.includes(":") ? endpoint.split(":")[1] : "") ?? "1883") || 1883; const adminUser = env.MESH_MQTT_ADMIN_USER ?? env.MESH_PROVISION_ADMIN_USER ?? "mesh-admin"; const adminPassword = env.MESH_MQTT_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); if (!host || !adminPassword) { throw new Error( "mosquitto host or admin password is not set — mosquitto's own code cannot reach the broker", ); } return new MosquittoClient({ host, port, adminUser, adminPassword: adminPassword ?? "" }); } get host(): string { return this.conn.host; } get port(): number { return this.conn.port; } /** * Run one `mosquitto_ctrl dynsec ` command against the broker as the admin client and return * its stdout. Connects over MQTT with the verified connect flags `-h`/`-p`/`-u`/`-P`. A failed * command rejects — a failure here is an error, not a success with a warning. * * The exit code cannot carry that verdict: `mosquitto_ctrl` 2.0.x exits 0 from its dynsec * subcommands **even when they fail** — a "Client not found", an "already exists", a rejected * "Connection error: Not authorized", an "Unable to connect" all return status 0 and report the * failure only as a line of text, on stdout or stderr (verified live against 2.0.11). Trusting the * exit code is exactly how a `createClient` the broker refused reads back as a provisioned * consumer. So the combined output is scanned for the tool's error markers and a match is raised as * the failure it is. * * The admin password rides on argv (`-P`): mosquitto_ctrl 2.x exposes no password env var and no * password file for a broker connection — its only non-interactive mechanism is `-P`, its only * other mechanism an interactive prompt. This is a real mosquitto limitation, not a choice; unlike * psql's PGPASSWORD there is nothing cleaner to reach for. The exposure is momentary and confined * to this single-purpose runtime container; see the module README. */ async ctl(...args: string[]): Promise { const base = [ "-h", this.conn.host, "-p", String(this.conn.port), "-u", this.conn.adminUser, "-P", this.conn.adminPassword, ]; const { stdout, stderr } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], { maxBuffer: 16 << 20, }); const failure = ctlError(`${stdout}\n${stderr}`); if (failure) { throw new Error(`mosquitto_ctrl dynsec ${args[0] ?? ""} failed: ${failure}`); } return stdout; } /** Whether a dynsec client with this username already exists. */ async clientExists(username: string): Promise { try { await this.ctl("getClient", username); return true; } catch (err) { // Only "not found" means the client is genuinely absent. Any other failure — an auth // rejection, an unreachable broker — must NOT be read as "absent": that would send us down // the createClient path and bury the real error. Re-raise anything that is not a clean miss. if (/not\s*found|does not exist|no such/i.test(String(err))) return false; throw err; } } /** * Whether this client already carries this role. mosquitto_ctrl has no idempotent addClientRole: * re-binding a role the client already has does not report "already exists" — it reports a bare * "Internal error" (verified against 2.0.11), indistinguishable from a genuine fault, so it cannot * be swallowed by message. Instead the binding is checked first. `getClient` lists each assigned * role under its "Roles:" heading as ` (priority: N)`; the role is matched as a whole token. */ async clientHasRole(username: string, role: string): Promise { let out: string; try { out = await this.ctl("getClient", username); } catch { return false; // no such client (or unreadable) — it certainly has no role } return new RegExp(`(^|\\s)${escapeRegExp(role)}\\s+\\(priority`, "m").test(out); } /** * Create (or reset to a known state) a client scoped to one topic namespace, idempotently. The * client is confined to `/#` by a same-named role: it may publish to, subscribe to and * receive on exactly its own subtree and nothing else — the MQTT analog of redis's keyspace-scoped * ACL user. Called again for an existing client, it resets the password and re-asserts the ACLs. */ async createScopedClient(username: string, password: string, topicPrefix: string): Promise { const role = username; // one role per client, named for it const pattern = `${topicPrefix}/#`; if (await this.clientExists(username)) { await this.ctl("setClientPassword", username, password); } else { await this.ctl("createClient", username, "-p", password); } // A role carrying exactly this client's topic ACLs. createRole, addRoleACL and addClientRole are // all one-shot: each rejects with an "already exists" when re-run against a role/ACL/binding it // created on a previous reconcile. That rejection is the intended terminal state — the ACL is // deterministic (`/#`, allow), so re-adding the identical entry is a no-op — so it is // swallowed. (Until the exit code was fixed this was invisible: the tool returned 0 and the // rejection was lost; now it surfaces, and each of these adds must tolerate its own idempotent // re-run explicitly.) await ignoreExisting(this.ctl("createRole", role)); for (const acl of ["publishClientSend", "publishClientReceive", "subscribePattern"]) { // allow (1) this client to send to, receive on, and subscribe under its own subtree. await ignoreExisting(this.ctl("addRoleACL", role, acl, pattern, "allow")); } // Bind the role only when it is not already bound — addClientRole is the one call whose // idempotent re-run cannot be recognised by message (see clientHasRole). if (!(await this.clientHasRole(username, role))) { await this.ctl("addClientRole", username, role); } } /** Remove a client and the per-client role created for it, idempotently. */ async deleteScopedClient(username: string): Promise { await ignoreMissing(this.ctl("deleteClient", username)); await ignoreMissing(this.ctl("deleteRole", username)); } /** The dynsec client list, parsed from `listClients`. */ async listClients(): Promise { const out = await this.ctl("listClients"); return out .split(/\r?\n/) .map((l) => l.trim()) .filter((l) => l.length > 0); } /** * Write the Dynamic Security bootstrap file offline, creating the admin client the plugin loads at * broker startup. This is a one-time seed, NOT part of the reconcile loop: run once before the * broker first starts, against the same path the broker's `plugin_opt_config_file` names. It must * never be a host-reconciled managed file — see the module README and novox/nox issue 011. */ async initBootstrapFile(configFile: string): Promise { // `dynsec init [admin-password]` is an offline file operation — it does // not connect to the broker. The password is a positional argument (omitting it prompts). await run("mosquitto_ctrl", ["dynsec", "init", configFile, this.conn.adminUser, this.conn.adminPassword], { maxBuffer: 16 << 20, }); } } /** Generate a URL-safe password with no argv- or MQTT-hostile characters. */ export function generatePassword(): string { return randomBytes(24).toString("base64url"); } /** Escape a string for literal use inside a RegExp. */ function escapeRegExp(s: string): string { return s.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); } /** * Find the failure `mosquitto_ctrl` reported in text while still exiting 0. Its dynsec subcommands * surface errors in three shapes, on stdout or stderr: * - `: Error: ` e.g. "createClient: Error: Client already exists" * - `Connection error: ` e.g. "Connection error: Not authorized" * - `Unable to connect ()` e.g. "Unable to connect (Lookup error.)." * The only other line it prints unprompted is the "running without encryption" warning, which * carries none of these markers. Returns the offending line, or undefined when the output is clean. */ function ctlError(output: string): string | undefined { for (const raw of output.split(/\r?\n/)) { const line = raw.trim(); if (!line) continue; if (/(^|:\s)Error:/i.test(line) || /^Connection error:/i.test(line) || /^Unable to connect/i.test(line)) { return line; } } return undefined; } /** Swallow a "already exists" failure so create paths are idempotent; rethrow anything else. */ async function ignoreExisting(p: Promise): Promise { try { await p; } catch (err) { if (!/exist/i.test(String(err))) throw err; } } /** Swallow a "not found" failure so delete paths are idempotent; rethrow anything else. */ async function ignoreMissing(p: Promise): Promise { try { await p; } catch (err) { if (!/not\s*found|does not exist|no such/i.test(String(err))) throw err; } } function readSecretFile(path: string | undefined): string | undefined { if (!path) return undefined; try { return readFileSync(path, "utf8").trim(); } catch { return undefined; } }