1 | "use strict";
|
2 | Object.defineProperty(exports, "__esModule", { value: true });
|
3 | exports.writableFork = void 0;
|
4 | const stream_1 = require("stream");
|
5 | const __1 = require("../..");
|
6 |
|
7 |
|
8 |
|
9 |
|
10 |
|
11 |
|
12 | function writableFork(chains, opt) {
|
13 | const readables = [];
|
14 | const allChainsDone = Promise.all(chains.map(async (chain) => {
|
15 | const readable = __1.readableCreate();
|
16 | readables.push(readable);
|
17 | return await __1._pipeline([readable, ...chain]);
|
18 | })).catch(err => {
|
19 | console.error(err);
|
20 | throw err;
|
21 | });
|
22 | return new stream_1.Writable({
|
23 | objectMode: true,
|
24 | ...opt,
|
25 | write(chunk, _encoding, cb) {
|
26 |
|
27 |
|
28 | readables.forEach(readable => readable.push(chunk));
|
29 | cb();
|
30 | },
|
31 | async final(cb) {
|
32 | try {
|
33 |
|
34 | readables.forEach(readable => readable.push(null));
|
35 | console.log(`writableFork.final is waiting for all chains to be done`);
|
36 | await allChainsDone;
|
37 | console.log(`writableFork.final all chains done`);
|
38 | cb();
|
39 | }
|
40 | catch (err) {
|
41 | cb(err);
|
42 | }
|
43 | },
|
44 | });
|
45 | }
|
46 | exports.writableFork = writableFork;
|