Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 65 additions & 6 deletions lib/internal/webstreams/readablestream.js
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ const {
extractSizeAlgorithm,
getNonWritablePropertyDescriptor,
isBrandCheck,
isNonThenable,
kEmptyQueue,
kResolvedPromise,
kState,
Expand Down Expand Up @@ -252,6 +253,16 @@ class ReadableStream {
*/
constructor(source = kEmptyObject, strategy = kEmptyObject) {
markTransferMode(this, false, true);
// Empty-argument `new ReadableStream()`: no source, no strategy, and
// no controller. Reads never deliver data, so skip those allocations
// until getReader/cancel/error first need a default controller.
// Subclasses that call those methods after super() materialize the
// controller in the subclass constructor; see
// ensureEmptyDefaultController.
if (source === kEmptyObject && strategy === kEmptyObject) {
this[kState] = createReadableStreamState();
return;
}
validateObject(source, 'source', kValidateObjectAllowObjects);
validateObject(strategy, 'strategy', kValidateObjectAllowObjectsAndNull);
this[kState] = createReadableStreamState();
Expand Down Expand Up @@ -301,8 +312,14 @@ class ReadableStream {
// only default controllers were wired here; byte stream controllers
// keep the previous no-op behavior.
const controller = this[kState].controller;
if (controller === undefined) {
if (this[kState].state === 'readable') {
readableStreamError(this, error);
}
return;
}
if (isReadableStreamDefaultController(controller))
controller.error(error);
readableStreamDefaultControllerError(controller, error);
}

// Used by the internal stream interop (end-of-stream). Materialized
Expand Down Expand Up @@ -351,6 +368,10 @@ class ReadableStream {
return PromiseReject(
new ERR_INVALID_STATE.TypeError('ReadableStream is locked'));
}
// Only materialize the deferred empty controller when cancel will
// actually run cancel steps. closed/errored streams return immediately.
if (this[kState].state === 'readable')
ensureEmptyDefaultController(this);
return readableStreamCancel(this, reason);
}

Expand Down Expand Up @@ -1422,6 +1443,7 @@ function createReadableStreamState() {
return {
__proto__: null,
closedPromise: undefined,
controller: undefined,
disturbed: false,
reader: undefined,
state: 'readable',
Expand Down Expand Up @@ -2580,6 +2602,7 @@ function setupReadableStreamBYOBReader(reader, stream) {
function setupReadableStreamDefaultReader(reader, stream) {
if (isReadableStreamLocked(stream))
throw new ERR_INVALID_STATE.TypeError('ReadableStream is locked');
ensureEmptyDefaultController(stream);
readableStreamReaderGenericInitialize(reader, stream);
reader[kState].readRequests = kEmptyQueue;
}
Expand Down Expand Up @@ -2729,7 +2752,8 @@ function readableStreamDefaultControllerPull(controller) {
// The pull algorithm may be a raw callback (a wrapped user source.pull
// returns its result uncoerced; a synchronous throw surfaces here) or an
// internal algorithm that always returns a promise; thenAlgorithmResult
// handles both.
// handles both. Non-thenable results react on kResolvedPromise so each
// pull is still separated by a microtask, matching the spec.
let result;
try {
result = controller[kState].pullAlgorithm(controller);
Expand Down Expand Up @@ -2796,6 +2820,43 @@ function readableStreamDefaultControllerPullSteps(controller, readRequest) {
readableStreamDefaultControllerPull(controller);
}

// Materialize the deferred default controller for `new ReadableStream()`.
//
// started is true immediately: the empty-argument start algorithm is a
// no-op, so there is no initial pull and nothing can observe an unstarted
// controller without first calling getReader/cancel/pipeTo/tee/values,
// all of which come through here. That is also why this still matches
// WPT: those tests either pass a source (leaving this path) or wait for
// start, which is already complete for a no-op start.
//
// Subclasses that call cancel(), getReader(), pipeTo(), tee(), or
// values() in the constructor body after super() will materialize the
// controller before the subclass constructor finishes. Passing a source
// (for example to install start/pull) leaves the empty-argument path
// and creates the controller during super() as usual.
function ensureEmptyDefaultController(stream) {
if (stream[kState].controller !== undefined)
return stream[kState].controller;
const controller = new ReadableStreamDefaultController(kSkipThrow);
controller[kState] = {
cancelAlgorithm: nonOpCancel,
closeRequested: false,
highWaterMark: 1,
pullAgain: false,
pullAlgorithm: nonOpCallback,
pulling: false,
pullFulfilled: undefined,
pullRejected: undefined,
queue: kEmptyQueue,
queueTotalSize: 0,
started: true,
sizeAlgorithm: defaultSizeAlgorithm,
stream,
};
stream[kState].controller = controller;
return controller;
}

function setupReadableStreamDefaultController(
stream,
controller,
Expand Down Expand Up @@ -2824,8 +2885,7 @@ function setupReadableStreamDefaultController(

const startResult = startAlgorithm();

if (startResult === null ||
(typeof startResult !== 'object' && typeof startResult !== 'function')) {
if (isNonThenable(startResult)) {
// Non-thenable start result: fulfillment is guaranteed and no .then
// lookup on the result is observable, so run the post-start step
// directly at the exact microtask position the promise reaction
Expand Down Expand Up @@ -3708,8 +3768,7 @@ function setupReadableByteStreamController(

const startResult = startAlgorithm();

if (startResult === null ||
(typeof startResult !== 'object' && typeof startResult !== 'function')) {
if (isNonThenable(startResult)) {
// See setupReadableStreamDefaultController.
queueMicrotask(() => {
controller[kState].started = true;
Expand Down
12 changes: 9 additions & 3 deletions lib/internal/webstreams/transformstream.js
Original file line number Diff line number Diff line change
Expand Up @@ -123,9 +123,15 @@ class TransformStream {
writableStrategy = kEmptyObject,
readableStrategy = kEmptyObject) {
markTransferMode(this, false, true);
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
validateObject(readableStrategy, 'readableStrategy', kValidateObjectAllowObjectsAndNull);
if (transformer !== kEmptyObject) {
validateObject(transformer, 'transformer', kValidateObjectAllowObjects);
}
if (writableStrategy !== kEmptyObject) {
validateObject(writableStrategy, 'writableStrategy', kValidateObjectAllowObjectsAndNull);
}
if (readableStrategy !== kEmptyObject) {
validateObject(readableStrategy, 'readableStrategy', kValidateObjectAllowObjectsAndNull);
}
const readableType = transformer?.readableType;
const writableType = transformer?.writableType;
const start = transformer?.start;
Expand Down
23 changes: 19 additions & 4 deletions lib/internal/webstreams/util.js
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,16 @@ function cloneAsUint8Array(view) {
);
}

// True when `value` cannot be a thenable: null, undefined, or a
// non-object non-function primitive. Objects and functions are treated
// as maybe-thenable without looking up `.then` (that lookup is
// observable). Proxies of objects/functions take the maybe-thenable
// path; a Proxy around a primitive is still an object.
function isNonThenable(value) {
return value === null ||
(typeof value !== 'object' && typeof value !== 'function');
}

function canCopyArrayBuffer(toBuffer, toIndex, fromBuffer, fromIndex, count) {
return toBuffer !== fromBuffer &&
!ArrayBufferPrototypeGetDetached(toBuffer) &&
Expand Down Expand Up @@ -333,6 +343,11 @@ function enqueueValueWithSize(controller, value, size) {
// each known call-site arity gets its own wrapper. The exact number of
// arguments passed through to the user callback is observable and must be
// preserved.
//
// Cold algorithms (cancel/close/abort/flush/transform) stay `async` so
// a user thenable is adopted with the same microtask count as before.
// Pull/write use the raw-callback contract instead (see
// createRawCallback*) and route results through thenAlgorithmResult().
function createPromiseCallbackNoParams(name, fn, thisArg) {
validateFunction(fn, name);
return async () => FunctionPrototypeCall(fn, thisArg);
Expand Down Expand Up @@ -364,8 +379,7 @@ const kResolvedPromise = PromiseResolve();
// matches the spec's "a promise resolved with" conversion (identity for
// native promises).
function thenAlgorithmResult(result, onFulfilled, onRejected) {
if (result === null ||
(typeof result !== 'object' && typeof result !== 'function')) {
if (isNonThenable(result)) {
PromisePrototypeThen(kResolvedPromise, onFulfilled);
} else {
PromisePrototypeThen(PromiseResolve(result), onFulfilled, onRejected);
Expand All @@ -389,7 +403,8 @@ function isPromisePending(promise) {
}

// Shared shapes for lazily-materialized { promise, resolve, reject }
// records whose settlement is already known.
// records whose settlement is already known. Each call mints a fresh
// promise so public slots (writer.ready / writer.closed) stay distinct.
function resolvedRecord() {
return {
promise: PromiseResolve(),
Expand Down Expand Up @@ -455,6 +470,7 @@ module.exports = {
extractSizeAlgorithm,
getNonWritablePropertyDescriptor,
isBrandCheck,
isNonThenable,
isPromisePending,
kEmptyQueue,
kResolvedPromise,
Expand All @@ -465,7 +481,6 @@ module.exports = {
nonOpCallback,
nonOpCancel,
nonOpFlush,

peekQueueValue,
rejectedHandledRecord,
resetQueue,
Expand Down
51 changes: 34 additions & 17 deletions lib/internal/webstreams/writablestream.js
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ const {
extractSizeAlgorithm,
getNonWritablePropertyDescriptor,
isBrandCheck,
isNonThenable,
isPromisePending,
kEmptyQueue,
kState,
Expand Down Expand Up @@ -183,6 +184,15 @@ class WritableStream {
*/
constructor(sink = kEmptyObject, strategy = kEmptyObject) {
markTransferMode(this, false, true);
if (sink === kEmptyObject && strategy === kEmptyObject) {
this[kState] = createWritableStreamState();
setupWritableStreamDefaultControllerFromSink(
this,
sink,
1,
defaultSizeAlgorithm);
return;
}
validateObject(sink, 'sink', kValidateObjectAllowObjects);
validateObject(strategy, 'strategy', kValidateObjectAllowObjectsAndNull);
const type = sink?.type;
Expand Down Expand Up @@ -532,7 +542,7 @@ class WritableStreamDefaultController {
get signal() {
if (!isWritableStreamDefaultController(this))
throw new ERR_INVALID_THIS('WritableStreamDefaultController');
return this[kState].abortController.signal;
return (this[kState].abortController ??= new AbortController()).signal;
}

/**
Expand Down Expand Up @@ -707,7 +717,9 @@ function writableStreamAbort(stream, reason) {
if (state === 'closed' || state === 'errored')
return PromiseResolve();

controller[kState].abortController.abort(reason);
// Materialize lazily so construction stays cheap, but abort() must
// still abort the same signal later observed via controller.signal.
(controller[kState].abortController ??= new AbortController()).abort(reason);

state = stream[kState].state;
if (state === 'closed' || state === 'errored')
Expand Down Expand Up @@ -1169,6 +1181,21 @@ function writableStreamDefaultControllerWrite(controller, chunk, chunkSize) {
writableStreamDefaultControllerAdvanceQueueIfNeeded(controller);
}

function writableStreamDefaultControllerCompleteWrite(controller) {
const stream = controller[kState].stream;
writableStreamFinishInFlightWrite(stream);
const streamState = stream[kState];
const {
state,
} = streamState;
assert(state === 'writable' || state === 'erroring');
dequeueValue(controller);
if (!streamState.closeQueuedOrInFlight &&
state === 'writable') {
writableStreamUpdateBackpressure(controller, streamState);
}
}

function writableStreamDefaultControllerProcessWrite(controller, chunk) {
const {
stream,
Expand All @@ -1181,17 +1208,7 @@ function writableStreamDefaultControllerProcessWrite(controller, chunk) {
// so they are created once on the first write and reused for every
// subsequent write instead of allocating two fresh closures per chunk.
controller[kState].writeFulfilled = () => {
writableStreamFinishInFlightWrite(stream);
const streamState = stream[kState];
const {
state,
} = streamState;
assert(state === 'writable' || state === 'erroring');
dequeueValue(controller);
if (!streamState.closeQueuedOrInFlight &&
state === 'writable') {
writableStreamUpdateBackpressure(controller, streamState);
}
writableStreamDefaultControllerCompleteWrite(controller);
writableStreamDefaultControllerAdvanceQueueIfNeeded(controller);
};
controller[kState].writeRejected = (error) => {
Expand All @@ -1204,7 +1221,8 @@ function writableStreamDefaultControllerProcessWrite(controller, chunk) {
// The write algorithm may be a raw callback (a wrapped user sink.write
// returns its result uncoerced; a synchronous throw surfaces here) or an
// internal algorithm that always returns a promise; thenAlgorithmResult
// handles both.
// handles both. Non-thenable results react on kResolvedPromise so each
// write completion is still separated by a microtask.
let result;
try {
result = writeAlgorithm(chunk, controller);
Expand Down Expand Up @@ -1373,7 +1391,7 @@ function setupWritableStreamDefaultController(
highWaterMark,
queue: kEmptyQueue,
queueTotalSize: 0,
abortController: new AbortController(),
abortController: undefined,
sizeAlgorithm,
started: false,
stream,
Expand All @@ -1387,8 +1405,7 @@ function setupWritableStreamDefaultController(

const startResult = startAlgorithm();

if (startResult === null ||
(typeof startResult !== 'object' && typeof startResult !== 'function')) {
if (isNonThenable(startResult)) {
// Non-thenable start result: fulfillment is guaranteed and no .then
// lookup on the result is observable, so run the post-start step
// directly at the exact microtask position the promise reaction
Expand Down
Loading