import { MaybeAsync, createPipeline, Middleware } from 'farrow-pipeline'; import type { AsyncWorker, AsyncWorkers } from './async'; import type { RunWorkflowOptions } from './sync'; const PARALLEL_WORKFLOW_SYMBOL = Symbol('PARALLEL_WORKFLOW_SYMBOL'); export type ParallelWorkflow = { run: (input: I, options?: RunWorkflowOptions) => Promise; use: (...I: AsyncWorkers) => ParallelWorkflow; [PARALLEL_WORKFLOW_SYMBOL]: true; }; export type ParallelWorkflow2Worker> = W extends ParallelWorkflow ? AsyncWorker : never; export type ParallelWorkflowRecord = Record>; export type ParallelWorkflows2Workers< // eslint-disable-next-line @typescript-eslint/no-invalid-void-type PS extends ParallelWorkflowRecord | void, > = { [K in keyof PS]: PS[K] extends ParallelWorkflow ? ParallelWorkflow2Worker : // 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 ParallelWorkflows2AsyncWorkers< // eslint-disable-next-line @typescript-eslint/no-invalid-void-type PS extends ParallelWorkflowRecord | void, > = { [K in keyof PS]: PS[K] extends ParallelWorkflow ? ParallelWorkflow2Worker : // 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 RunnerFromParallelWorkflow> = W extends ParallelWorkflow ? ParallelWorkflow['run'] : never; export type ParallelWorkflows2Runners< // eslint-disable-next-line @typescript-eslint/no-invalid-void-type PS extends ParallelWorkflowRecord | void, > = { [K in keyof PS]: PS[K] extends ParallelWorkflow ? RunnerFromParallelWorkflow : // 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 isParallelWorkflow = ( input: any, ): input is ParallelWorkflow => Boolean(input?.[PARALLEL_WORKFLOW_SYMBOL]); export const createParallelWorkflow = < // eslint-disable-next-line @typescript-eslint/no-invalid-void-type I = void, O = unknown, >(): ParallelWorkflow => { const pipeline = createPipeline[]>(); const use: ParallelWorkflow['use'] = (...input) => { pipeline.use(...input.map(mapParallelWorkerToAsyncMiddleware)); return workflow; }; const run: ParallelWorkflow['run'] = async (input, options) => // eslint-disable-next-line promise/prefer-await-to-then Promise.all(pipeline.run(input, { ...options, onLast: () => [] })).then( result => result.filter(Boolean), ); const workflow: ParallelWorkflow = { ...pipeline, run, use, [PARALLEL_WORKFLOW_SYMBOL]: true as const, }; return workflow; }; const mapParallelWorkerToAsyncMiddleware = (worker: AsyncWorker): Middleware[]> => (input, next) => [worker(input), ...next(input)];