Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions lib/internal/quic/quic.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
47 changes: 47 additions & 0 deletions test/parallel/test-quic-stream-uni-delayed-reader.mjs
Original file line number Diff line number Diff line change
@@ -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();
Loading