Skip to content

Commit b1596c2

Browse files
mcollinaaduh95
authored andcommitted
stream: cut per-chunk allocations in pipeTo
readableStreamPipeTo allocated, for every chunk written to the destination, a { promise, resolve, reject } write request record that it immediately marked as handled, and drove its loop with an async step()/run() pair whose implicit promises cost one allocation and one reaction per iteration. The parked-read path additionally allocated a read request object, a PromiseWithResolvers record, and a microtask closure per chunk; this is the steady state for pipeThrough, since a TransformStream's readable side has a high water mark of zero. Replace the per-write records with a single per-pipe tracker that the write request queue holds once per pending write and whose resolve()/reject() methods maintain a pending-write count, drive the pump loop with plain callbacks instead of async functions, and reuse one read request and one forwarding function across all chunks, the same pattern tee uses since c543cfb. Benchmark results (benchmark/compare.js --runs 20): webstreams/pipe-to.js +29.9% to +35.8% across all 16 configurations (all 99.9% confidence); a pipeThrough(TransformStream) passthrough loop improves ~17%; every other webstreams benchmark is unchanged. Signed-off-by: Matteo Collina <hello@matteocollina.com> PR-URL: #64890 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Yagiz Nizipli <yagiz@nizipli.com>
1 parent f0fc82d commit b1596c2

3 files changed

Lines changed: 310 additions & 61 deletions

File tree

lib/internal/webstreams/readablestream.js

