Skip to content

Commit cea5388

Browse files
anonrigcursoragent
authored andcommitted
stream: address webstreams review blockers
Stop sharing kResolvedPromise on writer.ready/closed and cancel(). Restore async wrappers for cancel/close/abort/flush/transform so thenable results keep the previous microtask count. Remove the write-queue drain loop so each write stays one microtask apart. Drop the native webstreams binding. isNonThenable and cloneAsUint8Array stay in JS so a detached buffer still throws TypeError. Initialize the deferred controller field and skip materializing it on cancel of a non-readable empty stream. Assisted-by: Grok Signed-off-by: Yagiz Nizipli <yagiz@nizipli.com>
1 parent 5999da3 commit cea5388

11 files changed

Lines changed: 147 additions & 265 deletions

File tree

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -114,7 +114,6 @@ const {
114114
isNonThenable,
115115
kEmptyQueue,
116116
kResolvedPromise,
117-
promiseFromAlgorithmResult,
118117
kState,
119118
kType,
120119
lazyTransfer,
@@ -314,8 +313,9 @@ class ReadableStream {
314313
// keep the previous no-op behavior.
315314
const controller = this[kState].controller;
316315
if (controller === undefined) {
317-
if (this[kState].state === 'readable')
316+
if (this[kState].state === 'readable') {
318317
readableStreamError(this, error);
318+
}
319319
return;
320320
}
321321
if (isReadableStreamDefaultController(controller))
@@ -368,7 +368,10 @@ class ReadableStream {
368368
return PromiseReject(
369369
new ERR_INVALID_STATE.TypeError('ReadableStream is locked'));
370370
}
371-
ensureEmptyDefaultController(this);
371+
// Only materialize the deferred empty controller when cancel will
372+
// actually run cancel steps. closed/errored streams return immediately.
373+
if (this[kState].state === 'readable')
374+
ensureEmptyDefaultController(this);
372375
return readableStreamCancel(this, reason);
373376
}
374377

@@ -1440,6 +1443,7 @@ function createReadableStreamState() {
14401443
return {
14411444
__proto__: null,
14421445
closedPromise: undefined,
1446+
controller: undefined,
14431447
disturbed: false,
14441448
reader: undefined,
14451449
state: 'readable',
@@ -2783,7 +2787,7 @@ function readableStreamDefaultControllerCancelSteps(controller, reason) {
27832787
resetQueue(controller);
27842788
const result = controller[kState].cancelAlgorithm(reason);
27852789
readableStreamDefaultControllerClearAlgorithms(controller);
2786-
return promiseFromAlgorithmResult(result);
2790+
return result;
27872791
}
27882792

27892793
function readableStreamDefaultControllerPullSteps(controller, readRequest) {
@@ -3642,7 +3646,7 @@ function readableByteStreamControllerCancelSteps(controller, reason) {
36423646
resetQueue(controller);
36433647
const result = controller[kState].cancelAlgorithm(reason);
36443648
readableByteStreamControllerClearAlgorithms(controller);
3645-
return promiseFromAlgorithmResult(result);
3649+
return result;
36463650
}
36473651

36483652
// Dequeues the first chunk of the byte queue as a Uint8Array view,

‎lib/internal/webstreams/transformstream.js‎

Lines changed: 10 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,6 @@ const {
5454
kType,
5555
nonOpCancel,
5656
nonOpFlush,
57-
delayedAlgorithmResult,
5857
} = require('internal/webstreams/util');
5958

6059
const {
@@ -124,8 +123,9 @@ class TransformStream {
124123
writableStrategy = kEmptyObject,
125124
readableStrategy = kEmptyObject) {
126125
markTransferMode(this, false, true);
127-
if (transformer !== kEmptyObject)
126+
if (transformer !== kEmptyObject) {
128127
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
128+
}
129129
if (writableStrategy !== kEmptyObject) {
130130
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
131131
}
@@ -354,7 +354,7 @@ const isTransformStream =
354354
const isTransformStreamDefaultController =
355355
isBrandCheck('TransformStreamDefaultController');
356356

357-
function defaultTransformAlgorithm(chunk, controller) {
357+
async function defaultTransformAlgorithm(chunk, controller) {
358358
transformStreamDefaultControllerEnqueue(controller, chunk);
359359
}
360360

@@ -595,16 +595,15 @@ async function transformStreamDefaultSinkAbortAlgorithm(stream, reason) {
595595

596596
const { promise, resolve, reject } = PromiseWithResolvers();
597597
controller[kState].finishPromise = promise;
598-
const cancelPromise =
599-
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
598+
const cancelPromise = controller[kState].cancelAlgorithm(reason);
600599
transformStreamDefaultControllerClearAlgorithms(controller);
601600

602601
PromisePrototypeThen(
603602
cancelPromise,
604603
() => {
605-
if (readable[kState].state === 'errored') {
604+
if (readable[kState].state === 'errored')
606605
reject(readable[kState].storedError);
607-
} else {
606+
else {
608607
readableStreamDefaultControllerError(readable[kState].controller, reason);
609608
resolve();
610609
}
@@ -629,8 +628,7 @@ function transformStreamDefaultSinkCloseAlgorithm(stream) {
629628
}
630629
const { promise, resolve, reject } = PromiseWithResolvers();
631630
controller[kState].finishPromise = promise;
632-
const flushPromise =
633-
delayedAlgorithmResult(controller[kState].flushAlgorithm(controller));
631+
const flushPromise = controller[kState].flushAlgorithm(controller);
634632
transformStreamDefaultControllerClearAlgorithms(controller);
635633
PromisePrototypeThen(
636634
flushPromise,
@@ -667,16 +665,15 @@ function transformStreamDefaultSourceCancelAlgorithm(stream, reason) {
667665

668666
const { promise, resolve, reject } = PromiseWithResolvers();
669667
controller[kState].finishPromise = promise;
670-
const cancelPromise =
671-
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
668+
const cancelPromise = controller[kState].cancelAlgorithm(reason);
672669
transformStreamDefaultControllerClearAlgorithms(controller);
673670

674671
PromisePrototypeThen(
675672
cancelPromise,
676673
() => {
677-
if (writable[kState].state === 'errored') {
674+
if (writable[kState].state === 'errored')
678675
reject(writable[kState].storedError);
679-
} else {
676+
else {
680677
writableStreamDefaultControllerErrorIfNeeded(
681678
writable[kState].controller,
682679
reason);

‎lib/internal/webstreams/util.js‎

Lines changed: 39 additions & 63 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ const {
44
Array,
55
ArrayBufferPrototypeGetByteLength,
66
ArrayBufferPrototypeGetDetached,
7+
ArrayBufferPrototypeSlice,
78
AsyncIteratorPrototype,
89
DataViewPrototypeGetBuffer,
910
DataViewPrototypeGetByteLength,
@@ -19,6 +20,7 @@ const {
1920
TypedArrayPrototypeGetBuffer,
2021
TypedArrayPrototypeGetByteLength,
2122
TypedArrayPrototypeGetByteOffset,
23+
Uint8Array,
2224
} = primordials;
2325

2426
const {
@@ -31,11 +33,6 @@ const {
3133
copyArrayBuffer,
3234
} = internalBinding('buffer');
3335

34-
const {
35-
isNonThenable,
36-
cloneAsUint8Array: nativeCloneAsUint8Array,
37-
} = internalBinding('webstreams');
38-
3936
const {
4037
inspect,
4138
} = require('util');
@@ -131,7 +128,22 @@ function ArrayBufferViewGetByteOffset(view) {
131128
}
132129

133130
function cloneAsUint8Array(view) {
134-
return nativeCloneAsUint8Array(view);
131+
const buffer = ArrayBufferViewGetBuffer(view);
132+
const byteOffset = ArrayBufferViewGetByteOffset(view);
133+
const byteLength = ArrayBufferViewGetByteLength(view);
134+
return new Uint8Array(
135+
ArrayBufferPrototypeSlice(buffer, byteOffset, byteOffset + byteLength),
136+
);
137+
}
138+
139+
// True when `value` cannot be a thenable: null, undefined, or a
140+
// non-object non-function primitive. Objects and functions are treated
141+
// as maybe-thenable without looking up `.then` (that lookup is
142+
// observable). Proxies of objects/functions take the maybe-thenable
143+
// path; a Proxy around a primitive is still an object.
144+
function isNonThenable(value) {
145+
return value === null ||
146+
(typeof value !== 'object' && typeof value !== 'function');
135147
}
136148

137149
function canCopyArrayBuffer(toBuffer, toIndex, fromBuffer, fromIndex, count) {
@@ -332,19 +344,13 @@ function enqueueValueWithSize(controller, value, size) {
332344
// arguments passed through to the user callback is observable and must be
333345
// preserved.
334346
//
335-
// These are intentionally not `async` functions and not `Promise.try`.
336-
// Both always allocate a Promise, even when the user callback is
337-
// synchronous and returns a non-thenable. Callers use `isNonThenable()`
338-
// (or `PromisePrototypeThen` for thenables) to settle the result.
347+
// Cold algorithms (cancel/close/abort/flush/transform) stay `async` so
348+
// a user thenable is adopted with the same microtask count as before.
349+
// Pull/write use the raw-callback contract instead (see
350+
// createRawCallback*) and route results through thenAlgorithmResult().
339351
function createPromiseCallbackNoParams(name, fn, thisArg) {
340352
validateFunction(fn, name);
341-
return () => {
342-
try {
343-
return FunctionPrototypeCall(fn, thisArg);
344-
} catch (error) {
345-
return PromiseReject(error);
346-
}
347-
};
353+
return async () => FunctionPrototypeCall(fn, thisArg);
348354
}
349355

350356
// Raw variants that skip the async wrapper's implicit result promise.
@@ -382,24 +388,12 @@ function thenAlgorithmResult(result, onFulfilled, onRejected) {
382388

383389
function createPromiseCallback1Param(name, fn, thisArg) {
384390
validateFunction(fn, name);
385-
return (arg) => {
386-
try {
387-
return FunctionPrototypeCall(fn, thisArg, arg);
388-
} catch (error) {
389-
return PromiseReject(error);
390-
}
391-
};
391+
return async (arg) => FunctionPrototypeCall(fn, thisArg, arg);
392392
}
393393

394394
function createPromiseCallback2Params(name, fn, thisArg) {
395395
validateFunction(fn, name);
396-
return (arg1, arg2) => {
397-
try {
398-
return FunctionPrototypeCall(fn, thisArg, arg1, arg2);
399-
} catch (error) {
400-
return PromiseReject(error);
401-
}
402-
};
396+
return async (arg1, arg2) => FunctionPrototypeCall(fn, thisArg, arg1, arg2);
403397
}
404398

405399
function isPromisePending(promise) {
@@ -408,31 +402,12 @@ function isPromisePending(promise) {
408402
return details?.[0] === kPending;
409403
}
410404

411-
// Convert a promise-returning algorithm's raw result into a Promise. A
412-
// value that cannot be a thenable (null, undefined, or a non-object
413-
// non-function primitive) becomes the shared resolved promise. Objects
414-
// and functions go through PromiseResolve so a `.then` lookup, if any,
415-
// stays observable.
416-
function promiseFromAlgorithmResult(result) {
417-
if (isNonThenable(result))
418-
return kResolvedPromise;
419-
return PromiseResolve(result);
420-
}
421-
422-
// Cancel/flush/abort only: insert an extra microtask so "upon fulfillment"
423-
// of an already-settled user promise runs after start-settlement reactions
424-
// queued during construction. Pull/write must not use this.
425-
function delayedAlgorithmResult(result) {
426-
if (isNonThenable(result))
427-
return kResolvedPromise;
428-
return PromisePrototypeThen(kResolvedPromise, () => result);
429-
}
430-
431405
// Shared shapes for lazily-materialized { promise, resolve, reject }
432-
// records whose settlement is already known.
406+
// records whose settlement is already known. Each call mints a fresh
407+
// promise so public slots (writer.ready / writer.closed) stay distinct.
433408
function resolvedRecord() {
434409
return {
435-
promise: kResolvedPromise,
410+
promise: PromiseResolve(),
436411
resolve: undefined,
437412
reject: undefined,
438413
};
@@ -456,13 +431,16 @@ function setPromiseHandled(promise) {
456431
PromisePrototypeThen(promise, undefined, () => {});
457432
}
458433

459-
// Shared no-op. Start/pull/write use the raw-callback contract (see
460-
// createRawCallback*): a non-thenable return takes the allocation-free
461-
// path in thenAlgorithmResult(). Cancel/flush/abort wrap the result
462-
// with promiseFromAlgorithmResult/delayedAlgorithmResult, so a sync
463-
// no-op is equivalent to the previous async empty functions.
434+
async function nonOpFlush() {}
435+
436+
// Shared non-op for the start/pull/write algorithm callbacks, which all
437+
// follow the raw-callback contract (see createRawCallback*): the
438+
// non-thenable return takes the allocation-free fast path in
439+
// thenAlgorithmResult().
464440
function nonOpCallback() {}
465441

442+
async function nonOpCancel() {}
443+
466444
let transfer;
467445
function lazyTransfer() {
468446
if (transfer === undefined)
@@ -478,7 +456,6 @@ module.exports = {
478456
Queue,
479457
canCopyArrayBuffer,
480458
cloneAsUint8Array,
481-
isNonThenable,
482459
copyArrayBuffer,
483460
createPromiseCallbackNoParams,
484461
createPromiseCallback1Param,
@@ -493,6 +470,7 @@ module.exports = {
493470
extractSizeAlgorithm,
494471
getNonWritablePropertyDescriptor,
495472
isBrandCheck,
473+
isNonThenable,
496474
isPromisePending,
497475
kEmptyQueue,
498476
kResolvedPromise,
@@ -501,10 +479,8 @@ module.exports = {
501479
lazyTransfer,
502480
materializeQueue,
503481
nonOpCallback,
504-
nonOpCancel: nonOpCallback,
505-
nonOpFlush: nonOpCallback,
506-
promiseFromAlgorithmResult,
507-
delayedAlgorithmResult,
482+
nonOpCancel,
483+
nonOpFlush,
508484
peekQueueValue,
509485
rejectedHandledRecord,
510486
resetQueue,

0 commit comments

Comments
 (0)