import { MaybeAsync, createAsyncPipeline, Middleware } from 'farrow-pipeline'; import type { RunWorkflowOptions } from './sync'; const ASYNC_WORKFLOW_SYMBOL = Symbol('ASYNC_WORKFLOW_SYMBOL'); export type AsyncWorker = (I: I) => MaybeAsync; export type AsyncWorkers = AsyncWorker[]; export type AsyncWorkflow = { run: (input: I, options?: RunWorkflowOptions) => MaybeAsync; use: (...I: AsyncWorkers) => AsyncWorkflow; [ASYNC_WORKFLOW_SYMBOL]: true; }; export type AsyncWorkflow2AsyncWorker> = W extends AsyncWorkflow ? AsyncWorker : never; export type AsyncWorkflowRecord = Record>; // eslint-disable-next-line @typescript-eslint/no-invalid-void-type export type AsyncWorkflows2AsyncWorkers = { [K in keyof PS]: PS[K] extends AsyncWorkflow ? AsyncWorkflow2AsyncWorker : // eslint-disable-next-line @typescript-eslint/no-invalid-void-type PS[K] extends void ? // eslint-disable-next-line @typescript-eslint/no-invalid-void-type void : never; }; export type RunnerFromAsyncWorkflow> = W extends AsyncWorkflow ? AsyncWorkflow['run'] : never; // eslint-disable-next-line @typescript-eslint/no-invalid-void-type export type AsyncWorkflows2Runners = { [K in keyof PS]: PS[K] extends AsyncWorkflow ? RunnerFromAsyncWorkflow : // eslint-disable-next-line @typescript-eslint/no-invalid-void-type PS[K] extends void ? // eslint-disable-next-line @typescript-eslint/no-invalid-void-type void : never; }; export const isAsyncWorkflow = (input: any): input is AsyncWorkflow => Boolean(input?.[ASYNC_WORKFLOW_SYMBOL]); // eslint-disable-next-line @typescript-eslint/no-invalid-void-type export const createAsyncWorkflow = (): AsyncWorkflow< I, O > => { const pipeline = createAsyncPipeline(); const use: AsyncWorkflow['use'] = (...input) => { pipeline.use(...input.map(mapAsyncWorkerToAsyncMiddleware)); return workflow; }; const run: AsyncWorkflow['run'] = async (input, options) => { const result = pipeline.run(input, { ...options, onLast: () => [] }); if (isPromise(result)) { // eslint-disable-next-line @typescript-eslint/no-shadow,promise/prefer-await-to-then return result.then(result => result.filter(Boolean)); } else { return result.filter(Boolean); } }; const workflow: AsyncWorkflow = { ...pipeline, use, run, [ASYNC_WORKFLOW_SYMBOL]: true as const, }; return workflow; }; const mapAsyncWorkerToAsyncMiddleware = (worker: AsyncWorker): Middleware> => async (input, next) => [await worker(input), ...(await next(input))]; function isPromise(obj: any): obj is Promise { /* eslint-disable promise/prefer-await-to-then */ return ( Boolean(obj) && (typeof obj === 'object' || typeof obj === 'function') && typeof obj.then === 'function' ); /* eslint-enable promise/prefer-await-to-then */ }