import Lacuna from '../lacuna'; import { BLOCKING_ERROR_CODES } from './constants'; import type { ServerError, ServerRequest, ServerResponse } from './server'; /** * Everything a subscriber needs to decide what to do about a response: what * was asked for, what came back, and a way to ask again. */ export interface ResponseEvent { /** The request that produced this response. */ request: ServerRequest; /** The resolved response - exactly what the caller is about to receive. */ response: ServerResponse; /** 0 on the original call, incremented once per retry. */ attempt: number; /** * Re-issues `request`. The resulting response is published as a fresh event * (with `attempt + 1`), so a handler that retries on some condition is * called again if that condition persists - no loop needed on its side. * * For a blocking error code, returning this from a handler makes the retry's * response the one the original caller receives. */ retry: () => Promise>; } /** * A response subscriber. Return nothing to observe only; return a * `ServerResponse` (typically `await event.retry()`) to replace what the * original caller receives. Replacement only applies to blocking error codes * (see BLOCKING_ERROR_CODES) - those are the only responses the transport * waits on. */ export type ResponseHandler = ( event: ResponseEvent ) => void | ServerResponse | Promise>; /** Unsubscribes a previously registered handler. Safe to call more than once. */ export type Unsubscribe = () => void; export const isBlockingError = (error: ServerError | undefined) => error !== undefined && BLOCKING_ERROR_CODES.includes(error.code); class Responses { lacuna: Lacuna; handlers = new Set(); constructor(lacuna: Lacuna) { this.lacuna = lacuna; } subscribe(handler: ResponseHandler): Unsubscribe { this.handlers.add(handler); return () => { this.handlers.delete(handler); }; } /** * Publishes an event without waiting for the handlers. Used for every * response the caller doesn't need protecting from, so the happy path pays * no handler latency. */ dispatch(event: ResponseEvent) { for (const handler of [...this.handlers]) { try { Promise.resolve(handler(event)).catch((e) => this.warn(e)); } catch (e) { this.warn(e); } } } /** * Publishes an event and waits for each handler in turn, in registration * order. A handler returning a response replaces the current one, and each * subsequent handler is given that replacement rather than the original. * * Every handler is called even after the blocking error has been resolved - * subscribers are promised *all* responses, and one that only observes * would otherwise silently miss exactly the attempts that mattered. Passing * the updated response along is what stops a second recovery handler from * redundantly retrying: handlers key off `response.error`, which by then is * gone. */ async dispatchBlocking(event: ResponseEvent): Promise> { let response = event.response; for (const handler of [...this.handlers]) { try { const handled = await handler({ ...event, response }); if (handled) response = handled as ServerResponse; } catch (e) { this.warn(e); } } return response; } private warn(e: unknown) { this.lacuna.log.warn('A response handler failed', e as object); } } export default Responses;