Skip to content

pipeline deadlock with readable tail? #40685

Description

@ronag

#40653 (comment)

@lpinca feel free to edit issue and fill out with your concern. Otherwise I will dig into this at some later point.


The pipeline() callback might never be called if the destination is a Transform stream that is not read. Here is an example.

const stream = require('stream');

const chunk = Buffer.alloc(1024);
let bytesRead = 0;

const readable = new stream.Readable({
  read() {
    if (bytesRead === 100 * 1024) {
      readable.push(null);
    } else {
      bytesRead += chunk.length;
      readable.push(chunk);
    }
  }
});

const passThrough = new stream.PassThrough();

stream.pipeline(readable, passThrough, function (err) {
  if (err) {
    throw err;
  }

  console.log('done');
});

In the callback version this is not a big issue because pipeline() returns the last stream, so the user could resume the returned stream.

pipeline(readable, passThrough, fn).resume();

In the promisified variant, pipeline() is just a wrapper that returns a promise that is fulfilled when the callback of the callback variant is called.

However the user could not be able to resume the returned stream because they might not have direct control on it. For example if the last entry is a generator.

const stream = require('stream');
const streamPromises = require('stream/promises');

const chunk = Buffer.alloc(1024);
let bytesRead = 0;

const readable = new stream.Readable({
  read() {
    if (bytesRead === 100 * 1024) {
      readable.push(null);
    } else {
      bytesRead += chunk.length;
      readable.push(chunk);
    }
  }
});

async function* passThrough(source) {
  for await (const chunk of source) {
    yield chunk;
  }
}

async function run() {
  await streamPromises.pipeline(readable, passThrough);
}

run()
  .then(function () {
    console.log('done');
  })
  .catch(console.error);

In this case stream.pipeline() returns an internal PassThrough proxy that the user cannot resume so the promise returned by streamPromises.pipeline() might never be fulfilled.

A possible workaround is to resume the returned stream in the promisified variant.

diff --git a/lib/stream/promises.js b/lib/stream/promises.js
index 0db01a8b20..5bbadd43c5 100644
--- a/lib/stream/promises.js
+++ b/lib/stream/promises.js
@@ -23,13 +23,17 @@ function pipeline(...streams) {
       signal = options.signal;
     }
 
-    pl(streams, (err, value) => {
+    const stream = pl(streams, (err, value) => {
       if (err) {
         reject(err);
       } else {
         resolve(value);
       }
     }, { signal });
+
+    if (stream.readable) {
+      stream.resume();
+    }
   });
 }

Activity

  1. added
    streamIssues and PRs related to Node.js streams.
    on Nov 1, 2021
  2. romayalon commented on Jan 16, 2022

    @romayalon

    I ran into the same issue...
    any news regarding a fix?

  3. trivikr commented on May 20, 2026

    @trivikr
    Member

    Minimal repro in node 26.1.0

    import { Readable } from 'node:stream';
    import { pipeline } from 'node:stream/promises';
    
    const chunk = Buffer.alloc(1024);
    let bytesRead = 0;
    let pushedEOF = false;
    
    const readable = new Readable({
      read() {
        while (bytesRead < 100 * 1024) {
          bytesRead += chunk.length;
          readable.push(chunk);
        }
    
        if (!pushedEOF) {
          pushedEOF = true;
          readable.push(null);
        }
      },
    });
    
    async function* passThrough(source) {
      for await (const chunk of source) {
        yield chunk;
      }
    }
    
    const result = await Promise.race([
      pipeline(readable, passThrough).then(() => 'resolved'),
      new Promise((resolve) => setTimeout(() => {
        resolve(`timed out; pushed EOF: ${pushedEOF}`);
      }, 1000)),
    ]);
    
    console.log(result);

    Observed

    timed out; pushed EOF: true

    Expected

    resolved

    Since the source has pushed all chunks and EOF

  4. github-actions commented on Aug 19, 2026

    @github-actions
    Contributor

    This issue has been marked as stale due to 90 days of inactivity.
    It will be automatically closed in 30 days if no further activity occurs. If this is still relevant, please leave a comment or update it to keep it open.

  5. added
    staleIssues and PRs marked stale due to inactivity and scheduled for automatic closure.
    on Aug 19, 2026
  6. github-actions commented on Sep 19, 2026

    @github-actions
    Contributor

    This issue has been automatically closed after 30 days of inactivity following its stale status (no activity for a total of 120 days).
    If this is still relevant, feel free to reopen it or leave a comment with additional details so we can continue the discussion.

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    staleIssues and PRs marked stale due to inactivity and scheduled for automatic closure.streamIssues and PRs related to Node.js streams.

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions