| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 8fe5807 commit 0ddfc6a
3 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -97,6 +97,7 @@ const { | |||
| 97 | 97 | ArrayBufferViewGetByteLength, | |
| 98 | 98 | ArrayBufferViewGetByteOffset, | |
| 99 | 99 | AsyncIterator, | |
| 100 | + Queue, | ||
| 100 | 101 | canCopyArrayBuffer, | |
| 101 | 102 | cloneAsUint8Array, | |
| 102 | 103 | copyArrayBuffer, | |
@@ -2192,9 +2193,9 @@ function readableStreamCancel(stream, reason) { | |||
| 2192 | 2193 | reader, | |
| 2193 | 2194 | } = stream[kState]; | |
| 2194 | 2195 | if (reader !== undefined && readableStreamHasBYOBReader(stream)) { | |
| 2195 | - for (let n = 0; n < reader[kState].readIntoRequests.length; n++) | ||
| 2196 | - reader[kState].readIntoRequests[n][kClose](); | ||
| 2197 | - reader[kState].readIntoRequests = []; | ||
| 2196 | + const readIntoRequests = reader[kState].readIntoRequests; | ||
| 2197 | + while (readIntoRequests.length) | ||
| 2198 | + readIntoRequests.shift()[kClose](); | ||
| 2198 | 2199 | } | |
| 2199 | 2200 | ||
| 2200 | 2201 | return PromisePrototypeThen( | |
@@ -2216,9 +2217,9 @@ function readableStreamClose(stream) { | |||
| 2216 | 2217 | reader[kState].close?.resolve(); | |
| 2217 | 2218 | ||
| 2218 | 2219 | if (readableStreamHasDefaultReader(stream)) { | |
| 2219 | - for (let n = 0; n < reader[kState].readRequests.length; n++) | ||
| 2220 | - reader[kState].readRequests[n][kClose](); | ||
| 2221 | - reader[kState].readRequests = []; | ||
| 2220 | + const readRequests = reader[kState].readRequests; | ||
| 2221 | + while (readRequests.length) | ||
| 2222 | + readRequests.shift()[kClose](); | ||
| 2222 | 2223 | } | |
| 2223 | 2224 | } | |
| 2224 | 2225 | ||
@@ -2246,14 +2247,14 @@ function readableStreamError(stream, error) { | |||
| 2246 | 2247 | } | |
| 2247 | 2248 | ||
| 2248 | 2249 | if (readableStreamHasDefaultReader(stream)) { | |
| 2249 | - for (let n = 0; n < reader[kState].readRequests.length; n++) | ||
| 2250 | - reader[kState].readRequests[n][kError](error); | ||
| 2251 | - reader[kState].readRequests = []; | ||
| 2250 | + const readRequests = reader[kState].readRequests; | ||
| 2251 | + while (readRequests.length) | ||
| 2252 | + readRequests.shift()[kError](error); | ||
| 2252 | 2253 | } else { | |
| 2253 | 2254 | assert(readableStreamHasBYOBReader(stream)); | |
| 2254 | - for (let n = 0; n < reader[kState].readIntoRequests.length; n++) | ||
| 2255 | - reader[kState].readIntoRequests[n][kError](error); | ||
| 2256 | - reader[kState].readIntoRequests = []; | ||
| 2255 | + const readIntoRequests = reader[kState].readIntoRequests; | ||
| 2256 | + while (readIntoRequests.length) | ||
| 2257 | + readIntoRequests.shift()[kError](error); | ||
| 2257 | 2258 | } | |
| 2258 | 2259 | } | |
| 2259 | 2260 | ||
@@ -2297,7 +2298,7 @@ function readableStreamFulfillReadRequest(stream, chunk, done) { | |||
| 2297 | 2298 | reader, | |
| 2298 | 2299 | } = stream[kState]; | |
| 2299 | 2300 | assert(reader[kState].readRequests.length); | |
| 2300 | - const readRequest = ArrayPrototypeShift(reader[kState].readRequests); | ||
| 2301 | + const readRequest = reader[kState].readRequests.shift(); | ||
| 2301 | 2302 | ||
| 2302 | 2303 | // TODO(@jasnell): It's not clear under what exact conditions done | |
| 2303 | 2304 | // will be true here. The spec requires this check but none of the | |
@@ -2315,7 +2316,7 @@ function readableStreamFulfillReadIntoRequest(stream, chunk, done) { | |||
| 2315 | 2316 | reader, | |
| 2316 | 2317 | } = stream[kState]; | |
| 2317 | 2318 | assert(reader[kState].readIntoRequests.length); | |
| 2318 | - const readIntoRequest = ArrayPrototypeShift(reader[kState].readIntoRequests); | ||
| 2319 | + const readIntoRequest = reader[kState].readIntoRequests.shift(); | ||
| 2319 | 2320 | if (done) | |
| 2320 | 2321 | readIntoRequest[kClose](chunk); | |
| 2321 | 2322 | else | |
@@ -2325,15 +2326,21 @@ function readableStreamFulfillReadIntoRequest(stream, chunk, done) { | |||
| 2325 | 2326 | function readableStreamAddReadRequest(stream, readRequest) { | |
| 2326 | 2327 | assert(readableStreamHasDefaultReader(stream)); | |
| 2327 | 2328 | assert(stream[kState].state === 'readable'); | |
| 2328 | - ArrayPrototypePush(stream[kState].reader[kState].readRequests, readRequest); | ||
| 2329 | + const readerState = stream[kState].reader[kState]; | ||
| 2330 | + let readRequests = readerState.readRequests; | ||
| 2331 | + if (readRequests === kEmptyQueue) | ||
| 2332 | + readRequests = readerState.readRequests = new Queue(); | ||
| 2333 | + readRequests.push(readRequest); | ||
| 2329 | 2334 | } | |
| 2330 | 2335 | ||
| 2331 | 2336 | function readableStreamAddReadIntoRequest(stream, readIntoRequest) { | |
| 2332 | 2337 | assert(readableStreamHasBYOBReader(stream)); | |
| 2333 | 2338 | assert(stream[kState].state !== 'errored'); | |
| 2334 | - ArrayPrototypePush( | ||
| 2335 | - stream[kState].reader[kState].readIntoRequests, | ||
| 2336 | - readIntoRequest); | ||
| 2339 | + const readerState = stream[kState].reader[kState]; | ||
| 2340 | + let readIntoRequests = readerState.readIntoRequests; | ||
| 2341 | + if (readIntoRequests === kEmptyQueue) | ||
| 2342 | + readIntoRequests = readerState.readIntoRequests = new Queue(); | ||
| 2343 | + readIntoRequests.push(readIntoRequest); | ||
| 2337 | 2344 | } | |
| 2338 | 2345 | ||
| 2339 | 2346 | function readableStreamReaderGenericCancel(reader, reason) { | |
@@ -2405,10 +2412,9 @@ function readableStreamDefaultReaderRelease(reader) { | |||
| 2405 | 2412 | } | |
| 2406 | 2413 | ||
| 2407 | 2414 | function readableStreamDefaultReaderErrorReadRequests(reader, e) { | |
| 2408 | - for (let n = 0; n < reader[kState].readRequests.length; ++n) { | ||
| 2409 | - reader[kState].readRequests[n][kError](e); | ||
| 2410 | - } | ||
| 2411 | - reader[kState].readRequests = []; | ||
| 2415 | + const readRequests = reader[kState].readRequests; | ||
| 2416 | + while (readRequests.length) | ||
| 2417 | + readRequests.shift()[kError](e); | ||
| 2412 | 2418 | } | |
| 2413 | 2419 | ||
| 2414 | 2420 | function readableStreamBYOBReaderRelease(reader) { | |
@@ -2420,10 +2426,9 @@ function readableStreamBYOBReaderRelease(reader) { | |||
| 2420 | 2426 | } | |
| 2421 | 2427 | ||
| 2422 | 2428 | function readableStreamBYOBReaderErrorReadIntoRequests(reader, e) { | |
| 2423 | - for (let n = 0; n < reader[kState].readIntoRequests.length; ++n) { | ||
| 2424 | - reader[kState].readIntoRequests[n][kError](e); | ||
| 2425 | - } | ||
| 2426 | - reader[kState].readIntoRequests = []; | ||
| 2429 | + const readIntoRequests = reader[kState].readIntoRequests; | ||
| 2430 | + while (readIntoRequests.length) | ||
| 2431 | + readIntoRequests.shift()[kError](e); | ||
| 2427 | 2432 | } | |
| 2428 | 2433 | ||
| 2429 | 2434 | function readableStreamReaderGenericRelease(reader) { | |
@@ -2495,14 +2500,18 @@ function setupReadableStreamBYOBReader(reader, stream) { | |||
| 2495 | 2500 | if (!isReadableByteStreamController(controller)) | |
| 2496 | 2501 | throw new ERR_INVALID_ARG_VALUE('stream', stream, 'must be a byte stream'); | |
| 2497 | 2502 | readableStreamReaderGenericInitialize(reader, stream); | |
| 2498 | - reader[kState].readIntoRequests = []; | ||
| 2503 | + // The read-request queues use the same ring buffer as [[queue]], drained | ||
| 2504 | + // from a moving head rather than with ArrayPrototypeShift. Start from the | ||
| 2505 | + // shared immutable empty queue so acquiring a reader allocates no request | ||
| 2506 | + // storage until a read actually parks. | ||
| 2507 | + reader[kState].readIntoRequests = kEmptyQueue; | ||
| 2499 | 2508 | } | |
| 2500 | 2509 | ||
| 2501 | 2510 | function setupReadableStreamDefaultReader(reader, stream) { | |
| 2502 | 2511 | if (isReadableStreamLocked(stream)) | |
| 2503 | 2512 | throw new ERR_INVALID_STATE.TypeError('ReadableStream is locked'); | |
| 2504 | 2513 | readableStreamReaderGenericInitialize(reader, stream); | |
| 2505 | - reader[kState].readRequests = []; | ||
| 2514 | + reader[kState].readRequests = kEmptyQueue; | ||
| 2506 | 2515 | } | |
| 2507 | 2516 | ||
| 2508 | 2517 | function readableStreamDefaultControllerClose(controller) { | |
@@ -3147,7 +3156,7 @@ function readableByteStreamControllerEnqueue(controller, chunk) { | |||
| 3147 | 3156 | } | |
| 3148 | 3157 | const transferredView = | |
| 3149 | 3158 | new Uint8Array(transferredBuffer, byteOffset, byteLength); | |
| 3150 | - const readRequest = ArrayPrototypeShift(readRequests); | ||
| 3159 | + const readRequest = readRequests.shift(); | ||
| 3151 | 3160 | readRequest[kChunk](transferredView); | |
| 3152 | 3161 | } | |
| 3153 | 3162 | } else if (readableStreamHasBYOBReader(stream)) { | |
@@ -3520,7 +3529,7 @@ function readableByteStreamControllerProcessReadRequestsUsingQueue(controller) { | |||
| 3520 | 3529 | } | |
| 3521 | 3530 | readableByteStreamControllerFillReadRequestFromQueue( | |
| 3522 | 3531 | controller, | |
| 3523 | - ArrayPrototypeShift(reader[kState].readRequests), | ||
| 3532 | + reader[kState].readRequests.shift(), | ||
| 3524 | 3533 | ); | |
| 3525 | 3534 | } | |
| 3526 | 3535 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -404,6 +404,7 @@ module.exports = { | |||
| 404 | 404 | ArrayBufferViewGetByteLength, | |
| 405 | 405 | ArrayBufferViewGetByteOffset, | |
| 406 | 406 | AsyncIterator, | |
| 407 | + Queue, | ||
| 407 | 408 | canCopyArrayBuffer, | |
| 408 | 409 | cloneAsUint8Array, | |
| 409 | 410 | copyArrayBuffer, | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1,8 +1,6 @@ | |||
| 1 | 1 | 'use strict'; | |
| 2 | 2 | ||
| 3 | 3 | const { | |
| 4 | - ArrayPrototypePush, | ||
| 5 | - ArrayPrototypeShift, | ||
| 6 | 4 | FunctionPrototypeBind, | |
| 7 | 5 | FunctionPrototypeCall, | |
| 8 | 6 | ObjectDefineProperties, | |
@@ -55,6 +53,7 @@ const { | |||
| 55 | 53 | } = require('internal/worker/js_transferable'); | |
| 56 | 54 | ||
| 57 | 55 | const { | |
| 56 | + Queue, | ||
| 58 | 57 | createPromiseCallbackNoParams, | |
| 59 | 58 | createPromiseCallback1Param, | |
| 60 | 59 | createPromiseCallback2Params, | |
@@ -613,7 +612,10 @@ function createWritableStreamState() { | |||
| 613 | 612 | controller: undefined, | |
| 614 | 613 | state: 'writable', | |
| 615 | 614 | 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, | ||
| 617 | 619 | writer: undefined, | |
| 618 | 620 | transfer: { | |
| 619 | 621 | __proto__: null, | |
@@ -823,7 +825,7 @@ function writableStreamRejectCloseAndClosedPromiseIfNeeded(stream) { | |||
| 823 | 825 | function writableStreamMarkFirstWriteRequestInFlight(stream) { | |
| 824 | 826 | assert(stream[kState].inFlightWriteRequest.promise === undefined); | |
| 825 | 827 | assert(stream[kState].writeRequests.length); | |
| 826 | - const writeRequest = ArrayPrototypeShift(stream[kState].writeRequests); | ||
| 828 | + const writeRequest = stream[kState].writeRequests.shift(); | ||
| 827 | 829 | stream[kState].inFlightWriteRequest = writeRequest; | |
| 828 | 830 | } | |
| 829 | 831 | ||
@@ -901,9 +903,9 @@ function writableStreamFinishErroring(stream) { | |||
| 901 | 903 | stream[kState].state = 'errored'; | |
| 902 | 904 | stream[kState].controller[kError](); | |
| 903 | 905 | 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); | ||
| 907 | 909 | ||
| 908 | 910 | if (stream[kState].pendingAbortRequest.abort.promise === undefined) { | |
| 909 | 911 | writableStreamRejectCloseAndClosedPromiseIfNeeded(stream); | |
@@ -952,7 +954,11 @@ function writableStreamAddWriteRequest(stream) { | |||
| 952 | 954 | // PromiseWithResolvers() already returns a { promise, resolve, reject } | |
| 953 | 955 | // record, so push it as-is instead of rebuilding an identical object. | |
| 954 | 956 | 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); | ||
| 956 | 962 | return writeRequest.promise; | |
| 957 | 963 | } | |
| 958 | 964 | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments