From 6789a1a077accad000981ef776d1fe7063791d79 Mon Sep 17 00:00:00 2001 From: David Date: Mon, 28 Sep 2026 19:34:39 +0000 Subject: [PATCH] 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 --- dist/cjs/Alex2Node.js | 70 +++++++++++++----- dist/cjs/compat/AlexaErrorResponse.js | 51 ++++++++++--- dist/cjs/compat/AlexaStatusMessage.js | 60 ++++++++++----- dist/cjs/device/Device.js | 69 +++++++++++++----- dist/cjs/index.js | 6 +- dist/cjs/topics.js | 17 +++++ dist/cjs/transport.js | 83 +++++++++++++++++++++ dist/esm/Alex2Node.d.ts | 5 +- dist/esm/Alex2Node.js | 37 +++++----- dist/esm/compat/AlexaErrorResponse.d.ts | 6 +- dist/esm/compat/AlexaErrorResponse.js | 18 ++--- dist/esm/compat/AlexaStatusMessage.d.ts | 6 +- dist/esm/compat/AlexaStatusMessage.js | 27 ++----- dist/esm/device/Device.d.ts | 14 ++-- dist/esm/device/Device.js | 36 ++++----- dist/esm/index.d.ts | 3 + dist/esm/index.js | 3 + dist/esm/topics.d.ts | 8 ++ dist/esm/topics.js | 10 +++ dist/esm/transport.d.ts | 48 ++++++++++++ dist/esm/transport.js | 77 ++++++++++++++++++++ dist/types/Alex2Node.d.ts | 5 +- dist/types/compat/AlexaErrorResponse.d.ts | 6 +- dist/types/compat/AlexaStatusMessage.d.ts | 6 +- dist/types/device/Device.d.ts | 14 ++-- dist/types/index.d.ts | 3 + dist/types/topics.d.ts | 8 ++ dist/types/transport.d.ts | 48 ++++++++++++ src/Alex2Node.ts | 34 +++++---- src/compat/AlexaErrorResponse.ts | 16 ++-- src/compat/AlexaStatusMessage.ts | 25 ++----- src/device/Device.ts | 36 +++++---- src/index.ts | 5 ++ src/topics.ts | 14 ++++ src/transport.ts | 89 +++++++++++++++++++++++ test/compat/messages.test.js | 24 +++--- test/connection.test.js | 31 ++++++-- test/helpers/endpoint.js | 2 +- test/transport.test.js | 65 +++++++++++++++++ 39 files changed, 839 insertions(+), 246 deletions(-) create mode 100644 dist/cjs/topics.js create mode 100644 dist/cjs/transport.js create mode 100644 dist/esm/topics.d.ts create mode 100644 dist/esm/topics.js create mode 100644 dist/esm/transport.d.ts create mode 100644 dist/esm/transport.js create mode 100644 dist/types/topics.d.ts create mode 100644 dist/types/transport.d.ts create mode 100644 src/topics.ts create mode 100644 src/transport.ts create mode 100644 test/transport.test.js diff --git a/dist/cjs/Alex2Node.js b/dist/cjs/Alex2Node.js index f060321..330a918 100644 --- a/dist/cjs/Alex2Node.js +++ b/dist/cjs/Alex2Node.js @@ -1,4 +1,37 @@ "use strict"; +var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + var desc = Object.getOwnPropertyDescriptor(m, k); + if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) { + desc = { enumerable: true, get: function() { return m[k]; } }; + } + Object.defineProperty(o, k2, desc); +}) : (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + o[k2] = m[k]; +})); +var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) { + Object.defineProperty(o, "default", { enumerable: true, value: v }); +}) : function(o, v) { + o["default"] = v; +}); +var __importStar = (this && this.__importStar) || (function () { + var ownKeys = function(o) { + ownKeys = Object.getOwnPropertyNames || function (o) { + var ar = []; + for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k; + return ar; + }; + return ownKeys(o); + }; + return function (mod) { + if (mod && mod.__esModule) return mod; + var result = {}; + if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]); + __setModuleDefault(result, mod); + return result; + }; +})(); var __importDefault = (this && this.__importDefault) || function (mod) { return (mod && mod.__esModule) ? mod : { "default": mod }; }; @@ -11,6 +44,8 @@ const validate_js_1 = require("./device/validate.js"); const events_1 = require("events"); const enums_js_1 = require("./compat/enums.js"); const types_js_1 = require("./registry/types.js"); +const topics = __importStar(require("./topics.js")); +const transport_js_1 = require("./transport.js"); exports.DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; // What addDevice() copies from the definition to the device const DESCRIBED = [ @@ -39,6 +74,8 @@ class Alex2MQTT extends events_1.EventEmitter { this.rootTopic = rootTopic; this.debugLogging = debugLogging; this.client = null; + // One for the life of the bridge: the devices keep it over disconnect() and connect() + this.publisher = new transport_js_1.MqttPublisher(() => this.client); this.devices = []; // The lines of check() that were logged. Discovery comes every few minutes: a line is logged once. this.logged = new Set(); @@ -76,8 +113,6 @@ class Alex2MQTT extends events_1.EventEmitter { ...(this.options.mqtt || {}), }; this.client = mqtt_1.default.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; this.log("Connected to MQTT broker"); @@ -100,12 +135,12 @@ class Alex2MQTT extends events_1.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); + 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 ${topic + "_r"}`, deviceArray); + this.log(`Discovery payloads published to ${result.topic}`, deviceArray); this.emit("discover", deviceArray.length); }); } @@ -179,15 +214,13 @@ class Alex2MQTT extends events_1.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) { - 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_js_1.default(this.client, this.rootTopic, name, endpointId, (categories ?? [])); + const device = new Device_js_1.default(null, this.rootTopic, name, endpointId, (categories ?? [])); for (const [field, value] of Object.entries(described)) { if (!DESCRIBED.includes(field)) throw new types_js_1.DeclarationError({ endpointId }, `${field} is not a field of an endpoint`); @@ -202,9 +235,6 @@ class Alex2MQTT extends events_1.EventEmitter { return device; } registerDevice(name, endpointId, displayCategory) { - 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`; @@ -216,7 +246,7 @@ class Alex2MQTT extends events_1.EventEmitter { } this.log(`Creating new device with endpoint: ${endpointId}`); const normalizedCategory = Array.isArray(displayCategory) ? displayCategory : [displayCategory || enums_js_1.DisplayCategory.LIGHT]; - const device = new Device_js_1.default(this.client, this.rootTopic, name, endpointId, normalizedCategory); + const device = new Device_js_1.default(null, this.rootTopic, name, endpointId, normalizedCategory); this.register(device, false); return device; } @@ -224,6 +254,7 @@ class Alex2MQTT extends events_1.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); } /** Forget a device (its listeners with it). Returns false when there was none. */ @@ -231,16 +262,21 @@ class Alex2MQTT extends events_1.EventEmitter { 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() { for (const d of this.devices) - d.removeAllListeners(); + this.release(d); this.devices = []; } + // A device the bridge forgot hears no directive and publishes nothing + release(device) { + device.removeAllListeners(); + device.publisher = null; + } getDevices() { return this.devices.slice(); } diff --git a/dist/cjs/compat/AlexaErrorResponse.js b/dist/cjs/compat/AlexaErrorResponse.js index 25a5d8e..4fb9359 100644 --- a/dist/cjs/compat/AlexaErrorResponse.js +++ b/dist/cjs/compat/AlexaErrorResponse.js @@ -1,16 +1,51 @@ "use strict"; +var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + var desc = Object.getOwnPropertyDescriptor(m, k); + if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) { + desc = { enumerable: true, get: function() { return m[k]; } }; + } + Object.defineProperty(o, k2, desc); +}) : (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + o[k2] = m[k]; +})); +var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) { + Object.defineProperty(o, "default", { enumerable: true, value: v }); +}) : function(o, v) { + o["default"] = v; +}); +var __importStar = (this && this.__importStar) || (function () { + var ownKeys = function(o) { + ownKeys = Object.getOwnPropertyNames || function (o) { + var ar = []; + for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k; + return ar; + }; + return ownKeys(o); + }; + return function (mod) { + if (mod && mod.__esModule) return mod; + var result = {}; + if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]); + __setModuleDefault(result, mod); + return result; + }; +})(); Object.defineProperty(exports, "__esModule", { value: true }); exports.AlexaErrorResponse = void 0; const crypto_1 = require("crypto"); const build_js_1 = require("../messages/build.js"); const errors_js_1 = require("../messages/errors.js"); +const topics = __importStar(require("../topics.js")); +const transport_js_1 = require("../transport.js"); /** An ErrorResponse as 1.x builds it: device.getErrorMessage(token), setErrorMessage(), send(). */ class AlexaErrorResponse { - constructor(correlationToken, rootTopic, endpointId, mqttClient) { + constructor(correlationToken, rootTopic, endpointId, publisher) { this.correlationToken = correlationToken; this.rootTopic = rootTopic; this.endpointId = endpointId; - this.mqttClient = mqttClient; + this.publisher = publisher; // Sent when setErrorMessage() was not called: 1.x published an empty payload, which is no ErrorResponse to Alexa this.error = errors_js_1.AlexaErrors.of("INTERNAL_ERROR", "the device answered with an error and did not say which"); // One id for the message, however often toJSON() is called @@ -39,16 +74,8 @@ class AlexaErrorResponse { } /** Resolves with the topic published to, or "" when the publish failed. Never rejects. */ send(sendAsync = false) { - 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 (0, transport_js_1.send)(this.publisher, answer(this.rootTopic, this.endpointId), () => this.toJSON(), this.onPublishError); } } exports.AlexaErrorResponse = AlexaErrorResponse; diff --git a/dist/cjs/compat/AlexaStatusMessage.js b/dist/cjs/compat/AlexaStatusMessage.js index e31382a..00929fa 100644 --- a/dist/cjs/compat/AlexaStatusMessage.js +++ b/dist/cjs/compat/AlexaStatusMessage.js @@ -1,9 +1,44 @@ "use strict"; +var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + var desc = Object.getOwnPropertyDescriptor(m, k); + if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) { + desc = { enumerable: true, get: function() { return m[k]; } }; + } + Object.defineProperty(o, k2, desc); +}) : (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + o[k2] = m[k]; +})); +var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) { + Object.defineProperty(o, "default", { enumerable: true, value: v }); +}) : function(o, v) { + o["default"] = v; +}); +var __importStar = (this && this.__importStar) || (function () { + var ownKeys = function(o) { + ownKeys = Object.getOwnPropertyNames || function (o) { + var ar = []; + for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k; + return ar; + }; + return ownKeys(o); + }; + return function (mod) { + if (mod && mod.__esModule) return mod; + var result = {}; + if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]); + __setModuleDefault(result, mod); + return result; + }; +})(); Object.defineProperty(exports, "__esModule", { value: true }); exports.AlexaStatusMessage = void 0; const crypto_1 = require("crypto"); const build_js_1 = require("../messages/build.js"); const StateBuilder_js_1 = require("../messages/StateBuilder.js"); +const topics = __importStar(require("../topics.js")); +const transport_js_1 = require("../transport.js"); const enums_js_1 = require("./enums.js"); // The 1.x helpers report every temperature in Celsius, whatever scale the device works in function celsius(scale, value) { @@ -14,11 +49,11 @@ function celsius(scale, value) { * device.getChangeReport(), the add...Prop() calls, send(). The values are not checked. */ class AlexaStatusMessage { - constructor(correlationToken, rootTopic, endpointId, mqttClient, isResponse = false, isDeferred = false, changeCause = null) { + constructor(correlationToken, rootTopic, endpointId, publisher, isResponse = false, isDeferred = false, changeCause = null) { this.correlationToken = correlationToken; this.rootTopic = rootTopic; this.endpointId = endpointId; - this.mqttClient = mqttClient; + this.publisher = publisher; this.isResponse = isResponse; this.isDeferred = isDeferred; this.changeCause = changeCause; @@ -109,24 +144,9 @@ class AlexaStatusMessage { * to). Never rejects (1.5.1 did, so an un-caught send() could kill the host). */ send(sendAsync = false) { - 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) => { - 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; - try { - payload = JSON.stringify(this.toJSON()); - } - catch (err) { - return failed(err); - } - 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 (0, transport_js_1.send)(this.publisher, topic, () => this.toJSON(), this.onPublishError); } } exports.AlexaStatusMessage = AlexaStatusMessage; diff --git a/dist/cjs/device/Device.js b/dist/cjs/device/Device.js index dacfd63..4f5485c 100644 --- a/dist/cjs/device/Device.js +++ b/dist/cjs/device/Device.js @@ -1,4 +1,37 @@ "use strict"; +var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + var desc = Object.getOwnPropertyDescriptor(m, k); + if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) { + desc = { enumerable: true, get: function() { return m[k]; } }; + } + Object.defineProperty(o, k2, desc); +}) : (function(o, m, k, k2) { + if (k2 === undefined) k2 = k; + o[k2] = m[k]; +})); +var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) { + Object.defineProperty(o, "default", { enumerable: true, value: v }); +}) : function(o, v) { + o["default"] = v; +}); +var __importStar = (this && this.__importStar) || (function () { + var ownKeys = function(o) { + ownKeys = Object.getOwnPropertyNames || function (o) { + var ar = []; + for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k; + return ar; + }; + return ownKeys(o); + }; + return function (mod) { + if (mod && mod.__esModule) return mod; + var result = {}; + if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]); + __setModuleDefault(result, mod); + return result; + }; +})(); Object.defineProperty(exports, "__esModule", { value: true }); const events_1 = require("events"); const AlexaErrorResponse_js_1 = require("../compat/AlexaErrorResponse.js"); @@ -10,12 +43,14 @@ const Alexa_js_1 = require("../registry/interfaces/Alexa.js"); const EndpointHealth_js_1 = require("../registry/interfaces/EndpointHealth.js"); const schema_js_1 = require("../registry/schema.js"); const types_js_1 = require("../registry/types.js"); +const topics = __importStar(require("../topics.js")); +const transport_js_1 = require("../transport.js"); const Capability_js_1 = require("./Capability.js"); const validate_js_1 = require("./validate.js"); class Device extends events_1.EventEmitter { - constructor(mqttClient, rootTopic, name, endpointId, displayCategory, description = "Alexa to Node.js bridge", manufacturerName = "Alex2Node", manufacturer = "Alex2Node", model = "Alex2Node_v1.0.0") { + /** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */ + constructor(client, rootTopic, name, endpointId, displayCategory, description = "Alexa to Node.js bridge", manufacturerName = "Alex2Node", manufacturer = "Alex2Node", model = "Alex2Node_v1.0.0") { super(); - this.mqttClient = mqttClient; this.rootTopic = rootTopic; this.name = name; this.endpointId = endpointId; @@ -39,10 +74,12 @@ class Device extends events_1.EventEmitter { */ this.endpointHealth = false; this.capabilities = []; - } - /** Internal (1.5.2): Alex2MQTT.connect() re-binds every registered device to its new broker client after a disconnect(). */ - setMqttClient(client) { - this.mqttClient = client; + /** + * 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. + */ + this.publisher = null; } getName() { return this.name; @@ -66,12 +103,12 @@ class Device extends events_1.EventEmitter { return this.description; } getErrorMessage(correlationToken) { - const msg = new AlexaErrorResponse_js_1.AlexaErrorResponse(correlationToken, this.rootTopic, this.endpointId, this.mqttClient); + const msg = new AlexaErrorResponse_js_1.AlexaErrorResponse(correlationToken, this.rootTopic, this.endpointId, this.publisher); msg.onPublishError = this.onPublishError; return msg; } getStatusMessage(correlationToken, isResponse = false, isDeferred = false) { - const msg = new AlexaStatusMessage_js_1.AlexaStatusMessage(correlationToken, this.rootTopic, this.endpointId, this.mqttClient, isResponse, isDeferred); + const msg = new AlexaStatusMessage_js_1.AlexaStatusMessage(correlationToken, this.rootTopic, this.endpointId, this.publisher, isResponse, isDeferred); msg.onPublishError = this.onPublishError; return msg; } @@ -82,7 +119,7 @@ class Device extends events_1.EventEmitter { * Without a changed property send() publishes nothing and resolves with "": Alex2MQTT would drop the report. */ getChangeReport(cause = "PHYSICAL_INTERACTION") { - const msg = new AlexaStatusMessage_js_1.AlexaStatusMessage("", this.rootTopic, this.endpointId, this.mqttClient, false, false, cause); + const msg = new AlexaStatusMessage_js_1.AlexaStatusMessage("", this.rootTopic, this.endpointId, this.publisher, false, false, cause); msg.onPublishError = this.onPublishError; return msg; } @@ -91,17 +128,9 @@ class Device extends events_1.EventEmitter { * (1.5.1). Resolves with the topic published to, or "" when the publish failed (never rejects, 1.5.2). */ sendSceneResponse(correlationToken, activated, cause = "VOICE_INTERACTION", sendAsync = false) { - const payload = (0, build_js_1.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 = () => (0, build_js_1.sceneEvent)({ endpointId: this.endpointId, correlationToken, activated, cause }); + return (0, transport_js_1.send)(this.publisher, answer(this.rootTopic, this.endpointId), build, this.onPublishError); } /** The capabilities the device declared, in the order it declared them. */ getCapabilities() { diff --git a/dist/cjs/index.js b/dist/cjs/index.js index 774c919..c3fcd76 100644 --- a/dist/cjs/index.js +++ b/dist/cjs/index.js @@ -36,7 +36,7 @@ var __importDefault = (this && this.__importDefault) || function (mod) { return (mod && mod.__esModule) ? mod : { "default": mod }; }; Object.defineProperty(exports, "__esModule", { value: true }); -exports.AlexaErrorResponse = exports.AlexaStatusMessage = exports.ThermostatMode = exports.TemperatureSensorScale = exports.PowerState = exports.DisplayCategory = exports.AlexaInterfaceType = exports.AlexaErrorType = exports.AlexaActions = exports.ActionMapping = exports.AlexaInterface = exports.property = exports.StateBuilder = exports.MessageError = exports.AlexaErrors = exports.AlexaError = exports.messages = exports.DisplayCategories = exports.States = exports.Actions = exports.Units = exports.Assets = exports.SemanticsBuilder = exports.semantics = exports.text = exports.asset = exports.EndpointHealth = exports.PowerController = exports.ToggleController = exports.TemperatureSensor = exports.RangeController = exports.ModeController = exports.BrightnessController = exports.Alexa = exports.SchemaError = exports.DeclarationError = exports.registry = exports.Capability = exports.Device = exports.DEFAULT_HOST = exports.Alex2MQTT = void 0; +exports.AlexaErrorResponse = exports.AlexaStatusMessage = exports.ThermostatMode = exports.TemperatureSensorScale = exports.PowerState = exports.DisplayCategory = exports.AlexaInterfaceType = exports.AlexaErrorType = exports.AlexaActions = exports.ActionMapping = exports.AlexaInterface = exports.MemoryPublisher = exports.topics = exports.property = exports.StateBuilder = exports.MessageError = exports.AlexaErrors = exports.AlexaError = exports.messages = exports.DisplayCategories = exports.States = exports.Actions = exports.Units = exports.Assets = exports.SemanticsBuilder = exports.semantics = exports.text = exports.asset = exports.EndpointHealth = exports.PowerController = exports.ToggleController = exports.TemperatureSensor = exports.RangeController = exports.ModeController = exports.BrightnessController = exports.Alexa = exports.SchemaError = exports.DeclarationError = exports.registry = exports.Capability = exports.Device = exports.DEFAULT_HOST = exports.Alex2MQTT = void 0; // Relative specifiers carry ".js": Node's ES module loader resolves no extension, and TypeScript maps it back to the .ts. var Alex2Node_js_1 = require("./Alex2Node.js"); Object.defineProperty(exports, "Alex2MQTT", { enumerable: true, get: function () { return __importDefault(Alex2Node_js_1).default; } }); @@ -80,6 +80,10 @@ Object.defineProperty(exports, "AlexaErrors", { enumerable: true, get: function Object.defineProperty(exports, "MessageError", { enumerable: true, get: function () { return index_js_5.MessageError; } }); Object.defineProperty(exports, "StateBuilder", { enumerable: true, get: function () { return index_js_5.StateBuilder; } }); Object.defineProperty(exports, "property", { enumerable: true, get: function () { return index_js_5.property; } }); +// Publishing: the topics of the Alex2MQTT contract, and a publisher that needs no broker for tests +exports.topics = __importStar(require("./topics.js")); +var transport_js_1 = require("./transport.js"); +Object.defineProperty(exports, "MemoryPublisher", { enumerable: true, get: function () { return transport_js_1.MemoryPublisher; } }); // 1.x var AlexaInterface_js_1 = require("./compat/AlexaInterface.js"); Object.defineProperty(exports, "AlexaInterface", { enumerable: true, get: function () { return AlexaInterface_js_1.AlexaInterface; } }); diff --git a/dist/cjs/topics.js b/dist/cjs/topics.js new file mode 100644 index 0000000..45f2d53 --- /dev/null +++ b/dist/cjs/topics.js @@ -0,0 +1,17 @@ +"use strict"; +// 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; +/** Where the bridge answers a discovery request with its endpoints. */ +const discoverReply = (root) => `${root}/discover_r`; +exports.discoverReply = discoverReply; +/** 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; +/** Where the answer goes once a DeferredResponse was sent for the directive. */ +const deferred = (root, endpointId) => `${root}/${endpointId}/deferredResponse`; +exports.deferred = deferred; +/** Where every ChangeReport of the root goes: the backend adds the user's token and posts it to Alexa. */ +const changeReport = (root) => `${root}/changeReport`; +exports.changeReport = changeReport; diff --git a/dist/cjs/transport.js b/dist/cjs/transport.js new file mode 100644 index 0000000..0d97477 --- /dev/null +++ b/dist/cjs/transport.js @@ -0,0 +1,83 @@ +"use strict"; +Object.defineProperty(exports, "__esModule", { value: true }); +exports.MemoryPublisher = exports.MqttPublisher = void 0; +exports.send = send; +const asError = (err) => (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. + */ +class MqttPublisher { + constructor(client) { + this.client = client; + } + publish(topic, message) { + const client = this.client(); + return new Promise((resolve) => { + const failed = (err) => 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); + } + }); + } +} +exports.MqttPublisher = MqttPublisher; +/** + * 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: "//alexaResponce", message: { event, context } } + */ +class MemoryPublisher { + constructor() { + /** Oldest first. The message is what a subscriber gets after JSON.parse. */ + this.published = []; + /** Set it and every publish fails with this error, as a publish to a broker that is gone does. */ + this.failWith = null; + } + publish(topic, message) { + 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 }); + } +} +exports.MemoryPublisher = MemoryPublisher; +/** + * 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. + */ +async function send(publisher, topic, build, report) { + let result; + 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 ""; +} diff --git a/dist/esm/Alex2Node.d.ts b/dist/esm/Alex2Node.d.ts index 983e964..e7d45dc 100644 --- a/dist/esm/Alex2Node.d.ts +++ b/dist/esm/Alex2Node.d.ts @@ -38,6 +38,7 @@ declare class Alex2MQTT extends EventEmitter { private rootTopic; private debugLogging; private client; + private readonly publisher; private devices; private MqttHost; private options; @@ -62,7 +63,8 @@ declare 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; registerDevice(name: string, endpointId: string, displayCategory?: DisplayCategory | DisplayCategory[] | null): Device; @@ -71,6 +73,7 @@ declare class Alex2MQTT extends EventEmitter { unregisterDevice(endpointId: string): boolean; /** Forget every device. */ clearDevices(): void; + private release; getDevices(): Device[]; getDevice(endpointId: string): Device | undefined; getRootTopic(): string; diff --git a/dist/esm/Alex2Node.js b/dist/esm/Alex2Node.js index 6838e74..0be84a6 100644 --- a/dist/esm/Alex2Node.js +++ b/dist/esm/Alex2Node.js @@ -5,6 +5,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"; export const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; // What addDevice() copies from the definition to the device const DESCRIBED = [ @@ -33,6 +35,8 @@ class Alex2MQTT extends EventEmitter { this.rootTopic = rootTopic; this.debugLogging = debugLogging; this.client = null; + // One for the life of the bridge: the devices keep it over disconnect() and connect() + this.publisher = new MqttPublisher(() => this.client); this.devices = []; // The lines of check() that were logged. Discovery comes every few minutes: a line is logged once. this.logged = new Set(); @@ -70,8 +74,6 @@ class Alex2MQTT extends EventEmitter { ...(this.options.mqtt || {}), }; 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; this.log("Connected to MQTT broker"); @@ -94,12 +96,12 @@ 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); + 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 ${topic + "_r"}`, deviceArray); + this.log(`Discovery payloads published to ${result.topic}`, deviceArray); this.emit("discover", deviceArray.length); }); } @@ -173,15 +175,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) { - 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 ?? [])); + const device = new Device(null, this.rootTopic, name, endpointId, (categories ?? [])); for (const [field, value] of Object.entries(described)) { if (!DESCRIBED.includes(field)) throw new DeclarationError({ endpointId }, `${field} is not a field of an endpoint`); @@ -196,9 +196,6 @@ class Alex2MQTT extends EventEmitter { return device; } registerDevice(name, endpointId, displayCategory) { - 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`; @@ -210,7 +207,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; } @@ -218,6 +215,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); } /** Forget a device (its listeners with it). Returns false when there was none. */ @@ -225,16 +223,21 @@ class Alex2MQTT extends EventEmitter { 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() { for (const d of this.devices) - d.removeAllListeners(); + this.release(d); this.devices = []; } + // A device the bridge forgot hears no directive and publishes nothing + release(device) { + device.removeAllListeners(); + device.publisher = null; + } getDevices() { return this.devices.slice(); } diff --git a/dist/esm/compat/AlexaErrorResponse.d.ts b/dist/esm/compat/AlexaErrorResponse.d.ts index 389b6cc..696f1ea 100644 --- a/dist/esm/compat/AlexaErrorResponse.d.ts +++ b/dist/esm/compat/AlexaErrorResponse.d.ts @@ -1,17 +1,17 @@ -import type { MqttClient } from "mqtt"; import type { ErrorResponseMessage } from "../messages/types.js"; +import type { Publisher } from "../transport.js"; /** An ErrorResponse as 1.x builds it: device.getErrorMessage(token), setErrorMessage(), send(). */ export declare class AlexaErrorResponse { private readonly correlationToken; private readonly rootTopic; private readonly endpointId; - private readonly mqttClient; + private readonly publisher; private error; private namespace?; private readonly messageId; /** Where a failed publish is reported (set by the Device that built this message, 1.5.2): send() never rejects. */ onPublishError?: (err: Error) => void; - constructor(correlationToken: string, rootTopic: string, endpointId: string, mqttClient: MqttClient); + constructor(correlationToken: string, rootTopic: string, endpointId: string, publisher: Publisher | null); /** * The error to answer with. otherParams are the fields the type adds to the payload (validRange). The header * carries the namespace the type is documented under, Alexa.ThermostatController for THERMOSTAT_IS_OFF; diff --git a/dist/esm/compat/AlexaErrorResponse.js b/dist/esm/compat/AlexaErrorResponse.js index 84cb1eb..a91cfb5 100644 --- a/dist/esm/compat/AlexaErrorResponse.js +++ b/dist/esm/compat/AlexaErrorResponse.js @@ -1,13 +1,15 @@ import { randomUUID } from "crypto"; import { errorResponse } from "../messages/build.js"; import { AlexaErrors } from "../messages/errors.js"; +import * as topics from "../topics.js"; +import { send } from "../transport.js"; /** An ErrorResponse as 1.x builds it: device.getErrorMessage(token), setErrorMessage(), send(). */ export class AlexaErrorResponse { - constructor(correlationToken, rootTopic, endpointId, mqttClient) { + constructor(correlationToken, rootTopic, endpointId, publisher) { this.correlationToken = correlationToken; this.rootTopic = rootTopic; this.endpointId = endpointId; - this.mqttClient = mqttClient; + this.publisher = publisher; // Sent when setErrorMessage() was not called: 1.x published an empty payload, which is no ErrorResponse to Alexa this.error = AlexaErrors.of("INTERNAL_ERROR", "the device answered with an error and did not say which"); // One id for the message, however often toJSON() is called @@ -36,15 +38,7 @@ export class AlexaErrorResponse { } /** Resolves with the topic published to, or "" when the publish failed. Never rejects. */ send(sendAsync = false) { - 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); } } diff --git a/dist/esm/compat/AlexaStatusMessage.d.ts b/dist/esm/compat/AlexaStatusMessage.d.ts index f4dee17..2f9e51a 100644 --- a/dist/esm/compat/AlexaStatusMessage.d.ts +++ b/dist/esm/compat/AlexaStatusMessage.d.ts @@ -1,5 +1,5 @@ -import type { MqttClient } from "mqtt"; import type { ChangeCause, Property } from "../messages/types.js"; +import type { Publisher } from "../transport.js"; import { TemperatureSensorScale } from "./enums.js"; import type { EndpointHealth, PowerController } from "./enums.js"; /** @@ -10,7 +10,7 @@ export declare class AlexaStatusMessage { private readonly correlationToken; private readonly rootTopic; private readonly endpointId; - private readonly mqttClient; + private readonly publisher; private readonly isResponse; private readonly isDeferred; private readonly changeCause; @@ -19,7 +19,7 @@ export declare class AlexaStatusMessage { private estimatedDeferralInSeconds?; /** Where a failed publish is reported (set by the Device that built this message, 1.5.2): send() never rejects. */ onPublishError?: (err: Error) => void; - constructor(correlationToken: string, rootTopic: string, endpointId: string, mqttClient: MqttClient, isResponse?: boolean, isDeferred?: boolean, changeCause?: ChangeCause | null); + constructor(correlationToken: string, rootTopic: string, endpointId: string, publisher: Publisher | null, isResponse?: boolean, isDeferred?: boolean, changeCause?: ChangeCause | null); private addProperty; /** ChangeReport: the add*Prop calls that follow describe what CHANGED (the default for a change report). */ changed(): this; diff --git a/dist/esm/compat/AlexaStatusMessage.js b/dist/esm/compat/AlexaStatusMessage.js index 75e1f4c..e84e73a 100644 --- a/dist/esm/compat/AlexaStatusMessage.js +++ b/dist/esm/compat/AlexaStatusMessage.js @@ -1,6 +1,8 @@ import { randomUUID } from "crypto"; import { changeReport, deferredResponse, response } from "../messages/build.js"; import { StateBuilder } from "../messages/StateBuilder.js"; +import * as topics from "../topics.js"; +import { send } from "../transport.js"; import { AlexaInterfaceType, TemperatureSensorScale } from "./enums.js"; // The 1.x helpers report every temperature in Celsius, whatever scale the device works in function celsius(scale, value) { @@ -11,11 +13,11 @@ function celsius(scale, value) { * device.getChangeReport(), the add...Prop() calls, send(). The values are not checked. */ export class AlexaStatusMessage { - constructor(correlationToken, rootTopic, endpointId, mqttClient, isResponse = false, isDeferred = false, changeCause = null) { + constructor(correlationToken, rootTopic, endpointId, publisher, isResponse = false, isDeferred = false, changeCause = null) { this.correlationToken = correlationToken; this.rootTopic = rootTopic; this.endpointId = endpointId; - this.mqttClient = mqttClient; + this.publisher = publisher; this.isResponse = isResponse; this.isDeferred = isDeferred; this.changeCause = changeCause; @@ -106,23 +108,8 @@ export class AlexaStatusMessage { * to). Never rejects (1.5.1 did, so an un-caught send() could kill the host). */ send(sendAsync = false) { - 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) => { - 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; - try { - payload = JSON.stringify(this.toJSON()); - } - catch (err) { - return failed(err); - } - 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); } } diff --git a/dist/esm/device/Device.d.ts b/dist/esm/device/Device.d.ts index f404968..5f0d570 100644 --- a/dist/esm/device/Device.d.ts +++ b/dist/esm/device/Device.d.ts @@ -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"; @@ -8,6 +7,7 @@ import type { AlexaInterfaceType } from "../compat/enums.js"; import type { ChangeCause } from "../messages/types.js"; import type { DisplayCategoryName } from "../registry/catalog.js"; import type { Directives, InterfaceDescriptor, Properties } from "../registry/types.js"; +import type { Publisher } from "../transport.js"; import { Capability } from "./Capability.js"; import type { AnyCapability, CapabilityJson, Declaration } from "./Capability.js"; import type { EndpointFields } from "./validate.js"; @@ -40,7 +40,6 @@ export interface EndpointJson extends EndpointFields { } type DeclarationArguments = {} extends Declaration ? [options?: Declaration] : [options: Declaration]; declare class Device extends EventEmitter { - private mqttClient; private rootTopic; name: string; endpointId: string; @@ -67,9 +66,14 @@ declare class Device extends EventEmitter { private capabilities; /** Where a failed publish from this device or a message it built is reported; Alex2MQTT sets it (1.5.2). */ onPublishError?: (err: Error) => void; - constructor(mqttClient: MqttClient, rootTopic: string, name: string, endpointId: string, displayCategory: Array | null, description?: string, manufacturerName?: string, manufacturer?: string, model?: string); - /** Internal (1.5.2): Alex2MQTT.connect() re-binds every registered device to its new broker client after a disconnect(). */ - setMqttClient(client: MqttClient): 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. + */ + publisher: Publisher | null; + /** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */ + constructor(client: unknown, rootTopic: string, name: string, endpointId: string, displayCategory: Array | null, description?: string, manufacturerName?: string, manufacturer?: string, model?: string); getName(): string; setName(name: string): void; getEndpointId(): string; diff --git a/dist/esm/device/Device.js b/dist/esm/device/Device.js index 87f334c..8137755 100644 --- a/dist/esm/device/Device.js +++ b/dist/esm/device/Device.js @@ -8,12 +8,14 @@ import { Alexa } from "../registry/interfaces/Alexa.js"; import { EndpointHealth } from "../registry/interfaces/EndpointHealth.js"; import { SchemaError } from "../registry/schema.js"; import { DeclarationError } from "../registry/types.js"; +import * as topics from "../topics.js"; +import { send } from "../transport.js"; import { Capability, commonOptions } from "./Capability.js"; import { checkCapability, checkCapabilityCount, checkEndpoint } from "./validate.js"; class Device extends EventEmitter { - constructor(mqttClient, rootTopic, name, endpointId, displayCategory, description = "Alexa to Node.js bridge", manufacturerName = "Alex2Node", manufacturer = "Alex2Node", model = "Alex2Node_v1.0.0") { + /** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */ + constructor(client, rootTopic, name, endpointId, displayCategory, description = "Alexa to Node.js bridge", manufacturerName = "Alex2Node", manufacturer = "Alex2Node", model = "Alex2Node_v1.0.0") { super(); - this.mqttClient = mqttClient; this.rootTopic = rootTopic; this.name = name; this.endpointId = endpointId; @@ -37,10 +39,12 @@ class Device extends EventEmitter { */ this.endpointHealth = false; this.capabilities = []; - } - /** Internal (1.5.2): Alex2MQTT.connect() re-binds every registered device to its new broker client after a disconnect(). */ - setMqttClient(client) { - this.mqttClient = client; + /** + * 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. + */ + this.publisher = null; } getName() { return this.name; @@ -64,12 +68,12 @@ class Device extends EventEmitter { return this.description; } getErrorMessage(correlationToken) { - const msg = new AlexaErrorResponse(correlationToken, this.rootTopic, this.endpointId, this.mqttClient); + const msg = new AlexaErrorResponse(correlationToken, this.rootTopic, this.endpointId, this.publisher); msg.onPublishError = this.onPublishError; return msg; } getStatusMessage(correlationToken, isResponse = false, isDeferred = false) { - const msg = new AlexaStatusMessage(correlationToken, this.rootTopic, this.endpointId, this.mqttClient, isResponse, isDeferred); + const msg = new AlexaStatusMessage(correlationToken, this.rootTopic, this.endpointId, this.publisher, isResponse, isDeferred); msg.onPublishError = this.onPublishError; return msg; } @@ -80,7 +84,7 @@ class Device extends EventEmitter { * Without a changed property send() publishes nothing and resolves with "": Alex2MQTT would drop the report. */ getChangeReport(cause = "PHYSICAL_INTERACTION") { - 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; } @@ -89,17 +93,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, activated, cause = "VOICE_INTERACTION", sendAsync = false) { - 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() { diff --git a/dist/esm/index.d.ts b/dist/esm/index.d.ts index 8a97471..a718bc1 100644 --- a/dist/esm/index.d.ts +++ b/dist/esm/index.d.ts @@ -13,6 +13,9 @@ export type { ActionId, AssetId, DisplayCategoryName, StateId, UnitOfMeasure, Ac export * as messages from "./messages/index.js"; export { AlexaError, AlexaErrors, MessageError, StateBuilder, property } from "./messages/index.js"; export type { ChangeCause, ChangeReportMessage, DeferredResponseMessage, ErrorResponseMessage, Header, ProactiveEventMessage, Property, PropertyOptions, ResponseMessage, SceneEventMessage, } from "./messages/index.js"; +export * as topics from "./topics.js"; +export { MemoryPublisher } from "./transport.js"; +export type { Publisher, PublishResult } from "./transport.js"; export { AlexaInterface } from "./compat/AlexaInterface.js"; export type { SupportedMode } from "./compat/AlexaInterface.js"; export { ActionMapping } from "./compat/ActionMapping.js"; diff --git a/dist/esm/index.js b/dist/esm/index.js index e522d0b..fa09b61 100644 --- a/dist/esm/index.js +++ b/dist/esm/index.js @@ -12,6 +12,9 @@ export { ASSETS as Assets, UNITS_OF_MEASURE as Units, ACTIONS as Actions, STATES // The messages: what the bridge publishes, built from plain values export * as messages from "./messages/index.js"; export { AlexaError, AlexaErrors, MessageError, StateBuilder, property } 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"; // 1.x export { AlexaInterface } from "./compat/AlexaInterface.js"; export { ActionMapping } from "./compat/ActionMapping.js"; diff --git a/dist/esm/topics.d.ts b/dist/esm/topics.d.ts new file mode 100644 index 0000000..14bb390 --- /dev/null +++ b/dist/esm/topics.d.ts @@ -0,0 +1,8 @@ +/** Where the bridge answers a discovery request with its endpoints. */ +export declare const discoverReply: (root: string) => string; +/** 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. */ +export declare const deferred: (root: string, endpointId: string) => string; +/** Where every ChangeReport of the root goes: the backend adds the user's token and posts it to Alexa. */ +export declare const changeReport: (root: string) => string; diff --git a/dist/esm/topics.js b/dist/esm/topics.js new file mode 100644 index 0000000..ba81060 --- /dev/null +++ b/dist/esm/topics.js @@ -0,0 +1,10 @@ +// 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) => `${root}/discover_r`; +/** 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. */ +export const deferred = (root, endpointId) => `${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) => `${root}/changeReport`; diff --git a/dist/esm/transport.d.ts b/dist/esm/transport.d.ts new file mode 100644 index 0000000..6c2ad3c --- /dev/null +++ b/dist/esm/transport.d.ts @@ -0,0 +1,48 @@ +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; +} +/** + * 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 declare class MqttPublisher implements Publisher { + private readonly client; + constructor(client: () => MqttClient | null); + publish(topic: string, message: object): Promise; +} +/** + * 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: "//alexaResponce", message: { event, context } } + */ +export declare class MemoryPublisher implements Publisher { + /** Oldest first. The message is what a subscriber gets after JSON.parse. */ + 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. */ + failWith: Error | null; + publish(topic: string, message: object): Promise; +} +/** + * 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 declare function send(publisher: Publisher | null, topic: string, build: () => object, report?: (err: Error) => void): Promise; diff --git a/dist/esm/transport.js b/dist/esm/transport.js new file mode 100644 index 0000000..403cef3 --- /dev/null +++ b/dist/esm/transport.js @@ -0,0 +1,77 @@ +const asError = (err) => (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 { + constructor(client) { + this.client = client; + } + publish(topic, message) { + const client = this.client(); + return new Promise((resolve) => { + const failed = (err) => 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: "//alexaResponce", message: { event, context } } + */ +export class MemoryPublisher { + constructor() { + /** Oldest first. The message is what a subscriber gets after JSON.parse. */ + this.published = []; + /** Set it and every publish fails with this error, as a publish to a broker that is gone does. */ + this.failWith = null; + } + publish(topic, message) { + 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, topic, build, report) { + let result; + 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 ""; +} diff --git a/dist/types/Alex2Node.d.ts b/dist/types/Alex2Node.d.ts index 983e964..e7d45dc 100644 --- a/dist/types/Alex2Node.d.ts +++ b/dist/types/Alex2Node.d.ts @@ -38,6 +38,7 @@ declare class Alex2MQTT extends EventEmitter { private rootTopic; private debugLogging; private client; + private readonly publisher; private devices; private MqttHost; private options; @@ -62,7 +63,8 @@ declare 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; registerDevice(name: string, endpointId: string, displayCategory?: DisplayCategory | DisplayCategory[] | null): Device; @@ -71,6 +73,7 @@ declare class Alex2MQTT extends EventEmitter { unregisterDevice(endpointId: string): boolean; /** Forget every device. */ clearDevices(): void; + private release; getDevices(): Device[]; getDevice(endpointId: string): Device | undefined; getRootTopic(): string; diff --git a/dist/types/compat/AlexaErrorResponse.d.ts b/dist/types/compat/AlexaErrorResponse.d.ts index 389b6cc..696f1ea 100644 --- a/dist/types/compat/AlexaErrorResponse.d.ts +++ b/dist/types/compat/AlexaErrorResponse.d.ts @@ -1,17 +1,17 @@ -import type { MqttClient } from "mqtt"; import type { ErrorResponseMessage } from "../messages/types.js"; +import type { Publisher } from "../transport.js"; /** An ErrorResponse as 1.x builds it: device.getErrorMessage(token), setErrorMessage(), send(). */ export declare class AlexaErrorResponse { private readonly correlationToken; private readonly rootTopic; private readonly endpointId; - private readonly mqttClient; + private readonly publisher; private error; private namespace?; private readonly messageId; /** Where a failed publish is reported (set by the Device that built this message, 1.5.2): send() never rejects. */ onPublishError?: (err: Error) => void; - constructor(correlationToken: string, rootTopic: string, endpointId: string, mqttClient: MqttClient); + constructor(correlationToken: string, rootTopic: string, endpointId: string, publisher: Publisher | null); /** * The error to answer with. otherParams are the fields the type adds to the payload (validRange). The header * carries the namespace the type is documented under, Alexa.ThermostatController for THERMOSTAT_IS_OFF; diff --git a/dist/types/compat/AlexaStatusMessage.d.ts b/dist/types/compat/AlexaStatusMessage.d.ts index f4dee17..2f9e51a 100644 --- a/dist/types/compat/AlexaStatusMessage.d.ts +++ b/dist/types/compat/AlexaStatusMessage.d.ts @@ -1,5 +1,5 @@ -import type { MqttClient } from "mqtt"; import type { ChangeCause, Property } from "../messages/types.js"; +import type { Publisher } from "../transport.js"; import { TemperatureSensorScale } from "./enums.js"; import type { EndpointHealth, PowerController } from "./enums.js"; /** @@ -10,7 +10,7 @@ export declare class AlexaStatusMessage { private readonly correlationToken; private readonly rootTopic; private readonly endpointId; - private readonly mqttClient; + private readonly publisher; private readonly isResponse; private readonly isDeferred; private readonly changeCause; @@ -19,7 +19,7 @@ export declare class AlexaStatusMessage { private estimatedDeferralInSeconds?; /** Where a failed publish is reported (set by the Device that built this message, 1.5.2): send() never rejects. */ onPublishError?: (err: Error) => void; - constructor(correlationToken: string, rootTopic: string, endpointId: string, mqttClient: MqttClient, isResponse?: boolean, isDeferred?: boolean, changeCause?: ChangeCause | null); + constructor(correlationToken: string, rootTopic: string, endpointId: string, publisher: Publisher | null, isResponse?: boolean, isDeferred?: boolean, changeCause?: ChangeCause | null); private addProperty; /** ChangeReport: the add*Prop calls that follow describe what CHANGED (the default for a change report). */ changed(): this; diff --git a/dist/types/device/Device.d.ts b/dist/types/device/Device.d.ts index f404968..5f0d570 100644 --- a/dist/types/device/Device.d.ts +++ b/dist/types/device/Device.d.ts @@ -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"; @@ -8,6 +7,7 @@ import type { AlexaInterfaceType } from "../compat/enums.js"; import type { ChangeCause } from "../messages/types.js"; import type { DisplayCategoryName } from "../registry/catalog.js"; import type { Directives, InterfaceDescriptor, Properties } from "../registry/types.js"; +import type { Publisher } from "../transport.js"; import { Capability } from "./Capability.js"; import type { AnyCapability, CapabilityJson, Declaration } from "./Capability.js"; import type { EndpointFields } from "./validate.js"; @@ -40,7 +40,6 @@ export interface EndpointJson extends EndpointFields { } type DeclarationArguments = {} extends Declaration ? [options?: Declaration] : [options: Declaration]; declare class Device extends EventEmitter { - private mqttClient; private rootTopic; name: string; endpointId: string; @@ -67,9 +66,14 @@ declare class Device extends EventEmitter { private capabilities; /** Where a failed publish from this device or a message it built is reported; Alex2MQTT sets it (1.5.2). */ onPublishError?: (err: Error) => void; - constructor(mqttClient: MqttClient, rootTopic: string, name: string, endpointId: string, displayCategory: Array | null, description?: string, manufacturerName?: string, manufacturer?: string, model?: string); - /** Internal (1.5.2): Alex2MQTT.connect() re-binds every registered device to its new broker client after a disconnect(). */ - setMqttClient(client: MqttClient): 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. + */ + publisher: Publisher | null; + /** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */ + constructor(client: unknown, rootTopic: string, name: string, endpointId: string, displayCategory: Array | null, description?: string, manufacturerName?: string, manufacturer?: string, model?: string); getName(): string; setName(name: string): void; getEndpointId(): string; diff --git a/dist/types/index.d.ts b/dist/types/index.d.ts index 8a97471..a718bc1 100644 --- a/dist/types/index.d.ts +++ b/dist/types/index.d.ts @@ -13,6 +13,9 @@ export type { ActionId, AssetId, DisplayCategoryName, StateId, UnitOfMeasure, Ac export * as messages from "./messages/index.js"; export { AlexaError, AlexaErrors, MessageError, StateBuilder, property } from "./messages/index.js"; export type { ChangeCause, ChangeReportMessage, DeferredResponseMessage, ErrorResponseMessage, Header, ProactiveEventMessage, Property, PropertyOptions, ResponseMessage, SceneEventMessage, } from "./messages/index.js"; +export * as topics from "./topics.js"; +export { MemoryPublisher } from "./transport.js"; +export type { Publisher, PublishResult } from "./transport.js"; export { AlexaInterface } from "./compat/AlexaInterface.js"; export type { SupportedMode } from "./compat/AlexaInterface.js"; export { ActionMapping } from "./compat/ActionMapping.js"; diff --git a/dist/types/topics.d.ts b/dist/types/topics.d.ts new file mode 100644 index 0000000..14bb390 --- /dev/null +++ b/dist/types/topics.d.ts @@ -0,0 +1,8 @@ +/** Where the bridge answers a discovery request with its endpoints. */ +export declare const discoverReply: (root: string) => string; +/** 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. */ +export declare const deferred: (root: string, endpointId: string) => string; +/** Where every ChangeReport of the root goes: the backend adds the user's token and posts it to Alexa. */ +export declare const changeReport: (root: string) => string; diff --git a/dist/types/transport.d.ts b/dist/types/transport.d.ts new file mode 100644 index 0000000..6c2ad3c --- /dev/null +++ b/dist/types/transport.d.ts @@ -0,0 +1,48 @@ +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; +} +/** + * 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 declare class MqttPublisher implements Publisher { + private readonly client; + constructor(client: () => MqttClient | null); + publish(topic: string, message: object): Promise; +} +/** + * 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: "//alexaResponce", message: { event, context } } + */ +export declare class MemoryPublisher implements Publisher { + /** Oldest first. The message is what a subscriber gets after JSON.parse. */ + 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. */ + failWith: Error | null; + publish(topic: string, message: object): Promise; +} +/** + * 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 declare function send(publisher: Publisher | null, topic: string, build: () => object, report?: (err: Error) => void): Promise; diff --git a/src/Alex2Node.ts b/src/Alex2Node.ts index ed2c782..c05f6f5 100644 --- a/src/Alex2Node.ts +++ b/src/Alex2Node.ts @@ -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(); } diff --git a/src/compat/AlexaErrorResponse.ts b/src/compat/AlexaErrorResponse.ts index b466662..e9c94f7 100644 --- a/src/compat/AlexaErrorResponse.ts +++ b/src/compat/AlexaErrorResponse.ts @@ -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 { - 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); } } diff --git a/src/compat/AlexaStatusMessage.ts b/src/compat/AlexaStatusMessage.ts index 9c23f3b..7d3d399 100644 --- a/src/compat/AlexaStatusMessage.ts +++ b/src/compat/AlexaStatusMessage.ts @@ -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 { - 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); } } diff --git a/src/device/Device.ts b/src/device/Device.ts index 3a9cb26..8515648 100644 --- a/src/device/Device.ts +++ b/src/device/Device.ts @@ -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 { - 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[] { diff --git a/src/index.ts b/src/index.ts index efe8fd3..1a1b201 100644 --- a/src/index.ts +++ b/src/index.ts @@ -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"; diff --git a/src/topics.ts b/src/topics.ts new file mode 100644 index 0000000..25b7e3f --- /dev/null +++ b/src/topics.ts @@ -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`; diff --git a/src/transport.ts b/src/transport.ts new file mode 100644 index 0000000..ad56aaf --- /dev/null +++ b/src/transport.ts @@ -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; +} + +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 { + 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: "//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 { + 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 { + 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 ""; +} diff --git a/test/compat/messages.test.js b/test/compat/messages.test.js index 3cc03f0..4d1016c 100644 --- a/test/compat/messages.test.js +++ b/test/compat/messages.test.js @@ -5,7 +5,7 @@ const { test } = require("node:test"); const assert = require("node:assert/strict"); const { Alex2MQTT, AlexaErrorResponse, AlexaErrorType, AlexaInterfaceType, AlexaStatusMessage, DisplayCategory, EndpointHealth, - PowerController, TemperatureSensorScale, + MemoryPublisher, PowerController, TemperatureSensorScale, } = require("alex2node"); const { setup, sleep, UUID_V4 } = require("../helpers/harness.js"); @@ -28,12 +28,6 @@ const property = (namespace, name, value, more = {}) => ( { namespace, name, value, timeOfSample: "