Node-RED's one broker node pointed at zurag.be:1884, where nothing listens. nodered now requires
mqtt-topic (asking for every topic: flows follow the devices' own) and a run-once `mqtt` step —
declared last, restarted when the binding, credential or settings change — points the mesh's broker
nodes at the bound broker through Node-RED's admin API with the module's api-token: the node the
step makes itself when none is named, or the ones an assignment names in `mqtt.brokers`. Only host,
port, TLS and the login change; the broker is asked first whether it takes the login; the deploy is
against the revision read ("nodes", so only that node restarts) and a digest makes a rerun a no-op.
A broker node nobody named is never touched. settings.js keeps `mqtt` and `topics` out of Node-RED.
233 lines
11 KiB
TypeScript
233 lines
11 KiB
TypeScript
// 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<string, unknown>;
|
|
}
|
|
|
|
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<void>;
|
|
}
|
|
|
|
export interface Marks {
|
|
get(name: string): Promise<string | undefined>;
|
|
set(name: string, digest: string): Promise<void>;
|
|
}
|
|
|
|
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<Outcome[]> {
|
|
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;
|
|
}
|