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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 52 additions & 0 deletions API.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,12 +91,14 @@ Stream.bytes() // Collect as Uint8Array
Stream.text() // Collect as string
Stream.arrayBuffer() // Collect as ArrayBuffer
Stream.array() // Collect as Uint8Array[]
Stream.dump() // Read to completion, retain nothing

// Sync Consumers (Terminal)
Stream.bytesSync() // Sync collect as Uint8Array
Stream.textSync() // Sync collect as string
Stream.arrayBufferSync() // Sync collect as ArrayBuffer
Stream.arraySync() // Sync collect as Uint8Array[]
Stream.dumpSync() // Sync read to completion, retain nothing

// Multi-Consumer
Stream.broadcast() // Push-model multi-consumer
Expand Down Expand Up @@ -813,6 +815,55 @@ for (const chunk of chunks) {
}
```

### `Stream.dump(source, options?)`

Read a source to completion and discard everything it yields. Every other
consumer retains what it reads; `dump()` retains nothing, so its peak memory is
one batch no matter how much the source produces.

This exists because reading is not only how you obtain data, it is also what
releases a source's backpressure budget. For some sources, reading is what releases
resources held on the producer's behalf. A source whose payload you do not want
still has to be read rather than abandoned. `bytes(source)` achieves that too,
but allocates the entire payload in order to throw it away.

```typescript
function dump(
source: any, // Any input Stream.from() can normalize
options?: ConsumeOptions
): Promise<undefined>
```

**Options:**
- `signal?: AbortSignal` - Cancellation signal
- `limit?: number` - Max bytes (throws `RangeError` if exceeded)

There is no default `limit`: a source is read to completion unless the caller
asks for a bound. When `limit` is absent no byte accounting is performed.

Ending for any reason other than normal completion such as a source error, an abort,
or exceeding `limit`, rejects and releases the source. A partial read is never
reported as success.

**Example:**
```typescript
// Read and discard, retaining nothing.
await Stream.dump(source);

// Replaces the discard-loop idiom.
for await (const _ of source) { } // before
await Stream.dump(source); // after

// Observe without retaining, by combining with tap().
let total = 0;
await Stream.dump(Stream.pull(source, Stream.tap((chunks) => {
if (chunks !== null) for (const c of chunks) total += c.byteLength;
})));

// Bound the work when the source may be unexpectedly large.
await Stream.dump(source, { limit: 1024 * 1024 });
```

### Sync Variants

Synchronous versions for use with sync sources. Same algorithms as async
Expand All @@ -827,6 +878,7 @@ Stream.bytesSync(source, options?: ConsumeSyncOptions)
Stream.textSync(source, options?: TextConsumeSyncOptions)
Stream.arrayBufferSync(source, options?: ConsumeSyncOptions)
Stream.arraySync(source, options?: ConsumeSyncOptions)
Stream.dumpSync(source, options?: ConsumeSyncOptions)
```

---
Expand Down
21 changes: 21 additions & 0 deletions docs/REQUIREMENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -329,6 +329,27 @@ Terminal consumers that collect streams into memory.
| ARRAY-009 | array() respects byte limit | ✅ |
| ARRAY-010 | array() preserves chunk boundaries | ✅ |

### 5.5 Stream.dump() / Stream.dumpSync()

| ID | Requirement | Status |
|----|-------------|--------|
| DUMP-001 | dump() reads an async source to completion | ✅ |
| DUMP-002 | dump() reads a sync source to completion | ✅ |
| DUMP-003 | dump() fulfills with undefined | ✅ |
| DUMP-004 | dump() retains no data (peak memory is one batch) | ✅ |
| DUMP-005 | dump() handles an empty source | ✅ |
| DUMP-006 | dump() rejects if the source errors mid-stream | ✅ |
| DUMP-007 | dump() respects AbortSignal | ✅ |
| DUMP-008 | dump() rejects if an already-aborted signal is passed | ✅ |
| DUMP-009 | dump() respects byte limit | ✅ |
| DUMP-010 | dump() performs no byte accounting when limit is absent | ✅ |
| DUMP-011 | dump() releases the source on abrupt completion | ✅ |
| DUMP-012 | dumpSync() reads a sync source to completion | ✅ |
| DUMP-013 | dumpSync() returns undefined | ✅ |
| DUMP-014 | dumpSync() throws if the source throws mid-stream | ✅ |
| DUMP-015 | dumpSync() respects byte limit | ✅ |
| DUMP-016 | dumpSync() throws TypeError on an async-only source | ✅ |

