// eslint-disable-next-line @typescript-eslint/triple-slash-reference /// import { Adapter, AdapterEnhancer, AdapterOptions, AdapterState, ReadyState, SealedAdapter } from '../types/AdapterTypes'; import Observable, { Subscription } from 'core-js/features/observable'; import createAsyncIterableQueue, { AsyncIterableQueue } from './utils/createAsyncIterableQueue'; import EventTarget, { Event } from 'event-target-shim'; import sealAdapter from './sealAdapter'; const DEFAULT_ENHANCER: AdapterEnhancer = (next) => (options) => next(options); export default function createAdapter( options: AdapterOptions = {}, enhancer: AdapterEnhancer = DEFAULT_ENHANCER ): SealedAdapter { const mutableAdapterState: TAdapterState = {} as TAdapterState; let sealed = false; let activeSubscription: Subscription; const adapter = enhancer((): Adapter => { const eventTarget = new EventTarget(); const ingressQueues: AsyncIterableQueue[] = []; let readyStatePropertyValue = ReadyState.CONNECTING; return { addEventListener: eventTarget.addEventListener.bind(eventTarget), dispatchEvent: eventTarget.dispatchEvent.bind(eventTarget), removeEventListener: eventTarget.removeEventListener.bind(eventTarget), activities: ({ signal } = {}): AsyncIterable => { const queue = createAsyncIterableQueue({ signal }); ingressQueues.push(queue); signal && signal.addEventListener('abort', () => { const index = ingressQueues.indexOf(queue); ~index || ingressQueues.splice(index, 1); }); return queue.iterable; }, close: () => { ingressQueues.forEach((ingressQueue) => ingressQueue.end()); ingressQueues.splice(0, Infinity); activeSubscription && activeSubscription.unsubscribe(); activeSubscription = null; }, // Egress middleware API egress: (): Promise => { return Promise.resolve(); }, getState: (name: keyof TAdapterState) => { return mutableAdapterState[name]; }, getReadyState: () => readyStatePropertyValue, // Ingress middleware API ingress: (activity) => { ingressQueues.forEach((ingressQueue) => ingressQueue.push(activity)); }, setState: (name: keyof TAdapterState, value: any) => { if (sealed && !(name in mutableAdapterState)) { throw new Error(`Cannot set config "${String(name)}" because it was not set before being sealed.`); } // TODO: Fix this typing // mutableAdapterState[name] = value; (mutableAdapterState as any)[name] = value; }, setReadyState: (readyState: ReadyState) => { if (readyState === readyStatePropertyValue) { return; } if (readyStatePropertyValue === ReadyState.CLOSED) { throw new Error('Cannot change "readyState" after it is CLOSED.'); } else if ( readyState !== ReadyState.CLOSED && readyState !== ReadyState.CONNECTING && readyState !== ReadyState.OPEN ) { throw new Error('"readyState" must be either CLOSED, CONNECTING or OPEN.'); } readyStatePropertyValue = readyState; if (readyState === ReadyState.CLOSED) { activeSubscription && activeSubscription.unsubscribe(); activeSubscription = null; } eventTarget.dispatchEvent(new Event(readyState === ReadyState.OPEN ? 'open' : 'error')); }, subscribe: (observable: Observable | false) => { activeSubscription && activeSubscription.unsubscribe(); activeSubscription = null; if (!observable) { return; } let subscription: Subscription; observable.subscribe({ start(thisSubscription: Subscription) { activeSubscription = thisSubscription; subscription = thisSubscription; }, complete() { if (activeSubscription === subscription) { activeSubscription = null; } }, error() { if (activeSubscription === subscription) { activeSubscription = null; } // TODO: Propagate the error to fail the adapter. // ingressQueues.forEach(ingressQueue => ingressQueue.push(error)); }, next(value: TActivity) { adapter.ingress(value); } }); } }; })(options); if (Object.getPrototypeOf(adapter) !== Object.prototype) { throw new Error('Object returned from enhancer must not be a class object.'); } const sealedAdapter = sealAdapter(adapter, mutableAdapterState); sealed = true; return sealedAdapter; }