Hacktoberfest 2026: as issues que os mantenedores marcaram para outubro, abertas e boas para iniciantes. Ver issues do Hacktoberfest

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

Aberta
#65,063 0 comentários 0 reações 0 responsáveis Ver no GitHub

Ninguém assumiu esta issue ainda.

Avaliação

Dificuldade
3/5
Tempo estimado
1-2 dias
Facilidade para iniciantes
68/100
Tipo de issue
Bug
Clareza
Claramente especificada
Status de atividade
Pouca atividade
Stack de tecnologia
javascript
Domínio
backend

Direção de pesquisa

Comece em lib/internal/streams/pipeline.js, em pipelineImpl() e finishImpl(), e depois revise as garantias de stream.pipeline() em doc/api/stream.md. Reproduza os casos síncronos ERR_STREAM_UNABLE_TO_PIPE e ERR_INVALID_RETURN_VALUE da issue e compare-os com o caso de controle assíncrono ENOENT. Considera-se concluído quando os streams adotados forem destruídos, os descritores de arquivo forem liberados e o listener de AbortSignal do chamador for removido.

Escrita pelo modelo de indexação a partir do texto da issue.

Descrição

Version

v24.11.1, and main (735a09f999a)

Platform

Reproduced on Windows 11 x64; the code path is platform independent.

Subsystem

stream

What steps will reproduce the bug?

stream.pipeline() wires the streams together in a loop, and that loop can throw synchronously. The most common case is ERR_STREAM_UNABLE_TO_PIPE, raised when the destination is already closed or destroyed. When it throws, every stream pipeline() had already taken ownership of is left undestroyed, so its resources leak. For a fs.ReadStream source that means a leaked file descriptor.

The trigger is an ordinary production condition: piping to a destination that has already gone away, e.g. pipeline(fs.createReadStream(file), res) after the HTTP client disconnected.

import { pipeline, PassThrough, Writable } from 'node:stream';
import { pipeline as pipelinePromise } from 'node:stream/promises';
import { once, getEventListeners } from 'node:events';
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';

const tmp = fs.mkdtempSync(path.join(os.tmpdir(), 'pipeline-leak-'));
const file = path.join(tmp, 'data.bin');
fs.writeFileSync(file, Buffer.alloc(4096, 'x'));

const openStream = async () => {
  const rs = fs.createReadStream(file);
  await once(rs, 'open');          // ensure the fd is really allocated
  return rs;
};
const deadWritable = () => {
  const w = new Writable({ write(c, e, cb) { cb(); } });
  w.destroy();                     // destination already gone
  return w;
};

// 1. callback form
{
  const sources = [];
  for (let i = 0; i < 50; i++) {
    const rs = await openStream();
    sources.push(rs);
    try {
      pipeline(rs, new PassThrough(), deadWritable(), () => {});
    } catch (err) {
      if (i === 0) console.log('callback form throws :', err.code);
    }
  }
  await new Promise((r) => setTimeout(r, 150));
  console.log('  sources undestroyed :', sources.filter((s) => !s.destroyed).length, '/ 50  (expected 0)');
  console.log('  fds still open      :', sources.filter((s) => s.fd != null).length, '/ 50  (expected 0)');
  for (const s of sources) s.destroy();
}

// 2. promise form, plus the caller's AbortSignal
{
  const ac = new AbortController();  // long-lived, e.g. a server shutdown signal
  const sources = [];
  for (let i = 0; i < 50; i++) {
    const rs = await openStream();
    sources.push(rs);
    try {
      await pipelinePromise(rs, new PassThrough(), deadWritable(), { signal: ac.signal });
    } catch (err) {
      if (i === 0) console.log('promise form rejects :', err.code);
    }
  }
  await new Promise((r) => setTimeout(r, 150));
  console.log('  sources undestroyed :', sources.filter((s) => !s.destroyed).length, '/ 50  (expected 0)');
  console.log('  fds still open      :', sources.filter((s) => s.fd != null).length, '/ 50  (expected 0)');
  console.log('  abort listeners     :', getEventListeners(ac.signal, 'abort').length, '/ 50  (expected 0)');
  for (const s of sources) s.destroy();
}

// 3. control: an ordinary asynchronous failure does clean up correctly
{
  const rs = fs.createReadStream(path.join(tmp, 'nope'));
  const mid = new PassThrough();
  await new Promise((resolve) => pipeline(rs, mid, new PassThrough(), () => resolve()));
  console.log('control (ENOENT)     : rs.destroyed =', rs.destroyed, '| mid.destroyed =', mid.destroyed, ' (both expected true)');
}

fs.rmSync(tmp, { recursive: true, force: true });
How often does it reproduce? Is there a required condition?

Every time. The only condition is that pipeline() throws while wiring the streams up, after at least one stream has already been wired.

What is the expected behavior? Why is that the expected behavior?

All the streams should be destroyed and the file descriptors released, and the listener added to the caller's AbortSignal should be removed.

doc/api/stream.md states:

stream.pipeline() closes all the streams when an error is raised.

and

stream.pipeline() will call stream.destroy(err) on all streams except: Readable streams which have emitted 'end' or 'close'; Writable streams which have emitted 'finish' or 'close'.

The sources here emitted none of those events. The control case in the reproduction shows that an asynchronous pipeline error does destroy everything, so the two paths disagree.

What do you see instead?
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)
Additional information

The wiring loop in pipelineImpl() (lib/internal/streams/pipeline.js) pushes a destroy function into destroys for each stream it adopts, and finishImpl() is the only code that drains destroys, disposes the AbortSignal listener and calls ac.abort().

The loop is not wrapped in try/finally, and it can throw at six places:

  • ERR_STREAM_UNABLE_TO_PIPE when the next stream is already closed or destroyed
  • ERR_INVALID_RETURN_VALUE (x3) when a transform function returns something that is not iterable
  • ERR_INVALID_ARG_TYPE (x2) when a value cannot be piped into the next stream

When any of these fire, finishImpl() never runs, so nothing in destroys is ever called.

The ERR_INVALID_RETURN_VALUE path leaks in exactly the same way:

const source = fs.createReadStream(file);
pipeline(source, () => 42, new PassThrough(), () => {});  // throws
// source is left open

I have a fix and will open a PR shortly.

Linguagem predominante
JavaScript
Estrelas
122k
Forks
37.4k
Merge médio
4d 3h
PRs com merge (30d)
279

Guia de contribuição

Abrir o guia de contribuição

Primeiros passos

  1. Leia a issue inteira e depois o guia de contribuição do projeto.
  2. Comente na issue dizendo que vai assumir — evita que duas pessoas façam o mesmo trabalho.
  3. Faça um fork do repositório e trabalhe em uma branch.
  4. Abra um pull request que referencie o número da issue.

Mais de nodejs/node

Todas as issues de nodejs/node

Issues semelhantes

Mais issues de JavaScript

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.