Skip to content

Commit 4889fb0

Browse files
anonrignodejs-github-bot
authored andcommitted
stream: skip write() checks in flowing pipe
pipe() installs one 'data' listener that calls dest.write() for every chunk. That repeats encoding, mode, and end checks that stay the same for a synchronous buffer write. When that listener is still the only one, hand the Buffer to the same synchronous write path without those checks. A second listener, a non-buffer chunk, or a busy writable still goes through emit('data'). On top of the flowing-read fast path, benchmark/streams/pipe.js is about 31% faster (20 runs). Object-mode pipe and readable-readall 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 236b53f commit 4889fb0

2 files changed

Lines changed: 82 additions & 1 deletion

File tree

‎lib/internal/streams/readable.js‎

Lines changed: 45 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -301,6 +301,10 @@ function ReadableState(options, stream, isDuplex) {
301301
this.length = 0;
302302
// Chunk prefetched by flowSync(), kept off the buffer array.
303303
this.fastChunk = null;
304+
// The sole pipe() 'data' listener, when there is exactly one destination.
305+
// flowSync() writes buffers straight to that destination.
306+
this.pipeOnData = null;
307+
this.pipePause = null;
304308
this.pipes = [];
305309

306310
// Should close be emitted on destroy. Defaults to true.
@@ -985,6 +989,16 @@ Readable.prototype.pipe = function(dest, pipeOpts) {
985989
}
986990

987991
state.pipes.push(dest);
992+
// Only a single pipe destination can skip emit('data') and call write()
993+
// directly. A second destination, or an extra 'data' listener, must go
994+
// through emit so every listener still runs.
995+
if (state.pipes.length === 1) {
996+
state.pipeOnData = ondata;
997+
state.pipePause = pause;
998+
} else {
999+
state.pipeOnData = null;
1000+
state.pipePause = null;
1001+
}
9881002
debug('pipe count=%d opts=%j', state.pipes.length, pipeOpts);
9891003

9901004
const doEnd = (!pipeOpts || pipeOpts.end !== false) &&
@@ -1175,6 +1189,8 @@ Readable.prototype.unpipe = function(dest) {
11751189
// remove all.
11761190
const dests = state.pipes;
11771191
state.pipes = [];
1192+
state.pipeOnData = null;
1193+
state.pipePause = null;
11781194
this.pause();
11791195

11801196
for (let i = 0; i < dests.length; i++)
@@ -1188,6 +1204,8 @@ Readable.prototype.unpipe = function(dest) {
11881204
return this;
11891205

11901206
state.pipes.splice(index, 1);
1207+
state.pipeOnData = null;
1208+
state.pipePause = null;
11911209
if (state.pipes.length === 0)
11921210
this.pause();
11931211

@@ -1380,6 +1398,32 @@ const kFastFlowNeed = kConstructed | kFlowing | kDataListening;
13801398
const kFastFlowBlock = kObjectMode | kDecoder | kEnded | kDestroyed |
13811399
kErrored | kPaused | kReading | kSync;
13821400

1401+
let writeKnownBuffer;
1402+
1403+
// Pipe's only listener is ondata(), which calls dest.write(). Skip emit
1404+
// and the general write() checks for a single Buffer in that steady state.
1405+
function deliverFlowChunk(stream, state, chunk) {
1406+
const pipeOnData = state.pipeOnData;
1407+
const events = stream._events;
1408+
if (pipeOnData !== null && events !== undefined && events.data === pipeOnData) {
1409+
writeKnownBuffer ??= require('internal/streams/writable').writeKnownBuffer;
1410+
const dest = state.pipes[0];
1411+
let ret;
1412+
try {
1413+
ret = writeKnownBuffer(dest, chunk);
1414+
} catch (error) {
1415+
dest.destroy(error);
1416+
return;
1417+
}
1418+
if (ret === undefined)
1419+
stream.emit('data', chunk);
1420+
else if (ret === false && state.pipePause !== null)
1421+
state.pipePause();
1422+
return;
1423+
}
1424+
stream.emit('data', chunk);
1425+
}
1426+
13831427
// Returns true when this call owned the flowing loop, including any
13841428
// fallback to read() after the fast path stops.
13851429
function flowSync(stream, state) {
@@ -1425,7 +1469,7 @@ function flowSync(stream, state) {
14251469

14261470
if ((state[kState] & (kErrorEmitted | kCloseEmitted)) === 0) {
14271471
state[kState] |= kDataEmitted;
1428-
stream.emit('data', current);
1472+
deliverFlowChunk(stream, state, current);
14291473
}
14301474

14311475
// Nested read() moved fastChunk into the buffer and may have refilled.

‎lib/internal/streams/writable.js‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -586,6 +586,43 @@ function writeOrBuffer(stream, state, chunk, encoding, callback) {
586586
return ret && (state[kState] & (kDestroyed | kErrored)) === 0;
587587
}
588588

589+
// Steady state of a flowing pipe into a byte-mode Writable: one Buffer,
590+
// nothing queued, and no user write callback. Returns undefined when the
591+
// caller must use write() instead. Otherwise the same boolean as write().
592+
const kWriteFlowBlock = kObjectMode | kDestroyed | kErrored | kSync |
593+
kEnding | kFinished | kWriting | kCorked | kBuffered | kEnded |
594+
kNeedDrain | kWriteCb | kExpectWriteCb | kBufferProcessing |
595+
kFinalCalled | kPrefinished | kOnFinished | kErrorEmitted;
596+
597+
function writeKnownBuffer(stream, chunk) {
598+
const state = stream._writableState;
599+
if (state == null || state.length !== 0 || !(chunk instanceof Buffer))
600+
return undefined;
601+
602+
const bits = state[kState];
603+
if ((bits & kConstructed) === 0 || (bits & kWriteFlowBlock) !== 0)
604+
return undefined;
605+
606+
const len = chunk.length;
607+
state.pendingcb++;
608+
state.length = len;
609+
state.writelen = len;
610+
state[kState] = bits | kWriting | kSync | kExpectWriteCb;
611+
stream._write(chunk, 'buffer', state.onwrite);
612+
state[kState] &= ~kSync;
613+
614+
const ret = state.length < state.highWaterMark || state.length === 0;
615+
if (!ret)
616+
state[kState] |= kNeedDrain;
617+
return ret && (state[kState] & (kDestroyed | kErrored)) === 0;
618+
}
619+
620+
ObjectDefineProperty(Writable, 'writeKnownBuffer', {
621+
__proto__: null,
622+
value: writeKnownBuffer,
623+
enumerable: false,
624+
});
625+
589626
function doWrite(stream, state, writev, len, chunk, encoding, cb) {
590627
state.writelen = len;
591628
if (cb !== nop) {

0 commit comments

Comments
 (0)