Something went wrong. Try again.
TypeScript client for accessing TLE Community.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110import 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<T = any> { /** The request that produced this response. */ request: ServerRequest;
/** The resolved response - exactly what the caller is about to receive. */ response: ServerResponse<T>;
/** 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<ServerResponse<T>>;}
/** * 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<any> | Promise<void | ServerResponse<any>>;
/** 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<ResponseHandler>();
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<T>(event: ResponseEvent<T>): Promise<ServerResponse<T>> { let response = event.response;
for (const handler of [...this.handlers]) { try { const handled = await handler({ ...event, response }); if (handled) response = handled as ServerResponse<T>; } 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;