stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

stream: clean up when pipeline throws synchronously - #65064

Open
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak
Open

stream: clean up when pipeline throws synchronously#65064
shani-singh1 wants to merge 1 commit into
nodejs:mainfrom
shani-singh1:stream-pipeline-sync-throw-leak

Conversation

@shani-singh1

Copy link
Copy Markdown
Contributor

pipelineImpl() wires the streams together in a loop that can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the next stream is already closed or destroyed, which happens routinely when a destination goes away first (for example pipeline(fs.createReadStream(file), res) after the HTTP client disconnected).

Each stream the loop adopts registers a destroy function in destroys. finishImpl() is the only code that drains destroys, disposes the listener added to the caller's AbortSignal and calls ac.abort(), and it never runs when the loop throws. Every stream already wired up is therefore left undestroyed and its resources leak. For an fs.ReadStream source that is a leaked file descriptor.

The loop has six synchronous throw sites: one ERR_STREAM_UNABLE_TO_PIPE, three ERR_INVALID_RETURN_VALUE and two ERR_INVALID_ARG_TYPE. This wraps the loop so the same teardown runs before the error propagates. The error is still thrown, so the observable failure mode is unchanged.

Most of the diff is the re-indentation of the existing loop. Reviewing with git diff -w shows the actual change is 13 lines.

Before

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 50 / 50 (expected 0)
fds still open : 50 / 50 (expected 0)
abort listeners : 50 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

After

callback form throws : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
promise form rejects : ERR_STREAM_UNABLE_TO_PIPE
sources undestroyed : 0 / 50 (expected 0)
fds still open : 0 / 50 (expected 0)
abort listeners : 0 / 50 (expected 0)
control (ENOENT) : rs.destroyed = true | mid.destroyed = true (both expected true)

The full reproduction is in the linked issue. The added test fails on main (5 failing assertions) and passes with this change.

I also checked the change against a set of ordinary pipeline() usages (happy path, fs read to writable, asynchronous mid-stream error, ENOENT source, async generator transform, abort via an outer signal) and the behaviour is identical before and after. The only observable differences are on the two synchronous-throw paths, where the streams are now destroyed.

Fixes: #65063

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

Review requested:

  • @nodejs/streams

@nodejs-github-botnodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to Node.js streams. labels Aug 6, 2026
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899CompareAugust 6, 2026 11:17
@ronag
ronag requested a lite review from CopilotAugust 6, 2026 12:28

@ronagronag left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would say that ownership is not taken until pipeline succeeds...

CopilotAI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

This PR fixes a resource-leak in stream.pipeline() when pipelineImpl() throws synchronously during its wiring loop (e.g. ERR_STREAM_UNABLE_TO_PIPE), ensuring already-adopted streams are torn down and abort listeners are removed before the error propagates.

Changes:

  • Wrap the stream-wiring loop in pipelineImpl() with synchronous-throw cleanup that drains destroys, disposes the outer AbortSignal listener, and aborts the internal controller.
  • Add a regression test covering synchronous-throw cleanup for ERR_STREAM_UNABLE_TO_PIPE and ERR_INVALID_RETURN_VALUE, including abort-listener disposal in the promise form.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

FileDescription
lib/internal/streams/pipeline.jsAdds teardown on synchronous throws during the wiring loop to prevent leaked streams/fds and leaked abort listeners.
test/parallel/test-stream-pipeline-sync-throw-cleanup.jsNew regression test asserting cleanup happens on synchronous throw paths (fd release, stream destroyed, abort listener removed).

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +410 to 421
} catch (err) {
// The loop above can throw synchronously (e.g. ERR_STREAM_UNABLE_TO_PIPE)
// after some streams have already been wired up. Those streams are
// registered in `destroys` but `finishImpl()` never runs, so tear them
// down here before propagating, otherwise their resources leak.
while (destroys.length) {
destroys.shift()(err);
}
disposable?.[SymbolDispose]();
ac.abort();
throw err;
}
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Thanks @ronag, that is a fair point and I think it splits the PR into two separate questions.

