Alex2Node/dist/esm/Alex2Node.js
David e0fe9d2812 bridge: subscribe to discover and +/alexaDirective only
The bridge subscribed to <root>/#, so the broker sent back every discover_r, alexaResponce and changeReport
the bridge published, and the backend's <endpointId>/alexaDirective_e. It now subscribes to <root>/discover
and <root>/+/alexaDirective, both named in src/topics.ts; a directive topic is recognised by the root and the
last segment instead of by counting three segments.
The test reads the subscriptions from the broker, publishes seven other topics under the root and one
directive for each of two endpoints: the bridge receives the two directives and the discovery request,
nothing else. Tests: 148 -> 150 (the count in the body of 6789a1a, 108 -> 113, should read 143 -> 148).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-28 19:36:09 +00:00

253 lines
12 KiB
JavaScript

// mqtt reaches Node as CommonJS: the default import works from both builds, a named one needs Node to detect it.
import mqtt from "mqtt";
import Device from "./device/Device.js";
import { checkEndpoint } from "./device/validate.js";
import { EventEmitter } from "events";
import { DisplayCategory } from "./compat/enums.js";
import { DeclarationError } from "./registry/types.js";
import * as topics from "./topics.js";
import { MqttPublisher } from "./transport.js";
export const 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 <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 }
*/
class Alex2MQTT extends 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;
// One for the life of the bridge: the devices keep it over disconnect() and connect()
this.publisher = new MqttPublisher(() => this.client);
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 || 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.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) => {
this.log(`MQTT Message Received`, { topic, payload: message.toString() });
if (topic === topics.discover(this.rootTopic)) {
this.log("Discovery request received, getting device json...");
const deviceArray = this.describeDevices();
this.lastDiscoveryAt = new Date().toISOString();
void this.publisher.publish(topics.discoverReply(this.rootTopic), deviceArray).then((result) => {
if (!result.ok) {
this.fail(result.error);
return;
}
this.log(`Discovery payloads published to ${result.topic}`, deviceArray);
this.emit("discover", deviceArray.length);
});
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)
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. 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(null, this.rootTopic, name, endpointId, (categories ?? []));
for (const [field, value] of Object.entries(described)) {
if (!DESCRIBED.includes(field))
throw new DeclarationError({ endpointId }, `${field} is not a field of an endpoint`);
if (value !== undefined)
Object.assign(device, { [field]: value });
}
checkEndpoint(device.getJSON());
const existing = this.getDevice(endpointId);
if (existing)
throw new 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 || DisplayCategory.LIGHT];
const device = new Device(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.onPublishError = (err) => this.fail(err); // a failed publish is an "error" event (when listened to), never a rejected send()
device.publisher = this.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;
}
}
export default Alex2MQTT;