Commit fe1a203457f for nodejs
commit fe1a203457f5a7c1dc2890372c4dc6ae1da08252
Author: James M Snell <jasnell@gmail.com>
Date: Sun Aug 30 18:33:21 2026 +0000
stream,quic: correct stream/iter writer implementations
Signed-off-by: James M Snell <jasnell@gmail.com>
Assisted-by: Opencode
PR-URL: https://github.com/nodejs/node/pull/66079
Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
diff --git a/lib/internal/quic/quic.js b/lib/internal/quic/quic.js
index 632369fc908..b1584cd0dfb 100644
--- a/lib/internal/quic/quic.js
+++ b/lib/internal/quic/quic.js
@@ -2223,10 +2223,41 @@ class QuicStream {
const handle = this.#handle;
const stream = this;
let closed = false;
+ let ending = false;
let errored = false;
let error = null;
let totalBytesWritten = 0;
let drainWakeup = null;
+ let pendingWrite = null;
+ let pendingEnd = null;
+
+ function raceWithSignal(promise, signal) {
+ if (signal === undefined) return promise;
+ signal.throwIfAborted();
+
+ const deferred = PromiseWithResolvers();
+ const onAbort = () => deferred.reject(signal.reason);
+ signal.addEventListener('abort', onAbort, {
+ __proto__: null,
+ once: true,
+ });
+ PromisePrototypeThen(
+ promise,
+ (value) => {
+ signal.removeEventListener('abort', onAbort);
+ deferred.resolve(value);
+ },
+ (reason) => {
+ signal.removeEventListener('abort', onAbort);
+ deferred.reject(reason);
+ });
+ return deferred.promise;
+ }
+
+ function waitForDrain(signal) {
+ drainWakeup ??= PromiseWithResolvers();
+ return raceWithSignal(drainWakeup.promise, signal);
+ }
// Drain callback - C++ fires this when send buffer has space
stream[kDrain] = () => {
@@ -2251,16 +2282,16 @@ class QuicStream {
};
// A note on backpressure handling: per the stream/iter spec, the default
- // backpressure policy for writers is strict, meaning that if the stream
- // signals backpressure additional writes are rejected until the buffer has
- // capacity again.
+ // backpressure policy for writers is strict. One async write may wait for
+ // capacity; additional writes are rejected until it settles.
- function writeConvertedSync(chunk) {
+ function writeConvertedSync(chunk, token) {
+ const isPending = pendingWrite === token && token !== undefined;
// If the stream is closed, errored, or write-ended, we cannot accept
// more data. Refuse the sync write.
- // If a drain is already pending, another operation is waiting
- // for capacity. Refuse the sync write.
- if (closed || errored || stream.#inner.state.writeEnded || drainWakeup != null) {
+ if (closed || errored || stream.#inner.state.writeEnded ||
+ (ending && !isPending) ||
+ (pendingWrite !== null && !isPending)) {
return false;
}
const len = TypedArrayPrototypeGetByteLength(chunk);
@@ -2292,25 +2323,52 @@ class QuicStream {
async function writeAsync(chunk, signal) {
if (errored) throw error;
- if (closed || stream.#inner.state.writeEnded) {
+ if (closed || ending || stream.#inner.state.writeEnded) {
throw new ERR_INVALID_STATE('Writer is closed');
}
signal?.throwIfAborted();
- // If a drain is already pending, another operation is waiting
- // for capacity. Under strict policy, reject immediately.
- // Later, if we add support for other backpressure policies,
- // we could instead await the existing drain before proceeding.
- if (drainWakeup != null) {
+
+ if (writeConvertedSync(chunk)) return;
+ await writeWhenDrained(chunk, writeConvertedSync, signal);
+ }
+
+ async function writeWhenDrained(chunks, writeConverted, signal) {
+ if (pendingWrite !== null) {
+ throw new ERR_INVALID_STATE.RangeError(
+ 'Backpressure violation: too many pending writes. ' +
+ 'Await each write() call to respect backpressure.');
+ }
+ if (stream.#inner.state.writeDesiredSize !== 0) {
throw new ERR_INVALID_STATE('Stream write buffer is full');
}
- if (!writeConvertedSync(chunk)) {
- throw new ERR_INVALID_STATE('Stream write buffer is full');
+ const done = PromiseWithResolvers();
+ const token = { __proto__: null, done };
+ pendingWrite = token;
+ try {
+ while (true) {
+ await waitForDrain(signal);
+ if (errored) throw error;
+ if (closed || stream.#inner.state.writeEnded) {
+ throw new ERR_INVALID_STATE('Writer is closed');
+ }
+ signal?.throwIfAborted();
+ if (writeConverted(chunks, token)) return;
+ if (stream.#inner.state.writeDesiredSize !== 0) {
+ throw new ERR_INVALID_STATE('Stream write buffer is full');
+ }
+ }
+ } finally {
+ if (pendingWrite === token) pendingWrite = null;
+ done.resolve();
}
}
- function writevConvertedSync(chunks) {
- if (closed || errored || stream.#inner.state.writeEnded || drainWakeup != null) {
+ function writevConvertedSync(chunks, token) {
+ const isPending = pendingWrite === token && token !== undefined;
+ if (closed || errored || stream.#inner.state.writeEnded ||
+ (ending && !isPending) ||
+ (pendingWrite !== null && !isPending)) {
return false;
}
let len = 0;
@@ -2335,22 +2393,13 @@ class QuicStream {
async function writevAsync(chunks, signal) {
if (errored) throw error;
- if (closed || stream.#inner.state.writeEnded) {
+ if (closed || ending || stream.#inner.state.writeEnded) {
throw new ERR_INVALID_STATE('Writer is closed');
}
signal?.throwIfAborted();
- // If a drain is already pending, another operation is waiting
- // for capacity. Under strict policy, reject immediately.
- // Later, if we add support for other backpressure policies,
- // we could instead await the existing drain before proceeding.
- if (drainWakeup != null) {
- throw new ERR_INVALID_STATE('Stream write buffer is full');
- }
-
- if (!writevConvertedSync(chunks)) {
- throw new ERR_INVALID_STATE('Stream write buffer is full');
- }
+ if (writevConvertedSync(chunks)) return;
+ await writeWhenDrained(chunks, writevConvertedSync, signal);
}
function endSync() {
@@ -2364,8 +2413,8 @@ class QuicStream {
// If we're already closed, just return the total bytes written.
if (closed) return totalBytesWritten;
- // If we are waiting for drain to complete, we cannot end synchronously.
- if (drainWakeup != null) return -1;
+ // Accepted writes and existing drain waiters must settle before ending.
+ if (ending || pendingWrite !== null || drainWakeup !== null) return -1;
// Fantastic, we can end synchronously!
handle.endWrite();
@@ -2381,13 +2430,7 @@ class QuicStream {
async function endAsync(signal) {
if (errored) throw error;
if (closed) return totalBytesWritten;
- if (signal !== undefined) {
- signal.throwIfAborted();
- // TODO(@jasnell): The stream/iter spec allows individual sync end
- // calls to be canceled via an AbortSignal. We currently do not support
- // this, but we can add before the impl is graduated from experimental.
- // At most we do here is check for signal abort at the start of the call.
- }
+ signal?.throwIfAborted();
// Per the streams/iter spec, endSync and end follow a try-fallback
// pattern. That is, callers should try endSync first and if it returns
@@ -2406,16 +2449,28 @@ class QuicStream {
if (n >= 0) return n;
if (errored) throw error;
- drainWakeup ??= PromiseWithResolvers();
- try {
- await drainWakeup.promise;
- } finally {
- drainWakeup = null;
+ if (pendingEnd === null) {
+ ending = true;
+ pendingEnd = finishEnd();
+ }
+ return raceWithSignal(pendingEnd, signal);
+ }
+
+ async function finishEnd() {
+ if (pendingWrite !== null) {
+ await pendingWrite.done.promise;
}
if (errored) throw error;
- const result = endSync();
+
+ if (drainWakeup !== null) {
+ await drainWakeup.promise;
+ }
if (errored) throw error;
- return result;
+
+ if (!stream.#inner.state.writeEnded) handle.endWrite();
+ closed = true;
+ ending = false;
+ return totalBytesWritten;
}
function fail(reason) {
@@ -2448,7 +2503,9 @@ class QuicStream {
const writer = {
__proto__: null,
get canWrite() {
- if (closed || errored || stream.#inner.state.writeEnded) return null;
+ if (closed || ending || errored || stream.#inner.state.writeEnded) {
+ return null;
+ }
return stream.#inner.state.writeDesiredSize > 0;
},
writeSync,
@@ -2459,7 +2516,7 @@ class QuicStream {
end,
fail,
[drainableProtocol]() {
- if (closed || errored) return null;
+ if (closed || ending || errored) return null;
// If a drain is already pending, return the existing promise.
if (drainWakeup != null) return drainWakeup.promise;
if (stream.#inner.state.writeDesiredSize > 0) return null;
@@ -2467,6 +2524,7 @@ class QuicStream {
return drainWakeup.promise;
},
[SymbolAsyncDispose]() {
+ if (ending) return pendingEnd;
if (!closed && !errored) fail();
return PromiseResolve();
},
diff --git a/lib/internal/streams/iter/transform.js b/lib/internal/streams/iter/transform.js
index c4f08b8cd32..2cd06957b80 100644
--- a/lib/internal/streams/iter/transform.js
+++ b/lib/internal/streams/iter/transform.js
@@ -12,10 +12,13 @@ const {
ArrayPrototypeMap,
ArrayPrototypePush,
ArrayPrototypeShift,
+ FunctionPrototypeCall,
MathMax,
NumberIsNaN,
ObjectEntries,
ObjectKeys,
+ PromisePrototypeThen,
+ PromiseResolve,
PromiseWithResolvers,
StringPrototypeStartsWith,
SymbolAsyncIterator,
@@ -44,6 +47,7 @@ const {
validateObject,
} = require('internal/validators');
const binding = internalBinding('zlib');
+const { markPromiseAsHandled } = internalBinding('util');
const constants = internalBinding('constants').zlib;
const {
@@ -358,6 +362,14 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
writeAvailIn = availInAfter;
writeAvailOutBefore = chunkSize - outOffset;
+ if (pendingBytes >= BATCH_HWM) {
+ const resolve = resolveWrite;
+ resolveWrite = undefined;
+ rejectWrite = undefined;
+ resolve(false);
+ return;
+ }
+
handle.write(writeFlush,
writeInput, writeInOff, writeAvailIn,
outBuf, outOffset, writeAvailOutBefore);
@@ -369,7 +381,7 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
const resolve = resolveWrite;
resolveWrite = undefined;
rejectWrite = undefined;
- if (resolve) resolve();
+ if (resolve) resolve(true);
}
// onError: called by C++ when the engine encounters an error.
@@ -403,24 +415,28 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
};
signal.addEventListener('abort', onAbort, { __proto__: null, once: true });
- // Dispatch input to the threadpool and return a promise.
- function processInputAsync(input, flushFlag) {
+ function continueInputAsync() {
const { promise, resolve, reject } = PromiseWithResolvers();
resolveWrite = resolve;
rejectWrite = reject;
+ writeAvailOutBefore = chunkSize - outOffset;
+
+ handle.write(writeFlush,
+ writeInput, writeInOff, writeAvailIn,
+ outBuf, outOffset, writeAvailOutBefore);
+ return promise;
+ }
+
+ // Dispatch input to the threadpool and return a promise.
+ function processInputAsync(input, flushFlag) {
writeInput = input;
writeFlush = flushFlag;
writeInOff = 0;
writeAvailIn = TypedArrayPrototypeGetByteLength(input);
- writeAvailOutBefore = chunkSize - outOffset;
// Keep input alive while the threadpool references it.
handle.buffer = input;
-
- handle.write(flushFlag,
- input, 0, writeAvailIn,
- outBuf, outOffset, writeAvailOutBefore);
- return promise;
+ return continueInputAsync();
}
function drainBatch() {
@@ -443,14 +459,42 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
return batch;
}
+ async function* processInputBatches(input, flushFlag) {
+ let complete = await processInputAsync(input, flushFlag);
+ while (!complete) {
+ while (pendingBytes >= BATCH_HWM) {
+ yield drainBatch();
+ }
+ signal?.throwIfAborted();
+ complete = await continueInputAsync();
+ }
+ }
+
let finalized = false;
const iter = source[SymbolAsyncIterator]();
+ const nextMethod = iter.next;
+ let returnMethod;
+ try {
+ returnMethod = iter.return;
+ } catch {
+ // Ignore cleanup-method access failures.
+ }
+
+ function readNext() {
+ const result = PromisePrototypeThen(
+ PromiseResolve(FunctionPrototypeCall(nextMethod, iter)),
+ undefined,
+ undefined);
+ markPromiseAsHandled(result);
+ return result;
+ }
+
try {
// Manually iterate the source so we can pre-read: calling
// iter.next() starts the upstream read + transform on libuv
// before we await the current compression on the threadpool.
- let nextResult = iter.next();
+ let nextResult = readNext();
while (true) {
const { value: chunks, done } = await nextResult;
@@ -462,21 +506,27 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
// Flush signal - finalize the engine.
if (!finalized) {
finalized = true;
- await processInputAsync(kEmpty, finishFlag);
+ for await (const batch of
+ processInputBatches(kEmpty, finishFlag)) {
+ yield batch;
+ }
while (pending.length > 0) {
yield drainBatch();
}
}
- nextResult = iter.next();
+ nextResult = readNext();
continue;
}
// Pre-read: start upstream I/O + transform for the NEXT batch
// while we compress the current batch on the threadpool.
- nextResult = iter.next();
+ nextResult = readNext();
for (let i = 0; i < chunks.length; i++) {
- await processInputAsync(chunks[i], processFlag);
+ for await (const batch of
+ processInputBatches(chunks[i], processFlag)) {
+ yield batch;
+ }
}
if (pendingBytes >= BATCH_HWM) {
@@ -492,7 +542,10 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
// Source ended - finalize if not already done by a null signal.
if (!finalized && !signal.aborted) {
finalized = true;
- await processInputAsync(kEmpty, finishFlag);
+ for await (const batch of
+ processInputBatches(kEmpty, finishFlag)) {
+ yield batch;
+ }
while (pending.length > 0) {
yield drainBatch();
}
@@ -500,9 +553,18 @@ function makeZlibTransform(createHandleFn, processFlag, finishFlag) {
} finally {
signal.removeEventListener('abort', onAbort);
handle.close();
- // Close the upstream iterator so its finally blocks run promptly
- // rather than waiting for GC.
- try { await iter.return?.(); } catch { /* Intentional no-op. */ }
+ // Do not await return(): an async generator queues it behind any
+ // outstanding next(), which is allowed to remain pending forever.
+ if (typeof returnMethod === 'function') {
+ try {
+ const returned = FunctionPrototypeCall(returnMethod, iter);
+ const cleanup = PromisePrototypeThen(
+ PromiseResolve(returned), undefined, undefined);
+ markPromiseAsHandled(cleanup);
+ } catch {
+ // Cleanup errors cannot replace the transform's completion.
+ }
+ }
}
},
};
diff --git a/test/parallel/test-quic-writer-abort-signal.mjs b/test/parallel/test-quic-writer-abort-signal.mjs
index 6b0ba9f68ae..edd3415c415 100644
--- a/test/parallel/test-quic-writer-abort-signal.mjs
+++ b/test/parallel/test-quic-writer-abort-signal.mjs
@@ -1,4 +1,4 @@
-// Flags: --experimental-quic --no-warnings
+// Flags: --experimental-quic --experimental-stream-iter --no-warnings
// Test: write with aborted signal rejects immediately.
@@ -27,7 +27,7 @@ const serverEndpoint = await listen(mustCall((serverSession) => {
const clientSession = await connect(serverEndpoint.address);
await clientSession.opened;
-const stream = await clientSession.createBidirectionalStream();
+const stream = await clientSession.createBidirectionalStream({ budget: 1024 });
const w = stream.writer;
// Create an already-aborted signal.
@@ -39,9 +39,24 @@ await assert.rejects(
{ message: 'already aborted' },
);
-// The writer should still be usable for normal writes.
-w.writeSync(encoder.encode('ok'));
-w.endSync();
+// A pre-aborted end leaves the writer open.
+await assert.rejects(
+ w.end({ signal }),
+ { message: 'already aborted' },
+);
+
+// Once end has started, aborting that call rejects only the operation. The
+// stream still closes after the accepted pending write reaches capacity.
+assert.strictEqual(w.writeSync(new Uint8Array(1024)), true);
+const pendingWrite = w.write(encoder.encode('data'));
+const controller = new AbortController();
+const reason = new Error('end aborted while draining');
+const abortedEnd = w.end({ signal: controller.signal });
+controller.abort(reason);
+
+await assert.rejects(abortedEnd, (error) => error === reason);
+await pendingWrite;
+assert.strictEqual(await w.end(), 1028);
for await (const _ of stream) { /* drain */ } // eslint-disable-line no-unused-vars
await Promise.all([stream.closed, serverDone.promise]);
diff --git a/test/parallel/test-quic-writer-write-rejects.mjs b/test/parallel/test-quic-writer-write-rejects.mjs
index c29c63c075b..b262db692df 100644
--- a/test/parallel/test-quic-writer-write-rejects.mjs
+++ b/test/parallel/test-quic-writer-write-rejects.mjs
@@ -1,8 +1,6 @@
// Flags: --experimental-quic --experimental-stream-iter --no-warnings
-// Test: write() rejects when flow-controlled.
-// The async write() method rejects with ERR_INVALID_STATE when the
-// chunk exceeds capacity (canWrite is false).
+// Test: strict backpressure allows one pending write and rejects the next.
import { hasQuic, skip, mustCall } from '../common/index.mjs';
import assert from 'node:assert';
@@ -41,16 +39,19 @@ assert.strictEqual(w.writeSync(new Uint8Array(1024)), true);
// canWrite should now be false.
assert.strictEqual(w.canWrite, false);
-// Async write() should reject when buffer is full.
+// The first async write waits for capacity.
+const pendingWrite = w.write(new Uint8Array(512));
+
+// A second pending write violates strict backpressure.
await assert.rejects(
- w.write(new Uint8Array(512)),
- { code: 'ERR_INVALID_STATE' },
+ w.writev([new Uint8Array(512)]),
+ { code: 'ERR_INVALID_STATE', name: 'RangeError' },
);
-// Wait for drain, then write should succeed.
+// The drain event admits the pending write.
const drain = w[dp]();
assert.ok(drain instanceof Promise);
-await drain;
+await Promise.all([drain, pendingWrite]);
assert.ok(w.canWrite === true);
// Now write succeeds.
diff --git a/test/parallel/test-stream-iter-transform-buffering.js b/test/parallel/test-stream-iter-transform-buffering.js
new file mode 100644
index 00000000000..7e90479d70c
--- /dev/null
+++ b/test/parallel/test-stream-iter-transform-buffering.js
@@ -0,0 +1,31 @@
+// Flags: --experimental-stream-iter --expose-gc
+'use strict';
+
+const common = require('../common');
+const assert = require('assert');
+const { gzipSync } = require('zlib');
+const { decompressGzip } = require('zlib/iter');
+const { from } = require('stream/iter');
+
+async function testDecompressionOutputIsBounded() {
+ let input = Buffer.alloc(32 * 1024 * 1024, 0x61);
+ const compressed = gzipSync(input);
+ input = null;
+ for (let i = 0; i < 3; i++) globalThis.gc();
+
+ const before = process.memoryUsage().arrayBuffers;
+ const controller = new AbortController();
+ const iterator = decompressGzip().transform(from(compressed), {
+ signal: controller.signal,
+ })[Symbol.asyncIterator]();
+ const first = await iterator.next();
+ for (let i = 0; i < 3; i++) globalThis.gc();
+ const retained = process.memoryUsage().arrayBuffers - before;
+
+ assert.strictEqual(first.done, false);
+ assert.ok(retained < 8 * 1024 * 1024,
+ `decompression retained ${retained} bytes before backpressure`);
+ await iterator.return();
+}
+
+testDecompressionOutputIsBounded().then(common.mustCall());
diff --git a/test/parallel/test-stream-iter-transform-coverage.js b/test/parallel/test-stream-iter-transform-coverage.js
index d10766badd4..7e49c6311dc 100644
--- a/test/parallel/test-stream-iter-transform-coverage.js
+++ b/test/parallel/test-stream-iter-transform-coverage.js
@@ -6,6 +6,7 @@
const common = require('../common');
const assert = require('assert');
+const { setTimeout } = require('timers/promises');
const {
from,
pull,
@@ -55,6 +56,52 @@ async function testEarlyConsumerExit() {
// If we get here without hanging or crashing, cleanup worked
}
+async function testPrefetchedRejectionIsHandled() {
+ const reason = new Error('prefetched read failed');
+ let reads = 0;
+ const source = {
+ __proto__: null,
+ [Symbol.asyncIterator]() { return this; },
+ next() {
+ if (reads++ === 0) {
+ return Promise.resolve({
+ __proto__: null,
+ done: false,
+ value: [Buffer.alloc(1024 * 1024, 0x61)],
+ });
+ }
+ return Promise.reject(reason);
+ },
+ return() { return { __proto__: null, done: true }; },
+ };
+ const controller = new AbortController();
+ const transformed = compressGzip().transform(source, {
+ signal: controller.signal,
+ });
+
+ await assert.rejects(bytes(transformed), (error) => error === reason);
+}
+
+async function testReturnDoesNotAwaitPrefetchedRead() {
+ const never = new Promise(() => {});
+ async function* source() {
+ yield [Buffer.alloc(1024 * 1024, 0x61)];
+ await never;
+ }
+ const controller = new AbortController();
+ const iterator = compressGzip().transform(source(), {
+ signal: controller.signal,
+ })[Symbol.asyncIterator]();
+
+ const first = await iterator.next();
+ assert.strictEqual(first.done, false);
+ const returned = await Promise.race([
+ iterator.return().then(() => true),
+ setTimeout(100, false),
+ ]);
+ assert.strictEqual(returned, true);
+}
+
// Gzip with explicit strategy option
async function testGzipWithStrategy() {
const input = 'strategy test data '.repeat(100);
@@ -110,6 +157,8 @@ async function testInvalidChunkSize() {
Promise.all([
testAbortMidCompression(),
testEarlyConsumerExit(),
+ testPrefetchedRejectionIsHandled(),
+ testReturnDoesNotAwaitPrefetchedRead(),
testGzipWithStrategy(),
testDeflateWithFixedStrategy(),
testBrotliWithParams(),