// 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 { connect as tcpConnect } from "node:net"; import { readFileSync } from "node:fs"; import { execFile } from "node:child_process"; import { basename, dirname } from "node:path"; import { promisify } from "node:util"; import { missingAcls, parseRoleAcls, staleAcls, wantedAcls } from "./topics.js"; 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; /** * The broker's own container, when `mosquitto_ctrl` is run inside it rather than from this * machine's packages. The broker's image carries the tool at the broker's version, and inside it * the broker listens on 127.0.0.1:1883 whatever port the machine publishes. */ readonly container?: string; /** The broker's image, to seed the security file before the broker has ever started. */ readonly image?: 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 ?? "", container: env.MESH_MQTT_CTRL_CONTAINER || undefined, image: env.MESH_MQTT_CTRL_IMAGE || undefined, }); } get host(): string { return this.conn.host; } get port(): number { return this.conn.port; } get adminUser(): string { return this.conn.adminUser; } /** * Run one `mosquitto_ctrl dynsec ` command against the broker as the admin client and return * its stdout. A failed command rejects — a failure here is an error, not a success with a warning. * * **No secret rides on a command line** (novox/hq issue 282). The admin's name and password reach * `mosquitto_ctrl` as an options file (`-o`), and a password a command sets reaches it at its own * prompt, both on stdin: a small shell, given no secret, reads the two admin lines into a 0600 file * that it removes as it exits, and hands the rest of stdin to the prompts. A password on argv is * kept by everything that records a command line — the container runtime records every exec's in * its event stream, which `docker events` shows to anyone who may ask the runtime — and the * earlier `-P ` put the broker's admin password there on every call. Verified against * mosquitto_ctrl 2.1.2: `-o` takes the connect options from a file, and `createClient` without * `-p` and `setClientPassword` without a password prompt for it twice, reading stdin. * `assertNoSecret` refuses, before anything runs, an argv that still carries one. * * The exit code cannot carry the verdict: `mosquitto_ctrl` 2.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. */ async ctl(...args: string[]): Promise { return this.ctlPrompted([], args); } /** `ctl`, answering the command's password prompt with `password` — typed twice, as asked. */ async ctlWithPassword(password: string, ...args: string[]): Promise { return this.ctlPrompted([password, password], args); } private async ctlPrompted(prompted: string[], args: string[]): Promise { const inside = this.conn.container !== undefined; const connect = ["-h", inside ? "127.0.0.1" : this.conn.host, "-p", inside ? "1883" : String(this.conn.port)]; oneLine("the admin user", this.conn.adminUser, true); oneLine("the admin password", this.conn.adminPassword, true); for (const p of prompted) oneLine("a client password", p, false); const input = [this.conn.adminUser, this.conn.adminPassword, ...prompted].map((l) => l + "\n").join(""); let stdout: string; let stderr: string; try { const [command, argv] = this.ctrl([...connect, "dynsec", ...args]); this.assertNoSecret([command, ...argv], prompted); ({ stdout, stderr } = await runWithInput(command, argv, input)); } catch (err) { if (err instanceof SecretOnArgvError) throw err; // Said without the argv or the input, so a failure carries nothing it was given. const e = err as { code?: unknown; signal?: unknown; stderr?: string; stdout?: string }; const detail = `${e.stderr ?? ""}${e.stdout ?? ""}`.trim().slice(0, 500); throw new Error(`mosquitto_ctrl dynsec ${args[0] ?? ""} could not run (${e.code ?? e.signal ?? "error"}): ${detail}`); } const failure = ctlError(`${stdout}\n${stderr}`); if (failure) { throw new Error(`mosquitto_ctrl dynsec ${args[0] ?? ""} failed: ${failure}`); } return stdout; } /** * Refuse to run a command whose argv carries a secret this client holds: the admin password, or a * password it was asked to set (novox/hq issue 282). Said by name, never by value. */ assertNoSecret(argv: readonly string[], others: readonly string[] = []): void { const secrets: Array<[string, string]> = [["the broker's admin password", this.conn.adminPassword]]; for (const o of others) secrets.push(["a client password", o]); for (const [name, value] of secrets) { if (value && argv.some((a) => a.includes(value))) { throw new SecretOnArgvError(name, argv[0] ?? ""); } } } /** 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 granted exactly these topic filters, idempotently. * The grant is a same-named role carrying, for every filter, publish, receive and subscribe — and * nothing else: an ACL the role carries that the filters no longer name is removed, so narrowing a * consumer's `topics` narrows what it may do. By default the filters are the consumer's own * subtree, `/#` (see topics.ts). Called again for an existing client, it resets the password * and re-asserts the ACLs. * * Only the role named for this client is ever changed. A client or role the mesh did not make — * a device carried from the predecessor's password file, its `legacy-full-access` role — is never * read, changed or removed here. */ async createScopedClient(username: string, password: string, filters: readonly string[]): Promise { const role = username; // one role per client, named for it if (await this.clientExists(username)) { await this.ctlWithPassword(password, "setClientPassword", username); // A disabled client is refused like a wrong password, so the check the provisioner runs // reports it lost; applying again must enable it, or the two would disagree for ever. if (/Disabled:\s*true/i.test(await this.ctl("getClient", username))) { await this.ctl("enableClient", username); } } else { await this.ctlWithPassword(password, "createClient", username); } // createRole and addRoleACL are one-shot: each rejects with an "already exists" when re-run // against a role/ACL it created on a previous reconcile. That rejection is the intended terminal // state, so it is swallowed. await ignoreExisting(this.ctl("createRole", role)); const wanted = wantedAcls(filters); const current = parseRoleAcls(await this.ctl("getRole", role)); for (const acl of missingAcls(current, wanted)) { await ignoreExisting(this.ctl("addRoleACL", role, acl.type, acl.topic, "allow")); } // What the consumer no longer asks for — added before it narrowed its topics — is taken away. for (const acl of staleAcls(current, wanted)) { await ignoreMissing(this.ctl("removeRoleACL", role, acl.type, acl.topic)); } // 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); } } /** * Whether a consumer's client accepts exactly this password, still carries its own role, and that * role grants exactly these filters. Read-only. The password is checked the way the consumer is * checked, by an MQTT CONNECT as it, and the broker's CONNACK code is the answer: 0 accepted, * 4 bad credentials, 5 not authorised. Nothing rides on argv. An unreachable broker rejects * (novox/hq issue 120). */ async holdsClient(username: string, password: string, filters: readonly string[]): Promise { const code = await mqttConnack(this.conn.host, this.conn.port, username, password); if (code === 4 || code === 5) return false; if (code !== 0) throw new Error(`mosquitto refused ${username} with CONNACK ${code}`); // The role, asked directly: only "not found" means absent. Any other failure to ask rejects, // unlike clientHasRole, which reads every failure as "no role". let client: string; let role: string; try { client = await this.ctl("getClient", username); role = await this.ctl("getRole", username); } catch (err) { if (/not\s*found|does not exist|no such/i.test(String(err))) return false; throw err; } if (!new RegExp(`(^|\\s)${escapeRegExp(username)}\\s+\\(priority`, "m").test(client)) return false; const current = parseRoleAcls(role); const wanted = wantedAcls(filters); return missingAcls(current, wanted).length === 0 && staleAcls(current, wanted).length === 0; } /** 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 novox/nox issue 011. */ async initBootstrapFile(configFile: string): Promise { // `dynsec init ` is an offline file operation — it does not connect to // the broker. Without the optional password argument it prompts for the password twice, and the // answers are given on stdin (novox/hq issue 282): an argument would be on the command line of // the container below, which the runtime keeps in what it says about the container. oneLine("the admin password", this.conn.adminPassword, true); const input = `${this.conn.adminPassword}\n${this.conn.adminPassword}\n`; let command: string; let argv: string[]; if (this.conn.image) { // Before the broker has ever started there is no container to enter: a throwaway one from the // broker's own image writes the file into the directory the broker will mount. command = "docker"; argv = [ "run", "--rm", "-i", "--entrypoint", "mosquitto_ctrl", "-v", `${dirname(configFile)}:/mosquitto/data`, this.conn.image, "dynsec", "init", `/mosquitto/data/${basename(configFile)}`, this.conn.adminUser, ]; } else { command = "mosquitto_ctrl"; argv = ["dynsec", "init", configFile, this.conn.adminUser]; } this.assertNoSecret([command, ...argv]); try { await runWithInput(command, argv, input); } catch (err) { const e = err as { code?: unknown; signal?: unknown; stderr?: string; stdout?: string }; const detail = `${e.stderr ?? ""}${e.stdout ?? ""}`.trim().slice(0, 500); throw new Error(`mosquitto_ctrl dynsec init could not run (${e.code ?? e.signal ?? "error"}): ${detail}`); } } /** * Whether the broker takes this admin password: the CONNACK of a connect as the admin client — 0 * accepted, 4 or 5 refused — or undefined when the broker cannot be reached. Nothing rides on argv. */ async adminAccepted(password: string): Promise { let code: number; try { code = await mqttConnack(this.conn.host, this.conn.port, this.conn.adminUser, password); } catch { return undefined; } if (code === 0) return true; if (code === 4 || code === 5) return false; throw new Error(`mosquitto answered the admin client with CONNACK ${code}`); } /** The same client, connecting with another admin password. */ withAdminPassword(adminPassword: string): MosquittoClient { return new MosquittoClient({ ...this.conn, adminPassword }); } /** * How `mosquitto_ctrl` is run here: inside the broker's container when one is named, through the * shell that turns the first two lines of stdin into its options file (OPTIONS_SHELL). */ ctrl(argv: string[]): [string, string[]] { if (this.conn.container) { return ["docker", ["exec", "-i", this.conn.container, "sh", "-c", OPTIONS_SHELL, "mosquitto_ctrl", ...argv]]; } return ["sh", ["-c", OPTIONS_SHELL, "mosquitto_ctrl", ...argv]]; } } /** Generate a URL-safe password with no argv- or MQTT-hostile characters. */ export function generatePassword(): string { return randomBytes(24).toString("base64url"); } /** * The shell `mosquitto_ctrl` runs under: it reads the admin's name and password from the first two * lines of stdin into an options file only it can read, removed as it exits, and runs the tool with * `-o` naming it — so neither is ever an argument. The rest of stdin is left to the tool's own * password prompts: `read` takes a pipe one byte at a time and so stops at the second line. The * script holds no secret; it is the same text on every call. */ export const OPTIONS_SHELL = [ "umask 077", 'f=$(mktemp) || exit 1', "trap 'rm -f \"$f\"' EXIT", 'IFS= read -r u || exit 1', 'IFS= read -r p || exit 1', `printf '%s\\n' "-u $u" "-P $p" > "$f"`, 'mosquitto_ctrl -o "$f" "$@"', ].join("\n"); /** * Why a dynsec command would put a password on a command line, or undefined (novox/hq issue 282). * A command line is recorded by the container runtime for anyone who may ask it; a password is set * through the provisioner, which answers mosquitto_ctrl's prompt on stdin, never through this tool. */ export function passwordOnArgv(parts: readonly string[]): string | undefined { const [verb, ...rest] = parts; if (verb === "init") return "init writes a new store with a password; the module's bootstrap seeds the store"; if (verb === "createClient" && rest.includes("-p")) { return "createClient -p puts the password on a command line; a client is created by provisioning it"; } if (verb === "setClientPassword" && rest.length > 1) { return "setClientPassword with a password puts it on a command line; a client's password is the mesh's to set"; } return undefined; } /** A secret found on a command line about to run: named, never quoted. */ export class SecretOnArgvError extends Error { constructor(secret: string, program: string) { super(`refused to run ${program}: its command line would carry ${secret}, and a command line is recorded ` + `where anyone who may ask the container runtime reads it (novox/hq issue 282)`); this.name = "SecretOnArgvError"; } } /** * Refuse a value that cannot be handed over as one line of stdin — or, for the admin's, as one word * of an options file, which `mosquitto_ctrl` splits at spaces. Named, never quoted. */ function oneLine(what: string, value: string, oneWord: boolean): void { if (/[\r\n]/.test(value) || (oneWord && /\s/.test(value))) { throw new Error(`${what} ${oneWord ? "holds whitespace" : "spans lines"} and cannot be handed to mosquitto_ctrl on stdin`); } } /** Run a program with this text on its stdin, answering its stdout and stderr. */ function runWithInput(command: string, argv: string[], input: string): Promise<{ stdout: string; stderr: string }> { const p = run(command, argv, { maxBuffer: 16 << 20, timeout: 30_000 }); p.child.stdin?.on("error", () => { /* a program that exits before reading all of it is judged by its output */ }); p.child.stdin?.end(input); return p; } /** 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; } } /** * Connect once over MQTT 3.1.1 with a username and password, return the broker's CONNACK return code, * and disconnect. A clean session under a throwaway client id, so no consumer session is taken over. */ export function mqttConnack(host: string, port: number, username: string, password: string): Promise { const str = (v: string): Buffer => { const b = Buffer.from(v, "utf8"); const len = Buffer.alloc(2); len.writeUInt16BE(b.length); return Buffer.concat([len, b]); }; const variable = Buffer.concat([str("MQTT"), Buffer.from([4, 0xc2, 0, 10])]); // level 4; user+pass+clean; keepalive 10s const payload = Buffer.concat([str(`mesh-holds-${randomBytes(6).toString("hex")}`), str(username), str(password)]); let remaining = variable.length + payload.length; const lenBytes: number[] = []; do { let byte = remaining % 128; remaining = Math.floor(remaining / 128); if (remaining > 0) byte |= 0x80; lenBytes.push(byte); } while (remaining > 0); const packet = Buffer.concat([Buffer.from([0x10, ...lenBytes]), variable, payload]); return new Promise((resolve, reject) => { const socket = tcpConnect({ host, port }); let buf = Buffer.alloc(0); const timer = setTimeout(() => { socket.destroy(); reject(new Error(`no CONNACK from ${host}:${port} within 10s`)); }, 10_000); socket.on("connect", () => socket.write(packet)); socket.on("data", (chunk) => { buf = Buffer.concat([buf, chunk]); if (buf.length < 4) return; clearTimeout(timer); if (buf[0] !== 0x20) { socket.destroy(); reject(new Error(`unexpected MQTT packet 0x${buf[0].toString(16)} instead of CONNACK`)); return; } const code = buf[3]; if (code === 0) socket.end(Buffer.from([0xe0, 0])); // DISCONNECT else socket.destroy(); resolve(code); }); socket.on("error", (err) => { clearTimeout(timer); reject(err); }); }); }