File
Blob: src/node/internal/streams_pipeline.js
| 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 | |
| 29 | import { |
| 30 | isIterable, |
| 31 | isReadable, |
| 32 | isReadableNodeStream, |
| 33 | isNodeStream, |
| 34 | } from 'node-internal:streams_util'; |
| 35 | import { eos } from 'node-internal:streams_end_of_stream'; |
| 36 | import { destroyer as destroyerImpl } from 'node-internal:streams_destroy'; |
| 37 | import { once } from 'node-internal:internal_http_util'; |
| 38 | |
| 39 | import { nextTick } from 'node-internal:internal_process'; |
| 40 | import { PassThrough } from 'node-internal:streams_transform'; |
| 41 | import { Duplex } from 'node-internal:streams_duplex'; |
| 42 | import { Readable, from } from 'node-internal:streams_readable'; |
| 43 | import { |
| 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'; |
| 52 | import { |
| 53 | validateFunction, |
| 54 | validateAbortSignal, |
| 55 | } from 'node-internal:validators'; |
| 56 | |
| 57 | function 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 | |
| 82 | function 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 | |
| 90 | function 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 | |
| 104 | async function* fromReadable(val) { |
| 105 | yield* Readable.prototype[Symbol.asyncIterator].call(val); |
| 106 | } |
| 107 | |
| 108 | async 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 | |
| 166 | export function pipeline(...streams) { |
| 167 | return pipelineImpl(streams, once(popCallback(streams))); |
| 168 | } |
| 169 | |
| 170 | export 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 | |
| 372 | export 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 | } |