Skip to content

Commit 3413ab9

Browse files
mcollinaaduh95
authored andcommitted
stream: use the ring buffer for WHATWG stream request queues
The [[queue]] backing the controllers became a ring buffer, but the read and write request queues (readRequests, readIntoRequests, writeRequests) were left as plain arrays consumed with ArrayPrototypeShift, which is O(n) and, even for the single pending request of the await-each regime, far slower than an indexed head advance. Back them with the same Queue, materialized lazily from the shared empty queue so acquiring a reader or constructing a writer allocates no request storage until a read or write actually parks. pipe-to: +4.7% to +9.4% (all 16 configs, ***) readable-read type=byob: +2.3% (**) parked read loop: +15%, write loop: +11% (local harness) Follow-up to #64312. Signed-off-by: Matteo Collina <hello@matteocollina.com> PR-URL: #64431 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Antoine du Hamel <duhamelantoine1995@gmail.com>
1 parent 3194059 commit 3413ab9

3 files changed

Lines changed: 54 additions & 38 deletions

File tree

lib/internal/webstreams/readablestream.js

Lines changed: 39 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -101,6 +101,7 @@ const {
101101
ArrayBufferViewGetByteLength,
102102
ArrayBufferViewGetByteOffset,
103103
AsyncIterator,
104+
Queue,
104105
canCopyArrayBuffer,
105106
cloneAsUint8Array,
106107
copyArrayBuffer,
@@ -2196,9 +2197,9 @@ function readableStreamCancel(stream, reason) {
21962197
reader,
21972198
} = stream[kState];
21982199
if (reader !== undefined && readableStreamHasBYOBReader(stream)) {
2199-
for (let n = 0; n < reader[kState].readIntoRequests.length; n++)
2200-
reader[kState].readIntoRequests[n][kClose]();
2201-
reader[kState].readIntoRequests = [];
2200+
const readIntoRequests = reader[kState].readIntoRequests;
2201+
while (readIntoRequests.length)
2202+
readIntoRequests.shift()[kClose]();
22022203
}
22032204

22042205
return PromisePrototypeThen(
@@ -2220,9 +2221,9 @@ function readableStreamClose(stream) {
22202221
reader[kState].close?.resolve();
22212222

22222223
if (readableStreamHasDefaultReader(stream)) {
2223-
for (let n = 0; n < reader[kState].readRequests.length; n++)
2224-
reader[kState].readRequests[n][kClose]();
2225-
reader[kState].readRequests = [];
2224+
const readRequests = reader[kState].readRequests;
2225+
while (readRequests.length)
2226+
readRequests.shift()[kClose]();
22262227
}
22272228
}
22282229

@@ -2250,14 +2251,14 @@ function readableStreamError(stream, error) {
22502251
}
22512252

22522253
if (readableStreamHasDefaultReader(stream)) {
2253-
for (let n = 0; n < reader[kState].readRequests.length; n++)
2254-
reader[kState].readRequests[n][kError](error);
2255-
reader[kState].readRequests = [];
2254+
const readRequests = reader[kState].readRequests;
2255+
while (readRequests.length)
2256+
readRequests.shift()[kError](error);
22562257
} else {
22572258
assert(readableStreamHasBYOBReader(stream));
2258-
for (let n = 0; n < reader[kState].readIntoRequests.length; n++)
2259-
reader[kState].readIntoRequests[n][kError](error);
2260-
reader[kState].readIntoRequests = [];
2259+
const readIntoRequests = reader[kState].readIntoRequests;
2260+
while (readIntoRequests.length)
2261+
readIntoRequests.shift()[kError](error);
22612262
}
22622263
}
22632264

@@ -2301,7 +2302,7 @@ function readableStreamFulfillReadRequest(stream, chunk, done) {
23012302
reader,
23022303
} = stream[kState];
23032304
assert(reader[kState].readRequests.length);
2304-
const readRequest = ArrayPrototypeShift(reader[kState].readRequests);
2305+
const readRequest = reader[kState].readRequests.shift();
23052306

23062307
// TODO(@jasnell): It's not clear under what exact conditions done
23072308
// will be true here. The spec requires this check but none of the
@@ -2319,7 +2320,7 @@ function readableStreamFulfillReadIntoRequest(stream, chunk, done) {
23192320
reader,
23202321
} = stream[kState];
23212322
assert(reader[kState].readIntoRequests.length);
2322-
const readIntoRequest = ArrayPrototypeShift(reader[kState].readIntoRequests);
2323+
const readIntoRequest = reader[kState].readIntoRequests.shift();
23232324
if (done)
23242325
readIntoRequest[kClose](chunk);
23252326
else
@@ -2329,15 +2330,21 @@ function readableStreamFulfillReadIntoRequest(stream, chunk, done) {
23292330
function readableStreamAddReadRequest(stream, readRequest) {
23302331
assert(readableStreamHasDefaultReader(stream));
23312332
assert(stream[kState].state === 'readable');
2332-
ArrayPrototypePush(stream[kState].reader[kState].readRequests, readRequest);
2333+
const readerState = stream[kState].reader[kState];
2334+
let readRequests = readerState.readRequests;
2335+
if (readRequests === kEmptyQueue)
2336+
readRequests = readerState.readRequests = new Queue();
2337+
readRequests.push(readRequest);
23332338
}
23342339

23352340
function readableStreamAddReadIntoRequest(stream, readIntoRequest) {
23362341
assert(readableStreamHasBYOBReader(stream));
23372342
assert(stream[kState].state !== 'errored');
2338-
ArrayPrototypePush(
2339-
stream[kState].reader[kState].readIntoRequests,
2340-
readIntoRequest);
2343+
const readerState = stream[kState].reader[kState];
2344+
let readIntoRequests = readerState.readIntoRequests;
2345+
if (readIntoRequests === kEmptyQueue)
2346+
readIntoRequests = readerState.readIntoRequests = new Queue();
2347+
readIntoRequests.push(readIntoRequest);
23412348
}
23422349

23432350
function readableStreamReaderGenericCancel(reader, reason) {
@@ -2409,10 +2416,9 @@ function readableStreamDefaultReaderRelease(reader) {
24092416
}
24102417

24112418
function readableStreamDefaultReaderErrorReadRequests(reader, e) {
2412-
for (let n = 0; n < reader[kState].readRequests.length; ++n) {
2413-
reader[kState].readRequests[n][kError](e);
2414-
}
2415-
reader[kState].readRequests = [];
2419+
const readRequests = reader[kState].readRequests;
2420+
while (readRequests.length)
2421+
readRequests.shift()[kError](e);
24162422
}
24172423

24182424
function readableStreamBYOBReaderRelease(reader) {
@@ -2424,10 +2430,9 @@ function readableStreamBYOBReaderRelease(reader) {
24242430
}
24252431

24262432
function readableStreamBYOBReaderErrorReadIntoRequests(reader, e) {
2427-
for (let n = 0; n < reader[kState].readIntoRequests.length; ++n) {
2428-
reader[kState].readIntoRequests[n][kError](e);
2429-
}
2430-
reader[kState].readIntoRequests = [];
2433+
const readIntoRequests = reader[kState].readIntoRequests;
2434+
while (readIntoRequests.length)
2435+
readIntoRequests.shift()[kError](e);
24312436
}
24322437

24332438
function readableStreamReaderGenericRelease(reader) {
@@ -2499,14 +2504,18 @@ function setupReadableStreamBYOBReader(reader, stream) {
24992504
if (!isReadableByteStreamController(controller))
25002505
throw new ERR_INVALID_ARG_VALUE('stream', stream, 'must be a byte stream');
25012506
readableStreamReaderGenericInitialize(reader, stream);
2502-
reader[kState].readIntoRequests = [];
2507+
// The read-request queues use the same ring buffer as [[queue]], drained
2508+
// from a moving head rather than with ArrayPrototypeShift. Start from the
2509+
// shared immutable empty queue so acquiring a reader allocates no request
2510+
// storage until a read actually parks.
2511+
reader[kState].readIntoRequests = kEmptyQueue;
25032512
}
25042513

25052514
function setupReadableStreamDefaultReader(reader, stream) {
25062515
if (isReadableStreamLocked(stream))
25072516
throw new ERR_INVALID_STATE.TypeError('ReadableStream is locked');
25082517
readableStreamReaderGenericInitialize(reader, stream);
2509-
reader[kState].readRequests = [];
2518+
reader[kState].readRequests = kEmptyQueue;
25102519
}
25112520

25122521
function readableStreamDefaultControllerClose(controller) {
@@ -3151,7 +3160,7 @@ function readableByteStreamControllerEnqueue(controller, chunk) {
31513160
}
31523161
const transferredView =
31533162
new Uint8Array(transferredBuffer, byteOffset, byteLength);
3154-
const readRequest = ArrayPrototypeShift(readRequests);
3163+
const readRequest = readRequests.shift();
31553164
readRequest[kChunk](transferredView);
31563165
}
31573166
} else if (readableStreamHasBYOBReader(stream)) {
@@ -3524,7 +3533,7 @@ function readableByteStreamControllerProcessReadRequestsUsingQueue(controller) {
35243533
}
35253534
readableByteStreamControllerFillReadRequestFromQueue(
35263535
controller,
3527-
ArrayPrototypeShift(reader[kState].readRequests),
3536+
reader[kState].readRequests.shift(),
35283537
);
35293538
}
35303539
}

lib/internal/webstreams/util.js

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -404,6 +404,7 @@ module.exports = {
404404
ArrayBufferViewGetByteLength,
405405
ArrayBufferViewGetByteOffset,
406406
AsyncIterator,
407+
Queue,
407408
canCopyArrayBuffer,
408409
cloneAsUint8Array,
409410
copyArrayBuffer,

lib/internal/webstreams/writablestream.js

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,6 @@
11
'use strict';
22

33
const {
4-
ArrayPrototypePush,
5-
ArrayPrototypeShift,
64
FunctionPrototypeBind,
75
FunctionPrototypeCall,
86
ObjectDefineProperties,
@@ -55,6 +53,7 @@ const {
5553
} = require('internal/worker/js_transferable');
5654

5755
const {
56+
Queue,
5857
createPromiseCallbackNoParams,
5958
createPromiseCallback1Param,
6059
createPromiseCallback2Params,
@@ -613,7 +612,10 @@ function createWritableStreamState() {
613612
controller: undefined,
614613
state: 'writable',
615614
storedError: undefined,
616-
writeRequests: [],
615+
// Ring-buffer request queue, materialized lazily on the first pending
616+
// write (see writableStreamAddWriteRequest) so construction allocates
617+
// no request storage.
618+
writeRequests: kEmptyQueue,
617619
writer: undefined,
618620
transfer: {
619621
__proto__: null,
@@ -823,7 +825,7 @@ function writableStreamRejectCloseAndClosedPromiseIfNeeded(stream) {
823825
function writableStreamMarkFirstWriteRequestInFlight(stream) {
824826
assert(stream[kState].inFlightWriteRequest.promise === undefined);
825827
assert(stream[kState].writeRequests.length);
826-
const writeRequest = ArrayPrototypeShift(stream[kState].writeRequests);
828+
const writeRequest = stream[kState].writeRequests.shift();
827829
stream[kState].inFlightWriteRequest = writeRequest;
828830
}
829831

@@ -901,9 +903,9 @@ function writableStreamFinishErroring(stream) {
901903
stream[kState].state = 'errored';
902904
stream[kState].controller[kError]();
903905
const storedError = stream[kState].storedError;
904-
for (let n = 0; n < stream[kState].writeRequests.length; n++)
905-
stream[kState].writeRequests[n].reject(storedError);
906-
stream[kState].writeRequests = [];
906+
const writeRequests = stream[kState].writeRequests;
907+
while (writeRequests.length)
908+
writeRequests.shift().reject(storedError);
907909

908910
if (stream[kState].pendingAbortRequest.abort.promise === undefined) {
909911
writableStreamRejectCloseAndClosedPromiseIfNeeded(stream);
@@ -952,7 +954,11 @@ function writableStreamAddWriteRequest(stream) {
952954
// PromiseWithResolvers() already returns a { promise, resolve, reject }
953955
// record, so push it as-is instead of rebuilding an identical object.
954956
const writeRequest = PromiseWithResolvers();
955-
ArrayPrototypePush(stream[kState].writeRequests, writeRequest);
957+
const streamState = stream[kState];
958+
let writeRequests = streamState.writeRequests;
959+
if (writeRequests === kEmptyQueue)
960+
writeRequests = streamState.writeRequests = new Queue();
961+
writeRequests.push(writeRequest);
956962
return writeRequest.promise;
957963
}
958964

0 commit comments

Comments
 (0)