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: "