diff --git a/lib/internal/streams/pipeline.js b/lib/internal/streams/pipeline.js index 546ff579d8eb..afaf2b60dafd 100644 --- a/lib/internal/streams/pipeline.js +++ b/lib/internal/streams/pipeline.js @@ -217,6 +217,7 @@ function pipelineImpl(streams, callback, opts) { const destroys = []; let finishCount = 0; + let syncThrow = false; function finish(err) { finishImpl(err, --finishCount === 0); @@ -227,6 +228,14 @@ function pipelineImpl(streams, callback, opts) { } function finishImpl(err, final) { + // If pipeline() threw synchronously while wiring the streams together, + // the caller already received the error. Any finish callbacks that fire + // afterwards (from stages already wired, which incremented finishCount) + // must not invoke the completion callback — doing so would report the + // same failure twice, once as an exception and once as a success. + if (syncThrow) { + return; + } if (err && (!error || error.code === 'ERR_STREAM_PREMATURE_CLOSE' || error.name === 'AbortError')) { error = err; } @@ -251,160 +260,175 @@ 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(); - } - - if (end) { - const { destroy, cleanup } = destroyer(stream, reading, writing); - destroys.push(destroy); - - 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); - } - } - 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); + 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 (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); + } + } + 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); - } - } 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 }); + 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 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); + } + } 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, 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 }); + } else { + throw new ERR_INVALID_ARG_TYPE( + 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); + } + ret = stream; } 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); + ret = Duplex.from(stream); } - } 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 }); - } else { - throw new ERR_INVALID_ARG_TYPE( - 'val', ['Readable', 'Iterable', 'AsyncIterable', 'ReadableStream', 'TransformStream'], ret); } - ret = stream; - } else { - ret = Duplex.from(stream); + } catch (err) { + // A synchronous throw while wiring the streams (e.g. an invalid + // stream/function in the middle) leaves stages already wired alive: + // their finish callbacks would otherwise report success to the caller, + // who already received the exception. Destroy the wired stages, mark + // the sync throw so no finish callback can fire, and rethrow. + syncThrow = true; + while (destroys.length) { + destroys.shift()(err); } + disposable?.[SymbolDispose](); + ac.abort(); + throw err; } if (signal?.aborted || outerSignal?.aborted) { diff --git a/test/parallel/test-stream-pipeline.js b/test/parallel/test-stream-pipeline.js index 8ee197b44e0b..44afe0a7038d 100644 --- a/test/parallel/test-stream-pipeline.js +++ b/test/parallel/test-stream-pipeline.js @@ -1774,3 +1774,28 @@ tmpdir.refresh(); }), ); } +{ + // https://github.com/nodejs/node/issues/65127 + // When pipeline() throws synchronously while wiring streams together, + // it must NOT also invoke the callback — and never with success. + // A node stream wired before the throwing stage increments finishCount, + // so a naive teardown would fire the callback with no error (double + // report: the caller already received the exception). + const r = Readable.from(['a']); + const t = new Transform({ + transform(chunk, encoding, callback) { + callback(null, chunk); + }, + }); + let threw = false; + try { + pipeline(r, t, () => 42, common.mustNotCall()); + } catch (err) { + threw = true; + assert.strictEqual(err.code, 'ERR_INVALID_RETURN_VALUE'); + } + assert.strictEqual(threw, true); + // Give any stray finish callbacks a chance to fire; if the bug + // regresses, common.mustNotCall() fails the test. + setTimeout(() => {}, 100); +}