// 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`)); } }); }); };