// node/pipeline.ts /* * Copyright (c) 2021-2026 Check Digit, LLC * * This code is licensed under the MIT license (see LICENSE.txt for details). */ import stream, { Duplex, Readable, Writable } from 'node:stream'; import { promisify } from 'node:util'; import debug from 'debug'; import type { Asyncerator } from '../asyncerator.ts'; const log = debug('asyncerator:pipeline'); const promisifiedPipeline = promisify(stream.pipeline) as ( ...streams: unknown[] ) => Promise; type CallbackPipeline = ( ...streams: [...unknown[], (error?: NodeJS.ErrnoException | null) => void] ) => NodeJS.WritableStream; const callbackPipeline = stream.pipeline as unknown as CallbackPipeline; export type PipelineSource = | string | Readable | Iterable | AsyncIterable | Asyncerator; export type PipelineTransformer = Duplex | ((input: Asyncerator) => Asyncerator); export interface PipelineOptions { signal: AbortSignal; } /* eslint-disable max-params */ /** * Overloads connect element types between function stages and infer the return type from the sink. * The overloads cover up to ten transforms. Node stream stages do not preserve these element-type guarantees * because their chunk types are not generic. Options are exposed only for promise-returning function sinks. */ // zero transforms export default function ( source: PipelineSource, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, sink: Writable, ): Promise; // 1 transform export default function ( source: PipelineSource, transform1: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, sink: Writable, ): Promise; // 2 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, sink: Writable, ): Promise; // 3 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, sink: Writable, ): Promise; // 4 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, sink: Writable, ): Promise; // 5 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, sink: Writable, ): Promise; // 6 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, sink: Writable, ): Promise; // 7 transforms export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, sink: Writable, ): Promise; // 8 transforms export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, sink: Writable, ): Promise; // 9 transforms export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, T8, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, T8, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function ( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, sink: Writable, ): Promise; // 10 transforms export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, T8, T9, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, transform10: PipelineTransformer, sink: (input: Asyncerator) => Promise, options?: PipelineOptions, ): Promise; export default function < Source, Sink, TransformSink, T1, T2, T3, T4, T5, T6, T7, T8, T9, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, transform10: PipelineTransformer, sink: Duplex | ((input: Asyncerator) => AsyncIterable), ): Readable; export default function < Source, TransformSink, T1, T2, T3, T4, T5, T6, T7, T8, T9, >( source: PipelineSource, transform1: PipelineTransformer, transform2: PipelineTransformer, transform3: PipelineTransformer, transform4: PipelineTransformer, transform5: PipelineTransformer, transform6: PipelineTransformer, transform7: PipelineTransformer, transform8: PipelineTransformer, transform9: PipelineTransformer, transform10: PipelineTransformer, sink: Writable, ): Promise; /* eslint-enable max-params */ /** * Wrap Node's callback-based stream.pipeline with overloads for asyncerator operators. * Return a promise for an async function or writable-only sink, or a Readable for a duplex or async generator sink. * Node also provides node:stream/promises.pipeline, which always returns a promise. * * These overloads cover the library's supported composition patterns and do not expose every native pipeline option. * * @param argumentList */ export default function ( ...argumentList: unknown[] ): Promise | Readable { let options: PipelineOptions | undefined = argumentList.at( -1, ) as PipelineOptions; if (!( Object.keys(options).length === 1 && Object.keys(options)[0] === 'signal' )) { options = undefined; } const sink = argumentList[ argumentList.length - (options === undefined ? 1 : 2) ] as object; /** * The sink is an async function, so return a promise */ if (sink.constructor.name === 'AsyncFunction') { return promisifiedPipeline(...argumentList) as Promise; } /** * The sink is WritableStream, so return a promise. Reject on error. */ if ((sink as Writable).writable && !(sink as Readable).readable) { return new Promise((resolve, reject) => { callbackPipeline(...argumentList, (error) => { if (error === undefined || error === null) { resolve(); } else { reject(error); } }); }); } /** * The sink is a transform, i.e., AsyncGenerator-like, or a Duplex stream. In this case, we return a Readable. */ return callbackPipeline(...argumentList, (error) => { if (error) { log('error', error); } else { log('complete'); } }) as unknown as Readable; }