Skip to content

Commit 823ee2b

Browse files
ronagBethGriggs
authored andcommitted
stream: reset flowing state if no 'readable' or 'data' listeners
If we don't have any 'readable' or 'data' listeners and we are not about to resume. Then reset flowing state to initial null state. PR-URL: #31036Fixes: #24474 Reviewed-By: Luigi Pinca <luigipinca@gmail.com> Reviewed-By: Matteo Collina <matteo.collina@gmail.com> Reviewed-By: Rich Trott <rtrott@gmail.com>
1 parent b12b930 commit 823ee2b

3 files changed

Lines changed: 58 additions & 5 deletions

File tree

‎lib/_stream_readable.js‎

Lines changed: 21 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ const {
2828
ObjectDefineProperty,
2929
ObjectSetPrototypeOf,
3030
SymbolAsyncIterator,
31+
Symbol
3132
}=primordials;
3233

3334
module.exports=Readable;
@@ -51,6 +52,8 @@ const {
5152
ERR_STREAM_UNSHIFT_AFTER_END_EVENT
5253
}=require('internal/errors').codes;
5354

55+
constkPaused=Symbol('kPaused');
56+
5457
// Lazy loaded to improve the startup performance.
5558
letStringDecoder;
5659
letcreateReadableStreamAsyncIterator;
@@ -127,7 +130,7 @@ function ReadableState(options, stream, isDuplex) {
127130
this.emittedReadable=false;
128131
this.readableListening=false;
129132
this.resumeScheduled=false;
130-
this.paused=true;
133+
this[kPaused]=null;
131134

132135
// Should close be emitted on destroy. Defaults to true.
133136
this.emitClose=!options||options.emitClose!==false;
@@ -159,6 +162,16 @@ function ReadableState(options, stream, isDuplex) {
159162
}
160163
}
161164

165+
// Legacy property for `paused`
166+
ObjectDefineProperty(ReadableState.prototype,'paused',{
167+
get(){
168+
returnthis[kPaused]!==false;
169+
},
170+
set(value){
171+
this[kPaused]=!!value;
172+
}
173+
});
174+
162175
functionReadable(options){
163176
if(!(thisinstanceofReadable))
164177
returnnewReadable(options);
@@ -348,7 +361,8 @@ function chunkInvalid(state, chunk) {
348361

349362

350363
Readable.prototype.isPaused=function(){
351-
returnthis._readableState.flowing===false;
364+
conststate=this._readableState;
365+
returnstate[kPaused]===true||state.flowing===false;
352366
};
353367

354368
// Backwards compatibility.
@@ -947,14 +961,16 @@ function updateReadableListening(self) {
947961
conststate=self._readableState;
948962
state.readableListening=self.listenerCount('readable')>0;
949963

950-
if(state.resumeScheduled&&!state.paused){
964+
if(state.resumeScheduled&&state[kPaused]===false){
951965
// Flowing needs to be set to true now, otherwise
952966
// the upcoming resume will not flow.
953967
state.flowing=true;
954968

955969
// Crude way to check if we should resume
956970
}elseif(self.listenerCount('data')>0){
957971
self.resume();
972+
}elseif(!state.readableListening){
973+
state.flowing=null;
958974
}
959975
}
960976

@@ -975,7 +991,7 @@ Readable.prototype.resume = function() {
975991
state.flowing=!state.readableListening;
976992
resume(this,state);
977993
}
978-
state.paused=false;
994+
state[kPaused]=false;
979995
returnthis;
980996
};
981997

@@ -1006,7 +1022,7 @@ Readable.prototype.pause = function() {
10061022
this._readableState.flowing=false;
10071023
this.emit('pause');
10081024
}
1009-
this._readableState.paused=true;
1025+
this._readableState[kPaused]=true;
10101026
returnthis;
10111027
};
10121028

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
'use strict';
2+
constcommon=require('../common');
3+
4+
const{ Readable }=require('stream');
5+
6+
constreadable=newReadable({
7+
read(){}
8+
});
9+
10+
functionread(){}
11+
12+
readable.setEncoding('utf8');
13+
readable.on('readable',read);
14+
readable.removeListener('readable',read);
15+
16+
process.nextTick(function(){
17+
readable.on('data',common.mustCall());
18+
readable.push('hello');
19+
});

‎test/parallel/test-stream-readable-pause-and-resume.js‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
'use strict';
22

33
constcommon=require('../common');
4+
constassert=require('assert');
45
const{ Readable }=require('stream');
56

67
letticks=18;
@@ -38,3 +39,20 @@ function readAndPause() {
3839

3940
rs.on('data',ondata);
4041
}
42+
43+
{
44+
constreadable=newReadable({
45+
read(){}
46+
});
47+
48+
functionread(){}
49+
50+
readable.setEncoding('utf8');
51+
readable.on('readable',read);
52+
readable.removeListener('readable',read);
53+
readable.pause();
54+
55+
process.nextTick(function(){
56+
assert(readable.isPaused());
57+
});
58+
}

0 commit comments

Comments
 (0)