Skip to content

Commit 638793e

Browse files
stream: add dump()/dumpSync() for stream/iter
Every consumer in node:stream/iter retains what it reads, so there is no way to read a streamable to completion while keeping nothing. To clear a stream, callers write `for await (const _ of source) {}`. This is apparent throughout multiple test suites, including the stream/iter tests themselves. These are especially useful for QUIC where a receiver doesn't want the payload, but still has to read it to completion to relieve backpressure. bytes() does that too, but allocates the whole payload to discard it. dump() pulls every batch and drops it, so peak memory is one batch regardless of volume. It takes the same signal and limit options as the other consumers, rejects if the source errors mid-stream, and fulfills with undefined. dumpSync() is the synchronous form. The name follows undici's body.dump(). drain() was the obvious alternative but this module already exports ondrain() and drainableProtocol() for writer-side backpressure, which is the opposite side of the pipe. Assisted-by: Claude Opus 5 Signed-off-by: Ethan Arrowood <ethan@arrowood.dev>
1 parent 1e9fd95 commit 638793e

6 files changed

Lines changed: 593 additions & 0 deletions

File tree

‎doc/api/quic.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -343,6 +343,10 @@ Only one async iterator can be obtained per stream. The stream is also
343343
compatible with `node:stream/iter` utilities such as `Stream.bytes()`,
344344
`Stream.text()`, and `Stream.pipeTo()`.
345345

346+
Consuming a stream is what returns flow-control credit to the peer, so a
347+
stream whose payload is not wanted should still be read to completion. Use
348+
`Stream.dump()` to read the stream without retaining any of it.
349+
346350
### Datagrams
347351

348352
In addition to streams, QUIC supports unreliable datagrams ([RFC 9221][]) for

‎doc/api/stream_iter.md‎

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1025,6 +1025,81 @@ added:
10251025

10261026
Synchronous version of [`bytes()`][].
10271027

1028+
### `dump(source[, options])`
1029+
1030+
<!-- YAML
1031+
added: REPLACEME
1032+
-->
1033+
1034+
* `source` {AsyncIterable|Iterable} whose chunks must be {Uint8Array\[]}
1035+
* `options` {Object}
1036+
* `signal` {AbortSignal}
1037+
* `limit` {number} Maximum number of bytes to consume. If the total bytes
1038+
read exceeds limit, an `ERR_OUT_OF_RANGE` error is thrown
1039+
* Returns: {Promise} Fulfills with `undefined`.
1040+
1041+
Read a source to completion, discarding every chunk. Unlike the other
1042+
consumers, `dump()` retains nothing. Memory tops-out at a single batch no
1043+
matter how much data the source produces.
1044+
1045+
Use this to consume a stream when the content doesn't matter. For example, A
1046+
QUIC stream only returns flow-control credit to the peer as its data is
1047+
consumed, so a receiver that does not want the payload must still read it to
1048+
completion.
1049+
1050+
If the source errors part-way through, the returned promise rejects with that
1051+
error.
1052+
1053+
There is no default `limit`. `dump()` reads until the source is exhausted
1054+
unless a limit is specified. When a limit is configured and the source exceeds
1055+
it, the promise rejects and the source is cancelled. A partial dump is never
1056+
reported as success.
1057+
1058+
```mjs
1059+
import { dump, from, pull, tap } from 'node:stream/iter';
1060+
1061+
// Count the bytes flowing through a stream without retaining any of them.
1062+
let bytesSeen = 0;
1063+
const counter = tap((chunks) => {
1064+
for (const chunk of chunks) bytesSeen += chunk.byteLength;
1065+
});
1066+
1067+
await dump(pull(from('hello world'), counter));
1068+
console.log(bytesSeen); // 11
1069+
```
1070+
1071+
```cjs
1072+
const { dump, from, pull, tap } = require('node:stream/iter');
1073+
1074+
async function run() {
1075+
// Count the bytes flowing through a stream without retaining any of them.
1076+
let bytesSeen = 0;
1077+
const counter = tap((chunks) => {
1078+
for (const chunk of chunks) bytesSeen += chunk.byteLength;
1079+
});
1080+
1081+
await dump(pull(from('hello world'), counter));
1082+
console.log(bytesSeen); // 11
1083+
}
1084+
1085+
run().catch(console.error);
1086+
```
1087+
1088+
### `dumpSync(source[, options])`
1089+
1090+
<!-- YAML
1091+
added: REPLACEME
1092+
-->
1093+
1094+
* `source` {Iterable} whose chunks must be {Uint8Array\[]}
1095+
* `options` {Object}
1096+
* `limit` {number} Maximum number of bytes to consume. If the total bytes
1097+
read exceeds limit, an `ERR_OUT_OF_RANGE` error is thrown
1098+
* Returns: {undefined}
1099+
1100+
Synchronous version of [`dump()`][]. Throws `ERR_INVALID_ARG_TYPE` if `source`
1101+
is not synchronously iterable.
1102+
10281103
### `text(source[, options])`
10291104

10301105
<!-- YAML
@@ -2167,6 +2242,7 @@ console.log(textSync(stream)); // 'hello world'
21672242
[`array()`]: #arraysource-options
21682243
[`arrayBuffer()`]: #arraybuffersource-options
21692244
[`bytes()`]: #bytessource-options
2245+
[`dump()`]: #drainsource-options
21702246
[`from()`]: #frominput
21712247
[`fromSync()`]: #fromsyncinput
21722248
[`node:zlib/iter`]: zlib.md#iterable-compression

