A broker that mirrors another delivers every message twice, and the bridge answered both copies of a discovery request. A request is now remembered by its messageId for 60 s, as a directive is, and its copy is dropped with a log line. The messageId is read from the header of the Discover directive or from the request itself when that is the header alone. A request without a messageId is answered every time, as before. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
401 lines
19 KiB
TypeScript
401 lines
19 KiB
TypeScript
// From a message on <root>/<endpointId>/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<string, string> };
|
|
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<T = unknown> {
|
|
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 <root>/<endpointId>/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<string, unknown> }): Promise<PublishResult>;
|
|
/** Answer ReportState with a StateReport, its properties as for respond(). */
|
|
report(fill?: Fill): Promise<PublishResult>;
|
|
/** Send a DeferredResponse now, when the answer takes longer than the 7 s Alex2MQTT waits for it. */
|
|
defer(estimatedDeferralInSeconds?: number): Promise<PublishResult>;
|
|
/** Answer with an ErrorResponse. An error that is not an AlexaError is sent as INTERNAL_ERROR. */
|
|
error(type: string, message: string, extra?: Record<string, unknown>): Promise<PublishResult>;
|
|
error(err: Error): Promise<PublishResult>;
|
|
}
|
|
|
|
/** What a handler returns is awaited and otherwise ignored; what it throws answers the directive. */
|
|
export type DirectiveHandler<T = unknown> = (ctx: DirectiveContext<T>) => 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<typeof setTimeout>);
|
|
},
|
|
};
|
|
|
|
/** 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<any> {
|
|
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<string, unknown> } = {}): Promise<PublishResult> {
|
|
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<PublishResult> {
|
|
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<PublishResult> {
|
|
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: string | Error, message = "", extra: Record<string, unknown> = {}): Promise<PublishResult> {
|
|
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 }));
|
|
}
|
|
|
|
// An interface of the device documents a DeferredResponse for this directive of another one
|
|
private deferredByAnother(): boolean {
|
|
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
|
|
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<PublishResult> {
|
|
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<string, Exchange>();
|
|
private readonly arrived = new Map<string, number>();
|
|
private readonly deferralNoted = new Set<string>();
|
|
|
|
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<PublishResult> {
|
|
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`);
|
|
}
|
|
|
|
/**
|
|
* true for a message that came before, within REMEMBERED_MS: a broker that mirrors another delivers each message
|
|
* twice. endpointId is "" for a message to the root, a discovery request.
|
|
*/
|
|
arrivedBefore(endpointId: string, messageId: string, now = Date.now()): boolean {
|
|
this.forget(now);
|
|
const id = key(endpointId, messageId);
|
|
if (this.arrived.has(id)) return true;
|
|
this.arrived.set(id, now);
|
|
return false;
|
|
}
|
|
|
|
/** The text of a message on the directive topic of endpointId. Resolves when the handlers are done. */
|
|
async dispatch(endpointId: string, text: string): Promise<void> {
|
|
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();
|
|
if (typeof messageId === "string" && this.arrivedBefore(endpointId, messageId, receivedAt)) {
|
|
log(`the directive ${messageId} for ${endpointId} is dropped: it arrived before`);
|
|
return;
|
|
}
|
|
|
|
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<any> {
|
|
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<void> {
|
|
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<PublishResult> {
|
|
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);
|
|
}
|
|
}
|
|
}
|