Skip to content

Commit 7f0def5

Browse files
anonrigcursoragent
authored andcommitted
stream: speed up WHATWG web streams
Avoid per-chunk async wrappers for sync pull/write/start and complete pipeTo writes without one microtask per chunk. Add a native webstreams binding with a Fast API isNonThenable check on the data plane and a memcpy clone for byte views. Empty stream construction skips redundant validation and lazily creates the writable AbortController, materializing it on abort() so controller.signal still reflects the abort reason. Use the shared kResolvedPromise on the pull/write hot path instead of allocating PromiseResolve(). Assisted-by: Grok Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com>
1 parent 9f0ce45 commit 7f0def5

12 files changed

Lines changed: 598 additions & 63 deletions

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 63 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -111,8 +111,10 @@ const {
111111
extractSizeAlgorithm,
112112
getNonWritablePropertyDescriptor,
113113
isBrandCheck,
114+
isNonThenable,
114115
kEmptyQueue,
115116
kResolvedPromise,
117+
promiseFromAlgorithmResult,
116118
kState,
117119
kType,
118120
lazyTransfer,
@@ -252,6 +254,16 @@ class ReadableStream {
252254
*/
253255
constructor(source = kEmptyObject, strategy = kEmptyObject) {
254256
markTransferMode(this, false, true);
257+
// Empty-argument `new ReadableStream()`: no source, no strategy, and
258+
// no controller. Reads never deliver data, so skip those allocations
259+
// until getReader/cancel/error first need a default controller.
260+
// Subclasses that call those methods after super() materialize the
261+
// controller in the subclass constructor; see
262+
// ensureEmptyDefaultController.
263+
if (source === kEmptyObject && strategy === kEmptyObject) {
264+
this[kState] = createReadableStreamState();
265+
return;
266+
}
255267
validateObject(source, 'source', kValidateObjectAllowObjects);
256268
validateObject(strategy, 'strategy', kValidateObjectAllowObjectsAndNull);
257269
this[kState] = createReadableStreamState();
@@ -301,8 +313,13 @@ class ReadableStream {
301313
// only default controllers were wired here; byte stream controllers
302314
// keep the previous no-op behavior.
303315
const controller = this[kState].controller;
316+
if (controller === undefined) {
317+
if (this[kState].state === 'readable')
318+
readableStreamError(this, error);
319+
return;
320+
}
304321
if (isReadableStreamDefaultController(controller))
305-
controller.error(error);
322+
readableStreamDefaultControllerError(controller, error);
306323
}
307324

308325
// Used by the internal stream interop (end-of-stream). Materialized
@@ -351,6 +368,7 @@ class ReadableStream {
351368
return PromiseReject(
352369
new ERR_INVALID_STATE.TypeError('ReadableStream is locked'));
353370
}
371+
ensureEmptyDefaultController(this);
354372
return readableStreamCancel(this, reason);
355373
}
356374

@@ -2580,6 +2598,7 @@ function setupReadableStreamBYOBReader(reader, stream) {
25802598
function setupReadableStreamDefaultReader(reader, stream) {
25812599
if (isReadableStreamLocked(stream))
25822600
throw new ERR_INVALID_STATE.TypeError('ReadableStream is locked');
2601+
ensureEmptyDefaultController(stream);
25832602
readableStreamReaderGenericInitialize(reader, stream);
25842603
reader[kState].readRequests = kEmptyQueue;
25852604
}
@@ -2729,7 +2748,8 @@ function readableStreamDefaultControllerPull(controller) {
27292748
// The pull algorithm may be a raw callback (a wrapped user source.pull
27302749
// returns its result uncoerced; a synchronous throw surfaces here) or an
27312750
// internal algorithm that always returns a promise; thenAlgorithmResult
2732-
// handles both.
2751+
// handles both. Non-thenable results react on kResolvedPromise so each
2752+
// pull is still separated by a microtask, matching the spec.
27332753
let result;
27342754
try {
27352755
result = controller[kState].pullAlgorithm(controller);
@@ -2763,7 +2783,7 @@ function readableStreamDefaultControllerCancelSteps(controller, reason) {
27632783
resetQueue(controller);
27642784
const result = controller[kState].cancelAlgorithm(reason);
27652785
readableStreamDefaultControllerClearAlgorithms(controller);
2766-
return result;
2786+
return promiseFromAlgorithmResult(result);
27672787
}
27682788

27692789
function readableStreamDefaultControllerPullSteps(controller, readRequest) {
@@ -2796,6 +2816,43 @@ function readableStreamDefaultControllerPullSteps(controller, readRequest) {
27962816
readableStreamDefaultControllerPull(controller);
27972817
}
27982818

2819+
// Materialize the deferred default controller for `new ReadableStream()`.
2820+
//
2821+
// started is true immediately: the empty-argument start algorithm is a
2822+
// no-op, so there is no initial pull and nothing can observe an unstarted
2823+
// controller without first calling getReader/cancel/pipeTo/tee/values,
2824+
// all of which come through here. That is also why this still matches
2825+
// WPT: those tests either pass a source (leaving this path) or wait for
2826+
// start, which is already complete for a no-op start.
2827+
//
2828+
// Subclasses that call cancel(), getReader(), pipeTo(), tee(), or
2829+
// values() in the constructor body after super() will materialize the
2830+
// controller before the subclass constructor finishes. Passing a source
2831+
// (for example to install start/pull) leaves the empty-argument path
2832+
// and creates the controller during super() as usual.
2833+
function ensureEmptyDefaultController(stream) {
2834+
if (stream[kState].controller !== undefined)
2835+
return stream[kState].controller;
2836+
const controller = new ReadableStreamDefaultController(kSkipThrow);
2837+
controller[kState] = {
2838+
cancelAlgorithm: nonOpCancel,
2839+
closeRequested: false,
2840+
highWaterMark: 1,
2841+
pullAgain: false,
2842+
pullAlgorithm: nonOpCallback,
2843+
pulling: false,
2844+
pullFulfilled: undefined,
2845+
pullRejected: undefined,
2846+
queue: kEmptyQueue,
2847+
queueTotalSize: 0,
2848+
started: true,
2849+
sizeAlgorithm: defaultSizeAlgorithm,
2850+
stream,
2851+
};
2852+
stream[kState].controller = controller;
2853+
return controller;
2854+
}
2855+
27992856
function setupReadableStreamDefaultController(
28002857
stream,
28012858
controller,
@@ -2824,8 +2881,7 @@ function setupReadableStreamDefaultController(
28242881

28252882
const startResult = startAlgorithm();
28262883

2827-
if (startResult === null ||
2828-
(typeof startResult !== 'object' && typeof startResult !== 'function')) {
2884+
if (isNonThenable(startResult)) {
28292885
// Non-thenable start result: fulfillment is guaranteed and no .then
28302886
// lookup on the result is observable, so run the post-start step
28312887
// directly at the exact microtask position the promise reaction
@@ -3586,7 +3642,7 @@ function readableByteStreamControllerCancelSteps(controller, reason) {
35863642
resetQueue(controller);
35873643
const result = controller[kState].cancelAlgorithm(reason);
35883644
readableByteStreamControllerClearAlgorithms(controller);
3589-
return result;
3645+
return promiseFromAlgorithmResult(result);
35903646
}
35913647

35923648
// Dequeues the first chunk of the byte queue as a Uint8Array view,
@@ -3708,8 +3764,7 @@ function setupReadableByteStreamController(
37083764

37093765
const startResult = startAlgorithm();
37103766

3711-
if (startResult === null ||
3712-
(typeof startResult !== 'object' && typeof startResult !== 'function')) {
3767+
if (isNonThenable(startResult)) {
37133768
// See setupReadableStreamDefaultController.
37143769
queueMicrotask(() => {
37153770
controller[kState].started = true;

‎lib/internal/webstreams/transformstream.js‎

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@ const {
5858
kType,
5959
nonOpCancel,
6060
nonOpFlush,
61+
delayedAlgorithmResult,
6162
} = require('internal/webstreams/util');
6263

6364
const {
@@ -127,9 +128,14 @@ class TransformStream {
127128
writableStrategy = kEmptyObject,
128129
readableStrategy = kEmptyObject) {
129130
markTransferMode(this, false, true);
130-
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
131-
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
132-
validateObject(readableStrategy, 'readableStrategy', kValidateObjectAllowObjectsAndNull);
131+
if (transformer !== kEmptyObject)
132+
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
133+
if (writableStrategy !== kEmptyObject) {
134+
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
135+
}
136+
if (readableStrategy !== kEmptyObject) {
137+
validateObject(readableStrategy, 'readableStrategy', kValidateObjectAllowObjectsAndNull);
138+
}
133139
const readableType = transformer?.readableType;
134140
const writableType = transformer?.writableType;
135141
const start = transformer?.start;
@@ -648,15 +654,16 @@ async function transformStreamDefaultSinkAbortAlgorithm(stream, reason) {
648654

649655
const { promise, resolve, reject } = PromiseWithResolvers();
650656
controller[kState].finishPromise = promise;
651-
const cancelPromise = controller[kState].cancelAlgorithm(reason);
657+
const cancelPromise =
658+
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
652659
transformStreamDefaultControllerClearAlgorithms(controller);
653660

654661
PromisePrototypeThen(
655662
cancelPromise,
656663
() => {
657-
if (readable[kState].state === 'errored')
664+
if (readable[kState].state === 'errored') {
658665
reject(readable[kState].storedError);
659-
else {
666+
} else {
660667
readableStreamDefaultControllerError(readable[kState].controller, reason);
661668
resolve();
662669
}
@@ -681,7 +688,8 @@ function transformStreamDefaultSinkCloseAlgorithm(stream) {
681688
}
682689
const { promise, resolve, reject } = PromiseWithResolvers();
683690
controller[kState].finishPromise = promise;
684-
const flushPromise = controller[kState].flushAlgorithm(controller);
691+
const flushPromise =
692+
delayedAlgorithmResult(controller[kState].flushAlgorithm(controller));
685693
transformStreamDefaultControllerClearAlgorithms(controller);
686694
PromisePrototypeThen(
687695
flushPromise,
@@ -724,15 +732,16 @@ function transformStreamDefaultSourceCancelAlgorithm(stream, reason) {
724732

725733
const { promise, resolve, reject } = PromiseWithResolvers();
726734
controller[kState].finishPromise = promise;
727-
const cancelPromise = controller[kState].cancelAlgorithm(reason);
735+
const cancelPromise =
736+
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
728737
transformStreamDefaultControllerClearAlgorithms(controller);
729738

730739
PromisePrototypeThen(
731740
cancelPromise,
732741
() => {
733-
if (writable[kState].state === 'errored')
742+
if (writable[kState].state === 'errored') {
734743
reject(writable[kState].storedError);
735-
else {
744+
} else {
736745
writableStreamDefaultControllerErrorIfNeeded(
737746
writable[kState].controller,
738747
reason);

0 commit comments

Comments
 (0)