1.5.1: broker errors as events that never crash the host (error only when listened to; connect/offline/reconnect/close; connected flag), configurable broker host + mqtt options + log hook, Alexa.ChangeReport via Device.getChangeReport (<root>/changeReport), SceneController discovery + sendSceneResponse, addCapability options (proactivelyReported), unregisterDevice/clearDevices/getDevices/disconnect, discover/directive events, promise-returning quiet send(); tests on an in-process broker (aedes); readme

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
David 2026-09-22 19:29:38 +00:00
parent e96850bdd6
commit 717af632de
16 changed files with 2080 additions and 225 deletions

View file

@ -1,122 +1,143 @@
import mqtt, { MqttClient } from "mqtt";
import mqtt, { MqttClient, IClientOptions } from "mqtt";
import Device from "./Device";
import { EventEmitter } from "events";
import { DisplayCategory } from "./DisplayCategory";
/** Optional settings for the bridge (1.5.1). Everything has the 1.4.0 behaviour as its default. */
export interface Alex2MQTTOptions {
/** The broker URL. Default: the public Alex2MQTT broker, mqtt://Alex2MQTT.stormysdream.club:1883. */
host?: string;
/** Extra mqtt.js client options (reconnectPeriod, connectTimeout, clientId, ...). Merged over the defaults. */
mqtt?: IClientOptions;
/** Where log lines go. Default: console (only when debugLogging is on). */
log?: (message: string, detail?: unknown) => void;
}
export const 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 <root>/#
* "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 {
private client: MqttClient | null = null;
private devices: Device[] = [];
private MqttHost = "mqtt://Alex2MQTT.stormysdream.club";
private MqttHost: string;
private options: Alex2MQTTOptions;
/** true while the broker connection is up. */
public connected = false;
/** ISO time of the last discovery request answered, null before the first. */
public lastDiscoveryAt: string | null = null;
constructor(
private username: string,
private password: string,
private rootTopic: string,
private debugLogging: boolean
private debugLogging: boolean = false,
options: Alex2MQTTOptions = {}
) {
super(); // Initialize EventEmitter
this.options = options || {};
this.MqttHost = this.options.host || DEFAULT_HOST;
}
private log(message: string, detail?: unknown): void {
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. */
private fail(err: Error): void {
this.log("error: " + err.message);
if (this.listenerCount("error") > 0) this.emit("error", err);
}
connect(): void {
const options = {
if (this.client) return;
const options: IClientOptions = {
username: this.username,
password: this.password,
reconnectPeriod: 1000,
connectTimeout: 10000,
...(this.options.mqtt || {}),
};
this.client = mqtt.connect(this.MqttHost, options);
this.client.on("connect", () => {
if (this.debugLogging) {
console.log("[Alex2Node.ts] Connected to MQTT broker");
}
this.connected = true;
this.log("Connected to MQTT broker");
this.client!.subscribe(this.rootTopic + "/#", (err) => {
if (this.debugLogging && !err) {
console.log("[Alex2Node.ts] Subscribed to topics");
}
if (err) {
console.error("[Alex2Node.ts] MQTT connection error:", err);
this.emit("error", 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: Error) => { this.connected = false; this.fail(err); });
this.client.on("message", (topic: string, message: Buffer) => {
if (this.debugLogging) {
console.log(`[Alex2Node.ts] MQTT Message Received`);
console.log(` Topic: ${topic}`);
console.log(` Payload: ${message.toString()}`);
}
this.log(`MQTT Message Received`, { topic, payload: message.toString() });
if (topic == `${this.rootTopic}/discover`) {
if (this.debugLogging) {
console.log(
"[Alex2Node.ts] Discovery request received, getting device json..."
);
}
this.log("Discovery request received, getting device json...");
const deviceArray = this.devices.map((device) => device.getJSON());
if (deviceArray.length > 0) {
this.client!.publish(
topic + "_r",
JSON.stringify(deviceArray),
(err) => {
if (err) {
console.error("[Alex2Node.ts] MQTT connection error:", err);
this.emit("error", err);
} else if (this.debugLogging) {
console.log(
`[Alex2Node.ts] Discovery payloads published to ${
topic + "_r"
}`
);
console.log(JSON.stringify(deviceArray, null, 2));
}
}
);
}
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 [root, endpointId, directiveType] = topic.split("/");
const [, endpointId, directiveType] = topic.split("/");
if (directiveType != "alexaDirective") {
return;
}
const device = this.devices.find((d) => d.endpointId === endpointId);
if (!device) {
if (this.debugLogging) {
console.warn(
`[Alex2Node.ts] No device found for endpointId: ${endpointId}`
);
}
this.log(`No device found for endpointId: ${endpointId}`);
return;
}
let payload;
try {
payload = JSON.parse(message.toString());
} catch (e) {
console.error(`[Alex2Node.ts] Failed to parse payload JSON`, e);
this.log(`Failed to parse payload JSON`, e);
return;
}
if (
payload.header &&
payload.header.namespace === "Alexa" &&
payload.header.name === "ReportState"
) {
if (this.debugLogging) {
console.log(
`[Alex2Node.ts] ReportState directive for device ${endpointId}`
);
}
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);
}
}
});
this.client.on("error", (err: Error) => {
if (this.debugLogging) {
console.error("[Alex2Node.ts] Connection failed:", err);
}
console.error("[Alex2Node.ts] MQTT connection error:", err);
this.emit("error", err);
}
/** Close the broker connection (resolves once closed). The devices stay registered; connect() again reuses them. */
disconnect(): Promise<void> {
return new Promise((resolve) => {
const c = this.client;
if (!c) return resolve();
this.client = null;
this.connected = false;
c.end(true, {}, () => resolve());
});
}
@ -128,17 +149,13 @@ class Alex2MQTT extends EventEmitter {
if (!this.client) {
throw new Error("Must call connect before creating devices");
}
if (this.debugLogging) {
console.log(
`[Alex2Node.ts] Creating new device with endpoint: ${endpointId}`
);
}
this.log(`Creating new device with endpoint: ${endpointId}`);
const normalizedCategory =
displayCategory === null
? [DisplayCategory.LIGHT]
: Array.isArray(displayCategory)
? displayCategory
: [displayCategory||DisplayCategory.LIGHT];
: [displayCategory || DisplayCategory.LIGHT];
const device = new Device(
this.client!,
@ -150,6 +167,32 @@ class Alex2MQTT extends EventEmitter {
this.devices.push(device);
return device;
}
/** Forget a device (its listeners with it). Returns false when there was none. */
unregisterDevice(endpointId: string): boolean {
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(): void {
for (const d of this.devices) d.removeAllListeners();
this.devices = [];
}
getDevices(): Device[] {
return this.devices.slice();
}
getDevice(endpointId: string): Device | undefined {
return this.devices.find((d) => d.endpointId === endpointId);
}
getRootTopic(): string {
return this.rootTopic;
}
getHost(): string {
return this.MqttHost;
}
}
export default Alex2MQTT;

View file

@ -125,7 +125,10 @@ export class AlexaErrorResponse {
};
}
send(sendAsync = false) {
toJSON(): { event: any } {
return { event: this.event };
}
send(sendAsync = false): Promise<string> {
const payload = {
event: this.event,
};
@ -133,16 +136,8 @@ export class AlexaErrorResponse {
sendAsync ? "deferredResponse" : "alexaResponce"
}`; //Yes this should be response but it is incorrect in both Alex2MQTT and Alex2ESP so for consistency is is wrong here too
const payloadStr = JSON.stringify(payload);
this.mqttClient.publish(topic, payloadStr, (err) => {
if (err) {
console.error(
"[AlexaErrorResponse] Failed to publish status message:",
err
);
} else {
console.log(`[AlexaErrorResponse] Sent message to ${topic}`);
}
return new Promise((resolve, reject) => {
this.mqttClient.publish(topic, payloadStr, (err) => (err ? reject(err) : resolve(topic)));
});
}
}

View file

@ -256,6 +256,7 @@ export class AlexaInterface {
case AlexaInterfaceType.REMOTE_VIDEO_PLAYER:
case AlexaInterfaceType.RTC_SESSION_CONTROLLER:
case AlexaInterfaceType.SCENE_CONTROLLER:
return []; // scenes have no reportable properties: Activate/Deactivate answer with ActivationStarted (Device.sendSceneResponse)
case AlexaInterfaceType.SECURITY_PANEL_CONTROLLER:
case AlexaInterfaceType.SEEK_CONTROLLER:
case AlexaInterfaceType.SIMPLE_EVENT_SOURCE:
@ -290,7 +291,12 @@ export class AlexaInterface {
supported: this.getProps().map(name => ({ name })),
},
};
if(this.type==AlexaInterfaceType.THERMOSTAT_CONTROLLER){
if (this.type == AlexaInterfaceType.SCENE_CONTROLLER) {
// Alexa.SceneController v3 carries no properties block: supportsDeactivation + proactivelyReported at the top level
delete doc.properties;
doc.supportsDeactivation = true;
doc.proactivelyReported = this.proactivelyReported;
} else if(this.type==AlexaInterfaceType.THERMOSTAT_CONTROLLER){
doc.configuration={
"supportedModes": ["HEAT", "COOL", "AUTO", "OFF"],
"supportsScheduling": false

View file

@ -47,8 +47,15 @@ interface ContextProperty {
instance?: string;
}
/** Why a ChangeReport is sent (Alexa.ChangeReport payload.change.cause.type). */
export type ChangeCause = "APP_INTERACTION" | "PHYSICAL_INTERACTION" | "PERIODIC_POLL" | "RULE_TRIGGER" | "VOICE_INTERACTION";
export class AlexaStatusMessage {
private context: { properties: ContextProperty[] } = { properties: [] };
/** ChangeReport only: the properties that changed (payload.change.properties); the rest go to context. */
private changeProps: ContextProperty[] = [];
private changeCause: ChangeCause | null = null;
private target: "context" | "change" = "context";
private event: {
header: DirectiveHeader;
endpoint: DirectiveEndpoint;
@ -65,16 +72,21 @@ export class AlexaStatusMessage {
endpointId: string,
mqttClient: MqttClient,
isResponse: boolean = false,
isDeferred: boolean = false
isDeferred: boolean = false,
changeCause: ChangeCause | null = null
) {
this.rootTopic = rootTopic;
this.endpointId = endpointId;
this.mqttClient = mqttClient;
this.isDeferred = isDeferred;
this.changeCause = changeCause;
if (changeCause) this.target = "change";
this.event = {
header: {
namespace: "Alexa",
name: isDeferred
name: changeCause
? "ChangeReport"
: isDeferred
? "DeferredResponse"
: isResponse
? "Response"
@ -112,9 +124,33 @@ export class AlexaStatusMessage {
uncertaintyInMilliseconds,
};
if (instance) prop.instance = instance;
this.context.properties.push(prop);
(this.target === "change" ? this.changeProps : this.context.properties).push(prop);
return this;
}
/** ChangeReport: the add*Prop calls that follow describe what CHANGED (the default for a change report). */
public changed(): this {
this.target = "change";
return this;
}
/** ChangeReport: the add*Prop calls that follow describe the other, unchanged properties (context). */
public unchanged(): this {
this.target = "context";
return this;
}
/** True for a ChangeReport (Device.getChangeReport). */
public isChangeReport(): boolean {
return this.changeCause !== null;
}
/** The message as it will be published (for tests and logging). */
public toJSON(): { event: any; context: { properties: ContextProperty[] } | null } {
if (this.changeCause) {
return {
event: { ...this.event, payload: { change: { cause: { type: this.changeCause }, properties: this.changeProps } } },
context: this.context,
};
}
return { context: this.isDeferred ? null : this.context, event: this.event };
}
public addModeControllerProp(instance: string, value: string, uncertaintyInMs = 0): this {
return this.addProperty(
AlexaInterfaceType.MODE_CONTROLLER,
@ -250,23 +286,18 @@ export class AlexaStatusMessage {
return this;
}
public send(sendAsync: boolean): void {
const payload = {
context: this.isDeferred ? null : this.context,
event: this.event,
};
const topic = `${this.rootTopic}/${this.endpointId}/${sendAsync?"deferredResponse":"alexaResponce"}`; //Yes this should be response but it is incorrect in both Alex2MQTT and Alex2ESP so for consistency is is wrong here too
const payloadStr = JSON.stringify(payload);
this.mqttClient.publish(topic, payloadStr, (err) => {
if (err) {
console.error(
"[AlexaStatusMessage] Failed to publish status message:",
err
);
} else {
console.log(`[AlexaStatusMessage] Sent message to ${topic}`);
}
/**
* Publish: a Response/StateReport to <root>/<endpoint>/alexaResponce (sendAsync: deferredResponse), a ChangeReport
* to <root>/changeReport (Alex2MQTT adds the user's token and posts it to the Alexa event gateway). Resolves with
* the topic; rejects on a publish error (1.4.0 only logged).
*/
public send(sendAsync: boolean = false): Promise<string> {
const payloadStr = JSON.stringify(this.toJSON());
const topic = this.changeCause
? `${this.rootTopic}/changeReport`
: `${this.rootTopic}/${this.endpointId}/${sendAsync ? "deferredResponse" : "alexaResponce"}`; //Yes this should be response but it is incorrect in both Alex2MQTT and Alex2ESP so for consistency is is wrong here too
return new Promise((resolve, reject) => {
this.mqttClient.publish(topic, payloadStr, (err) => (err ? reject(err) : resolve(topic)));
});
}
}

View file

@ -2,7 +2,8 @@ import { AlexaInterface, AlexaInterfaceType } from "./AlexaInterface";
import { DisplayCategory } from "./DisplayCategory";
import { EventEmitter } from "events";
import { MqttClient } from "mqtt";
import { AlexaStatusMessage } from "./AlexaStatusMessage";
import { AlexaStatusMessage, ChangeCause } from "./AlexaStatusMessage";
import { v4 as uuidv4 } from "uuid";
import { AlexaErrorResponse } from "./AlexaErrorResponse";
class Device extends EventEmitter {
@ -72,6 +73,35 @@ class Device extends EventEmitter {
isDeferred
);
}
/**
* A proactive Alexa.ChangeReport (1.5.1): add the changed properties (the default target), optionally
* .unchanged() then the others, and .send() - it goes to <root>/changeReport, which Alex2MQTT forwards to the Alexa
* event gateway with the user's token. Alexa only accepts it for capabilities registered with proactivelyReported.
*/
getChangeReport(cause: ChangeCause = "PHYSICAL_INTERACTION"): AlexaStatusMessage {
return new AlexaStatusMessage("", this.rootTopic, this.endpointId, this.mqttClient, false, false, cause);
}
/**
* Alexa.SceneController: answer an Activate / Deactivate directive with ActivationStarted / DeactivationStarted
* (1.5.1). Resolves with the topic published to.
*/
sendSceneResponse(correlationToken: string, activated: boolean, cause: ChangeCause = "VOICE_INTERACTION", sendAsync = false): Promise<string> {
const payload = {
context: {},
event: {
header: { namespace: "Alexa.SceneController", name: activated ? "ActivationStarted" : "DeactivationStarted", messageId: uuidv4(), correlationToken, payloadVersion: "3" },
endpoint: { endpointId: this.endpointId },
payload: { cause: { type: cause }, timestamp: new Date().toISOString() },
},
};
const topic = `${this.rootTopic}/${this.endpointId}/${sendAsync ? "deferredResponse" : "alexaResponce"}`;
return new Promise((resolve, reject) => {
this.mqttClient.publish(topic, JSON.stringify(payload), (err) => (err ? reject(err) : resolve(topic)));
});
}
getCapabilities(): AlexaInterface[] {
return this.capabilities.slice();
}
setManufacturerName(name: string): void {
this.manufacturerName = name;
}
@ -100,15 +130,9 @@ class Device extends EventEmitter {
return this.softwareVersion;
}
addCapability(type: AlexaInterfaceType): AlexaInterface {
// const existing = this.capabilities.find(
// (iface) => iface.getType() === type
// );
// if (existing) {
// return existing;
// }
const newCapability = new AlexaInterface(type);
/** Add a capability. 1.5.1: options.retrievable / proactivelyReported / instance (a change report needs proactivelyReported). */
addCapability(type: AlexaInterfaceType, options: { retrievable?: boolean; proactivelyReported?: boolean; instance?: string } = {}): AlexaInterface {
const newCapability = new AlexaInterface(type, options.retrievable ?? true, options.proactivelyReported ?? false, options.instance ?? "");
this.capabilities.push(newCapability);
// console.log(

View file

@ -1,6 +1,9 @@
export { default as Alex2MQTT } from './Alex2Node';
export { AlexaInterfaceType } from './AlexaInterface';
export { ActionMapping, AlexaActions } from './ActionMapping';
export {DisplayCategory} from "./DisplayCategory";
export { PowerController,EndpointHealth,TemperatureSensorScale,ThermostatMode } from './AlexaStatusMessage';
export { AlexaErrorType,AlexaErrorResponse} from './AlexaErrorResponse';
export { default as Alex2MQTT, DEFAULT_HOST } from "./Alex2Node";
export type { Alex2MQTTOptions } from "./Alex2Node";
export { default as Device } from "./Device";
export { AlexaInterface, AlexaInterfaceType } from "./AlexaInterface";
export { ActionMapping, AlexaActions } from "./ActionMapping";
export { DisplayCategory } from "./DisplayCategory";
export { AlexaStatusMessage, PowerController, EndpointHealth, TemperatureSensorScale, ThermostatMode } from "./AlexaStatusMessage";
export type { ChangeCause } from "./AlexaStatusMessage";
export { AlexaErrorType, AlexaErrorResponse } from "./AlexaErrorResponse";