diff --git a/lib/internal/streams/pipeline.js b/lib/internal/streams/pipeline.js index 546ff579d8eb..6d2434e9d260 100644 --- a/lib/internal/streams/pipeline.js +++ b/lib/internal/streams/pipeline.js @@ -251,160 +251,172 @@ function pipelineImpl(streams, callback, opts) { } let ret; - for (let i = 0; i < streams.length; i++) { - const stream = streams[i]; - const reading = i < streams.length - 1; - const writing = i > 0; - const next = i + 1 < streams.length ? streams[i + 1] : null; - const end = reading || opts?.end !== false; - const isLastStream = i === streams.length - 1; - - if (isNodeStream(stream)) { - if (next !== null && (next?.closed || next?.destroyed)) { - throw new ERR_STREAM_UNABLE_TO_PIPE(); - } + try { + for (let i = 0; i < streams.length; i++) { + const stream = streams[i]; + const reading = i < streams.length - 1; + const writing = i > 0; + const next = i + 1 < streams.length ? streams[i + 1] : null; + const end = reading || opts?.end !== false; + const isLastStream = i === streams.length - 1; + + if (isNodeStream(stream)) { + if (next !== null && (next?.closed || next?.destroyed)) { + throw new ERR_STREAM_UNABLE_TO_PIPE(); + } - if (end) { - const { destroy, cleanup } = destroyer(stream, reading, writing); - destroys.push(destroy); + if (end) { + const { destroy, cleanup } = destroyer(stream, reading, writing); + destroys.push(destroy); - if (isReadable(stream) && isLastStream) { - lastStreamCleanup.push(cleanup); + if (isReadable(stream) && isLastStream) { + lastStreamCleanup.push(cleanup); + } } - } - // Catch stream errors that occur after pipe/pump has completed. - function onError(err) { - if ( - err && - err.name !== 'AbortError' && - err.code !== 'ERR_STREAM_PREMATURE_CLOSE' - ) { - finishOnlyHandleError(err); + // Catch stream errors that occur after pipe/pump has completed. + function onError(err) { + if ( + err && + err.name !== 'AbortError' && + err.code !== 'ERR_STREAM_PREMATURE_CLOSE' + ) { + finishOnlyHandleError(err); + } } - } - stream.on('error', onError); - if (isReadable(stream) && isLastStream) { - lastStreamCleanup.push(() => { - stream.removeListener('error', onError); - }); - } - } - - if (i === 0) { - if (typeof stream === 'function') { - ret = stream({ signal }); - if (!isIterable(ret)) { - throw new ERR_INVALID_RETURN_VALUE( - 'Iterable, AsyncIterable or Stream', 'source', ret); + stream.on('error', onError); + if (isReadable(stream) && isLastStream) { + lastStreamCleanup.push(() => { + stream.removeListener('error', onError); + }); } - } else if (isIterable(stream) || isReadableNodeStream(stream) || isTransformStream(stream)) { - ret = stream; - } else { - ret = Duplex.from(stream); - } - } else if (typeof stream === 'function') { - if (isTransformStream(ret)) { - ret = makeAsyncIterable(ret?.readable); - } else { - ret = makeAsyncIterable(ret); } - ret = stream(ret, { signal }); - if (reading) { - if (!isIterable(ret, true)) { - throw new ERR_INVALID_RETURN_VALUE( - 'AsyncIterable', `transform[${i - 1}]`, ret); + if (i === 0) { + if (typeof stream === 'function') { + ret = stream({ signal }); + if (!isIterable(ret)) { + throw new ERR_INVALID_RETURN_VALUE( + 'Iterable, AsyncIterable or Stream', 'source', ret); + } + } else if (isIterable(stream) || isReadableNodeStream(stream) || isTransformStream(stream)) { + ret = stream; + } else { + ret = Duplex.from(stream); } - } else { - PassThrough ??= require('internal/streams/passthrough'); - - // If the last argument to pipeline is not a stream - // we must create a proxy stream so that pipeline(...) - // always returns a stream which can be further - // composed through `.pipe(stream)`. - - const pt = new PassThrough({ - objectMode: true, - }); + } else if (typeof stream === 'function') { + if (isTransformStream(ret)) { + ret = makeAsyncIterable(ret?.readable); + } else { + ret = makeAsyncIterable(ret); + } + ret = stream(ret, { signal }); - // Handle Promises/A+ spec, `then` could be a getter that throws on - // second use. - const then = ret?.then; - if (typeof then === 'function') { - finishCount++; - then.call(ret, - (val) => { - value = val; - if (val != null) { - pt.write(val); - } - if (end) { - pt.end(); - } - process.nextTick(finish); - }, (err) => { - pt.destroy(err); - process.nextTick(finish, err); - }, - ); - } else if (isIterable(ret, true)) { - finishCount++; - pumpToNode(ret, pt, finish, { end }); - } else if (isReadableStream(ret) || isTransformStream(ret)) { + if (reading) { + if (!isIterable(ret, true)) { + throw new ERR_INVALID_RETURN_VALUE( + 'AsyncIterable', `transform[${i - 1}]`, ret); + } + } else { + PassThrough ??= require('internal/streams/passthrough'); + + // If the last argument to pipeline is not a stream + // we must create a proxy stream so that pipeline(...) + // always returns a stream which can be further + // composed through `.pipe(stream)`. + + const pt = new PassThrough({ + objectMode: true, + }); + + // Handle Promises/A+ spec, `then` could be a getter that throws on + // second use. + const then = ret?.then; + if (typeof then === 'function') { + finishCount++; + then.call(ret, + (val) => { + value = val; + if (val != null) { + pt.write(val); + } + if (end) { + pt.end(); + } + process.nextTick(finish); + }, (err) => { + pt.destroy(err); + process.nextTick(finish, err); + }, + ); + } else if (isIterable(ret, true)) { + finishCount++; + pumpToNode(ret, pt, finish, { end }); + } else if (isReadableStream(ret) || isTransformStream(ret)) { + const toRead = ret.readable || ret; + finishCount++; + pumpToNode(toRead, pt, finish, { end }); + } else { + throw new ERR_INVALID_RETURN_VALUE( + 'AsyncIterable or Promise', 'destination', ret); + } + + ret = pt; + + const { destroy, cleanup } = destroyer(ret, false, true); + destroys.push(destroy); + if (isLastStream) { + lastStreamCleanup.push(cleanup); + } + } + } else if (isNodeStream(stream)) { + if (isReadableNodeStream(ret)) { + finishCount += 2; + const cleanup = pipe(ret, stream, finish, finishOnlyHandleError, { end }); + if (isReadable(stream) && isLastStream) { + lastStreamCleanup.push(cleanup); + } + } else if (isTransformStream(ret) || isReadableStream(ret)) { const toRead = ret.readable || ret; finishCount++; - pumpToNode(toRead, pt, finish, { end }); + pumpToNode(toRead, stream, finish, { end }); + } else if (isIterable(ret)) { + finishCount++; + pumpToNode(ret, stream, finish, { end }); } else { - throw new ERR_INVALID_RETURN_VALUE( - 'AsyncIterable or Promise', 'destination', ret); + throw new ERR_INVALID_ARG_TYPE( + 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); } - - ret = pt; - - const { destroy, cleanup } = destroyer(ret, false, true); - destroys.push(destroy); - if (isLastStream) { - lastStreamCleanup.push(cleanup); - } - } - } else if (isNodeStream(stream)) { - if (isReadableNodeStream(ret)) { - finishCount += 2; - const cleanup = pipe(ret, stream, finish, finishOnlyHandleError, { end }); - if (isReadable(stream) && isLastStream) { - lastStreamCleanup.push(cleanup); + ret = stream; + } else if (isWebStream(stream)) { + if (isReadableNodeStream(ret)) { + finishCount++; + pumpToWeb(makeAsyncIterable(ret), stream, finish, { end }); + } else if (isReadableStream(ret) || isIterable(ret)) { + finishCount++; + pumpToWeb(ret, stream, finish, { end }); + } else if (isTransformStream(ret)) { + finishCount++; + pumpToWeb(ret.readable, stream, finish, { end }); + } else { + throw new ERR_INVALID_ARG_TYPE( + 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); } - } else if (isTransformStream(ret) || isReadableStream(ret)) { - const toRead = ret.readable || ret; - finishCount++; - pumpToNode(toRead, stream, finish, { end }); - } else if (isIterable(ret)) { - finishCount++; - pumpToNode(ret, stream, finish, { end }); - } else { - throw new ERR_INVALID_ARG_TYPE( - 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); - } - ret = stream; - } else if (isWebStream(stream)) { - if (isReadableNodeStream(ret)) { - finishCount++; - pumpToWeb(makeAsyncIterable(ret), stream, finish, { end }); - } else if (isReadableStream(ret) || isIterable(ret)) { - finishCount++; - pumpToWeb(ret, stream, finish, { end }); - } else if (isTransformStream(ret)) { - finishCount++; - pumpToWeb(ret.readable, stream, finish, { end }); + ret = stream; } else { - throw new ERR_INVALID_ARG_TYPE( - 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); + ret = Duplex.from(stream); } - ret = stream; - } else { - ret = Duplex.from(stream); } + } catch (err) { + // The loop above can throw synchronously, e.g. ERR_STREAM_UNABLE_TO_PIPE. + // `finishImpl()` never runs on this path, so the listener that was added to + // the caller's `AbortSignal` would stay attached for the lifetime of that + // signal. Dispose of it before propagating. + // The streams themselves are deliberately left alone: ownership is not + // taken until the pipeline has been established, so destroying them + // remains the caller's responsibility. + disposable?.[SymbolDispose](); + throw err; } if (signal?.aborted || outerSignal?.aborted) { diff --git a/test/parallel/test-stream-pipeline-sync-throw-cleanup.js b/test/parallel/test-stream-pipeline-sync-throw-cleanup.js new file mode 100644 index 000000000000..9c549556e4a5 --- /dev/null +++ b/test/parallel/test-stream-pipeline-sync-throw-cleanup.js @@ -0,0 +1,67 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { pipeline, PassThrough, Readable } = require('stream'); +const { pipeline: pipelinePromise } = require('stream/promises'); +const { getEventListeners } = require('events'); + +// pipeline() can throw synchronously while it is still wiring the streams +// together, e.g. ERR_STREAM_UNABLE_TO_PIPE when the destination has already +// been destroyed. On that path finishImpl() never runs, so the listener added +// to the caller's AbortSignal has to be disposed of here instead, otherwise it +// stays attached for the lifetime of the signal. +// +// The streams themselves are intentionally not destroyed, see the +// ERR_INVALID_RETURN_VALUE cases in test-stream-pipeline.js. + +// Callback form. `pipeline()` does not accept options, so this only checks +// that the throw still propagates and no listener is left behind on the +// streams' behalf. +{ + const source = new Readable({ read() {} }); + const dest = new PassThrough(); + dest.destroy(); + + assert.throws(() => { + pipeline(source, new PassThrough(), dest, common.mustNotCall()); + }, { code: 'ERR_STREAM_UNABLE_TO_PIPE' }); +} + +// Promise form with a caller-owned signal: the listener must be removed. +{ + const ac = new AbortController(); + assert.strictEqual(getEventListeners(ac.signal, 'abort').length, 0); + + const source = new Readable({ read() {} }); + const dest = new PassThrough(); + dest.destroy(); + + assert.rejects( + pipelinePromise(source, new PassThrough(), dest, { signal: ac.signal }), + { code: 'ERR_STREAM_UNABLE_TO_PIPE' }, + ).then(common.mustCall(() => { + assert.strictEqual(getEventListeners(ac.signal, 'abort').length, 0); + })); +} + +// The same signal reused across several failed calls must not accumulate +// listeners. +{ + const ac = new AbortController(); + const pending = []; + + for (let i = 0; i < 10; i++) { + const dest = new PassThrough(); + dest.destroy(); + pending.push(assert.rejects( + pipelinePromise(new Readable({ read() {} }), new PassThrough(), dest, + { signal: ac.signal }), + { code: 'ERR_STREAM_UNABLE_TO_PIPE' }, + )); + } + + Promise.all(pending).then(common.mustCall(() => { + assert.strictEqual(getEventListeners(ac.signal, 'abort').length, 0); + })); +}