diff --git a/dist/cjs/Alex2Node.js b/dist/cjs/Alex2Node.js index 330a918..559a668 100644 --- a/dist/cjs/Alex2Node.js +++ b/dist/cjs/Alex2Node.js @@ -57,7 +57,7 @@ const DESCRIBED = [ * 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 /# + * "connect" connected (or reconnected) and subscribed to /discover and /+/alexaDirective * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed @@ -116,7 +116,7 @@ class Alex2MQTT extends events_1.EventEmitter { this.client.on("connect", () => { this.connected = true; this.log("Connected to MQTT broker"); - this.client.subscribe(this.rootTopic + "/#", (err) => { + this.client.subscribe(topics.subscriptions(this.rootTopic), (err) => { if (err) { this.fail(err); return; @@ -131,7 +131,7 @@ class Alex2MQTT extends events_1.EventEmitter { 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 == `${this.rootTopic}/discover`) { + if (topic === topics.discover(this.rootTopic)) { this.log("Discovery request received, getting device json..."); const deviceArray = this.describeDevices(); this.lastDiscoveryAt = new Date().toISOString(); @@ -143,35 +143,34 @@ class Alex2MQTT extends events_1.EventEmitter { this.log(`Discovery payloads published to ${result.topic}`, deviceArray); this.emit("discover", deviceArray.length); }); + return; } - else if (topic.split("/").length == 3) { - const [, endpointId, directiveType] = topic.split("/"); - if (directiveType != "alexaDirective") { - 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); - } + // 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); } }); } diff --git a/dist/cjs/topics.js b/dist/cjs/topics.js index 45f2d53..ba276f5 100644 --- a/dist/cjs/topics.js +++ b/dist/cjs/topics.js @@ -2,10 +2,28 @@ // The topics of the Alex2MQTT contract, named once. The backend and Alex2ESP use the same names, so none of them // can change here alone. Object.defineProperty(exports, "__esModule", { value: true }); -exports.changeReport = exports.deferred = exports.response = exports.discoverReply = void 0; +exports.changeReport = exports.deferred = exports.response = exports.subscriptions = exports.directive = exports.discoverReply = exports.discover = void 0; +exports.directiveEndpoint = directiveEndpoint; +/** Where the backend asks for the endpoints of the root. */ +const discover = (root) => `${root}/discover`; +exports.discover = discover; /** Where the bridge answers a discovery request with its endpoints. */ const discoverReply = (root) => `${root}/discover_r`; exports.discoverReply = discoverReply; +/** Where the backend publishes the directives for one endpoint. */ +const directive = (root, endpointId) => `${root}/${endpointId}/alexaDirective`; +exports.directive = directive; +/** + * What the bridge subscribes to: the discovery requests and the directives of every endpoint. Not /#, which + * also delivers what the bridge itself publishes and the backend's /alexaDirective_e. + */ +const subscriptions = (root) => [(0, exports.discover)(root), (0, exports.directive)(root, "+")]; +exports.subscriptions = subscriptions; +/** The endpointId of a directive topic, null for any other topic. */ +function directiveEndpoint(root, topic) { + const [, endpointId, ...rest] = topic.startsWith(`${root}/`) ? topic.slice(root.length).split("/") : []; + return endpointId && rest.length === 1 && rest[0] === "alexaDirective" ? endpointId : null; +} /** Where the answer to a directive goes. "alexaResponce" is how the backend and Alex2ESP spell it. */ const response = (root, endpointId) => `${root}/${endpointId}/alexaResponce`; exports.response = response; diff --git a/dist/esm/Alex2Node.d.ts b/dist/esm/Alex2Node.d.ts index e7d45dc..25e9bb4 100644 --- a/dist/esm/Alex2Node.d.ts +++ b/dist/esm/Alex2Node.d.ts @@ -23,7 +23,7 @@ export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; * 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 /# + * "connect" connected (or reconnected) and subscribed to /discover and /+/alexaDirective * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed diff --git a/dist/esm/Alex2Node.js b/dist/esm/Alex2Node.js index 0be84a6..ec92742 100644 --- a/dist/esm/Alex2Node.js +++ b/dist/esm/Alex2Node.js @@ -18,7 +18,7 @@ const DESCRIBED = [ * 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 /# + * "connect" connected (or reconnected) and subscribed to /discover and /+/alexaDirective * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed @@ -77,7 +77,7 @@ class Alex2MQTT extends EventEmitter { this.client.on("connect", () => { this.connected = true; this.log("Connected to MQTT broker"); - this.client.subscribe(this.rootTopic + "/#", (err) => { + this.client.subscribe(topics.subscriptions(this.rootTopic), (err) => { if (err) { this.fail(err); return; @@ -92,7 +92,7 @@ class Alex2MQTT extends EventEmitter { 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 == `${this.rootTopic}/discover`) { + if (topic === topics.discover(this.rootTopic)) { this.log("Discovery request received, getting device json..."); const deviceArray = this.describeDevices(); this.lastDiscoveryAt = new Date().toISOString(); @@ -104,35 +104,34 @@ class Alex2MQTT extends EventEmitter { this.log(`Discovery payloads published to ${result.topic}`, deviceArray); this.emit("discover", deviceArray.length); }); + return; } - else if (topic.split("/").length == 3) { - const [, endpointId, directiveType] = topic.split("/"); - if (directiveType != "alexaDirective") { - 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); - } + // 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); } }); } diff --git a/dist/esm/topics.d.ts b/dist/esm/topics.d.ts index 14bb390..7ff3912 100644 --- a/dist/esm/topics.d.ts +++ b/dist/esm/topics.d.ts @@ -1,5 +1,16 @@ +/** Where the backend asks for the endpoints of the root. */ +export declare const discover: (root: string) => string; /** Where the bridge answers a discovery request with its endpoints. */ export declare const discoverReply: (root: string) => string; +/** Where the backend publishes the directives for one endpoint. */ +export declare const directive: (root: string, endpointId: string) => string; +/** + * What the bridge subscribes to: the discovery requests and the directives of every endpoint. Not /#, which + * also delivers what the bridge itself publishes and the backend's /alexaDirective_e. + */ +export declare const subscriptions: (root: string) => string[]; +/** The endpointId of a directive topic, null for any other topic. */ +export declare function directiveEndpoint(root: string, topic: string): string | null; /** Where the answer to a directive goes. "alexaResponce" is how the backend and Alex2ESP spell it. */ export declare const response: (root: string, endpointId: string) => string; /** Where the answer goes once a DeferredResponse was sent for the directive. */ diff --git a/dist/esm/topics.js b/dist/esm/topics.js index ba81060..e00d375 100644 --- a/dist/esm/topics.js +++ b/dist/esm/topics.js @@ -1,7 +1,21 @@ // 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 backend asks for the endpoints of the root. */ +export const discover = (root) => `${root}/discover`; /** Where the bridge answers a discovery request with its endpoints. */ export const discoverReply = (root) => `${root}/discover_r`; +/** Where the backend publishes the directives for one endpoint. */ +export const directive = (root, endpointId) => `${root}/${endpointId}/alexaDirective`; +/** + * What the bridge subscribes to: the discovery requests and the directives of every endpoint. Not /#, which + * also delivers what the bridge itself publishes and the backend's /alexaDirective_e. + */ +export const subscriptions = (root) => [discover(root), directive(root, "+")]; +/** The endpointId of a directive topic, null for any other topic. */ +export function directiveEndpoint(root, topic) { + const [, endpointId, ...rest] = topic.startsWith(`${root}/`) ? topic.slice(root.length).split("/") : []; + return endpointId && rest.length === 1 && rest[0] === "alexaDirective" ? endpointId : null; +} /** Where the answer to a directive goes. "alexaResponce" is how the backend and Alex2ESP spell it. */ export const response = (root, endpointId) => `${root}/${endpointId}/alexaResponce`; /** Where the answer goes once a DeferredResponse was sent for the directive. */ diff --git a/dist/types/Alex2Node.d.ts b/dist/types/Alex2Node.d.ts index e7d45dc..25e9bb4 100644 --- a/dist/types/Alex2Node.d.ts +++ b/dist/types/Alex2Node.d.ts @@ -23,7 +23,7 @@ export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; * 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 /# + * "connect" connected (or reconnected) and subscribed to /discover and /+/alexaDirective * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed diff --git a/dist/types/topics.d.ts b/dist/types/topics.d.ts index 14bb390..7ff3912 100644 --- a/dist/types/topics.d.ts +++ b/dist/types/topics.d.ts @@ -1,5 +1,16 @@ +/** Where the backend asks for the endpoints of the root. */ +export declare const discover: (root: string) => string; /** Where the bridge answers a discovery request with its endpoints. */ export declare const discoverReply: (root: string) => string; +/** Where the backend publishes the directives for one endpoint. */ +export declare const directive: (root: string, endpointId: string) => string; +/** + * What the bridge subscribes to: the discovery requests and the directives of every endpoint. Not /#, which + * also delivers what the bridge itself publishes and the backend's /alexaDirective_e. + */ +export declare const subscriptions: (root: string) => string[]; +/** The endpointId of a directive topic, null for any other topic. */ +export declare function directiveEndpoint(root: string, topic: string): string | null; /** Where the answer to a directive goes. "alexaResponce" is how the backend and Alex2ESP spell it. */ export declare const response: (root: string, endpointId: string) => string; /** Where the answer goes once a DeferredResponse was sent for the directive. */ diff --git a/src/Alex2Node.ts b/src/Alex2Node.ts index c05f6f5..ee09778 100644 --- a/src/Alex2Node.ts +++ b/src/Alex2Node.ts @@ -38,7 +38,7 @@ const DESCRIBED = [ * 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 /# + * "connect" connected (or reconnected) and subscribed to /discover and /+/alexaDirective * "offline" the connection dropped; mqtt.js reconnects on its own (reconnectPeriod, default 1 s) * "reconnect" a reconnect attempt starts * "close" the connection closed @@ -101,7 +101,7 @@ class Alex2MQTT extends EventEmitter { this.client.on("connect", () => { this.connected = true; this.log("Connected to MQTT broker"); - this.client!.subscribe(this.rootTopic + "/#", (err) => { + this.client!.subscribe(topics.subscriptions(this.rootTopic), (err) => { if (err) { this.fail(err); return; } this.log("Subscribed to topics"); this.emit("connect"); @@ -114,7 +114,7 @@ class Alex2MQTT extends EventEmitter { this.client.on("message", (topic: string, message: Buffer) => { this.log(`MQTT Message Received`, { topic, payload: message.toString() }); - if (topic == `${this.rootTopic}/discover`) { + if (topic === topics.discover(this.rootTopic)) { this.log("Discovery request received, getting device json..."); const deviceArray = this.describeDevices(); this.lastDiscoveryAt = new Date().toISOString(); @@ -123,31 +123,30 @@ class Alex2MQTT extends EventEmitter { this.log(`Discovery payloads published to ${result.topic}`, deviceArray); this.emit("discover", deviceArray.length); }); - } else if (topic.split("/").length == 3) { - const [, endpointId, directiveType] = topic.split("/"); - if (directiveType != "alexaDirective") { - 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); - } + 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); } }); } diff --git a/src/topics.ts b/src/topics.ts index 25b7e3f..586dd1c 100644 --- a/src/topics.ts +++ b/src/topics.ts @@ -1,9 +1,27 @@ // 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 backend asks for the endpoints of the root. */ +export const discover = (root: string): string => `${root}/discover`; + /** Where the bridge answers a discovery request with its endpoints. */ export const discoverReply = (root: string): string => `${root}/discover_r`; +/** Where the backend publishes the directives for one endpoint. */ +export const directive = (root: string, endpointId: string): string => `${root}/${endpointId}/alexaDirective`; + +/** + * What the bridge subscribes to: the discovery requests and the directives of every endpoint. Not /#, which + * also delivers what the bridge itself publishes and the backend's /alexaDirective_e. + */ +export const subscriptions = (root: string): string[] => [discover(root), directive(root, "+")]; + +/** The endpointId of a directive topic, null for any other topic. */ +export function directiveEndpoint(root: string, topic: string): string | null { + const [, endpointId, ...rest] = topic.startsWith(`${root}/`) ? topic.slice(root.length).split("/") : []; + return endpointId && rest.length === 1 && rest[0] === "alexaDirective" ? endpointId : null; +} + /** 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`; diff --git a/test/subscriptions.test.js b/test/subscriptions.test.js new file mode 100644 index 0000000..9e6f028 --- /dev/null +++ b/test/subscriptions.test.js @@ -0,0 +1,63 @@ +"use strict"; +// What the bridge subscribes to: /discover and /+/alexaDirective. 1.x subscribed to /#, which also +// delivered everything the bridge published itself and the backend's /alexaDirective_e. +const { test } = require("node:test"); +const assert = require("node:assert/strict"); +const { Alex2MQTT, DisplayCategory, PowerController, topics } = require("alex2node"); +const { broker, watcher, connected, directive, until, sleep } = require("./helpers/harness.js"); + +test("topics: the two filters, and the endpointId of a directive topic", () => { + assert.deepEqual(topics.subscriptions("root"), ["root/discover", "root/+/alexaDirective"]); + assert.equal(topics.directive("root", "lamp-1"), "root/lamp-1/alexaDirective"); + assert.equal(topics.directiveEndpoint("root", "root/lamp-1/alexaDirective"), "lamp-1"); + for (const other of [ + "root/lamp-1/alexaDirective_e", "root/lamp-1/alexaResponce", "root/alexaDirective", "root/a/b/alexaDirective", + "root//alexaDirective", "root2/lamp-1/alexaDirective", "root/discover", + ]) { + assert.equal(topics.directiveEndpoint("root", other), null, other); + } +}); + +test("the bridge subscribes to discover and +/alexaDirective: two endpoints get their directives, no other topic reaches the bridge", async () => { + const { aedes, url } = await broker(); + const subscribed = []; + aedes.on("subscribe", (subscriptions, client) => { + if (client.id === "bridge-under-test") subscribed.push(...subscriptions.map((s) => s.topic)); + }); + const alexa = await watcher(url, "root/#"); + const received = []; + const log = (line, detail) => { if (line === "MQTT Message Received") received.push(detail.topic); }; + const bridge = await connected(new Alex2MQTT("u", "p", "root", false, { host: url, log, mqtt: { clientId: "bridge-under-test" } })); + assert.deepEqual(subscribed.sort(), ["root/+/alexaDirective", "root/discover"]); + + const heard = []; + for (const id of ["lamp-1", "lamp-2"]) { + const lamp = bridge.registerDevice(id, id, DisplayCategory.LIGHT); + const answer = (d) => { + heard.push(`${id} ${d.header.name}`); + return lamp.getStatusMessage(d.header.correlationToken, true).addPowerControllerProp(PowerController.ON).send(); + }; + lamp.on("Event", answer); + lamp.on("ReportState", answer); + } + bridge.on("directive", ({ endpointId, name }) => heard.push(`bridge ${endpointId} ${name}`)); + + // Under the root, and none of them for the bridge: the last three are what the bridge and the backend publish + const turnOn = directive("Alexa.PowerController", "TurnOn", "lamp-1", "not-for-the-bridge"); + for (const topic of ["root/lamp-1/alexaDirective_e", "root/lamp-1/status", "root/lamp-1/alexaDirective/more", "root/alexaDirective"]) { + alexa.publish(topic, turnOn); + } + alexa.publish("root/lamp-1/alexaResponce", { event: {} }); + alexa.publish("root/changeReport", { event: {} }); + alexa.publish("root/discover_r", []); + alexa.send("root", directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + alexa.send("root", directive("Alexa.PowerController", "TurnOff", "lamp-2", "ct-2")); + alexa.publish("root/discover", ""); + + const answered = (id) => alexa.on(`root/${id}/alexaResponce`).some((m) => m.event.header); + await until(() => answered("lamp-1") && answered("lamp-2") && alexa.on("root/discover_r").length === 2, 3000, "both responses and the discovery answer"); + await sleep(100); // time for the echo of the bridge's own publishes, which must not come + + assert.deepEqual(heard, ["bridge lamp-1 TurnOn", "lamp-1 TurnOn", "bridge lamp-2 TurnOff", "lamp-2 TurnOff"]); + assert.deepEqual(received, ["root/lamp-1/alexaDirective", "root/lamp-2/alexaDirective", "root/discover"]); +});