Alex2Node/dist/esm/Alex2Node.js
David feff828ec1 dispatch: typed handlers, respond/defer/error, automatic ErrorResponse, watchdog
src/dispatcher.ts routes a directive to capability.on(name | "*"), device.onDirective() or
device.onReportState(), with the payload checked by the descriptor of the interface. The
DirectiveContext answers with respond/report/defer/error; respond() and report() start from
device.state(). A handler or 1.x listener that throws or rejects is answered with INTERNAL_ERROR
(an AlexaError with itself) and reported through "error", never as an unhandled rejection.
Undeclared interface, unknown directive, AdjustMode on an unordered mode and a missing handler
get INVALID_DIRECTIVE, a bad payload INVALID_VALUE; a device with an "Event" or "ReportState"
listener keeps the answer to itself. Every answer passes the dispatcher: the first one per
correlationToken is published, a second is refused. No answer within answerWithinMs (6500,
0 = off, unref'd timer) sends INTERNAL_ERROR and emits "unanswered". A messageId that arrived
in the last 60 s is dropped. New: "unknownEndpoint", options answerUnknownEndpoints, publisher,
timers, and bridge.receive() to run a bridge without a broker. 168 tests pass (18 new).

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

262 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 { Dispatcher } from "./dispatcher.js";
import * as topics from "./topics.js";
import { asError, 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 }
* "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 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 || DEFAULT_HOST;
this.publisher = this.options.publisher ?? new MqttPublisher(() => this.client);
this.dispatcher = new 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.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)) {
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(asError(err));
}
}
async answerDiscovery() {
this.log("Discovery request received, getting device json...");
const deviceArray = this.describeDevices();
this.lastDiscoveryAt = new Date().toISOString();
const result = await this.publisher.publish(topics.discoverReply(this.rootTopic), deviceArray);
if (!result.ok) {
this.fail(result.error);
return;
}
this.log(`Discovery payloads published to ${result.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(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()
// 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;
}
}
export default Alex2MQTT;