Alex2Node/dist/cjs/dispatcher.js
David 3a07c861b1 registry: SceneController events, DoorbellEventSource, SimpleEventSource, TimeHoldController, InventoryLevelSensor, WakeOnLANController
Six descriptors written from their pages replace the last stubs of tiers 1 and 2. A scene answers Activate and
Deactivate through ctx.respond() with ActivationStarted and DeactivationStarted, the time and the cause filled
in. device.raise(descriptor, name, payload) publishes DoorbellPress and the Event of a button on <root>/event
with the endpoint and a new messageId; it throws a MessageError for an interface or instance the device did
not declare, an event that answers a directive, a payload that does not fit and a message over 16000 bytes.
TurnOn of a device with WakeOnLANController is deferred without the warning.

On the wire: a doorbell has no properties object and proactivelyReported on the capability; SimpleEventSource
is version 1.0, InventoryLevelSensor and WakeOnLANController version 3 (1.5.2: 1). A scene declared without
options is announced as before. Alex2MQTT has no topic yet for the WakeUp event.
24 examples of the six pages are saved as fixtures. 267 tests pass, 241 before.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-28 21:01:53 +00:00

324 lines
16 KiB
JavaScript

"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.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, (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 }));
}
// 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_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;