Skip to content

Commit 3882929

Browse files
mcollinaaduh95
authored andcommitted
stream: skip unobserved 'readable' emission at EOF
When a stream reaches its end via push(null) while nobody is observing it (no 'readable' listener, not flowing, no pending readable need), do not schedule the deferred 'readable' emission: it fires into the void and costs a tick. Record that the end-of-stream notification is owed instead; attaching a 'readable' listener later emits it, and read() and resume() reach the end of the stream on their own. As a side effect, a 'readable' listener attached more than a tick after an unobserved end now receives the owed 'readable' before 'end', where it previously received only 'end'. For an HTTP server this removes one nextTick and one dead emit per request whose body the handler never reads. Signed-off-by: Matteo Collina <hello@matteocollina.com> PR-URL: #65749 Reviewed-By: Robert Nagy <ronagy@icloud.com> Reviewed-By: Paolo Insogna <paolo@cowtech.it> Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 96d4c23 commit 3882929

2 files changed

Lines changed: 99 additions & 4 deletions

File tree

‎lib/internal/streams/readable.js‎

Lines changed: 23 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -131,6 +131,7 @@ const kHasPaused = 1 << 25;
131131
const kPaused = 1 << 26;
132132
const kDataListening = 1 << 27;
133133
const kEndScheduled = 1 << 28;
134+
const kEofReadablePending = 1 << 29;
134135

135136
// TODO(benjamingr) it is likely slower to do it this way than with free functions
136137
function makeBitMapDescriptor(bit) {
@@ -670,7 +671,7 @@ Readable.prototype.read = function(n) {
670671
state.highWaterMark = computeNewHighWaterMark(n);
671672

672673
if (n !== 0)
673-
state[kState] &= ~kEmittedReadable;
674+
state[kState] &= ~(kEmittedReadable | kEofReadablePending);
674675

675676
// If we're doing read(0) to trigger a readable event, but we
676677
// already have a bunch of data in the buffer, then just trigger
@@ -813,7 +814,16 @@ function onEofChunk(stream, state) {
813814
// If we are sync, wait until next tick to emit the data.
814815
// Otherwise we risk emitting data in the flow()
815816
// the readable code triggers during a read() call.
816-
emitReadable(stream);
817+
if ((state[kState] & (kFlowing | kNeedReadable)) !== 0 ||
818+
stream.listenerCount('readable') > 0) {
819+
emitReadable(stream);
820+
} else {
821+
// Nobody is observing the stream: do not schedule the 'readable'
822+
// emission at all. A 'readable' listener attached later redeems it
823+
// (see Readable.prototype.on), and read() and resume() reach the
824+
// end of the stream on their own.
825+
state[kState] |= kEofReadablePending;
826+
}
817827
} else {
818828
// Emit 'readable' now to make sure it gets picked up.
819829
state[kState] &= ~kNeedReadable;
@@ -1162,8 +1172,17 @@ Readable.prototype.on = function(ev, fn) {
11621172
debug('on readable');
11631173
if (state.length) {
11641174
emitReadable(this);
1165-
} else if ((state[kState] & kReading) === 0) {
1166-
process.nextTick(nReadingNextTick, this);
1175+
} else {
1176+
if ((state[kState] & kEofReadablePending) !== 0) {
1177+
// The end-of-stream 'readable' emission was skipped because
1178+
// nobody was listening when the stream ended (see onEofChunk):
1179+
// emit it now.
1180+
state[kState] &= ~kEofReadablePending;
1181+
emitReadable(this);
1182+
}
1183+
if ((state[kState] & kReading) === 0) {
1184+
process.nextTick(nReadingNextTick, this);
1185+
}
11671186
}
11681187
}
11691188
}
Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
'use strict';
2+
const common = require('../common');
3+
const { Readable } = require('stream');
4+
const assert = require('assert');
5+
6+
// When a stream reaches its end while nobody is observing it, the
7+
// end-of-stream 'readable' emission is not scheduled. It is emitted
8+
// later if a 'readable' listener is attached, and read()/resume()
9+
// still reach 'end' on their own.
10+
11+
{
12+
// Listener attached synchronously after push(null) gets 'readable'
13+
// and then 'end'.
14+
const r = new Readable({ read() {} });
15+
r.push(null);
16+
let readableEmitted = false;
17+
r.on('readable', common.mustCall(() => {
18+
readableEmitted = true;
19+
assert.strictEqual(r.read(), null);
20+
}));
21+
r.on('end', common.mustCall(() => {
22+
assert.strictEqual(readableEmitted, true);
23+
}));
24+
}
25+
26+
{
27+
// Listener attached one macrotask after the unobserved end still gets
28+
// the owed 'readable' before 'end'.
29+
const r = new Readable({ read() {} });
30+
r.push(null);
31+
setImmediate(common.mustCall(() => {
32+
let readableEmitted = false;
33+
r.on('readable', common.mustCall(() => {
34+
readableEmitted = true;
35+
assert.strictEqual(r.read(), null);
36+
}));
37+
r.on('end', common.mustCall(() => {
38+
assert.strictEqual(readableEmitted, true);
39+
}));
40+
}));
41+
}
42+
43+
{
44+
// A stream that ends unobserved still emits 'end' when resumed later.
45+
const r = new Readable({ read() {} });
46+
r.push(null);
47+
setImmediate(common.mustCall(() => {
48+
r.resume();
49+
r.on('end', common.mustCall());
50+
}));
51+
}
52+
53+
{
54+
// read() after an unobserved end consumes the owed notification: a
55+
// 'readable' listener attached afterwards does not receive it, but
56+
// 'end' is still emitted.
57+
const r = new Readable({ read() {} });
58+
r.push(null);
59+
setImmediate(common.mustCall(() => {
60+
assert.strictEqual(r.read(), null);
61+
r.on('end', common.mustCall());
62+
}));
63+
}
64+
65+
{
66+
// A 'data' listener attached after an unobserved end still gets 'end'.
67+
const r = new Readable({ read() {} });
68+
r.push('x');
69+
r.push(null);
70+
setImmediate(common.mustCall(() => {
71+
r.on('data', common.mustCall((chunk) => {
72+
assert.strictEqual(chunk.toString(), 'x');
73+
}));
74+
r.on('end', common.mustCall());
75+
}));
76+
}

0 commit comments

Comments
 (0)