Skip to content

Commit 26102c2

Browse files
anonrigaduh95
authored andcommitted
stream: speed up flowing pipe of buffers
flow() pulls one already-buffered chunk and calls _read() for the next one on every iteration. That goes through the general read() path, which updates a holey buffer array and then pulls the chunk back out. While a synchronous byte-mode flow is in progress, keep that prefetched chunk on the readable state and emit it directly. _read() of the next chunk still runs before 'data', and a nested read() moves the chunk back onto the buffer. benchmark/streams/pipe.js is about 77% faster (15 runs). pipe-object-mode, readable-readall, and readable-bigread stay within noise. Assisted-by: a closed-source coding agent Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com> PR-URL: #66182 Reviewed-By: Matteo Collina <matteo.collina@gmail.com> Reviewed-By: Robert Nagy <ronagy@icloud.com> Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day> Reviewed-By: Zeyu "Alex" Yang <himself65@outlook.com>
1 parent 605938e commit 26102c2

2 files changed

Lines changed: 126 additions & 2 deletions

File tree

‎lib/internal/streams/destroy.js‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -182,6 +182,7 @@ function undestroy() {
182182
r.errored = null;
183183
r.errorEmitted = false;
184184
r.reading = false;
185+
r.fastChunk = null;
185186
r.ended = r.readable === false;
186187
r.endEmitted = r.readable === false;
187188
}

‎lib/internal/streams/readable.js‎

Lines changed: 125 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,9 @@ const kPaused = 1 << 26;
139139
const kDataListening = 1 << 27;
140140
const kEndScheduled = 1 << 28;
141141
const kEofReadablePending = 1 << 29;
142+
// Set only while flowSync() is inside _read(). push() then keeps the
143+
// chunk on state.fastChunk instead of the buffer array.
144+
const kFastPush = 1 << 30;
142145

