Alex2Node/dist/esm/dispatcher.js
David feff828ec1 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>
2026-09-28 19:49:41 +00:00

282 lines
14 KiB
JavaScript

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);
}
}
}