// Node-RED's MQTT broker config node, pointed at the broker the mesh bound — `mqtt-topic`. // // **Why a step.** Node-RED keeps a broker as a config node in its flows (`flows.json`) and the login // and password in its encrypted credentials file, both of them Node-RED's to write. So this reads the // binding and the pair credential and makes the broker node say the same thing through Node-RED's // admin API — `GET /flows`, then `POST /flows` with the changed node and its `credentials`, deployed // as "nodes" so only what changed restarts — with the module's own `api-token`. // // **Which broker nodes are the mesh's.** Never guessed: a flow may talk to a broker that has nothing // to do with this mesh. The step owns the node it creates itself (id `mesh-mqtt-topic`, "mesh: // mqtt-topic") and the ones an assignment names in settings (`mqtt.brokers`: node ids — how ace's // existing broker node, which every one of its MQTT flows uses, is handed over). With none named and // none made yet, it makes one, so a fresh Node-RED has a broker the flows can pick. // // **Only the connection, and only when it differs.** Host, port, TLS off (the broker serves plain // MQTT), login, password. Every other field of the node — client id, keepalive, birth/close/will // messages — is left as it is. Node-RED never hands a stored password back, so the step keeps a // digest of what it last wrote: equal host/port/login and an equal digest is "already as the mesh // says". // // **Nothing loses its connection without someone seeing it.** The broker is asked first whether it // takes the delivered login; if not, nothing is written and the step fails saying why. // // Pure logic over two seams (Node-RED, the broker), tested against fakes (test/mqtt.test.ts). import { createHash } from "node:crypto"; import type { Probe } from "./probe.js"; export const PROVISION = "mqtt-topic"; /** The id and name of the broker node the step makes when none is named. */ export const MESH_BROKER_ID = "mesh-mqtt-topic"; export const MESH_BROKER_NAME = "mesh: mqtt-topic"; /** What the mesh wrote at `binds.mqtt-topic`. */ export interface Binding { provision?: string; from?: string; at?: string; as?: string; serves?: Record; } export type Outcome = | { what: string; result: "unchanged"; note?: string } | { what: string; result: "written"; fields: string[]; note?: string } | { what: string; result: "refused"; problem: string }; /** A flow node; a broker config node carries `broker`, `port`, `usetls`. */ export interface FlowNode { id: string; type: string; [key: string]: unknown; } /** Node-RED's admin API, as the step uses it. */ export interface NodeRed { /** The whole flow configuration and its revision (API v2). */ flows(): Promise<{ rev: string; flows: FlowNode[] }>; /** A node's stored credentials as Node-RED shows them: the user, and only whether a password is set. */ credentials(type: string, id: string): Promise<{ user?: string; has_password?: boolean }>; /** Deploy the configuration against the revision it was read at; "nodes" restarts only what changed. */ deploy(rev: string, flows: FlowNode[]): Promise; } export interface Marks { get(name: string): Promise; set(name: string, digest: string): Promise; } export interface Deps { nodered: NodeRed; probe: Probe; marks: Marks; } export interface Wanted { host: string; port: number; user: string; password: string; } export function digest(...parts: (string | number)[]): string { return createHash("sha256").update(parts.map(String).join("\u0000")).digest("hex"); } function isLoopback(host: string): boolean { const h = host.toLowerCase(); return h === "localhost" || h === "::1" || h === "[::1]" || /^127\./.test(h); } /** * The broker and login the mesh says Node-RED uses. A loopback `at` — what the mesh hands a machine * that is not on the private network — is refused: from Node-RED's own container it is Node-RED. */ export function wanted(binding: Binding | undefined, credential: string | undefined): { ok: true; want: Wanted } | { ok: false; problem: string } { if (!binding) return { ok: false, problem: `no binding for ${PROVISION} was delivered — the mesh writes it before this step runs` }; const host = typeof binding.at === "string" ? binding.at.trim() : ""; if (!host) return { ok: false, problem: `the ${PROVISION} binding names no host (at)` }; if (isLoopback(host)) { return { ok: false, problem: `the ${PROVISION} binding says the broker is at ${host}, which from Node-RED's own container is Node-RED ` + `itself; put the machine on the private network so the broker has an address Node-RED can dial`, }; } const port = Number(binding.serves?.port); if (!Number.isInteger(port) || port <= 0 || port > 65535) { return { ok: false, problem: `the ${PROVISION} binding serves no usable port (${String(binding.serves?.port)})` }; } const scheme = binding.serves?.scheme; if (scheme !== undefined && scheme !== "mqtt") return { ok: false, problem: `the ${PROVISION} binding serves scheme ${String(scheme)}; this step writes plain MQTT` }; const user = typeof binding.as === "string" ? binding.as.trim() : ""; if (!user) return { ok: false, problem: `the ${PROVISION} binding names no login (as)` }; const password = (credential ?? "").replace(/\n$/, ""); if (!password) return { ok: false, problem: `the ${PROVISION} credential is empty or was not delivered` }; return { ok: true, want: { host, port, user, password } }; } /** The broker node ids an assignment named in settings (`mqtt.brokers`), or none. */ export function namedBrokers(settings: unknown): string[] { const brokers = (settings as { mqtt?: { brokers?: unknown } } | undefined)?.mqtt?.brokers; return Array.isArray(brokers) ? brokers.filter((b): b is string => typeof b === "string" && b.length > 0) : []; } /** A new broker node, Node-RED 5's defaults, pointed at the broker. */ export function newBrokerNode(want: Wanted): FlowNode { return { id: MESH_BROKER_ID, type: "mqtt-broker", name: MESH_BROKER_NAME, broker: want.host, port: String(want.port), clientid: "", autoConnect: true, usetls: false, protocolVersion: "4", keepalive: "60", cleansession: true, autoUnsubscribe: true, birthTopic: "", birthQos: "0", birthRetain: "false", birthPayload: "", birthMsg: {}, closeTopic: "", closeQos: "0", closeRetain: "false", closePayload: "", closeMsg: {}, willTopic: "", willQos: "0", willRetain: "false", willPayload: "", willMsg: {}, userProps: "", sessionExpiry: "", }; } const markFor = (id: string, w: Wanted): string => digest("nodered-mqtt", id, w.host, w.port, w.user, w.password); function scrub(err: unknown, secret: string): string { let text = err instanceof Error ? err.message : String(err); for (const form of new Set([secret, encodeURIComponent(secret)])) text = text.split(form).join("***"); return text; } /** * Bring the mesh's broker nodes in line with the binding: one outcome per node. Never throws. A node * the settings name that is not in the flows is refused (the others are still put right). */ export async function reconcileBrokers(deps: Deps, binding: Binding | undefined, credential: string | undefined, named: readonly string[]): Promise { const w = wanted(binding, credential); if ("problem" in w) return [{ what: "mqtt", result: "refused", problem: w.problem }]; const want = w.want; let note: string | undefined; try { const probe = await deps.probe(want.host, want.port, want.user, want.password, "#"); if (probe.connack === 4 || probe.connack === 5) { return [{ what: "mqtt", result: "refused", problem: `the broker at ${want.host}:${want.port} does not (yet) take the login ${want.user} with the delivered password ` + `(CONNACK ${probe.connack}); mosquitto's provisioner creates it from the grant — nothing was written`, }]; } if (probe.connack !== 0) return [{ what: "mqtt", result: "refused", problem: `the broker at ${want.host}:${want.port} answered CONNACK ${probe.connack}; nothing was written` }]; if (probe.suback === 0x80) note = `warning: ${want.user} may not subscribe to every topic; flows subscribing outside its grant will get nothing`; } catch (err) { return [{ what: "mqtt", result: "refused", problem: `the broker at ${want.host}:${want.port} could not be asked: ${scrub(err, want.password)}; nothing was written` }]; } const outcomes: Outcome[] = []; for (let attempt = 0; attempt < 2; attempt++) { outcomes.length = 0; try { const { rev, flows } = await deps.nodered.flows(); const targets = named.length > 0 ? [...named] : [MESH_BROKER_ID]; const changed: { id: string; fields: string[] }[] = []; for (const id of targets) { let node = flows.find((n) => n.id === id); if (node && node.type !== "mqtt-broker") { outcomes.push({ what: `broker ${id}`, result: "refused", problem: `node ${id} is a ${node.type}, not an mqtt-broker` }); continue; } if (!node) { if (id !== MESH_BROKER_ID) { outcomes.push({ what: `broker ${id}`, result: "refused", problem: `the settings name broker node ${id}, and Node-RED's flows have no such node` }); continue; } node = newBrokerNode(want); flows.push(node); node.credentials = { user: want.user, password: want.password }; changed.push({ id, fields: ["node"] }); continue; } const fields: string[] = []; if (String(node.broker ?? "") !== want.host) fields.push("broker"); if (Number(node.port ?? 0) !== want.port) fields.push("port"); if (node.usetls === true) fields.push("usetls"); const creds = await deps.nodered.credentials("mqtt-broker", id); if ((creds.user ?? "") !== want.user) fields.push("user"); if (!creds.has_password || (await deps.marks.get(`broker-${id}`)) !== markFor(id, want)) fields.push("password"); if (fields.length === 0) { outcomes.push(note ? { what: `broker ${id}`, result: "unchanged", note } : { what: `broker ${id}`, result: "unchanged" }); continue; } node.broker = want.host; node.port = String(want.port); node.usetls = false; node.credentials = { user: want.user, password: want.password }; changed.push({ id, fields }); } if (changed.length > 0) { await deps.nodered.deploy(rev, flows); for (const c of changed) { await deps.marks.set(`broker-${c.id}`, markFor(c.id, want)); outcomes.push({ what: `broker ${c.id}`, result: "written", fields: c.fields, ...(note ? { note } : {}) }); } } return outcomes; } catch (err) { // A deploy against a revision someone else changed meanwhile (409) is read again once. if (attempt === 0 && /\b409\b/.test(String(err))) continue; return [...outcomes, { what: "mqtt", result: "refused", problem: scrub(err, want.password) }]; } } return outcomes; }