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{pipelineaspipelinePromise}from'node:stream/promises';import{once,getEventListeners}from'node:events';importfsfrom'node:fs';importosfrom'node:os';importpathfrom'node:path';consttmp=fs.mkdtempSync(path.join(os.tmpdir(),'pipeline-leak-'));constfile=path.join(tmp,'data.bin');fs.writeFileSync(file,Buffer.alloc(4096,'x'));constopenStream=async()=>{constrs=fs.createReadStream(file);awaitonce(rs,'open');// ensure the fd is really allocatedreturnrs;};constdeadWritable=()=>{constw=newWritable({write(c,e,cb){cb();}});w.destroy();// destination already gonereturnw;};// 1. callback form{constsources=[];for(leti=0;i<50;i++){constrs=awaitopenStream();sources.push(rs);try{pipeline(rs,newPassThrough(),deadWritable(),()=>{});}catch(err){if(i===0)console.log('callback form throws :',err.code);}}awaitnewPromise((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(constsofsources)s.destroy();}// 2. promise form, plus the caller's AbortSignal{constac=newAbortController();// long-lived, e.g. a server shutdown signalconstsources=[];for(leti=0;i<50;i++){constrs=awaitopenStream();sources.push(rs);try{awaitpipelinePromise(rs,newPassThrough(),deadWritable(),{signal: ac.signal});}catch(err){if(i===0)console.log('promise form rejects :',err.code);}}awaitnewPromise((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(constsofsources)s.destroy();}// 3. control: an ordinary asynchronous failure does clean up correctly{constrs=fs.createReadStream(path.join(tmp,'nope'));constmid=newPassThrough();awaitnewPromise((resolve)=>pipeline(rs,mid,newPassThrough(),()=>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 destroyedERR_INVALID_RETURN_VALUE (x3) when a transform function returns something that is not iterableERR_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:
constsource=fs.createReadStream(file);pipeline(source,()=>42,newPassThrough(),()=>{});// throws// source is left openI have a fix and will open a PR shortly.
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 isERR_STREAM_UNABLE_TO_PIPE, raised when the destination is already closed or destroyed. When it throws, every streampipeline()had already taken ownership of is left undestroyed, so its resources leak. For afs.ReadStreamsource 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.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
AbortSignalshould be removed.doc/api/stream.mdstates:and
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?
Additional information
The wiring loop in
pipelineImpl()(lib/internal/streams/pipeline.js) pushes a destroy function intodestroysfor each stream it adopts, andfinishImpl()is the only code that drainsdestroys, disposes theAbortSignallistener and callsac.abort().The loop is not wrapped in
try/finally, and it can throw at six places:ERR_STREAM_UNABLE_TO_PIPEwhen the next stream is already closed or destroyedERR_INVALID_RETURN_VALUE(x3) when a transform function returns something that is not iterableERR_INVALID_ARG_TYPE(x2) when a value cannot be piped into the next streamWhen any of these fire,
finishImpl()never runs, so nothing indestroysis ever called.The
ERR_INVALID_RETURN_VALUEpath leaks in exactly the same way:I have a fix and will open a PR shortly.