Skip to content
Open
15 changes: 12 additions & 3 deletions doc/api/stream_iter.md
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,12 @@ Both forms receive an `options` parameter with the following property:
can check `signal.aborted` or listen for the `'abort'` event to perform
early cleanup.

In `pull()`, stateless transforms receive a new `options` object for every
call, and stateful transforms one for the pipeline, so a transform can modify
its `options` without affecting other transforms. The object does not inherit
from `Object.prototype`. Transforms passed to [`pullSync()`][] receive no
`options`.

The flush signal (`null`) is sent after the source ends, giving transforms
a chance to emit trailing data (e.g., compression footers).

Expand Down Expand Up @@ -714,7 +720,7 @@ added:

* `source` {Iterable} The sync data source.
* `...transforms` {Function|Object} Zero or more sync transforms.
* `writer` {Object} Destination with `write(chunk)` method.
* `writer` {Object} Destination with a `writeSync(chunk)` method.
* `options` {Object}
* `failOnIncompleteClose` {boolean} If `true`, call `writer.fail()` when
`writer.endSync()` cannot close the writer synchronously. Ignored when
Expand All @@ -727,8 +733,11 @@ added:
Synchronous version of [`pipeTo()`][]. The `source`, all transforms, and the
`writer` must be synchronous. Cannot accept async iterables or promises.

