Files
mesh-catalog/modules/mosquitto/client.ts
T
jochen 0cb0f814b4 Every credential provider says whether it still holds a consumer
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.
2026-09-26 01:09:52 +02:00

318 lines
14 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,
];
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<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);
} 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}`);
return this.clientHasRole(username, username);
}
/** 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);
});
});
}