Skip to content

Commit b44b865

Browse files
jasnelladuh95
authored andcommitted
stream: make stream/iter toWritable apply backpressure
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode PR-URL: #66079 Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
1 parent 6299b09 commit b44b865

3 files changed

Lines changed: 26 additions & 15 deletions

File tree

‎doc/api/stream_iter.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1757,9 +1757,9 @@ non-Error reason is wrapped in an `ERR_FALSY_VALUE_REJECTION` or
17571757
`ERR_OPERATION_FAILED` error before it is passed to the callback. The error's
17581758
`reason` property contains the original value.
17591759

1760-
The Writable's `highWaterMark` is set to `Number.MAX_SAFE_INTEGER` to
1761-
effectively disable its internal buffering, allowing the underlying Writer
1762-
to manage backpressure directly.
1760+
The Writable uses the default classic stream `highWaterMark`. Classic stream
1761+
backpressure bounds writes waiting to reach the underlying Writer, while the
1762+
Writer controls completion of the active `_write()` or `_writev()` operation.
17631763

17641764
```mjs
17651765
import { push, toWritable } from 'node:stream/iter';

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

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,6 @@
1414
const {
1515
ArrayPrototypePush,
1616
FunctionPrototypeCall,
17-
NumberMAX_SAFE_INTEGER,
1817
Promise,
1918
PromisePrototypeThen,
2019
PromiseReject,
@@ -937,11 +936,6 @@ function toWritable(writer) {
937936

938937
const writableOptions = {
939938
__proto__: null,
940-
// Use MAX_SAFE_INTEGER to effectively disable the Writable's
941-
// internal buffering. The underlying stream/iter Writer has its
942-
// own backpressure handling; we want _write to be called
943-
// immediately so the Writer can manage flow control directly.
944-
highWaterMark: NumberMAX_SAFE_INTEGER,
945939
write: _write,
946940
final: _final,
947941
destroy: _destroy,

‎test/parallel/test-stream-iter-writable-from.js‎

Lines changed: 23 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66

77
const common = require('../common');
88
const assert = require('assert');
9+
const { once } = require('events');
910
const { setImmediate, setTimeout } = require('timers/promises');
1011
const {
1112
push,
@@ -603,18 +604,33 @@ async function testDestroyWithoutFail() {
603604
}
604605

605606
// =============================================================================
606-
// Custom highWaterMark option
607+
// Classic Writable backpressure
607608
// =============================================================================
608609

609-
function testHighWaterMarkIsMaxSafeInt() {
610+
function testUsesBoundedHighWaterMark() {
610611
const writer = {
611612
write(chunk) { return Promise.resolve(); },
612613
};
613614

614-
// HWM is set to MAX_SAFE_INTEGER to disable Writable's internal
615-
// buffering. The underlying Writer manages backpressure directly.
616615
const writable = toWritable(writer);
617-
assert.strictEqual(writable.writableHighWaterMark, Number.MAX_SAFE_INTEGER);
616+
assert.ok(writable.writableHighWaterMark > 0);
617+
assert.ok(writable.writableHighWaterMark < Number.MAX_SAFE_INTEGER);
618+
}
619+
620+
async function testAppliesClassicBackpressure() {
621+
let resolveWrite;
622+
const writable = toWritable({
623+
write: common.mustCall(() => new Promise((resolve) => {
624+
resolveWrite = resolve;
625+
})),
626+
});
627+
const chunk = Buffer.alloc(writable.writableHighWaterMark);
628+
629+
assert.strictEqual(writable.write(chunk), false);
630+
const finished = once(writable, 'finish');
631+
resolveWrite();
632+
writable.end();
633+
await finished;
618634
}
619635

620636
// =============================================================================
@@ -700,10 +716,11 @@ async function testEndThrowsSyncPropagation() {
700716

701717
testInvalidWriterThrows();
702718
testNoWritevWithoutWriterWritev();
703-
testHighWaterMarkIsMaxSafeInt();
719+
testUsesBoundedHighWaterMark();
704720

705721
Promise.all([
706722
testBasicWrite(),
723+
testAppliesClassicBackpressure(),
707724
testFalsyWriterRejectionBecomesClassicError(),
708725
testClassicWrapperReusePreservesErrorIdentity(),
709726
testWriteDelegatesToWriter(),

0 commit comments

Comments
 (0)