stream: cut per-chunk overhead in WHATWG streams · nodejs/node@c81f894 · GitHub
Skip to content

Commit c81f894

Browse files
mcollinaaduh95
authored andcommitted
stream: cut per-chunk overhead in WHATWG streams
Consolidate the spec's per-chunk predicate chains (CanCloseOrEnqueue, IsLocked, HasDefaultReader, GetNumReadRequests, GetDesiredSize and the writable-side equivalents) into single passes over the controller and stream state, mirror "close queued or in flight" as a boolean flag maintained at the few close-request transition sites, and materialize the TransformStream [[backpressureChangePromise]] record lazily on first observation so backpressure flips nobody is waiting on allocate nothing. Assisted-by: Claude Code Signed-off-by: Matteo Collina <hello@matteocollina.com> PR-URL: #64252 Reviewed-By: Antoine du Hamel <duhamelantoine1995@gmail.com> Reviewed-By: Yagiz Nizipli <yagiz@nizipli.com> Reviewed-By: Filip Skokan <panva.ip@gmail.com>
1 parent d940f02 commit c81f894

4 files changed

Lines changed: 119 additions & 98 deletions

File tree

lib/internal/webstreams/readablestream.js

Lines changed: 47 additions & 36 deletions

lib/internal/webstreams/transformstream.js

Lines changed: 24 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -258,12 +258,7 @@ function InternalTransferredTransformStream() {
258258
readable: undefined,
259259
writable: undefined,
260260
backpressure: undefined,
261-
backpressureChange: {
262-
__proto__: null,
263-
promise: undefined,
264-
resolve: undefined,
265-
reject: undefined,
266-
},
261+
backpressureChange: undefined,
267262
controller: undefined,
268263
};
269264
}
@@ -390,12 +385,7 @@ function initializeTransformStream(
390385
writable,
391386
controller: undefined,
392387
backpressure: undefined,
393-
backpressureChange: {
394-
__proto__: null,
395-
promise: undefined,
396-
resolve: undefined,
397-
reject: undefined,
398-
},
388+
backpressureChange: undefined,
399389
};
400390

401391
transformStreamSetBackpressure(stream, true);
@@ -429,12 +419,27 @@ function transformStreamUnblockWrite(stream) {
429419
transformStreamSetBackpressure(stream, false);
430420
}
431421

422+
// The spec's [[backpressureChangePromise]] is only ever observed by the
423+
// source pull algorithm (settles when backpressure next becomes true) and
424+
// by a sink write arriving while backpressure is set (settles when
425+
// backpressure next becomes false). Instead of allocating a fresh promise
426+
// record on every flip, the record is materialized lazily on first
427+
// observation and dropped once settled; flips nobody is waiting on
428+
// allocate nothing.
429+
function transformStreamBackpressureChangePromise(stream) {
430+
const state = stream[kState];
431+
return (state.backpressureChange ??= PromiseWithResolvers()).promise;
432+
}
433+
432434
function transformStreamSetBackpressure(stream, backpressure) {
433-
assert(stream[kState].backpressure !== backpressure);
434-
if (stream[kState].backpressureChange.promise !== undefined)
435-
stream[kState].backpressureChange.resolve?.();
436-
stream[kState].backpressureChange = PromiseWithResolvers();
437-
stream[kState].backpressure = backpressure;
435+
const state = stream[kState];
436+
assert(state.backpressure !== backpressure);
437+
const backpressureChange = state.backpressureChange;
438+
if (backpressureChange !== undefined) {
439+
state.backpressureChange = undefined;
440+
backpressureChange.resolve();
441+
}
442+
state.backpressure = backpressure;
438443
}
439444

440445
function setupTransformStreamDefaultController(
@@ -554,7 +559,7 @@ function transformStreamDefaultSinkWriteAlgorithm(stream, chunk) {
554559
} = stream[kState];
555560
assert(writable[kState].state === 'writable');
556561
if (stream[kState].backpressure) {
557-
const backpressureChange = stream[kState].backpressureChange.promise;
562+
const backpressureChange = transformStreamBackpressureChangePromise(stream);
558563
return PromisePrototypeThen(
559564
backpressureChange,
560565
() => {
@@ -638,9 +643,8 @@ function transformStreamDefaultSinkCloseAlgorithm(stream) {
638643

639644
function transformStreamDefaultSourcePullAlgorithm(stream) {
640645
assert(stream[kState].backpressure);
641-
assert(stream[kState].backpressureChange.promise !== undefined);
642646
transformStreamSetBackpressure(stream, false);
643-
return stream[kState].backpressureChange.promise;
647+
return transformStreamBackpressureChangePromise(stream);
644648
}
645649

646650
function transformStreamDefaultSourceCancelAlgorithm(stream, reason) {

lib/internal/webstreams/util.js

Lines changed: 18 additions & 19 deletions

0 commit comments

Comments
 (0)