2.0 ignored the first argument of the Device constructor in silence and had removed setMqttClient(), so a 1.x call threw a TypeError. Both are deprecated shims now, as the design lists them: setMqttClient() is back and does nothing, a client given to the constructor is ignored, and each says once in a process what to do. The line goes to the log hook of the bridge; a device on no bridge has none and emits a DeprecationWarning (ALEX2NODE_DEVICE_CLIENT, ALEX2NODE_SET_MQTT_CLIENT). null and undefined as the client warn nothing: the bridge passes null. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
320 lines
15 KiB
JavaScript
320 lines
15 KiB
JavaScript
"use strict";
|
|
var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) {
|
|
if (k2 === undefined) k2 = k;
|
|
var desc = Object.getOwnPropertyDescriptor(m, k);
|
|
if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) {
|
|
desc = { enumerable: true, get: function() { return m[k]; } };
|
|
}
|
|
Object.defineProperty(o, k2, desc);
|
|
}) : (function(o, m, k, k2) {
|
|
if (k2 === undefined) k2 = k;
|
|
o[k2] = m[k];
|
|
}));
|
|
var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) {
|
|
Object.defineProperty(o, "default", { enumerable: true, value: v });
|
|
}) : function(o, v) {
|
|
o["default"] = v;
|
|
});
|
|
var __importStar = (this && this.__importStar) || (function () {
|
|
var ownKeys = function(o) {
|
|
ownKeys = Object.getOwnPropertyNames || function (o) {
|
|
var ar = [];
|
|
for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k;
|
|
return ar;
|
|
};
|
|
return ownKeys(o);
|
|
};
|
|
return function (mod) {
|
|
if (mod && mod.__esModule) return mod;
|
|
var result = {};
|
|
if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]);
|
|
__setModuleDefault(result, mod);
|
|
return result;
|
|
};
|
|
})();
|
|
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");
|
|
const dispatcher_js_1 = require("./dispatcher.js");
|
|
const topics = __importStar(require("./topics.js"));
|
|
const transport_js_1 = require("./transport.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 messageId of a discovery request: in the header of the Discover directive, or beside the namespace and the
|
|
// name when the request is the header alone. undefined for a request without one, which is answered every time.
|
|
function requestId(text) {
|
|
let request;
|
|
try {
|
|
request = JSON.parse(text);
|
|
}
|
|
catch {
|
|
return undefined;
|
|
}
|
|
const header = request?.directive?.header ?? request?.header ?? request;
|
|
const { messageId } = (typeof header === "object" && header !== null ? header : {});
|
|
return typeof messageId === "string" ? messageId : undefined;
|
|
}
|
|
/**
|
|
* 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 <root>/discover and <root>/+/alexaDirective
|
|
* "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 }
|
|
* "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name }
|
|
* "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR:
|
|
* { endpointId, namespace, name, correlationToken }
|
|
*/
|
|
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;
|
|
this.publisher = this.options.publisher ?? new transport_js_1.MqttPublisher(() => this.client);
|
|
this.dispatcher = new dispatcher_js_1.Dispatcher({
|
|
rootTopic,
|
|
publisher: this.publisher,
|
|
getDevice: (endpointId) => this.getDevice(endpointId),
|
|
log: (message, detail) => this.log(message, detail),
|
|
fail: (err) => this.fail(err),
|
|
emit: (event, ...args) => { this.emit(event, ...args); },
|
|
answerWithinMs: this.options.answerWithinMs,
|
|
answerUnknownEndpoints: this.options.answerUnknownEndpoints,
|
|
timers: this.options.timers,
|
|
});
|
|
}
|
|
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);
|
|
this.client.on("connect", () => {
|
|
this.connected = true;
|
|
this.log("Connected to MQTT broker");
|
|
this.client.subscribe(topics.subscriptions(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) => { void this.receive(topic, message); });
|
|
}
|
|
/**
|
|
* What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to
|
|
* the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test
|
|
* calls it in place of the broker:
|
|
*
|
|
* await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive));
|
|
*/
|
|
async receive(topic, payload) {
|
|
const text = payload.toString();
|
|
this.log(`MQTT Message Received`, { topic, payload: text });
|
|
try {
|
|
if (topic === topics.discover(this.rootTopic)) {
|
|
const messageId = requestId(text);
|
|
if (messageId !== undefined && this.dispatcher.arrivedBefore("", messageId)) {
|
|
this.log(`the discovery request ${messageId} is dropped: it arrived before`);
|
|
return;
|
|
}
|
|
await this.answerDiscovery();
|
|
return;
|
|
}
|
|
// Another topic can still arrive: a broker may deliver what an earlier session of this client id subscribed to
|
|
const endpointId = topics.directiveEndpoint(this.rootTopic, topic);
|
|
if (endpointId !== null)
|
|
await this.dispatcher.dispatch(endpointId, text);
|
|
}
|
|
catch (err) {
|
|
// A listener of "discover", "directive" or "unknownEndpoint" threw
|
|
this.fail((0, transport_js_1.asError)(err));
|
|
}
|
|
}
|
|
async answerDiscovery() {
|
|
this.log("Discovery request received, getting device json...");
|
|
const deviceArray = this.describeDevices();
|
|
this.lastDiscoveryAt = new Date().toISOString();
|
|
// send(), as for every message of a device: a publisher that throws or rejects is reported like one that fails
|
|
const topic = await (0, transport_js_1.send)(this.publisher, topics.discoverReply(this.rootTopic), () => deviceArray, (err) => this.fail(err));
|
|
if (topic === "")
|
|
return;
|
|
this.log(`Discovery payloads published to ${topic}`, deviceArray);
|
|
this.emit("discover", deviceArray.length);
|
|
}
|
|
// 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. A device can be declared before connect(): what it sends before the bridge has
|
|
* a connection is not published, and the "error" event says so.
|
|
*/
|
|
addDevice(definition) {
|
|
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(null, 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) {
|
|
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(null, this.rootTopic, name, endpointId, normalizedCategory);
|
|
this.register(device, false);
|
|
return device;
|
|
}
|
|
register(device, endpointHealth) {
|
|
device.alexaInterface = this.options.alexaInterface !== false;
|
|
device.endpointHealth = endpointHealth;
|
|
device.log = (message) => this.log(message);
|
|
device.onPublishError = (err) => this.fail(err); // a failed publish is an "error" event (when listened to), never a rejected send()
|
|
// Through the dispatcher, which lets one answer to a directive pass
|
|
device.publisher = this.dispatcher.publisher;
|
|
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.release(this.devices[i]);
|
|
this.devices.splice(i, 1);
|
|
return true;
|
|
}
|
|
/** Forget every device. */
|
|
clearDevices() {
|
|
for (const d of this.devices)
|
|
this.release(d);
|
|
this.devices = [];
|
|
}
|
|
// A device the bridge forgot hears no directive and publishes nothing
|
|
release(device) {
|
|
device.removeAllListeners();
|
|
device.publisher = null;
|
|
}
|
|
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;
|