185 lines
7.9 KiB
TypeScript
185 lines
7.9 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 { 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 non-zero
|
|
* exit rejects — a failed command is an error here, not a success with a warning.
|
|
*
|
|
* 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 } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], {
|
|
maxBuffer: 16 << 20,
|
|
});
|
|
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 {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* 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 fails if it already exists; that
|
|
// is fine — the setRoleACL calls below assert the intended state either way.
|
|
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 this.ctl("addRoleACL", role, acl, pattern, "allow");
|
|
}
|
|
await ignoreExisting(this.ctl("addClientRole", username, role));
|
|
}
|
|
|
|
/** 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");
|
|
}
|
|
|
|
/** 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;
|
|
}
|
|
}
|