bridge: publish through a Publisher; devices no longer hold the broker client

src/transport.ts has the Publisher interface, MqttPublisher (the client the bridge has at the time) and
MemoryPublisher (tests without a broker); src/topics.ts names the four topics the library publishes to.
Device, AlexaStatusMessage, AlexaErrorResponse and sendSceneResponse shared three copies of the
publish-and-report code: they now call one send() that resolves the topic or "" and never rejects.
registerDevice() and addDevice() work before connect(); a send() without a connection resolves "" and the
"error" event says to call connect(). unregisterDevice() and clearDevices() take the publisher from the device.
Device.setMqttClient() is gone, the first constructor argument of Device is ignored, and the message classes
take a Publisher where they took the client. Tests: 108 -> 113.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
David 2026-09-28 19:34:39 +00:00
parent 46fa06728c
commit 6789a1a077
39 changed files with 839 additions and 246 deletions

View file

@ -7,6 +7,8 @@ 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";
/** Optional settings for the bridge (1.5.1). Everything has the 1.4.0 behaviour as its default. */
export interface Alex2MQTTOptions {
@ -47,6 +49,8 @@ const DESCRIBED = [
*/
class Alex2MQTT extends EventEmitter {
private client: MqttClient | null = null;
// One for the life of the bridge: the devices keep it over disconnect() and connect()
private readonly publisher = new MqttPublisher(() => this.client);
private devices: Device[] = [];
private MqttHost: string;
private options: Alex2MQTTOptions;
@ -93,7 +97,6 @@ class Alex2MQTT extends EventEmitter {
};
this.client = mqtt.connect(this.MqttHost, options);
for (const d of this.devices) d.setMqttClient(this.client); // after a disconnect(): the devices publish through the new client
this.client.on("connect", () => {
this.connected = true;
@ -115,9 +118,9 @@ class Alex2MQTT extends EventEmitter {
this.log("Discovery request received, getting device json...");
const deviceArray = this.describeDevices();
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);
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);
});
} else if (topic.split("/").length == 3) {
@ -187,15 +190,13 @@ class Alex2MQTT extends EventEmitter {
*
* 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.
* 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: EndpointDefinition): Device {
if (!this.client) {
throw new Error("Must call connect before creating devices");
}
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(this.client, this.rootTopic, name, endpointId, (categories ?? []) as DisplayCategory[]);
const device = new Device(null, this.rootTopic, name, endpointId, (categories ?? []) as DisplayCategory[]);
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 });
@ -212,9 +213,6 @@ class Alex2MQTT extends EventEmitter {
endpointId: string,
displayCategory?: DisplayCategory | DisplayCategory[] | null
): Device {
if (!this.client) {
throw new Error("Must call connect before creating devices");
}
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`;
@ -223,7 +221,7 @@ class Alex2MQTT extends EventEmitter {
}
this.log(`Creating new device with endpoint: ${endpointId}`);
const normalizedCategory = Array.isArray(displayCategory) ? displayCategory : [displayCategory || DisplayCategory.LIGHT];
const device = new Device(this.client, this.rootTopic, name, endpointId, normalizedCategory);
const device = new Device(null, this.rootTopic, name, endpointId, normalizedCategory);
this.register(device, false);
return device;
}
@ -232,6 +230,7 @@ class Alex2MQTT extends EventEmitter {
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);
}
@ -239,15 +238,20 @@ class Alex2MQTT extends EventEmitter {
unregisterDevice(endpointId: string): boolean {
const i = this.devices.findIndex((d) => d.endpointId === endpointId);
if (i === -1) return false;
this.devices[i].removeAllListeners();
this.release(this.devices[i]);
this.devices.splice(i, 1);
return true;
}
/** Forget every device. */
clearDevices(): void {
for (const d of this.devices) d.removeAllListeners();
for (const d of this.devices) this.release(d);
this.devices = [];
}
// A device the bridge forgot hears no directive and publishes nothing
private release(device: Device): void {
device.removeAllListeners();
device.publisher = null;
}
getDevices(): Device[] {
return this.devices.slice();
}

View file

@ -1,9 +1,11 @@
import { randomUUID } from "crypto";
import type { MqttClient } from "mqtt";
import { errorResponse } from "../messages/build.js";
import { AlexaErrors } from "../messages/errors.js";
import type { AlexaError } from "../messages/errors.js";
import type { ErrorResponseMessage } from "../messages/types.js";
import * as topics from "../topics.js";
import { send } from "../transport.js";
import type { Publisher } from "../transport.js";
/** An ErrorResponse as 1.x builds it: device.getErrorMessage(token), setErrorMessage(), send(). */
export class AlexaErrorResponse {
@ -19,7 +21,7 @@ export class AlexaErrorResponse {
private readonly correlationToken: string,
private readonly rootTopic: string,
private readonly endpointId: string,
private readonly mqttClient: MqttClient
private readonly publisher: Publisher | null
) {}
/**
@ -52,13 +54,7 @@ export class AlexaErrorResponse {
/** Resolves with the topic published to, or "" when the publish failed. Never rejects. */
send(sendAsync = false): Promise<string> {
const topic = `${this.rootTopic}/${this.endpointId}/${sendAsync ? "deferredResponse" : "alexaResponce"}`; // "alexaResponce" is how Alex2MQTT and Alex2ESP spell the topic
return new Promise((resolve) => {
this.mqttClient.publish(topic, JSON.stringify(this.toJSON()), (err) => {
if (!err) return resolve(topic);
if (this.onPublishError) this.onPublishError(err); // -> the bridge's "error" event (when somebody listens)
resolve(""); // 1.5.1 rejected here, and an un-caught send() then killed the host on any broker hiccup
});
});
const answer = sendAsync ? topics.deferred : topics.response;
return send(this.publisher, answer(this.rootTopic, this.endpointId), () => this.toJSON(), this.onPublishError);
}
}

View file

@ -1,8 +1,10 @@
import { randomUUID } from "crypto";
import type { MqttClient } from "mqtt";
import { changeReport, deferredResponse, response } from "../messages/build.js";
import { StateBuilder } from "../messages/StateBuilder.js";
import type { ChangeCause, Property } from "../messages/types.js";
import * as topics from "../topics.js";
import { send } from "../transport.js";
import type { Publisher } from "../transport.js";
import { AlexaInterfaceType, TemperatureSensorScale } from "./enums.js";
import type { EndpointHealth, PowerController } from "./enums.js";
@ -27,7 +29,7 @@ export class AlexaStatusMessage {
private readonly correlationToken: string,
private readonly rootTopic: string,
private readonly endpointId: string,
private readonly mqttClient: MqttClient,
private readonly publisher: Publisher | null,
private readonly isResponse: boolean = false,
private readonly isDeferred: boolean = false,
private readonly changeCause: ChangeCause | null = null
@ -139,21 +141,8 @@ export class AlexaStatusMessage {
* to). Never rejects (1.5.1 did, so an un-caught send() could kill the host).
*/
public send(sendAsync: boolean = false): Promise<string> {
const topic = this.changeCause
? `${this.rootTopic}/changeReport`
: `${this.rootTopic}/${this.endpointId}/${sendAsync ? "deferredResponse" : "alexaResponce"}`; // "alexaResponce" is how Alex2MQTT and Alex2ESP spell the topic
return new Promise((resolve) => {
const failed = (err: Error): void => {
if (this.onPublishError) this.onPublishError(err); // -> the bridge's "error" event (when somebody listens)
resolve(""); // 1.5.1 rejected here, and an un-caught send() then killed the host on any broker hiccup
};
let payload: string;
try {
payload = JSON.stringify(this.toJSON());
} catch (err) {
return failed(err as Error);
}
this.mqttClient.publish(topic, payload, (err) => (err ? failed(err) : resolve(topic)));
});
const answer = sendAsync ? topics.deferred : topics.response;
const topic = this.changeCause ? topics.changeReport(this.rootTopic) : answer(this.rootTopic, this.endpointId);
return send(this.publisher, topic, () => this.toJSON(), this.onPublishError);
}
}

View file

@ -1,5 +1,4 @@
import { EventEmitter } from "events";
import type { MqttClient } from "mqtt";
import { AlexaErrorResponse } from "../compat/AlexaErrorResponse.js";
import { AlexaInterface } from "../compat/AlexaInterface.js";
import { AlexaStatusMessage } from "../compat/AlexaStatusMessage.js";
@ -13,6 +12,9 @@ import { EndpointHealth } from "../registry/interfaces/EndpointHealth.js";
import { SchemaError } from "../registry/schema.js";
import { DeclarationError } from "../registry/types.js";
import type { Directives, EndpointView, InterfaceDescriptor, Label, Properties } from "../registry/types.js";
import * as topics from "../topics.js";
import { send } from "../transport.js";
import type { Publisher } from "../transport.js";
import { Capability, commonOptions } from "./Capability.js";
import type { AnyCapability, CapabilityJson, Declaration } from "./Capability.js";
import { checkCapability, checkCapabilityCount, checkEndpoint } from "./validate.js";
@ -72,9 +74,16 @@ class Device extends EventEmitter {
private capabilities: AnyCapability[] = [];
/** Where a failed publish from this device or a message it built is reported; Alex2MQTT sets it (1.5.2). */
public onPublishError?: (err: Error) => void;
/**
* What the device and the messages it builds publish through. The bridge sets it when the device is registered
* and clears it when the device is unregistered; a MemoryPublisher here tests a device without a broker. While it
* is null every send() resolves "" and reports why.
*/
public publisher: Publisher | null = null;
/** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */
constructor(
private mqttClient: MqttClient,
client: unknown,
private rootTopic: string,
public name: string,
public endpointId: string,
@ -87,11 +96,6 @@ class Device extends EventEmitter {
super();
}
/** Internal (1.5.2): Alex2MQTT.connect() re-binds every registered device to its new broker client after a disconnect(). */
setMqttClient(client: MqttClient): void {
this.mqttClient = client;
}
getName(): string {
return this.name;
}
@ -124,7 +128,7 @@ class Device extends EventEmitter {
correlationToken,
this.rootTopic,
this.endpointId,
this.mqttClient
this.publisher
);
msg.onPublishError = this.onPublishError;
return msg;
@ -138,7 +142,7 @@ class Device extends EventEmitter {
correlationToken,
this.rootTopic,
this.endpointId,
this.mqttClient,
this.publisher,
isResponse,
isDeferred
);
@ -152,7 +156,7 @@ class Device extends EventEmitter {
* Without a changed property send() publishes nothing and resolves with "": Alex2MQTT would drop the report.
*/
getChangeReport(cause: ChangeCause = "PHYSICAL_INTERACTION"): AlexaStatusMessage {
const msg = new AlexaStatusMessage("", this.rootTopic, this.endpointId, this.mqttClient, false, false, cause);
const msg = new AlexaStatusMessage("", this.rootTopic, this.endpointId, this.publisher, false, false, cause);
msg.onPublishError = this.onPublishError;
return msg;
}
@ -161,15 +165,9 @@ class Device extends EventEmitter {
* (1.5.1). Resolves with the topic published to, or "" when the publish failed (never rejects, 1.5.2).
*/
sendSceneResponse(correlationToken: string, activated: boolean, cause: ChangeCause = "VOICE_INTERACTION", sendAsync = false): Promise<string> {
const payload = sceneEvent({ endpointId: this.endpointId, correlationToken, activated, cause });
const topic = `${this.rootTopic}/${this.endpointId}/${sendAsync ? "deferredResponse" : "alexaResponce"}`;
return new Promise((resolve) => {
this.mqttClient.publish(topic, JSON.stringify(payload), (err) => {
if (!err) return resolve(topic);
if (this.onPublishError) this.onPublishError(err); // -> the bridge's "error" event (when somebody listens)
resolve(""); // 1.5.1 rejected here, and an un-caught send() then killed the host on any broker hiccup
});
});
const answer = sendAsync ? topics.deferred : topics.response;
const build = () => sceneEvent({ endpointId: this.endpointId, correlationToken, activated, cause });
return send(this.publisher, answer(this.rootTopic, this.endpointId), build, this.onPublishError);
}
/** The capabilities the device declared, in the order it declared them. */
getCapabilities(): AnyCapability[] {

View file

@ -34,6 +34,11 @@ export type {
Property, PropertyOptions, ResponseMessage, SceneEventMessage,
} from "./messages/index.js";
// Publishing: the topics of the Alex2MQTT contract, and a publisher that needs no broker for tests
export * as topics from "./topics.js";
export { MemoryPublisher } from "./transport.js";
export type { Publisher, PublishResult } from "./transport.js";
// 1.x
export { AlexaInterface } from "./compat/AlexaInterface.js";
export type { SupportedMode } from "./compat/AlexaInterface.js";

14
src/topics.ts Normal file
View file

@ -0,0 +1,14 @@
// The topics of the Alex2MQTT contract, named once. The backend and Alex2ESP use the same names, so none of them
// can change here alone.
/** Where the bridge answers a discovery request with its endpoints. */
export const discoverReply = (root: string): string => `${root}/discover_r`;
/** Where the answer to a directive goes. "alexaResponce" is how the backend and Alex2ESP spell it. */
export const response = (root: string, endpointId: string): string => `${root}/${endpointId}/alexaResponce`;
/** Where the answer goes once a DeferredResponse was sent for the directive. */
export const deferred = (root: string, endpointId: string): string => `${root}/${endpointId}/deferredResponse`;
/** Where every ChangeReport of the root goes: the backend adds the user's token and posts it to Alexa. */
export const changeReport = (root: string): string => `${root}/changeReport`;

89
src/transport.ts Normal file
View file

@ -0,0 +1,89 @@
import type { MqttClient } from "mqtt";
/** What became of a publish. A failure is a value: publish() never rejects. */
export type PublishResult =
| { ok: true; topic: string }
| { ok: false; topic: string; error: Error };
/** Where the bridge, its devices and their messages publish. The message is sent as JSON. */
export interface Publisher {
publish(topic: string, message: object): Promise<PublishResult>;
}
const asError = (err: unknown): Error => (err instanceof Error ? err : new Error(String(err)));
/**
* Publishes to the broker through the client the bridge has at the time: a new one after disconnect() and connect(),
* none before connect() and after disconnect(). Without a client the publish is refused. With a client that lost
* its connection mqtt.js keeps the message and sends it after the reconnect, and the publish resolves then.
*/
export class MqttPublisher implements Publisher {
constructor(private readonly client: () => MqttClient | null) {}
publish(topic: string, message: object): Promise<PublishResult> {
const client = this.client();
return new Promise((resolve) => {
const failed = (err: unknown): void => resolve({ ok: false, topic, error: asError(err) });
if (!client) {
failed(new Error(`nothing was published to ${topic}: the bridge is not connected, call connect() first`));
return;
}
try {
client.publish(topic, JSON.stringify(message), (err) => (err ? failed(err) : resolve({ ok: true, topic })));
} catch (err) {
failed(err);
}
});
}
}
/**
* Keeps what it is asked to publish, for tests and dry runs without a broker:
*
* const sent = new MemoryPublisher();
* device.publisher = sent;
* await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send();
* sent.published[0] // { topic: "<root>/<endpointId>/alexaResponce", message: { event, context } }
*/
export class MemoryPublisher implements Publisher {
/** Oldest first. The message is what a subscriber gets after JSON.parse. */
public readonly published: Array<{ topic: string; message: any }> = [];
/** Set it and every publish fails with this error, as a publish to a broker that is gone does. */
public failWith: Error | null = null;
publish(topic: string, message: object): Promise<PublishResult> {
if (this.failWith) return Promise.resolve({ ok: false, topic, error: this.failWith });
try {
this.published.push({ topic, message: JSON.parse(JSON.stringify(message)) });
} catch (err) {
return Promise.resolve({ ok: false, topic, error: asError(err) });
}
return Promise.resolve({ ok: true, topic });
}
}
/**
* send() as 1.x promises it: resolves with the topic, or with "" when nothing was published, and never rejects
* (1.5.1 rejected, and a send() nobody caught killed the host on any broker hiccup). What went wrong goes to report:
* the message cannot be built, the device is on no bridge, or the publish failed.
*/
export async function send(
publisher: Publisher | null,
topic: string,
build: () => object,
report?: (err: Error) => void
): Promise<string> {
let result: PublishResult;
try {
const message = build();
if (!publisher) {
throw new Error(`nothing was published to ${topic}: the device is on no bridge, register it with addDevice() or registerDevice()`);
}
result = await publisher.publish(topic, message);
} catch (err) {
result = { ok: false, topic, error: asError(err) };
}
if (result.ok) return topic;
if (report) report(result.error);
return "";
}