// 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;