Destroying the streams. I accept the ownership argument. If pipeline() never got established then the caller still holds the streams and can clean them up, so destroying them is arguably not pipeline()'s call to make.

The AbortSignal listener, which I think is separate. When options.signal is passed, pipelineImpl() attaches a listener to it via addAbortListener() before the wiring loop and only disposes it in finishImpl(), which never runs on this path. The signal belongs to the caller, the listener is pipeline()'s internal abort closure, and the caller has no handle on it to remove. Measured over 20 failed calls sharing one long lived signal:

 main with the patch
listeners left on the caller's AbortSignal 20/20 0/20
'error' listeners pipeline left on sources 40 40
sources destroyed 0/20 20/20

So even setting the destroy question aside, pipeline() currently leaves two of its own artefacts behind on objects the caller owns: the abort listener, and the onError listeners it attached to the streams it had already reached.

Would you prefer I narrow this to only undoing pipeline()'s own side effects, i.e. dispose the abort listener and remove the listeners it added, and leave the streams untouched for the caller to destroy? That respects ownership staying with the caller and still stops the leak. Happy to redo it that way.

One doc question while we are here: stream.md currently says "stream.pipeline() closes all the streams when an error is raised". If the synchronous throw path is deliberately not covered by that, it may be worth spelling out that the caller is responsible for cleanup when pipeline() throws rather than calling back.

On the Copilot comment about finishCount: that behaviour is pre-existing rather than introduced here. pipeline(Readable.from(['a']), transform, () => 42, cb) throws ERR_INVALID_RETURN_VALUE and still invokes cb on main today, with err undefined. I get the same result with and without this patch, so it looks like a separate bug. Happy to open a separate issue for it.

@mcollinamcollina left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lgtm

@mcollinamcollina added the request-ci Add this label to start a Jenkins CI on a PR. label Aug 6, 2026
`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 <teamdeveloperworld@gmail.com>
@shani-singh1
shani-singh1force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9fCompareAugust 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

You were right, and the test suite says so explicitly. CI caught it:

=== release test-stream-pipeline ===
AssertionError [ERR_ASSERTION]: Expected values to be strictly equal:
at test/parallel/test-stream-pipeline.js:890:12

That line is assert.strictEqual(s.destroyed, false) inside the ERR_INVALID_RETURN_VALUE block, so "ownership is not taken until pipeline succeeds" is already asserted behaviour and my patch was contradicting it. Sorry for the noise.

I have narrowed this to only the part that is not about ownership: disposing the listener that pipelineImpl() adds to the caller's AbortSignal. The streams are now left completely untouched.

 main narrowed
listeners left on the caller's AbortSignal 20/20 0/20
sources destroyed 0/20 0/20

I also re-ran the three ERR_INVALID_RETURN_VALUE blocks from test-stream-pipeline.js against the new patch and all three pass, including line 890 which was failing.

The test in this PR now only covers the signal listener, since the not-destroyed behaviour is already covered by the existing cases.

@shani-singh1

Copy link
Copy Markdown
ContributorAuthor

Filed the finishCount behaviour separately as #65127, since it is pre-existing and independent of this PR. It reproduces on v24.11.1 and on main at aed4eaf89dd, with a control showing the callback is correctly not invoked when nothing incremented finishCount before the throw.

This PR is unchanged and still only disposes the AbortSignal listener.

@nodejs-github-bot

Copy link
Copy Markdown
Collaborator

@panvapanva removed the request-ci Add this label to start a Jenkins CI on a PR. label Aug 28, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ciPRs that need a full CI run.streamIssues and PRs related to Node.js streams.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

stream.pipeline() leaks file descriptors when it throws synchronously (ERR_STREAM_UNABLE_TO_PIPE)

6 participants

@shani-singh1@nodejs-github-bot@mcollina@ronag@panva