Save
This commit is contained in:
@@ -19,3 +19,4 @@ export { replace } from "./replace";
|
||||
export { split } from "./split";
|
||||
export { stringify } from "./stringify";
|
||||
export { unbatch } from "./unbatch";
|
||||
export { compose } from "./compose";
|
||||
|
||||
45
src/functions/compose.ts
Normal file
45
src/functions/compose.ts
Normal file
@@ -0,0 +1,45 @@
|
||||
import { Transform, Writable, Pipe, WritableOptions } from "stream";
|
||||
|
||||
class Compose extends Writable implements Pipe {
|
||||
private head: Writable | Transform;
|
||||
private tail: Writable | Transform;
|
||||
constructor(
|
||||
streams: Array<Transform | Writable>,
|
||||
options?: WritableOptions,
|
||||
) {
|
||||
super(options);
|
||||
if (streams.length < 2) {
|
||||
throw new Error("Cannot compose 1 or less streams");
|
||||
}
|
||||
this.head = streams[0];
|
||||
for (let i = 1; i < streams.length; i++) {
|
||||
streams[i - 1].pipe(streams[i]);
|
||||
}
|
||||
this.tail = streams[streams.length - 1];
|
||||
}
|
||||
|
||||
public pipe<T extends NodeJS.WritableStream>(
|
||||
destination: T,
|
||||
options: { end?: boolean } | undefined,
|
||||
) {
|
||||
return this.tail.pipe(
|
||||
destination,
|
||||
options,
|
||||
);
|
||||
}
|
||||
|
||||
public _write(chunk: any, enc: string, cb: any) {
|
||||
this.head.write(chunk.toString ? chunk.toString() : chunk, cb);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Return a Readable stream of readable streams concatenated together
|
||||
* @param streams Readable streams to concatenate
|
||||
*/
|
||||
export function compose(
|
||||
streams: Array<Transform | Writable>,
|
||||
options?: WritableOptions,
|
||||
): Compose {
|
||||
return new Compose(streams, options);
|
||||
}
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Readable, Writable, Transform, Duplex } from "stream";
|
||||
import { Readable, Writable, WritableOptions, Transform, Duplex } from "stream";
|
||||
import { ChildProcess } from "child_process";
|
||||
import * as baseFunctions from "./baseFunctions";
|
||||
|
||||
@@ -287,3 +287,13 @@ export function accumulatorBy<T, S extends FlushStrategy>(
|
||||
) {
|
||||
return baseFunctions.accumulatorBy(batchRate, flushStrategy, iteratee);
|
||||
}
|
||||
|
||||
export function compose(
|
||||
streams: Array<Writable | Transform>,
|
||||
options?: WritableOptions,
|
||||
) {
|
||||
return baseFunctions.compose(
|
||||
streams,
|
||||
options,
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user