dispatch: typed handlers, respond/defer/error, automatic ErrorResponse, watchdog
src/dispatcher.ts routes a directive to capability.on(name | "*"), device.onDirective() or device.onReportState(), with the payload checked by the descriptor of the interface. The DirectiveContext answers with respond/report/defer/error; respond() and report() start from device.state(). A handler or 1.x listener that throws or rejects is answered with INTERNAL_ERROR (an AlexaError with itself) and reported through "error", never as an unhandled rejection. Undeclared interface, unknown directive, AdjustMode on an unordered mode and a missing handler get INVALID_DIRECTIVE, a bad payload INVALID_VALUE; a device with an "Event" or "ReportState" listener keeps the answer to itself. Every answer passes the dispatcher: the first one per correlationToken is published, a second is refused. No answer within answerWithinMs (6500, 0 = off, unref'd timer) sends INTERNAL_ERROR and emits "unanswered". A messageId that arrived in the last 60 s is dropped. New: "unknownEndpoint", options answerUnknownEndpoints, publisher, timers, and bridge.receive() to run a bridge without a broker. 168 tests pass (18 new). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
e0fe9d2812
commit
feff828ec1
40 changed files with 2149 additions and 136 deletions
93
dist/cjs/Alex2Node.js
vendored
93
dist/cjs/Alex2Node.js
vendored
|
|
@ -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. */
|
||||
|
|
|
|||
17
dist/cjs/device/Capability.js
vendored
17
dist/cjs/device/Capability.js
vendored
|
|
@ -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;
|
||||
|
|
|
|||
26
dist/cjs/device/Device.js
vendored
26
dist/cjs/device/Device.js
vendored
|
|
@ -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();
|
||||
|
|
|
|||
319
dist/cjs/dispatcher.js
vendored
Normal file
319
dist/cjs/dispatcher.js
vendored
Normal file
|
|
@ -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;
|
||||
6
dist/cjs/messages/build.js
vendored
6
dist/cjs/messages/build.js
vendored
|
|
@ -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). */
|
||||
|
|
|
|||
4
dist/cjs/topics.js
vendored
4
dist/cjs/topics.js
vendored
|
|
@ -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;
|
||||
|
|
|
|||
12
dist/cjs/transport.js
vendored
12
dist/cjs/transport.js
vendored
|
|
@ -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: "<root>/<endpointId>/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;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue