Skip to content
File

Blob: src/node/internal/streams_compose.js

javascript190 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4//
5// Copyright Joyent, Inc. and other Node contributors.
6//
7// Permission is hereby granted, free of charge, to any person obtaining a
8// copy of this software and associated documentation files (the
9// "Software"), to deal in the Software without restriction, including
10// without limitation the rights to use, copy, modify, merge, publish,
11// distribute, sublicense, and/or sell copies of the Software, and to permit
12// persons to whom the Software is furnished to do so, subject to the
13// following conditions:
14//
15// The above copyright notice and this permission notice shall be included
16// in all copies or substantial portions of the Software.
17//
18// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS
19// OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF
20// MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN
21// NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM,
22// DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR
23// OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE
24// USE OR OTHER DEALINGS IN THE SOFTWARE.
25 
26import { pipeline } from 'node-internal:streams_pipeline';
27import { Duplex } from 'node-internal:streams_duplex';
28import {
29 Readable as ReadableConstructor,
30 from,
31} from 'node-internal:streams_readable';
32import {
33 isNodeStream,
34 isReadable,
35 isWritable,
36} from 'node-internal:streams_util';
37import { destroyer } from 'node-internal:streams_destroy';
38import {
39 AbortError,
40 ERR_INVALID_ARG_VALUE,
41 ERR_MISSING_ARGS,
42} from 'node-internal:internal_errors';
43 
44export function compose(...streams) {
45 if (streams.length === 0) {
46 throw new ERR_MISSING_ARGS('streams');
47 }
48 if (streams.length === 1) {
49 return from(Duplex, streams[0]);
50 }
51 const orgStreams = [...streams];
52 if (typeof streams[0] === 'function') {
53 streams[0] = from(Duplex, streams[0]);
54 }
55 if (typeof streams[streams.length - 1] === 'function') {
56 const idx = streams.length - 1;
57 streams[idx] = from(Duplex, streams[idx]);
58 }
59 for (let n = 0; n < streams.length; ++n) {
60 if (!isNodeStream(streams[n])) {
61 // TODO(ronag): Add checks for non streams.
62 continue;
63 }
64 if (n < streams.length - 1 && !isReadable(streams[n])) {
65 throw new ERR_INVALID_ARG_VALUE(
66 `streams[${n}]`,
67 orgStreams[n],
68 'must be readable'
69 );
70 }
71 if (n > 0 && !isWritable(streams[n])) {
72 throw new ERR_INVALID_ARG_VALUE(
73 `streams[${n}]`,
74 orgStreams[n],
75 'must be writable'
76 );
77 }
78 }
79 let ondrain;
80 let onfinish;
81 let onreadable;
82 let onclose;
83 let d;
84 function onfinished(err) {
85 const cb = onclose;
86 onclose = null;
87 if (cb) {
88 cb(err);
89 } else if (err) {
90 d.destroy(err);
91 } else if (!readable && !writable) {
92 d.destroy();
93 }
94 }
95 const head = streams[0];
96 const tail = pipeline(streams, onfinished);
97 const writable = !!isWritable(head);
98 const readable = !!isReadable(tail);
99 
100 // TODO(ronag): Avoid double buffering.
101 // Implement Writable/Readable/Duplex traits.
102 // See, https://github.com/nodejs/node/pull/33515.
103 d = new Duplex({
104 // TODO (ronag): highWaterMark?
105 writableObjectMode: !!(
106 head !== null &&
107 head !== undefined &&
108 head.writableObjectMode
109 ),
110 readableObjectMode: !!(
111 tail !== null &&
112 tail !== undefined &&
113 tail.writableObjectMode
114 ),
115 writable,
116 readable,
117 });
118 if (writable) {
119 const w = head;
120 d._write = function (chunk, encoding, callback) {
121 if (head.write(chunk, encoding)) {
122 callback();
123 } else {
124 ondrain = callback;
125 }
126 };
127 d._final = function (callback) {
128 w.end();
129 onfinish = callback;
130 };
131 w.on('drain', function () {
132 if (ondrain) {
133 const cb = ondrain;
134 ondrain = null;
135 cb();
136 }
137 });
138 tail.on('finish', function () {
139 if (onfinish) {
140 const cb = onfinish;
141 onfinish = null;
142 cb();
143 }
144 });
145 }
146 if (readable) {
147 tail.on('readable', function () {
148 if (onreadable) {
149 const cb = onreadable;
150 onreadable = null;
151 cb();
152 }
153 });
154 tail.on('end', function () {
155 d.push(null);
156 });
157 d._read = function (_n) {
158 while (true) {
159 const buf = tail.read();
160 if (buf === null) {
161 onreadable = d._read;
162 return;
163 }
164 if (!d.push(buf)) {
165 return;
166 }
167 }
168 };
169 }
170 d._destroy = function (err, callback) {
171 if (!err && onclose !== null) {
172 err = new AbortError();
173 }
174 onreadable = null;
175 ondrain = null;
176 onfinish = null;
177 if (onclose === null) {
178 callback(err);
179 } else {
180 onclose = callback;
181 destroyer(tail, err);
182 }
183 };
184 return d;
185}
186 
187ReadableConstructor.prototype.compose = function (...streams) {
188 return compose(this, ...streams);
189};