Files
mesh-catalog/modules/mosquitto/client.ts
T
jochen 6fd93afc6c Review fixes: holds and create agree, and no password leaves a check
create re-enables what holds refuses (mssql login, mosquitto client,
mailu mailbox, gitea user) and clears an expired postgres password, so
no disabled account loops. mssql and mongodb checks take the password
from the environment, never argv; mosquitto_ctrl failures no longer
repeat -P. mosquitto reads 'could not ask' as an error, not absence.
mailu checks existence and enabled only: its imap passdb cannot verify
a password. mssql checks the user's SID; gitea pages teams at 50.
2026-09-26 01:24:32 +02:00

343 lines
16 KiB
TypeScript

// 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 { 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 <args>` 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<string> {
const base = [
"-h", this.conn.host,
"-p", String(this.conn.port),
"-u", this.conn.adminUser,
"-P", this.conn.adminPassword,
];
let stdout: string;
let stderr: string;
try {
({ stdout, stderr } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], {
maxBuffer: 16 << 20,
timeout: 30_000,
}));
} catch (err) {
// A failed run's message repeats its argv, the admin password (-P) included; say what failed
// without it.
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;
}
/** Whether a dynsec client with this username already exists. */
async clientExists(username: string): Promise<boolean> {
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 `<role> (priority: N)`; the role is matched as a whole token.
*/
async clientHasRole(username: string, role: string): Promise<boolean> {
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 `<prefix>/#` 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<void> {
const role = username; // one role per client, named for it
const pattern = `${topicPrefix}/#`;
if (await this.clientExists(username)) {
await this.ctl("setClientPassword", username, password);
// 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.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 (`<prefix>/#`, 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);
}
}
/**
* Whether a consumer's client accepts exactly this password and still carries its own role.
* 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): Promise<boolean> {
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 out: string;
try {
out = await this.ctl("getClient", username);
} catch (err) {
if (/not\s*found|does not exist|no such/i.test(String(err))) return false;
throw err;
}
return new RegExp(`(^|\\s)${escapeRegExp(username)}\\s+\\(priority`, "m").test(out);
}
/** Remove a client and the per-client role created for it, idempotently. */
async deleteScopedClient(username: string): Promise<void> {
await ignoreMissing(this.ctl("deleteClient", username));
await ignoreMissing(this.ctl("deleteRole", username));
}
/** The dynsec client list, parsed from `listClients`. */
async listClients(): Promise<string[]> {
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<void> {
// `dynsec init <file> <admin-username> [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:
* - `<command>: Error: <message>` e.g. "createClient: Error: Client already exists"
* - `Connection error: <message>` e.g. "Connection error: Not authorized"
* - `Unable to connect (<message>)` 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<string>): Promise<void> {
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<string>): Promise<void> {
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.
*/
function mqttConnack(host: string, port: number, username: string, password: string): Promise<number> {
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);
});
});
}