Skip to content

Commit 83a8457

Browse files
committed
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 7f0def5 commit 83a8457

11 files changed

Lines changed: 139 additions & 256 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: 9 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,6 @@ const {
5858
kType,
5959
nonOpCancel,
6060
nonOpFlush,
61-
delayedAlgorithmResult,
6261
} = require('internal/webstreams/util');
6362

6463
const {
@@ -128,8 +127,9 @@ class TransformStream {
128127
writableStrategy = kEmptyObject,
129128
readableStrategy = kEmptyObject) {
130129
markTransferMode(this, false, true);
131-
if (transformer !== kEmptyObject)
130+
if (transformer !== kEmptyObject) {
132131
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
132+
}
133133
if (writableStrategy !== kEmptyObject) {
134134
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
135135
}
@@ -654,16 +654,15 @@ async function transformStreamDefaultSinkAbortAlgorithm(stream, reason) {
654654

655655
const { promise, resolve, reject } = PromiseWithResolvers();
656656
controller[kState].finishPromise = promise;
657-
const cancelPromise =
658-
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
657+
const cancelPromise = controller[kState].cancelAlgorithm(reason);
659658
transformStreamDefaultControllerClearAlgorithms(controller);
660659

661660
PromisePrototypeThen(
662661
cancelPromise,
663662
() => {
664-
if (readable[kState].state === 'errored') {
663+
if (readable[kState].state === 'errored')
665664
reject(readable[kState].storedError);
666-
} else {
665+
else {
667666
readableStreamDefaultControllerError(readable[kState].controller, reason);
668667
resolve();
669668
}
@@ -688,8 +687,7 @@ function transformStreamDefaultSinkCloseAlgorithm(stream) {
688687
}
689688
const { promise, resolve, reject } = PromiseWithResolvers();
690689
controller[kState].finishPromise = promise;
691-
const flushPromise =
692-
delayedAlgorithmResult(controller[kState].flushAlgorithm(controller));
690+
const flushPromise = controller[kState].flushAlgorithm(controller);
693691
transformStreamDefaultControllerClearAlgorithms(controller);
694692
PromisePrototypeThen(
695693
flushPromise,
@@ -732,16 +730,15 @@ function transformStreamDefaultSourceCancelAlgorithm(stream, reason) {
732730

733731
const { promise, resolve, reject } = PromiseWithResolvers();
734732
controller[kState].finishPromise = promise;
735-
const cancelPromise =
736-
delayedAlgorithmResult(controller[kState].cancelAlgorithm(reason));
733+
const cancelPromise = controller[kState].cancelAlgorithm(reason);
737734
transformStreamDefaultControllerClearAlgorithms(controller);
738735

739736
PromisePrototypeThen(
740737
cancelPromise,
741738
() => {
742-
if (writable[kState].state === 'errored') {
739+
if (writable[kState].state === 'errored')
743740
reject(writable[kState].storedError);
744-
} else {
741+
else {
745742
writableStreamDefaultControllerErrorIfNeeded(
746743
writable[kState].controller,
747744
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.
@@ -390,24 +396,12 @@ function thenAlgorithmResult(result, onFulfilled, onRejected) {
390396

391397
function createPromiseCallback1Param(name, fn, thisArg) {
392398
validateFunction(fn, name);
393-
return (arg) => {
394-
try {
395-
return FunctionPrototypeCall(fn, thisArg, arg);
396-
} catch (error) {
397-
return PromiseReject(error);
398-
}
399-
};
399+
return async (arg) => FunctionPrototypeCall(fn, thisArg, arg);
400400
}
401401

402402
function createPromiseCallback2Params(name, fn, thisArg) {
403403
validateFunction(fn, name);
404-
return (arg1, arg2) => {
405-
try {
406-
return FunctionPrototypeCall(fn, thisArg, arg1, arg2);
407-
} catch (error) {
408-
return PromiseReject(error);
409-
}
410-
};
404+
return async (arg1, arg2) => FunctionPrototypeCall(fn, thisArg, arg1, arg2);
411405
}
412406

413407
function isPromisePending(promise) {
@@ -416,31 +410,12 @@ function isPromisePending(promise) {
416410
return details?.[0] === kPending;
417411
}
418412

419-
// Convert a promise-returning algorithm's raw result into a Promise. A
420-
// value that cannot be a thenable (null, undefined, or a non-object
421-
// non-function primitive) becomes the shared resolved promise. Objects
422-
// and functions go through PromiseResolve so a `.then` lookup, if any,
423-
// stays observable.
424-
function promiseFromAlgorithmResult(result) {
425-
if (isNonThenable(result))
426-
return kResolvedPromise;
427-
return PromiseResolve(result);
428-
}
429-
430-
// Cancel/flush/abort only: insert an extra microtask so "upon fulfillment"
431-
// of an already-settled user promise runs after start-settlement reactions
432-
// queued during construction. Pull/write must not use this.
433-
function delayedAlgorithmResult(result) {
434-
if (isNonThenable(result))
435-
return kResolvedPromise;
436-
return PromisePrototypeThen(kResolvedPromise, () => result);
437-
}
438-
439413
// Shared shapes for lazily-materialized { promise, resolve, reject }
440-
// records whose settlement is already known.
414+
// records whose settlement is already known. Each call mints a fresh
415+
// promise so public slots (writer.ready / writer.closed) stay distinct.
441416
function resolvedRecord() {
442417
return {
443-
promise: kResolvedPromise,
418+
promise: PromiseResolve(),
444419
resolve: undefined,
445420
reject: undefined,
446421
};
@@ -464,13 +439,16 @@ function setPromiseHandled(promise) {
464439
PromisePrototypeThen(promise, undefined, () => {});
465440
}
466441

467-
// Shared no-op. Start/pull/write use the raw-callback contract (see
468-
// createRawCallback*): a non-thenable return takes the allocation-free
469-
// path in thenAlgorithmResult(). Cancel/flush/abort wrap the result
470-
// with promiseFromAlgorithmResult/delayedAlgorithmResult, so a sync
471-
// no-op is equivalent to the previous async empty functions.
442+
async function nonOpFlush() {}
443+
444+
// Shared non-op for the start/pull/write algorithm callbacks, which all
445+
// follow the raw-callback contract (see createRawCallback*): the
446+
// non-thenable return takes the allocation-free fast path in
447+
// thenAlgorithmResult().
472448
function nonOpCallback() {}
473449

450+
async function nonOpCancel() {}
451+
474452
let transfer;
475453
function lazyTransfer() {
476454
if (transfer === undefined)
@@ -486,7 +464,6 @@ module.exports = {
486464
Queue,
487465
canCopyArrayBuffer,
488466
cloneAsUint8Array,
489-
isNonThenable,
490467
copyArrayBuffer,
491468
createPromiseCallbackNoParams,
492469
createPromiseCallback1Param,
@@ -501,6 +478,7 @@ module.exports = {
501478
extractSizeAlgorithm,
502479
getNonWritablePropertyDescriptor,
503480
isBrandCheck,
481+
isNonThenable,
504482
isPromisePending,
505483
kEmptyQueue,
506484
kParkedAlgorithmResult,
@@ -510,10 +488,8 @@ module.exports = {
510488
lazyTransfer,
511489
materializeQueue,
512490
nonOpCallback,
513-
nonOpCancel: nonOpCallback,
514-
nonOpFlush: nonOpCallback,
515-
promiseFromAlgorithmResult,
516-
delayedAlgorithmResult,
491+
nonOpCancel,
492+
nonOpFlush,
517493
peekQueueValue,
518494
rejectedHandledRecord,
519495
resetQueue,

0 commit comments

Comments
 (0)