143146
// TODO(benjamingr) it is likely slower to do it this way than with free functions
144147
function makeBitMapDescriptor(bit) {
@@ -296,6 +299,8 @@ function ReadableState(options, stream, isDuplex) {
296299
this.buffer = [];
297300
this.bufferIndex = 0;
298301
this.length = 0;
302+
// Chunk prefetched by flowSync(), kept off the buffer array.
303+
this.fastChunk = null;
299304
this.pipes = [];
300305

301306
// Should close be emitted on destroy. Defaults to true.
@@ -400,9 +405,29 @@ Readable.prototype[SymbolAsyncDispose] = async function() {
400405
// similar to how Writable.write() returns true if you should
401406
// write() some more.
402407
Readable.prototype.push = function(chunk, encoding) {
408+
const state = this._readableState;
409+
410+
// flowSync() is inside _read() and wants a single buffer parked on
411+
// fastChunk. A second push, a string, or EOF drops back to the buffer.
412+
if ((state[kState] & kFastPush) !== 0 && encoding == null &&
413+
state.fastChunk == null && chunk instanceof Buffer && chunk.length > 0) {
414+
state.fastChunk = chunk;
415+
state.length = chunk.length;
416+
state[kState] &= ~kReading;
417+
return chunk.length < state.highWaterMark;
418+
}
419+
if ((state[kState] & kFastPush) !== 0) {
420+
state[kState] &= ~kFastPush;
421+
if (state.fastChunk != null) {
422+
const first = state.fastChunk;
423+
state.fastChunk = null;
424+
state.length = 0;
425+
readableAddChunkPushByteMode(this, state, first);
426+
}
427+
}
428+
403429
debug('push', chunk);
404430

405-
const state = this._readableState;
406431
return (state[kState] & kObjectMode) === 0 ?
407432
readableAddChunkPushByteMode(this, state, chunk, encoding) :
408433
readableAddChunkPushObjectMode(this, state, chunk, encoding);
@@ -601,6 +626,8 @@ Readable.prototype.isPaused = function() {
601626
// Backwards compatibility.
602627
Readable.prototype.setEncoding = function(enc) {
603628
const state = this._readableState;
629+
if (state.fastChunk != null)
630+
materializeFastChunk(state);
604631

605632
const decoder = new StringDecoder(enc);
606633
state.decoder = decoder;
@@ -667,6 +694,11 @@ function howMuchToRead(n, state) {
667694

668695
// You can override either this method, or the async _read(n) below.
669696
Readable.prototype.read = function(n) {
697+
// A nested read() during flowSync()'s 'data' event must see the
698+
// prefetched chunk. Null for every read that is not inside that loop.
699+
if (this._readableState.fastChunk != null)
700+
materializeFastChunk(this._readableState);
701+
670702
debug('read', n);
671703
// Same as parseInt(undefined, 10), however V8 7.3 performance regressed
672704
// in this scenario, so we are doing it manually.
@@ -1336,9 +1368,91 @@ Readable.prototype.pause = function() {
13361368
function flow(stream) {
13371369
const state = stream._readableState;
13381370
debug('flow');
1371+
// Byte-mode pipe sits in read() to pull one already-buffered chunk and
1372+
// refill. That read is most of the per-chunk cost. flowSync() keeps the
1373+
// same prefetch order without the buffer array or the general read path.
1374+
if (flowSync(stream, state))
1375+
return;
13391376
while ((state[kState] & kFlowing) !== 0 && stream.read() !== null);
13401377
}
13411378

1379+
const kFastFlowNeed = kConstructed | kFlowing | kDataListening;
1380+
const kFastFlowBlock = kObjectMode | kDecoder | kEnded | kDestroyed |
1381+
kErrored | kPaused | kReading | kSync;
1382+
1383+
// Returns true when this call owned the flowing loop, including any
1384+
// fallback to read() after the fast path stops.
1385+
function flowSync(stream, state) {
1386+
const bits = state[kState];
1387+
if ((bits & kFastFlowNeed) !== kFastFlowNeed ||
1388+
(bits & kFastFlowBlock) !== 0 ||
1389+
!(state.highWaterMark > 0) ||
1390+
state.fastChunk != null) {
1391+
return false;
1392+
}
1393+
1394+
if (state.length !== 0) {
1395+
const buf = state.buffer;
1396+
const idx = state.bufferIndex;
1397+
const chunk = buf[idx];
1398+
// Only the one-chunk prefetch left by read(0) / the previous read.
1399+
if (buf.length !== idx + 1 || chunk == null || chunk.length !== state.length)
1400+
return false;
1401+
buf.length = 0;
1402+
state.bufferIndex = 0;
1403+
state.fastChunk = chunk;
1404+
}
1405+
1406+
while ((state[kState] & kFlowing) !== 0) {
1407+
const current = state.fastChunk;
1408+
if (current == null)
1409+
break;
1410+
if ((state[kState] & kFastFlowBlock) !== 0)
1411+
break;
1412+
1413+
// _read() of the next chunk runs before 'data', matching read().
1414+
state.fastChunk = null;
1415+
state.length = 0;
1416+
state[kState] |= kReading | kSync | kFastPush;
1417+
try {
1418+
stream._read(state.highWaterMark);
1419+
} catch (err) {
1420+
state[kState] &= ~(kSync | kFastPush);
1421+
errorOrDestroy(stream, err);
1422+
break;
1423+
}
1424+
state[kState] &= ~(kSync | kFastPush);
1425+
1426+
if ((state[kState] & (kErrorEmitted | kCloseEmitted)) === 0) {
1427+
state[kState] |= kDataEmitted;
1428+
stream.emit('data', current);
1429+
}
1430+
1431+
// Nested read() moved fastChunk into the buffer and may have refilled.
1432+
if (state.fastChunk == null && state.length !== 0)
1433+
break;
1434+
}
1435+
1436+
if (state.fastChunk != null)
1437+
materializeFastChunk(state);
1438+
1439+
while ((state[kState] & kFlowing) !== 0 && stream.read() !== null);
1440+
return true;
1441+
}
1442+
1443+
// Put fastChunk at the head of the buffer. state.length already counts it.
1444+
function materializeFastChunk(state) {
1445+
const chunk = state.fastChunk;
1446+
if (chunk == null)
1447+
return;
1448+
state.fastChunk = null;
1449+
if (state.bufferIndex > 0) {
1450+
state.buffer[--state.bufferIndex] = chunk;
1451+
} else {
1452+
state.buffer.unshift(chunk);
1453+
}
1454+
}
1455+
13421456
// Wrap an old-style stream as the async data source.
13431457
// This is *not* part of the readable stream interface.
13441458
// It is an ugly unfortunate mess of history.
@@ -1724,7 +1838,16 @@ ObjectDefineProperties(Readable.prototype, {
17241838
__proto__: null,
17251839
enumerable: false,
17261840
get: function() {
1727-
return this._readableState?.buffer;
1841+
const state = this._readableState;
1842+
if (state == null)
1843+
return undefined;
1844+
if (state.fastChunk == null)
1845+
return state.buffer;
1846+
if (state.bufferIndex === state.buffer.length)
1847+
return [state.fastChunk];
1848+
const out = state.buffer.slice(state.bufferIndex);
1849+
out.push(state.fastChunk);
1850+
return out;
17281851
},
17291852
},
17301853

0 commit comments

Comments
 (0)