diff --git a/doc/api/fs.md b/doc/api/fs.md index d6eb4208ef0a..7feb4026f5de 100644 --- a/doc/api/fs.md +++ b/doc/api/fs.md @@ -412,8 +412,9 @@ Return the file contents as an async iterable using the chunks (default 128 KB). If transforms are provided, they are applied via [`stream/iter pull()`][]. -The file handle is locked while the iterable is being consumed and unlocked -when iteration completes, an error occurs, or the consumer breaks. +The file handle is locked from the first read of the iterable, and unlocked +when iteration completes, an error occurs, or the consumer breaks. An iterable +that is never read does not lock the file handle. This function is only available when the `--experimental-stream-iter` flag is enabled. @@ -487,8 +488,8 @@ Synchronous counterpart of [`filehandle.pull()`][]. Returns a sync iterable that reads the file using synchronous I/O on the main thread. Reads are performed in `chunkSize`-byte chunks (default 128 KB). -The file handle is locked while the iterable is being consumed. Unlike the -async `pull()`, this method does not support `AbortSignal` since all +The file handle is locked from the first read of the iterable until iteration +ends, as with [`filehandle.pull()`][]. Unlike the async `pull()`, this method does not support `AbortSignal` since all operations are synchronous. This function is only available when the `--experimental-stream-iter` flag is @@ -1134,6 +1135,11 @@ The writer supports both `Symbol.asyncDispose` and `Symbol.dispose`: for it to complete. * `using w = fh.writer()` — calls `fail()` unconditionally. +Async writes (`write()` and `writev()`) that are started without awaiting the +previous one are performed one at a time, in the order they were called, so +they never overlap in the file. A queued write is not performed if the writer +fails, or its `signal` aborts, before its turn. + The `writeSync()` and `writevSync()` methods enable the try-sync fast path used by [`stream/iter pipeTo()`][]. When the reader's chunk size matches the writer's `chunkSize`, all writes in a `pipeTo()` pipeline complete diff --git a/doc/api/stream_iter.md b/doc/api/stream_iter.md index 44a6d848a218..3348fc79bdc8 100644 --- a/doc/api/stream_iter.md +++ b/doc/api/stream_iter.md @@ -18,6 +18,13 @@ functions or objects with a `transform` method. Data flows in **batches** ({Uint8Array\[]} per iteration) to amortize the cost of async operations. +The module implements the WinterTC [Iterable Streams API][] draft. The +classic stream interop functions ([`fromReadable()`][], [`fromWritable()`][], +[`toReadable()`][], [`toReadableSync()`][] and [`toWritable()`][]), +[`Broadcast.from()`][], [`Share.from()`][], [`SyncShare.fromSync()`][] and the +protocol symbols exported by `Stream` are Node.js extensions that are not part +of the draft. + ```mjs import { from, pull, text } from 'node:stream/iter'; import { compressGzip, decompressGzip } from 'node:zlib/iter'; @@ -407,6 +414,18 @@ converted to a `USVString` and then UTF-8 encoded. `writev()` and chunks. Writer option dictionaries treat `null` as an empty dictionary and ignore unknown members. +Arguments are converted before the write itself starts. If the conversion runs +user code (for example a `toString()` method, or the iterator of a `writev()` +argument) that writes to the same writer, those writes are ordered before the +write whose argument is being converted, and they count against the same +backpressure limits. + +After `end()` or `endSync()` has been called, and until all buffered data has +been consumed, the writer is _closing_. While closing, `canWrite` is `null`, +`write()` and `writev()` reject with a `TypeError`, `writeSync()` and +`writevSync()` return `false`, `endSync()` returns `-1`, and calling `end()` +again returns the same promise as the first call. + Each async method has a synchronous `*Sync` counterpart designed for a try-fallback pattern: attempt the fast synchronous path first, and fall back to the async version only when the synchronous call indicates it could not @@ -697,6 +716,10 @@ added: * `...transforms` {Function|Object} Zero or more sync transforms. * `writer` {Object} Destination with `write(chunk)` method. * `options` {Object} + * `failOnIncompleteClose` {boolean} If `true`, call `writer.fail()` when + `writer.endSync()` cannot close the writer synchronously. Ignored when + `preventFail` is `true`. This option is a Node.js extension. + **Default:** `false`. * `preventClose` {boolean} **Default:** `false`. * `preventFail` {boolean} **Default:** `false`. * Returns: {number} Total bytes written. @@ -707,6 +730,15 @@ Synchronous version of [`pipeTo()`][]. The `source`, all transforms, and the The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`, `endSync`) and `fail()` for this to work. +`pipeToSync()` never falls back to the asynchronous writer methods. If +`writer.endSync()` returns `-1` because the writer cannot close synchronously +(for example, a `push()` writer whose consumer has not read all of the data +yet), `pipeToSync()` throws `ERR_INVALID_STATE`. All of the data was accepted +by then, so by default the writer is not failed: it can still be closed, for +example with `await writer.end()`. If the writer cannot be closed any other +way (for example, it has no `end()` method), or the caller will not close it, +set `failOnIncompleteClose` to fail it with the thrown error instead. + ### `pull(source[, ...transforms][, options])` -* `callback` {Function} `(chunks) => void` Called with each batch and with - `null` when the source ends. +* `callback` {Function} `(chunks, options) => void` Called with each batch and + with `null` when the source ends. `options.signal` is the pipeline's + {AbortSignal}. * Returns: {Function} A stateless transform. Create a pass-through transform that observes batches without modifying them. @@ -1286,7 +1332,7 @@ added: --> * `options` {Object} - * `budget` {number} Buffer size in bytes. Must be >= 16384. + * `budget` {number} Buffer size in bytes. Must be a positive integer. **Default:** `65536`. * `backpressure` {string} `'strict'`, `'unbounded'`, `'drop-oldest'`, or `'drop-newest'`. **Default:** `'strict'`. @@ -1354,6 +1400,10 @@ run().catch(console.error); Cancel the broadcast. If `reason` is provided, all consumers reject with that exact reason. If it is omitted, consumers complete normally. +Cancelling also closes the paired writer: afterwards its `canWrite` is `null` +and `write()` rejects with a `TypeError`. This lets a [`Broadcast.from()`][] +pump stop pulling from its source. + #### `broadcast.consumerCount` * {number} @@ -1400,7 +1450,7 @@ added: * `source` {AsyncIterable} The source to share. * `options` {Object} - * `budget` {number} Buffer size in bytes. Must be >= 16384. + * `budget` {number} Buffer size in bytes. Must be a positive integer. **Default:** `65536`. * `backpressure` {string} `'strict'`, `'unbounded'`, `'drop-oldest'`, or `'drop-newest'`. **Default:** `'strict'`. @@ -1411,6 +1461,28 @@ Create a pull-model multi-consumer shared stream. Unlike `broadcast()`, the source is only read when a consumer pulls. Multiple consumers share a single buffer. +A consumer created with `share.pull()` starts reading at the oldest entry still +in the buffer. Entries are released once every consumer has read them. When +every consumer has detached, the buffered data is kept for consumers that +attach later, and the source is not closed. Call `share.cancel()` (or dispose +the share) to release the source once it is no longer needed. + +With `'strict'` backpressure, a consumer that needs to pull from the source +while the buffer is at or above `budget` is rejected with `ERR_OUT_OF_RANGE` +and detached; further reads from that consumer complete with `{ done: true }`. +Detaching keeps a consumer that is not retried (for example, one read with +`for await...of`, which does not call `return()` when a read rejects) from +holding buffered data and blocking the other consumers. + +With `'unbounded'`, such a consumer waits until the slowest consumer releases +budget. With `'drop-newest'`, the entry pulled from the source is discarded +and the consumer then waits in the same way, so in both cases a stalled +consumer also stalls the consumers that are ahead of it. Only `'drop-oldest'` +lets consumers that are ahead continue, by discarding the oldest buffered +entries that the slowest consumer has not read yet. A batch pulled from the +source that is larger than `budget` is split into smaller entries first, so +eviction keeps the newest chunks that fit within the budget. + ```mjs import { from, share, text } from 'node:stream/iter'; @@ -1506,20 +1578,18 @@ added: * `source` {Iterable} The sync source to share. * `options` {Object} - * `budget` {number} Must be >= 16384. + * `budget` {number} Must be a positive integer. **Default:** `65536`. - * `backpressure` {string} `'strict'`, `'drop-oldest'`, or `'drop-newest'`. + * `backpressure` {string} `'strict'` or `'drop-oldest'`. **Default:** `'strict'`. * Returns: {SyncShare} Synchronous version of [`share()`][]. -Because there is no way to wait in a synchronous context, `'unbounded'` is not -supported and throws `ERR_INVALID_ARG_VALUE`. With `'drop-newest'`, a consumer -that reaches the end of the buffer while the budget is exhausted discards a -single entry from the source and then returns `{ done: true }` without a -value; the consumer is not detached, so it can resume once the slowest -consumer advances and releases budget. +A synchronous consumer cannot wait for the slowest consumer to release budget, +and the slowest consumer cannot advance while another consumer's read is +running. `'unbounded'` and `'drop-newest'` are therefore not supported and +throw `ERR_INVALID_ARG_VALUE`. ### Class: `SyncShare` @@ -2243,12 +2313,18 @@ const stream = fromSync(new Greeting('world')); console.log(textSync(stream)); // 'hello world' ``` +[Iterable Streams API]: https://iter-streams.proposal.wintertc.org/ [`--experimental-stream-iter`]: cli.md#--experimental-stream-iter +[`Broadcast.from()`]: #broadcastfrominput-options +[`Share.from()`]: #static-method-sharefrominput-options +[`SyncShare.fromSync()`]: #static-method-syncsharefromsyncinput-options [`array()`]: #arraysource-options [`arrayBuffer()`]: #arraybuffersource-options [`bytes()`]: #bytessource-options [`from()`]: #frominput +[`fromReadable()`]: #fromreadablereadable [`fromSync()`]: #fromsyncinput +[`fromWritable()`]: #fromwritablewritable-options [`node:zlib/iter`]: zlib.md#iterable-compression [`ondrain()`]: #ondraindrainable [`pipeTo()`]: #pipetosource-transforms-writer-options @@ -2260,3 +2336,6 @@ console.log(textSync(stream)); // 'hello world' [`tap()`]: #tapcallback [`text()`]: #textsource-options [`toAsyncStreamable`]: #streamtoasyncstreamable +[`toReadable()`]: #toreadablesource-options +[`toReadableSync()`]: #toreadablesyncsource-options +[`toWritable()`]: #towritablewriter diff --git a/lib/internal/errors.js b/lib/internal/errors.js index 0e80f682817d..ddd07a98209d 100644 --- a/lib/internal/errors.js +++ b/lib/internal/errors.js @@ -1848,7 +1848,7 @@ E('ERR_STREAM_UNABLE_TO_PIPE', 'Cannot pipe to a closed or destroyed stream', Er E('ERR_STREAM_UNSHIFT_AFTER_END_EVENT', 'stream.unshift() after end event', Error); E('ERR_STREAM_WRAP', 'Stream has StringDecoder set or is in objectMode', Error); -E('ERR_STREAM_WRITE_AFTER_END', 'write after end', Error); +E('ERR_STREAM_WRITE_AFTER_END', 'write after end', Error, TypeError); E('ERR_SYNTHETIC', 'JavaScript Callstack', Error); E('ERR_SYSTEM_ERROR', 'A system error occurred', SystemError, HideStackFramesError); E('ERR_TEST_FAILURE', function(error, failureType) { diff --git a/lib/internal/fs/promises.js b/lib/internal/fs/promises.js index efa981c55e31..c21b1e5740fc 100644 --- a/lib/internal/fs/promises.js +++ b/lib/internal/fs/promises.js @@ -480,6 +480,18 @@ if (getOptionValue('--experimental-stream-iter')) { const kNullPrototo = { __proto__: null }; const kDefaultChunkSize = 131072; const kNone = -1; + + // Called when iteration of a pull()/pullSync() iterable actually starts. + function lockForIteration(handle, fd) { + if (handle[kFd] === kNone || handle[kFd] !== fd) + throw new ERR_INVALID_STATE('The FileHandle is closed'); + if (handle[kClosePromise]) + throw new ERR_INVALID_STATE('The FileHandle is closing'); + if (handle[kLocked]) + throw new ERR_INVALID_STATE('The FileHandle is locked'); + handle[kLocked] = true; + } + /** * Return the file contents as an AsyncIterable using the * new streams pull model. Optional transforms and options (including @@ -527,11 +539,13 @@ if (getOptionValue('--experimental-stream-iter')) { validateAbortSignal(signal, 'options.signal'); } - this[kLocked] = true; - const source = { __proto__: null, async *[SymbolAsyncIterator]() { + // The handle is locked only while the iterable is actually being + // consumed. Locking eagerly in pull() would leave the handle locked + // forever if the returned iterable were never iterated. + lockForIteration(handle, fd); handle[kRef](); try { if (signal) { @@ -641,10 +655,6 @@ if (getOptionValue('--experimental-stream-iter')) { validateInteger(readSize, 'options.chunkSize', 1); } - this[kLocked] = true; - - handle[kRef](); - function cleanup() { handle[kLocked] = false; handle[kUnref](); @@ -656,15 +666,24 @@ if (getOptionValue('--experimental-stream-iter')) { const source = { __proto__: null, [SymbolIterator]() { + let started = false; let done = false; return { __proto__: null, next() { - if (done || remaining === 0) { - if (!done) { - done = true; - cleanup(); - } + if (done) { + return { done: true, value: undefined }; + } + if (!started) { + // Lock lazily, on the first read, so that an iterable that is + // never consumed does not leave the handle locked forever. + lockForIteration(handle, fd); + handle[kRef](); + started = true; + } + if (remaining === 0) { + done = true; + cleanup(); return { done: true, value: undefined }; } const toRead = remaining > 0 ? @@ -692,7 +711,7 @@ if (getOptionValue('--experimental-stream-iter')) { return() { if (!done) { done = true; - cleanup(); + if (started) cleanup(); } return { done: true, value: undefined }; }, @@ -751,6 +770,11 @@ if (getOptionValue('--experimental-stream-iter')) { let asyncPending = 0; let released = false; const pendingWrites = new SafeSet(); + // Async writes are issued one at a time, in call order. Without this, + // overlapping write() calls would race on the file position: when no + // `start` is given every write targets the fd's current (shared) + // position, and partial writes are completed in follow-up syscalls. + let writeQueue = PromiseResolve(); validateBoolean(autoClose, 'options.autoClose'); @@ -866,6 +890,22 @@ if (getOptionValue('--experimental-stream-iter')) { }); } + function ignoreWriteQueueResult() {} + + // Run `start` once every previously queued write has settled. + function enqueueWrite(start, signal) { + const operation = PromisePrototypeThen(writeQueue, () => { + // A failed writer must not touch the file again, and a write whose + // signal aborted while it was queued must not be started. + if (errored) throw error; + signal?.throwIfAborted(); + return start(); + }); + writeQueue = PromisePrototypeThen(operation, ignoreWriteQueueResult, + ignoreWriteQueueResult); + return operation; + } + function trackOperation(operation) { const { promise, resolve, reject } = PromiseWithResolvers(); const pending = { __proto__: null, reject }; @@ -950,8 +990,9 @@ if (getOptionValue('--experimental-stream-iter')) { if (bytesRemaining > 0) bytesRemaining -= chunk.byteLength; const position = pos; if (pos >= 0) pos += chunk.byteLength; - return trackOperation( - writeAll(chunk, 0, chunk.byteLength, position, signal)); + return trackOperation(enqueueWrite( + () => writeAll(chunk, 0, chunk.byteLength, position, signal), + signal)); }, writev(chunks, options = kNullPrototo) { @@ -983,7 +1024,8 @@ if (getOptionValue('--experimental-stream-iter')) { if (bytesRemaining > 0) bytesRemaining -= totalSize; const position = pos; if (pos >= 0) pos += totalSize; - return trackOperation(writevAll(chunks, position, signal)); + return trackOperation(enqueueWrite( + () => writevAll(chunks, position, signal), signal)); }, writeSync(chunk) { diff --git a/lib/internal/quic/quic.js b/lib/internal/quic/quic.js index 4922ce562751..df3649010c74 100644 --- a/lib/internal/quic/quic.js +++ b/lib/internal/quic/quic.js @@ -2324,7 +2324,7 @@ class QuicStream { async function writeAsync(chunk, signal) { if (errored) throw error; if (closed || ending || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); @@ -2350,7 +2350,7 @@ class QuicStream { await waitForDrain(signal); if (errored) throw error; if (closed || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); if (writeConverted(chunks, token)) return; @@ -2394,7 +2394,7 @@ class QuicStream { async function writevAsync(chunks, signal) { if (errored) throw error; if (closed || ending || stream.#inner.state.writeEnded) { - throw new ERR_INVALID_STATE('Writer is closed'); + throw new ERR_INVALID_STATE.TypeError('Writer is closed'); } signal?.throwIfAborted(); diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index e32f8aa544d2..cdf25021bbc5 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -40,9 +40,6 @@ const { ERR_INVALID_STATE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { broadcastProtocol, @@ -72,6 +69,7 @@ const { toWriterUint8Array, validateBatchEntry, yieldAbortable, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -895,7 +893,7 @@ function broadcast(options = { __proto__: null }) { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); const opts = { __proto__: null, diff --git a/lib/internal/streams/iter/classic.js b/lib/internal/streams/iter/classic.js index 59eb5be5819b..2c1fa6407149 100644 --- a/lib/internal/streams/iter/classic.js +++ b/lib/internal/streams/iter/classic.js @@ -912,7 +912,7 @@ function fromWritable(writable, options = kNullPrototype) { syncWritableError(); if (errored) return PromiseReject(error); if (!isWritable()) { - return PromiseReject(new ERR_STREAM_WRITE_AFTER_END()); + return PromiseReject(new ERR_STREAM_WRITE_AFTER_END.TypeError()); } if (signal?.aborted) return PromiseReject(signal.reason); @@ -944,7 +944,7 @@ function fromWritable(writable, options = kNullPrototype) { syncWritableError(); if (errored) return PromiseReject(error); if (!isWritable()) { - return PromiseReject(new ERR_STREAM_WRITE_AFTER_END()); + return PromiseReject(new ERR_STREAM_WRITE_AFTER_END.TypeError()); } if (signal?.aborted) return PromiseReject(signal.reason); if (chunks.length === 0) return PromiseResolve(); diff --git a/lib/internal/streams/iter/consumers.js b/lib/internal/streams/iter/consumers.js index 9e9a8fde4ec0..b118d6e3e7f4 100644 --- a/lib/internal/streams/iter/consumers.js +++ b/lib/internal/streams/iter/consumers.js @@ -53,9 +53,9 @@ const { const { concatBytes, - createBatchEntry, getProtocolMethod, - validateBatchEntry, + recordChunk, + validateRecordedChunks, yieldAbortable, } = require('internal/streams/iter/utils'); @@ -94,17 +94,6 @@ function isMergeOptions(value) { // Shared chunk collection helpers // ============================================================================= -function flattenBatchEntries(entries) { - const chunks = []; - for (let i = 0; i < entries.length; i++) { - const batch = validateBatchEntry(entries[i]); - for (let j = 0; j < batch.length; j++) { - ArrayPrototypePush(chunks, batch[j]); - } - } - return chunks; -} - /** * Collect chunks from a sync source into an array. * @param {Iterable} source @@ -114,23 +103,20 @@ function flattenBatchEntries(entries) { function collectSync(source, limit) { // Normalize source via fromSync() - accepts strings, ArrayBuffers, protocols, etc. const normalized = fromSync(source); - const entries = []; + const chunks = []; + const checks = []; let totalBytes = 0; for (const batch of normalized) { - const entry = createBatchEntry(batch); - if (limit !== undefined) { - for (let i = 0; i < entry.views.length; i++) { - totalBytes += entry.views[i].byteLength; - if (totalBytes > limit) { - throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); - } + for (let i = 0; i < batch.length; i++) { + totalBytes += recordChunk(chunks, checks, batch[i]); + if (limit !== undefined && totalBytes > limit) { + throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); } } - ArrayPrototypePush(entries, entry); } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } /** @@ -145,14 +131,17 @@ async function collectAsync(source, signal, limit) { // Normalize source via from() - accepts strings, ArrayBuffers, protocols, etc. const normalized = from(source); - const entries = []; + const chunks = []; + const checks = []; // Fast path: no signal and no limit if (!signal && limit === undefined) { for await (const batch of normalized) { - ArrayPrototypePush(entries, createBatchEntry(batch)); + for (let i = 0; i < batch.length; i++) { + recordChunk(chunks, checks, batch[i]); + } } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } // Slow path: with signal or limit checks @@ -161,19 +150,15 @@ async function collectAsync(source, signal, limit) { for await (const batch of iterable) { signal?.throwIfAborted(); - const entry = createBatchEntry(batch); - if (limit !== undefined) { - for (let i = 0; i < entry.views.length; i++) { - totalBytes += entry.views[i].byteLength; - if (totalBytes > limit) { - throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); - } + for (let i = 0; i < batch.length; i++) { + totalBytes += recordChunk(chunks, checks, batch[i]); + if (limit !== undefined && totalBytes > limit) { + throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes); } } - ArrayPrototypePush(entries, entry); } - return flattenBatchEntries(entries); + return validateRecordedChunks(chunks, checks); } /** diff --git a/lib/internal/streams/iter/from.js b/lib/internal/streams/iter/from.js index cac8fd3565b6..e434ff9c253b 100644 --- a/lib/internal/streams/iter/from.js +++ b/lib/internal/streams/iter/from.js @@ -61,6 +61,11 @@ const { // allocate output for the entire batch at once. const FROM_BATCH_SIZE = 128; const kNormalizationCancelled = Symbol('kNormalizationCancelled'); +// Yielded by normalizeAsyncValue() (only when `emitFlush` is true) right +// before it waits on a promise or on a nested async iterable. Callers that +// batch chunks yield whatever they have collected so far, so that chunks that +// are already available are not held back until the wait completes. +const kFlushBatch = Symbol('kFlushBatch'); function createNormalizationContext() { return { @@ -444,14 +449,15 @@ function yieldNormalizationAbortable(source, context) { * @yields {Uint8Array} */ async function* normalizeAsyncValue( - value, allowNestedAsyncStreamables = true, context) { + value, allowNestedAsyncStreamables = true, context, emitFlush = false) { throwIfNormalizationCancelled(context); // Handle promises first if (isPromise(value)) { + if (emitFlush) yield kFlushBatch; const resolved = await waitForNormalization(value, context); yield* normalizeAsyncValue( - resolved, allowNestedAsyncStreamables, context); + resolved, allowNestedAsyncStreamables, context, emitFlush); return; } @@ -478,13 +484,15 @@ async function* normalizeAsyncValue( if (asyncStreamableMethod !== undefined) { const result = FunctionPrototypeCall(asyncStreamableMethod, value); if (isPromise(result)) { + if (emitFlush) yield kFlushBatch; yield* normalizeAsyncValue( await waitForNormalization(result, context), allowNestedAsyncStreamables, - context); + context, + emitFlush); } else { yield* normalizeAsyncValue( - result, allowNestedAsyncStreamables, context); + result, allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -493,7 +501,8 @@ async function* normalizeAsyncValue( const streamableMethod = getProtocolMethod(value, toStreamable); if (streamableMethod !== undefined) { const result = FunctionPrototypeCall(streamableMethod, value); - yield* normalizeAsyncValue(result, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + result, allowNestedAsyncStreamables, context, emitFlush); return; } @@ -501,7 +510,7 @@ async function* normalizeAsyncValue( if (ArrayIsArray(value)) { for (let i = 0; i < value.length; i++) { yield* normalizeAsyncValue( - value[i], allowNestedAsyncStreamables, context); + value[i], allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -510,8 +519,11 @@ async function* normalizeAsyncValue( // have both) if (isAsyncIterable(value)) { const iterable = yieldNormalizationAbortable(value, context); + if (emitFlush) yield kFlushBatch; for await (const item of iterable) { - yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + item, allowNestedAsyncStreamables, context, emitFlush); + if (emitFlush) yield kFlushBatch; } return; } @@ -519,7 +531,8 @@ async function* normalizeAsyncValue( // Handle sync iterables if (isSyncIterable(value)) { for (const item of value) { - yield* normalizeAsyncValue(item, allowNestedAsyncStreamables, context); + yield* normalizeAsyncValue( + item, allowNestedAsyncStreamables, context, emitFlush); } return; } @@ -547,9 +560,7 @@ async function* normalizeAsyncSource(source, context) { for await (const value of iterable) { // Fast path 1: value is already a Uint8Array[] batch if (isUint8ArrayBatch(value)) { - if (value.length > 0) { - yield value; - } + yield* yieldBoundedBatch(value); continue; } // Fast path 2: value is a single Uint8Array (very common) @@ -557,9 +568,19 @@ async function* normalizeAsyncSource(source, context) { yield [value]; continue; } - // Slow path: normalize the value + // Slow path: normalize the value. Chunks are batched, but whatever has + // been collected is yielded before waiting on async content (such as a + // nested async iterable), so data is not held back until it ends. let batch = []; - for await (const chunk of normalizeAsyncValue(value, true, context)) { + for await (const chunk of + normalizeAsyncValue(value, true, context, true)) { + if (chunk === kFlushBatch) { + if (batch.length > 0) { + yield batch; + batch = []; + } + continue; + } ArrayPrototypePush(batch, chunk); if (batch.length === FROM_BATCH_SIZE) { yield batch; @@ -604,7 +625,15 @@ async function* normalizeAsyncSource(source, context) { batch = []; } let asyncBatch = []; - for await (const chunk of normalizeAsyncValue(value, false, context)) { + for await (const chunk of + normalizeAsyncValue(value, false, context, true)) { + if (chunk === kFlushBatch) { + if (asyncBatch.length > 0) { + yield asyncBatch; + asyncBatch = []; + } + continue; + } ArrayPrototypePush(asyncBatch, chunk); if (asyncBatch.length === FROM_BATCH_SIZE) { yield asyncBatch; @@ -709,7 +738,7 @@ function fromSync(input) { if (!isIterable && isAsyncIterable(input)) { throw new ERR_INVALID_ARG_TYPE( 'input', - 'a synchronous input (not AsyncIterable)', + 'a synchronous input, not an async iterable', input, ); } @@ -719,7 +748,7 @@ function fromSync(input) { typeof input.then === 'function') { throw new ERR_INVALID_ARG_TYPE( 'input', - 'a synchronous input (not Promise)', + 'a synchronous input, not a promise', input, ); } diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js index 70354edb10b3..176d8dbfc83a 100644 --- a/lib/internal/streams/iter/pull.js +++ b/lib/internal/streams/iter/pull.js @@ -12,6 +12,7 @@ const { ArrayPrototypeSlice, FunctionPrototypeCall, PromisePrototypeThen, + PromiseReject, PromiseResolve, SymbolAsyncIterator, SymbolIterator, @@ -850,37 +851,83 @@ function pull(source, ...args) { }); const { signal } = options; const normalized = from(source); - signal?.throwIfAborted(); return { __proto__: null, [SymbolAsyncIterator]() { - const controller = new AbortController(); - const iteratorSignal = signal === undefined ? - controller.signal : AbortSignal.any([signal, controller.signal]); - - async function* pipeline() { - yield* createAsyncPipeline(normalized, transforms, iteratorSignal); + if (signal === undefined) { + const controller = new AbortController(); + async function* pipeline() { + yield* createAsyncPipeline(normalized, transforms, controller.signal); + } + const iterator = pipeline(); + return { + __proto__: null, + next(value) { + return iterator.next(value); + }, + return(value) { + controller.abort(lazyDOMException('Aborted', 'AbortError')); + return iterator.return(value); + }, + throw(error) { + abortSignal(controller.signal, error); + return iterator.throw(error); + }, + [SymbolAsyncIterator]() { + return this; + }, + }; } - const iterator = pipeline(); + return createAbortablePullIterator(normalized, transforms, signal); + }, + }; +} - return { - __proto__: null, - next(value) { - return iterator.next(value); - }, - return(value) { - controller.abort(lazyDOMException('Aborted', 'AbortError')); - return iterator.return(value); - }, - throw(error) { - abortSignal(controller.signal, error); - return iterator.throw(error); - }, - [SymbolAsyncIterator]() { - return this; - }, - }; +// Once `signal` aborts the pipeline, the pull that observed it rejects and +// so does every later pull, with the abort reason. That includes the case of +// an already-aborted signal, where the pipeline is never started. A plain +// async generator would instead complete after throwing, so later pulls +// would report a clean end of the stream. +function createAbortablePullIterator(source, transforms, signal) { + let aborted = signal.aborted; + let controller; + let iterator; + if (!aborted) { + controller = new AbortController(); + const iteratorSignal = AbortSignal.any([signal, controller.signal]); + async function* pipeline() { + yield* createAsyncPipeline(source, transforms, iteratorSignal); + } + iterator = pipeline(); + } + + function onRejected(error) { + if (signal.aborted) aborted = true; + throw error; + } + + return { + __proto__: null, + next(value) { + if (aborted) return PromiseReject(signal.reason); + return PromisePrototypeThen(iterator.next(value), undefined, onRejected); + }, + return(value) { + if (aborted) { + return PromiseResolve({ __proto__: null, done: true, value }); + } + controller.abort(lazyDOMException('Aborted', 'AbortError')); + return iterator.return(value); + }, + throw(error) { + if (aborted) return PromiseReject(error); + abortSignal(controller.signal, error); + return PromisePrototypeThen(iterator.throw(error), undefined, + onRejected); + }, + [SymbolAsyncIterator]() { + return this; }, }; } @@ -1001,6 +1048,7 @@ function pipeToSync(source, ...args) { normalized; let totalBytes = 0; + let closedSync = true; try { for (const batch of pipeline) { @@ -1029,14 +1077,25 @@ function pipeToSync(source, ...args) { } if (!options.preventClose) { - if (FunctionPrototypeCall(endSync, writer) < 0) { - throw new ERR_INVALID_STATE( - 'Writer could not be closed synchronously'); - } + closedSync = FunctionPrototypeCall(endSync, writer) >= 0; } } catch (error) { if (!options.preventFail) { - writer.fail?.(error); + failWriterQuietly(writer, error); + } + throw error; + } + + // endSync() returning -1 only means that the writer cannot close + // synchronously; every chunk was accepted. pipeToSync() never falls back to + // the async end(), so report it. By default the writer is left as it is so + // that the caller can still close it (e.g. with `await writer.end()`); + // `failOnIncompleteClose` fails it instead, for callers that cannot. + if (!closedSync) { + const error = new ERR_INVALID_STATE( + 'Writer could not be closed synchronously'); + if (options.failOnIncompleteClose && !options.preventFail) { + failWriterQuietly(writer, error); } throw error; } @@ -1044,6 +1103,17 @@ function pipeToSync(source, ...args) { return totalBytes; } +// Call writer.fail(error) on a best-effort basis. The error that made the +// pipe fail is what the caller must see; an exception from fail() itself +// must not replace it. +function failWriterQuietly(writer, error) { + try { + writer.fail?.(error); + } catch { + // Ignored, see above. + } +} + /** * Write an async source through transforms to a writer. * @param {AsyncIterable|Iterable} source @@ -1061,7 +1131,7 @@ async function pipeTo(source, ...args) { function failWriter(error) { if (!options.preventFail) { - writer.fail?.(error); + failWriterQuietly(writer, error); } } diff --git a/lib/internal/streams/iter/push.js b/lib/internal/streams/iter/push.js index 31b437c85b75..764a3afac92a 100644 --- a/lib/internal/streams/iter/push.js +++ b/lib/internal/streams/iter/push.js @@ -22,9 +22,6 @@ const { ERR_INVALID_STATE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { drainableProtocol, @@ -40,6 +37,7 @@ const { getWriterSignal, parsePullArgs, validateBatchEntry, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -132,7 +130,7 @@ class PushQueue { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); this.#budget = budget; this.#backpressure = backpressure; this.#signal = signal; diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index b53da4608b23..6e9881b682a4 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -44,7 +44,9 @@ const { getMinCursor, onSignalAbort, parsePullArgs, + splitBatchEntry, validateBatchEntry, + validateBudget, } = require('internal/streams/iter/utils'); const { converters, @@ -59,12 +61,10 @@ const { ERR_INVALID_ARG_TYPE, ERR_INVALID_ARG_VALUE, ERR_INVALID_RETURN_VALUE, + ERR_INVALID_STATE, ERR_OUT_OF_RANGE, }, } = require('internal/errors'); -const { - validateInteger, -} = require('internal/validators'); const { markPromiseAsHandled } = internalBinding('util'); @@ -427,9 +427,7 @@ class ShareImpl { if (result.done) { this.#sourceExhausted = true; } else if (!discard) { - const entry = createBatchEntry(result.value); - this.#buffer.push(entry); - this.#bufferedBytes += entry.byteLength; + this.#bufferBatch(result.value); } } catch (error) { this.#sourceError = error; @@ -447,7 +445,31 @@ class ShareImpl { })(); } + #bufferBatch(batch) { + const entry = createBatchEntry(batch); + // 'drop-oldest' evicts whole entries. A single pulled batch can be much + // larger than the budget (for example when from() combines many values + // of a sync source), and evicting it would discard every chunk in it, + // including ones that no consumer has read. Split such batches so that + // eviction keeps the newest chunks that fit the budget. + const entries = this.#options.backpressure === 'drop-oldest' ? + splitBatchEntry(entry, this.#options.budget) : undefined; + if (entries === undefined) { + this.#buffer.push(entry); + } else { + for (let i = 0; i < entries.length; i++) { + this.#buffer.push(entries[i]); + } + } + this.#bufferedBytes += entry.byteLength; + } + #tryTrimBuffer() { + // Retain buffered data for consumers that attach while none are active. + // Without this, the last consumer detaching would discard entries that + // a late-joining consumer has not read and cannot get back from the + // source. + if (this.#consumers.size === 0) return; if (this.#cachedMinCursorConsumers === 0) { this.#recomputeMinCursor(); } @@ -525,6 +547,7 @@ class SyncShareImpl { #sourceExhausted = false; #sourceError = kNoShareError; #cancelled = false; + #pulling = false; #cachedMinCursor = 0; #cachedMinCursorConsumers = 0; /** Cumulative byte size of buffered entries */ @@ -606,14 +629,24 @@ class SyncShareImpl { return { __proto__: null, done: true, value: undefined }; } - // Check buffer limit - let dropped = false; + // Check buffer limit. 'unbounded' and 'drop-newest' are rejected + // by shareSync(). if (self.#bufferedBytes >= self.#options.budget) { switch (self.#options.backpressure) { - case 'strict': - throw new ERR_OUT_OF_RANGE( + case 'strict': { + const error = new ERR_OUT_OF_RANGE( 'buffered bytes', `< ${self.#options.budget}`, self.#bufferedBytes); + // Detach before throwing, as the async share does. Neither + // for...of nor a transform pipeline calls return() when + // next() throws, so a consumer left registered here would + // keep its cursor forever and wedge the other consumers. + state.detached = true; + if (self.#deleteConsumer(state)) { + self.#tryTrimBuffer(); + } + throw error; + } case 'drop-oldest': while (self.#bufferedBytes >= self.#options.budget && self.#buffer.length > 0) { @@ -629,21 +662,10 @@ class SyncShareImpl { } self.#recomputeMinCursor(); break; - case 'drop-newest': - // Discarding does not reclaim budget, and the slowest - // consumer cannot advance while this synchronous next() is - // running, so at most one entry may be dropped per call. - // Looping here would spin forever on an unbounded source - // and drain a finite one in a single call. - self.#pullFromSource(true); - dropped = true; - break; } } - if (!dropped) { - self.#pullFromSource(); - } + self.#pullFromSource(); if (self.#sourceError !== kNoShareError) { state.detached = true; @@ -722,9 +744,18 @@ class SyncShareImpl { this.cancel(); } - #pullFromSource(discard = false) { + #pullFromSource() { if (this.#sourceExhausted || this.#cancelled) return; + // A source that reads from its own share would re-enter its own + // iterator. Fail that read without touching the share's state; the + // source sees the error and may handle it. + if (this.#pulling) { + throw new ERR_INVALID_STATE( + 'shareSync() source cannot be read while it is producing a value'); + } + + this.#pulling = true; try { this.#sourceIterator ||= this.#source[SymbolIterator](); @@ -732,18 +763,42 @@ class SyncShareImpl { if (result.done) { this.#sourceExhausted = true; - } else if (!discard) { - const entry = createBatchEntry(result.value); - this.#buffer.push(entry); - this.#bufferedBytes += entry.byteLength; + } else { + this.#bufferBatch(result.value); } } catch (error) { this.#sourceError = error; this.#sourceExhausted = true; + } finally { + this.#pulling = false; + } + } + + #bufferBatch(batch) { + const entry = createBatchEntry(batch); + // 'drop-oldest' evicts whole entries. A single pulled batch can be much + // larger than the budget (for example when from() combines many values + // of a sync source), and evicting it would discard every chunk in it, + // including ones that no consumer has read. Split such batches so that + // eviction keeps the newest chunks that fit the budget. + const entries = this.#options.backpressure === 'drop-oldest' ? + splitBatchEntry(entry, this.#options.budget) : undefined; + if (entries === undefined) { + this.#buffer.push(entry); + } else { + for (let i = 0; i < entries.length; i++) { + this.#buffer.push(entries[i]); + } } + this.#bufferedBytes += entry.byteLength; } #tryTrimBuffer() { + // Retain buffered data for consumers that attach while none are active. + // Without this, the last consumer detaching would discard entries that + // a late-joining consumer has not read and cannot get back from the + // source. + if (this.#consumers.size === 0) return; if (this.#cachedMinCursorConsumers === 0) { this.#recomputeMinCursor(); } @@ -808,7 +863,7 @@ function share(source, options = { __proto__: null }) { backpressure = 'strict', signal, } = options; - validateInteger(budget, 'options.budget', 16384); + validateBudget(budget); const opts = { __proto__: null, @@ -837,11 +892,15 @@ function shareSync(source, options = { __proto__: null }) { budget = kMultiConsumerDefaultBudget, backpressure = 'strict', } = options; - validateInteger(budget, 'options.budget', 16384); - if (backpressure === 'unbounded') { + validateBudget(budget); + // A synchronous consumer can neither wait for the slowest consumer to + // release budget ('unbounded') nor keep pulling and discarding until it + // does ('drop-newest'): the slowest consumer cannot advance while the + // call is running. + if (backpressure === 'unbounded' || backpressure === 'drop-newest') { throw new ERR_INVALID_ARG_VALUE( 'options.backpressure', backpressure, - 'unbounded is not supported by shareSync()'); + `${backpressure} is not supported by shareSync()`); } const opts = { diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js index a6069c9ed3b7..2f72a78c0d21 100644 --- a/lib/internal/streams/iter/utils.js +++ b/lib/internal/streams/iter/utils.js @@ -4,6 +4,8 @@ const { Array, ArrayBufferPrototypeGetByteLength, ArrayBufferPrototypeGetDetached, + ArrayBufferPrototypeGetResizable, + ArrayPrototypePush, ArrayPrototypeSlice, PromiseResolve, PromiseWithResolvers, @@ -33,6 +35,7 @@ const { const { isSharedArrayBuffer, isUint8Array } = require('internal/util/types'); const { + validateInteger, validateOneOf, } = require('internal/validators'); const { @@ -234,6 +237,60 @@ function validateByteView(snapshot) { return value; } +/** + * Validate an explicit `budget` option. The spec only requires the default + * budget to be at least 16384 bytes; any positive explicit budget is valid. + * @param {unknown} budget + */ +function validateBudget(budget) { + validateInteger(budget, 'options.budget', 1); +} + +/** + * Append a chunk to `chunks` for later concatenation, and the information + * needed to detect that it was resized or detached in the meantime to + * `checks`. This provides the same guarantee as snapshotByteView() and + * validateByteView() without allocating a snapshot for every chunk in the + * common case: a non-empty view of a fixed-length, non-shared ArrayBuffer can + * only change by the buffer being detached, which makes its byteLength 0, so + * its byteLength is all that needs to be recorded. + * @param {Uint8Array[]} chunks + * @param {Array} checks + * @param {Uint8Array} value + * @returns {number} The byteLength of `value`. + */ +function recordChunk(chunks, checks, value) { + const buffer = TypedArrayPrototypeGetBuffer(value); + const byteLength = TypedArrayPrototypeGetByteLength(value); + ArrayPrototypePush(chunks, value); + if (byteLength === 0 || isSharedArrayBuffer(buffer) || + ArrayBufferPrototypeGetResizable(buffer)) { + ArrayPrototypePush(checks, snapshotByteView(value)); + } else { + ArrayPrototypePush(checks, byteLength); + } + return byteLength; +} + +/** + * Validate chunks recorded with recordChunk(). + * @param {Uint8Array[]} chunks + * @param {Array} checks + * @returns {Uint8Array[]} `chunks` + */ +function validateRecordedChunks(chunks, checks) { + for (let i = 0; i < chunks.length; i++) { + const check = checks[i]; + if (typeof check !== 'number') { + validateByteView(check); + } else if (TypedArrayPrototypeGetByteLength(chunks[i]) !== check) { + throw new ERR_INVALID_STATE.TypeError( + 'Byte view was resized or detached after being accepted'); + } + } + return chunks; +} + function createBatchEntry(chunks) { const views = new Array(chunks.length); let byteLength = 0; @@ -245,6 +302,35 @@ function createBatchEntry(chunks) { return { __proto__: null, views, byteLength }; } +/** + * Split a batch entry into consecutive entries that are each smaller than + * `limit` bytes, preserving chunk order. A single chunk of `limit` bytes or + * more forms an entry on its own. Returns `undefined` if no split is needed. + * @param {{ views: object[], byteLength: number }} entry + * @param {number} limit + * @returns {Array<{ views: object[], byteLength: number }>|undefined} + */ +function splitBatchEntry(entry, limit) { + const { views } = entry; + if (entry.byteLength < limit || views.length < 2) return undefined; + const entries = []; + let current = []; + let byteLength = 0; + for (let i = 0; i < views.length; i++) { + const view = views[i]; + if (current.length > 0 && byteLength + view.byteLength >= limit) { + ArrayPrototypePush(entries, + { __proto__: null, views: current, byteLength }); + current = []; + byteLength = 0; + } + ArrayPrototypePush(current, view); + byteLength += view.byteLength; + } + ArrayPrototypePush(entries, { __proto__: null, views: current, byteLength }); + return entries; +} + function validateBatchEntry(entry) { const chunks = new Array(entry.views.length); for (let i = 0; i < entry.views.length; i++) { @@ -346,7 +432,10 @@ function toWriterUint8Array(chunk) { * @returns {boolean} */ function getProtocolMethod(value, symbol) { - if (value === null || typeof value !== 'object' || !(symbol in value)) { + // Functions are objects too, and may implement the protocols. + if (value === null || + (typeof value !== 'object' && typeof value !== 'function') || + !(symbol in value)) { return undefined; } const method = value[symbol]; @@ -456,6 +545,8 @@ module.exports = { concatBytes, convertChunks, createBatchEntry, + recordChunk, + splitBatchEntry, getProtocolMethod, getWriterSignal, getMinCursor, @@ -469,6 +560,8 @@ module.exports = { toWriterUint8Array, validateBackpressure, validateBatchEntry, + validateBudget, + validateRecordedChunks, validateByteView, yieldAbortable, }; diff --git a/lib/internal/streams/iter/webidl.js b/lib/internal/streams/iter/webidl.js index 273a52443bf2..ef61abb77944 100644 --- a/lib/internal/streams/iter/webidl.js +++ b/lib/internal/streams/iter/webidl.js @@ -104,6 +104,13 @@ converters.PipeToOptions = createDictionaryConverter('PipeToOptions', [ ]); converters.PipeToSyncOptions = createDictionaryConverter( 'PipeToSyncOptions', [ + // Node.js extension. + { + __proto__: null, + key: 'failOnIncompleteClose', + converter: baseConverters.boolean, + defaultValue: () => false, + }, { __proto__: null, key: 'preventClose', diff --git a/test/parallel/test-fs-promises-file-handle-pull.js b/test/parallel/test-fs-promises-file-handle-pull.js index cdfaf273f933..055a1e395298 100644 --- a/test/parallel/test-fs-promises-file-handle-pull.js +++ b/test/parallel/test-fs-promises-file-handle-pull.js @@ -148,17 +148,30 @@ async function testPullLocking() { const fh = await open(filePath, 'r'); try { - // First pull locks the handle + // The handle is locked once the first iterable starts being consumed. const readable = fh.pull(); + const other = fh.pull(); + const iter = readable[Symbol.asyncIterator](); + const first = await iter.next(); + assert.strictEqual(first.done, false); - // Second pull while locked should throw + // Consuming a second iterable while locked should fail, and so should + // creating new consumers. + await assert.rejects( + other[Symbol.asyncIterator]().next(), + { code: 'ERR_INVALID_STATE' }, + ); assert.throws( () => fh.pull(), { code: 'ERR_INVALID_STATE' }, ); + assert.throws( + () => fh.writer(), + { code: 'ERR_INVALID_STATE' }, + ); - // Consume the first stream to unlock - await text(readable); + // Finish consuming the first stream to unlock + while (!(await iter.next()).done); // Now it should be usable again const readable2 = fh.pull(); @@ -417,6 +430,36 @@ async function testPullSyncArgumentValidation() { } } +// ============================================================================= +// An iterable that is never consumed must not lock the handle +// ============================================================================= + +async function testPullUnconsumedDoesNotLock() { + const filePath = path.join(tmpDir, 'pull-unconsumed.txt'); + fs.writeFileSync(filePath, 'unconsumed'); + + const fh = await open(filePath, 'r'); + try { + fh.pull(); + fh.pull((chunks) => chunks); + assert.strictEqual(await text(fh.pull()), 'unconsumed'); + const w = fh.writer(); + assert.strictEqual(w.endSync(), 0); + } finally { + await fh.close(); + } +} + +async function testPullAfterCloseRejectsOnIteration() { + const filePath = path.join(tmpDir, 'pull-closed-before-iter.txt'); + fs.writeFileSync(filePath, 'data'); + + const fh = await open(filePath, 'r'); + const readable = fh.pull(); + await fh.close(); + await assert.rejects(text(readable), { code: 'ERR_INVALID_STATE' }); +} + Promise.all([ testBasicPull(), testPullBinary(), @@ -437,4 +480,6 @@ Promise.all([ testPullChunkSize(), testPullChunkSizeSmall(), testPullSyncArgumentValidation(), + testPullUnconsumedDoesNotLock(), + testPullAfterCloseRejectsOnIteration(), ]).then(common.mustCall()); diff --git a/test/parallel/test-fs-promises-file-handle-pullsync.js b/test/parallel/test-fs-promises-file-handle-pullsync.js index 20c429972573..82bb5d82393f 100644 --- a/test/parallel/test-fs-promises-file-handle-pullsync.js +++ b/test/parallel/test-fs-promises-file-handle-pullsync.js @@ -473,6 +473,31 @@ async function testPullArgumentValidation() { // Run all tests // ============================================================================= +// ============================================================================= +// An iterable that is never consumed must not lock the handle +// ============================================================================= + +async function testPullSyncUnconsumedDoesNotLock() { + const filePath = path.join(tmpDir, 'pullsync-unconsumed.txt'); + fs.writeFileSync(filePath, 'unconsumed'); + + const fh = await open(filePath, 'r'); + try { + fh.pullSync(); + // Returning an iterator that never started must not unlock or close + // the handle on behalf of another consumer. + fh.pullSync()[Symbol.iterator]().return(); + const iter = fh.pullSync()[Symbol.iterator](); + assert.strictEqual(iter.next().done, false); + assert.throws(() => fh.writer(), { code: 'ERR_INVALID_STATE' }); + iter.return(); + const w = fh.writer(); + assert.strictEqual(w.endSync(), 0); + } finally { + await fh.close(); + } +} + Promise.all([ testBasicPullSync(), testLargeFile(), @@ -495,4 +520,5 @@ Promise.all([ testPullSyncChunkSize(), testWriterChunkSize(), testPullArgumentValidation(), + testPullSyncUnconsumedDoesNotLock(), ]).then(common.mustCall()); diff --git a/test/parallel/test-fs-promises-file-handle-writer.js b/test/parallel/test-fs-promises-file-handle-writer.js index a636844f8e57..1074c8156aad 100644 --- a/test/parallel/test-fs-promises-file-handle-writer.js +++ b/test/parallel/test-fs-promises-file-handle-writer.js @@ -1133,6 +1133,69 @@ async function testWriterWebIDLConversion() { // Run all tests // ============================================================================= +// ============================================================================= +// Overlapping (un-awaited) writes must land in call order +// ============================================================================= + +function makeTestData(size) { + const data = Buffer.allocUnsafe(size); + let seed = 0x9e3779b9; + for (let i = 0; i < size; i++) { + seed ^= seed << 13; seed ^= seed >>> 17; seed ^= seed << 5; + data[i] = seed & 0xff; + } + return data; +} + +async function testConcurrentWritesPreserveOrder() { + const data = makeTestData(512 * 1024); + for (const options of [{}, { start: 0 }]) { + for (const useWritev of [false, true]) { + const filePath = path.join( + tmpDir, + `writer-concurrent-${options.start ?? 'none'}-${useWritev}.bin`); + const fh = await open(filePath, 'w'); + const w = fh.writer(options); + const writes = []; + let offset = 0; + let i = 0; + while (offset < data.length) { + const size = 1 + ((i++ * 7919) % 9000); + const end = Math.min(offset + size, data.length); + const chunk = data.subarray(offset, end); + writes.push(useWritev ? + w.writev([chunk.subarray(0, 1), chunk.subarray(1)]) : + w.write(chunk)); + offset = end; + } + await Promise.all(writes); + assert.strictEqual(await w.end(), data.length); + await fh.close(); + assert.deepStrictEqual(fs.readFileSync(filePath), data); + } + } +} + +async function testFailStopsQueuedWrites() { + const filePath = path.join(tmpDir, 'writer-fail-queued.bin'); + const fh = await open(filePath, 'w'); + const w = fh.writer(); + const reason = new Error('stop'); + const writes = []; + for (let i = 0; i < 10; i++) { + writes.push(w.write(Buffer.alloc(1024, i))); + } + w.fail(reason); + const results = await Promise.allSettled(writes); + for (const result of results) { + assert.strictEqual(result.status, 'rejected'); + assert.strictEqual(result.reason, reason); + } + await fh.close(); + // None of the queued writes may reach the file after fail(). + assert.strictEqual(fs.statSync(filePath).size, 0); +} + Promise.all([ testBasicWrite(), testBasicWritev(), @@ -1187,4 +1250,6 @@ Promise.all([ testWriterLimitAndStart(), testWriterArgumentValidation(), testWriterWebIDLConversion(), + testConcurrentWritesPreserveOrder(), + testFailStopsQueuedWrites(), ]).then(common.mustCall()); diff --git a/test/parallel/test-quic-stream-writer-api.mjs b/test/parallel/test-quic-stream-writer-api.mjs index 6ccd52b046a4..2b9184d68df5 100644 --- a/test/parallel/test-quic-stream-writer-api.mjs +++ b/test/parallel/test-quic-stream-writer-api.mjs @@ -103,6 +103,11 @@ await clientSession.opened; { code: 'ERR_INVALID_ARG_TYPE' }, ); assert.strictEqual(w.endSync(), 12); + // Writes after the writer is closed reject with a TypeError. + await assert.rejects(w.write('closed'), + { code: 'ERR_INVALID_STATE', name: 'TypeError' }); + await assert.rejects(w.writev(['closed']), + { code: 'ERR_INVALID_STATE', name: 'TypeError' }); for await (const _ of stream) { /* drain */ } // eslint-disable-line no-unused-vars await stream.closed; } diff --git a/test/parallel/test-stream-iter-from-async.js b/test/parallel/test-stream-iter-from-async.js index 29ecdf425767..1b58fd6603ca 100644 --- a/test/parallel/test-stream-iter-from-async.js +++ b/test/parallel/test-stream-iter-from-async.js @@ -99,12 +99,61 @@ async function testFromBoundsNestedAsyncIterable() { const iterator = from(source())[Symbol.asyncIterator](); const first = await iterator.next(); assert.strictEqual(first.done, false); - assert.strictEqual(first.value.length, 128); + assert.ok(first.value.length > 0); + assert.ok(first.value.length <= 128); await iterator.return(); assert.strictEqual(nestedClosed, true); } +async function testFromDoesNotHoldBackNestedAsyncIterable() { + // Chunks from a nested async iterable must be delivered as they become + // available, not held back until the nested iterable produces more data + // or ends. + const { promise: release, resolve } = Promise.withResolvers(); + async function* nested() { + yield new Uint8Array([1]); + await release; + yield new Uint8Array([2]); + } + + async function* source() { + yield nested(); + yield [new Uint8Array([3]), Promise.resolve(new Uint8Array([4]))]; + } + + const iterator = from(source())[Symbol.asyncIterator](); + assert.deepStrictEqual(await iterator.next(), + { done: false, value: [new Uint8Array([1])] }); + resolve(); + const rest = []; + for (let r = await iterator.next(); !r.done; r = await iterator.next()) { + for (const chunk of r.value) rest.push(chunk[0]); + } + assert.deepStrictEqual(rest, [2, 3, 4]); +} + +async function testFromBoundsPreBatchedAsyncValues() { + // An async source yielding an already-batched Uint8Array[] larger than the + // batch bound is split, like the same batch from a sync source. + const big = Array.from({ length: 300 }, (_, i) => new Uint8Array([i & 0xff])); + async function* source() { + yield big; + } + const sizes = []; + for await (const batch of from(source())) sizes.push(batch.length); + assert.deepStrictEqual(sizes, [128, 128, 44]); + + // Batches within the bound are still passed through as-is. + const small = [new Uint8Array([1]), new Uint8Array([2])]; + async function* smallSource() { + yield small; + } + for await (const batch of from(smallSource())) { + assert.strictEqual(batch, small); + } +} + async function testFromSyncIterableAsAsync() { // Sync iterable passed to from() should work function* gen() { @@ -458,12 +507,37 @@ function testFromUndefinedThrows() { assert.throws(() => from(undefined), { code: 'ERR_INVALID_ARG_TYPE' }); } +async function testFromFunctionWithProtocols() { + // Functions are objects and may implement the protocols. + function asyncSource() {} + asyncSource[Symbol.for('Stream.toAsyncStreamable')] = + async () => 'async-function'; + assert.strictEqual(await text(from(asyncSource)), 'async-function'); + + function syncSource() {} + syncSource[Symbol.for('Stream.toStreamable')] = () => 'sync-function'; + assert.strictEqual(await text(from(syncSource)), 'sync-function'); + + async function* nested() { + yield asyncSource; + yield syncSource; + } + assert.strictEqual(await text(from(nested())), + 'async-functionsync-function'); + + // A function without a protocol is still rejected. + assert.throws(() => from(() => {}), { code: 'ERR_INVALID_ARG_TYPE' }); +} + Promise.all([ testFromString(), testFromAsyncGenerator(), testFromAsyncIteratorResultShapes(), testFromSourceErrorDoesNotWaitForReturn(), testFromBoundsNestedAsyncIterable(), + testFromFunctionWithProtocols(), + testFromDoesNotHoldBackNestedAsyncIterable(), + testFromBoundsPreBatchedAsyncValues(), testFromSyncIterableAsAsync(), testFromSyncIterableAwaitsPromiseValues(), testFromSyncIterableRejectsNestedAsyncIterable(), diff --git a/test/parallel/test-stream-iter-from-sync.js b/test/parallel/test-stream-iter-from-sync.js index 50d85dc94a98..133597914c45 100644 --- a/test/parallel/test-stream-iter-from-sync.js +++ b/test/parallel/test-stream-iter-from-sync.js @@ -183,7 +183,10 @@ function testFromSyncIgnoresAsyncStreamable() { // Explicit async iterable rejected function testFromSyncRejectsAsyncIterable() { async function* gen() { yield [new TextEncoder().encode('a')]; } - assert.throws(() => fromSync(gen()), { code: 'ERR_INVALID_ARG_TYPE' }); + assert.throws(() => fromSync(gen()), { + code: 'ERR_INVALID_ARG_TYPE', + message: /must be a synchronous input, not an async iterable\./, + }); } function testFromSyncPrefersIteratorForDualIterable() { @@ -212,8 +215,10 @@ function testFromSyncPrefersIteratorForThenableIterable() { // Promise rejected function testFromSyncRejectsPromise() { - assert.throws(() => fromSync(Promise.resolve('hello')), - { code: 'ERR_INVALID_ARG_TYPE' }); + assert.throws(() => fromSync(Promise.resolve('hello')), { + code: 'ERR_INVALID_ARG_TYPE', + message: /must be a synchronous input, not a promise\./, + }); } // DataView input should be converted to Uint8Array (zero-copy) @@ -239,6 +244,17 @@ function testFromSyncUndefinedThrows() { assert.throws(() => fromSync(undefined), { code: 'ERR_INVALID_ARG_TYPE' }); } +function testFromSyncFunctionWithToStreamable() { + // Functions are objects and may implement the protocol. + function source() {} + source[Symbol.for('Stream.toStreamable')] = () => 'from-function'; + assert.strictEqual(textSync(fromSync(source)), 'from-function'); + // ...also when nested inside another source. + assert.strictEqual(textSync(fromSync([source, '!'])), 'from-function!'); + // A function without a protocol is still rejected. + assert.throws(() => fromSync(() => {}), { code: 'ERR_INVALID_ARG_TYPE' }); +} + Promise.all([ testFromSyncString(), testFromSyncUint8Array(), @@ -260,4 +276,5 @@ Promise.all([ testFromSyncPrefersIteratorForThenableIterable(), testFromSyncRejectsPromise(), testFromSyncDataView(), + testFromSyncFunctionWithToStreamable(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-from-writable-lifecycle.js b/test/parallel/test-stream-iter-from-writable-lifecycle.js index cf6732c4184f..d5be23fdd03d 100644 --- a/test/parallel/test-stream-iter-from-writable-lifecycle.js +++ b/test/parallel/test-stream-iter-from-writable-lifecycle.js @@ -202,7 +202,7 @@ async function testAlreadyFinished() { const writer = fromWritable(writable); assert.strictEqual(writer.canWrite, null); await assert.rejects(writer.write('a'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }); + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }); assert.strictEqual(await writer.end(), 0); assertNoTerminalListeners(writable); } @@ -215,7 +215,7 @@ async function testAlreadyDestroyed() { const writer = fromWritable(writable); assert.strictEqual(writer.canWrite, null); await assert.rejects(writer.write('a'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }); + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }); assertNoTerminalListeners(writable); } diff --git a/test/parallel/test-stream-iter-pipeto-edge.js b/test/parallel/test-stream-iter-pipeto-edge.js index 13223a226a1f..0c5448db5a8c 100644 --- a/test/parallel/test-stream-iter-pipeto-edge.js +++ b/test/parallel/test-stream-iter-pipeto-edge.js @@ -5,7 +5,9 @@ const common = require('../common'); const assert = require('assert'); -const { pipeToSync, fromSync } = require('stream/iter'); +const { + pipeTo, pipeToSync, fromSync, push, text, +} = require('stream/iter'); // pipeToSync cannot complete when endSync() requires async fallback. async function testPipeToSyncEndSyncFailure() { @@ -22,6 +24,29 @@ async function testPipeToSyncEndSyncFailure() { assert.strictEqual(endCalled, false); } +// The data was accepted, so endSync() returning -1 does not fail the writer, +// even without preventFail, and the caller can still close it. +async function testPipeToSyncEndSyncFailureDoesNotFailWriter() { + const writer = { + writeSync() { return true; }, + endSync: common.mustCall(() => -1), + end: common.mustNotCall(), + fail: common.mustNotCall(), + }; + assert.throws(() => pipeToSync(fromSync('data'), writer), + { code: 'ERR_INVALID_STATE' }); + + // A push() writer whose consumer has not read yet cannot close + // synchronously. After the throw, the data is intact and the writer can be + // ended asynchronously. + const { writer: pushWriter, readable } = push(); + assert.throws(() => pipeToSync(fromSync(['abc', 'def']), pushWriter), + { code: 'ERR_INVALID_STATE' }); + const result = text(readable); + assert.strictEqual(await pushWriter.end(), 6); + assert.strictEqual(await result, 'abcdef'); +} + // pipeToSync requires endSync() when closing is enabled. async function testPipeToSyncNoEndSync() { let writeCalled = false; @@ -68,9 +93,95 @@ async function testPipeToSyncPreventClose() { assert.strictEqual(endCalled, false); } +// An exception thrown by writer.fail() must not replace the error that made +// the pipe fail. +async function testFailThrowingDoesNotMaskError() { + const cause = new Error('write failed'); + const syncWriter = { + writeSync() { throw cause; }, + endSync: common.mustNotCall(), + fail: common.mustCall((error) => { + assert.strictEqual(error, cause); + throw new Error('fail() threw'); + }), + }; + assert.throws(() => pipeToSync(fromSync('data'), syncWriter), + (error) => error === cause); + + const asyncWriter = { + async write() { throw cause; }, + end: common.mustNotCall(), + fail: common.mustCall((error) => { + assert.strictEqual(error, cause); + throw new Error('fail() threw'); + }), + }; + await assert.rejects(pipeTo(fromSync('data'), asyncWriter), + (error) => error === cause); + + // Same for the already-aborted signal path of pipeTo(). + const signal = AbortSignal.abort(); + await assert.rejects( + pipeTo(fromSync('data'), { + write: common.mustNotCall(), + fail() { throw new Error('fail() threw'); }, + }, { signal }), + (error) => error === signal.reason); +} + +// failOnIncompleteClose fails a writer that cannot be closed synchronously, +// e.g. a sync-only writer that has no end(). +async function testPipeToSyncFailOnIncompleteClose() { + let failReason; + const writer = { + writeSync() { return true; }, + endSync: common.mustCall(() => -1), + fail: common.mustCall((reason) => { failReason = reason; }), + }; + assert.throws( + () => pipeToSync(fromSync('data'), writer, { failOnIncompleteClose: true }), + (error) => { + assert.strictEqual(error.code, 'ERR_INVALID_STATE'); + assert.strictEqual(error, failReason); + return true; + }); + + // preventFail takes precedence. + assert.throws( + () => pipeToSync(fromSync('data'), { + writeSync() { return true; }, + endSync: common.mustCall(() => -1), + fail: common.mustNotCall(), + }, { failOnIncompleteClose: true, preventFail: true }), + { code: 'ERR_INVALID_STATE' }); + + // It has no effect when the writer closes synchronously. + assert.strictEqual(pipeToSync(fromSync('data'), { + writeSync() { return true; }, + endSync: common.mustCall(() => 4), + fail: common.mustNotCall(), + }, { failOnIncompleteClose: true }), 4); + + // A push() writer is failed with the error, so its consumer sees it. + const { writer: pushWriter, readable } = push(); + let thrown; + assert.throws(() => { + try { + pipeToSync(fromSync('abc'), pushWriter, { failOnIncompleteClose: true }); + } catch (error) { + thrown = error; + throw error; + } + }, { code: 'ERR_INVALID_STATE' }); + await assert.rejects(text(readable), (error) => error === thrown); +} + Promise.all([ testPipeToSyncEndSyncFailure(), + testPipeToSyncEndSyncFailureDoesNotFailWriter(), + testPipeToSyncFailOnIncompleteClose(), testPipeToSyncNoEndSync(), testPipeToSyncPreventFail(), testPipeToSyncPreventClose(), + testFailThrowingDoesNotMaskError(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-pull-async.js b/test/parallel/test-stream-iter-pull-async.js index 452ff3b17348..79cb1eaba9ef 100644 --- a/test/parallel/test-stream-iter-pull-async.js +++ b/test/parallel/test-stream-iter-pull-async.js @@ -65,14 +65,72 @@ async function testPullStatefulTransformReceiver() { } async function testPullWithAbortSignal() { + let started = false; async function* gen() { + started = true; yield [new Uint8Array([1])]; } - assert.throws( - () => pull(gen(), { signal: AbortSignal.abort() }), - { name: 'AbortError' }, - ); + // An already-aborted signal does not make pull() throw; the returned + // iterable rejects instead, without reading from the source. + const signal = AbortSignal.abort(); + const iterator = pull(gen(), { signal })[Symbol.asyncIterator](); + await assert.rejects(iterator.next(), (error) => error === signal.reason); + await assert.rejects(iterator.next(), (error) => error === signal.reason); + assert.strictEqual(started, false); + assert.deepStrictEqual(await iterator.return(), + { __proto__: null, done: true, value: undefined }); + + await assert.rejects(text(pull(gen(), (chunks) => chunks, { signal })), + (error) => error === signal.reason); + assert.strictEqual(started, false); +} + +async function testPullKeepsRejectingAfterAbort() { + for (const transforms of [[], [(chunks) => chunks]]) { + // Abort while a read is pending. + { + const ac = new AbortController(); + const reason = new Error('stop'); + async function* gen() { + yield [new Uint8Array([1])]; + await new Promise(() => {}); + } + const iterator = + pull(gen(), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + const pending = iterator.next(); + ac.abort(reason); + await assert.rejects(pending, (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + } + // Abort between reads. + { + const ac = new AbortController(); + const reason = new Error('stop'); + async function* gen() { + yield [new Uint8Array([1])]; + yield [new Uint8Array([2])]; + } + const iterator = + pull(gen(), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + ac.abort(reason); + await assert.rejects(iterator.next(), (error) => error === reason); + await assert.rejects(iterator.next(), (error) => error === reason); + } + // Aborting after the pipeline completed does not change the result. + { + const ac = new AbortController(); + const iterator = + pull(from('x'), ...transforms, { signal: ac.signal })[Symbol.asyncIterator](); + assert.strictEqual((await iterator.next()).done, false); + assert.strictEqual((await iterator.next()).done, true); + ac.abort(); + assert.strictEqual((await iterator.next()).done, true); + } + } } async function testPullNormalizesSourceAtCallTime() { @@ -98,7 +156,7 @@ async function testPullNormalizesSourceAtCallTime() { assert.strictEqual(iteratorCalls, 1); } -function testPullPreAbortOrdering() { +async function testPullPreAbortOrdering() { const reason = new Error('already aborted'); let protocolCalls = 0; const source = { @@ -109,8 +167,11 @@ function testPullPreAbortOrdering() { }; const signal = AbortSignal.abort(reason); - assert.throws(() => pull(source, { signal }), (error) => error === reason); + // Source conversion still happens when pull() is called; the abort is + // reported when the result is read. + const result = pull(source, { signal }); assert.strictEqual(protocolCalls, 1); + await assert.rejects(text(result), (error) => error === reason); assert.throws( () => pull(null, { signal }), { code: 'ERR_INVALID_ARG_TYPE' }, @@ -528,6 +589,7 @@ async function testTransformOptionsNotShared() { testPullStatefulTransform(), testPullStatefulTransformReceiver(), testPullWithAbortSignal(), + testPullKeepsRejectingAfterAbort(), testPullNormalizesSourceAtCallTime(), testPullPreAbortOrdering(), testPullChainedTransforms(), diff --git a/test/parallel/test-stream-iter-push-backpressure.js b/test/parallel/test-stream-iter-push-backpressure.js index 9fb3325f3423..58ba2d331dcb 100644 --- a/test/parallel/test-stream-iter-push-backpressure.js +++ b/test/parallel/test-stream-iter-push-backpressure.js @@ -150,6 +150,22 @@ async function testStrictPendingQueueOverflow() { await iter.return(); } +async function testSmallBudget() { + // An explicit budget below the 16384-byte default is honored. + const { writer, readable } = push({ budget: 4 }); + assert.strictEqual(writer.writeSync('abc'), true); + assert.strictEqual(writer.canWrite, true); + // The buffer may overshoot the budget by one write. + assert.strictEqual(writer.writeSync('defgh'), true); + assert.strictEqual(writer.canWrite, false); + assert.strictEqual(writer.writeSync('i'), false); + const write = writer.write('ij'); + const result = text(readable); + await write; + writer.endSync(); + assert.strictEqual(await result, 'abcdefghij'); +} + Promise.all([ testStrictBackpressure(), testDropOldest(), @@ -157,4 +173,5 @@ Promise.all([ testBlockBackpressure(), testBlockWriteSyncDoesNotEnqueue(), testStrictPendingQueueOverflow(), + testSmallBudget(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-resizable-buffers.js b/test/parallel/test-stream-iter-resizable-buffers.js index b66973393cac..91f7429ec12f 100644 --- a/test/parallel/test-stream-iter-resizable-buffers.js +++ b/test/parallel/test-stream-iter-resizable-buffers.js @@ -167,6 +167,35 @@ async function testPipeRejectsWriterResize() { ); } +async function testConsumersRejectDetachedViews() { + // Views of fixed-length buffers are tracked without a full snapshot; they + // must still be rejected when detached after being accepted. + const asyncBuffer = new ArrayBuffer(2); + async function* asyncSource() { + yield [new Uint8Array(asyncBuffer), Uint8Array.of(1)]; + asyncBuffer.transfer(); + } + await assert.rejects(array(asyncSource()), kResizeError); + const limitedBuffer = new ArrayBuffer(2); + async function* limitedSource() { + yield [new Uint8Array(limitedBuffer)]; + limitedBuffer.transfer(); + } + await assert.rejects(array(limitedSource(), { limit: 10 }), kResizeError); + + const syncBuffer = new ArrayBuffer(2); + function* syncSource() { + yield [new Uint8Array(syncBuffer, 1)]; + syncBuffer.transfer(); + } + assert.throws(() => arraySync(syncSource()), kResizeError); + + // Unchanged views are returned as-is. + const chunk = new Uint8Array(4); + const [result] = arraySync([[chunk]]); + assert.strictEqual(result, chunk); +} + Promise.all([ testBufferedViewMutationRejected(), testDropOldestUsesAcceptedByteLength(), @@ -174,5 +203,6 @@ Promise.all([ testBroadcastRejectsResizedBufferedView(), testShareRejectsResizedBufferedView(), testConsumersRejectResizedViews(), + testConsumersRejectDetachedViews(), testPipeRejectsWriterResize(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-share-async.js b/test/parallel/test-stream-iter-share-async.js index ab6bc0914798..8966c82c74e6 100644 --- a/test/parallel/test-stream-iter-share-async.js +++ b/test/parallel/test-stream-iter-share-async.js @@ -380,6 +380,64 @@ async function testShareLateJoiningConsumer() { assert.strictEqual(data2, ''); } +async function testShareRetainsBufferWhenAllConsumersDetach() { + // Data that a consumer had not read yet must stay available to consumers + // that attach after every previous consumer has detached. + const enc = new TextEncoder(); + async function* gen() { + yield [enc.encode('a')]; + yield [enc.encode('b')]; + yield [enc.encode('c')]; + } + const shared = share(gen(), { budget: 16384 }); + const c1 = shared.pull()[Symbol.asyncIterator](); + const c2 = shared.pull()[Symbol.asyncIterator](); + assert.deepStrictEqual((await c1.next()).value, [enc.encode('a')]); + // c2 has not read 'a' yet; it is the last consumer to detach. + await c1.return(); + await c2.return(); + assert.strictEqual(shared.consumerCount, 0); + + assert.strictEqual(await text(shared.pull()), 'abc'); +} + +async function testShareDropOldestSplitsOversizedBatches() { + // from() combines the values of this sync generator into a single batch + // that is much larger than the budget. Evicting that batch as a whole would + // leave the slower consumer with nothing at all. + function* source() { + for (let i = 0; i < 50; i++) { + const chunk = new Uint8Array(4096); + chunk[0] = i; + yield chunk; + } + } + const shared = share(source(), { + budget: 65536, + backpressure: 'drop-oldest', + }); + const fast = shared.pull()[Symbol.asyncIterator](); + const slow = shared.pull()[Symbol.asyncIterator](); + + const fastSeen = []; + for (let r = await fast.next(); !r.done; r = await fast.next()) { + for (const chunk of r.value) fastSeen.push(chunk[0]); + } + assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i)); + + const slowSeen = []; + for (let r = await slow.next(); !r.done; r = await slow.next()) { + for (const chunk of r.value) slowSeen.push(chunk[0]); + } + // The slow consumer lost the oldest chunks but keeps an in-order suffix + // that fits the budget. + assert.ok(slowSeen.length > 0); + assert.ok(slowSeen.length * 4096 < 65536); + assert.deepStrictEqual( + slowSeen, + Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i)); +} + async function testShareConsumerBreak() { // Verify that a consumer breaking mid-iteration detaches properly const enc = new TextEncoder(); @@ -474,6 +532,8 @@ Promise.all([ testShareSourceError(), testShareSourceErrorFollowsBufferedData(), testShareLateJoiningConsumer(), + testShareRetainsBufferWhenAllConsumersDetach(), + testShareDropOldestSplitsOversizedBatches(), testShareConsumerBreak(), testShareMultipleConsumersConcurrentPull(), testShareConsumerConcurrentNextCalls(), diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js index fd0184df9bce..a843e7311dd1 100644 --- a/test/parallel/test-stream-iter-share-sync.js +++ b/test/parallel/test-stream-iter-share-sync.js @@ -146,72 +146,168 @@ function testShareSyncRejectsUnbounded() { ); } -function testShareSyncDropNewest() { - let pulls = 0; - function* source() { - for (let i = 0; i < 4; i++) { - pulls++; - const chunk = new Uint8Array(16384); - chunk[0] = i; - yield [chunk]; - } - } +function testShareSyncRejectsDropNewest() { + // A synchronous consumer can neither wait for the slowest consumer nor keep + // discarding until it advances, so 'drop-newest' is rejected like + // 'unbounded'. + assert.throws( + () => shareSync(fromSync('data'), { backpressure: 'drop-newest' }), + { code: 'ERR_INVALID_ARG_VALUE' }, + ); +} - const shared = shareSync(source(), { - budget: 16384, - backpressure: 'drop-newest', - }); - const fast = shared.pull()[Symbol.iterator](); - const slow = shared.pull()[Symbol.iterator](); +// shareSync() accepts string source directly (normalized via fromSync()) +function testShareSyncRetainsBufferWhenAllConsumersDetach() { + // Data that a consumer had not read yet must stay available to consumers + // that attach after every previous consumer has detached. + const enc = new TextEncoder(); + function* gen() { + yield [enc.encode('a')]; + yield [enc.encode('b')]; + yield [enc.encode('c')]; + } + const shared = shareSync(gen(), { budget: 16384 }); + const c1 = shared.pull()[Symbol.iterator](); + const c2 = shared.pull()[Symbol.iterator](); + assert.deepStrictEqual(c1.next().value, [enc.encode('a')]); + c1.return(); + c2.return(); + assert.strictEqual(shared.consumerCount, 0); - assert.strictEqual(fast.next().value[0][0], 0); + assert.strictEqual(textSync(shared.pull()), 'abc'); +} - // The budget is exhausted and the slow consumer cannot advance while this - // call is running, so exactly one entry is dropped and no value is - // available. The consumer is not detached. - assert.strictEqual(fast.next().done, true); - assert.strictEqual(pulls, 2); +function testShareSyncStrictBackpressureDetaches() { + for (const transformed of [false, true]) { + function* source() { + for (let i = 0; i < 10; i++) { + yield [new Uint8Array(16384)]; + } + } + const shared = shareSync(source(), { + budget: 32768, + backpressure: 'strict', + }); + const consumer = transformed ? + shared.pull((chunks) => chunks) : shared.pull(); + const fast = consumer[Symbol.iterator](); + // This consumer prevents the buffer from being trimmed. + const slow = shared.pull()[Symbol.iterator](); + + fast.next(); + fast.next(); + assert.throws(() => fast.next(), { code: 'ERR_OUT_OF_RANGE' }); + // The rejected consumer is detached, as with the async share. + assert.strictEqual(shared.consumerCount, 1); + assert.strictEqual(fast.next().done, true); - // The slow consumer still sees the buffered entry, which releases budget. - assert.strictEqual(slow.next().value[0][0], 0); + // The detached consumer no longer pins the buffer, so the remaining + // consumer can read the whole source. + let count = 0; + while (!slow.next().done) count++; + assert.strictEqual(count, 10); + } +} - // Entry 1 was dropped for every consumer, so both resume at entry 2. - assert.strictEqual(slow.next().value[0][0], 2); - assert.strictEqual(pulls, 3); - assert.strictEqual(fast.next().value[0][0], 2); +function testShareSyncStrictForOfDoesNotWedgeOthers() { + // for...of does not call return() when next() throws. The consumer that + // hit the budget must still not keep the other consumers from reading. + function* source() { + for (let i = 0; i < 20; i++) yield [new Uint8Array(8192)]; + } + const shared = shareSync(source(), { budget: 16384, backpressure: 'strict' }); + const slow = shared.pull()[Symbol.iterator](); + assert.throws(() => { + // eslint-disable-next-line no-unused-vars + for (const _ of shared.pull()) { /* consume */ } + }, { code: 'ERR_OUT_OF_RANGE' }); + assert.strictEqual(shared.consumerCount, 1); + let count = 0; + while (!slow.next().done) count++; + assert.strictEqual(count, 20); } -// Regression test: a full buffer must not spin pulling-and-discarding from an -// unbounded source, since discarding never reclaims budget. -function testShareSyncDropNewestUnboundedSource() { - let pulls = 0; +function testShareSyncDropOldestSplitsOversizedBatches() { + // fromSync() combines the values of this generator into a single batch + // that is much larger than the budget. Evicting that batch as a whole would + // leave the slower consumer with nothing at all. function* source() { - for (;;) { - pulls++; - yield [new Uint8Array(16384)]; + for (let i = 0; i < 50; i++) { + const chunk = new Uint8Array(4096); + chunk[0] = i; + yield chunk; } } - const shared = shareSync(source(), { - budget: 16384, - backpressure: 'drop-newest', + budget: 65536, + backpressure: 'drop-oldest', }); const fast = shared.pull()[Symbol.iterator](); - shared.pull(); + const slow = shared.pull()[Symbol.iterator](); - assert.strictEqual(fast.next().done, false); - assert.strictEqual(pulls, 1); + const fastSeen = []; + for (let r = fast.next(); !r.done; r = fast.next()) { + for (const chunk of r.value) fastSeen.push(chunk[0]); + } + assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i)); - // Each blocked call drops at most one entry and returns without a value. - for (let i = 0; i < 3; i++) { - assert.strictEqual(fast.next().done, true); - assert.strictEqual(pulls, 2 + i); + const slowSeen = []; + for (let r = slow.next(); !r.done; r = slow.next()) { + for (const chunk of r.value) slowSeen.push(chunk[0]); } + // The slow consumer lost the oldest chunks but keeps an in-order suffix + // that fits the budget. + assert.ok(slowSeen.length > 0); + assert.ok(slowSeen.length * 4096 < 65536); + assert.deepStrictEqual( + slowSeen, + Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i)); +} - shared.cancel(); +function testShareSyncReentrantSourceRead() { + const enc = new TextEncoder(); + let sibling; + let reentrantError; + function* source() { + yield [enc.encode('a')]; + try { + sibling.next(); + } catch (err) { + reentrantError = err; + } + yield [enc.encode('b')]; + } + const shared = shareSync(source()); + const c1 = shared.pull()[Symbol.iterator](); + sibling = shared.pull()[Symbol.iterator](); + + assert.deepStrictEqual(c1.next().value, [enc.encode('a')]); + assert.deepStrictEqual(sibling.next().value, [enc.encode('a')]); + // Pulling 'b' runs the source, which tries to read its own share. + assert.deepStrictEqual(c1.next().value, [enc.encode('b')]); + assert.strictEqual(reentrantError?.code, 'ERR_INVALID_STATE'); + // The failed re-entrant read left the share intact. + assert.deepStrictEqual(sibling.next().value, [enc.encode('b')]); + assert.strictEqual(c1.next().done, true); + assert.strictEqual(sibling.next().done, true); +} + +function testShareSyncReentrantSourceReadUncaught() { + let sibling; + function* source() { + yield [new Uint8Array(1)]; + sibling.next(); + } + const shared = shareSync(source()); + const c1 = shared.pull()[Symbol.iterator](); + sibling = shared.pull()[Symbol.iterator](); + c1.next(); + sibling.next(); + // The error escapes the source, so it becomes the share's source error. + assert.throws(() => c1.next(), { code: 'ERR_INVALID_STATE' }); + assert.throws(() => sibling.next(), { code: 'ERR_INVALID_STATE' }); } -// shareSync() accepts string source directly (normalized via fromSync()) function testShareSyncStringSource() { const shared = shareSync('hello-sync-share'); const result = textSync(shared.pull()); @@ -227,7 +323,12 @@ Promise.all([ testShareSyncCancelWithFalsyReason(), testShareSyncSourceError(), testShareSyncRejectsUnbounded(), - testShareSyncDropNewest(), - testShareSyncDropNewestUnboundedSource(), + testShareSyncRejectsDropNewest(), + testShareSyncDropOldestSplitsOversizedBatches(), + testShareSyncReentrantSourceRead(), + testShareSyncReentrantSourceReadUncaught(), testShareSyncStringSource(), + testShareSyncRetainsBufferWhenAllConsumersDetach(), + testShareSyncStrictBackpressureDetaches(), + testShareSyncStrictForOfDoesNotWedgeOthers(), ]).then(common.mustCall()); diff --git a/test/parallel/test-stream-iter-validation.js b/test/parallel/test-stream-iter-validation.js index b130a92ae3de..16ff12c060e7 100644 --- a/test/parallel/test-stream-iter-validation.js +++ b/test/parallel/test-stream-iter-validation.js @@ -21,15 +21,18 @@ const { // push() validation // ============================================================================= -// Budget must be integer >= 16384 +// Budget must be an integer >= 1 assert.throws(() => push({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => push({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected +// WebIDL conversion truncates fractions: 0.5 becomes 0, 1.5 becomes 1. +assert.throws(() => push({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.strictEqual(push({ budget: 1.5 }).writer.canWrite, true); +// Values < 1 are rejected assert.throws(() => push({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => push({ budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => push({ budget: -100 }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => push({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum accepted value +// Only the default must be at least 16384; smaller explicit budgets are valid +assert.strictEqual(push({ budget: 1 }).writer.canWrite, true); +assert.strictEqual(push({ budget: 16383 }).writer.canWrite, true); assert.strictEqual(push({ budget: 16384 }).writer.canWrite, true); // MAX_SAFE_INTEGER is accepted assert.strictEqual(push({ budget: Number.MAX_SAFE_INTEGER }).writer.canWrite, @@ -87,11 +90,11 @@ assert.throws(() => duplex({ b: 'bad' }), { code: 'ERR_INVALID_ARG_TYPE' }); // Budget validation (cascades through to push()) assert.throws(() => duplex({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => duplex({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => duplex({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => duplex({ budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected (both directions) +// Values < 1 are rejected (both directions) assert.throws(() => duplex({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); // MAX_SAFE_INTEGER is accepted { @@ -129,17 +132,16 @@ assert.throws(() => pullSync(fromSync('a'), 42), { code: 'ERR_INVALID_ARG_TYPE' // ============================================================================= assert.throws(() => broadcast({ budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => broadcast({ budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => broadcast({ budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => broadcast({ budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// Values < 16384 are rejected +// Values < 1 are rejected assert.throws(() => broadcast({ budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => broadcast({ budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => broadcast({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum accepted value +// Small explicit budgets are accepted { - const bc = broadcast({ budget: 16384 }); + const bc = broadcast({ budget: 1 }); bc.broadcast.push(); assert.strictEqual(bc.writer.canWrite, true); bc.writer.endSync(); @@ -235,7 +237,7 @@ assert.throws( assert.throws(() => share(42), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => share(from('a'), { budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => share(from('a'), { budget: 1.5 }), { code: 'ERR_OUT_OF_RANGE' }); +assert.throws(() => share(from('a'), { budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); @@ -250,28 +252,28 @@ assert.throws(() => share(from('a'), { backpressure: 'bad' }), { code: 'ERR_INVA assert.strictEqual(shared.consumerCount, 0); } -// share() values < 16384 are rejected +// share() values < 1 are rejected assert.throws(() => share(from('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum, MAX_SAFE_INTEGER is accepted -share(from('a'), { budget: 16384 }).cancel(); +// 1 is the minimum, MAX_SAFE_INTEGER is accepted +share(from('a'), { budget: 1 }).cancel(); share(from('a'), { budget: Number.MAX_SAFE_INTEGER }).cancel(); assert.throws(() => shareSync(42), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => shareSync(fromSync('a'), { budget: 'bad' }), { code: 'ERR_OUT_OF_RANGE' }); -assert.throws(() => shareSync(fromSync('a'), { budget: 1.5 }), +assert.throws(() => shareSync(fromSync('a'), { budget: 0.5 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => shareSync(fromSync('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), { code: 'ERR_OUT_OF_RANGE' }); -// shareSync() values < 16384 are rejected +// shareSync() values < 1 are rejected assert.throws(() => shareSync(fromSync('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => shareSync(fromSync('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); -// 16384 is the minimum, MAX_SAFE_INTEGER is accepted -shareSync(fromSync('a'), { budget: 16384 }).cancel(); +// 1 is the minimum, MAX_SAFE_INTEGER is accepted +shareSync(fromSync('a'), { budget: 1 }).cancel(); shareSync(fromSync('a'), { budget: Number.MAX_SAFE_INTEGER }).cancel(); // Share.from / SyncShare.fromSync reject non-iterable diff --git a/test/parallel/test-stream-iter-writable-interop.js b/test/parallel/test-stream-iter-writable-interop.js index af2e05aae94c..f00aae924259 100644 --- a/test/parallel/test-stream-iter-writable-interop.js +++ b/test/parallel/test-stream-iter-writable-interop.js @@ -718,7 +718,7 @@ async function testWriteAfterEnd() { await assert.rejects( writer.write('should fail'), - { code: 'ERR_STREAM_WRITE_AFTER_END' }, + { code: 'ERR_STREAM_WRITE_AFTER_END', name: 'TypeError' }, ); }