"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; const mqtt_1 = __importDefault(require("mqtt")); const Device_1 = __importDefault(require("./Device")); const events_1 = require("events"); const DisplayCategory_1 = require("./DisplayCategory"); exports.DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; /** * 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 = []; /** 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 = Object.assign({ username: this.username, password: this.password, reconnectPeriod: 1000, connectTimeout: 10000 }, (this.options.mqtt || {})); this.client = mqtt_1.default.connect(this.MqttHost, options); 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.devices.map((device) => device.getJSON()); 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); } } }); } /** 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()); }); } registerDevice(name, endpointId, displayCategory) { if (!this.client) { throw new Error("Must call connect before creating devices"); } this.log(`Creating new device with endpoint: ${endpointId}`); const normalizedCategory = displayCategory === null ? [DisplayCategory_1.DisplayCategory.LIGHT] : Array.isArray(displayCategory) ? displayCategory : [displayCategory || DisplayCategory_1.DisplayCategory.LIGHT]; const device = new Device_1.default(this.client, this.rootTopic, name, endpointId, normalizedCategory); this.devices.push(device); return 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;