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(),