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.deferredByAnother()) 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 })); } // An interface of the device documents a DeferredResponse for this directive of another one deferredByAnother() { const { namespace, name } = this; return this.device.getCapabilities().some(({ descriptor }) => descriptor.defers?.some((directive) => directive.namespace === namespace && directive.name === name)); } // 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); } } }