diff --git a/modules/nodered/Dockerfile b/modules/nodered/Dockerfile index afa2f8c..3031ce6 100644 --- a/modules/nodered/Dockerfile +++ b/modules/nodered/Dockerfile @@ -13,7 +13,7 @@ ARG RUNTIME_BASE FROM ${BUILD_BASE} AS build WORKDIR /app/modules/nodered COPY . . -RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts \ +RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts mqtt/probe.ts mqtt/connection.ts mqtt/index.ts \ --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist FROM ${RUNTIME_BASE} @@ -22,3 +22,6 @@ COPY --from=build /app/modules/nodered/dist /app/modules/nodered/dist # provider's provisioner runs its reconcile loop in the same process, with the broker connected — # the convention novox/hq issues 060/061 settled. ENV MESH_TOOL_MODULES=/app/modules/nodered/dist/tools/index.js +# NOT dist/mqtt/index.js: that is a step the host runs to completion, named by the `mqtt` +# container's args as `mesh-tools run …` (novox/hq ADR 0052). Listed here it would run inside the +# serving sidecar too, and exit it. diff --git a/modules/nodered/client.ts b/modules/nodered/client.ts index 0cb05fc..552d814 100644 --- a/modules/nodered/client.ts +++ b/modules/nodered/client.ts @@ -82,6 +82,35 @@ export class NodeRedClient { return modules.map((m) => ({ name: m.name, version: m.version, types: m.types ?? [] })); } + /** The whole flow configuration with its revision (API v2), for a deploy that must not clobber + * a change made meanwhile. */ + async flowsWithRev(): Promise<{ rev: string; flows: any[] }> { + const body = await this.req("/flows", { headers: this.headers({ "Node-RED-API-Version": "v2" }) }); + return { rev: String(body?.rev ?? ""), flows: Array.isArray(body?.flows) ? body.flows : [] }; + } + + /** A node's stored credentials as Node-RED shows them: plain fields, and `has_` for secret ones. */ + async credentials(type: string, id: string): Promise<{ user?: string; has_password?: boolean }> { + return (await this.req(`/credentials/${encodeURIComponent(type)}/${encodeURIComponent(id)}`, { headers: this.headers() })) ?? {}; + } + + /** + * Deploy the flow configuration read at `rev`. Node-RED answers 409 when the flows changed since, + * rather than overwriting what someone deployed in between. A node carrying `credentials` has them + * stored (encrypted) and counts as changed, so a "nodes" deploy restarts it and nothing else. + */ + async deployFlowsAt(rev: string, config: any[], type = "nodes"): Promise { + await this.req("/flows", { + method: "POST", + headers: this.headers({ + "Content-Type": "application/json", + "Node-RED-API-Version": "v2", + "Node-RED-Deployment-Type": type, + }), + body: JSON.stringify({ rev, flows: config }), + }); + } + /** * Replace the whole flow configuration and deploy. Returns the new revision. `type` maps to * Node-RED's deployment types — "full" (default), "nodes", or "flows". diff --git a/modules/nodered/module.json b/modules/nodered/module.json index a0e8c27..d5c8ec7 100644 --- a/modules/nodered/module.json +++ b/modules/nodered/module.json @@ -40,13 +40,18 @@ "mode": "0700", "owner": "1000:1000" }, + { + "id": "written", + "type": "directory", + "mode": "0700" + }, { "id": "settings-code", "type": "file", "path": "${dir:state}/settings.js", "mode": "0600", "owner": "1000:1000", - "content": "// Node-RED's settings, written by the mesh from the nodered module. What an assignment may change\n// is settings.json beside this file (merged key by key); the credentials are the mesh's secrets and\n// reach Node-RED only through this file. The flows' own credentials stay encrypted in the user\n// directory under the key Node-RED keeps there (.config.runtime.json), which is data, not this.\nconst fs = require(\"fs\");\nconst path = require(\"path\");\nconst crypto = require(\"crypto\");\n\nconst ADMIN_PASSWORD = \"${secret:admin}\";\nconst API_TOKEN = \"${secret:api-token}\";\nconst ADMIN = { username: \"admin\", permissions: \"*\" };\n\nconst settings = JSON.parse(fs.readFileSync(path.join(__dirname, \"settings.json\"), \"utf8\"));\n// The mesh's keys, not Node-RED's: endpoints lands in every merged file; timeZone is the\n// assignment's way to set the zone flows schedule and format in.\nif (settings.timeZone) process.env.TZ = settings.timeZone;\ndelete settings.timeZone;\ndelete settings.endpoints;\n\nfunction same(a, b) {\n const x = crypto.createHash(\"sha256\").update(String(a)).digest();\n const y = crypto.createHash(\"sha256\").update(String(b)).digest();\n return crypto.timingSafeEqual(x, y);\n}\n\n// The admin secret is a password, or, accepted from an existing install, the bcrypt hash its\n// settings held, so the password people already use keeps working.\nfunction passwordMatches(given) {\n if (/^\\$2[aby]\\$\\d\\d\\$/.test(ADMIN_PASSWORD)) return require(\"bcryptjs\").compare(String(given), ADMIN_PASSWORD);\n return Promise.resolve(same(given, ADMIN_PASSWORD));\n}\n\nmodule.exports = Object.assign(settings, {\n uiPort: 1880,\n adminAuth: {\n type: \"credentials\",\n users: (username) => Promise.resolve(username === ADMIN.username ? ADMIN : null),\n authenticate: (username, password) =>\n username === ADMIN.username\n ? passwordMatches(password).then((ok) => (ok ? ADMIN : null))\n : Promise.resolve(null),\n // The module's own tools call the admin API with this bearer token.\n tokens: (token) => Promise.resolve(same(token, API_TOKEN) ? { username: \"mesh\", permissions: \"*\" } : null),\n },\n});\n" + "content": "// Node-RED's settings, written by the mesh from the nodered module. What an assignment may change\n// is settings.json beside this file (merged key by key); the credentials are the mesh's secrets and\n// reach Node-RED only through this file. The flows' own credentials stay encrypted in the user\n// directory under the key Node-RED keeps there (.config.runtime.json), which is data, not this.\nconst fs = require(\"fs\");\nconst path = require(\"path\");\nconst crypto = require(\"crypto\");\n\nconst ADMIN_PASSWORD = \"${secret:admin}\";\nconst API_TOKEN = \"${secret:api-token}\";\nconst ADMIN = { username: \"admin\", permissions: \"*\" };\n\nconst settings = JSON.parse(fs.readFileSync(path.join(__dirname, \"settings.json\"), \"utf8\"));\n// The mesh's keys, not Node-RED's: endpoints lands in every merged file; timeZone is the\n// assignment's way to set the zone flows schedule and format in; mqtt names the broker nodes the\n// module's MQTT step keeps pointed at the mesh's broker; topics is what nodered asks the broker for.\nif (settings.timeZone) process.env.TZ = settings.timeZone;\ndelete settings.timeZone;\ndelete settings.endpoints;\ndelete settings.mqtt;\ndelete settings.topics;\n\nfunction same(a, b) {\n const x = crypto.createHash(\"sha256\").update(String(a)).digest();\n const y = crypto.createHash(\"sha256\").update(String(b)).digest();\n return crypto.timingSafeEqual(x, y);\n}\n\n// The admin secret is a password, or, accepted from an existing install, the bcrypt hash its\n// settings held, so the password people already use keeps working.\nfunction passwordMatches(given) {\n if (/^\\$2[aby]\\$\\d\\d\\$/.test(ADMIN_PASSWORD)) return require(\"bcryptjs\").compare(String(given), ADMIN_PASSWORD);\n return Promise.resolve(same(given, ADMIN_PASSWORD));\n}\n\nmodule.exports = Object.assign(settings, {\n uiPort: 1880,\n adminAuth: {\n type: \"credentials\",\n users: (username) => Promise.resolve(username === ADMIN.username ? ADMIN : null),\n authenticate: (username, password) =>\n username === ADMIN.username\n ? passwordMatches(password).then((ok) => (ok ? ADMIN : null))\n : Promise.resolve(null),\n // The module's own tools call the admin API with this bearer token.\n tokens: (token) => Promise.resolve(same(token, API_TOKEN) ? { username: \"mesh\", permissions: \"*\" } : null),\n },\n});\n" }, { "id": "settings", @@ -107,19 +112,59 @@ "runtime-config" ], "artifact": "runtime" + }, + { + "id": "mqtt", + "type": "container", + "name": "mesh-nodered-mqtt", + "network": "host", + "run-once": true, + "volumes": [ + "/var/lib/mesh/nodered/config.json:/run/config/config.json:ro", + "${dir:written}:/var/lib/nodered-provisions", + "${dir:state}/mqtt-topic.json:/run/provisions/mqtt-topic.json:ro", + "${dir:state}/mqtt-topic.secret:/run/provisions/mqtt-topic.secret:ro", + "${dir:state}/settings.json:/run/provisions/settings.json:ro" + ], + "env": { + "MESH_NODERED_URL": "http://127.0.0.1:${port:1880}", + "MESH_NODERED_CONFIG_FILE": "/run/config/config.json", + "MESH_PROVISIONS_DIR": "/run/provisions", + "MESH_WRITTEN_DIR": "/var/lib/nodered-provisions" + }, + "args": [ + "run", + "/app/modules/nodered/dist/mqtt/index.js" + ], + "restart-on": [ + "bound-mqtt-topic", + "secret-mqtt-topic", + "settings" + ], + "artifact": "runtime" } ], "requires": [ + "mqtt-topic", "route" ], "contributes": { + "mqtt-topic": { + "topics": [ + "#" + ] + }, "route": { "label": "nodered", "endpoint": "web" } }, "binds": { - "route": "${dir:state}/route.json" + "route": "${dir:state}/route.json", + "mqtt-topic": "${dir:state}/mqtt-topic.json" + }, + "secrets": { + "mqtt-topic": "${dir:state}/mqtt-topic.secret" }, "build": { "on": [ diff --git a/modules/nodered/mqtt/connection.ts b/modules/nodered/mqtt/connection.ts new file mode 100644 index 0000000..c720441 --- /dev/null +++ b/modules/nodered/mqtt/connection.ts @@ -0,0 +1,232 @@ +// 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; +} diff --git a/modules/nodered/mqtt/index.ts b/modules/nodered/mqtt/index.ts new file mode 100644 index 0000000..df257f4 --- /dev/null +++ b/modules/nodered/mqtt/index.ts @@ -0,0 +1,102 @@ +// nodered's MQTT step — run once by the host after Node-RED starts, and again whenever the +// `mqtt-topic` binding, its pair credential or the settings change (the container's `restart-on`, +// novox/hq ADR 0099). It points the mesh's broker config nodes at the broker the mesh bound, through +// Node-RED's admin API (connection.ts). It connects to no mesh broker. +// +// Exits non-zero when anything could not be put right, so the node reports the step failed and the +// host runs it again on the next apply. Declared last in the manifest, so its failing gates nothing +// else of nodered's (novox/hq ADR 0136). Never prints a password. + +import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import { join } from "node:path"; + +import { NodeRedClient } from "../client.js"; +import { namedBrokers, reconcileBrokers, type Binding, type Marks } from "./connection.js"; +import { probeBroker } from "./probe.js"; + +const dir = process.env.MESH_PROVISIONS_DIR ?? "/run/provisions"; +const writtenDir = process.env.MESH_WRITTEN_DIR ?? "/var/lib/nodered-provisions"; +const waitSeconds = Number(process.env.MESH_NODERED_WAIT_SECONDS ?? "180"); + +const readIfThere = (path: string): Promise => readFile(path, "utf8").catch(() => undefined); +const parse = (raw: string | undefined): T | undefined => { + if (raw === undefined) return undefined; + try { + return JSON.parse(raw) as T; + } catch { + return undefined; + } +}; + +const marks: Marks = { + async get(name) { + return (await readIfThere(join(writtenDir, `${name}.digest`)))?.trim() || undefined; + }, + async set(name, value) { + await mkdir(writtenDir, { recursive: true, mode: 0o700 }); + const path = join(writtenDir, `${name}.digest`); + await writeFile(`${path}.tmp`, `${value}\n`, { mode: 0o600 }); + await rename(`${path}.tmp`, path); + }, +}; + +let client: NodeRedClient; +try { + client = NodeRedClient.fromEnv(); +} catch (err) { + console.error(`[nodered-mqtt] ${err instanceof Error ? err.message : String(err)}`); + process.exit(1); +} + +/** Node-RED answers the admin API once its flows are loaded and the token is good. */ +async function ready(): Promise { + const until = Date.now() + waitSeconds * 1000; + for (;;) { + try { + await client.flowsWithRev(); + return true; + } catch (err) { + if (/\b(401|403)\b/.test(String(err))) { + console.error("[nodered-mqtt] Node-RED refuses the api-token — settings.js and this step disagree"); + return false; + } + } + if (Date.now() >= until) return false; + await new Promise((r) => setTimeout(r, 2000)); + } +} + +if (!(await ready())) { + console.error(`[nodered-mqtt] Node-RED's admin API did not answer at ${client.baseUrl} within ${waitSeconds}s`); + process.exit(1); +} + +const binding = parse(await readIfThere(join(dir, "mqtt-topic.json"))); +const secret = await readIfThere(join(dir, "mqtt-topic.secret")); +const settings = parse(await readIfThere(join(dir, "settings.json"))); + +const outcomes = await reconcileBrokers( + { + nodered: { + flows: () => client.flowsWithRev(), + credentials: (type, id) => client.credentials(type, id), + deploy: (rev, flows) => client.deployFlowsAt(rev, flows, "nodes"), + }, + probe: probeBroker, + marks, + }, + binding, + secret, + namedBrokers(settings), +); + +let failed = 0; +for (const o of outcomes) { + if (o.result === "unchanged") console.log(`[nodered-mqtt] ${o.what}: already as the mesh says${o.note ? ` — ${o.note}` : ""}`); + else if (o.result === "written") console.log(`[nodered-mqtt] ${o.what}: wrote ${o.fields.join(", ")}${o.note ? ` — ${o.note}` : ""}`); + else { + failed++; + console.error(`[nodered-mqtt] ${o.what}: ${o.problem}`); + } +} +process.exitCode = failed > 0 ? 1 : 0; diff --git a/modules/nodered/mqtt/probe.ts b/modules/nodered/mqtt/probe.ts new file mode 100644 index 0000000..1c1b43c --- /dev/null +++ b/modules/nodered/mqtt/probe.ts @@ -0,0 +1,117 @@ +// Ask the broker, before Node-RED is told anything, whether it takes the login and password the +// mesh delivered — and whether that login may subscribe to every topic, as flows expect. +// +// One MQTT 3.1.1 session: CONNECT (clean, a throwaway client id, so no flow's session is taken +// over), read the CONNACK, optionally SUBSCRIBE once and read the SUBACK, DISCONNECT. No dependency: +// the handful of bytes MQTT needs for this are written here. + +import { randomBytes } from "node:crypto"; +import { connect } from "node:net"; + +export interface ProbeResult { + /** 0 accepted; 4 bad username or password; 5 not authorised. */ + connack: number; + /** The SUBACK return code for the filter asked about: 0–2 granted, 0x80 refused. */ + suback?: number; +} + +export type Probe = (host: string, port: number, username: string, password: string, subscribe?: string) => Promise; + +function str(v: string): Buffer { + const b = Buffer.from(v, "utf8"); + const len = Buffer.alloc(2); + len.writeUInt16BE(b.length); + return Buffer.concat([len, b]); +} + +function packet(type: number, body: Buffer): Buffer { + let remaining = body.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); + return Buffer.concat([Buffer.from([type, ...lenBytes]), body]); +} + +/** The first complete packet in `buf`: its type byte, its body, and how many bytes it took. */ +export function firstPacket(buf: Buffer): { type: number; body: Buffer; used: number } | undefined { + if (buf.length < 2) return undefined; + let length = 0; + let multiplier = 1; + let i = 1; + for (;;) { + if (i >= buf.length) return undefined; + const byte = buf[i++]; + length += (byte & 0x7f) * multiplier; + if ((byte & 0x80) === 0) break; + multiplier *= 128; + if (i > 4) throw new Error("malformed MQTT remaining length"); + } + if (buf.length < i + length) return undefined; + return { type: buf[0], body: buf.subarray(i, i + length), used: i + length }; +} + +export const probeBroker: Probe = (host, port, username, password, subscribe) => { + const connectBody = Buffer.concat([ + str("MQTT"), + Buffer.from([4, 0xc2, 0, 10]), // level 4 (3.1.1); username + password + clean session; keepalive 10s + str(`mesh-probe-${randomBytes(6).toString("hex")}`), + str(username), + str(password), + ]); + return new Promise((resolve, reject) => { + const socket = connect({ host, port }); + let buf = Buffer.alloc(0); + const result: ProbeResult = { connack: -1 }; + const timer = setTimeout(() => { + socket.destroy(); + reject(new Error(`no answer from the broker at ${host}:${port} within 10s`)); + }, 10_000); + const finish = (): void => { + clearTimeout(timer); + if (result.connack === 0) socket.end(Buffer.from([0xe0, 0])); + else socket.destroy(); + resolve(result); + }; + socket.on("connect", () => socket.write(packet(0x10, connectBody))); + socket.on("data", (chunk) => { + buf = Buffer.concat([buf, chunk]); + for (;;) { + let p; + try { + p = firstPacket(buf); + } catch (err) { + clearTimeout(timer); + socket.destroy(); + reject(err); + return; + } + if (!p) return; + buf = buf.subarray(p.used); + const kind = p.type >> 4; + if (kind === 2) { + result.connack = p.body[1] ?? -1; + if (result.connack !== 0 || !subscribe) return finish(); + // SUBSCRIBE, packet id 1, one filter at QoS 0. + socket.write(packet(0x82, Buffer.concat([Buffer.from([0, 1]), str(subscribe), Buffer.from([0])]))); + } else if (kind === 9) { + result.suback = p.body[2]; + return finish(); + } + } + }); + socket.on("error", (err) => { + clearTimeout(timer); + reject(err); + }); + socket.on("close", () => { + if (result.connack === -1) { + clearTimeout(timer); + reject(new Error(`the broker at ${host}:${port} closed the connection without answering`)); + } + }); + }); +}; diff --git a/modules/nodered/package.json b/modules/nodered/package.json index b8459ea..29df536 100644 --- a/modules/nodered/package.json +++ b/modules/nodered/package.json @@ -1,9 +1,14 @@ { "name": "@novox/module-nodered", "version": "0.1.0", - "description": "nodered — flow-based automation. Its client and tools live here (novox/hq ADR 0039).", + "description": "nodered \u2014 flow-based automation. Its client and tools live here (novox/hq ADR 0039).", "type": "module", "private": true, + "scripts": { + "build": "tsc client.ts tools/index.ts mqtt/probe.ts mqtt/connection.ts mqtt/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist", + "typecheck": "tsc -p tsconfig.json", + "test": "node --test --experimental-strip-types 'test/*.test.ts'" + }, "dependencies": { "@novox/mesh-sdk": "^0.1.0" }, diff --git a/modules/nodered/test/mqtt.test.ts b/modules/nodered/test/mqtt.test.ts new file mode 100644 index 0000000..59a1f6a --- /dev/null +++ b/modules/nodered/test/mqtt.test.ts @@ -0,0 +1,124 @@ +// What holds nodered's MQTT step (mqtt/connection.ts): the broker nodes the mesh owns — the one it +// makes, or the ones settings name — are made to use the broker and login the mesh bound, only after +// the broker takes that login; every other field of a node is kept; nothing is deployed when nothing +// differs; a node that is not named is never touched; a loopback broker address is refused. +// +// Node-RED and the broker are fakes answering as the real ones do (admin API v2 of nodered/node-red +// 5.0.7, the build ace runs). + +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { MESH_BROKER_ID, namedBrokers, reconcileBrokers, type Binding, type FlowNode, type Marks, type NodeRed } from "../mqtt/connection.ts"; +import type { Probe } from "../mqtt/probe.ts"; + +const MINTED = "mesh-minted-password"; +const binding = (at = "ace.internal"): Binding => ({ provision: "mqtt-topic", from: "ace", at, as: "mesh_ace_nodered", serves: { scheme: "mqtt", port: 1883 } }); + +/** ace's flows, reduced: its one broker node (dead, zurag.be:1884) and a node that uses it. */ +function aceFlows(): FlowNode[] { + return [ + { id: "2b0aece9c5f3b307", type: "mqtt-broker", name: "MQTT Broker", broker: "zurag.be", port: "1884", clientid: "", usetls: false, protocolVersion: "4", keepalive: "60" }, + { id: "fe0cae96f1e3ae4d", type: "mqtt in", topic: "stat/sonoff_office_light_switch/RESULT", broker: "2b0aece9c5f3b307", z: "t" }, + { id: "other-broker", type: "mqtt-broker", name: "someone else's", broker: "test.mosquitto.org", port: "1883" }, + ]; +} + +function fakeNodeRed(flows: FlowNode[], creds: Record = {}) { + let rev = "r1"; + const deploys: FlowNode[][] = []; + const nodered: NodeRed = { + async flows() { + return { rev, flows: structuredClone(flows) }; + }, + async credentials(_type, id) { + const c = creds[id] ?? {}; + return { user: c.user, has_password: Boolean(c.password) }; + }, + async deploy(at, next) { + if (at !== rev) throw new Error("Node-RED /flows: 409 version_mismatch"); + for (const n of next) { + if (n.credentials) creds[n.id] = { ...(creds[n.id] ?? {}), ...(n.credentials as object) }; + } + flows.splice(0, flows.length, ...next.map(({ credentials: _c, ...n }) => n as FlowNode)); + deploys.push(next); + rev = `r${deploys.length + 1}`; + }, + }; + return { nodered, deploys, flows, creds }; +} + +const marks = (): Marks & { store: Map } => { + const store = new Map(); + return { store, get: async (k) => store.get(k), set: async (k, v) => void store.set(k, v) }; +}; +const takes: Probe = async (_h, _p, user, pass) => ({ connack: user === "mesh_ace_nodered" && pass === MINTED ? 0 : 5, suback: 0 }); + +test("ace: the named broker node is moved to the bound broker and login; the other broker is not touched", async () => { + const f = fakeNodeRed(aceFlows(), { "2b0aece9c5f3b307": { user: "luffy", password: "old" }, "other-broker": { user: "x", password: "y" } }); + const m = marks(); + const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: m }, binding(), `${MINTED}\n`, ["2b0aece9c5f3b307"]); + assert.deepEqual(out, [{ what: "broker 2b0aece9c5f3b307", result: "written", fields: ["broker", "port", "user", "password"] }]); + const node = f.flows.find((n) => n.id === "2b0aece9c5f3b307"); + assert.deepEqual(node, { ...aceFlows()[0], broker: "ace.internal", port: "1883", usetls: false }); + assert.deepEqual(f.creds["2b0aece9c5f3b307"], { user: "mesh_ace_nodered", password: MINTED }); + assert.deepEqual(f.flows.find((n) => n.id === "other-broker"), aceFlows()[2]); + assert.deepEqual(f.creds["other-broker"], { user: "x", password: "y" }); + // Only the changed node carried credentials in the deploy. + assert.deepEqual(f.deploys[0].filter((n) => n.credentials).map((n) => n.id), ["2b0aece9c5f3b307"]); + + // Again: nothing differs, nothing is deployed. + const again = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: m }, binding(), MINTED, ["2b0aece9c5f3b307"]); + assert.deepEqual(again, [{ what: "broker 2b0aece9c5f3b307", result: "unchanged" }]); + assert.equal(f.deploys.length, 1); +}); + +test("fresh: with nothing named, the step makes its own broker node", async () => { + const f = fakeNodeRed([]); + const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding(), MINTED, []); + assert.deepEqual(out, [{ what: `broker ${MESH_BROKER_ID}`, result: "written", fields: ["node"] }]); + assert.equal(f.flows[0].type, "mqtt-broker"); + assert.equal(f.flows[0].broker, "ace.internal"); + assert.deepEqual(f.creds[MESH_BROKER_ID], { user: "mesh_ace_nodered", password: MINTED }); +}); + +test("a login the broker does not take is never written", async () => { + const f = fakeNodeRed(aceFlows()); + const out = await reconcileBrokers({ nodered: f.nodered, probe: async () => ({ connack: 5 }), marks: marks() }, binding(), MINTED, ["2b0aece9c5f3b307"]); + assert.equal(out[0].result, "refused"); + assert.equal(f.deploys.length, 0); +}); + +test("a named node that is not there, or a loopback broker, is refused", async () => { + const f = fakeNodeRed(aceFlows()); + const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding(), MINTED, ["gone"]); + assert.match((out[0] as { problem: string }).problem, /no such node/); + const lo = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding("127.0.0.1"), MINTED, []); + assert.match((lo[0] as { problem: string }).problem, /Node-RED itself/); + assert.equal(f.deploys.length, 0); +}); + +test("a deploy that raced another is read again once", async () => { + const f = fakeNodeRed(aceFlows()); + let first = true; + const racing: NodeRed = { + ...f.nodered, + async flows() { + const got = await f.nodered.flows(); + if (first) { + first = false; + return { ...got, rev: "stale" }; + } + return got; + }, + }; + const out = await reconcileBrokers({ nodered: racing, probe: takes, marks: marks() }, binding(), MINTED, ["2b0aece9c5f3b307"]); + assert.equal(out[0].result, "written"); + assert.equal(f.deploys.length, 1); +}); + +test("settings name broker nodes under mqtt.brokers", () => { + assert.deepEqual(namedBrokers({ mqtt: { brokers: ["a", "", 3, "b"] } }), ["a", "b"]); + assert.deepEqual(namedBrokers({ endpoints: {} }), []); + assert.deepEqual(namedBrokers(undefined), []); +}); diff --git a/modules/nodered/tsconfig.json b/modules/nodered/tsconfig.json index 426d382..c36997c 100644 --- a/modules/nodered/tsconfig.json +++ b/modules/nodered/tsconfig.json @@ -8,5 +8,11 @@ "skipLibCheck": true, "noEmit": true }, - "include": ["client.ts", "tools/index.ts"] + "include": [ + "client.ts", + "tools/index.ts", + "mqtt/probe.ts", + "mqtt/connection.ts", + "mqtt/index.ts" + ] }