Skip to content
File

Blob: src/node/internal/streams_pipeline.js

javascript433 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 
26/* TODO: the following is adopted code, enabling linting one day */
27/* eslint-disable */
28 
29import {
30 isIterable,
31 isReadable,
32 isReadableNodeStream,
33 isNodeStream,
34} from 'node-internal:streams_util';
35import { eos } from 'node-internal:streams_end_of_stream';
36import { destroyer as destroyerImpl } from 'node-internal:streams_destroy';
37import { once } from 'node-internal:internal_http_util';
38 
39import { nextTick } from 'node-internal:internal_process';
40import { PassThrough } from 'node-internal:streams_transform';
41import { Duplex } from 'node-internal:streams_duplex';
42import { Readable, from } from 'node-internal:streams_readable';
43import {
44 aggregateTwoErrors,
45 ERR_INVALID_ARG_TYPE,
46 ERR_INVALID_RETURN_VALUE,
47 ERR_MISSING_ARGS,
48 ERR_STREAM_DESTROYED,
49 ERR_STREAM_PREMATURE_CLOSE,
50 AbortError,
51} from 'node-internal:internal_errors';
52import {
53 validateFunction,
54 validateAbortSignal,
55} from 'node-internal:validators';
56 
57function destroyer(stream, reading, writing) {
58 let finished = false;
59 stream.on('close', () => {
60 finished = true;
61 });
62 const cleanup = eos(
63 stream,
64 {
65 readable: reading,
66 writable: writing,
67 },
68 (err) => {
69 finished = !err;
70 }
71 );
72 return {
73 destroy: (err) => {
74 if (finished) return;
75 finished = true;
76 destroyerImpl(stream, err || new ERR_STREAM_DESTROYED('pipe'));
77 },
78 cleanup,
79 };
80}
81 
82function popCallback(streams) {
83 // Streams should never be an empty array. It should always contain at least
84 // a single stream. Therefore optimize for the average case instead of
85 // checking for length === 0 as well.
86 validateFunction(streams[streams.length - 1], 'streams[stream.length - 1]');
87 return streams.pop();
88}
89 
90function makeAsyncIterable(val) {
91 if (isIterable(val)) {
92 return val;
93 } else if (isReadableNodeStream(val)) {
94 // Legacy streams are not Iterable.
95 return fromReadable(val);
96 }
97 throw new ERR_INVALID_ARG_TYPE(
98 'val',
99 ['Readable', 'Iterable', 'AsyncIterable'],
100 val
101 );
102}
103 
104async function* fromReadable(val) {
105 yield* Readable.prototype[Symbol.asyncIterator].call(val);
106}
107 
108async function pump(iterable, writable, finish, { end }) {
109 let error;
110 let onresolve = null;
111 const resume = (err) => {
112 if (err) {
113 error = err;
114 }
115 if (onresolve) {
116 const callback = onresolve;
117 onresolve = null;
118 callback();
119 }
120 };
121 const wait = () => {
122 return new Promise((resolve, reject) => {
123 if (error) {
124 reject(error);
125 } else {
126 onresolve = () => {
127 if (error) {
128 reject(error);
129 } else {
130 resolve();
131 }
132 };
133 }
134 });
135 };
136 writable.on('drain', resume);
137 const cleanup = eos(
138 writable,
139 {
140 readable: false,
141 },
142 resume
143 );
144 try {
145 if (writable.writableNeedDrain) {
146 await wait();
147 }
148 for await (const chunk of iterable) {
149 if (!writable.write(chunk)) {
150 await wait();
151 }
152 }
153 if (end) {
154 writable.end();
155 }
156 await wait();
157 finish();
158 } catch (err) {
159 finish(error !== err ? aggregateTwoErrors(error, err) : err);
160 } finally {
161 cleanup();
162 writable.off('drain', resume);
163 }
164}
165 
166export function pipeline(...streams) {
167 return pipelineImpl(streams, once(popCallback(streams)));
168}
169 
170export function pipelineImpl(streams, callback, opts) {
171 if (streams.length === 1 && Array.isArray(streams[0])) {
172 streams = streams[0];
173 }
174 if (streams.length < 2) {
175 throw new ERR_MISSING_ARGS('streams');
176 }
177 const ac = new AbortController();
178 const signal = ac.signal;
179 const outerSignal = opts?.signal;
180 
181 // Need to cleanup event listeners if last stream is readable
182 // https://github.com/nodejs/node/issues/35452
183 const lastStreamCleanup = [];
184 validateAbortSignal(outerSignal, 'options.signal');
185 function abort() {
186 finishImpl(new AbortError());
187 }
188 outerSignal === null || outerSignal === undefined
189 ? undefined
190 : outerSignal.addEventListener('abort', abort);
191 let error;
192 let value;
193 const destroys = [];
194 let finishCount = 0;
195 function finish(err) {
196 finishImpl(err, --finishCount === 0);
197 }
198 function finishImpl(err, final) {
199 if (err && (!error || error.code === 'ERR_STREAM_PREMATURE_CLOSE')) {
200 error = err;
201 }
202 if (!error && !final) {
203 return;
204 }
205 while (destroys.length) {
206 destroys.shift()(error);
207 }
208 outerSignal === null || outerSignal === undefined
209 ? undefined
210 : outerSignal.removeEventListener('abort', abort);
211 ac.abort();
212 if (final) {
213 if (!error) {
214 lastStreamCleanup.forEach((fn) => fn());
215 }
216 nextTick(callback, error, value);
217 }
218 }
219 let ret;
220 for (let i = 0; i < streams.length; i++) {
221 const stream = streams[i];
222 const reading = i < streams.length - 1;
223 const writing = i > 0;
224 const end =
225 reading ||
226 (opts === null || opts === undefined ? undefined : opts.end) !== false;
227 const isLastStream = i === streams.length - 1;
228 if (isNodeStream(stream)) {
229 if (end) {
230 const { destroy, cleanup } = destroyer(stream, reading, writing);
231 destroys.push(destroy);
232 if (isReadable(stream) && isLastStream) {
233 lastStreamCleanup.push(cleanup);
234 }
235 }
236 
237 // Catch stream errors that occur after pipe/pump has completed.
238 function onError(err) {
239 if (
240 err &&
241 err.name !== 'AbortError' &&
242 err.code !== 'ERR_STREAM_PREMATURE_CLOSE'
243 ) {
244 finish(err);
245 }
246 }
247 stream.on('error', onError);
248 if (isReadable(stream) && isLastStream) {
249 lastStreamCleanup.push(() => {
250 stream.removeListener('error', onError);
251 });
252 }
253 }
254 if (i === 0) {
255 if (typeof stream === 'function') {
256 ret = stream({
257 signal,
258 });
259 if (!isIterable(ret)) {
260 throw new ERR_INVALID_RETURN_VALUE(
261 'Iterable, AsyncIterable or Stream',
262 'source',
263 ret
264 );
265 }
266 } else if (isIterable(stream) || isReadableNodeStream(stream)) {
267 ret = stream;
268 } else {
269 ret = from(Duplex, stream);
270 }
271 } else if (typeof stream === 'function') {
272 ret = makeAsyncIterable(ret);
273 ret = stream(ret, {
274 signal,
275 });
276 if (reading) {
277 if (!isIterable(ret, true)) {
278 throw new ERR_INVALID_RETURN_VALUE(
279 'AsyncIterable',
280 `transform[${i - 1}]`,
281 ret
282 );
283 }
284 } else {
285 let _ret;
286 // If the last argument to pipeline is not a stream
287 // we must create a proxy stream so that pipeline(...)
288 // always returns a stream which can be further
289 // composed through `.pipe(stream)`.
290 
291 const pt = new PassThrough({
292 objectMode: true,
293 });
294 
295 // Handle Promises/A+ spec, `then` could be a getter that throws on
296 // second use.
297 const then =
298 (_ret = ret) === null || _ret === undefined ? undefined : _ret.then;
299 if (typeof then === 'function') {
300 finishCount++;
301 then.call(
302 ret,
303 (val) => {
304 value = val;
305 if (val != null) {
306 pt.write(val);
307 }
308 if (end) {
309 pt.end();
310 }
311 nextTick(finish);
312 },
313 (err) => {
314 pt.destroy(err);
315 nextTick(finish, err);
316 }
317 );
318 } else if (isIterable(ret, true)) {
319 finishCount++;
320 pump(ret, pt, finish, {
321 end,
322 });
323 } else {
324 throw new ERR_INVALID_RETURN_VALUE(
325 'AsyncIterable or Promise',
326 'destination',
327 ret
328 );
329 }
330 ret = pt;
331 const { destroy, cleanup } = destroyer(ret, false, true);
332 destroys.push(destroy);
333 if (isLastStream) {
334 lastStreamCleanup.push(cleanup);
335 }
336 }
337 } else if (isNodeStream(stream)) {
338 if (isReadableNodeStream(ret)) {
339 finishCount += 2;
340 const cleanup = pipe(ret, stream, finish, {
341 end,
342 });
343 if (isReadable(stream) && isLastStream) {
344 lastStreamCleanup.push(cleanup);
345 }
346 } else if (isIterable(ret)) {
347 finishCount++;
348 pump(ret, stream, finish, {
349 end,
350 });
351 } else {
352 throw new ERR_INVALID_ARG_TYPE(
353 'val',
354 ['Readable', 'Iterable', 'AsyncIterable'],
355 ret
356 );
357 }
358 ret = stream;
359 } else {
360 ret = from(Duplex, stream);
361 }
362 }
363 if (
364 (signal !== null && signal !== undefined && signal.aborted) ||
365 (outerSignal !== null && outerSignal !== undefined && outerSignal.aborted)
366 ) {
367 nextTick(abort);
368 }
369 return ret;
370}
371 
372export function pipe(src, dst, finish, { end }) {
373 let ended = false;
374 dst.on('close', () => {
375 if (!ended) {
376 // Finish if the destination closes before the source has completed.
377 finish(new ERR_STREAM_PREMATURE_CLOSE());
378 }
379 });
380 src.pipe(dst, {
381 end,
382 });
383 if (end) {
384 // Compat. Before node v10.12.0 stdio used to throw an error so
385 // pipe() did/does not end() stdio destinations.
386 // Now they allow it but "secretly" don't close the underlying fd.
387 src.once('end', () => {
388 ended = true;
389 dst.end();
390 });
391 } else {
392 finish();
393 }
394 eos(
395 src,
396 {
397 readable: true,
398 writable: false,
399 },
400 (err) => {
401 const rState = src._readableState;
402 if (
403 err &&
404 err.code === 'ERR_STREAM_PREMATURE_CLOSE' &&
405 rState &&
406 rState.ended &&
407 !rState.errored &&
408 !rState.errorEmitted
409 ) {
410 // Some readable streams will emit 'close' before 'end'. However, since
411 // this is on the readable side 'end' should still be emitted if the
412 // stream has been ended and no error emitted. This should be allowed in
413 // favor of backwards compatibility. Since the stream is piped to a
414 // destination this should not result in any observable difference.
415 // We don't need to check if this is a writable premature close since
416 // eos will only fail with premature close on the reading side for
417 // duplex streams.
418 src.once('end', finish).once('error', finish);
419 } else {
420 finish(err);
421 }
422 }
423 );
424 return eos(
425 dst,
426 {
427 readable: false,
428 writable: true,
429 },
430 finish
431 );
432}