From c4f4c9f0c0cecba82cc5640c619a11437c6c01bd Mon Sep 17 00:00:00 2001 From: Shani Singh Date: Thu, 6 Aug 2026 07:18:52 +0530 Subject: [PATCH] stream: dispose abort listener when pipeline throws `pipelineImpl()` adds a listener to the caller's `AbortSignal` before it wires the streams together, and only disposes of it in `finishImpl()`. The wiring loop can throw synchronously, for example `ERR_STREAM_UNABLE_TO_PIPE` when the destination is already destroyed, and `finishImpl()` never runs on that path, so the listener stays attached for the lifetime of the signal. A long lived signal reused across many failed calls accumulates one listener per call. Dispose of it before propagating the error. The streams themselves are left untouched, since ownership is not taken until the pipeline has been established, which is the behaviour the existing `ERR_INVALID_RETURN_VALUE` cases in `test-stream-pipeline.js` assert. Signed-off-by: Shani Singh --- lib/internal/streams/pipeline.js | 284 +++++++++--------- ...test-stream-pipeline-sync-throw-cleanup.js | 67 +++++ 2 files changed, 215 insertions(+), 136 deletions(-) create mode 100644 test/parallel/test-stream-pipeline-sync-throw-cleanup.js 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); + })); +}