Skip to content

Commit c438250

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 eaba4cd commit c438250

4 files changed

Lines changed: 119 additions & 98 deletions

File tree

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 47 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -2476,21 +2476,26 @@ function readableStreamDefaultControllerClose(controller) {
24762476
}
24772477

24782478
functionreadableStreamDefaultControllerEnqueue(controller,chunk){
2479-
if(!readableStreamDefaultControllerCanCloseOrEnqueue(controller))
2479+
// Equivalent to readableStreamDefaultControllerCanCloseOrEnqueue()
2480+
// followed by isReadableStreamLocked() and
2481+
// readableStreamGetNumReadRequests(), but with the state loaded once:
2482+
// this runs for every enqueued chunk.
2483+
constcontrollerState=controller[kState];
2484+
conststream=controllerState.stream;
2485+
if(controllerState.closeRequested||stream[kState].state!=='readable')
24802486
return;
24812487

2482-
const{
2483-
stream,
2484-
}=controller[kState];
2485-
2486-
if(isReadableStreamLocked(stream)&&
2487-
readableStreamGetNumReadRequests(stream)){
2488+
constreader=stream[kState].reader;
2489+
if(reader!==undefined&&
2490+
reader[kState]!==undefined&&
2491+
reader[kType]==='ReadableStreamDefaultReader'&&
2492+
reader[kState].readRequests.length){
24882493
readableStreamFulfillReadRequest(stream,chunk,false);
24892494
}else{
24902495
try{
24912496
constchunkSize=
24922497
FunctionPrototypeCall(
2493-
controller[kState].sizeAlgorithm,
2498+
controllerState.sizeAlgorithm,
24942499
undefined,
24952500
chunk);
24962501
enqueueValueWithSize(controller,chunk,chunkSize);
@@ -2529,22 +2534,27 @@ function readableStreamDefaultControllerGetDesiredSize(controller) {
25292534
}
25302535

25312536
functionreadableStreamDefaultControllerShouldCallPull(controller){
2532-
const{
2533-
stream,
2534-
}=controller[kState];
2535-
if(!readableStreamDefaultControllerCanCloseOrEnqueue(controller)||
2536-
!controller[kState].started)
2537+
// Single-pass version of the spec's predicate chain (CanCloseOrEnqueue,
2538+
// IsLocked, HasDefaultReader, GetNumReadRequests, GetDesiredSize): this
2539+
// runs at least once per chunk on every default-stream path. The
2540+
// desired-size computation is inlined because the stream state is
2541+
// already known to be 'readable' here.
2542+
constcontrollerState=controller[kState];
2543+
conststream=controllerState.stream;
2544+
if(controllerState.closeRequested||
2545+
stream[kState].state!=='readable'||
2546+
!controllerState.started)
25372547
returnfalse;
25382548

2539-
if(isReadableStreamLocked(stream)&&
2540-
readableStreamGetNumReadRequests(stream)){
2549+
constreader=stream[kState].reader;
2550+
if(reader!==undefined&&
2551+
reader[kState]!==undefined&&
2552+
reader[kType]==='ReadableStreamDefaultReader'&&
2553+
reader[kState].readRequests.length){
25412554
returntrue;
25422555
}
25432556

2544-
constdesiredSize=readableStreamDefaultControllerGetDesiredSize(controller);
2545-
assert(desiredSize!==null);
2546-
2547-
returndesiredSize>0;
2557+
returncontrollerState.highWaterMark-controllerState.queueTotalSize>0;
25482558
}
25492559

25502560
functionreadableStreamDefaultControllerCallPullIfNeeded(controller){
@@ -2794,28 +2804,29 @@ function readableByteStreamControllerGetDesiredSize(controller) {
27942804
}
27952805

27962806
functionreadableByteStreamControllerShouldCallPull(controller){
2797-
const{
2798-
stream,
2799-
}=controller[kState];
2807+
// Single-pass version of the spec's predicate chain (HasDefaultReader,
2808+
// GetNumReadRequests, HasBYOBReader, GetNumReadIntoRequests,
2809+
// GetDesiredSize): this runs at least once per chunk on every byte
2810+
// stream path. The desired-size computation is inlined because the
2811+
// stream state is already known to be 'readable' here.
2812+
constcontrollerState=controller[kState];
2813+
conststream=controllerState.stream;
28002814
if(stream[kState].state!=='readable'||
2801-
controller[kState].closeRequested||
2802-
!controller[kState].started){
2815+
controllerState.closeRequested||
2816+
!controllerState.started){
28032817
returnfalse;
28042818
}
2805-
if(readableStreamHasDefaultReader(stream)&&
2806-
readableStreamGetNumReadRequests(stream)>0){
2807-
returntrue;
2808-
}
2809-
2810-
if(readableStreamHasBYOBReader(stream)&&
2811-
readableStreamGetNumReadIntoRequests(stream)>0){
2812-
returntrue;
2819+
constreader=stream[kState].reader;
2820+
if(reader!==undefined&&reader[kState]!==undefined){
2821+
consttype=reader[kType];
2822+
if(type==='ReadableStreamDefaultReader'){
2823+
if(reader[kState].readRequests.length)returntrue;
2824+
}elseif(type==='ReadableStreamBYOBReader'){
2825+
if(reader[kState].readIntoRequests.length)returntrue;
2826+
}
28132827
}
28142828

2815-
constdesiredSize=readableByteStreamControllerGetDesiredSize(controller);
2816-
assert(desiredSize!==null);
2817-
2818-
returndesiredSize>0;
2829+
returncontrollerState.highWaterMark-controllerState.queueTotalSize>0;
28192830
}
28202831

28212832
functionreadableByteStreamControllerHandleQueueDrain(controller){

‎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+
functiontransformStreamBackpressureChangePromise(stream){
430+
conststate=stream[kState];
431+
return(state.backpressureChange??=PromiseWithResolvers()).promise;
432+
}
433+
432434
functiontransformStreamSetBackpressure(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+
conststate=stream[kState];
436+
assert(state.backpressure!==backpressure);
437+
constbackpressureChange=state.backpressureChange;
438+
if(backpressureChange!==undefined){
439+
state.backpressureChange=undefined;
440+
backpressureChange.resolve();
441+
}
442+
state.backpressure=backpressure;
438443
}
439444

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

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

646650
functiontransformStreamDefaultSourceCancelAlgorithm(stream,reason){

‎lib/internal/webstreams/util.js‎

Lines changed: 18 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -134,45 +134,44 @@ function isBrandCheck(brand) {
134134
};
135135
}
136136

137+
// The queue helpers below run once per chunk on the hot paths of every
138+
// default readable/writable stream, so they load the controller state a
139+
// single time and don't assert the existence of the queue fields (both
140+
// are unconditionally initialized during controller setup and only ever
141+
// replaced wholesale).
137142
functiondequeueValue(controller){
138-
assert(controller[kState].queue!==undefined);
139-
assert(controller[kState].queueTotalSize!==undefined);
140-
assert(controller[kState].queue.length);
143+
conststate=controller[kState];
144+
assert(state.queue.length);
141145
const{
142146
value,
143147
size,
144-
}=ArrayPrototypeShift(controller[kState].queue);
145-
controller[kState].queueTotalSize=
146-
MathMax(0,controller[kState].queueTotalSize-size);
148+
}=ArrayPrototypeShift(state.queue);
149+
state.queueTotalSize=MathMax(0,state.queueTotalSize-size);
147150
returnvalue;
148151
}
149152

150153
functionresetQueue(controller){
151-
assert(controller[kState].queue!==undefined);
152-
assert(controller[kState].queueTotalSize!==undefined);
153-
controller[kState].queue=[];
154-
controller[kState].queueTotalSize=0;
154+
conststate=controller[kState];
155+
state.queue=[];
156+
state.queueTotalSize=0;
155157
}
156158

157159
functionpeekQueueValue(controller){
158-
assert(controller[kState].queue!==undefined);
159-
assert(controller[kState].queueTotalSize!==undefined);
160-
assert(controller[kState].queue.length);
161-
returncontroller[kState].queue[0].value;
160+
conststate=controller[kState];
161+
assert(state.queue.length);
162+
returnstate.queue[0].value;
162163
}
163164

164165
functionenqueueValueWithSize(controller,value,size){
165-
assert(controller[kState].queue!==undefined);
166-
assert(controller[kState].queueTotalSize!==undefined);
166+
conststate=controller[kState];
167167
constcoercedSize=+size;
168168
if(NumberIsNaN(coercedSize)||
169169
coercedSize<0||
170170
coercedSize===Infinity){
171171
thrownewERR_INVALID_ARG_VALUE.RangeError('size',size);
172172
}
173-
size=coercedSize;
174-
ArrayPrototypePush(controller[kState].queue,{ value, size });
175-
controller[kState].queueTotalSize+=size;
173+
ArrayPrototypePush(state.queue,{ value,size: coercedSize});
174+
state.queueTotalSize+=coercedSize;
176175
}
177176

178177
// Arity-specialized variants of the promise-callback wrapper. The generic

0 commit comments

Comments
 (0)