diff --git a/dist/cjs/Alex2Node.js b/dist/cjs/Alex2Node.js index 559a668..bf0f8cd 100644 --- a/dist/cjs/Alex2Node.js +++ b/dist/cjs/Alex2Node.js @@ -44,6 +44,7 @@ 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 dispatcher_js_1 = require("./dispatcher.js"); const topics = __importStar(require("./topics.js")); const transport_js_1 = require("./transport.js"); exports.DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; @@ -65,6 +66,9 @@ const DESCRIBED = [ * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } + * "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name } + * "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR: + * { endpointId, namespace, name, correlationToken } */ class Alex2MQTT extends events_1.EventEmitter { constructor(username, password, rootTopic, debugLogging = false, options = {}) { @@ -74,8 +78,6 @@ 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(); @@ -85,6 +87,18 @@ class Alex2MQTT extends events_1.EventEmitter { this.lastDiscoveryAt = null; this.options = options || {}; this.MqttHost = this.options.host || exports.DEFAULT_HOST; + this.publisher = this.options.publisher ?? new transport_js_1.MqttPublisher(() => this.client); + this.dispatcher = new dispatcher_js_1.Dispatcher({ + rootTopic, + publisher: this.publisher, + getDevice: (endpointId) => this.getDevice(endpointId), + log: (message, detail) => this.log(message, detail), + fail: (err) => this.fail(err), + emit: (event, ...args) => { this.emit(event, ...args); }, + answerWithinMs: this.options.answerWithinMs, + answerUnknownEndpoints: this.options.answerUnknownEndpoints, + timers: this.options.timers, + }); } log(message, detail) { if (this.options.log) @@ -129,50 +143,44 @@ class Alex2MQTT extends events_1.EventEmitter { this.client.on("reconnect", () => { this.log("reconnecting"); this.emit("reconnect"); }); this.client.on("close", () => { this.connected = false; this.emit("close"); }); this.client.on("error", (err) => { this.connected = false; this.fail(err); }); - this.client.on("message", (topic, message) => { - this.log(`MQTT Message Received`, { topic, payload: message.toString() }); + this.client.on("message", (topic, message) => { void this.receive(topic, message); }); + } + /** + * What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to + * the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test + * calls it in place of the broker: + * + * await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive)); + */ + async receive(topic, payload) { + const text = payload.toString(); + this.log(`MQTT Message Received`, { topic, payload: text }); + try { if (topic === topics.discover(this.rootTopic)) { - this.log("Discovery request received, getting device json..."); - const deviceArray = this.describeDevices(); - this.lastDiscoveryAt = new Date().toISOString(); - 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); - }); + await this.answerDiscovery(); return; } // Another topic can still arrive: a broker may deliver what an earlier session of this client id subscribed to const endpointId = topics.directiveEndpoint(this.rootTopic, topic); - if (endpointId === null) - return; - const device = this.devices.find((d) => d.endpointId === endpointId); - if (!device) { - this.log(`No device found for endpointId: ${endpointId}`); - return; - } - let payload; - try { - payload = JSON.parse(message.toString()); - } - catch (e) { - this.log(`Failed to parse payload JSON`, e); - return; - } - if (!payload || !payload.header) - return; - this.emit("directive", { endpointId, namespace: payload.header.namespace, name: payload.header.name }); - if (payload.header.namespace === "Alexa" && payload.header.name === "ReportState") { - this.log(`ReportState directive for device ${endpointId}`); - device.emit("ReportState", payload); - } - else { - device.emit("Event", payload, payload.header.namespace); - } - }); + if (endpointId !== null) + await this.dispatcher.dispatch(endpointId, text); + } + catch (err) { + // A listener of "discover", "directive" or "unknownEndpoint" threw + this.fail((0, transport_js_1.asError)(err)); + } + } + async answerDiscovery() { + this.log("Discovery request received, getting device json..."); + const deviceArray = this.describeDevices(); + this.lastDiscoveryAt = new Date().toISOString(); + const result = await this.publisher.publish(topics.discoverReply(this.rootTopic), deviceArray); + if (!result.ok) { + this.fail(result.error); + return; + } + this.log(`Discovery payloads published to ${result.topic}`, deviceArray); + this.emit("discover", deviceArray.length); } // The endpoint objects of a discovery answer. What check() says about a device is logged, each line once. A device // that cannot be described is left out and reported as an error: the others are still announced. @@ -253,7 +261,8 @@ 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; + // Through the dispatcher, which lets one answer to a directive pass + device.publisher = this.dispatcher.publisher; this.devices.push(device); } /** Forget a device (its listeners with it). Returns false when there was none. */ diff --git a/dist/cjs/device/Capability.js b/dist/cjs/device/Capability.js index 82bda5e..67b3f52 100644 --- a/dist/cjs/device/Capability.js +++ b/dist/cjs/device/Capability.js @@ -2,6 +2,7 @@ Object.defineProperty(exports, "__esModule", { value: true }); exports.Capability = exports.commonOptions = void 0; const schema_js_1 = require("../registry/schema.js"); +const types_js_1 = require("../registry/types.js"); /** The same as a schema. The friendly names are checked with the instance, by the rules of the interface. */ exports.commonOptions = schema_js_1.s.object({ instance: schema_js_1.s.optional(schema_js_1.s.string()), @@ -18,6 +19,7 @@ class Capability { this.endpointId = ""; /** What the bridge logs about the declaration, once, when it answers a discovery. */ this.notes = []; + this.handlers = new Map(); this.instance = declared.instance ?? ""; this.friendlyNames = declared.friendlyNames ?? []; this.retrievable = declared.retrievable ?? true; @@ -28,6 +30,21 @@ class Capability { get namespace() { return this.descriptor.namespace; } + on(name, handler) { + const names = Object.keys(this.descriptor.directives); + // An interface the library only names (tier 3) describes no directive: any name is taken + if (name !== "*" && this.descriptor.tier !== 3 && !names.includes(name)) { + throw new types_js_1.DeclarationError(this, `${name} is not a directive of the interface, which has ${names.join(", ") || "none"}`); + } + if (typeof handler !== "function") + throw new types_js_1.DeclarationError(this, `the handler of ${name} is a function`); + this.handlers.set(name, handler); + return this; + } + /** The handler on(name) registered, or the one of "*". */ + handlerFor(name) { + return this.handlers.get(name) ?? this.handlers.get("*"); + } /** Mappings of "open", "close", "raise", "lower" and of states, when the declaration has any. */ get semantics() { const { semantics } = this.options; diff --git a/dist/cjs/device/Device.js b/dist/cjs/device/Device.js index 4f5485c..97d6a3d 100644 --- a/dist/cjs/device/Device.js +++ b/dist/cjs/device/Device.js @@ -132,6 +132,32 @@ class Device extends events_1.EventEmitter { 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); } + /** + * How the device reports its state, every retrievable property of it: + * + * blinds.state((s) => s.set(lift, "rangeValue", motor.position).health("OK")); + * + * ReportState is answered with it, and it is the context of every ctx.respond(), where the handler sets only what + * the directive changed. Alexa wants the whole state in both (alexa-response.html, "Synchronous response"). + */ + state(fill) { + this.stateProvider = fill; + return this; + } + /** Answer ReportState in a handler of its own, for a state that has to be read from the device first. */ + onReportState(handler) { + this.reportStateHandler = handler; + return this; + } + /** The handler for the directives of declared capabilities that have no handler of their own. */ + onDirective(handler) { + this.directiveHandler = handler; + return this; + } + /** The capability declared for the interface, under the instance for a generic controller. */ + capability(namespace, instance = "") { + return this.capabilities.find((capability) => capability.namespace === namespace && capability.instance === instance); + } /** The capabilities the device declared, in the order it declared them. */ getCapabilities() { return this.capabilities.slice(); diff --git a/dist/cjs/dispatcher.js b/dist/cjs/dispatcher.js new file mode 100644 index 0000000..9e5f340 --- /dev/null +++ b/dist/cjs/dispatcher.js @@ -0,0 +1,319 @@ +"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.Dispatcher = exports.ANSWER_WITHIN_MS = void 0; +const build_js_1 = require("./messages/build.js"); +const errors_js_1 = require("./messages/errors.js"); +const StateBuilder_js_1 = require("./messages/StateBuilder.js"); +const schema_js_1 = require("./registry/schema.js"); +const topics = __importStar(require("./topics.js")); +const transport_js_1 = require("./transport.js"); +// unref: a directive that waits for its answer does not keep a process alive that is otherwise done +const processTimers = { + setTimeout(run, ms) { + const timer = setTimeout(run, ms); + timer.unref(); + return timer; + }, + clearTimeout(timer) { + clearTimeout(timer); + }, +}; +/** Alex2MQTT gives up on a directive after DIRECTIVE_BUDGET_MS: the answer of the watchdog has to be there before. */ +exports.ANSWER_WITHIN_MS = topics.DIRECTIVE_BUDGET_MS - 500; +// How long a directive is remembered: its messageId, to drop the copy of a mirroring broker, and that it was +// answered, to drop a second answer. A deferred answer may take this long and still find its directive. +const REMEMBERED_MS = 60000; +// A directive no handler can be called for. A device with a 1.x listener answers it by itself. +class RoutingError extends errors_js_1.AlexaError { +} +const key = (endpointId, id) => `${endpointId}\n${id}`; +class Context { + constructor(dispatcher, device, exchange, raw) { + this.dispatcher = dispatcher; + this.device = device; + this.exchange = exchange; + this.raw = raw; + this.payload = raw.payload ?? {}; + } + get endpointId() { return this.exchange.endpointId; } + get namespace() { return this.exchange.namespace; } + get name() { return this.exchange.name; } + get instance() { return this.raw.header.instance ?? ""; } + get correlationToken() { return this.exchange.correlationToken; } + get receivedAt() { return this.exchange.receivedAt; } + get answered() { return this.exchange.answered; } + get deferred() { return this.exchange.deferred; } + get reportState() { return this.namespace === "Alexa" && this.name === "ReportState"; } + async respond(fill, options = {}) { + const { endpointId, correlationToken } = this; + const own = this.capability?.descriptor.responseFor?.(this.name); + const payload = own?.payload ? own.payload.parse(options.payload ?? {}, "payload") : options.payload; + const context = this.state(fill); + return this.answer((0, build_js_1.response)({ endpointId, correlationToken, namespace: own?.namespace, name: own?.name, payload, context })); + } + async report(fill) { + if (!this.reportState) { + throw new Error(`report() answers ReportState with a StateReport: answer ${this.namespace}.${this.name} with respond()`); + } + const { endpointId, correlationToken } = this; + return this.answer((0, build_js_1.stateReport)({ endpointId, correlationToken, context: this.state(fill) })); + } + defer(estimatedDeferralInSeconds) { + const { endpointId, correlationToken } = this; + if (this.capability && !this.capability.descriptor.deferrable) + this.dispatcher.noteDeferral(this.namespace); + // Always where the first answer goes: a second DeferredResponse is refused there + const topic = topics.response(this.dispatcher.rootTopic, endpointId); + return this.dispatcher.send(topic, (0, build_js_1.deferredResponse)({ endpointId, correlationToken, estimatedDeferralInSeconds })); + } + error(typeOrError, message = "", extra = {}) { + const { endpointId, correlationToken } = this; + let error; + if (typeOrError instanceof errors_js_1.AlexaError) + error = typeOrError; + else if (typeOrError instanceof Error) + error = new errors_js_1.AlexaError("INTERNAL_ERROR", typeOrError.message); + else + error = new errors_js_1.AlexaError(typeOrError, message, extra); + const { type, alexaMessage, namespace } = error; + return this.answer((0, build_js_1.errorResponse)({ endpointId, correlationToken, type, message: alexaMessage, extra: error.extra, namespace })); + } + // The state of the device, and over it what the handler says about this directive + state(fill) { + const whole = new StateBuilder_js_1.StateBuilder(); + this.device.stateProvider?.(whole); + const own = new StateBuilder_js_1.StateBuilder(); + fill?.(own); + return [...whole.context.filter((property) => !own.context.some((other) => (0, build_js_1.sameProperty)(property, other))), ...own.context]; + } + answer(message) { + const topic = this.deferred ? topics.deferred : topics.response; + return this.dispatcher.send(topic(this.dispatcher.rootTopic, this.endpointId), message); + } +} +class Dispatcher { + constructor(options) { + this.options = options; + /** + * What the devices of the bridge publish through. An answer to a directive that waits is let through once, and a + * DeferredResponse once before it; everything else passes. + */ + this.publisher = { publish: (topic, message) => this.publish(topic, message) }; + // Both oldest first. The directives by endpoint and correlationToken, when they arrived by endpoint and messageId. + this.exchanges = new Map(); + this.arrived = new Map(); + this.deferralNoted = new Set(); + const { answerWithinMs = exports.ANSWER_WITHIN_MS } = options; + if (typeof answerWithinMs !== "number" || !(answerWithinMs >= 0)) { + throw new Error(`answerWithinMs is ${String(answerWithinMs)}: pass the milliseconds a handler has to answer in, or 0 for no limit`); + } + this.rootTopic = options.rootTopic; + this.answerWithinMs = answerWithinMs; + this.timers = options.timers ?? processTimers; + } + /** Publish an answer; one that is not published is reported to the bridge. */ + async send(topic, message) { + const result = await this.publish(topic, message); + if (!result.ok) + this.options.fail(result.error); + return result; + } + /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ + noteDeferral(namespace) { + if (this.deferralNoted.has(namespace)) + return; + this.deferralNoted.add(namespace); + this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); + } + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ + async dispatch(endpointId, text) { + const { log, emit } = this.options; + let raw; + try { + raw = JSON.parse(text); + } + catch (err) { + // Without a correlationToken there is nothing to answer + log(`the directive for ${endpointId} is dropped: it is not JSON`, err); + return; + } + if (typeof raw !== "object" || raw === null || typeof raw.header !== "object" || raw.header === null) { + log(`the directive for ${endpointId} is dropped: it has no header`); + return; + } + const { namespace, name, messageId } = raw.header; + const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; + const receivedAt = Date.now(); + this.forget(receivedAt); + if (typeof messageId === "string") { + if (this.arrived.has(key(endpointId, messageId))) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; + } + this.arrived.set(key(endpointId, messageId), receivedAt); + } + const device = this.options.getDevice(endpointId); + if (!device) { + log(`No device found for endpointId: ${endpointId}`); + emit("unknownEndpoint", { endpointId, namespace, name }); + // Off by default: the endpoint may belong to another bridge on the same root topic + if (this.options.answerUnknownEndpoints) { + const message = `${endpointId} is not an endpoint of this bridge`; + await this.send(topics.response(this.rootTopic, endpointId), (0, build_js_1.errorResponse)({ endpointId, correlationToken, type: "NO_SUCH_ENDPOINT", message })); + } + return; + } + const exchange = { endpointId, namespace, name, correlationToken, receivedAt, answered: false, deferred: false }; + if (correlationToken) { + this.exchanges.set(key(endpointId, correlationToken), exchange); + this.watch(exchange); + } + const ctx = new Context(this, device, exchange, raw); + emit("directive", { endpointId, namespace, name }); + // The 1.x listeners first, called one by one: emit() would let what one of them throws out, and lose what an + // async one rejects with + const event = ctx.reportState ? "ReportState" : "Event"; + const listeners = device.rawListeners(event); + const running = listeners.map((listener) => this.isolated(ctx, () => listener.call(device, raw, ...(ctx.reportState ? [] : [namespace])))); + running.push(this.isolated(ctx, () => this.route(device, ctx)(ctx), listeners.length > 0)); + await Promise.all(running); + } + // The handler for the directive, the payload of ctx checked for it + route(device, ctx) { + const { endpointId, namespace, name, instance } = ctx; + const refused = (type, problem) => new RoutingError(type, problem); + if (ctx.reportState) { + if (device.reportStateHandler) + return device.reportStateHandler; + if (device.stateProvider) + return (context) => context.report(); + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ReportState: register one with onReportState() or state() of the device`); + } + const capability = device.capability(namespace, instance); + const declared = `${namespace}${instance ? ` "${instance}"` : ""}`; + if (!capability) + throw refused("INVALID_DIRECTIVE", `${declared} is not declared on ${endpointId}`); + ctx.capability = capability; + const directive = capability.descriptor.directives[name]; + // An interface the library only names describes no directive: its handlers get the payload as it arrived + if (!directive && capability.descriptor.tier !== 3) + throw refused("INVALID_DIRECTIVE", `${name} is not a directive of ${namespace}`); + if (directive?.when && !directive.when(capability)) { + throw refused("INVALID_DIRECTIVE", `${declared} as ${endpointId} declares it does not take ${name}`); + } + const handler = capability.handlerFor(name) ?? device.directiveHandler; + if (!handler) { + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ${namespace}.${name}: register one with on("${name}", handler) of the capability`); + } + try { + if (directive) + ctx.payload = directive.payload.parse(ctx.payload, "payload"); + } + catch (err) { + throw err instanceof schema_js_1.SchemaError ? refused("INVALID_VALUE", err.message) : err; + } + return handler; + } + // Runs a listener or a handler. What it throws or rejects with answers the directive and goes no further. + async isolated(ctx, run, answersItself = false) { + try { + await run(); + } + catch (thrown) { + if (thrown instanceof RoutingError && answersItself) + return; + const err = (0, transport_js_1.asError)(thrown); + // An AlexaError is the answer the handler chose. Anything else is a fault in it, and so is an answer too many. + if (!(err instanceof errors_js_1.AlexaError) || ctx.answered) + this.options.fail(err); + if (!ctx.answered) + await ctx.error(err); + } + } + watch(exchange) { + const within = this.answerWithinMs; + if (within === 0) + return; + exchange.watchdog = this.timers.setTimeout(() => { + exchange.watchdog = undefined; + const { endpointId, namespace, name, correlationToken } = exchange; + const message = `no answer within ${within / 1000} s`; + this.options.emit("unanswered", { endpointId, namespace, name, correlationToken }); + this.options.fail(new Error(`${endpointId}: ${namespace}.${name} got ${message} and is answered with INTERNAL_ERROR. ` + + "Answer in the handler, or defer() first when the device takes longer")); + void this.send(topics.response(this.rootTopic, endpointId), (0, build_js_1.errorResponse)({ endpointId, correlationToken, type: "INTERNAL_ERROR", message })); + }, within); + } + publish(topic, message) { + const { event } = message; + const token = event?.header?.correlationToken; + const exchange = typeof token === "string" ? this.exchanges.get(key(event?.endpoint?.endpointId ?? "", token)) : undefined; + if (exchange) { + const deferral = event?.header?.name === "DeferredResponse"; + const { endpointId, namespace, name } = exchange; + if (exchange.answered || (deferral && exchange.deferred)) { + const was = exchange.answered ? "answered" : "deferred"; + const error = new Error(`nothing was published to ${topic}: ${namespace}.${name} for ${endpointId} was ${was} before, and Alex2MQTT takes one answer`); + return Promise.resolve({ ok: false, topic, error }); + } + if (deferral) + exchange.deferred = true; + else + exchange.answered = true; + if (exchange.watchdog !== undefined) { + this.timers.clearTimeout(exchange.watchdog); + exchange.watchdog = undefined; + } + } + return this.options.publisher.publish(topic, message); + } + // Drops what is older than REMEMBERED_MS. A directive the watchdog still waits for is kept. + forget(now) { + const old = (since) => now - since > REMEMBERED_MS; + for (const [id, since] of this.arrived) { + if (!old(since)) + break; + this.arrived.delete(id); + } + for (const [id, exchange] of this.exchanges) { + if (!old(exchange.receivedAt)) + break; + if (exchange.watchdog === undefined) + this.exchanges.delete(id); + } + } +} +exports.Dispatcher = Dispatcher; diff --git a/dist/cjs/messages/build.js b/dist/cjs/messages/build.js index fd6399a..a80af8f 100644 --- a/dist/cjs/messages/build.js +++ b/dist/cjs/messages/build.js @@ -1,6 +1,6 @@ "use strict"; Object.defineProperty(exports, "__esModule", { value: true }); -exports.MessageError = void 0; +exports.sameProperty = exports.MessageError = void 0; exports.response = response; exports.stateReport = stateReport; exports.deferredResponse = deferredResponse; @@ -78,7 +78,9 @@ function errorResponse(fields) { }, }; } +/** The same property of the same instance, whatever its value. */ const sameProperty = (a, b) => a.namespace === b.namespace && a.name === b.name && (a.instance ?? "") === (b.instance ?? ""); +exports.sameProperty = sameProperty; /** * A change of state nobody asked for. The header has no correlationToken (message-guide.html, "Header object"). * A property that changed is left out of the context: it is reported in one of the two (same page, "Context @@ -96,7 +98,7 @@ function changeReport(fields) { endpoint: { endpointId }, payload: { change: { cause: { type: fields.cause ?? "PHYSICAL_INTERACTION" }, properties: [...changed] } }, }, - context: { properties: context.filter((property) => !changed.some((other) => sameProperty(property, other))) }, + context: { properties: context.filter((property) => !changed.some((other) => (0, exports.sameProperty)(property, other))) }, }; } /** The answer to Activate and Deactivate of a scene (alexa-scenecontroller.html). */ diff --git a/dist/cjs/topics.js b/dist/cjs/topics.js index ba276f5..c49665f 100644 --- a/dist/cjs/topics.js +++ b/dist/cjs/topics.js @@ -2,7 +2,7 @@ // 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.subscriptions = exports.directive = exports.discoverReply = exports.discover = void 0; +exports.changeReport = exports.DIRECTIVE_BUDGET_MS = exports.deferred = exports.response = exports.subscriptions = exports.directive = exports.discoverReply = exports.discover = void 0; exports.directiveEndpoint = directiveEndpoint; /** Where the backend asks for the endpoints of the root. */ const discover = (root) => `${root}/discover`; @@ -30,6 +30,8 @@ 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; +/** How long Alex2MQTT waits for the answer to a directive before it tells Alexa that the endpoint did not answer. */ +exports.DIRECTIVE_BUDGET_MS = 7000; /** 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 index 0d97477..e27a393 100644 --- a/dist/cjs/transport.js +++ b/dist/cjs/transport.js @@ -1,8 +1,10 @@ "use strict"; Object.defineProperty(exports, "__esModule", { value: true }); -exports.MemoryPublisher = exports.MqttPublisher = void 0; +exports.MemoryPublisher = exports.MqttPublisher = exports.asError = void 0; exports.send = send; +/** What was thrown, as an Error. */ const asError = (err) => (err instanceof Error ? err : new Error(String(err))); +exports.asError = asError; /** * 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 @@ -15,7 +17,7 @@ class MqttPublisher { publish(topic, message) { const client = this.client(); return new Promise((resolve) => { - const failed = (err) => resolve({ ok: false, topic, error: asError(err) }); + const failed = (err) => resolve({ ok: false, topic, error: (0, exports.asError)(err) }); if (!client) { failed(new Error(`nothing was published to ${topic}: the bridge is not connected, call connect() first`)); return; @@ -34,7 +36,7 @@ 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; + * device.publisher = sent; // or new Alex2MQTT(..., { publisher: sent }) and bridge.receive() * await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send(); * sent.published[0] // { topic: "//alexaResponce", message: { event, context } } */ @@ -52,7 +54,7 @@ class MemoryPublisher { 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: false, topic, error: (0, exports.asError)(err) }); } return Promise.resolve({ ok: true, topic }); } @@ -73,7 +75,7 @@ async function send(publisher, topic, build, report) { result = await publisher.publish(topic, message); } catch (err) { - result = { ok: false, topic, error: asError(err) }; + result = { ok: false, topic, error: (0, exports.asError)(err) }; } if (result.ok) return topic; diff --git a/dist/esm/Alex2Node.d.ts b/dist/esm/Alex2Node.d.ts index 25e9bb4..513a727 100644 --- a/dist/esm/Alex2Node.d.ts +++ b/dist/esm/Alex2Node.d.ts @@ -3,6 +3,8 @@ import Device from "./device/Device.js"; import type { EndpointDefinition } from "./device/Device.js"; import { EventEmitter } from "events"; import { DisplayCategory } from "./compat/enums.js"; +import type { Timers } from "./dispatcher.js"; +import type { Publisher } from "./transport.js"; /** Optional settings for the bridge (1.5.1). Everything has the 1.4.0 behaviour as its default. */ export interface Alex2MQTTOptions { /** The broker URL. Default: the public Alex2MQTT broker, mqtt://Alex2MQTT.stormysdream.club:1883. */ @@ -16,6 +18,20 @@ export interface Alex2MQTTOptions { * { type: "AlexaInterface", interface: "Alexa", version: "3" }, which Amazon requires and 1.x left out. */ alexaInterface?: boolean; + /** + * How long a directive may wait for its answer, in milliseconds. After that the bridge answers INTERNAL_ERROR and + * emits "unanswered". Default: 6500, before Alex2MQTT gives up at 7000. 0: the bridge never answers for a handler. + */ + answerWithinMs?: number; + /** + * true: a directive for an endpoint the bridge does not have is answered with NO_SUCH_ENDPOINT. Default: it is + * left alone, because the endpoint may belong to another bridge on the same root topic. + */ + answerUnknownEndpoints?: boolean; + /** Publish here and not to the broker: a MemoryPublisher and receive() run a bridge that never connects. */ + publisher?: Publisher; + /** The timer of answerWithinMs, for a test that does not wait. Default: setTimeout, unref'd. */ + timers?: Timers; } export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; /** @@ -31,6 +47,9 @@ export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } + * "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name } + * "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR: + * { endpointId, namespace, name, correlationToken } */ declare class Alex2MQTT extends EventEmitter { private username; @@ -39,6 +58,7 @@ declare class Alex2MQTT extends EventEmitter { private debugLogging; private client; private readonly publisher; + private readonly dispatcher; private devices; private MqttHost; private options; @@ -52,6 +72,15 @@ declare class Alex2MQTT extends EventEmitter { /** Emit "error" only when somebody listens: an unhandled "error" event would crash the host process. */ private fail; connect(): void; + /** + * What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to + * the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test + * calls it in place of the broker: + * + * await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive)); + */ + receive(topic: string, payload: Buffer | string): Promise; + private answerDiscovery; private describeDevices; /** Close the broker connection (resolves once closed). The devices stay registered; connect() again reuses them. */ disconnect(): Promise; diff --git a/dist/esm/Alex2Node.js b/dist/esm/Alex2Node.js index ec92742..1f5ee6e 100644 --- a/dist/esm/Alex2Node.js +++ b/dist/esm/Alex2Node.js @@ -5,8 +5,9 @@ import { checkEndpoint } from "./device/validate.js"; import { EventEmitter } from "events"; import { DisplayCategory } from "./compat/enums.js"; import { DeclarationError } from "./registry/types.js"; +import { Dispatcher } from "./dispatcher.js"; import * as topics from "./topics.js"; -import { MqttPublisher } from "./transport.js"; +import { asError, MqttPublisher } from "./transport.js"; export const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; // What addDevice() copies from the definition to the device const DESCRIBED = [ @@ -26,6 +27,9 @@ const DESCRIBED = [ * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } + * "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name } + * "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR: + * { endpointId, namespace, name, correlationToken } */ class Alex2MQTT extends EventEmitter { constructor(username, password, rootTopic, debugLogging = false, options = {}) { @@ -35,8 +39,6 @@ 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(); @@ -46,6 +48,18 @@ class Alex2MQTT extends EventEmitter { this.lastDiscoveryAt = null; this.options = options || {}; this.MqttHost = this.options.host || DEFAULT_HOST; + this.publisher = this.options.publisher ?? new MqttPublisher(() => this.client); + this.dispatcher = new Dispatcher({ + rootTopic, + publisher: this.publisher, + getDevice: (endpointId) => this.getDevice(endpointId), + log: (message, detail) => this.log(message, detail), + fail: (err) => this.fail(err), + emit: (event, ...args) => { this.emit(event, ...args); }, + answerWithinMs: this.options.answerWithinMs, + answerUnknownEndpoints: this.options.answerUnknownEndpoints, + timers: this.options.timers, + }); } log(message, detail) { if (this.options.log) @@ -90,50 +104,44 @@ class Alex2MQTT extends EventEmitter { this.client.on("reconnect", () => { this.log("reconnecting"); this.emit("reconnect"); }); this.client.on("close", () => { this.connected = false; this.emit("close"); }); this.client.on("error", (err) => { this.connected = false; this.fail(err); }); - this.client.on("message", (topic, message) => { - this.log(`MQTT Message Received`, { topic, payload: message.toString() }); + this.client.on("message", (topic, message) => { void this.receive(topic, message); }); + } + /** + * What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to + * the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test + * calls it in place of the broker: + * + * await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive)); + */ + async receive(topic, payload) { + const text = payload.toString(); + this.log(`MQTT Message Received`, { topic, payload: text }); + try { if (topic === topics.discover(this.rootTopic)) { - this.log("Discovery request received, getting device json..."); - const deviceArray = this.describeDevices(); - this.lastDiscoveryAt = new Date().toISOString(); - 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); - }); + await this.answerDiscovery(); return; } // Another topic can still arrive: a broker may deliver what an earlier session of this client id subscribed to const endpointId = topics.directiveEndpoint(this.rootTopic, topic); - if (endpointId === null) - return; - const device = this.devices.find((d) => d.endpointId === endpointId); - if (!device) { - this.log(`No device found for endpointId: ${endpointId}`); - return; - } - let payload; - try { - payload = JSON.parse(message.toString()); - } - catch (e) { - this.log(`Failed to parse payload JSON`, e); - return; - } - if (!payload || !payload.header) - return; - this.emit("directive", { endpointId, namespace: payload.header.namespace, name: payload.header.name }); - if (payload.header.namespace === "Alexa" && payload.header.name === "ReportState") { - this.log(`ReportState directive for device ${endpointId}`); - device.emit("ReportState", payload); - } - else { - device.emit("Event", payload, payload.header.namespace); - } - }); + if (endpointId !== null) + await this.dispatcher.dispatch(endpointId, text); + } + catch (err) { + // A listener of "discover", "directive" or "unknownEndpoint" threw + this.fail(asError(err)); + } + } + async answerDiscovery() { + this.log("Discovery request received, getting device json..."); + const deviceArray = this.describeDevices(); + this.lastDiscoveryAt = new Date().toISOString(); + const result = await this.publisher.publish(topics.discoverReply(this.rootTopic), deviceArray); + if (!result.ok) { + this.fail(result.error); + return; + } + this.log(`Discovery payloads published to ${result.topic}`, deviceArray); + this.emit("discover", deviceArray.length); } // The endpoint objects of a discovery answer. What check() says about a device is logged, each line once. A device // that cannot be described is left out and reported as an error: the others are still announced. @@ -214,7 +222,8 @@ 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; + // Through the dispatcher, which lets one answer to a directive pass + device.publisher = this.dispatcher.publisher; this.devices.push(device); } /** Forget a device (its listeners with it). Returns false when there was none. */ diff --git a/dist/esm/device/Capability.d.ts b/dist/esm/device/Capability.d.ts index 3b05184..90135cc 100644 --- a/dist/esm/device/Capability.d.ts +++ b/dist/esm/device/Capability.d.ts @@ -1,3 +1,5 @@ +import type { DirectiveHandler } from "../dispatcher.js"; +import type { Infer } from "../registry/schema.js"; import type { Declared, Directives, InterfaceDescriptor, Label, Properties, Semantics } from "../registry/types.js"; /** What a declaration can set for any interface. */ export interface CommonOptions { @@ -66,10 +68,24 @@ export declare class Capability

, declared: CommonOptions & { options: O; }); get namespace(): string; + /** + * What the device does on a directive of the interface: + * + * lift.on("SetRangeValue", async (ctx) => { await motor.moveTo(ctx.payload.rangeValue); return ctx.respond(); }); + * + * "*" is the handler of every directive that has none of its own. A handler that throws answers the directive + * with an ErrorResponse: the AlexaError it threw, INTERNAL_ERROR for anything else. Throws a DeclarationError + * for a name that is not a directive of the interface. + */ + on(name: K, handler: DirectiveHandler>): this; + on(name: "*", handler: DirectiveHandler): this; + /** The handler on(name) registered, or the one of "*". */ + handlerFor(name: string): DirectiveHandler | undefined; /** Mappings of "open", "close", "raise", "lower" and of states, when the declaration has any. */ get semantics(): Semantics | undefined; /** The capability object for discovery, its fields in the order of the example on alexa-discovery-objects.html. */ diff --git a/dist/esm/device/Capability.js b/dist/esm/device/Capability.js index 1876be1..ea8f453 100644 --- a/dist/esm/device/Capability.js +++ b/dist/esm/device/Capability.js @@ -1,4 +1,5 @@ import { s } from "../registry/schema.js"; +import { DeclarationError } from "../registry/types.js"; /** The same as a schema. The friendly names are checked with the instance, by the rules of the interface. */ export const commonOptions = s.object({ instance: s.optional(s.string()), @@ -15,6 +16,7 @@ export class Capability { this.endpointId = ""; /** What the bridge logs about the declaration, once, when it answers a discovery. */ this.notes = []; + this.handlers = new Map(); this.instance = declared.instance ?? ""; this.friendlyNames = declared.friendlyNames ?? []; this.retrievable = declared.retrievable ?? true; @@ -25,6 +27,21 @@ export class Capability { get namespace() { return this.descriptor.namespace; } + on(name, handler) { + const names = Object.keys(this.descriptor.directives); + // An interface the library only names (tier 3) describes no directive: any name is taken + if (name !== "*" && this.descriptor.tier !== 3 && !names.includes(name)) { + throw new DeclarationError(this, `${name} is not a directive of the interface, which has ${names.join(", ") || "none"}`); + } + if (typeof handler !== "function") + throw new DeclarationError(this, `the handler of ${name} is a function`); + this.handlers.set(name, handler); + return this; + } + /** The handler on(name) registered, or the one of "*". */ + handlerFor(name) { + return this.handlers.get(name) ?? this.handlers.get("*"); + } /** Mappings of "open", "close", "raise", "lower" and of states, when the declaration has any. */ get semantics() { const { semantics } = this.options; diff --git a/dist/esm/device/Device.d.ts b/dist/esm/device/Device.d.ts index 5f0d570..f912268 100644 --- a/dist/esm/device/Device.d.ts +++ b/dist/esm/device/Device.d.ts @@ -5,6 +5,7 @@ import { AlexaStatusMessage } from "../compat/AlexaStatusMessage.js"; import { DisplayCategory } from "../compat/enums.js"; import type { AlexaInterfaceType } from "../compat/enums.js"; import type { ChangeCause } from "../messages/types.js"; +import type { DirectiveHandler, Fill } from "../dispatcher.js"; import type { DisplayCategoryName } from "../registry/catalog.js"; import type { Directives, InterfaceDescriptor, Properties } from "../registry/types.js"; import type { Publisher } from "../transport.js"; @@ -72,6 +73,12 @@ declare class Device extends EventEmitter { * is null every send() resolves "" and reports why. */ publisher: Publisher | null; + /** Set by state(). */ + stateProvider?: Fill; + /** Set by onReportState(). */ + reportStateHandler?: DirectiveHandler>; + /** Set by onDirective(). */ + directiveHandler?: DirectiveHandler; /** 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; @@ -95,6 +102,21 @@ declare 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, sendAsync?: boolean): Promise; + /** + * How the device reports its state, every retrievable property of it: + * + * blinds.state((s) => s.set(lift, "rangeValue", motor.position).health("OK")); + * + * ReportState is answered with it, and it is the context of every ctx.respond(), where the handler sets only what + * the directive changed. Alexa wants the whole state in both (alexa-response.html, "Synchronous response"). + */ + state(fill: Fill): this; + /** Answer ReportState in a handler of its own, for a state that has to be read from the device first. */ + onReportState(handler: DirectiveHandler>): this; + /** The handler for the directives of declared capabilities that have no handler of their own. */ + onDirective(handler: DirectiveHandler): this; + /** The capability declared for the interface, under the instance for a generic controller. */ + capability(namespace: string, instance?: string): AnyCapability | undefined; /** The capabilities the device declared, in the order it declared them. */ getCapabilities(): AnyCapability[]; setManufacturerName(name: string): void; diff --git a/dist/esm/device/Device.js b/dist/esm/device/Device.js index 8137755..8a208f4 100644 --- a/dist/esm/device/Device.js +++ b/dist/esm/device/Device.js @@ -97,6 +97,32 @@ class Device extends EventEmitter { const build = () => sceneEvent({ endpointId: this.endpointId, correlationToken, activated, cause }); return send(this.publisher, answer(this.rootTopic, this.endpointId), build, this.onPublishError); } + /** + * How the device reports its state, every retrievable property of it: + * + * blinds.state((s) => s.set(lift, "rangeValue", motor.position).health("OK")); + * + * ReportState is answered with it, and it is the context of every ctx.respond(), where the handler sets only what + * the directive changed. Alexa wants the whole state in both (alexa-response.html, "Synchronous response"). + */ + state(fill) { + this.stateProvider = fill; + return this; + } + /** Answer ReportState in a handler of its own, for a state that has to be read from the device first. */ + onReportState(handler) { + this.reportStateHandler = handler; + return this; + } + /** The handler for the directives of declared capabilities that have no handler of their own. */ + onDirective(handler) { + this.directiveHandler = handler; + return this; + } + /** The capability declared for the interface, under the instance for a generic controller. */ + capability(namespace, instance = "") { + return this.capabilities.find((capability) => capability.namespace === namespace && capability.instance === instance); + } /** The capabilities the device declared, in the order it declared them. */ getCapabilities() { return this.capabilities.slice(); diff --git a/dist/esm/dispatcher.d.ts b/dist/esm/dispatcher.d.ts new file mode 100644 index 0000000..770aead --- /dev/null +++ b/dist/esm/dispatcher.d.ts @@ -0,0 +1,103 @@ +import type Device from "./device/Device.js"; +import { StateBuilder } from "./messages/StateBuilder.js"; +import type { Publisher, PublishResult } from "./transport.js"; +/** A directive as Alex2MQTT publishes it: the directive of Alexa without endpoint.scope, the token of the user. */ +export interface RawDirective { + header: { + namespace: string; + name: string; + instance?: string; + messageId?: string; + correlationToken?: string; + payloadVersion?: string; + }; + endpoint?: { + endpointId?: string; + cookie?: Record; + }; + payload?: unknown; +} +/** Sets the properties of a message. */ +export type Fill = (state: StateBuilder) => void; +/** One directive, and the ways to answer it. Each method resolves with what became of the publish, none rejects. */ +export interface DirectiveContext { + readonly endpointId: string; + readonly namespace: string; + readonly name: string; + /** The instance of a generic controller, "" without one. */ + readonly instance: string; + /** "" when the directive came without one; it can then not be answered. */ + readonly correlationToken: string; + /** The payload, checked by the descriptor of the interface. As it arrived for an interface the library only names. */ + readonly payload: T; + readonly raw: RawDirective; + /** Date.now() when the directive arrived. */ + readonly receivedAt: number; + /** A Response, a StateReport or an ErrorResponse was sent. */ + readonly answered: boolean; + /** A DeferredResponse was sent: the answer goes to //deferredResponse. */ + readonly deferred: boolean; + /** + * Answer with a Response, or with the response the interface has of its own. The context is the state of + * device.state(), and what fill sets replaces the same property of it. + */ + respond(fill?: Fill, options?: { + payload?: Record; + }): Promise; + /** Answer ReportState with a StateReport, its properties as for respond(). */ + report(fill?: Fill): Promise; + /** Send a DeferredResponse now, when the answer takes longer than the 7 s Alex2MQTT waits for it. */ + defer(estimatedDeferralInSeconds?: number): Promise; + /** Answer with an ErrorResponse. An error that is not an AlexaError is sent as INTERNAL_ERROR. */ + error(type: string, message: string, extra?: Record): Promise; + error(err: Error): Promise; +} +/** What a handler returns is awaited and otherwise ignored; what it throws answers the directive. */ +export type DirectiveHandler = (ctx: DirectiveContext) => unknown; +/** The clock of the watchdog, replaced in tests. */ +export interface Timers { + setTimeout(run: () => void, ms: number): unknown; + clearTimeout(timer: unknown): void; +} +/** What the dispatcher needs of its bridge. */ +export interface DispatcherOptions { + rootTopic: string; + /** Where the answers go. */ + publisher: Publisher; + getDevice(endpointId: string): Device | undefined; + log(message: string, detail?: unknown): void; + /** Reports an error: the "error" event of the bridge. */ + fail(err: Error): void; + emit(event: string, ...args: unknown[]): void; + answerWithinMs?: number; + answerUnknownEndpoints?: boolean; + timers?: Timers; +} +/** Alex2MQTT gives up on a directive after DIRECTIVE_BUDGET_MS: the answer of the watchdog has to be there before. */ +export declare const ANSWER_WITHIN_MS: number; +export declare class Dispatcher { + private readonly options; + readonly rootTopic: string; + /** + * What the devices of the bridge publish through. An answer to a directive that waits is let through once, and a + * DeferredResponse once before it; everything else passes. + */ + readonly publisher: Publisher; + private readonly answerWithinMs; + private readonly timers; + private readonly exchanges; + private readonly arrived; + private readonly deferralNoted; + constructor(options: DispatcherOptions); + /** Publish an answer; one that is not published is reported to the bridge. */ + send(topic: string, message: object): Promise; + /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ + noteDeferral(namespace: string): void; + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ + dispatch(endpointId: string, text: string): Promise; + private route; + private isolated; + private watch; + private publish; + private forget; +} diff --git a/dist/esm/dispatcher.js b/dist/esm/dispatcher.js new file mode 100644 index 0000000..7c6e91f --- /dev/null +++ b/dist/esm/dispatcher.js @@ -0,0 +1,282 @@ +import { deferredResponse, errorResponse, response, sameProperty, stateReport } from "./messages/build.js"; +import { AlexaError } from "./messages/errors.js"; +import { StateBuilder } from "./messages/StateBuilder.js"; +import { SchemaError } from "./registry/schema.js"; +import * as topics from "./topics.js"; +import { asError } from "./transport.js"; +// unref: a directive that waits for its answer does not keep a process alive that is otherwise done +const processTimers = { + setTimeout(run, ms) { + const timer = setTimeout(run, ms); + timer.unref(); + return timer; + }, + clearTimeout(timer) { + clearTimeout(timer); + }, +}; +/** Alex2MQTT gives up on a directive after DIRECTIVE_BUDGET_MS: the answer of the watchdog has to be there before. */ +export const ANSWER_WITHIN_MS = topics.DIRECTIVE_BUDGET_MS - 500; +// How long a directive is remembered: its messageId, to drop the copy of a mirroring broker, and that it was +// answered, to drop a second answer. A deferred answer may take this long and still find its directive. +const REMEMBERED_MS = 60000; +// A directive no handler can be called for. A device with a 1.x listener answers it by itself. +class RoutingError extends AlexaError { +} +const key = (endpointId, id) => `${endpointId}\n${id}`; +class Context { + constructor(dispatcher, device, exchange, raw) { + this.dispatcher = dispatcher; + this.device = device; + this.exchange = exchange; + this.raw = raw; + this.payload = raw.payload ?? {}; + } + get endpointId() { return this.exchange.endpointId; } + get namespace() { return this.exchange.namespace; } + get name() { return this.exchange.name; } + get instance() { return this.raw.header.instance ?? ""; } + get correlationToken() { return this.exchange.correlationToken; } + get receivedAt() { return this.exchange.receivedAt; } + get answered() { return this.exchange.answered; } + get deferred() { return this.exchange.deferred; } + get reportState() { return this.namespace === "Alexa" && this.name === "ReportState"; } + async respond(fill, options = {}) { + const { endpointId, correlationToken } = this; + const own = this.capability?.descriptor.responseFor?.(this.name); + const payload = own?.payload ? own.payload.parse(options.payload ?? {}, "payload") : options.payload; + const context = this.state(fill); + return this.answer(response({ endpointId, correlationToken, namespace: own?.namespace, name: own?.name, payload, context })); + } + async report(fill) { + if (!this.reportState) { + throw new Error(`report() answers ReportState with a StateReport: answer ${this.namespace}.${this.name} with respond()`); + } + const { endpointId, correlationToken } = this; + return this.answer(stateReport({ endpointId, correlationToken, context: this.state(fill) })); + } + defer(estimatedDeferralInSeconds) { + const { endpointId, correlationToken } = this; + if (this.capability && !this.capability.descriptor.deferrable) + this.dispatcher.noteDeferral(this.namespace); + // Always where the first answer goes: a second DeferredResponse is refused there + const topic = topics.response(this.dispatcher.rootTopic, endpointId); + return this.dispatcher.send(topic, deferredResponse({ endpointId, correlationToken, estimatedDeferralInSeconds })); + } + error(typeOrError, message = "", extra = {}) { + const { endpointId, correlationToken } = this; + let error; + if (typeOrError instanceof AlexaError) + error = typeOrError; + else if (typeOrError instanceof Error) + error = new AlexaError("INTERNAL_ERROR", typeOrError.message); + else + error = new AlexaError(typeOrError, message, extra); + const { type, alexaMessage, namespace } = error; + return this.answer(errorResponse({ endpointId, correlationToken, type, message: alexaMessage, extra: error.extra, namespace })); + } + // The state of the device, and over it what the handler says about this directive + state(fill) { + const whole = new StateBuilder(); + this.device.stateProvider?.(whole); + const own = new StateBuilder(); + fill?.(own); + return [...whole.context.filter((property) => !own.context.some((other) => sameProperty(property, other))), ...own.context]; + } + answer(message) { + const topic = this.deferred ? topics.deferred : topics.response; + return this.dispatcher.send(topic(this.dispatcher.rootTopic, this.endpointId), message); + } +} +export class Dispatcher { + constructor(options) { + this.options = options; + /** + * What the devices of the bridge publish through. An answer to a directive that waits is let through once, and a + * DeferredResponse once before it; everything else passes. + */ + this.publisher = { publish: (topic, message) => this.publish(topic, message) }; + // Both oldest first. The directives by endpoint and correlationToken, when they arrived by endpoint and messageId. + this.exchanges = new Map(); + this.arrived = new Map(); + this.deferralNoted = new Set(); + const { answerWithinMs = ANSWER_WITHIN_MS } = options; + if (typeof answerWithinMs !== "number" || !(answerWithinMs >= 0)) { + throw new Error(`answerWithinMs is ${String(answerWithinMs)}: pass the milliseconds a handler has to answer in, or 0 for no limit`); + } + this.rootTopic = options.rootTopic; + this.answerWithinMs = answerWithinMs; + this.timers = options.timers ?? processTimers; + } + /** Publish an answer; one that is not published is reported to the bridge. */ + async send(topic, message) { + const result = await this.publish(topic, message); + if (!result.ok) + this.options.fail(result.error); + return result; + } + /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ + noteDeferral(namespace) { + if (this.deferralNoted.has(namespace)) + return; + this.deferralNoted.add(namespace); + this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); + } + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ + async dispatch(endpointId, text) { + const { log, emit } = this.options; + let raw; + try { + raw = JSON.parse(text); + } + catch (err) { + // Without a correlationToken there is nothing to answer + log(`the directive for ${endpointId} is dropped: it is not JSON`, err); + return; + } + if (typeof raw !== "object" || raw === null || typeof raw.header !== "object" || raw.header === null) { + log(`the directive for ${endpointId} is dropped: it has no header`); + return; + } + const { namespace, name, messageId } = raw.header; + const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; + const receivedAt = Date.now(); + this.forget(receivedAt); + if (typeof messageId === "string") { + if (this.arrived.has(key(endpointId, messageId))) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; + } + this.arrived.set(key(endpointId, messageId), receivedAt); + } + const device = this.options.getDevice(endpointId); + if (!device) { + log(`No device found for endpointId: ${endpointId}`); + emit("unknownEndpoint", { endpointId, namespace, name }); + // Off by default: the endpoint may belong to another bridge on the same root topic + if (this.options.answerUnknownEndpoints) { + const message = `${endpointId} is not an endpoint of this bridge`; + await this.send(topics.response(this.rootTopic, endpointId), errorResponse({ endpointId, correlationToken, type: "NO_SUCH_ENDPOINT", message })); + } + return; + } + const exchange = { endpointId, namespace, name, correlationToken, receivedAt, answered: false, deferred: false }; + if (correlationToken) { + this.exchanges.set(key(endpointId, correlationToken), exchange); + this.watch(exchange); + } + const ctx = new Context(this, device, exchange, raw); + emit("directive", { endpointId, namespace, name }); + // The 1.x listeners first, called one by one: emit() would let what one of them throws out, and lose what an + // async one rejects with + const event = ctx.reportState ? "ReportState" : "Event"; + const listeners = device.rawListeners(event); + const running = listeners.map((listener) => this.isolated(ctx, () => listener.call(device, raw, ...(ctx.reportState ? [] : [namespace])))); + running.push(this.isolated(ctx, () => this.route(device, ctx)(ctx), listeners.length > 0)); + await Promise.all(running); + } + // The handler for the directive, the payload of ctx checked for it + route(device, ctx) { + const { endpointId, namespace, name, instance } = ctx; + const refused = (type, problem) => new RoutingError(type, problem); + if (ctx.reportState) { + if (device.reportStateHandler) + return device.reportStateHandler; + if (device.stateProvider) + return (context) => context.report(); + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ReportState: register one with onReportState() or state() of the device`); + } + const capability = device.capability(namespace, instance); + const declared = `${namespace}${instance ? ` "${instance}"` : ""}`; + if (!capability) + throw refused("INVALID_DIRECTIVE", `${declared} is not declared on ${endpointId}`); + ctx.capability = capability; + const directive = capability.descriptor.directives[name]; + // An interface the library only names describes no directive: its handlers get the payload as it arrived + if (!directive && capability.descriptor.tier !== 3) + throw refused("INVALID_DIRECTIVE", `${name} is not a directive of ${namespace}`); + if (directive?.when && !directive.when(capability)) { + throw refused("INVALID_DIRECTIVE", `${declared} as ${endpointId} declares it does not take ${name}`); + } + const handler = capability.handlerFor(name) ?? device.directiveHandler; + if (!handler) { + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ${namespace}.${name}: register one with on("${name}", handler) of the capability`); + } + try { + if (directive) + ctx.payload = directive.payload.parse(ctx.payload, "payload"); + } + catch (err) { + throw err instanceof SchemaError ? refused("INVALID_VALUE", err.message) : err; + } + return handler; + } + // Runs a listener or a handler. What it throws or rejects with answers the directive and goes no further. + async isolated(ctx, run, answersItself = false) { + try { + await run(); + } + catch (thrown) { + if (thrown instanceof RoutingError && answersItself) + return; + const err = asError(thrown); + // An AlexaError is the answer the handler chose. Anything else is a fault in it, and so is an answer too many. + if (!(err instanceof AlexaError) || ctx.answered) + this.options.fail(err); + if (!ctx.answered) + await ctx.error(err); + } + } + watch(exchange) { + const within = this.answerWithinMs; + if (within === 0) + return; + exchange.watchdog = this.timers.setTimeout(() => { + exchange.watchdog = undefined; + const { endpointId, namespace, name, correlationToken } = exchange; + const message = `no answer within ${within / 1000} s`; + this.options.emit("unanswered", { endpointId, namespace, name, correlationToken }); + this.options.fail(new Error(`${endpointId}: ${namespace}.${name} got ${message} and is answered with INTERNAL_ERROR. ` + + "Answer in the handler, or defer() first when the device takes longer")); + void this.send(topics.response(this.rootTopic, endpointId), errorResponse({ endpointId, correlationToken, type: "INTERNAL_ERROR", message })); + }, within); + } + publish(topic, message) { + const { event } = message; + const token = event?.header?.correlationToken; + const exchange = typeof token === "string" ? this.exchanges.get(key(event?.endpoint?.endpointId ?? "", token)) : undefined; + if (exchange) { + const deferral = event?.header?.name === "DeferredResponse"; + const { endpointId, namespace, name } = exchange; + if (exchange.answered || (deferral && exchange.deferred)) { + const was = exchange.answered ? "answered" : "deferred"; + const error = new Error(`nothing was published to ${topic}: ${namespace}.${name} for ${endpointId} was ${was} before, and Alex2MQTT takes one answer`); + return Promise.resolve({ ok: false, topic, error }); + } + if (deferral) + exchange.deferred = true; + else + exchange.answered = true; + if (exchange.watchdog !== undefined) { + this.timers.clearTimeout(exchange.watchdog); + exchange.watchdog = undefined; + } + } + return this.options.publisher.publish(topic, message); + } + // Drops what is older than REMEMBERED_MS. A directive the watchdog still waits for is kept. + forget(now) { + const old = (since) => now - since > REMEMBERED_MS; + for (const [id, since] of this.arrived) { + if (!old(since)) + break; + this.arrived.delete(id); + } + for (const [id, exchange] of this.exchanges) { + if (!old(exchange.receivedAt)) + break; + if (exchange.watchdog === undefined) + this.exchanges.delete(id); + } + } +} diff --git a/dist/esm/index.d.ts b/dist/esm/index.d.ts index a718bc1..8828d39 100644 --- a/dist/esm/index.d.ts +++ b/dist/esm/index.d.ts @@ -16,6 +16,7 @@ export type { ChangeCause, ChangeReportMessage, DeferredResponseMessage, ErrorRe export * as topics from "./topics.js"; export { MemoryPublisher } from "./transport.js"; export type { Publisher, PublishResult } from "./transport.js"; +export type { DirectiveContext, DirectiveHandler, Fill, RawDirective, Timers } from "./dispatcher.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/messages/build.d.ts b/dist/esm/messages/build.d.ts index 3cc59b2..a5ba823 100644 --- a/dist/esm/messages/build.d.ts +++ b/dist/esm/messages/build.d.ts @@ -60,6 +60,8 @@ export interface ChangeReportFields extends Envelope { /** The other properties of the endpoint, as they are now. */ context?: readonly Property[]; } +/** The same property of the same instance, whatever its value. */ +export declare const sameProperty: (a: Property, b: Property) => boolean; /** * A change of state nobody asked for. The header has no correlationToken (message-guide.html, "Header object"). * A property that changed is left out of the context: it is reported in one of the two (same page, "Context diff --git a/dist/esm/messages/build.js b/dist/esm/messages/build.js index 55fc04f..7989381 100644 --- a/dist/esm/messages/build.js +++ b/dist/esm/messages/build.js @@ -66,7 +66,8 @@ export function errorResponse(fields) { }, }; } -const sameProperty = (a, b) => a.namespace === b.namespace && a.name === b.name && (a.instance ?? "") === (b.instance ?? ""); +/** The same property of the same instance, whatever its value. */ +export const sameProperty = (a, b) => a.namespace === b.namespace && a.name === b.name && (a.instance ?? "") === (b.instance ?? ""); /** * A change of state nobody asked for. The header has no correlationToken (message-guide.html, "Header object"). * A property that changed is left out of the context: it is reported in one of the two (same page, "Context diff --git a/dist/esm/topics.d.ts b/dist/esm/topics.d.ts index 7ff3912..3a6ab99 100644 --- a/dist/esm/topics.d.ts +++ b/dist/esm/topics.d.ts @@ -15,5 +15,7 @@ export declare function directiveEndpoint(root: string, topic: string): string | 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; +/** How long Alex2MQTT waits for the answer to a directive before it tells Alexa that the endpoint did not answer. */ +export declare const DIRECTIVE_BUDGET_MS = 7000; /** 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 index e00d375..f33deb3 100644 --- a/dist/esm/topics.js +++ b/dist/esm/topics.js @@ -20,5 +20,7 @@ export function directiveEndpoint(root, topic) { 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`; +/** How long Alex2MQTT waits for the answer to a directive before it tells Alexa that the endpoint did not answer. */ +export const DIRECTIVE_BUDGET_MS = 7000; /** 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 index 6c2ad3c..e0547bf 100644 --- a/dist/esm/transport.d.ts +++ b/dist/esm/transport.d.ts @@ -12,6 +12,8 @@ export type PublishResult = { export interface Publisher { publish(topic: string, message: object): Promise; } +/** What was thrown, as an Error. */ +export declare const asError: (err: unknown) => Error; /** * 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 @@ -26,7 +28,7 @@ export declare class MqttPublisher implements Publisher { * Keeps what it is asked to publish, for tests and dry runs without a broker: * * const sent = new MemoryPublisher(); - * device.publisher = sent; + * device.publisher = sent; // or new Alex2MQTT(..., { publisher: sent }) and bridge.receive() * await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send(); * sent.published[0] // { topic: "//alexaResponce", message: { event, context } } */ diff --git a/dist/esm/transport.js b/dist/esm/transport.js index 403cef3..1ddaa79 100644 --- a/dist/esm/transport.js +++ b/dist/esm/transport.js @@ -1,4 +1,5 @@ -const asError = (err) => (err instanceof Error ? err : new Error(String(err))); +/** What was thrown, as an Error. */ +export 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 @@ -29,7 +30,7 @@ export class MqttPublisher { * Keeps what it is asked to publish, for tests and dry runs without a broker: * * const sent = new MemoryPublisher(); - * device.publisher = sent; + * device.publisher = sent; // or new Alex2MQTT(..., { publisher: sent }) and bridge.receive() * await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send(); * sent.published[0] // { topic: "//alexaResponce", message: { event, context } } */ diff --git a/dist/types/Alex2Node.d.ts b/dist/types/Alex2Node.d.ts index 25e9bb4..513a727 100644 --- a/dist/types/Alex2Node.d.ts +++ b/dist/types/Alex2Node.d.ts @@ -3,6 +3,8 @@ import Device from "./device/Device.js"; import type { EndpointDefinition } from "./device/Device.js"; import { EventEmitter } from "events"; import { DisplayCategory } from "./compat/enums.js"; +import type { Timers } from "./dispatcher.js"; +import type { Publisher } from "./transport.js"; /** Optional settings for the bridge (1.5.1). Everything has the 1.4.0 behaviour as its default. */ export interface Alex2MQTTOptions { /** The broker URL. Default: the public Alex2MQTT broker, mqtt://Alex2MQTT.stormysdream.club:1883. */ @@ -16,6 +18,20 @@ export interface Alex2MQTTOptions { * { type: "AlexaInterface", interface: "Alexa", version: "3" }, which Amazon requires and 1.x left out. */ alexaInterface?: boolean; + /** + * How long a directive may wait for its answer, in milliseconds. After that the bridge answers INTERNAL_ERROR and + * emits "unanswered". Default: 6500, before Alex2MQTT gives up at 7000. 0: the bridge never answers for a handler. + */ + answerWithinMs?: number; + /** + * true: a directive for an endpoint the bridge does not have is answered with NO_SUCH_ENDPOINT. Default: it is + * left alone, because the endpoint may belong to another bridge on the same root topic. + */ + answerUnknownEndpoints?: boolean; + /** Publish here and not to the broker: a MemoryPublisher and receive() run a bridge that never connects. */ + publisher?: Publisher; + /** The timer of answerWithinMs, for a test that does not wait. Default: setTimeout, unref'd. */ + timers?: Timers; } export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; /** @@ -31,6 +47,9 @@ export declare const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } + * "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name } + * "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR: + * { endpointId, namespace, name, correlationToken } */ declare class Alex2MQTT extends EventEmitter { private username; @@ -39,6 +58,7 @@ declare class Alex2MQTT extends EventEmitter { private debugLogging; private client; private readonly publisher; + private readonly dispatcher; private devices; private MqttHost; private options; @@ -52,6 +72,15 @@ declare class Alex2MQTT extends EventEmitter { /** Emit "error" only when somebody listens: an unhandled "error" event would crash the host process. */ private fail; connect(): void; + /** + * What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to + * the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test + * calls it in place of the broker: + * + * await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive)); + */ + receive(topic: string, payload: Buffer | string): Promise; + private answerDiscovery; private describeDevices; /** Close the broker connection (resolves once closed). The devices stay registered; connect() again reuses them. */ disconnect(): Promise; diff --git a/dist/types/device/Capability.d.ts b/dist/types/device/Capability.d.ts index 3b05184..90135cc 100644 --- a/dist/types/device/Capability.d.ts +++ b/dist/types/device/Capability.d.ts @@ -1,3 +1,5 @@ +import type { DirectiveHandler } from "../dispatcher.js"; +import type { Infer } from "../registry/schema.js"; import type { Declared, Directives, InterfaceDescriptor, Label, Properties, Semantics } from "../registry/types.js"; /** What a declaration can set for any interface. */ export interface CommonOptions { @@ -66,10 +68,24 @@ export declare class Capability

, declared: CommonOptions & { options: O; }); get namespace(): string; + /** + * What the device does on a directive of the interface: + * + * lift.on("SetRangeValue", async (ctx) => { await motor.moveTo(ctx.payload.rangeValue); return ctx.respond(); }); + * + * "*" is the handler of every directive that has none of its own. A handler that throws answers the directive + * with an ErrorResponse: the AlexaError it threw, INTERNAL_ERROR for anything else. Throws a DeclarationError + * for a name that is not a directive of the interface. + */ + on(name: K, handler: DirectiveHandler>): this; + on(name: "*", handler: DirectiveHandler): this; + /** The handler on(name) registered, or the one of "*". */ + handlerFor(name: string): DirectiveHandler | undefined; /** Mappings of "open", "close", "raise", "lower" and of states, when the declaration has any. */ get semantics(): Semantics | undefined; /** The capability object for discovery, its fields in the order of the example on alexa-discovery-objects.html. */ diff --git a/dist/types/device/Device.d.ts b/dist/types/device/Device.d.ts index 5f0d570..f912268 100644 --- a/dist/types/device/Device.d.ts +++ b/dist/types/device/Device.d.ts @@ -5,6 +5,7 @@ import { AlexaStatusMessage } from "../compat/AlexaStatusMessage.js"; import { DisplayCategory } from "../compat/enums.js"; import type { AlexaInterfaceType } from "../compat/enums.js"; import type { ChangeCause } from "../messages/types.js"; +import type { DirectiveHandler, Fill } from "../dispatcher.js"; import type { DisplayCategoryName } from "../registry/catalog.js"; import type { Directives, InterfaceDescriptor, Properties } from "../registry/types.js"; import type { Publisher } from "../transport.js"; @@ -72,6 +73,12 @@ declare class Device extends EventEmitter { * is null every send() resolves "" and reports why. */ publisher: Publisher | null; + /** Set by state(). */ + stateProvider?: Fill; + /** Set by onReportState(). */ + reportStateHandler?: DirectiveHandler>; + /** Set by onDirective(). */ + directiveHandler?: DirectiveHandler; /** 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; @@ -95,6 +102,21 @@ declare 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, sendAsync?: boolean): Promise; + /** + * How the device reports its state, every retrievable property of it: + * + * blinds.state((s) => s.set(lift, "rangeValue", motor.position).health("OK")); + * + * ReportState is answered with it, and it is the context of every ctx.respond(), where the handler sets only what + * the directive changed. Alexa wants the whole state in both (alexa-response.html, "Synchronous response"). + */ + state(fill: Fill): this; + /** Answer ReportState in a handler of its own, for a state that has to be read from the device first. */ + onReportState(handler: DirectiveHandler>): this; + /** The handler for the directives of declared capabilities that have no handler of their own. */ + onDirective(handler: DirectiveHandler): this; + /** The capability declared for the interface, under the instance for a generic controller. */ + capability(namespace: string, instance?: string): AnyCapability | undefined; /** The capabilities the device declared, in the order it declared them. */ getCapabilities(): AnyCapability[]; setManufacturerName(name: string): void; diff --git a/dist/types/dispatcher.d.ts b/dist/types/dispatcher.d.ts new file mode 100644 index 0000000..770aead --- /dev/null +++ b/dist/types/dispatcher.d.ts @@ -0,0 +1,103 @@ +import type Device from "./device/Device.js"; +import { StateBuilder } from "./messages/StateBuilder.js"; +import type { Publisher, PublishResult } from "./transport.js"; +/** A directive as Alex2MQTT publishes it: the directive of Alexa without endpoint.scope, the token of the user. */ +export interface RawDirective { + header: { + namespace: string; + name: string; + instance?: string; + messageId?: string; + correlationToken?: string; + payloadVersion?: string; + }; + endpoint?: { + endpointId?: string; + cookie?: Record; + }; + payload?: unknown; +} +/** Sets the properties of a message. */ +export type Fill = (state: StateBuilder) => void; +/** One directive, and the ways to answer it. Each method resolves with what became of the publish, none rejects. */ +export interface DirectiveContext { + readonly endpointId: string; + readonly namespace: string; + readonly name: string; + /** The instance of a generic controller, "" without one. */ + readonly instance: string; + /** "" when the directive came without one; it can then not be answered. */ + readonly correlationToken: string; + /** The payload, checked by the descriptor of the interface. As it arrived for an interface the library only names. */ + readonly payload: T; + readonly raw: RawDirective; + /** Date.now() when the directive arrived. */ + readonly receivedAt: number; + /** A Response, a StateReport or an ErrorResponse was sent. */ + readonly answered: boolean; + /** A DeferredResponse was sent: the answer goes to //deferredResponse. */ + readonly deferred: boolean; + /** + * Answer with a Response, or with the response the interface has of its own. The context is the state of + * device.state(), and what fill sets replaces the same property of it. + */ + respond(fill?: Fill, options?: { + payload?: Record; + }): Promise; + /** Answer ReportState with a StateReport, its properties as for respond(). */ + report(fill?: Fill): Promise; + /** Send a DeferredResponse now, when the answer takes longer than the 7 s Alex2MQTT waits for it. */ + defer(estimatedDeferralInSeconds?: number): Promise; + /** Answer with an ErrorResponse. An error that is not an AlexaError is sent as INTERNAL_ERROR. */ + error(type: string, message: string, extra?: Record): Promise; + error(err: Error): Promise; +} +/** What a handler returns is awaited and otherwise ignored; what it throws answers the directive. */ +export type DirectiveHandler = (ctx: DirectiveContext) => unknown; +/** The clock of the watchdog, replaced in tests. */ +export interface Timers { + setTimeout(run: () => void, ms: number): unknown; + clearTimeout(timer: unknown): void; +} +/** What the dispatcher needs of its bridge. */ +export interface DispatcherOptions { + rootTopic: string; + /** Where the answers go. */ + publisher: Publisher; + getDevice(endpointId: string): Device | undefined; + log(message: string, detail?: unknown): void; + /** Reports an error: the "error" event of the bridge. */ + fail(err: Error): void; + emit(event: string, ...args: unknown[]): void; + answerWithinMs?: number; + answerUnknownEndpoints?: boolean; + timers?: Timers; +} +/** Alex2MQTT gives up on a directive after DIRECTIVE_BUDGET_MS: the answer of the watchdog has to be there before. */ +export declare const ANSWER_WITHIN_MS: number; +export declare class Dispatcher { + private readonly options; + readonly rootTopic: string; + /** + * What the devices of the bridge publish through. An answer to a directive that waits is let through once, and a + * DeferredResponse once before it; everything else passes. + */ + readonly publisher: Publisher; + private readonly answerWithinMs; + private readonly timers; + private readonly exchanges; + private readonly arrived; + private readonly deferralNoted; + constructor(options: DispatcherOptions); + /** Publish an answer; one that is not published is reported to the bridge. */ + send(topic: string, message: object): Promise; + /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ + noteDeferral(namespace: string): void; + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ + dispatch(endpointId: string, text: string): Promise; + private route; + private isolated; + private watch; + private publish; + private forget; +} diff --git a/dist/types/index.d.ts b/dist/types/index.d.ts index a718bc1..8828d39 100644 --- a/dist/types/index.d.ts +++ b/dist/types/index.d.ts @@ -16,6 +16,7 @@ export type { ChangeCause, ChangeReportMessage, DeferredResponseMessage, ErrorRe export * as topics from "./topics.js"; export { MemoryPublisher } from "./transport.js"; export type { Publisher, PublishResult } from "./transport.js"; +export type { DirectiveContext, DirectiveHandler, Fill, RawDirective, Timers } from "./dispatcher.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/messages/build.d.ts b/dist/types/messages/build.d.ts index 3cc59b2..a5ba823 100644 --- a/dist/types/messages/build.d.ts +++ b/dist/types/messages/build.d.ts @@ -60,6 +60,8 @@ export interface ChangeReportFields extends Envelope { /** The other properties of the endpoint, as they are now. */ context?: readonly Property[]; } +/** The same property of the same instance, whatever its value. */ +export declare const sameProperty: (a: Property, b: Property) => boolean; /** * A change of state nobody asked for. The header has no correlationToken (message-guide.html, "Header object"). * A property that changed is left out of the context: it is reported in one of the two (same page, "Context diff --git a/dist/types/topics.d.ts b/dist/types/topics.d.ts index 7ff3912..3a6ab99 100644 --- a/dist/types/topics.d.ts +++ b/dist/types/topics.d.ts @@ -15,5 +15,7 @@ export declare function directiveEndpoint(root: string, topic: string): string | 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; +/** How long Alex2MQTT waits for the answer to a directive before it tells Alexa that the endpoint did not answer. */ +export declare const DIRECTIVE_BUDGET_MS = 7000; /** 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 index 6c2ad3c..e0547bf 100644 --- a/dist/types/transport.d.ts +++ b/dist/types/transport.d.ts @@ -12,6 +12,8 @@ export type PublishResult = { export interface Publisher { publish(topic: string, message: object): Promise; } +/** What was thrown, as an Error. */ +export declare const asError: (err: unknown) => Error; /** * 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 @@ -26,7 +28,7 @@ export declare class MqttPublisher implements Publisher { * Keeps what it is asked to publish, for tests and dry runs without a broker: * * const sent = new MemoryPublisher(); - * device.publisher = sent; + * device.publisher = sent; // or new Alex2MQTT(..., { publisher: sent }) and bridge.receive() * await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send(); * sent.published[0] // { topic: "//alexaResponce", message: { event, context } } */ diff --git a/src/Alex2Node.ts b/src/Alex2Node.ts index ee09778..129cbab 100644 --- a/src/Alex2Node.ts +++ b/src/Alex2Node.ts @@ -7,8 +7,11 @@ import { checkEndpoint } from "./device/validate.js"; import { EventEmitter } from "events"; import { DisplayCategory } from "./compat/enums.js"; import { DeclarationError } from "./registry/types.js"; +import { Dispatcher } from "./dispatcher.js"; +import type { Timers } from "./dispatcher.js"; import * as topics from "./topics.js"; -import { MqttPublisher } from "./transport.js"; +import { asError, MqttPublisher } from "./transport.js"; +import type { Publisher } from "./transport.js"; /** Optional settings for the bridge (1.5.1). Everything has the 1.4.0 behaviour as its default. */ export interface Alex2MQTTOptions { @@ -23,6 +26,20 @@ export interface Alex2MQTTOptions { * { type: "AlexaInterface", interface: "Alexa", version: "3" }, which Amazon requires and 1.x left out. */ alexaInterface?: boolean; + /** + * How long a directive may wait for its answer, in milliseconds. After that the bridge answers INTERNAL_ERROR and + * emits "unanswered". Default: 6500, before Alex2MQTT gives up at 7000. 0: the bridge never answers for a handler. + */ + answerWithinMs?: number; + /** + * true: a directive for an endpoint the bridge does not have is answered with NO_SUCH_ENDPOINT. Default: it is + * left alone, because the endpoint may belong to another bridge on the same root topic. + */ + answerUnknownEndpoints?: boolean; + /** Publish here and not to the broker: a MemoryPublisher and receive() run a bridge that never connects. */ + publisher?: Publisher; + /** The timer of answerWithinMs, for a test that does not wait. Default: setTimeout, unref'd. */ + timers?: Timers; } export const DEFAULT_HOST = "mqtt://Alex2MQTT.stormysdream.club:1883"; @@ -46,11 +63,15 @@ const DESCRIBED = [ * unconditionally, and Node kills a process that has an unhandled "error" event) * "discover" (n) a discovery request was answered with n devices * "directive" (info) a directive was dispatched to a device: { endpointId, namespace, name } + * "unknownEndpoint" (info) a directive came for an endpoint that is not registered: { endpointId, namespace, name } + * "unanswered" (info) nothing answered a directive in time and the bridge answered INTERNAL_ERROR: + * { endpointId, namespace, name, correlationToken } */ 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 readonly publisher: Publisher; + private readonly dispatcher: Dispatcher; private devices: Device[] = []; private MqttHost: string; private options: Alex2MQTTOptions; @@ -71,6 +92,18 @@ class Alex2MQTT extends EventEmitter { super(); // Initialize EventEmitter this.options = options || {}; this.MqttHost = this.options.host || DEFAULT_HOST; + this.publisher = this.options.publisher ?? new MqttPublisher(() => this.client); + this.dispatcher = new Dispatcher({ + rootTopic, + publisher: this.publisher, + getDevice: (endpointId) => this.getDevice(endpointId), + log: (message, detail) => this.log(message, detail), + fail: (err) => this.fail(err), + emit: (event, ...args) => { this.emit(event, ...args); }, + answerWithinMs: this.options.answerWithinMs, + answerUnknownEndpoints: this.options.answerUnknownEndpoints, + timers: this.options.timers, + }); } private log(message: string, detail?: unknown): void { @@ -112,43 +145,41 @@ class Alex2MQTT extends EventEmitter { this.client.on("close", () => { this.connected = false; this.emit("close"); }); this.client.on("error", (err: Error) => { this.connected = false; this.fail(err); }); - this.client.on("message", (topic: string, message: Buffer) => { - this.log(`MQTT Message Received`, { topic, payload: message.toString() }); + this.client.on("message", (topic: string, message: Buffer) => { void this.receive(topic, message); }); + } + + /** + * What the bridge does with a message of the broker: it answers a discovery request, and it hands a directive to + * the handlers of its device. Resolves when they are done and never rejects. With the publisher option a test + * calls it in place of the broker: + * + * await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify(directive)); + */ + async receive(topic: string, payload: Buffer | string): Promise { + const text = payload.toString(); + this.log(`MQTT Message Received`, { topic, payload: text }); + try { if (topic === topics.discover(this.rootTopic)) { - this.log("Discovery request received, getting device json..."); - const deviceArray = this.describeDevices(); - this.lastDiscoveryAt = new Date().toISOString(); - 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); - }); + await this.answerDiscovery(); return; } // Another topic can still arrive: a broker may deliver what an earlier session of this client id subscribed to const endpointId = topics.directiveEndpoint(this.rootTopic, topic); - if (endpointId === null) return; - const device = this.devices.find((d) => d.endpointId === endpointId); - if (!device) { - this.log(`No device found for endpointId: ${endpointId}`); - return; - } - let payload; - try { - payload = JSON.parse(message.toString()); - } catch (e) { - this.log(`Failed to parse payload JSON`, e); - return; - } - if (!payload || !payload.header) return; - this.emit("directive", { endpointId, namespace: payload.header.namespace, name: payload.header.name }); - if (payload.header.namespace === "Alexa" && payload.header.name === "ReportState") { - this.log(`ReportState directive for device ${endpointId}`); - device.emit("ReportState", payload); - } else { - device.emit("Event", payload, payload.header.namespace); - } - }); + if (endpointId !== null) await this.dispatcher.dispatch(endpointId, text); + } catch (err) { + // A listener of "discover", "directive" or "unknownEndpoint" threw + this.fail(asError(err)); + } + } + + private async answerDiscovery(): Promise { + this.log("Discovery request received, getting device json..."); + const deviceArray = this.describeDevices(); + this.lastDiscoveryAt = new Date().toISOString(); + const result = await this.publisher.publish(topics.discoverReply(this.rootTopic), deviceArray); + if (!result.ok) { this.fail(result.error); return; } + this.log(`Discovery payloads published to ${result.topic}`, deviceArray); + this.emit("discover", deviceArray.length); } // The endpoint objects of a discovery answer. What check() says about a device is logged, each line once. A device @@ -229,7 +260,8 @@ 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; + // Through the dispatcher, which lets one answer to a directive pass + device.publisher = this.dispatcher.publisher; this.devices.push(device); } diff --git a/src/device/Capability.ts b/src/device/Capability.ts index 9ef7092..9d2b00a 100644 --- a/src/device/Capability.ts +++ b/src/device/Capability.ts @@ -1,4 +1,7 @@ +import type { DirectiveHandler } from "../dispatcher.js"; import { s } from "../registry/schema.js"; +import type { Infer } from "../registry/schema.js"; +import { DeclarationError } from "../registry/types.js"; import type { Declared, Directives, InterfaceDescriptor, Label, Properties, Semantics, } from "../registry/types.js"; @@ -66,6 +69,7 @@ export class Capability

>(); constructor(readonly descriptor: InterfaceDescriptor, declared: CommonOptions & { options: O }) { this.instance = declared.instance ?? ""; @@ -80,6 +84,33 @@ export class Capability

{ await motor.moveTo(ctx.payload.rangeValue); return ctx.respond(); }); + * + * "*" is the handler of every directive that has none of its own. A handler that throws answers the directive + * with an ErrorResponse: the AlexaError it threw, INTERNAL_ERROR for anything else. Throws a DeclarationError + * for a name that is not a directive of the interface. + */ + on(name: K, handler: DirectiveHandler>): this; + on(name: "*", handler: DirectiveHandler): this; + on(name: string, handler: DirectiveHandler): this { + const names = Object.keys(this.descriptor.directives); + // An interface the library only names (tier 3) describes no directive: any name is taken + if (name !== "*" && this.descriptor.tier !== 3 && !names.includes(name)) { + throw new DeclarationError(this, `${name} is not a directive of the interface, which has ${names.join(", ") || "none"}`); + } + if (typeof handler !== "function") throw new DeclarationError(this, `the handler of ${name} is a function`); + this.handlers.set(name, handler); + return this; + } + + /** The handler on(name) registered, or the one of "*". */ + handlerFor(name: string): DirectiveHandler | undefined { + return this.handlers.get(name) ?? this.handlers.get("*"); + } + /** Mappings of "open", "close", "raise", "lower" and of states, when the declaration has any. */ get semantics(): Semantics | undefined { const { semantics } = this.options as { semantics?: Semantics }; diff --git a/src/device/Device.ts b/src/device/Device.ts index 8515648..a5635b9 100644 --- a/src/device/Device.ts +++ b/src/device/Device.ts @@ -6,6 +6,7 @@ import { DisplayCategory } from "../compat/enums.js"; import type { AlexaInterfaceType } from "../compat/enums.js"; import { sceneEvent } from "../messages/build.js"; import type { ChangeCause } from "../messages/types.js"; +import type { DirectiveHandler, Fill } from "../dispatcher.js"; import type { DisplayCategoryName } from "../registry/catalog.js"; import { Alexa } from "../registry/interfaces/Alexa.js"; import { EndpointHealth } from "../registry/interfaces/EndpointHealth.js"; @@ -80,6 +81,12 @@ class Device extends EventEmitter { * is null every send() resolves "" and reports why. */ public publisher: Publisher | null = null; + /** Set by state(). */ + public stateProvider?: Fill; + /** Set by onReportState(). */ + public reportStateHandler?: DirectiveHandler>; + /** Set by onDirective(). */ + public directiveHandler?: DirectiveHandler; /** client is ignored: in 1.x it was the broker client, and a device could only be built after connect(). */ constructor( @@ -169,6 +176,32 @@ class Device extends EventEmitter { const build = () => sceneEvent({ endpointId: this.endpointId, correlationToken, activated, cause }); return send(this.publisher, answer(this.rootTopic, this.endpointId), build, this.onPublishError); } + /** + * How the device reports its state, every retrievable property of it: + * + * blinds.state((s) => s.set(lift, "rangeValue", motor.position).health("OK")); + * + * ReportState is answered with it, and it is the context of every ctx.respond(), where the handler sets only what + * the directive changed. Alexa wants the whole state in both (alexa-response.html, "Synchronous response"). + */ + state(fill: Fill): this { + this.stateProvider = fill; + return this; + } + /** Answer ReportState in a handler of its own, for a state that has to be read from the device first. */ + onReportState(handler: DirectiveHandler>): this { + this.reportStateHandler = handler; + return this; + } + /** The handler for the directives of declared capabilities that have no handler of their own. */ + onDirective(handler: DirectiveHandler): this { + this.directiveHandler = handler; + return this; + } + /** The capability declared for the interface, under the instance for a generic controller. */ + capability(namespace: string, instance = ""): AnyCapability | undefined { + return this.capabilities.find((capability) => capability.namespace === namespace && capability.instance === instance); + } /** The capabilities the device declared, in the order it declared them. */ getCapabilities(): AnyCapability[] { return this.capabilities.slice(); diff --git a/src/dispatcher.ts b/src/dispatcher.ts new file mode 100644 index 0000000..2619406 --- /dev/null +++ b/src/dispatcher.ts @@ -0,0 +1,386 @@ +// From a message on //alexaDirective to the handlers of the device, and to the one answer +// Alex2MQTT takes for it. Every answer of the bridge passes through here, those of the 1.x message classes too, so +// the dispatcher knows which directive is still waiting and answers for a handler that failed or stayed silent. +import type { AnyCapability } from "./device/Capability.js"; +import type Device from "./device/Device.js"; +import { deferredResponse, errorResponse, response, sameProperty, stateReport } from "./messages/build.js"; +import { AlexaError } from "./messages/errors.js"; +import { StateBuilder } from "./messages/StateBuilder.js"; +import type { Property } from "./messages/types.js"; +import { SchemaError } from "./registry/schema.js"; +import * as topics from "./topics.js"; +import { asError } from "./transport.js"; +import type { Publisher, PublishResult } from "./transport.js"; + +/** A directive as Alex2MQTT publishes it: the directive of Alexa without endpoint.scope, the token of the user. */ +export interface RawDirective { + header: { + namespace: string; + name: string; + instance?: string; + messageId?: string; + correlationToken?: string; + payloadVersion?: string; + }; + endpoint?: { endpointId?: string; cookie?: Record }; + payload?: unknown; +} + +/** Sets the properties of a message. */ +export type Fill = (state: StateBuilder) => void; + +/** One directive, and the ways to answer it. Each method resolves with what became of the publish, none rejects. */ +export interface DirectiveContext { + readonly endpointId: string; + readonly namespace: string; + readonly name: string; + /** The instance of a generic controller, "" without one. */ + readonly instance: string; + /** "" when the directive came without one; it can then not be answered. */ + readonly correlationToken: string; + /** The payload, checked by the descriptor of the interface. As it arrived for an interface the library only names. */ + readonly payload: T; + readonly raw: RawDirective; + /** Date.now() when the directive arrived. */ + readonly receivedAt: number; + /** A Response, a StateReport or an ErrorResponse was sent. */ + readonly answered: boolean; + /** A DeferredResponse was sent: the answer goes to //deferredResponse. */ + readonly deferred: boolean; + /** + * Answer with a Response, or with the response the interface has of its own. The context is the state of + * device.state(), and what fill sets replaces the same property of it. + */ + respond(fill?: Fill, options?: { payload?: Record }): Promise; + /** Answer ReportState with a StateReport, its properties as for respond(). */ + report(fill?: Fill): Promise; + /** Send a DeferredResponse now, when the answer takes longer than the 7 s Alex2MQTT waits for it. */ + defer(estimatedDeferralInSeconds?: number): Promise; + /** Answer with an ErrorResponse. An error that is not an AlexaError is sent as INTERNAL_ERROR. */ + error(type: string, message: string, extra?: Record): Promise; + error(err: Error): Promise; +} + +/** What a handler returns is awaited and otherwise ignored; what it throws answers the directive. */ +export type DirectiveHandler = (ctx: DirectiveContext) => unknown; + +/** The clock of the watchdog, replaced in tests. */ +export interface Timers { + setTimeout(run: () => void, ms: number): unknown; + clearTimeout(timer: unknown): void; +} + +// unref: a directive that waits for its answer does not keep a process alive that is otherwise done +const processTimers: Timers = { + setTimeout(run, ms) { + const timer = setTimeout(run, ms); + timer.unref(); + return timer; + }, + clearTimeout(timer) { + clearTimeout(timer as ReturnType); + }, +}; + +/** What the dispatcher needs of its bridge. */ +export interface DispatcherOptions { + rootTopic: string; + /** Where the answers go. */ + publisher: Publisher; + getDevice(endpointId: string): Device | undefined; + log(message: string, detail?: unknown): void; + /** Reports an error: the "error" event of the bridge. */ + fail(err: Error): void; + emit(event: string, ...args: unknown[]): void; + answerWithinMs?: number; + answerUnknownEndpoints?: boolean; + timers?: Timers; +} + +/** Alex2MQTT gives up on a directive after DIRECTIVE_BUDGET_MS: the answer of the watchdog has to be there before. */ +export const ANSWER_WITHIN_MS = topics.DIRECTIVE_BUDGET_MS - 500; + +// How long a directive is remembered: its messageId, to drop the copy of a mirroring broker, and that it was +// answered, to drop a second answer. A deferred answer may take this long and still find its directive. +const REMEMBERED_MS = 60_000; + +// A directive on its way to its answer +interface Exchange { + readonly endpointId: string; + readonly namespace: string; + readonly name: string; + readonly correlationToken: string; + readonly receivedAt: number; + answered: boolean; + deferred: boolean; + /** Set while the watchdog waits for the first answer. */ + watchdog?: unknown; +} + +// A directive no handler can be called for. A device with a 1.x listener answers it by itself. +class RoutingError extends AlexaError {} + +const key = (endpointId: string, id: string): string => `${endpointId}\n${id}`; + +class Context implements DirectiveContext { + payload: unknown; + /** The capability the directive is for; none for ReportState. */ + capability?: AnyCapability; + + constructor( + private readonly dispatcher: Dispatcher, + private readonly device: Device, + private readonly exchange: Exchange, + readonly raw: RawDirective + ) { + this.payload = raw.payload ?? {}; + } + + get endpointId(): string { return this.exchange.endpointId; } + get namespace(): string { return this.exchange.namespace; } + get name(): string { return this.exchange.name; } + get instance(): string { return this.raw.header.instance ?? ""; } + get correlationToken(): string { return this.exchange.correlationToken; } + get receivedAt(): number { return this.exchange.receivedAt; } + get answered(): boolean { return this.exchange.answered; } + get deferred(): boolean { return this.exchange.deferred; } + get reportState(): boolean { return this.namespace === "Alexa" && this.name === "ReportState"; } + + async respond(fill?: Fill, options: { payload?: Record } = {}): Promise { + const { endpointId, correlationToken } = this; + const own = this.capability?.descriptor.responseFor?.(this.name); + const payload = own?.payload ? own.payload.parse(options.payload ?? {}, "payload") : options.payload; + const context = this.state(fill); + return this.answer(response({ endpointId, correlationToken, namespace: own?.namespace, name: own?.name, payload, context })); + } + + async report(fill?: Fill): Promise { + if (!this.reportState) { + throw new Error(`report() answers ReportState with a StateReport: answer ${this.namespace}.${this.name} with respond()`); + } + const { endpointId, correlationToken } = this; + return this.answer(stateReport({ endpointId, correlationToken, context: this.state(fill) })); + } + + defer(estimatedDeferralInSeconds?: number): Promise { + const { endpointId, correlationToken } = this; + if (this.capability && !this.capability.descriptor.deferrable) this.dispatcher.noteDeferral(this.namespace); + // Always where the first answer goes: a second DeferredResponse is refused there + const topic = topics.response(this.dispatcher.rootTopic, endpointId); + return this.dispatcher.send(topic, deferredResponse({ endpointId, correlationToken, estimatedDeferralInSeconds })); + } + + error(typeOrError: string | Error, message = "", extra: Record = {}): Promise { + const { endpointId, correlationToken } = this; + let error: AlexaError; + if (typeOrError instanceof AlexaError) error = typeOrError; + else if (typeOrError instanceof Error) error = new AlexaError("INTERNAL_ERROR", typeOrError.message); + else error = new AlexaError(typeOrError, message, extra); + const { type, alexaMessage, namespace } = error; + return this.answer(errorResponse({ endpointId, correlationToken, type, message: alexaMessage, extra: error.extra, namespace })); + } + + // The state of the device, and over it what the handler says about this directive + private state(fill?: Fill): Property[] { + const whole = new StateBuilder(); + this.device.stateProvider?.(whole); + const own = new StateBuilder(); + fill?.(own); + return [...whole.context.filter((property) => !own.context.some((other) => sameProperty(property, other))), ...own.context]; + } + + private answer(message: object): Promise { + const topic = this.deferred ? topics.deferred : topics.response; + return this.dispatcher.send(topic(this.dispatcher.rootTopic, this.endpointId), message); + } +} + +export class Dispatcher { + readonly rootTopic: string; + /** + * What the devices of the bridge publish through. An answer to a directive that waits is let through once, and a + * DeferredResponse once before it; everything else passes. + */ + readonly publisher: Publisher = { publish: (topic, message) => this.publish(topic, message) }; + private readonly answerWithinMs: number; + private readonly timers: Timers; + // Both oldest first. The directives by endpoint and correlationToken, when they arrived by endpoint and messageId. + private readonly exchanges = new Map(); + private readonly arrived = new Map(); + private readonly deferralNoted = new Set(); + + constructor(private readonly options: DispatcherOptions) { + const { answerWithinMs = ANSWER_WITHIN_MS } = options; + if (typeof answerWithinMs !== "number" || !(answerWithinMs >= 0)) { + throw new Error(`answerWithinMs is ${String(answerWithinMs)}: pass the milliseconds a handler has to answer in, or 0 for no limit`); + } + this.rootTopic = options.rootTopic; + this.answerWithinMs = answerWithinMs; + this.timers = options.timers ?? processTimers; + } + + /** Publish an answer; one that is not published is reported to the bridge. */ + async send(topic: string, message: object): Promise { + const result = await this.publish(topic, message); + if (!result.ok) this.options.fail(result.error); + return result; + } + + /** defer() for an interface Amazon documents no DeferredResponse for: said once, the DeferredResponse is sent. */ + noteDeferral(namespace: string): void { + if (this.deferralNoted.has(namespace)) return; + this.deferralNoted.add(namespace); + this.options.log(`warning: Amazon documents no DeferredResponse for ${namespace}, Alexa may not wait for the answer`); + } + + /** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */ + async dispatch(endpointId: string, text: string): Promise { + const { log, emit } = this.options; + let raw: RawDirective; + try { + raw = JSON.parse(text); + } catch (err) { + // Without a correlationToken there is nothing to answer + log(`the directive for ${endpointId} is dropped: it is not JSON`, err); + return; + } + if (typeof raw !== "object" || raw === null || typeof raw.header !== "object" || raw.header === null) { + log(`the directive for ${endpointId} is dropped: it has no header`); + return; + } + const { namespace, name, messageId } = raw.header; + const correlationToken = typeof raw.header.correlationToken === "string" ? raw.header.correlationToken : ""; + const receivedAt = Date.now(); + this.forget(receivedAt); + if (typeof messageId === "string") { + if (this.arrived.has(key(endpointId, messageId))) { + log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`); + return; + } + this.arrived.set(key(endpointId, messageId), receivedAt); + } + + const device = this.options.getDevice(endpointId); + if (!device) { + log(`No device found for endpointId: ${endpointId}`); + emit("unknownEndpoint", { endpointId, namespace, name }); + // Off by default: the endpoint may belong to another bridge on the same root topic + if (this.options.answerUnknownEndpoints) { + const message = `${endpointId} is not an endpoint of this bridge`; + await this.send(topics.response(this.rootTopic, endpointId), errorResponse({ endpointId, correlationToken, type: "NO_SUCH_ENDPOINT", message })); + } + return; + } + + const exchange: Exchange = { endpointId, namespace, name, correlationToken, receivedAt, answered: false, deferred: false }; + if (correlationToken) { + this.exchanges.set(key(endpointId, correlationToken), exchange); + this.watch(exchange); + } + const ctx = new Context(this, device, exchange, raw); + emit("directive", { endpointId, namespace, name }); + + // The 1.x listeners first, called one by one: emit() would let what one of them throws out, and lose what an + // async one rejects with + const event = ctx.reportState ? "ReportState" : "Event"; + const listeners = device.rawListeners(event); + const running = listeners.map((listener) => this.isolated(ctx, () => listener.call(device, raw, ...(ctx.reportState ? [] : [namespace])))); + running.push(this.isolated(ctx, () => this.route(device, ctx)(ctx), listeners.length > 0)); + await Promise.all(running); + } + + // The handler for the directive, the payload of ctx checked for it + private route(device: Device, ctx: Context): DirectiveHandler { + const { endpointId, namespace, name, instance } = ctx; + const refused = (type: string, problem: string): RoutingError => new RoutingError(type, problem); + if (ctx.reportState) { + if (device.reportStateHandler) return device.reportStateHandler; + if (device.stateProvider) return (context) => context.report(); + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ReportState: register one with onReportState() or state() of the device`); + } + const capability = device.capability(namespace, instance); + const declared = `${namespace}${instance ? ` "${instance}"` : ""}`; + if (!capability) throw refused("INVALID_DIRECTIVE", `${declared} is not declared on ${endpointId}`); + ctx.capability = capability; + const directive = capability.descriptor.directives[name]; + // An interface the library only names describes no directive: its handlers get the payload as it arrived + if (!directive && capability.descriptor.tier !== 3) throw refused("INVALID_DIRECTIVE", `${name} is not a directive of ${namespace}`); + if (directive?.when && !directive.when(capability)) { + throw refused("INVALID_DIRECTIVE", `${declared} as ${endpointId} declares it does not take ${name}`); + } + const handler = capability.handlerFor(name) ?? device.directiveHandler; + if (!handler) { + throw refused("INVALID_DIRECTIVE", `${endpointId} has no handler for ${namespace}.${name}: register one with on("${name}", handler) of the capability`); + } + try { + if (directive) ctx.payload = directive.payload.parse(ctx.payload, "payload"); + } catch (err) { + throw err instanceof SchemaError ? refused("INVALID_VALUE", err.message) : err; + } + return handler; + } + + // Runs a listener or a handler. What it throws or rejects with answers the directive and goes no further. + private async isolated(ctx: Context, run: () => unknown, answersItself = false): Promise { + try { + await run(); + } catch (thrown) { + if (thrown instanceof RoutingError && answersItself) return; + const err = asError(thrown); + // An AlexaError is the answer the handler chose. Anything else is a fault in it, and so is an answer too many. + if (!(err instanceof AlexaError) || ctx.answered) this.options.fail(err); + if (!ctx.answered) await ctx.error(err); + } + } + + private watch(exchange: Exchange): void { + const within = this.answerWithinMs; + if (within === 0) return; + exchange.watchdog = this.timers.setTimeout(() => { + exchange.watchdog = undefined; + const { endpointId, namespace, name, correlationToken } = exchange; + const message = `no answer within ${within / 1000} s`; + this.options.emit("unanswered", { endpointId, namespace, name, correlationToken }); + this.options.fail(new Error( + `${endpointId}: ${namespace}.${name} got ${message} and is answered with INTERNAL_ERROR. ` + + "Answer in the handler, or defer() first when the device takes longer" + )); + void this.send(topics.response(this.rootTopic, endpointId), errorResponse({ endpointId, correlationToken, type: "INTERNAL_ERROR", message })); + }, within); + } + + private publish(topic: string, message: object): Promise { + const { event } = message as { event?: { header?: { name?: string; correlationToken?: unknown }; endpoint?: { endpointId?: string } } }; + const token = event?.header?.correlationToken; + const exchange = typeof token === "string" ? this.exchanges.get(key(event?.endpoint?.endpointId ?? "", token)) : undefined; + if (exchange) { + const deferral = event?.header?.name === "DeferredResponse"; + const { endpointId, namespace, name } = exchange; + if (exchange.answered || (deferral && exchange.deferred)) { + const was = exchange.answered ? "answered" : "deferred"; + const error = new Error(`nothing was published to ${topic}: ${namespace}.${name} for ${endpointId} was ${was} before, and Alex2MQTT takes one answer`); + return Promise.resolve({ ok: false, topic, error }); + } + if (deferral) exchange.deferred = true; + else exchange.answered = true; + if (exchange.watchdog !== undefined) { + this.timers.clearTimeout(exchange.watchdog); + exchange.watchdog = undefined; + } + } + return this.options.publisher.publish(topic, message); + } + + // Drops what is older than REMEMBERED_MS. A directive the watchdog still waits for is kept. + private forget(now: number): void { + const old = (since: number): boolean => now - since > REMEMBERED_MS; + for (const [id, since] of this.arrived) { + if (!old(since)) break; + this.arrived.delete(id); + } + for (const [id, exchange] of this.exchanges) { + if (!old(exchange.receivedAt)) break; + if (exchange.watchdog === undefined) this.exchanges.delete(id); + } + } +} diff --git a/src/index.ts b/src/index.ts index 1a1b201..d683bd3 100644 --- a/src/index.ts +++ b/src/index.ts @@ -39,6 +39,9 @@ export * as topics from "./topics.js"; export { MemoryPublisher } from "./transport.js"; export type { Publisher, PublishResult } from "./transport.js"; +// Directives: what a handler is given +export type { DirectiveContext, DirectiveHandler, Fill, RawDirective, Timers } from "./dispatcher.js"; + // 1.x export { AlexaInterface } from "./compat/AlexaInterface.js"; export type { SupportedMode } from "./compat/AlexaInterface.js"; diff --git a/src/messages/build.ts b/src/messages/build.ts index 5ce732a..6b2fad4 100644 --- a/src/messages/build.ts +++ b/src/messages/build.ts @@ -123,7 +123,8 @@ export interface ChangeReportFields extends Envelope { context?: readonly Property[]; } -const sameProperty = (a: Property, b: Property): boolean => +/** The same property of the same instance, whatever its value. */ +export const sameProperty = (a: Property, b: Property): boolean => a.namespace === b.namespace && a.name === b.name && (a.instance ?? "") === (b.instance ?? ""); /** diff --git a/src/topics.ts b/src/topics.ts index 586dd1c..018ff81 100644 --- a/src/topics.ts +++ b/src/topics.ts @@ -28,5 +28,8 @@ export const response = (root: string, endpointId: string): string => `${root}/$ /** Where the answer goes once a DeferredResponse was sent for the directive. */ export const deferred = (root: string, endpointId: string): string => `${root}/${endpointId}/deferredResponse`; +/** How long Alex2MQTT waits for the answer to a directive before it tells Alexa that the endpoint did not answer. */ +export const DIRECTIVE_BUDGET_MS = 7000; + /** 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 index ad56aaf..a2d3d3f 100644 --- a/src/transport.ts +++ b/src/transport.ts @@ -10,7 +10,8 @@ export interface Publisher { publish(topic: string, message: object): Promise; } -const asError = (err: unknown): Error => (err instanceof Error ? err : new Error(String(err))); +/** What was thrown, as an Error. */ +export 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(), @@ -41,7 +42,7 @@ export class MqttPublisher implements Publisher { * Keeps what it is asked to publish, for tests and dry runs without a broker: * * const sent = new MemoryPublisher(); - * device.publisher = sent; + * device.publisher = sent; // or new Alex2MQTT(..., { publisher: sent }) and bridge.receive() * await device.getStatusMessage(token, true).addPowerControllerProp(PowerController.ON).send(); * sent.published[0] // { topic: "//alexaResponce", message: { event, context } } */ diff --git a/test/dispatch/dispatch.test.js b/test/dispatch/dispatch.test.js new file mode 100644 index 0000000..afa1327 --- /dev/null +++ b/test/dispatch/dispatch.test.js @@ -0,0 +1,442 @@ +"use strict"; +// Dispatch without a broker: a bridge that publishes to a MemoryPublisher and is given its messages by receive(). +// The watchdog runs on timers the test fires by hand. +const { test } = require("node:test"); +const assert = require("node:assert/strict"); +const path = require("node:path"); +const { spawnSync } = require("node:child_process"); +const { + Alex2MQTT, AlexaErrors, AlexaInterfaceType, DeclarationError, MemoryPublisher, ModeController, PowerController, + PowerState, RangeController, BrightnessController, asset, text, +} = require("alex2node"); +const { directive } = require("../helpers/harness.js"); + +// Timers that wait for fire() +function manualTimers() { + const pending = []; + return { + pending, + setTimeout(run, ms) { + const timer = { run, ms }; + pending.push(timer); + return timer; + }, + clearTimeout(timer) { + const i = pending.indexOf(timer); + if (i !== -1) pending.splice(i, 1); + }, + fire() { + for (const timer of pending.splice(0)) timer.run(); + }, + }; +} + +function bridgeWithout(options = {}) { + const sent = new MemoryPublisher(); + const timers = manualTimers(); + const logged = []; + const bridge = new Alex2MQTT("u", "p", "root", false, { publisher: sent, timers, log: (line) => logged.push(line), ...options }); + const errors = []; + bridge.on("error", (err) => errors.push(err)); + /** What was published, as [the last part of the topic, the name of the event, its payload]. */ + const answers = () => sent.published.map(({ topic, message }) => [topic.split("/").pop(), message.event.header.name, message.event.payload]); + const receive = (message) => bridge.receive(`root/${message.endpoint.endpointId}/alexaDirective`, JSON.stringify(message)); + return { bridge, sent, timers, logged, errors, answers, receive }; +} + +const instanced = (message, instance) => ({ ...message, header: { ...message.header, instance } }); +const nextTurn = () => new Promise((resolve) => setImmediate(resolve)); + +// Blinds with a lift from 0 to 100 and a tilt in three unordered modes +function blinds(bridge) { + const device = bridge.addDevice({ endpointId: "blinds-1", name: "Blinds", categories: ["INTERIOR_BLIND"] }); + const lift = device.add(RangeController, { + instance: "Blind.Lift", friendlyNames: [asset("Alexa.Setting.Opening")], range: { min: 0, max: 100, precision: 1 }, + }); + const mode = (value) => ({ value, friendlyNames: [text(value, "en-US")] }); + const tilt = device.add(ModeController, { + instance: "Blind.Tilt", friendlyNames: [text("tilt", "en-US")], supportedModes: [mode("Up"), mode("Flat"), mode("Down")], + }); + return { device, lift, tilt }; +} + +test("a directive reaches the handler of its capability with a checked payload, and respond() answers with the state of the device under what the handler set", async () => { + const { bridge, sent, timers, receive } = bridgeWithout(); + const { device, lift, tilt } = blinds(bridge); + device.state((s) => s.set(lift, "rangeValue", 10).set(tilt, "mode", "Flat").health("OK")); + const seen = []; + lift.on("SetRangeValue", (ctx) => { + seen.push([ctx.endpointId, ctx.namespace, ctx.name, ctx.instance, ctx.correlationToken, ctx.payload, ctx.answered]); + return ctx.respond((s) => s.set(lift, "rangeValue", ctx.payload.rangeValue)); + }); + + await receive(instanced(directive("Alexa.RangeController", "SetRangeValue", "blinds-1", "ct-1", { rangeValue: 7 }), "Blind.Lift")); + assert.deepEqual(seen, [["blinds-1", "Alexa.RangeController", "SetRangeValue", "Blind.Lift", "ct-1", { rangeValue: 7 }, false]]); + assert.equal(sent.published.length, 1); + const [{ topic, message }] = sent.published; + assert.equal(topic, "root/blinds-1/alexaResponce"); + assert.deepEqual([message.event.header.namespace, message.event.header.name, message.event.header.correlationToken], ["Alexa", "Response", "ct-1"]); + assert.deepEqual(message.context.properties.map((p) => [p.name, p.instance, p.value]), [ + ["mode", "Blind.Tilt", "Flat"], + ["connectivity", undefined, { value: "OK" }], + ["rangeValue", "Blind.Lift", 7], + ]); + assert.deepEqual(timers.pending, [], "an answered directive leaves no timer"); +}); + +test("a directive no handler can be called for is answered with INVALID_DIRECTIVE", async () => { + const { bridge, answers, receive, errors } = bridgeWithout(); + const { lift, tilt } = blinds(bridge); + lift.on("SetRangeValue", (ctx) => ctx.respond()); + tilt.on("*", (ctx) => ctx.respond()); + const range = (name, payload) => directive("Alexa.RangeController", name, "blinds-1", `ct-${name}`, payload); + + await receive(directive("Alexa.PowerController", "TurnOn", "blinds-1", "ct-undeclared")); + await receive(instanced(range("SetRangeValue", { rangeValue: 7 }), "Blind.Height")); + await receive(instanced(range("SetRange", { rangeValue: 7 }), "Blind.Lift")); + await receive(instanced(range("AdjustRangeValue", { rangeValueDelta: 5, rangeValueDeltaDefault: false }), "Blind.Lift")); + await receive(instanced(directive("Alexa.ModeController", "AdjustMode", "blinds-1", "ct-unordered", { modeDelta: 1 }), "Blind.Tilt")); + assert.deepEqual(answers(), [ + ["alexaResponce", "ErrorResponse", { type: "INVALID_DIRECTIVE", message: "Alexa.PowerController is not declared on blinds-1" }], + ["alexaResponce", "ErrorResponse", { type: "INVALID_DIRECTIVE", message: "Alexa.RangeController \"Blind.Height\" is not declared on blinds-1" }], + ["alexaResponce", "ErrorResponse", { type: "INVALID_DIRECTIVE", message: "SetRange is not a directive of Alexa.RangeController" }], + ["alexaResponce", "ErrorResponse", { + type: "INVALID_DIRECTIVE", + message: "blinds-1 has no handler for Alexa.RangeController.AdjustRangeValue: register one with on(\"AdjustRangeValue\", handler) of the capability", + }], + ["alexaResponce", "ErrorResponse", { type: "INVALID_DIRECTIVE", message: "Alexa.ModeController \"Blind.Tilt\" as blinds-1 declares it does not take AdjustMode" }], + ]); + assert.deepEqual(errors, [], "the sender was wrong, not the bridge"); +}); + +test("a payload the interface does not describe is answered with INVALID_VALUE and the handler is not called", async () => { + const { bridge, answers, receive } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + let called = 0; + lamp.add(BrightnessController).on("SetBrightness", () => { called += 1; }); + + await receive(directive("Alexa.BrightnessController", "SetBrightness", "lamp-1", "ct-1", { brightness: "lots" })); + assert.equal(called, 0); + const [[, name, payload]] = answers(); + assert.equal(name, "ErrorResponse"); + assert.equal(payload.type, "INVALID_VALUE"); + assert.match(payload.message, /^payload\.brightness: expected .*got "lots"$/); +}); + +test("on(): \"*\" and onDirective() take what has no handler of its own; a name that is no directive is refused", async () => { + const { bridge, receive } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + const power = lamp.add(PowerController); + lamp.add(BrightnessController); + const called = []; + const by = (who) => (ctx) => { called.push(`${who} ${ctx.name}`); return ctx.respond(); }; + power.on("TurnOn", by("TurnOn")).on("*", by("*")); + lamp.onDirective(by("device")); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + await receive(directive("Alexa.PowerController", "TurnOff", "lamp-1", "ct-2")); + await receive(directive("Alexa.BrightnessController", "SetBrightness", "lamp-1", "ct-3", { brightness: 40 })); + assert.deepEqual(called, ["TurnOn TurnOn", "* TurnOff", "device SetBrightness"]); + assert.throws(() => power.on("Toggle", by("Toggle")), (err) => err instanceof DeclarationError + && err.message === "lamp-1: Alexa.PowerController: Toggle is not a directive of the interface, which has TurnOn, TurnOff"); +}); + +test("a handler that throws or rejects answers INTERNAL_ERROR, the bridge reports what it threw, and nothing is left unhandled", async () => { + const unhandled = []; + const note = (reason) => unhandled.push(reason); + process.on("unhandledRejection", note); + const { bridge, answers, receive, errors } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + const power = lamp.add(PowerController); + power.on("TurnOn", () => { throw new TypeError("relay is undefined"); }); + power.on("TurnOff", async () => { await nextTurn(); throw new Error("the relay did not answer"); }); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + await receive(directive("Alexa.PowerController", "TurnOff", "lamp-1", "ct-2")); + await nextTurn(); + process.off("unhandledRejection", note); + assert.deepEqual(answers(), [ + ["alexaResponce", "ErrorResponse", { type: "INTERNAL_ERROR", message: "relay is undefined" }], + ["alexaResponce", "ErrorResponse", { type: "INTERNAL_ERROR", message: "the relay did not answer" }], + ]); + assert.ok(errors[0] instanceof TypeError, "the error as it was thrown"); + assert.deepEqual(errors.map((err) => err.message), ["relay is undefined", "the relay did not answer"]); + assert.deepEqual(unhandled, []); +}); + +test("an AlexaError a handler throws is the answer, under the namespace of its type; ctx.error() sends the same", async () => { + const { bridge, sent, receive, errors } = bridgeWithout(); + const { lift } = blinds(bridge); + const outOfRange = AlexaErrors.valueOutOfRange("the lift goes from 0 to 100", { minimumValue: 0, maximumValue: 100 }); + lift.on("SetRangeValue", (ctx) => { + if (ctx.payload.rangeValue > 100) throw outOfRange; + if (ctx.payload.rangeValue === 50) return ctx.error(AlexaErrors.of("OBSTACLE_DETECTED", "something is in the way")); + return ctx.error("ENDPOINT_UNREACHABLE", "the motor is unplugged", { since: "today" }); + }); + const set = (rangeValue) => receive(instanced(directive("Alexa.RangeController", "SetRangeValue", "blinds-1", `ct-${rangeValue}`, { rangeValue }), "Blind.Lift")); + + await set(101); + await set(50); + await set(1); + assert.deepEqual(sent.published.map(({ message }) => [message.event.header.namespace, message.event.header.correlationToken, message.event.payload]), [ + ["Alexa", "ct-101", { type: "VALUE_OUT_OF_RANGE", message: "the lift goes from 0 to 100", validRange: { minimumValue: 0, maximumValue: 100 } }], + ["Alexa.Safety", "ct-50", { type: "OBSTACLE_DETECTED", message: "something is in the way" }], + ["Alexa", "ct-1", { type: "ENDPOINT_UNREACHABLE", message: "the motor is unplugged", since: "today" }], + ]); + assert.deepEqual(errors, [], "an error the handler chose is not a fault of the bridge"); +}); + +test("defer(): a DeferredResponse without a context now, the answer on deferredResponse, a second defer() refused", async () => { + const { bridge, sent, answers, receive, timers, logged, errors } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"], endpointHealth: false }); + const results = []; + lamp.add(PowerController).on("TurnOn", async (ctx) => { + results.push(await ctx.defer(5), [ctx.deferred, ctx.answered]); + assert.deepEqual(timers.pending, [], "the watchdog waits for the first answer only"); + results.push(await ctx.defer()); + results.push(await ctx.respond((s) => s.set(PowerController, "powerState", "ON")), [ctx.deferred, ctx.answered]); + }); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-2")); + assert.deepEqual(answers().slice(0, 2), [ + ["alexaResponce", "DeferredResponse", { estimatedDeferralInSeconds: 5 }], + ["deferredResponse", "Response", {}], + ]); + assert.equal("context" in sent.published[0].message, false); + assert.equal(sent.published.length, 4); + assert.deepEqual(results.slice(0, 5).map((result) => (Array.isArray(result) ? result : result.ok)), [true, [true, false], false, true, [true, true]]); + assert.equal( + results[2].error.message, + "nothing was published to root/lamp-1/alexaResponce: Alexa.PowerController.TurnOn for lamp-1 was deferred before, and Alex2MQTT takes one answer" + ); + assert.deepEqual(errors.map((err) => err.message), [results[2].error.message, results[7].error.message]); + const warnings = logged.filter((line) => line.startsWith("warning: Amazon documents no DeferredResponse")); + assert.deepEqual(warnings, ["warning: Amazon documents no DeferredResponse for Alexa.PowerController, Alexa may not wait for the answer"]); +}); + +test("the watchdog: a directive nobody answers gets INTERNAL_ERROR and the bridge emits \"unanswered\"; a late answer is refused", async () => { + const { bridge, answers, receive, timers, errors } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + let late; + lamp.add(PowerController).on("TurnOn", (ctx) => { late = ctx; }); + const unanswered = []; + bridge.on("unanswered", (info) => unanswered.push(info)); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + assert.deepEqual(answers(), []); + assert.deepEqual(timers.pending.map((timer) => timer.ms), [6500], "before Alex2MQTT gives up at 7 s"); + timers.fire(); + await nextTurn(); + assert.deepEqual(answers(), [["alexaResponce", "ErrorResponse", { type: "INTERNAL_ERROR", message: "no answer within 6.5 s" }]]); + assert.deepEqual(unanswered, [{ endpointId: "lamp-1", namespace: "Alexa.PowerController", name: "TurnOn", correlationToken: "ct-1" }]); + assert.match(errors[0].message, /^lamp-1: Alexa\.PowerController\.TurnOn got no answer within 6\.5 s/); + + const result = await late.respond(); + assert.equal(result.ok, false); + assert.match(result.error.message, /was answered before/); + assert.equal(answers().length, 1); +}); + +test("answerWithinMs: the time of the watchdog, 0 for none, anything else refused", async () => { + const quick = bridgeWithout({ answerWithinMs: 200 }); + const never = bridgeWithout({ answerWithinMs: 0 }); + for (const { bridge, receive } of [quick, never]) { + bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }).add(PowerController).on("TurnOn", () => {}); + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + } + assert.deepEqual(quick.timers.pending.map((timer) => timer.ms), [200]); + assert.deepEqual(never.timers.pending, []); + assert.throws( + () => bridgeWithout({ answerWithinMs: -1 }), + { message: "answerWithinMs is -1: pass the milliseconds a handler has to answer in, or 0 for no limit" } + ); +}); + +test("the timer of the watchdog does not keep the process alive", () => { + const script = ` + const { Alex2MQTT, MemoryPublisher, PowerController } = require(${JSON.stringify(path.join(__dirname, "..", ".."))}); + const bridge = new Alex2MQTT("u", "p", "root", false, { publisher: new MemoryPublisher(), answerWithinMs: 60000 }); + bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }).add(PowerController).on("TurnOn", () => {}); + const header = { namespace: "Alexa.PowerController", name: "TurnOn", messageId: "m-1", correlationToken: "ct-1", payloadVersion: "3" }; + bridge.receive("root/lamp-1/alexaDirective", JSON.stringify({ header, endpoint: { endpointId: "lamp-1" }, payload: {} })); + `; + const run = spawnSync(process.execPath, ["-e", script], { encoding: "utf8", timeout: 10000 }); + assert.equal(run.status, 0, `the process did not end by itself: ${run.error || run.stderr}`); +}); + +test("a directive that arrives twice is handled once, by its messageId", async () => { + const { bridge, answers, receive, logged } = bridgeWithout(); + const called = []; + for (const endpointId of ["lamp-1", "lamp-2"]) { + const lamp = bridge.addDevice({ endpointId, name: `Lamp ${endpointId.slice(-1)}`, categories: ["LIGHT"] }); + lamp.add(PowerController).on("TurnOn", (ctx) => { called.push(endpointId); return ctx.respond(); }); + } + const turnOn = directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1"); + + await receive(turnOn); + await receive(turnOn); + // The messageId is new to the other endpoint + await receive({ ...turnOn, endpoint: { endpointId: "lamp-2" } }); + assert.deepEqual(called, ["lamp-1", "lamp-2"]); + assert.equal(answers().length, 2); + assert.ok(logged.includes("the directive message-ct-1 for lamp-1 is dropped: it arrived before")); +}); + +test("a directive for an endpoint the bridge does not have: \"unknownEndpoint\", and NO_SUCH_ENDPOINT only when asked for", async () => { + const silent = bridgeWithout(); + const answering = bridgeWithout({ answerUnknownEndpoints: true }); + const unknown = []; + silent.bridge.on("unknownEndpoint", (info) => unknown.push(info)); + + await silent.receive(directive("Alexa.PowerController", "TurnOn", "lamp-9", "ct-1")); + await answering.receive(directive("Alexa.PowerController", "TurnOn", "lamp-9", "ct-1")); + assert.deepEqual(unknown, [{ endpointId: "lamp-9", namespace: "Alexa.PowerController", name: "TurnOn" }]); + assert.deepEqual(silent.answers(), []); + assert.deepEqual(silent.timers.pending, []); + assert.deepEqual(answering.answers(), [ + ["alexaResponce", "ErrorResponse", { type: "NO_SUCH_ENDPOINT", message: "lamp-9 is not an endpoint of this bridge" }], + ]); +}); + +test("what is not a directive is dropped: nothing published, nothing thrown", async () => { + const { bridge, answers, errors, logged, timers } = bridgeWithout(); + bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + await bridge.receive("root/lamp-1/alexaDirective", "{ not JSON"); + await bridge.receive("root/lamp-1/alexaDirective", "null"); + await bridge.receive("root/lamp-1/alexaDirective", JSON.stringify({ payload: {} })); + assert.deepEqual([answers(), errors, timers.pending], [[], [], []]); + assert.deepEqual(logged.filter((line) => line.includes("dropped")), [ + "the directive for lamp-1 is dropped: it is not JSON", + "the directive for lamp-1 is dropped: it has no header", + "the directive for lamp-1 is dropped: it has no header", + ]); +}); + +test("ReportState: answered from state() of the device, by onReportState(), or with INVALID_DIRECTIVE when there is neither", async () => { + const { bridge, sent, answers, receive } = bridgeWithout(); + const lamp = bridge.addDevice({ endpointId: "lamp-1", name: "Lamp", categories: ["LIGHT"] }); + lamp.add(PowerController); + const stateOf = () => sent.published[sent.published.length - 1].message.context.properties.map((p) => [p.name, p.value]); + + await receive(directive("Alexa", "ReportState", "lamp-1", "ct-1")); + assert.deepEqual(answers()[0].slice(0, 2), ["alexaResponce", "ErrorResponse"]); + assert.equal(answers()[0][2].type, "INVALID_DIRECTIVE"); + + lamp.state((s) => s.set(PowerController, "powerState", "OFF").health("OK")); + await receive(directive("Alexa", "ReportState", "lamp-1", "ct-2")); + assert.deepEqual(answers()[1], ["alexaResponce", "StateReport", {}]); + assert.deepEqual(stateOf(), [["powerState", "OFF"], ["connectivity", { value: "OK" }]]); + + lamp.onReportState((ctx) => ctx.report((s) => s.health("UNREACHABLE"))); + await receive(directive("Alexa", "ReportState", "lamp-1", "ct-3")); + assert.deepEqual(stateOf(), [["powerState", "OFF"], ["connectivity", { value: "UNREACHABLE" }]]); + + lamp.capability("Alexa.PowerController").on("TurnOn", (ctx) => ctx.report()); + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-4")); + assert.deepEqual(answers()[3][2], { + type: "INTERNAL_ERROR", + message: "report() answers ReportState with a StateReport: answer Alexa.PowerController.TurnOn with respond()", + }); +}); + +// The device of a 1.x caller: declared by name, one listener for every directive, the answer an AlexaStatusMessage +function legacyLamp(bridge, onEvent) { + const lamp = bridge.registerDevice("Lamp", "lamp-1", ["LIGHT"]); + lamp.addCapability("Alexa.PowerController", { proactivelyReported: true }); + lamp.addCapability(AlexaInterfaceType.ENDPOINT_HEALTH); + lamp.on("Event", onEvent); + lamp.on("ReportState", (request) => lamp.getStatusMessage(request.header.correlationToken).addPowerControllerProp(PowerState.OFF).send()); + return lamp; +} + +test("1.x: \"Event\" and \"ReportState\" listeners get every directive and their AlexaStatusMessage is the answer", async () => { + const { bridge, answers, receive, timers, errors } = bridgeWithout(); + const got = []; + const lamp = legacyLamp(bridge, (request, namespace) => { + got.push([namespace, request.header.name, request.payload]); + if (namespace !== "Alexa.PowerController") { + const error = lamp.getErrorMessage(request.header.correlationToken); + error.setErrorMessage("INVALID_DIRECTIVE", "a lamp has no thermostat"); + return error.send(); + } + return lamp.getStatusMessage(request.header.correlationToken, true).addPowerControllerProp(PowerState.ON).send(); + }); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + await receive(directive("Alexa.ThermostatController", "SetTargetTemperature", "lamp-1", "ct-2", { targetSetpoint: { value: 20, scale: "CELSIUS" } })); + await receive(directive("Alexa", "ReportState", "lamp-1", "ct-3")); + assert.deepEqual(got, [ + ["Alexa.PowerController", "TurnOn", {}], + ["Alexa.ThermostatController", "SetTargetTemperature", { targetSetpoint: { value: 20, scale: "CELSIUS" } }], + ]); + assert.deepEqual(answers(), [ + ["alexaResponce", "Response", {}], + ["alexaResponce", "ErrorResponse", { type: "INVALID_DIRECTIVE", message: "a lamp has no thermostat" }], + ["alexaResponce", "StateReport", {}], + ]); + assert.deepEqual([timers.pending, errors], [[], []], "each answer stopped its watchdog"); +}); + +test("1.x: a deferred answer, DeferredResponse first and the Response on deferredResponse later", async () => { + const { bridge, answers, receive, timers, errors } = bridgeWithout(); + let finish; + const lamp = legacyLamp(bridge, (request) => { + const token = request.header.correlationToken; + lamp.getStatusMessage(token, false, true).addEstimatedDeferralTime(20).send(false).catch(() => {}); + finish = () => lamp.getStatusMessage(token, true).addPowerControllerProp(PowerState.ON).send(true); + }); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + assert.deepEqual(timers.pending, []); + assert.equal(await finish(), "root/lamp-1/deferredResponse"); + assert.equal(await finish(), "", "the second answer is not published"); + assert.deepEqual(answers(), [ + ["alexaResponce", "DeferredResponse", { estimatedDeferralInSeconds: 20 }], + ["deferredResponse", "Response", {}], + ]); + assert.match(errors[0].message, /Alexa\.PowerController\.TurnOn for lamp-1 was answered before/); +}); + +test("1.x: a listener that stays silent, throws or rejects gets no INVALID_DIRECTIVE but the answer of the watchdog or INTERNAL_ERROR", async () => { + const unhandled = []; + const note = (reason) => unhandled.push(reason); + process.on("unhandledRejection", note); + const { bridge, answers, receive, timers, errors } = bridgeWithout(); + legacyLamp(bridge, async (request) => { + if (request.header.name === "TurnOff") throw new Error("the relay did not answer"); + }); + + await receive(directive("Alexa.SceneController", "Activate", "lamp-1", "ct-1")); + assert.deepEqual(answers(), [], "the listener owns the answer"); + timers.fire(); + await receive(directive("Alexa.PowerController", "TurnOff", "lamp-1", "ct-2")); + await nextTurn(); + process.off("unhandledRejection", note); + assert.deepEqual(answers(), [ + ["alexaResponce", "ErrorResponse", { type: "INTERNAL_ERROR", message: "no answer within 6.5 s" }], + ["alexaResponce", "ErrorResponse", { type: "INTERNAL_ERROR", message: "the relay did not answer" }], + ]); + assert.equal(errors.length, 2); + assert.deepEqual(unhandled, []); +}); + +test("1.x and a handler on one device: the listener is called first, and the first answer is the one published", async () => { + const { bridge, answers, receive, errors } = bridgeWithout(); + const order = []; + const lamp = legacyLamp(bridge, (request) => { + order.push("listener"); + return lamp.getStatusMessage(request.header.correlationToken, true).addPowerControllerProp(PowerState.ON).send(); + }); + let second; + lamp.capability("Alexa.PowerController").on("TurnOn", async (ctx) => { + order.push("handler"); + second = await ctx.respond(); + }); + + await receive(directive("Alexa.PowerController", "TurnOn", "lamp-1", "ct-1")); + assert.deepEqual(order, ["listener", "handler"]); + assert.equal(answers().length, 1); + assert.equal(second.ok, false); + assert.deepEqual(errors.map((err) => err.message), [second.error.message]); +}); diff --git a/test/fixtures/types.ts b/test/fixtures/types.ts index 38e61c4..f026ed5 100644 --- a/test/fixtures/types.ts +++ b/test/fixtures/types.ts @@ -157,3 +157,15 @@ messages.deferredResponse({ endpointId: "lamp-1", correlationToken: "ct" }).cont // @ts-expect-error the range is two numbers AlexaErrors.valueOutOfRange("too high", { minimumValue: 1 }); const errorNamespace: string = AlexaErrors.setpointsTooClose("too close", { value: 2, scale: "CELSIUS" }).namespace; + +// Handlers: the directives of the interface by name, the payload of each typed by its descriptor +lift.on("SetRangeValue", (ctx) => { + const target: number = ctx.payload.rangeValue; + return ctx.respond((s) => s.set(lift, "rangeValue", target)); +}); +lift.on("*", (ctx) => ctx.error("INVALID_DIRECTIVE", `${ctx.name} is not supported`)); +// @ts-expect-error not a directive of the interface +lift.on("SetRange", (ctx) => ctx.respond()); +// @ts-expect-error the payload of SetRangeValue has no brightness +lift.on("SetRangeValue", (ctx) => ctx.payload.brightness); +device.state((s) => s.health("OK")).onReportState((ctx) => ctx.report());