‎lib/internal/streams/iter/consumers.js‎

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
// New Streams API - Consumers & Utilities
44
//
55
// bytes(), text(), arrayBuffer() - collect entire stream
6+
// dump() - consume entire stream, retaining nothing
67
// tap(), tapSync() - observe without modifying
78
// merge() - temporal combining of sources
89
// ondrain() - backpressure drain utility
@@ -230,6 +231,20 @@ function validateSyncConsumerOptions(options) {
230231
validateBaseConsumerOptions(options);
231232
}
232233

234+
function validateSyncDumpOptions(options) {
235+
validateObject(options, 'options');
236+
if (options.limit !== undefined) {
237+
validateInteger(options.limit, 'options.limit', 0);
238+
}
239+
}
240+
241+
function validateDumpOptions(options) {
242+
validateSyncDumpOptions(options);
243+
if (options.signal !== undefined) {
244+
validateAbortSignal(options.signal, 'options.signal');
245+
}
246+
}
247+
233248
// =============================================================================
234249
// Sync Consumers
235250
// =============================================================================
@@ -285,6 +300,32 @@ function arraySync(source, options = kNullPrototype) {
285300
return collectSync(source, options.limit);
286301
}
287302

303+
/**
304+
* Read a sync source to completion, discarding every chunk.
305+
* @param {Iterable<Uint8Array[]>} source
306+
* @param {{ limit?: number }} [options]
307+
* @returns {undefined}
308+
*/
309+
function dumpSync(source, options = kNullPrototype) {
310+
validateSyncDumpOptions(options);
311+
312+
const limit = options.limit;
313+
const normalized = fromSync(source);
314+
let totalBytes = 0;
315+
316+
for (const batch of normalized) {
317+
// With no limit, just iterate through the stream completely.
318+
if (limit === undefined) continue;
319+
// Otherwise, calculate totalBytes and track the limit.
320+
for (let i = 0; i < batch.length; i++) {
321+
totalBytes += TypedArrayPrototypeGetByteLength(batch[i]);
322+
if (totalBytes > limit) {
323+
throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes);
324+
}
325+
}
326+
}
327+
}
328+
288329
// =============================================================================
289330
// Async Consumers
290331
// =============================================================================
@@ -341,6 +382,40 @@ async function array(source, options = kNullPrototype) {
341382
return collectAsync(source, options.signal, options.limit);
342383
}
343384

385+
/**
386+
* Read an async or sync source to completion, discarding every chunk.
387+
* @param {AsyncIterable<Uint8Array[]>|Iterable<Uint8Array[]>} source
388+
* @param {{ signal?: AbortSignal, limit?: number }} [options]
389+
* @returns {Promise<undefined>}
390+
*/
391+
async function dump(source, options = kNullPrototype) {
392+
validateDumpOptions(options);
393+
const signal = options.signal;
394+
const limit = options.limit;
395+
396+
signal?.throwIfAborted();
397+
398+
const abortableSource = signal && isAsyncIterable(source) ?
399+
yieldAbortable(source, signal) : source;
400+
const normalized = from(abortableSource);
401+
const iterable = signal ? yieldAbortable(normalized, signal) : normalized;
402+
403+
let totalBytes = 0;
404+
405+
for await (const batch of iterable) {
406+
signal?.throwIfAborted();
407+
// With no limit, just iterate through the stream completely.
408+
if (limit === undefined) continue;
409+
// Otherwise, calculate totalBytes and track the limit.
410+
for (let i = 0; i < batch.length; i++) {
411+
totalBytes += TypedArrayPrototypeGetByteLength(batch[i]);
412+
if (totalBytes > limit) {
413+
throw new ERR_OUT_OF_RANGE('totalBytes', `<= ${limit}`, totalBytes);
414+
}
415+
}
416+
}
417+
}
418+
344419
// =============================================================================
345420
// Tap Utilities
346421
// =============================================================================
@@ -610,6 +685,8 @@ module.exports = {
610685
arraySync,
611686
bytes,
612687
bytesSync,
688+
dump,
689+
dumpSync,
613690
merge,
614691
ondrain,
615692
tap,

‎lib/stream/iter.js‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,8 @@ const {
4444
arrayBufferSync,
4545
array,
4646
arraySync,
47+
dump,
48+
dumpSync,
4749
tap,
4850
tapSync,
4951
merge,
@@ -100,12 +102,14 @@ const Stream = ObjectFreeze({
100102
text,
101103
arrayBuffer,
102104
array,
105+
dump,
103106

104107
// Consumers (sync)
105108
bytesSync,
106109
textSync,
107110
arrayBufferSync,
108111
arraySync,
112+
dumpSync,
109113

110114
// Combining
111115
merge,
@@ -164,12 +168,14 @@ module.exports = {
164168
text,
165169
arrayBuffer,
166170
array,
171+
dump,
167172

168173
// Consumers (sync)
169174
bytesSync,
170175
textSync,
171176
arrayBufferSync,
172177
arraySync,
178+
dumpSync,
173179

174180
// Combining
175181
merge,

0 commit comments

Comments
 (0)