"use strict"; var __importDefault = (this && this.__importDefault) || function (mod) { return (mod && mod.__esModule) ? mod : { "default": mod }; }; Object.defineProperty(exports, "__esModule", { value: true }); exports.DEFAULT_HOST = void 0; // mqtt reaches Node as CommonJS: the default import works from both builds, a named one needs Node to detect it. const mqtt_1 = __importDefault(require("mqtt")); const Device_js_1 = __importDefault(require("./device/Device.js")); const validate_js_1 = require("./device/validate.js"); const events_1 = require("events"); const enums_js_1 = require("./compat/enums.js"); const types_js_1 = require("./registry/types.js"); exports.DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; // What addDevice() copies from the definition to the device const DESCRIBED = [ "description", "manufacturerName", "manufacturer", "model", "serialNumber", "firmwareVersion", "softwareVersion", "customIdentifier", "cookie", ]; /** * The Alexa-to-MQTT bridge: one broker connection for a user's root topic, a set of registered devices, discovery and * directive dispatch. * * Events (all optional to listen to - since 1.5.1 a broker outage never throws out of the library): * "connect" connected (or reconnected) and subscribed to /# * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed * "error" (err) a connection or publish error. Emitted ONLY when a listener is attached (1.4.0 emitted it * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } */ class Alex2MQTT extends events_1.EventEmitter { constructor(username, password, rootTopic, debugLogging = false, options = {}) { super(); // Initialize EventEmitter this.username = username; this.password = password; this.rootTopic = rootTopic; this.debugLogging = debugLogging; this.client = null; this.devices = []; // The lines of check() that were logged. Discovery comes every few minutes: a line is logged once. this.logged = new Set(); /** true while the broker connection is up. */ this.connected = false; /** ISO time of the last discovery request answered, null before the first. */ this.lastDiscoveryAt = null; this.options = options || {}; this.MqttHost = this.options.host || exports.DEFAULT_HOST; } log(message, detail) { if (this.options.log) this.options.log(message, detail); else if (this.debugLogging) { if (detail === undefined) console.log(`[Alex2Node.ts] ${message}`); else console.log(`[Alex2Node.ts] ${message}`, detail); } } /** Emit "error" only when somebody listens: an unhandled "error" event would crash the host process. */ fail(err) { this.log("error: " + err.message); if (this.listenerCount("error") > 0) this.emit("error", err); } connect() { if (this.client) return; const options = { username: this.username, password: this.password, reconnectPeriod: 1000, connectTimeout: 10000, ...(this.options.mqtt || {}), }; this.client = mqtt_1.default.connect(this.MqttHost, options); for (const d of this.devices) d.setMqttClient(this.client); // after a disconnect(): the devices publish through the new client this.client.on("connect", () => { this.connected = true; this.log("Connected to MQTT broker"); this.client.subscribe(this.rootTopic + "/#", (err) => { if (err) { this.fail(err); return; } this.log("Subscribed to topics"); this.emit("connect"); }); }); this.client.on("offline", () => { this.connected = false; this.log("offline"); this.emit("offline"); }); this.client.on("reconnect", () => { this.log("reconnecting"); this.emit("reconnect"); }); this.client.on("close", () => { this.connected = false; this.emit("close"); }); this.client.on("error", (err) => { this.connected = false; this.fail(err); }); this.client.on("message", (topic, message) => { this.log(`MQTT Message Received`, { topic, payload: message.toString() }); if (topic == `${this.rootTopic}/discover`) { this.log("Discovery request received, getting device json..."); const deviceArray = this.describeDevices(); this.lastDiscoveryAt = new Date().toISOString(); this.client.publish(topic + "_r", JSON.stringify(deviceArray), (err) => { if (err) { this.fail(err); return; } this.log(`Discovery payloads published to ${topic + "_r"}`, deviceArray); this.emit("discover", deviceArray.length); }); } else if (topic.split("/").length == 3) { const [, endpointId, directiveType] = topic.split("/"); if (directiveType != "alexaDirective") { return; } const device = this.devices.find((d) => d.endpointId === endpointId); if (!device) { this.log(`No device found for endpointId: ${endpointId}`); return; } let payload; try { payload = JSON.parse(message.toString()); } catch (e) { this.log(`Failed to parse payload JSON`, e); return; } if (!payload || !payload.header) return; this.emit("directive", { endpointId, namespace: payload.header.namespace, name: payload.header.name }); if (payload.header.namespace === "Alexa" && payload.header.name === "ReportState") { this.log(`ReportState directive for device ${endpointId}`); device.emit("ReportState", payload); } else { device.emit("Event", payload, payload.header.namespace); } } }); } // The endpoint objects of a discovery answer. What check() says about a device is logged, each line once. A device // that cannot be described is left out and reported as an error: the others are still announced. describeDevices() { const endpoints = []; for (const device of this.devices) { try { endpoints.push(device.getJSON()); for (const line of device.check()) { if (this.logged.has(line)) continue; this.logged.add(line); this.log(`warning: ${line}`); } } catch (err) { this.fail(new Error(`${device.endpointId} is not in the discovery answer: ${err instanceof Error ? err.message : err}`)); } } return endpoints; } /** Close the broker connection (resolves once closed). The devices stay registered; connect() again reuses them. */ disconnect() { return new Promise((resolve) => { const c = this.client; if (!c) return resolve(); this.client = null; this.connected = false; c.end(true, {}, () => resolve()); }); } /** * Declare a device: * * const blinds = bridge.addDevice({ endpointId: "bedroom-blinds", name: "Bedroom Blinds", * categories: ["INTERIOR_BLIND"], manufacturerName: "Acme", description: "Roller blind by Acme" }); * * Throws a DeclarationError when Alexa would reject the endpoint (an endpointId with a slash, a name with * punctuation) or when the endpointId is registered already. Discovery lists Alexa.EndpointHealth for the device * unless endpointHealth is false. */ addDevice(definition) { if (!this.client) { throw new Error("Must call connect before creating devices"); } const { endpointId, name, categories, endpointHealth, ...described } = definition; // Without categories the device would take LIGHT, the default of registerDevice: here it is a mistake const device = new Device_js_1.default(this.client, this.rootTopic, name, endpointId, (categories ?? [])); for (const [field, value] of Object.entries(described)) { if (!DESCRIBED.includes(field)) throw new types_js_1.DeclarationError({ endpointId }, `${field} is not a field of an endpoint`); if (value !== undefined) Object.assign(device, { [field]: value }); } (0, validate_js_1.checkEndpoint)(device.getJSON()); const existing = this.getDevice(endpointId); if (existing) throw new types_js_1.DeclarationError({ endpointId }, `is registered already, as "${existing.name}"`); this.register(device, endpointHealth !== false); return device; } registerDevice(name, endpointId, displayCategory) { if (!this.client) { throw new Error("Must call connect before creating devices"); } const existing = this.devices.find((d) => d.endpointId === endpointId); if (existing) { // 1.5.2: endpointIds are unique per root topic; 1.5.1 added a second device that never got a directive const warning = `warning: registerDevice("${endpointId}") is already registered as "${existing.name}", returning that device`; if (this.options.log) this.options.log(warning); else console.warn(`[Alex2Node.ts] ${warning}`); // visible without a log hook too return existing; } this.log(`Creating new device with endpoint: ${endpointId}`); const normalizedCategory = Array.isArray(displayCategory) ? displayCategory : [displayCategory || enums_js_1.DisplayCategory.LIGHT]; const device = new Device_js_1.default(this.client, this.rootTopic, name, endpointId, normalizedCategory); this.register(device, false); return device; } register(device, endpointHealth) { device.alexaInterface = this.options.alexaInterface !== false; device.endpointHealth = endpointHealth; device.onPublishError = (err) => this.fail(err); // a failed publish is an "error" event (when listened to), never a rejected send() this.devices.push(device); } /** Forget a device (its listeners with it). Returns false when there was none. */ unregisterDevice(endpointId) { const i = this.devices.findIndex((d) => d.endpointId === endpointId); if (i === -1) return false; this.devices[i].removeAllListeners(); this.devices.splice(i, 1); return true; } /** Forget every device. */ clearDevices() { for (const d of this.devices) d.removeAllListeners(); this.devices = []; } getDevices() { return this.devices.slice(); } getDevice(endpointId) { return this.devices.find((d) => d.endpointId === endpointId); } getRootTopic() { return this.rootTopic; } getHost() { return this.MqttHost; } } exports.default = Alex2MQTT;