stream: add drain()/drainSync() for stream/iter - #65598
Conversation
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.
drain() 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. drainSync() is the synchronous form.
Assisted-by: Claude Opus 5
Signed-off-by: Ethan Arrowood <ethan@arrowood.dev>
|
Review requested:
|
Replaces the `for await (const _ of source) {}` discard idiom with
drain(), which says what it does and does not need an eslint-disable
for the unused loop variable. 89 loops across 60 files, plus 32 now
redundant eslint-disable directives and comments removed.
Only files that already declare --experimental-stream-iter are
converted, so no test gains an experimental flag it did not already
opt into. That leaves 19 files with the old idiom; converting those
would change what those tests run under and belongs in a separate
discussion.
Loops that break early are left alone: those cancel the source rather
than reading it to completion, which is not what drain() does.
Four QUIC tests already bound a local `drain` to the writer-side
drainableProtocol promise. Those locals are renamed to drainPromise;
in test-quic-stream-writer-api.mjs the local shadowed the import and
broke at runtime rather than at parse time.
Assisted-by: Claude Opus 5
Signed-off-by: Ethan Arrowood <ethan@arrowood.dev>
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #65598 +/- ##
==========================================
- Coverage 90.07% 90.06% -0.02%
==========================================
Files 751 751
Lines 254881 255069 +188
Branches 48111 48177 +66
==========================================
+ Hits 229582 229720 +138
- Misses 16486 16547 +61
+ Partials 8813 8802 -11
🚀 New features to boost your workflow:
|
pimterry
left a comment
There was a problem hiding this comment.
Very onboard with the concept, implementation looks good, small point on a test case.
Larger question though: can we rename this? Drain is a separate related but different concept on the writer API. Giving them the same name is a bit confusing, and hits practical issues as well (like the various tests here that have to rename existing drain references that mean something else).
dump() or discard()?
| // drain: does not retain data | ||
| // ============================================================================= | ||
|
|
||
| async function testDrainDoesNotRetain() { |
There was a problem hiding this comment.
This doesn't really do what it says it does - as is it roughly duplicates testDrainConsumesToCompletion I think. It allocates the buffer once outside the source, so it can't check what's retained anyway. Needs a buffer per chunk with a WeakRef & GC dance if you want to test this I think.
|
I'm okay with considering a different name. Honestly, I stuck with Let the bikeshedding begin! |
The retention test allocated a single buffer outside the source and yielded it for every chunk, so it could not observe what drain() held on to: that buffer stayed reachable whether drain() retained anything or not. Its one assertion, that every chunk was handed out, already duplicated testDrainConsumesToCompletion. Allocate a distinct buffer per chunk, keep only WeakRefs to them, and use gcUntil() to check none survive the drain. Add array() as a positive control that retains every chunk, so the test fails if the measurement stops being able to observe retention rather than passing vacuously. bytes() does not work as that control: it concatenates into a new buffer, so the original chunks become collectable for it too. Requires --expose-gc. Assisted-by: Claude Opus 5 Signed-off-by: Ethan Arrowood <ethan@arrowood.dev>
|
|
Adds
drain()anddrainSync()tonode:stream/iter.This will enable us to replace
for awaitloops withawait drain(stream)across the repo.I used two commits: one for the feature, another for the migration so it's easy to review.
I only migrated tests that were already using iter streams. But there are more that could benefit from this util.
This api is inspired by things like Undici's
body.dump()method.Let me know what you think of this API addition!