diff --git a/lib/_stream_readable.js b/lib/_stream_readable.js index f8982b1f08..1f464ebfaf 100644 --- a/lib/_stream_readable.js +++ b/lib/_stream_readable.js @@ -305,9 +305,11 @@ Readable.prototype.pipe = function(dest, pipeOpts) { dest !== process.stdout && dest !== process.stderr) { src.once('end', onend); - dest.on('unpipe', function(readable) { - if (readable === src) + dest.on('unpipe', function unpipe(readable) { + if (readable === src) { src.removeListener('end', onend); + dest.removeListener('unpipe', unpipe); + } }); } @@ -324,9 +326,9 @@ Readable.prototype.pipe = function(dest, pipeOpts) { // too slow. var ondrain = pipeOnDrain(src); dest.on('drain', ondrain); - dest.on('unpipe', function(readable) { + dest.on('unpipe', function cleanup(readable) { if (readable === src) - dest.removeListener('drain', ondrain); + dest.removeListener('unpipe', cleanup); // if the reader is waiting for a drain event from this // specific writer, then it would cause it to never start @@ -339,23 +341,36 @@ Readable.prototype.pipe = function(dest, pipeOpts) { // if the dest has an error, then stop piping into it. // however, don't suppress the throwing behavior for this. - dest.once('error', function(er) { + function onerror(er) { unpipe(); if (dest.listeners('error').length === 0) dest.emit('error', er); - }); + } + dest.once('error', onerror); // if the dest emits close, then presumably there's no point writing // to it any more. - dest.on('close', unpipe); - dest.on('finish', function() { + dest.once('close', unpipe); + function onfinish() { dest.removeListener('close', unpipe); - }); + } + dest.once('finish', onfinish); function unpipe() { src.unpipe(dest); } + // cleanup event handlers once the pipe is broken + dest.once('unpipe', function cleanup(readable) { + if (readable !== src) return; + + dest.removeListener('close', unpipe); + dest.removeListener('finish', onfinish); + dest.removeListener('drain', ondrain); + dest.removeListener('error', onerror); + dest.removeListener('unpipe', cleanup); + }); + // tell the dest that it's being piped to dest.emit('pipe', src); diff --git a/test/simple/test-stream2-unpipe-leak.js b/test/simple/test-stream2-unpipe-leak.js new file mode 100644 index 0000000000..a560bfa0cb --- /dev/null +++ b/test/simple/test-stream2-unpipe-leak.js @@ -0,0 +1,63 @@ +// Copyright Joyent, Inc. and other Node contributors. +// +// Permission is hereby granted, free of charge, to any person obtaining a +// copy of this software and associated documentation files (the +// "Software"), to deal in the Software without restriction, including +// without limitation the rights to use, copy, modify, merge, publish, +// distribute, sublicense, and/or sell copies of the Software, and to permit +// persons to whom the Software is furnished to do so, subject to the +// following conditions: +// +// The above copyright notice and this permission notice shall be included +// in all copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS +// OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF +// MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN +// NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, +// DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR +// OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE +// USE OR OTHER DEALINGS IN THE SOFTWARE. + + +var common = require('../common.js'); +var assert = require('assert'); +var stream = require('stream'); + +var util = require('util'); + +function TestWriter() { + stream.Writable.call(this); +} +util.inherits(TestWriter, stream.Writable); + +TestWriter.prototype._write = function (buffer, callback) { + callback(null); +}; + +var dest = new TestWriter(); + +function TestReader() { + stream.Readable.call(this); +} +util.inherits(TestReader, stream.Readable); + +TestReader.prototype._read = function (size, callback) { + callback(new Buffer('hallo')); +}; + +var src = new TestReader(); + +for (var i = 0; i < 10; i++) { + src.pipe(dest); + src.unpipe(dest); +} + +assert.equal(src.listeners('end').length, 0); +assert.equal(src.listeners('readable').length, 0); + +assert.equal(dest.listeners('unpipe').length, 0); +assert.equal(dest.listeners('drain').length, 0); +assert.equal(dest.listeners('error').length, 0); +assert.equal(dest.listeners('close').length, 0); +assert.equal(dest.listeners('finish').length, 0);