Skip to content

Commit fbe89e6

Browse files
committed
streams: add stream.pipe
pipe is similar to pipeline however it supports stream composition. Refs: nodejs#32020
1 parent 4e17ffc commit fbe89e6

File tree

2 files changed

+100
-0
lines changed

2 files changed

+100
-0
lines changed

lib/internal/streams/pipelinify.js

+98
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
'use strict';
2+
3+
const pipeline = require('internal/streams/pipeline');
4+
const Duplex = require('internal/streams/duplex');
5+
const { destroyer } = require('internal/streams/destroy');
6+
7+
module.exports = function pipe(...streams) {
8+
let onclose;
9+
let ret;
10+
11+
const r = pipeline(streams, function(err) {
12+
if (onclose) {
13+
const cb = onclose;
14+
onclose = null;
15+
cb(err);
16+
} else {
17+
ret.destroy(err);
18+
}
19+
});
20+
const w = streams[0];
21+
22+
const writable = w.writable;
23+
const readable = r.readable;
24+
const objectMode = w.readableObjectMode;
25+
26+
ret = new Duplex({
27+
writable,
28+
readable,
29+
objectMode,
30+
highWaterMark: 1
31+
});
32+
33+
if (writable) {
34+
let ondrain;
35+
let onfinish;
36+
37+
ret._write = function(chunk, encoding, callback) {
38+
if (w.write(chunk, encoding)) {
39+
callback();
40+
} else {
41+
ondrain = callback;
42+
}
43+
};
44+
45+
ret._final = function(chunk, encoding, callback) {
46+
w.end(chunk, encoding);
47+
onfinish = callback;
48+
};
49+
50+
ret.on('drain', function () {
51+
if (ondrain) {
52+
const cb = ondrain;
53+
ondrain = null;
54+
cb();
55+
}
56+
});
57+
58+
ret.on('finish', function () {
59+
if (onfinish) {
60+
const cb = onfinish
61+
onfinish = null;
62+
cb();
63+
}
64+
});
65+
}
66+
67+
if (readable) {
68+
let onreadable;
69+
70+
r.on('readable', function () {
71+
if (onreadable) {
72+
const cb = onreadable
73+
onreadable = null;
74+
cb();
75+
}
76+
});
77+
78+
ret._read = function () {
79+
while (true) {
80+
const buf = r.read();
81+
82+
if (buf === null) {
83+
onreadable = ret._read;
84+
return;
85+
}
86+
87+
if (!ret.push(buf)) {
88+
return;
89+
}
90+
}
91+
};
92+
}
93+
94+
ret._destroy = function(err, callback) {
95+
onclose = callback;
96+
r.destroy(err);
97+
};
98+
}

lib/stream.js

+2
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ const {
3030
} = require('internal/util');
3131

3232
const pipeline = require('internal/streams/pipeline');
33+
const pipelinify = require('internal/streams/pipelinify');
3334
const eos = require('internal/streams/end-of-stream');
3435
const internalBuffer = require('internal/buffer');
3536

@@ -42,6 +43,7 @@ Stream.Duplex = require('internal/streams/duplex');
4243
Stream.Transform = require('internal/streams/transform');
4344
Stream.PassThrough = require('internal/streams/passthrough');
4445
Stream.pipeline = pipeline;
46+
Stream.pipelinify = pipelinify;
4547
const { addAbortSignal } = require('internal/streams/add-abort-signal');
4648
Stream.addAbortSignal = addAbortSignal;
4749
Stream.finished = eos;

0 commit comments

Comments
 (0)