Skip to content
Merged
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
157 changes: 92 additions & 65 deletions lib/internal/webstreams/readablestream.js
Original file line number Diff line number Diff line change
Expand Up @@ -46,10 +46,6 @@ const {
DOMException,
} = internalBinding('messaging');

const {
markPromiseAsHandled,
} = internalBinding('util');

const {
isArrayBufferView,
isDataView,
Expand Down Expand Up @@ -139,7 +135,7 @@ const {
writableStreamCloseQueuedOrInFlight,
writableStreamDefaultWriterCloseWithErrorPropagation,
writableStreamDefaultWriterRelease,
writableStreamDefaultWriterWrite,
writableStreamDefaultWriterWriteWithRequest,
writerClosedPromise,
writerReadyPromise,
} = require('internal/webstreams/writablestream');
Expand Down Expand Up @@ -1530,8 +1526,38 @@ function readableStreamPipeTo(

const promise = PromiseWithResolvers();

const state = {
currentWrite: PromiseResolve(),
// One shared write request tracks every chunk written to the
// destination, instead of a { promise, resolve, reject } record per
// write. `stall` is armed by waitForPendingWrites() during shutdown;
// `failed`/`failure` latch a write that could not proceed.
const writeTracker = {
// Non-undefined: queue entries are discriminated from kNilRequest
// by `promise === undefined`.
promise: null,
pending: 0,
failed: false,
failure: undefined,
stall: null,
resolve() {
if (--this.pending === 0 && this.stall !== null) {
const stall = this.stall;
this.stall = null;
if (this.failed)
stall.reject(this.failure);
else
stall.resolve();
}
},
reject(error) {
this.pending--;
this.failed = true;
this.failure = error;
if (this.stall !== null) {
const stall = this.stall;
this.stall = null;
stall.reject(error);
}
},
};

// The error here can be undefined. The rejected arg
Expand All @@ -1548,11 +1574,14 @@ function readableStreamPipeTo(
promise.resolve();
}

async function waitForCurrentWrite() {
const write = state.currentWrite;
await write;
if (write !== state.currentWrite)
await waitForCurrentWrite();
function waitForPendingWrites() {
if (writeTracker.pending === 0) {
return writeTracker.failed ?
PromiseReject(writeTracker.failure) :
PromiseResolve();
}
writeTracker.stall = PromiseWithResolvers();
return writeTracker.stall.promise;
}

function shutdownWithAnAction(action, rejected, originalError) {
Expand All @@ -1561,7 +1590,7 @@ function readableStreamPipeTo(
if (dest[kState].state === 'writable' &&
!writableStreamCloseQueuedOrInFlight(dest)) {
PromisePrototypeThen(
waitForCurrentWrite(),
waitForPendingWrites(),
complete,
(error) => finalize(true, error));
return;
Expand All @@ -1582,7 +1611,7 @@ function readableStreamPipeTo(
if (dest[kState].state === 'writable' &&
!writableStreamCloseQueuedOrInFlight(dest)) {
PromisePrototypeThen(
waitForCurrentWrite(),
waitForPendingWrites(),
() => finalize(rejected, error),
(error) => finalize(true, error));
return;
Expand Down Expand Up @@ -1639,25 +1668,46 @@ function readableStreamPipeTo(
PromisePrototypeThen(promise, action, () => {});
}

async function step() {
if (shuttingDown) return true;
// The pump loop is callback-driven to avoid per-iteration promise
// allocations. At most one read is in flight at a time, so one read
// request and one forwarding function are reused for every chunk;
// the chunk travels through `pendingChunk`.
let pendingChunk;
let readRequest;

// Ready promise rejection is handled by the destination-errored
// watcher.
function ignoreReadyRejection() {}

function forwardChunk() {
const chunk = pendingChunk;
pendingChunk = undefined;
writableStreamDefaultWriterWriteWithRequest(writer, chunk, writeTracker);
pump();
}

function pump() {
if (shuttingDown) return;

if (dest[kState].backpressure) {
await writerReadyPromise(writer).promise;
if (shuttingDown) return true;
PromisePrototypeThen(
writerReadyPromise(writer).promise,
pump,
ignoreReadyRejection);
return;
}

const controller = source[kState].controller;

// Fast path: batch reads when data is buffered in a default controller.
// This avoids creating PipeToReadableStreamReadRequest objects and
// reduces promise allocation overhead.
// This avoids parking read requests and reduces promise allocation
// overhead.
if (source[kState].state === 'readable' &&
isReadableStreamDefaultController(controller) &&
controller[kState].queue.length > 0) {

while (controller[kState].queue.length > 0) {
if (shuttingDown) return true;
if (shuttingDown) return;

const chunk = dequeueValue(controller);

Expand All @@ -1668,8 +1718,7 @@ function readableStreamPipeTo(

// Write the chunk - we're already in a separate microtask from enqueue
// because we awaited the writer ready promise above.
state.currentWrite = writableStreamDefaultWriterWrite(writer, chunk);
markPromiseAsHandled(state.currentWrite);
writableStreamDefaultWriterWriteWithRequest(writer, chunk, writeTracker);

// Check backpressure after each write
if (dest[kState].backpressure) {
Expand All @@ -1686,24 +1735,29 @@ function readableStreamPipeTo(

// Check if stream closed during batch
if (source[kState].state === 'closed') {
return true;
return;
}

// Yield to microtask queue between batches to allow events/signals to fire
return false;
// Yield to microtask queue between batches to allow events/signals
// to fire
queueMicrotask(pump);
return;
}

// Slow path: use read request for async reads
const promise = PromiseWithResolvers();
// eslint-disable-next-line no-use-before-define
readableStreamDefaultReaderRead(reader, new PipeToReadableStreamReadRequest(writer, state, promise));

return promise.promise;
}

async function run() {
// Run until step resolves as true
while (!await step());
// Slow path: park a lazily materialized read request. Close and
// error are handled by the source watchers.
readRequest ??= {
[kChunk](chunk) {
// Per spec, pipeTo must queue a microtask for the write to avoid
// synchronous write during enqueue(). See WHATWG Streams spec
// "ReadableStreamPipeTo" step 15's "chunk steps".
pendingChunk = chunk;
queueMicrotask(forwardChunk);
},
[kClose]() {},
[kError]() {},
};
readableStreamDefaultReaderRead(reader, readRequest);
}

if (signal !== undefined) {
Expand All @@ -1715,7 +1769,7 @@ function readableStreamPipeTo(
disposable = addAbortListener(signal, abortAlgorithm);
}

setPromiseHandled(run());
pump();

watchErrored(source, readerClosedPromise(reader).promise, (error) => {
if (!preventAbort) {
Expand Down Expand Up @@ -1760,33 +1814,6 @@ function readableStreamPipeTo(
return promise.promise;
}

class PipeToReadableStreamReadRequest {
constructor(writer, state, promise) {
this.writer = writer;
this.state = state;
this.promise = promise;
}

[kChunk](chunk) {
// Per spec, pipeTo must queue a microtask for the write to avoid
// synchronous write during enqueue(). See WHATWG Streams spec
// "ReadableStreamPipeTo" step 15's "chunk steps".
queueMicrotask(() => {
this.state.currentWrite = writableStreamDefaultWriterWrite(this.writer, chunk);
markPromiseAsHandled(this.state.currentWrite);
this.promise.resolve(false);
});
}

[kClose]() {
this.promise.resolve(true);
}

[kError](error) {
this.promise.reject(error);
}
}

function readableStreamTee(stream, cloneForBranch2) {
if (isReadableByteStreamController(stream[kState].controller)) {
return readableByteStreamTee(stream);
Expand Down
56 changes: 56 additions & 0 deletions lib/internal/webstreams/writablestream.js
Original file line number Diff line number Diff line change
Expand Up @@ -1001,6 +1001,61 @@ function writableStreamDefaultWriterWrite(writer, chunk) {
return promise;
}

// Variant of writableStreamDefaultWriterWrite for pipeTo: the caller
// provides a shared request object instead of a per-write promise record.
// `pending` is incremented before the controller write, which can settle
// requests synchronously when the stream starts erroring; precondition
// failures are latched on `failed`/`failure`.
function writableStreamDefaultWriterWriteWithRequest(writer, chunk, request) {
const writerState = writer[kState];
const stream = writerState.stream;
assert(stream !== undefined);
const streamState = stream[kState];
const {
controller,
} = streamState;
const chunkSize = writableStreamDefaultControllerGetChunkSize(
controller,
chunk);
if (stream !== writerState.stream) {
request.failed = true;
request.failure =
new ERR_INVALID_STATE.TypeError('Mismatched WritableStreams');
return;
}
const {
state,
} = streamState;

if (state === 'errored') {
request.failed = true;
request.failure = streamState.storedError;
return;
}

if (streamState.closeQueuedOrInFlight || state === 'closed') {
request.failed = true;
request.failure =
new ERR_INVALID_STATE.TypeError('WritableStream is closed');
return;
}

if (state === 'erroring') {
request.failed = true;
request.failure = streamState.storedError;
return;
}

assert(state === 'writable');

let writeRequests = streamState.writeRequests;
if (writeRequests === kEmptyQueue)
writeRequests = streamState.writeRequests = new Queue();
writeRequests.push(request);
request.pending++;
writableStreamDefaultControllerWrite(controller, chunk, chunkSize);
}

function writableStreamDefaultWriterRelease(writer) {
const {
stream,
Expand Down Expand Up @@ -1376,6 +1431,7 @@ module.exports = {
writableStreamCloseQueuedOrInFlight,
writableStreamAddWriteRequest,
writableStreamDefaultWriterWrite,
writableStreamDefaultWriterWriteWithRequest,
writableStreamDefaultWriterRelease,
writableStreamDefaultWriterGetDesiredSize,
writableStreamDefaultWriterEnsureReadyPromiseRejected,
Expand Down
Loading
Loading