diff --git a/API.md b/API.md index d19ae19..b371562 100644 --- a/API.md +++ b/API.md @@ -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 @@ -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 +``` + +**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 @@ -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) ``` --- diff --git a/docs/REQUIREMENTS.md b/docs/REQUIREMENTS.md index db5cda5..48064b4 100644 --- a/docs/REQUIREMENTS.md +++ b/docs/REQUIREMENTS.md @@ -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() diff --git a/index.bs b/index.bs index 0e7d8fb..09741a4 100644 --- a/index.bs +++ b/index.bs @@ -475,6 +475,8 @@ namespace Stream { ArrayBuffer arrayBufferSync(any source, optional ConsumeSyncOptions options = {}); Promise<sequence<Uint8Array>> array(any source, optional ConsumeOptions options = {}); sequence<Uint8Array> arraySync(any source, optional ConsumeSyncOptions options = {}); + Promise<undefined> dump(any source, optional ConsumeOptions options = {}); + undefined dumpSync(any source, optional ConsumeSyncOptions options = {}); /* Utilities */ StatelessTransformFn tap(any callback); @@ -843,7 +845,7 @@ The pipeToSync(source, ...args) 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} -------------------------------------------------------- @@ -919,6 +921,33 @@ The array(source, options) method collects all ch The arraySync(source, options) method performs the same algorithm synchronously. +`Stream.dump()` / `Stream.dumpSync()` {#dump-method} +-------------------------------------------------------- + +
+The dump(source, options) method reads |source| to completion and discards everything it reads. + +
    +
  1. Let |normalized| be the result of {{Stream/from()}} with |source|. +
  2. Let |signal| be |options|["{{ConsumeOptions/signal}}"] if present. +
  3. Let |limit| be |options|["{{ConsumeOptions/limit}}"] if present. +
  4. If |signal| is present and [=AbortSignal/aborted=], return [=a promise rejected with=] its abort reason. +
  5. Let |totalBytes| be 0. +
  6. 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. +
  7. Return undefined. +
+
+ +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 dumpSync(source, options) method performs the same algorithm synchronously using {{Stream/fromSync()}} for normalization, and returns undefined. + Utilities {#utilities} ===================== diff --git a/src/consumers.js b/src/consumers.js index 0e9e957..1306528 100644 --- a/src/consumers.js +++ b/src/consumers.js @@ -67,10 +67,12 @@ exports.bytesSync = bytesSync; exports.textSync = textSync; exports.arrayBufferSync = arrayBufferSync; exports.arraySync = arraySync; +exports.dumpSync = dumpSync; exports.bytes = bytes; exports.text = text; exports.arrayBuffer = arrayBuffer; exports.array = array; +exports.dump = dump; exports.tap = tap; exports.tapSync = tapSync; exports.ondrain = ondrain; @@ -176,6 +178,34 @@ function arraySync(source, options) { } return chunks; } +/** + * Read a sync source to completion, discarding everything it yields. + * + * Unlike the other consumers, nothing is retained: peak memory is one batch + * regardless of how much the source produces. Reading is also what releases a + * source's backpressure budget, so a source whose payload is not wanted still + * needs to be read rather than abandoned. + * + * @param source - Sync iterable yielding Uint8Array[] batches + * @param options - Optional limit + */ +function dumpSync(source, options) { + var limit = options === null || options === void 0 ? void 0 : options.limit; + var totalBytes = 0; + for (var _i = 0, source_3 = source; _i < source_3.length; _i++) { + var batch = source_3[_i]; + // Fast path: with no limit there is no reason to look at the chunks at all. + if (limit === undefined) + continue; + for (var _a = 0, batch_3 = batch; _a < batch_3.length; _a++) { + var chunk = batch_3[_a]; + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError("Stream exceeded byte limit of ".concat(limit)); + } + } + } +} // ============================================================================= // Async Consumers // ============================================================================= @@ -188,8 +218,8 @@ function arraySync(source, options) { */ function bytes(source, options) { return __awaiter(this, void 0, void 0, function () { - var signal, limit, chunks, batch, _i, batch_3, chunk, e_1_1, _a, source_3, batch, _b, batch_4, chunk, totalBytes, batch, _c, batch_5, chunk, e_2_1, _d, source_4, batch, _e, batch_6, chunk; - var _f, source_5, source_5_1, _g, source_6, source_6_1; + var signal, limit, chunks, batch, _i, batch_4, chunk, e_1_1, _a, source_4, batch, _b, batch_5, chunk, totalBytes, batch, _c, batch_6, chunk, e_2_1, _d, source_5, batch, _e, batch_7, chunk; + var _f, source_6, source_6_1, _g, source_7, source_7_1; var _h, e_1, _j, _k, _l, e_2, _m, _o; var _p, _q, _r; return __generator(this, function (_s) { @@ -207,16 +237,16 @@ function bytes(source, options) { _s.label = 1; case 1: _s.trys.push([1, 6, 7, 12]); - _f = true, source_5 = __asyncValues(source); + _f = true, source_6 = __asyncValues(source); _s.label = 2; - case 2: return [4 /*yield*/, source_5.next()]; + case 2: return [4 /*yield*/, source_6.next()]; case 3: - if (!(source_5_1 = _s.sent(), _h = source_5_1.done, !_h)) return [3 /*break*/, 5]; - _k = source_5_1.value; + if (!(source_6_1 = _s.sent(), _h = source_6_1.done, !_h)) return [3 /*break*/, 5]; + _k = source_6_1.value; _f = false; batch = _k; - for (_i = 0, batch_3 = batch; _i < batch_3.length; _i++) { - chunk = batch_3[_i]; + for (_i = 0, batch_4 = batch; _i < batch_4.length; _i++) { + chunk = batch_4[_i]; chunks.push(chunk); } _s.label = 4; @@ -230,8 +260,8 @@ function bytes(source, options) { return [3 /*break*/, 12]; case 7: _s.trys.push([7, , 10, 11]); - if (!(!_f && !_h && (_j = source_5.return))) return [3 /*break*/, 9]; - return [4 /*yield*/, _j.call(source_5)]; + if (!(!_f && !_h && (_j = source_6.return))) return [3 /*break*/, 9]; + return [4 /*yield*/, _j.call(source_6)]; case 8: _s.sent(); _s.label = 9; @@ -243,10 +273,10 @@ function bytes(source, options) { case 12: return [3 /*break*/, 14]; case 13: if ((0, from_js_1.isSyncIterable)(source)) { - for (_a = 0, source_3 = source; _a < source_3.length; _a++) { - batch = source_3[_a]; - for (_b = 0, batch_4 = batch; _b < batch_4.length; _b++) { - chunk = batch_4[_b]; + for (_a = 0, source_4 = source; _a < source_4.length; _a++) { + batch = source_4[_a]; + for (_b = 0, batch_5 = batch; _b < batch_5.length; _b++) { + chunk = batch_5[_b]; chunks.push(chunk); } } @@ -262,20 +292,20 @@ function bytes(source, options) { _s.label = 16; case 16: _s.trys.push([16, 21, 22, 27]); - _g = true, source_6 = __asyncValues(source); + _g = true, source_7 = __asyncValues(source); _s.label = 17; - case 17: return [4 /*yield*/, source_6.next()]; + case 17: return [4 /*yield*/, source_7.next()]; case 18: - if (!(source_6_1 = _s.sent(), _l = source_6_1.done, !_l)) return [3 /*break*/, 20]; - _o = source_6_1.value; + if (!(source_7_1 = _s.sent(), _l = source_7_1.done, !_l)) return [3 /*break*/, 20]; + _o = source_7_1.value; _g = false; batch = _o; // Check for abort on each iteration if (signal === null || signal === void 0 ? void 0 : signal.aborted) { throw (_q = signal.reason) !== null && _q !== void 0 ? _q : new DOMException('Aborted', 'AbortError'); } - for (_c = 0, batch_5 = batch; _c < batch_5.length; _c++) { - chunk = batch_5[_c]; + for (_c = 0, batch_6 = batch; _c < batch_6.length; _c++) { + chunk = batch_6[_c]; if (limit !== undefined) { totalBytes += chunk.byteLength; if (totalBytes > limit) { @@ -295,8 +325,8 @@ function bytes(source, options) { return [3 /*break*/, 27]; case 22: _s.trys.push([22, , 25, 26]); - if (!(!_g && !_l && (_m = source_6.return))) return [3 /*break*/, 24]; - return [4 /*yield*/, _m.call(source_6)]; + if (!(!_g && !_l && (_m = source_7.return))) return [3 /*break*/, 24]; + return [4 /*yield*/, _m.call(source_7)]; case 23: _s.sent(); _s.label = 24; @@ -308,14 +338,14 @@ function bytes(source, options) { case 27: return [3 /*break*/, 29]; case 28: if ((0, from_js_1.isSyncIterable)(source)) { - for (_d = 0, source_4 = source; _d < source_4.length; _d++) { - batch = source_4[_d]; + for (_d = 0, source_5 = source; _d < source_5.length; _d++) { + batch = source_5[_d]; // Check for abort on each iteration if (signal === null || signal === void 0 ? void 0 : signal.aborted) { throw (_r = signal.reason) !== null && _r !== void 0 ? _r : new DOMException('Aborted', 'AbortError'); } - for (_e = 0, batch_6 = batch; _e < batch_6.length; _e++) { - chunk = batch_6[_e]; + for (_e = 0, batch_7 = batch; _e < batch_7.length; _e++) { + chunk = batch_7[_e]; if (limit !== undefined) { totalBytes += chunk.byteLength; if (totalBytes > limit) { @@ -393,8 +423,8 @@ function arrayBuffer(source, options) { */ function array(source, options) { return __awaiter(this, void 0, void 0, function () { - var signal, limit, chunks, batch, _i, batch_7, chunk, e_3_1, _a, source_7, batch, _b, batch_8, chunk, totalBytes, batch, _c, batch_9, chunk, e_4_1, _d, source_8, batch, _e, batch_10, chunk; - var _f, source_9, source_9_1, _g, source_10, source_10_1; + var signal, limit, chunks, batch, _i, batch_8, chunk, e_3_1, _a, source_8, batch, _b, batch_9, chunk, totalBytes, batch, _c, batch_10, chunk, e_4_1, _d, source_9, batch, _e, batch_11, chunk; + var _f, source_10, source_10_1, _g, source_11, source_11_1; var _h, e_3, _j, _k, _l, e_4, _m, _o; var _p, _q, _r; return __generator(this, function (_s) { @@ -412,16 +442,16 @@ function array(source, options) { _s.label = 1; case 1: _s.trys.push([1, 6, 7, 12]); - _f = true, source_9 = __asyncValues(source); + _f = true, source_10 = __asyncValues(source); _s.label = 2; - case 2: return [4 /*yield*/, source_9.next()]; + case 2: return [4 /*yield*/, source_10.next()]; case 3: - if (!(source_9_1 = _s.sent(), _h = source_9_1.done, !_h)) return [3 /*break*/, 5]; - _k = source_9_1.value; + if (!(source_10_1 = _s.sent(), _h = source_10_1.done, !_h)) return [3 /*break*/, 5]; + _k = source_10_1.value; _f = false; batch = _k; - for (_i = 0, batch_7 = batch; _i < batch_7.length; _i++) { - chunk = batch_7[_i]; + for (_i = 0, batch_8 = batch; _i < batch_8.length; _i++) { + chunk = batch_8[_i]; chunks.push(chunk); } _s.label = 4; @@ -435,8 +465,8 @@ function array(source, options) { return [3 /*break*/, 12]; case 7: _s.trys.push([7, , 10, 11]); - if (!(!_f && !_h && (_j = source_9.return))) return [3 /*break*/, 9]; - return [4 /*yield*/, _j.call(source_9)]; + if (!(!_f && !_h && (_j = source_10.return))) return [3 /*break*/, 9]; + return [4 /*yield*/, _j.call(source_10)]; case 8: _s.sent(); _s.label = 9; @@ -448,10 +478,10 @@ function array(source, options) { case 12: return [3 /*break*/, 14]; case 13: if ((0, from_js_1.isSyncIterable)(source)) { - for (_a = 0, source_7 = source; _a < source_7.length; _a++) { - batch = source_7[_a]; - for (_b = 0, batch_8 = batch; _b < batch_8.length; _b++) { - chunk = batch_8[_b]; + for (_a = 0, source_8 = source; _a < source_8.length; _a++) { + batch = source_8[_a]; + for (_b = 0, batch_9 = batch; _b < batch_9.length; _b++) { + chunk = batch_9[_b]; chunks.push(chunk); } } @@ -467,20 +497,20 @@ function array(source, options) { _s.label = 16; case 16: _s.trys.push([16, 21, 22, 27]); - _g = true, source_10 = __asyncValues(source); + _g = true, source_11 = __asyncValues(source); _s.label = 17; - case 17: return [4 /*yield*/, source_10.next()]; + case 17: return [4 /*yield*/, source_11.next()]; case 18: - if (!(source_10_1 = _s.sent(), _l = source_10_1.done, !_l)) return [3 /*break*/, 20]; - _o = source_10_1.value; + if (!(source_11_1 = _s.sent(), _l = source_11_1.done, !_l)) return [3 /*break*/, 20]; + _o = source_11_1.value; _g = false; batch = _o; // Check for abort on each iteration if (signal === null || signal === void 0 ? void 0 : signal.aborted) { throw (_q = signal.reason) !== null && _q !== void 0 ? _q : new DOMException('Aborted', 'AbortError'); } - for (_c = 0, batch_9 = batch; _c < batch_9.length; _c++) { - chunk = batch_9[_c]; + for (_c = 0, batch_10 = batch; _c < batch_10.length; _c++) { + chunk = batch_10[_c]; if (limit !== undefined) { totalBytes += chunk.byteLength; if (totalBytes > limit) { @@ -500,8 +530,8 @@ function array(source, options) { return [3 /*break*/, 27]; case 22: _s.trys.push([22, , 25, 26]); - if (!(!_g && !_l && (_m = source_10.return))) return [3 /*break*/, 24]; - return [4 /*yield*/, _m.call(source_10)]; + if (!(!_g && !_l && (_m = source_11.return))) return [3 /*break*/, 24]; + return [4 /*yield*/, _m.call(source_11)]; case 23: _s.sent(); _s.label = 24; @@ -513,14 +543,14 @@ function array(source, options) { case 27: return [3 /*break*/, 29]; case 28: if ((0, from_js_1.isSyncIterable)(source)) { - for (_d = 0, source_8 = source; _d < source_8.length; _d++) { - batch = source_8[_d]; + for (_d = 0, source_9 = source; _d < source_9.length; _d++) { + batch = source_9[_d]; // Check for abort on each iteration if (signal === null || signal === void 0 ? void 0 : signal.aborted) { throw (_r = signal.reason) !== null && _r !== void 0 ? _r : new DOMException('Aborted', 'AbortError'); } - for (_e = 0, batch_10 = batch; _e < batch_10.length; _e++) { - chunk = batch_10[_e]; + for (_e = 0, batch_11 = batch; _e < batch_11.length; _e++) { + chunk = batch_11[_e]; if (limit !== undefined) { totalBytes += chunk.byteLength; if (totalBytes > limit) { @@ -540,6 +570,115 @@ function array(source, options) { }); }); } +/** + * Read an async or sync source to completion, discarding everything it yields. + * + * Unlike the other consumers, nothing is retained: peak memory is one batch + * regardless of how much the source produces. Reading is also what releases a + * source's backpressure budget, and for some sources what releases resources + * held on the producer's behalf, so a source whose payload is not wanted still + * needs to be read rather than abandoned. bytes() achieves the same thing but + * allocates the entire payload in order to throw it away. + * + * Ending for any reason other than normal completion - a source error, an + * abort, or exceeding the limit - rejects. A partial read is never reported as + * success. + * + * @param source - Iterable or async iterable yielding Uint8Array[] batches + * @param options - Optional signal and limit + * @returns Promise resolving to undefined + */ +function dump(source, options) { + return __awaiter(this, void 0, void 0, function () { + var signal, limit, totalBytes, batch, _i, batch_12, chunk, e_5_1, _a, source_12, batch, _b, batch_13, chunk; + var _c, source_13, source_13_1; + var _d, e_5, _e, _f; + var _g, _h, _j; + return __generator(this, function (_k) { + switch (_k.label) { + case 0: + signal = options === null || options === void 0 ? void 0 : options.signal; + limit = options === null || options === void 0 ? void 0 : options.limit; + // Check for abort + if (signal === null || signal === void 0 ? void 0 : signal.aborted) { + throw (_g = signal.reason) !== null && _g !== void 0 ? _g : new DOMException('Aborted', 'AbortError'); + } + totalBytes = 0; + if (!(0, from_js_1.isAsyncIterable)(source)) return [3 /*break*/, 13]; + _k.label = 1; + case 1: + _k.trys.push([1, 6, 7, 12]); + _c = true, source_13 = __asyncValues(source); + _k.label = 2; + case 2: return [4 /*yield*/, source_13.next()]; + case 3: + if (!(source_13_1 = _k.sent(), _d = source_13_1.done, !_d)) return [3 /*break*/, 5]; + _f = source_13_1.value; + _c = false; + batch = _f; + // Check for abort on each iteration + if (signal === null || signal === void 0 ? void 0 : signal.aborted) { + throw (_h = signal.reason) !== null && _h !== void 0 ? _h : new DOMException('Aborted', 'AbortError'); + } + if (limit === undefined) + return [3 /*break*/, 4]; + for (_i = 0, batch_12 = batch; _i < batch_12.length; _i++) { + chunk = batch_12[_i]; + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError("Stream exceeded byte limit of ".concat(limit)); + } + } + _k.label = 4; + case 4: + _c = true; + return [3 /*break*/, 2]; + case 5: return [3 /*break*/, 12]; + case 6: + e_5_1 = _k.sent(); + e_5 = { error: e_5_1 }; + return [3 /*break*/, 12]; + case 7: + _k.trys.push([7, , 10, 11]); + if (!(!_c && !_d && (_e = source_13.return))) return [3 /*break*/, 9]; + return [4 /*yield*/, _e.call(source_13)]; + case 8: + _k.sent(); + _k.label = 9; + case 9: return [3 /*break*/, 11]; + case 10: + if (e_5) throw e_5.error; + return [7 /*endfinally*/]; + case 11: return [7 /*endfinally*/]; + case 12: return [3 /*break*/, 14]; + case 13: + if ((0, from_js_1.isSyncIterable)(source)) { + for (_a = 0, source_12 = source; _a < source_12.length; _a++) { + batch = source_12[_a]; + // Check for abort on each iteration + if (signal === null || signal === void 0 ? void 0 : signal.aborted) { + throw (_j = signal.reason) !== null && _j !== void 0 ? _j : new DOMException('Aborted', 'AbortError'); + } + if (limit === undefined) + continue; + for (_b = 0, batch_13 = batch; _b < batch_13.length; _b++) { + chunk = batch_13[_b]; + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError("Stream exceeded byte limit of ".concat(limit)); + } + } + } + } + else { + throw new TypeError('Source must be iterable'); + } + _k.label = 14; + case 14: return [2 /*return*/]; + } + }); + }); +} /** * Create a pass-through transform that observes chunks without modifying them. * Useful for logging, hashing, metrics, etc. @@ -667,9 +806,9 @@ function merge() { return _a = {}, _a[Symbol.asyncIterator] = function () { return __asyncGenerator(this, arguments, function _a() { - var signal, _b, _c, _d, batch, e_5_1, states, startIterator, pending, _e, index, result, returnPromises; + var signal, _b, _c, _d, batch, e_6_1, states, startIterator, pending, _e, index, result, returnPromises; var _this = this; - var _f, e_5, _g, _h; + var _f, e_6, _g, _h; var _j, _k, _l; return __generator(this, function (_m) { switch (_m.label) { @@ -708,8 +847,8 @@ function merge() { return [3 /*break*/, 4]; case 9: return [3 /*break*/, 16]; case 10: - e_5_1 = _m.sent(); - e_5 = { error: e_5_1 }; + e_6_1 = _m.sent(); + e_6 = { error: e_6_1 }; return [3 /*break*/, 16]; case 11: _m.trys.push([11, , 14, 15]); @@ -720,7 +859,7 @@ function merge() { _m.label = 13; case 13: return [3 /*break*/, 15]; case 14: - if (e_5) throw e_5.error; + if (e_6) throw e_6.error; return [7 /*endfinally*/]; case 15: return [7 /*endfinally*/]; case 16: return [4 /*yield*/, __await(void 0)]; diff --git a/src/consumers.test.ts b/src/consumers.test.ts index 540fc93..a63ac00 100644 --- a/src/consumers.test.ts +++ b/src/consumers.test.ts @@ -6,6 +6,8 @@ import { describe, it } from 'node:test'; import * as assert from 'node:assert'; +import v8 from 'node:v8'; +import vm from 'node:vm'; import { bytes, bytesSync, @@ -15,6 +17,8 @@ import { arrayBufferSync, array, arraySync, + dump, + dumpSync, tap, tapSync, merge, @@ -30,6 +34,34 @@ async function* delayedSource(items: string[], delayMs: number) { } } +// Obtain a gc() without requiring the test runner to pass --expose-gc. +const gc: (() => void) | undefined = (() => { + const existing = (globalThis as { gc?: () => void }).gc; + if (existing) return existing; + try { + v8.setFlagsFromString('--expose-gc'); + const fn = vm.runInNewContext('gc') as () => void; + v8.setFlagsFromString('--no-expose-gc'); + return fn; + } catch { + return undefined; + } +})(); + +/** + * Run gc() until condition() holds, or give up after maxCount attempts. + * Returns whether the condition was met. + */ +async function gcUntil(condition: () => boolean, maxCount = 10) { + if (!gc) return false; + for (let i = 0; i < maxCount; i++) { + await new Promise((resolve) => setImmediate(resolve)); + gc(); + if (condition()) return true; + } + return false; +} + describe('bytesSync()', () => { it('should collect all bytes from source [BYTES-003]', () => { const source = fromSync('Hello, World!'); @@ -297,6 +329,207 @@ describe('array()', () => { }); }); +describe('dumpSync()', () => { + it('should read a sync source to completion [DUMP-002, DUMP-012]', () => { + let pulls = 0; + function* gen() { + for (let i = 0; i < 5; i++) { + pulls++; + yield [new Uint8Array([i])]; + } + } + dumpSync(gen()); + assert.strictEqual(pulls, 5); + }); + + it('should return undefined [DUMP-013]', () => { + assert.strictEqual(dumpSync(fromSync('hello')), undefined); + }); + + it('should handle an empty source [DUMP-005]', () => { + assert.strictEqual(dumpSync(fromSync([])), undefined); + }); + + it('should throw if the source throws mid-stream [DUMP-014]', () => { + function* failing() { + yield [new Uint8Array([1])]; + throw new Error('sync boom'); + } + assert.throws(() => dumpSync(failing()), /sync boom/); + }); + + it('should respect byte limit [DUMP-015]', () => { + const source = fromSync('Hello, World!'); + assert.throws( + () => dumpSync(source, { limit: 5 }), + /Stream exceeded byte limit of 5/ + ); + }); + + it('should allow data within limit [DUMP-015]', () => { + assert.strictEqual(dumpSync(fromSync('hello'), { limit: 100 }), undefined); + }); + + it('should throw TypeError on an async-only source [DUMP-016]', () => { + async function* asyncGen() { + yield [new Uint8Array([1])]; + } + assert.throws( + () => dumpSync(asyncGen() as unknown as Iterable), + TypeError + ); + }); +}); + +describe('dump()', () => { + it('should read an async source to completion [DUMP-001]', async () => { + let pulls = 0; + async function* gen() { + for (let i = 0; i < 5; i++) { + pulls++; + yield [new Uint8Array([i])]; + } + } + await dump(gen()); + assert.strictEqual(pulls, 5); + }); + + it('should fulfill with undefined [DUMP-003]', async () => { + assert.strictEqual(await dump(from('hello')), undefined); + }); + + it('should retain no data [DUMP-004]', async (t) => { + if (!gc) return t.skip('gc unavailable'); + + // Allocate a distinct buffer per chunk and keep only weak references. + // Whatever the consumer retained is still reachable afterwards, and + // whatever it discarded is collectable. A single pooled buffer would + // prove nothing: it stays reachable either way. + const kChunks = 64; + + async function measure( + consumer: (s: AsyncIterable) => Promise, + keepResult: boolean + ) { + const refs: WeakRef[] = []; + async function* source() { + for (let i = 0; i < kChunks; i++) { + const buf = new Uint8Array(64 * 1024); + refs.push(new WeakRef(buf)); + yield [buf]; + } + } + const result = await consumer(source()); + let live = kChunks; + const settled = await gcUntil(() => { + live = refs.filter((ref) => ref.deref() !== undefined).length; + return live === 0; + }); + // Keep the collecting consumer's result reachable across the gc, so + // what is measured is its retention rather than the result being + // dropped. + if (keepResult) assert.strictEqual((result as unknown[]).length, kChunks); + return { settled, live }; + } + + const dumped = await measure(dump, false); + assert.ok( + dumped.settled, + `dump() retained ${dumped.live}/${kChunks} chunks` + ); + + // Positive control: array() retains every chunk. Without this the test + // could pass simply because the collector reclaimed everything anyway. + // bytes() will not do here -- it concatenates into a new buffer, so the + // original chunks become collectable for it too. + const collected = await measure(array, true); + assert.strictEqual(collected.live, kChunks); + }); + + it('should handle an empty source [DUMP-005]', async () => { + assert.strictEqual(await dump(from([])), undefined); + }); + + it('should reject if the source errors mid-stream [DUMP-006]', async () => { + async function* failing() { + yield [new Uint8Array([1])]; + throw new Error('async boom'); + } + await assert.rejects(async () => await dump(failing()), /async boom/); + }); + + it('should respect AbortSignal [DUMP-007, DUMP-008]', async () => { + const controller = new AbortController(); + controller.abort(); + await assert.rejects( + async () => await dump(from('test'), { signal: controller.signal }), + /Abort/ + ); + }); + + it('should respect byte limit [DUMP-009]', async () => { + await assert.rejects( + async () => await dump(from('Hello, World!'), { limit: 5 }), + /Stream exceeded byte limit of 5/ + ); + }); + + it('should release the source on abrupt completion [DUMP-011]', async () => { + let returned = false; + const source = { + async *[Symbol.asyncIterator]() { + try { + while (true) yield [new Uint8Array(64)]; + } finally { + returned = true; + } + }, + }; + await assert.rejects( + async () => await dump(source, { limit: 128 }), + /Stream exceeded byte limit of 128/ + ); + assert.strictEqual(returned, true); + }); + + it('should accept a sync source [DUMP-002]', async () => { + assert.strictEqual(await dump(fromSync('sync source')), undefined); + }); + + it('should not inspect chunks when no limit is set [DUMP-010]', async () => { + // byteLength throws, so this only completes if dump() never reads it. + const hostile = Object.defineProperty(new Uint8Array(4), 'byteLength', { + get() { + throw new Error('byteLength must not be read without a limit'); + }, + }); + async function* gen() { + yield [hostile]; + } + assert.strictEqual(await dump(gen()), undefined); + + // With a limit it must inspect, and therefore surface the getter's throw. + async function* gen2() { + yield [hostile]; + } + await assert.rejects( + async () => await dump(gen2(), { limit: 10 }), + /byteLength must not be read/ + ); + }); + + it('should observe without retaining when combined with tap()', async () => { + let total = 0; + const counter = tap((chunks) => { + if (chunks !== null) { + for (const chunk of chunks) total += chunk.byteLength; + } + }); + await dump(pull(from('hello world'), counter)); + assert.strictEqual(total, 11); + }); +}); + describe('tapSync()', () => { it('should call callback with chunks [TAP-003]', () => { const observed: (Uint8Array[] | null)[] = []; diff --git a/src/consumers.ts b/src/consumers.ts index 9fe7cb0..d8710e6 100644 --- a/src/consumers.ts +++ b/src/consumers.ts @@ -143,6 +143,36 @@ export function arraySync( return chunks; } +/** + * Read a sync source to completion, discarding everything it yields. + * + * Unlike the other consumers, nothing is retained: peak memory is one batch + * regardless of how much the source produces. Reading is also what releases a + * source's backpressure budget, so a source whose payload is not wanted still + * needs to be read rather than abandoned. + * + * @param source - Sync iterable yielding Uint8Array[] batches + * @param options - Optional limit + */ +export function dumpSync( + source: Iterable, + options?: ConsumeSyncOptions +): undefined { + const limit = options?.limit; + + let totalBytes = 0; + for (const batch of source) { + // Fast path: with no limit there is no reason to look at the chunks at all. + if (limit === undefined) continue; + for (const chunk of batch) { + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError(`Stream exceeded byte limit of ${limit}`); + } + } + } +} + // ============================================================================= // Async Consumers // ============================================================================= @@ -355,6 +385,74 @@ export async function array( return chunks; } +/** + * Read an async or sync source to completion, discarding everything it yields. + * + * Unlike the other consumers, nothing is retained: peak memory is one batch + * regardless of how much the source produces. Reading is also what releases a + * source's backpressure budget, and for some sources what releases resources + * held on the producer's behalf, so a source whose payload is not wanted still + * needs to be read rather than abandoned. bytes() achieves the same thing but + * allocates the entire payload in order to throw it away. + * + * Ending for any reason other than normal completion - a source error, an + * abort, or exceeding the limit - rejects. A partial read is never reported as + * success. + * + * @param source - Iterable or async iterable yielding Uint8Array[] batches + * @param options - Optional signal and limit + * @returns Promise resolving to undefined + */ +export async function dump( + source: AsyncIterable | Iterable, + options?: ConsumeOptions +): Promise { + const signal = options?.signal; + const limit = options?.limit; + + // Check for abort + if (signal?.aborted) { + throw signal.reason ?? new DOMException('Aborted', 'AbortError'); + } + + let totalBytes = 0; + if (isAsyncIterable(source)) { + for await (const batch of source) { + // Check for abort on each iteration + if (signal?.aborted) { + throw signal.reason ?? new DOMException('Aborted', 'AbortError'); + } + + if (limit === undefined) continue; + + for (const chunk of batch) { + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError(`Stream exceeded byte limit of ${limit}`); + } + } + } + } else if (isSyncIterable(source)) { + for (const batch of source) { + // Check for abort on each iteration + if (signal?.aborted) { + throw signal.reason ?? new DOMException('Aborted', 'AbortError'); + } + + if (limit === undefined) continue; + + for (const chunk of batch) { + totalBytes += chunk.byteLength; + if (totalBytes > limit) { + throw new RangeError(`Stream exceeded byte limit of ${limit}`); + } + } + } + } else { + throw new TypeError('Source must be iterable'); + } +} + // ============================================================================= // Tap Utilities // ============================================================================= diff --git a/src/index.js b/src/index.js index 7c165e6..6da5020 100644 --- a/src/index.js +++ b/src/index.js @@ -85,11 +85,13 @@ exports.Stream = { text: consumers_js_1.text, arrayBuffer: consumers_js_1.arrayBuffer, array: consumers_js_1.array, + dump: consumers_js_1.dump, // Consumers (sync) bytesSync: consumers_js_1.bytesSync, textSync: consumers_js_1.textSync, arrayBufferSync: consumers_js_1.arrayBufferSync, arraySync: consumers_js_1.arraySync, + dumpSync: consumers_js_1.dumpSync, // Combining merge: consumers_js_1.merge, // Multi-consumer (push model) diff --git a/src/index.ts b/src/index.ts index 851ae39..6e31124 100644 --- a/src/index.ts +++ b/src/index.ts @@ -121,6 +121,8 @@ import { arrayBufferSync, array, arraySync, + dump, + dumpSync, tap, tapSync, merge, @@ -177,12 +179,14 @@ export const Stream = { text, arrayBuffer, array, + dump, // Consumers (sync) bytesSync, textSync, arrayBufferSync, arraySync, + dumpSync, // Combining merge,