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.
118 lines
4.1 KiB
TypeScript
118 lines
4.1 KiB
TypeScript
// 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`));
|
||
}
|
||
});
|
||
});
|
||
};
|