From e0fe9d2812b10784b1e32d0ca4794e45b7cc9313 Mon Sep 17 00:00:00 2001 From: David Date: Mon, 28 Sep 2026 19:36:09 +0000 Subject: [PATCH] bridge: subscribe to discover and +/alexaDirective only The bridge subscribed to /#, so the broker sent back every discover_r, alexaResponce and changeReport the bridge published, and the backend's /alexaDirective_e. It now subscribes to /discover and /+/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 --- dist/cjs/Alex2Node.js | 61 ++++++++++++++++++------------------ dist/cjs/topics.js | 20 +++++++++++- dist/esm/Alex2Node.d.ts | 2 +- dist/esm/Alex2Node.js | 61 ++++++++++++++++++------------------ dist/esm/topics.d.ts | 11 +++++++ dist/esm/topics.js | 14 +++++++++ dist/types/Alex2Node.d.ts | 2 +- dist/types/topics.d.ts | 11 +++++++ src/Alex2Node.ts | 55 ++++++++++++++++----------------- src/topics.ts | 18 +++++++++++ test/subscriptions.test.js | 63 ++++++++++++++++++++++++++++++++++++++ 11 files changed, 225 insertions(+), 93 deletions(-) create mode 100644 test/subscriptions.test.js 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"]); +});