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