diff --git a/dist/cjs/Alex2Node.js b/dist/cjs/Alex2Node.js index be918c3..0225dda 100644 --- a/dist/cjs/Alex2Node.js +++ b/dist/cjs/Alex2Node.js @@ -53,6 +53,20 @@ const DESCRIBED = [ "description", "manufacturerName", "manufacturer", "model", "serialNumber", "firmwareVersion", "softwareVersion", "customIdentifier", "cookie", ]; +// The messageId of a discovery request: in the header of the Discover directive, or beside the namespace and the +// name when the request is the header alone. undefined for a request without one, which is answered every time. +function requestId(text) { + let request; + try { + request = JSON.parse(text); + } + catch { + return undefined; + } + const header = request?.directive?.header ?? request?.header ?? request; + const { messageId } = (typeof header === "object" && header !== null ? header : {}); + return typeof messageId === "string" ? messageId : undefined; +} /** * The Alexa-to-MQTT bridge: one broker connection for a user's root topic, a set of registered devices, discovery and * directive dispatch. @@ -157,6 +171,11 @@ class Alex2MQTT extends events_1.EventEmitter { this.log(`MQTT Message Received`, { topic, payload: text }); try { if (topic === topics.discover(this.rootTopic)) { + const messageId = requestId(text); + if (messageId !== undefined && this.dispatcher.arrivedBefore("", messageId)) { + this.log(`the discovery request ${messageId} is dropped: it arrived before`); + return; + } await this.answerDiscovery(); return; } diff --git a/dist/cjs/dispatcher.js b/dist/cjs/dispatcher.js index bf626c8..4deccd4 100644 --- a/dist/cjs/dispatcher.js +++ b/dist/cjs/dispatcher.js @@ -163,6 +163,18 @@ class Dispatcher { this.deferralNoted.add(namespace); this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); } + /** + * true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message + * twice. endpointId is "" for a message to the root, a discovery request. + */ + arrivedBefore(endpointId, messageId, now = Date.now()) { + this.forget(now); + const id = key(endpointId, messageId); + if (this.arrived.has(id)) + return true; + this.arrived.set(id, now); + return false; + } /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ async dispatch(endpointId, text) { const { log, emit } = this.options; @@ -182,13 +194,9 @@ class Dispatcher { const { namespace, name, messageId } = raw.header; const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; const receivedAt = Date.now(); - this.forget(receivedAt); - if (typeof messageId === "string") { - if (this.arrived.has(key(endpointId, messageId))) { - log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); - return; - } - this.arrived.set(key(endpointId, messageId), receivedAt); + if (typeof messageId === "string" && this.arrivedBefore(endpointId, messageId, receivedAt)) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; } const device = this.options.getDevice(endpointId); if (!device) { diff --git a/dist/esm/Alex2Node.js b/dist/esm/Alex2Node.js index 3d8f175..5c0be0d 100644 --- a/dist/esm/Alex2Node.js +++ b/dist/esm/Alex2Node.js @@ -14,6 +14,20 @@ const DESCRIBED = [ "description", "manufacturerName", "manufacturer", "model", "serialNumber", "firmwareVersion", "softwareVersion", "customIdentifier", "cookie", ]; +// The messageId of a discovery request: in the header of the Discover directive, or beside the namespace and the +// name when the request is the header alone. undefined for a request without one, which is answered every time. +function requestId(text) { + let request; + try { + request = JSON.parse(text); + } + catch { + return undefined; + } + const header = request?.directive?.header ?? request?.header ?? request; + const { messageId } = (typeof header === "object" && header !== null ? header : {}); + return typeof messageId === "string" ? messageId : undefined; +} /** * The Alexa-to-MQTT bridge: one broker connection for a user's root topic, a set of registered devices, discovery and * directive dispatch. @@ -118,6 +132,11 @@ class Alex2MQTT extends EventEmitter { this.log(`MQTT Message Received`, { topic, payload: text }); try { if (topic === topics.discover(this.rootTopic)) { + const messageId = requestId(text); + if (messageId !== undefined && this.dispatcher.arrivedBefore("", messageId)) { + this.log(`the discovery request ${messageId} is dropped: it arrived before`); + return; + } await this.answerDiscovery(); return; } diff --git a/dist/esm/dispatcher.d.ts b/dist/esm/dispatcher.d.ts index 770aead..4354b47 100644 --- a/dist/esm/dispatcher.d.ts +++ b/dist/esm/dispatcher.d.ts @@ -93,6 +93,11 @@ export declare class Dispatcher { send(topic: string, message: object): Promise; /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ noteDeferral(namespace: string): void; + /** + * true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message + * twice. endpointId is "" for a message to the root, a discovery request. + */ + arrivedBefore(endpointId: string, messageId: string, now?: number): boolean; /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ dispatch(endpointId: string, text: string): Promise; private route; diff --git a/dist/esm/dispatcher.js b/dist/esm/dispatcher.js index 29d870e..4b1583b 100644 --- a/dist/esm/dispatcher.js +++ b/dist/esm/dispatcher.js @@ -127,6 +127,18 @@ export class Dispatcher { this.deferralNoted.add(namespace); this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); } + /** + * true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message + * twice. endpointId is "" for a message to the root, a discovery request. + */ + arrivedBefore(endpointId, messageId, now = Date.now()) { + this.forget(now); + const id = key(endpointId, messageId); + if (this.arrived.has(id)) + return true; + this.arrived.set(id, now); + return false; + } /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ async dispatch(endpointId, text) { const { log, emit } = this.options; @@ -146,13 +158,9 @@ export class Dispatcher { const { namespace, name, messageId } = raw.header; const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; const receivedAt = Date.now(); - this.forget(receivedAt); - if (typeof messageId === "string") { - if (this.arrived.has(key(endpointId, messageId))) { - log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); - return; - } - this.arrived.set(key(endpointId, messageId), receivedAt); + if (typeof messageId === "string" && this.arrivedBefore(endpointId, messageId, receivedAt)) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; } const device = this.options.getDevice(endpointId); if (!device) { diff --git a/dist/types/dispatcher.d.ts b/dist/types/dispatcher.d.ts index 770aead..4354b47 100644 --- a/dist/types/dispatcher.d.ts +++ b/dist/types/dispatcher.d.ts @@ -93,6 +93,11 @@ export declare class Dispatcher { send(topic: string, message: object): Promise; /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ noteDeferral(namespace: string): void; + /** + * true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message + * twice. endpointId is "" for a message to the root, a discovery request. + */ + arrivedBefore(endpointId: string, messageId: string, now?: number): boolean; /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ dispatch(endpointId: string, text: string): Promise; private route; diff --git a/src/Alex2Node.ts b/src/Alex2Node.ts index c7ced76..c08b4b1 100644 --- a/src/Alex2Node.ts +++ b/src/Alex2Node.ts @@ -50,6 +50,20 @@ const DESCRIBED = [ "customIdentifier", "cookie", ]; +// The messageId of a discovery request: in the header of the Discover directive, or beside the namespace and the +// name when the request is the header alone. undefined for a request without one, which is answered every time. +function requestId(text: string): string | undefined { + let request: { directive?: { header?: unknown }; header?: unknown } | null; + try { + request = JSON.parse(text); + } catch { + return undefined; + } + const header = request?.directive?.header ?? request?.header ?? request; + const { messageId } = (typeof header === "object" && header !== null ? header : {}) as { messageId?: unknown }; + return typeof messageId === "string" ? messageId : undefined; +} + /** * The Alexa-to-MQTT bridge: one broker connection for a user's root topic, a set of registered devices, discovery and * directive dispatch. @@ -160,6 +174,11 @@ class Alex2MQTT extends EventEmitter { this.log(`MQTT Message Received`, { topic, payload: text }); try { if (topic === topics.discover(this.rootTopic)) { + const messageId = requestId(text); + if (messageId !== undefined && this.dispatcher.arrivedBefore("", messageId)) { + this.log(`the discovery request ${messageId} is dropped: it arrived before`); + return; + } await this.answerDiscovery(); return; } diff --git a/src/dispatcher.ts b/src/dispatcher.ts index 0d8b22b..c851960 100644 --- a/src/dispatcher.ts +++ b/src/dispatcher.ts @@ -240,6 +240,18 @@ export class Dispatcher { this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); } + /** + * true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message + * twice. endpointId is "" for a message to the root, a discovery request. + */ + arrivedBefore(endpointId: string, messageId: string, now = Date.now()): boolean { + this.forget(now); + const id = key(endpointId, messageId); + if (this.arrived.has(id)) return true; + this.arrived.set(id, now); + return false; + } + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ async dispatch(endpointId: string, text: string): Promise { const { log, emit } = this.options; @@ -258,13 +270,9 @@ export class Dispatcher { const { namespace, name, messageId } = raw.header; const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; const receivedAt = Date.now(); - this.forget(receivedAt); - if (typeof messageId === "string") { - if (this.arrived.has(key(endpointId, messageId))) { - log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); - return; - } - this.arrived.set(key(endpointId, messageId), receivedAt); + if (typeof messageId === "string" && this.arrivedBefore(endpointId, messageId, receivedAt)) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; } const device = this.options.getDevice(endpointId); diff --git a/test/transport.test.js b/test/transport.test.js index 1e1bb92..c999145 100644 --- a/test/transport.test.js +++ b/test/transport.test.js @@ -89,3 +89,38 @@ test("a discovery answer that cannot be published: receive() resolves and the \" await bridge.receive("root/discover", ""); } }); + +test("a discovery request that arrives twice is answered once, by its messageId", async () => { + const sent = new MemoryPublisher(); + const logged = []; + const bridge = new Alex2MQTT("user", "password", "root", false, { publisher: sent, log: (message) => logged.push(message) }); + bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }).add(PowerController); + const header = (messageId) => ({ namespace: "Alexa.Discovery", name: "Discover", payloadVersion: "3", messageId }); + const answers = () => sent.published.filter(({ topic }) => topic === "root/discover_r").length; + + // The request as the Discover directive of Alexa, with its header alone, and the header by itself + const requests = [{ directive: { header: header("d-1"), payload: {} } }, { header: header("d-2") }, header("d-3")]; + for (const [i, request] of requests.entries()) { + await bridge.receive("root/discover", JSON.stringify(request)); + await bridge.receive("root/discover", JSON.stringify(request)); + assert.equal(answers(), i + 1, `d-${i + 1}`); + assert.ok(logged.includes(`the discovery request d-${i + 1} is dropped: it arrived before`)); + } + + // A request without a messageId cannot be told from the next one: each is answered + for (const request of ["", "{}", "not JSON", JSON.stringify({ namespace: "Alexa.Discovery", name: "Discover" })]) { + await bridge.receive("root/discover", request); + await bridge.receive("root/discover", request); + } + assert.equal(answers(), 3 + 8); + + // The messageId of a discovery request is not the one of a directive to an endpoint + const called = []; + bridge.getDevice("lamp-1").capability("Alexa.PowerController").on("TurnOn", (ctx) => { called.push(ctx.name); return ctx.respond(); }); + await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify({ + header: { namespace: "Alexa.PowerController", name: "TurnOn", messageId: "d-1", correlationToken: "ct", payloadVersion: "3" }, + endpoint: { endpointId: "lamp-1" }, + payload: {}, + })); + assert.deepEqual(called, ["TurnOn"]); +});