From 2526c29286eb7ac14b88ab1d9fc810729c2630a2 Mon Sep 17 00:00:00 2001 From: GiHoon1123 Date: Wed, 30 Sep 2026 16:58:30 +0900 Subject: [PATCH] quic: retain readers for readable streams Create readers when readable QuicStream instances are constructed so data received before the first async iterator pull survives native stream cleanup. Add a regression test for a peer-initiated unidirectional stream that sends data and FIN before the reader is pulled. Assisted-by: Codex Signed-off-by: GiHoon1123 --- lib/internal/quic/quic.js | 6 +++ .../test-quic-stream-uni-delayed-reader.mjs | 47 +++++++++++++++++++ 2 files changed, 53 insertions(+) create mode 100644 test/parallel/test-quic-stream-uni-delayed-reader.mjs diff --git a/lib/internal/quic/quic.js b/lib/internal/quic/quic.js index 4922ce562751..88a85c7428c7 100644 --- a/lib/internal/quic/quic.js +++ b/lib/internal/quic/quic.js @@ -1659,6 +1659,12 @@ class QuicStream { this.#handle = handle; handle[kOwner] = this; const inner = this.#inner; + // Keep a reader alive for readable streams before the native handle can + // be destroyed. A peer may send data and FIN before the consumer asks + // for its async iterator. + if (!isLocal || direction === kStreamDirectionBidirectional) { + inner.reader = handle.getReader(); + } inner.session = session; inner.direction = direction; inner.isLocal = isLocal; diff --git a/test/parallel/test-quic-stream-uni-delayed-reader.mjs b/test/parallel/test-quic-stream-uni-delayed-reader.mjs new file mode 100644 index 000000000000..45a233617161 --- /dev/null +++ b/test/parallel/test-quic-stream-uni-delayed-reader.mjs @@ -0,0 +1,47 @@ +// Flags: --experimental-quic --experimental-stream-iter --no-warnings + +// Test: data received before the async iterator is pulled remains readable +// after a peer-initiated unidirectional stream closes. + +import { hasQuic, skip, mustCall } from '../common/index.mjs'; +import assert from 'node:assert'; + +if (!hasQuic) { + skip('QUIC is not enabled'); +} + +const { listen, connect } = await import('../common/quic.mjs'); +const { bytes } = await import('stream/iter'); + +const encoder = new TextEncoder(); +const expected = encoder.encode('data before the first pull'); +const done = Promise.withResolvers(); + +const serverEndpoint = await listen(mustCall(async (serverSession) => { + await serverSession.opened; + const stream = await serverSession.createUnidirectionalStream({ + body: expected, + }); + await stream.closed; + serverSession.close(); +})); + +const clientSession = await connect(serverEndpoint.address); +await clientSession.opened; + +clientSession.onstream = mustCall(async (stream) => { + const iterator = stream[Symbol.asyncIterator](); + + // Delay the first pull until after the peer has sent FIN and the native + // stream handle has been closed. + await stream.closed; + + const received = await bytes(iterator); + assert.deepStrictEqual(received, expected); + clientSession.close(); + done.resolve(); +}); + +await done.promise; +await clientSession.closed; +await serverEndpoint.close();