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

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-bot nodejs-github-bot added needs-ci PRs that need a full CI run. stream Issues and PRs related to the stream subsystem. labels Aug 6, 2026
@shani-singh1
shani-singh1 force-pushed the stream-pipeline-sync-throw-leak branch from 9c3b6a8 to cb6d899 Compare August 6, 2026 11:17
@ronag
ronag requested a lite review from Copilot August 6, 2026 12:28

@ronag ronag 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...

Copilot AI 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.

File Description
lib/internal/streams/pipeline.js Adds 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.js New 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
Author

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.

@mcollina mcollina 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

@mcollina mcollina 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-singh1 force-pushed the stream-pipeline-sync-throw-leak branch from cb6d899 to c4f4c9f Compare August 6, 2026 21:41
@shani-singh1

Copy link
Copy Markdown
Author

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

needs-ci PRs that need a full CI run. request-ci Add this label to start a Jenkins CI on a PR. stream Issues and PRs related to the stream subsystem.

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)

5 participants