---

## 6. Stream.broadcast()
Expand Down
31 changes: 30 additions & 1 deletion index.bs
Original file line number Diff line number Diff line change
Expand Up @@ -475,6 +475,8 @@ namespace Stream {
ArrayBuffer arrayBufferSync(any source, optional ConsumeSyncOptions options = {});
Promise&lt;sequence&lt;Uint8Array>> array(any source, optional ConsumeOptions options = {});
sequence&lt;Uint8Array> arraySync(any source, optional ConsumeSyncOptions options = {});
Promise&lt;undefined> dump(any source, optional ConsumeOptions options = {});
undefined dumpSync(any source, optional ConsumeSyncOptions options = {});

/* Utilities */
StatelessTransformFn tap(any callback);
Expand Down Expand Up @@ -843,7 +845,7 @@ The <dfn method for="Stream">pipeToSync(source, ...args)</dfn> method is the syn
Consumers {#consumers}
======================

Consumer functions are terminal operations that collect an entire stream into memory. All consumers accept any input that {{Stream/from()}} can normalize. The optional `limit` parameter protects against unbounded memory growth; exceeding it throws a {{RangeError}}.
Consumer functions are terminal operations that read a stream to completion. All of them except {{Stream/dump()}} collect the stream into memory; {{Stream/dump()}} reads the stream and discards it, retaining nothing beyond the current batch. All consumers accept any input that {{Stream/from()}} can normalize. The optional `limit` parameter bounds the number of bytes a consumer will read; exceeding it throws a {{RangeError}}. For the collecting consumers this protects against unbounded memory growth, and for {{Stream/dump()}} against unbounded work.

`Stream.bytes()` / `Stream.bytesSync()` {#bytes-method}
--------------------------------------------------------
Expand Down Expand Up @@ -919,6 +921,33 @@ The <dfn method for="Stream">array(source, options)</dfn> method collects all ch

The <dfn method for="Stream">arraySync(source, options)</dfn> method performs the same algorithm synchronously.

`Stream.dump()` / `Stream.dumpSync()` {#dump-method}
--------------------------------------------------------

<div algorithm>
The <dfn method for="Stream">dump(source, options)</dfn> method reads |source| to completion and discards everything it reads.

<ol>
<li>Let |normalized| be the result of {{Stream/from()}} with |source|.
<li>Let |signal| be |options|["{{ConsumeOptions/signal}}"] if present.
<li>Let |limit| be |options|["{{ConsumeOptions/limit}}"] if present.
<li>If |signal| is present and [=AbortSignal/aborted=], return [=a promise rejected with=] its abort reason.
<li>Let |totalBytes| be 0.
<li>Asynchronously iterate |normalized|. Before each iteration step, if |signal| is present and [=AbortSignal/aborted=], stop iteration and reject with |signal|'s abort reason. If |limit| is present, then for each batch yielded, for each chunk, add its byte length to |totalBytes|, and if |totalBytes| exceeds |limit|, throw a {{RangeError}}. Discard each batch without retaining it.
<li>Return undefined.
</ol>
</div>

Note: Unlike the other consumers, {{Stream/dump()}} retains no data. Its peak memory is one batch regardless of how much the source yields, so it is the appropriate way to read a stream whose contents are not wanted. Callers that need the data should use {{Stream/bytes()}} or {{Stream/array()}} instead.

Note: Reading a source to completion is what releases its [=backpressure policy=] budget, and for some sources also releases resources held on the producer's behalf. A source whose payload is unwanted therefore still needs to be read rather than abandoned, which is the case {{Stream/dump()}} exists to serve. Where the source should instead be cancelled without being read, callers should terminate iteration early or cancel the source directly.

Note: When |limit| is absent, no byte accounting is performed. Implementations are not required to inspect chunk byte lengths in that case.

Terminating for any reason other than normal completion, such as an error from the source, an abort via |signal|, or exceeding |limit|, is an abrupt completion, so the iterator's `return()` method is invoked and the source is released. {{Stream/dump()}} never reports a partial read as success.

The <dfn method for="Stream">dumpSync(source, options)</dfn> method performs the same algorithm synchronously using {{Stream/fromSync()}} for normalization, and returns undefined.


Utilities {#utilities}
=====================
Expand Down
Loading
Loading