Compare commits

..
Author SHA1 Message Date
jschoubben 204bfbbaf9 dnsmasq: answer a container's query, which arrives on the runtime's bridge
dnsmasq admits a query by the interface it arrives on when told interface=,
and by the address it is sent to when told listen-address=. A container's
query is sent to the machine's private address but arrives on docker0, so
interface=mesh0 dropped it silently on every machine (novox/hq issue 110).
Name the address, not the interface.
2026-09-30 14:31:26 +02:00
jschoubben e0c09f46b6 Merge pull request 'dnsmasq: the runtime is reloaded when its dns file changes, and may then be restarted safely' (#175) from fix/110-the-runtime-reads-its-dns into main 2026-09-30 12:22:10 +00:00
jschoubben e0faf012be dnsmasq: the runtime is reloaded when its dns file changes, and may then be restarted safely
novox/hq 04-ISSUES/110. The module writes the runtime's `dns` key and
deliberately ordered no restart, because a restart stops every container.
It also ordered no reload, and the key holds only for containers created
after the runtime next starts. On two of four machines the runtime
predated the file — one since August — so every container there was
handed a public resolver and no mesh name resolved, while everything read
as fine. The third machine works only because its runtime happened to
restart later.

The file now also sets live-restore, which the runtime reads on a reload,
and the module declares the runtime reloaded when the file changes (ADR
0102: reload, don't restart). A reload still does not make `dns` take
effect; what it does is make the one restart that key needs keep every
container running. That restart stays the operator's, once per machine,
and is harmless from the second time on.

The mesh restarts nothing here. Undeclared, the runtime's unit goes back
to the state it was found in (ADR 0118).
2026-09-30 14:22:03 +02:00
mesh-admin 1b19c79d63 Merge pull request 'postgres: the provisioner dials the port the mesh gave the store' (#174) from fix/postgres-provisioner-dials-the-port-it-was-given into main 2026-09-30 12:06:25 +00:00
jschoubben 4ec2ae1f7f postgres: the provisioner dials the port the mesh gave the store
MESH_PROVISION_POSTGRES named 127.0.0.1:5432 literally; the seat twin
(${seat:mesh-store:5432}) corrected it only on the machine holding the
mesh-store seat. On any other machine — ace, where the module provides
postgres-database without the seat — the twin is empty and the sidecar
would have dialled whatever else holds 5432 (HAL's postgres). ${port:5432}
is the machine port on every node; on novox the twin still says the same
6852, so nothing moves there.
2026-09-30 14:01:53 +02:00
11 changed files with 27 additions and 720 deletions
File diff suppressed because one or more lines are too long
+1 -4
View File
@@ -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 mqtt/probe.ts mqtt/connection.ts mqtt/index.ts \
RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
@@ -22,6 +22,3 @@ 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.
+3 -38
View File
@@ -3,8 +3,7 @@
//
// Node-RED exposes a runtime admin API under its base URL: GET/POST /flows for the whole flow
// configuration, GET /nodes for installed node modules. A default install has no auth; when
// adminAuth is on, a bearer token is required — the module's settings accept the mesh-minted
// api-token, which the runtime config file carries as `token`.
// adminAuth is on, a bearer token (minted at /auth/token) is required.
import { readFileSync } from "node:fs";
@@ -82,35 +81,6 @@ 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_<field>` 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<void> {
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".
@@ -118,13 +88,8 @@ export class NodeRedClient {
async deployFlows(config: any[], type = "full"): Promise<{ rev?: string; nodeCount: number }> {
const body = await this.req("/flows", {
method: "POST",
// v2 answers { rev }; v1 answers 204 with no body, which req() cannot parse.
headers: this.headers({
"Content-Type": "application/json",
"Node-RED-API-Version": "v2",
"Node-RED-Deployment-Type": type,
}),
body: JSON.stringify({ flows: config }),
headers: this.headers({ "Content-Type": "application/json", "Node-RED-Deployment-Type": type }),
body: JSON.stringify(config),
});
return { rev: body?.rev, nodeCount: config.length };
}
+7 -85
View File
@@ -5,8 +5,6 @@
"flows.deployed"
],
"own-secrets": {
"admin": "/var/lib/mesh/nodered/admin",
"api-token": "/var/lib/mesh/nodered/api-token",
"broker": "/var/lib/mesh/nodered/broker"
},
"capabilities": [
@@ -28,63 +26,26 @@
"path": "/var/lib/mesh/nodered",
"mode": "0700"
},
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{
"id": "data",
"type": "directory",
"path": "/services/nodered/data",
"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; 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",
"type": "file",
"path": "${dir:state}/settings.json",
"mode": "0600",
"owner": "1000:1000",
"merge": "json",
"content": "{\n \"flowFile\": \"flows.json\",\n \"flowFilePretty\": true,\n \"diagnostics\": { \"enabled\": true, \"ui\": true },\n \"runtimeState\": { \"enabled\": false, \"ui\": false },\n \"logging\": { \"console\": { \"level\": \"info\", \"metrics\": false, \"audit\": false } },\n \"exportGlobalContextKeys\": false,\n \"externalModules\": {},\n \"editorTheme\": { \"projects\": { \"enabled\": false } },\n \"functionExternalModules\": true,\n \"debugMaxLength\": 1000,\n \"mqttReconnectTime\": 15000,\n \"serialReconnectTime\": 15000\n}\n"
},
{
"id": "server",
"type": "container",
"name": "nodered",
"image": "nodered/node-red@sha256:a649dd711d55490151a2c39a8e48ad0c44325488fbc0e66315f2d2e19e5e1ace",
"image": "nodered/node-red@sha256:02a2b92a41b73d2bc388238b86e4fcaab7fb5466373adb24e1df6aa5845265ff",
"env": {
"TZ": "Etc/UTC"
},
"ports": [
"1880"
],
"args": [
"--settings",
"/config/settings.js"
],
"volumes": [
"${dir:data}:/data",
"${dir:state}/settings.js:/config/settings.js:ro",
"${dir:state}/settings.json:/config/settings.json:ro"
],
"restart-on": [
"settings-code",
"settings"
"/services/nodered/data:/data"
]
},
{
@@ -92,7 +53,8 @@
"type": "file",
"path": "/var/lib/mesh/nodered/config.json",
"mode": "0600",
"content": "{\n \"token\": \"${secret:api-token}\"\n}\n"
"content": "{}\n",
"merge": "json"
},
{
"id": "runtime",
@@ -105,66 +67,26 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_NODERED_URL": "http://127.0.0.1:${port:1880}",
"MESH_NODERED_URL": "http://127.0.0.1:1880",
"MESH_NODERED_CONFIG_FILE": "/run/config/config.json"
},
"restart-on": [
"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",
"mqtt-topic": "${dir:state}/mqtt-topic.json"
},
"secrets": {
"mqtt-topic": "${dir:state}/mqtt-topic.secret"
"route": "/var/lib/mesh/nodered/route.json"
},
"build": {
"on": [
-232
View File
@@ -1,232 +0,0 @@
// 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;
}
-102
View File
@@ -1,102 +0,0 @@
// 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<string | undefined> => readFile(path, "utf8").catch(() => undefined);
const parse = <T>(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<boolean> {
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<Binding>(await readIfThere(join(dir, "mqtt-topic.json")));
const secret = await readIfThere(join(dir, "mqtt-topic.secret"));
const settings = parse<unknown>(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;
-117
View File
@@ -1,117 +0,0 @@
// 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<ProbeResult>;
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`));
}
});
});
};
+1 -6
View File
@@ -1,14 +1,9 @@
{
"name": "@novox/module-nodered",
"version": "0.1.0",
"description": "nodered \u2014 flow-based automation. Its client and tools live here (novox/hq ADR 0039).",
"description": "nodered — 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"
},
-124
View File
@@ -1,124 +0,0 @@
// 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<string, { user?: string; password?: string }> = {}) {
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<string, string> } => {
const store = new Map<string, string>();
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), []);
});
+1 -7
View File
@@ -8,11 +8,5 @@
"skipLibCheck": true,
"noEmit": true
},
"include": [
"client.ts",
"tools/index.ts",
"mqtt/probe.ts",
"mqtt/connection.ts",
"mqtt/index.ts"
]
"include": ["client.ts", "tools/index.ts"]
}
+1 -1
View File
@@ -105,7 +105,7 @@
"/var/lib/postgres/superuser.secret:/run/secrets/superuser:ro"
],
"env": {
"MESH_PROVISION_POSTGRES": "postgres://postgres@127.0.0.1:5432/postgres?sslmode=disable",
"MESH_PROVISION_POSTGRES": "postgres://postgres@127.0.0.1:${port:5432}/postgres?sslmode=disable",
"MESH_PROVISION_POSTGRES_PORT": "${seat:mesh-store:5432}",
"MESH_PROVISION_PASSWORD_FILE": "/run/secrets/superuser",
"MESH_BROKER_FILE": "/run/secrets/broker",