Lines changed: 92 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,7 @@ const {
135135
writableStreamCloseQueuedOrInFlight,
136136
writableStreamDefaultWriterCloseWithErrorPropagation,
137137
writableStreamDefaultWriterRelease,
138-
writableStreamDefaultWriterWrite,
138+
writableStreamDefaultWriterWriteWithRequest,
139139
writerClosedPromise,
140140
writerReadyPromise,
141141
} = require('internal/webstreams/writablestream');
@@ -1526,8 +1526,38 @@ function readableStreamPipeTo(
15261526

15271527
const promise = PromiseWithResolvers();
15281528

1529-
const state = {
1530-
currentWrite: PromiseResolve(),
1529+
// One shared write request tracks every chunk written to the
1530+
// destination, instead of a { promise, resolve, reject } record per
1531+
// write. `stall` is armed by waitForPendingWrites() during shutdown;
1532+
// `failed`/`failure` latch a write that could not proceed.
1533+
const writeTracker = {
1534+
// Non-undefined: queue entries are discriminated from kNilRequest
1535+
// by `promise === undefined`.
1536+
promise: null,
1537+
pending: 0,
1538+
failed: false,
1539+
failure: undefined,
1540+
stall: null,
1541+
resolve() {
1542+
if (--this.pending === 0 && this.stall !== null) {
1543+
const stall = this.stall;
1544+
this.stall = null;
1545+
if (this.failed)
1546+
stall.reject(this.failure);
1547+
else
1548+
stall.resolve();
1549+
}
1550+
},
1551+
reject(error) {
1552+
this.pending--;
1553+
this.failed = true;
1554+
this.failure = error;
1555+
if (this.stall !== null) {
1556+
const stall = this.stall;
1557+
this.stall = null;
1558+
stall.reject(error);
1559+
}
1560+
},
15311561
};
15321562

15331563
// The error here can be undefined. The rejected arg
@@ -1544,11 +1574,14 @@ function readableStreamPipeTo(
15441574
promise.resolve();
15451575
}
15461576

1547-
async function waitForCurrentWrite() {
1548-
const write = state.currentWrite;
1549-
await write;
1550-
if (write !== state.currentWrite)
1551-
await waitForCurrentWrite();
1577+
function waitForPendingWrites() {
1578+
if (writeTracker.pending === 0) {
1579+
return writeTracker.failed ?
1580+
PromiseReject(writeTracker.failure) :
1581+
PromiseResolve();
1582+
}
1583+
writeTracker.stall = PromiseWithResolvers();
1584+
return writeTracker.stall.promise;
15521585
}
15531586

15541587
function shutdownWithAnAction(action, rejected, originalError) {
@@ -1557,7 +1590,7 @@ function readableStreamPipeTo(
15571590
if (dest[kState].state === 'writable' &&
15581591
!writableStreamCloseQueuedOrInFlight(dest)) {
15591592
PromisePrototypeThen(
1560-
waitForCurrentWrite(),
1593+
waitForPendingWrites(),
15611594
complete,
15621595
(error) => finalize(true, error));
15631596
return;
@@ -1578,7 +1611,7 @@ function readableStreamPipeTo(
15781611
if (dest[kState].state === 'writable' &&
15791612
!writableStreamCloseQueuedOrInFlight(dest)) {
15801613
PromisePrototypeThen(
1581-
waitForCurrentWrite(),
1614+
waitForPendingWrites(),
15821615
() => finalize(rejected, error),
15831616
(error) => finalize(true, error));
15841617
return;
@@ -1635,25 +1668,46 @@ function readableStreamPipeTo(
16351668
PromisePrototypeThen(promise, action, () => {});
16361669
}
16371670

1638-
async function step() {
1639-
if (shuttingDown) return true;
1671+
// The pump loop is callback-driven to avoid per-iteration promise
1672+
// allocations. At most one read is in flight at a time, so one read
1673+
// request and one forwarding function are reused for every chunk;
1674+
// the chunk travels through `pendingChunk`.
1675+
let pendingChunk;
1676+
let readRequest;
1677+
1678+
// Ready promise rejection is handled by the destination-errored
1679+
// watcher.
1680+
function ignoreReadyRejection() {}
1681+
1682+
function forwardChunk() {
1683+
const chunk = pendingChunk;
1684+
pendingChunk = undefined;
1685+
writableStreamDefaultWriterWriteWithRequest(writer, chunk, writeTracker);
1686+
pump();
1687+
}
1688+
1689+
function pump() {
1690+
if (shuttingDown) return;
16401691

16411692
if (dest[kState].backpressure) {
1642-
await writerReadyPromise(writer).promise;
1643-
if (shuttingDown) return true;
1693+
PromisePrototypeThen(
1694+
writerReadyPromise(writer).promise,
1695+
pump,
1696+
ignoreReadyRejection);
1697+
return;
16441698
}
16451699

16461700
const controller = source[kState].controller;
16471701

16481702
// Fast path: batch reads when data is buffered in a default controller.
1649-
// This avoids creating PipeToReadableStreamReadRequest objects and
1650-
// reduces promise allocation overhead.
1703+
// This avoids parking read requests and reduces promise allocation
1704+
// overhead.
16511705
if (source[kState].state === 'readable' &&
16521706
isReadableStreamDefaultController(controller) &&
16531707
controller[kState].queue.length > 0) {
16541708

16551709
while (controller[kState].queue.length > 0) {
1656-
if (shuttingDown) return true;
1710+
if (shuttingDown) return;
16571711

16581712
const chunk = dequeueValue(controller);
16591713

@@ -1664,8 +1718,7 @@ function readableStreamPipeTo(
16641718

16651719
// Write the chunk - we're already in a separate microtask from enqueue
16661720
// because we awaited the writer ready promise above.
1667-
state.currentWrite = writableStreamDefaultWriterWrite(writer, chunk);
1668-
setPromiseHandled(state.currentWrite);
1721+
writableStreamDefaultWriterWriteWithRequest(writer, chunk, writeTracker);
16691722

16701723
// Check backpressure after each write
16711724
if (dest[kState].backpressure) {
@@ -1682,24 +1735,29 @@ function readableStreamPipeTo(
16821735

16831736
// Check if stream closed during batch
16841737
if (source[kState].state === 'closed') {
1685-
return true;
1738+
return;
16861739
}
16871740

1688-
// Yield to microtask queue between batches to allow events/signals to fire
1689-
return false;
1741+
// Yield to microtask queue between batches to allow events/signals
1742+
// to fire
1743+
queueMicrotask(pump);
1744+
return;
16901745
}
16911746

1692-
// Slow path: use read request for async reads
1693-
const promise = PromiseWithResolvers();
1694-
// eslint-disable-next-line no-use-before-define
1695-
readableStreamDefaultReaderRead(reader, new PipeToReadableStreamReadRequest(writer, state, promise));
1696-
1697-
return promise.promise;
1698-
}
1699-
1700-
async function run() {
1701-
// Run until step resolves as true
1702-
while (!await step());
1747+
// Slow path: park a lazily materialized read request. Close and
1748+
// error are handled by the source watchers.
1749+
readRequest ??= {
1750+
[kChunk](chunk) {
1751+
// Per spec, pipeTo must queue a microtask for the write to avoid
1752+
// synchronous write during enqueue(). See WHATWG Streams spec
1753+
// "ReadableStreamPipeTo" step 15's "chunk steps".
1754+
pendingChunk = chunk;
1755+
queueMicrotask(forwardChunk);
1756+
},
1757+
[kClose]() {},
1758+
[kError]() {},
1759+
};
1760+
readableStreamDefaultReaderRead(reader, readRequest);
17031761
}
17041762

17051763
if (signal !== undefined) {
@@ -1711,7 +1769,7 @@ function readableStreamPipeTo(
17111769
disposable = addAbortListener(signal, abortAlgorithm);
17121770
}
17131771

1714-
setPromiseHandled(run());
1772+
pump();
17151773

17161774
watchErrored(source, readerClosedPromise(reader).promise, (error) => {
17171775
if (!preventAbort) {
@@ -1756,33 +1814,6 @@ function readableStreamPipeTo(
17561814
return promise.promise;
17571815
}
17581816

1759-
class PipeToReadableStreamReadRequest {
1760-
constructor(writer, state, promise) {
1761-
this.writer = writer;
1762-
this.state = state;
1763-
this.promise = promise;
1764-
}
1765-
1766-
[kChunk](chunk) {
1767-
// Per spec, pipeTo must queue a microtask for the write to avoid
1768-
// synchronous write during enqueue(). See WHATWG Streams spec
1769-
// "ReadableStreamPipeTo" step 15's "chunk steps".
1770-
queueMicrotask(() => {
1771-
this.state.currentWrite = writableStreamDefaultWriterWrite(this.writer, chunk);
1772-
setPromiseHandled(this.state.currentWrite);
1773-
this.promise.resolve(false);
1774-
});
1775-
}
1776-
1777-
[kClose]() {
1778-
this.promise.resolve(true);
1779-
}
1780-
1781-
[kError](error) {
1782-
this.promise.reject(error);
1783-
}
1784-
}
1785-
17861817
function readableStreamTee(stream, cloneForBranch2) {
17871818
if (isReadableByteStreamController(stream[kState].controller)) {
17881819
return readableByteStreamTee(stream);

lib/internal/webstreams/writablestream.js

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1001,6 +1001,61 @@ function writableStreamDefaultWriterWrite(writer, chunk) {
10011001
return promise;
10021002
}
10031003

1004+
// Variant of writableStreamDefaultWriterWrite for pipeTo: the caller
1005+
// provides a shared request object instead of a per-write promise record.
1006+
// `pending` is incremented before the controller write, which can settle
1007+
// requests synchronously when the stream starts erroring; precondition
1008+
// failures are latched on `failed`/`failure`.
1009+
function writableStreamDefaultWriterWriteWithRequest(writer, chunk, request) {
1010+
const writerState = writer[kState];
1011+
const stream = writerState.stream;
1012+
assert(stream !== undefined);
1013+
const streamState = stream[kState];
1014+
const {
1015+
controller,
1016+
} = streamState;
1017+
const chunkSize = writableStreamDefaultControllerGetChunkSize(
1018+
controller,
1019+
chunk);
1020+
if (stream !== writerState.stream) {
1021+
request.failed = true;
1022+
request.failure =
1023+
new ERR_INVALID_STATE.TypeError('Mismatched WritableStreams');
1024+
return;
1025+
}
1026+
const {
1027+
state,
1028+
} = streamState;
1029+
1030+
if (state === 'errored') {
1031+
request.failed = true;
1032+
request.failure = streamState.storedError;
1033+
return;
1034+
}
1035+
1036+
if (streamState.closeQueuedOrInFlight || state === 'closed') {
1037+
request.failed = true;
1038+
request.failure =
1039+
new ERR_INVALID_STATE.TypeError('WritableStream is closed');
1040+
return;
1041+
}
1042+
1043+
if (state === 'erroring') {
1044+
request.failed = true;
1045+
request.failure = streamState.storedError;
1046+
return;
1047+
}
1048+
1049+
assert(state === 'writable');
1050+
1051+
let writeRequests = streamState.writeRequests;
1052+
if (writeRequests === kEmptyQueue)
1053+
writeRequests = streamState.writeRequests = new Queue();
1054+
writeRequests.push(request);
1055+
request.pending++;
1056+
writableStreamDefaultControllerWrite(controller, chunk, chunkSize);
1057+
}
1058+
10041059
function writableStreamDefaultWriterRelease(writer) {
10051060
const {
10061061
stream,
@@ -1376,6 +1431,7 @@ module.exports = {
13761431
writableStreamCloseQueuedOrInFlight,
13771432
writableStreamAddWriteRequest,
13781433
writableStreamDefaultWriterWrite,
1434+
writableStreamDefaultWriterWriteWithRequest,
13791435
writableStreamDefaultWriterRelease,
13801436
writableStreamDefaultWriterGetDesiredSize,
13811437
writableStreamDefaultWriterEnsureReadyPromiseRejected,

0 commit comments

Comments
 (0)