The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`,
`endSync`) and `fail()` for this to work.
The `writer` must have a `writeSync()` method. The other methods are
optional: `writevSync()` is used for batches of more than one chunk if it is
present, `endSync()` is called to close the writer (unless `preventClose` is
`true`), and `fail()` is called if the pipe fails (unless `preventFail` is
`true`). A writer without `endSync()` is not closed.

`pipeToSync()` never falls back to the asynchronous writer methods. If
`writer.endSync()` returns `-1` because the writer cannot close synchronously
Expand Down
35 changes: 19 additions & 16 deletions lib/internal/streams/iter/broadcast.js
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,11 @@ const {
} = require('internal/streams/iter/pull');

const {
IterResult,
PendingRequest,
PendingWrite,
kMultiConsumerDefaultBudget,
kNullOnceOption,
kResolvedPromise,
convertChunks,
createBatchEntry,
Expand Down Expand Up @@ -96,7 +100,7 @@ function raceEndWithSignal(promise, signal) {

const { promise: aborted, reject } = PromiseWithResolvers();
const onAbort = () => reject(signal.reason);
signal.addEventListener('abort', onAbort, { __proto__: null, once: true });
signal.addEventListener('abort', onAbort, kNullOnceOption);
if (signal.aborted) onAbort();

return SafePromisePrototypeFinally(
Expand Down Expand Up @@ -207,13 +211,13 @@ class BroadcastImpl {
const self = this;

const kDone = PromiseResolve(
{ __proto__: null, done: true, value: undefined });
new IterResult(true, undefined));

function detach() {
state.detached = true;
self.#waiters.delete(state);
if (state.resolve) {
state.resolve({ __proto__: null, done: true, value: undefined });
state.resolve(new IterResult(true, undefined));
}
self.#resolvePendingDone(state);
if (self.#deleteConsumer(state)) {
Expand All @@ -225,8 +229,7 @@ class BroadcastImpl {
return {
__proto__: null,
[SymbolAsyncIterator]() {
return {
__proto__: null,
return ObjectSetPrototypeOf({
next() {
if (state.detached) {
if (state.error !== kNoBroadcastError) {
Expand All @@ -246,7 +249,7 @@ class BroadcastImpl {
self.#tryTrimBuffer();
}
return PromiseResolve(
{ __proto__: null, done: false, value: chunk });
new IterResult(false, chunk));
}

if (self.#errored) {
Expand All @@ -264,7 +267,7 @@ class BroadcastImpl {
if (state.resolve) {
const { promise, resolve, reject } = PromiseWithResolvers();
ArrayPrototypePush(state.pending,
{ __proto__: null, resolve, reject });
new PendingRequest(resolve, reject));
return promise;
}

Expand All @@ -284,7 +287,7 @@ class BroadcastImpl {
detach();
return kDone;
},
};
}, null);
},
};
}
Expand All @@ -308,7 +311,7 @@ class BroadcastImpl {
if (hasReason) {
consumer.reject?.(reason);
} else {
consumer.resolve({ __proto__: null, done: true, value: undefined });
consumer.resolve(new IterResult(true, undefined));
}
consumer.resolve = null;
consumer.reject = null;
Expand Down Expand Up @@ -398,9 +401,9 @@ class BroadcastImpl {
--this.#cachedMinCursorConsumers === 0) {
this.#tryTrimBuffer();
}
consumer.resolve({ __proto__: null, done: false, value: chunk });
consumer.resolve(new IterResult(false, chunk));
} else {
consumer.resolve({ __proto__: null, done: true, value: undefined });
consumer.resolve(new IterResult(true, undefined));
this.#resolvePendingDone(consumer);
consumer.detached = true;
}
Expand Down Expand Up @@ -530,7 +533,7 @@ class BroadcastImpl {
const resolve = consumer.resolve;
consumer.resolve = null;
consumer.reject = null;
resolve({ __proto__: null, done: false, value: chunk });
resolve(new IterResult(false, chunk));
if (consumer.detached && this.#deleteConsumer(consumer)) {
this.#tryTrimBuffer();
} else if (this.#promotePending(consumer)) {
Expand Down Expand Up @@ -574,7 +577,7 @@ class BroadcastImpl {
}
while (consumer.pending.length > 0) {
ArrayPrototypeShift(consumer.pending).resolve(
{ __proto__: null, done: true, value: undefined });
new IterResult(true, undefined));
}
}

Expand Down Expand Up @@ -625,7 +628,7 @@ class BroadcastWriter {
if (canWrite === null) return null;
if (canWrite) return PromiseResolve(true);
const { promise, resolve, reject } = PromiseWithResolvers();
ArrayPrototypePush(this.#pendingDrains, { __proto__: null, resolve, reject });
ArrayPrototypePush(this.#pendingDrains, new PendingRequest(resolve, reject));
return promise;
}

Expand Down Expand Up @@ -802,7 +805,7 @@ class BroadcastWriter {
*/
#createPendingWrite(batch, signal) {
const { promise, resolve, reject } = PromiseWithResolvers();
const entry = { __proto__: null, batch, resolve, reject };
const entry = new PendingWrite(batch, resolve, reject);
this.#pendingWrites.push(entry);
if (signal) {
wireBroadcastWriteSignal(entry, signal, resolve, reject, this);
Expand Down Expand Up @@ -871,7 +874,7 @@ function wireBroadcastWriteSignal(entry, signal, resolve, reject, self) {
entry.batch = null;
reject(reason);
};
signal.addEventListener('abort', onAbort, { __proto__: null, once: true });
signal.addEventListener('abort', onAbort, kNullOnceOption);
}

// =============================================================================
Expand Down
45 changes: 26 additions & 19 deletions lib/internal/streams/iter/classic.js
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
const {
ArrayPrototypePush,
FunctionPrototypeCall,
ObjectFreeze,
Promise,
PromisePrototypeThen,
PromiseReject,
Expand Down Expand Up @@ -61,6 +62,8 @@ const {
} = require('internal/streams/iter/types');

const {
PendingRequest,
kNullOnceOption,
convertChunks,
getWriterSignal,
onSignalAbort,
Expand Down Expand Up @@ -101,10 +104,7 @@ function raceWithSignal(promise, signal) {
reject,
} = PromiseWithResolvers();
const onAbort = () => reject(signal.reason);
signal.addEventListener('abort', onAbort, {
__proto__: null,
once: true,
});
signal.addEventListener('abort', onAbort, kNullOnceOption);
PromisePrototypeThen(
promise,
(value) => {
Expand Down Expand Up @@ -522,6 +522,18 @@ function toReadableSync(source, options = kNullPrototype) {
// Cache: one Writer adapter per Writable instance.
const fromWritableCache = new SafeWeakMap();

// A write queued by fromWritable() until the writable can accept it. Like
// the other per-write records, constructed rather than created as a
// `{ __proto__: null, ... }` literal (a dictionary-mode object).
function QueuedWrite(chunks, resolve, reject) {
this.chunks = chunks;
this.resolve = resolve;
this.reject = reject;
this.signal = undefined;
this.onAbort = undefined;
}
QueuedWrite.prototype = ObjectFreeze({ __proto__: null });

/**
* Create a stream/iter Writer adapter from a classic Writable (or duck-type).
*
Expand Down Expand Up @@ -607,7 +619,9 @@ function fromWritable(writable, options = kNullPrototype) {
}

function removeDrainListenerIfIdle() {
// Keep listening while needsDrain is set: it is only cleared by 'drain'.
if (!drainListenerInstalled ||
needsDrain ||
pendingWrites.length !== 0 ||
drainWaiters.length !== 0) {
return;
Expand Down Expand Up @@ -643,7 +657,7 @@ function fromWritable(writable, options = kNullPrototype) {
if (preserveReason) {
pending[i].reject(reason);
} else {
pending[i].close();
pending[i].resolve(false);
}
}

Expand Down Expand Up @@ -687,7 +701,12 @@ function fromWritable(writable, options = kNullPrototype) {
for (let i = 0; i < chunks.length; i++) {
const bytes = chunks[i];
if (!writable.write(bytes)) {
// Listen for 'drain' now, even if nothing is waiting yet: the
// Writable can emit it before the next write or wait (for example
// when its write callback runs on a microtask), and needsDrain
// would never be cleared.
needsDrain = true;
installDrainListener();
ok = false;
}
totalBytes += TypedArrayPrototypeGetByteLength(bytes);
Expand Down Expand Up @@ -755,14 +774,7 @@ function fromWritable(writable, options = kNullPrototype) {

function queueWrite(chunks, signal) {
const { promise, resolve, reject } = PromiseWithResolvers();
const entry = {
__proto__: null,
chunks,
resolve,
reject,
signal: undefined,
onAbort: undefined,
};
const entry = new QueuedWrite(chunks, resolve, reject);
pendingWrites.push(entry);
installDrainListener();

Expand Down Expand Up @@ -1038,12 +1050,7 @@ function fromWritable(writable, options = kNullPrototype) {
return PromiseResolve(true);
}
const { promise, resolve, reject } = PromiseWithResolvers();
ArrayPrototypePush(drainWaiters, {
__proto__: null,
resolve,
reject,
close() { resolve(false); },
});
ArrayPrototypePush(drainWaiters, new PendingRequest(resolve, reject));
installDrainListener();
return promise;
};
Expand Down
33 changes: 17 additions & 16 deletions lib/internal/streams/iter/consumers.js
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ const {
ArrayPrototypeShift,
ArrayPrototypeSlice,
FunctionPrototypeCall,
ObjectFreeze,
Promise,
PromisePrototypeThen,
SafePromiseAllReturnVoid,
Expand Down Expand Up @@ -53,6 +54,7 @@ const {
} = require('internal/streams/iter/from');

const {
kNullOnceOption,
concatBytes,
getProtocolMethod,
recordChunk,
Expand Down Expand Up @@ -456,6 +458,16 @@ function ondrain(drainable) {

const kNoMergeError = Symbol('kNoMergeError');

// An entry in merge()'s ready queue: a value from `iterator`, or, for a
// source that failed, `reason` with no iterator. Created for every merged
// chunk, so constructed rather than created as a dictionary-mode literal.
function MergeEntry(iterator, value, reason) {
this.iterator = iterator;
this.value = value;
this.reason = reason;
}
MergeEntry.prototype = ObjectFreeze({ __proto__: null });

/**
* Merge multiple async iterables by yielding values in temporal order.
* @param {...(AsyncIterable<Uint8Array[]>|object)} args
Expand Down Expand Up @@ -517,10 +529,7 @@ function merge(...args) {
waitResolve = null;
}
};
signal.addEventListener('abort', onAbort, {
__proto__: null,
once: true,
});
signal.addEventListener('abort', onAbort, kNullOnceOption);
}

// Called when a source's .next() settles. Pushes the result into
Expand All @@ -531,12 +540,8 @@ function merge(...args) {
if (result.done) {
activeCount--;
} else {
ArrayPrototypePush(ready, {
__proto__: null,
kind: 'value',
iterator,
value: result.value,
});
ArrayPrototypePush(ready,
new MergeEntry(iterator, result.value, undefined));
}
if (waitResolve) {
waitResolve();
Expand All @@ -547,11 +552,7 @@ function merge(...args) {
const onRejected = (iterator, reason) => {
pendingPulls.delete(iterator);
if (stopped) return;
ArrayPrototypePush(ready, {
__proto__: null,
kind: 'error',
reason,
});
ArrayPrototypePush(ready, new MergeEntry(undefined, undefined, reason));
if (waitResolve) {
waitResolve();
waitResolve = null;
Expand Down Expand Up @@ -579,7 +580,7 @@ function merge(...args) {
// Drain ready queue synchronously
while (ready.length > 0) {
const item = ArrayPrototypeShift(ready);
if (item.kind === 'error') {
if (item.iterator === undefined) {
throw item.reason;
}
yield item.value;
Expand Down
4 changes: 2 additions & 2 deletions lib/internal/streams/iter/duplex.js
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ const {
const {
converters,
} = require('internal/streams/iter/webidl');
const { kNullOnceOption } = require('internal/streams/iter/utils');

/**
* Create a pair of connected duplex channels for bidirectional communication.
Expand Down Expand Up @@ -66,8 +67,7 @@ function duplex(options = { __proto__: null }) {
if (signal.aborted) {
abortBoth();
} else {
signal.addEventListener('abort', abortBoth,
{ __proto__: null, once: true });
signal.addEventListener('abort', abortBoth, kNullOnceOption);
}
}

Expand Down
Loading
Loading