Skip to content

Commit 168285e

Browse files
committed
quic: fix readable stream truncation on stop-sending, abort & timeout
Signed-off-by: Tim Perry <pimterry@gmail.com>
1 parent 7b6b21a commit 168285e

7 files changed

Lines changed: 153 additions & 83 deletions

File tree

‎lib/internal/blob.js‎

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -615,7 +615,12 @@ function createBlobReaderStream(reader) {
615615
const kMaxBatchChunks = 16;
616616

617617
async function* createBlobReaderIterable(reader, options = kEmptyObject) {
618-
const { getReadError } = options;
618+
// getEndError(status) lets the caller map a read ending to an error, or null.
619+
// It is consulted both on a read failure (status < 0, always an error - we
620+
// fall back to a generic one if the caller returns null) and on a clean EOS
621+
// (status 0, an error only if the caller returns one, e.g. a stream that
622+
// ended before its FIN).
623+
const { getEndError } = options;
619624
let wakeup = PromiseWithResolvers();
620625
let immediate;
621626
let fin = false;
@@ -649,8 +654,8 @@ async function* createBlobReaderIterable(reader, options = kEmptyObject) {
649654
break;
650655
}
651656
if (pullResult.status < 0) {
652-
error = typeof getReadError === 'function' ?
653-
getReadError(pullResult.status) :
657+
error = (typeof getEndError === 'function' &&
658+
getEndError(pullResult.status)) ||
654659
new ERR_INVALID_STATE('The reader is not readable');
655660
break;
656661
}
@@ -666,7 +671,11 @@ async function* createBlobReaderIterable(reader, options = kEmptyObject) {
666671
yield batch;
667672
}
668673

669-
if (eos) return;
674+
if (eos) {
675+
const eosError = typeof getEndError === 'function' ? getEndError(0) : null;
676+
if (eosError) throw eosError;
677+
return;
678+
}
670679
if (error) throw error;
671680

672681
if (blocked) {

‎lib/internal/quic/quic.js‎

Lines changed: 28 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1668,30 +1668,35 @@ class QuicStream {
16681668
// Non-readable stream (outbound-only unidirectional, or closed)
16691669
if (!inner.reader) return;
16701670

1671-
yield* createBlobReaderIterable(inner.reader, {
1672-
getReadError: () => {
1673-
// The read side ends for one of three reasons:
1674-
// * Clean FIN received from the peer (state.finReceived
1675-
// === true). The iterator stops without calling this;
1676-
// fall through to the generic state error if it does.
1677-
// * Peer sent us a RESET_STREAM. The C++ side records the
1678-
// code in state.resetCode regardless of whether the JS
1679-
// onreset handler was attached. state.finReceived stays
1680-
// false because no FIN was seen.
1681-
// * We aborted locally via stream.resetStream() or
1682-
// stream.stopSending(). Both paths run EndReadable in
1683-
// C++, setting state.readEnded without setting
1684-
// state.finReceived. There is no peer code to surface.
1685-
if (inner.state.readEnded && !inner.state.finReceived) {
1686-
const peerResetCode = inner.state.resetCode;
1687-
if (peerResetCode !== undefined && peerResetCode > 0n) {
1688-
return new ERR_QUIC_STREAM_RESET(Number(peerResetCode));
1689-
}
1690-
return new ERR_QUIC_STREAM_ABORTED(
1691-
'Stream aborted before FIN was received');
1671+
// Maps the read side's end state to the error to surface, or null if it
1672+
// ended cleanly. state.finReceived is set only when the peer explicitly
1673+
// sent a FIN, confirming we got the whole stream - the one and only clean
1674+
// ending. Returns null in that case. Without a FIN the read side is
1675+
// truncated, for one of these reasons:
1676+
// * Peer sent us a RESET_STREAM. The C++ side records the
1677+
// code in state.resetCode regardless of whether the JS
1678+
// onreset handler was attached. state.finReceived stays
1679+
// false because no FIN was seen.
1680+
// * We aborted locally via stream.stopSending(). This runs EndReadable
1681+
// in C++, setting state.readEnded but not state.finReceived. There
1682+
// is no peer code to surface.
1683+
// * The session was torn down before a FIN arrived (a local or peer
1684+
// connection close, an idle timeout, or any error), again leaving
1685+
// state.finReceived false.
1686+
const readTruncationError = () => {
1687+
if (inner.state.readEnded && !inner.state.finReceived) {
1688+
const peerResetCode = inner.state.resetCode;
1689+
if (peerResetCode !== undefined && peerResetCode > 0n) {
1690+
return new ERR_QUIC_STREAM_RESET(Number(peerResetCode));
16921691
}
1693-
return new ERR_INVALID_STATE('The stream is not readable');
1694-
},
1692+
return new ERR_QUIC_STREAM_ABORTED(
1693+
'Stream aborted before FIN was received');
1694+
}
1695+
return null;
1696+
};
1697+
1698+
yield* createBlobReaderIterable(inner.reader, {
1699+
getEndError: readTruncationError,
16951700
});
16961701
}
16971702

‎src/quic/streams.cc‎

Lines changed: 17 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1405,13 +1405,6 @@ BaseObjectPtr<Blob::Reader> Stream::get_reader() {
14051405
return reader;
14061406
}
14071407

1408-
void Stream::set_final_size(uint64_t final_size) {
1409-
DCHECK_IMPLIES(state()->fin_received == 1,
1410-
final_size <= STAT_GET(Stats, final_size));
1411-
state()->fin_received = 1;
1412-
STAT_SET(Stats, final_size, final_size);
1413-
}
1414-
14151408
void Stream::set_outbound(std::shared_ptr<DataQueue> source) {
14161409
if (!source || !is_writable()) return;
14171410
Debug(this, "Setting the outbound data source");
@@ -1607,13 +1600,20 @@ void Stream::EndWritable() {
16071600
state()->write_ended = 1;
16081601
}
16091602

1610-
void Stream::EndReadable(std::optional<uint64_t> maybe_final_size) {
1603+
void Stream::EndReadable(std::optional<uint64_t> maybe_final_size,
1604+
bool clean_fin) {
16111605
if (!is_readable()) return;
16121606
state()->read_ended = 1;
1607+
// fin_received marks a clean completion of the read side (a real FIN). Any
1608+
// unclean end (reset/abort/session teardown) truncates the stream, which
1609+
// the JS reader will expose as a read error later.
1610+
if (clean_fin) state()->fin_received = 1;
16131611
// Flush any accumulated data before capping so the reader can see it.
16141612
FlushAccumulation();
1615-
set_final_size(maybe_final_size.value_or(STAT_GET(Stats, bytes_received)));
1616-
inbound_->cap(STAT_GET(Stats, final_size));
1613+
const uint64_t final_size =
1614+
maybe_final_size.value_or(STAT_GET(Stats, bytes_received));
1615+
STAT_SET(Stats, final_size, final_size);
1616+
inbound_->cap(final_size);
16171617
// Notify the JS reader so it can see EOS. Pass fin=true so the
16181618
// wakeup promise resolves with a value the iterator can check to
16191619
// avoid waiting for another wakeup that will never come.
@@ -1642,8 +1642,9 @@ void Stream::Destroy(QuicError error) {
16421642
// End the writable before marking as destroyed.
16431643
EndWritable();
16441644

1645-
// Also end the readable side if it isn't already.
1646-
EndReadable();
1645+
// Also end the readable side if it isn't already. If not already ended,
1646+
// this will eventually surface as a error, since the data is truncated.
1647+
EndReadable(std::nullopt, /* clean_fin = */ false);
16471648

16481649
// We are going to release our reference to the outbound_ queue here.
16491650
outbound_.reset();
@@ -1690,7 +1691,7 @@ void Stream::ReceiveData(const uint8_t* data,
16901691
// end the readable side if this is the last bit of data we've received.
16911692
Debug(this, "Receiving %zu bytes of data", len);
16921693
if (state()->read_ended == 1 || len == 0) {
1693-
if (flags.fin) EndReadable();
1694+
if (flags.fin) EndReadable(std::nullopt, /* clean_fin = */ true);
16941695
return;
16951696
}
16961697

@@ -1763,7 +1764,7 @@ void Stream::ReceiveData(const uint8_t* data,
17631764

17641765
if (flags.fin) {
17651766
FlushAccumulation();
1766-
EndReadable();
1767+
EndReadable(std::nullopt, /* clean_fin = */ true);
17671768
} else if (reader_ && was_empty) {
17681769
// Notify the reader once when the accumulator transitions from empty
17691770
// to non-empty. This wakes the reader exactly once per accumulation
@@ -1794,7 +1795,7 @@ void Stream::ReceiveStreamReset(uint64_t final_size, QuicError error) {
17941795
final_size,
17951796
error);
17961797
state()->reset_code = error.code();
1797-
EndReadable(final_size);
1798+
EndReadable(final_size, /* clean_fin = */ false);
17981799
EmitReset(error);
17991800
}
18001801

@@ -1817,7 +1818,7 @@ void Stream::DoStreamReset(error_code code) {
18171818
}
18181819

18191820
void Stream::SendStopSending(error_code code) {
1820-
EndReadable();
1821+
EndReadable(std::nullopt, /* clean_fin = */ false);
18211822

18221823
if (!is_pending()) {
18231824
// If the stream is a local unidirectional there's nothing to do here.

‎src/quic/streams.h‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -309,7 +309,10 @@ class Stream final : public AsyncWrap,
309309
void Commit(size_t datalen, bool fin = false);
310310

311311
void EndWritable();
312-
void EndReadable(std::optional<uint64_t> maybe_final_size = std::nullopt);
312+
// clean_fin indicates the read side is ending because a real FIN frame was
313+
// received from the peer (as opposed to a reset, a local abort, or the
314+
// session being torn down). Anything else => truncated read.
315+
void EndReadable(std::optional<uint64_t> maybe_final_size, bool clean_fin);
313316
void EntryRead(size_t amount) override;
314317
void BeforePull() override;
315318

@@ -398,7 +401,6 @@ class Stream final : public AsyncWrap,
398401
// Gets a reader for the data received for this stream from the peer,
399402
BaseObjectPtr<Blob::Reader> get_reader();
400403

401-
void set_final_size(uint64_t amount);
402404
void set_outbound(std::shared_ptr<DataQueue> source);
403405

404406
// Streaming outbound support
Lines changed: 26 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -1,64 +1,53 @@
11
// Flags: --experimental-quic --experimental-stream-iter --no-warnings
22

3-
// Test: peer RESET_STREAM causes iterator to error.
4-
// When the server resets the stream, the client's async iterator
5-
// should throw or return early.
3+
// Test: a peer RESET_STREAM truncates the readable. The async iterator
4+
// delivers the data received before the reset, then throws
5+
// ERR_QUIC_STREAM_RESET (carrying the peer's code) at the end - rather than
6+
// ending cleanly.
67

78
import { hasQuic, skip, mustCall } from '../common/index.mjs';
8-
import * as assert from 'node:assert';
9+
import { setTimeout as delay } from 'node:timers/promises';
10+
import assert from 'node:assert';
911

1012
if (!hasQuic) {
1113
skip('QUIC is not enabled');
1214
}
1315

1416
const { listen, connect } = await import('../common/quic.mjs');
1517

16-
const encoder = new TextEncoder();
17-
18-
const serverReady = Promise.withResolvers();
19-
2018
const serverEndpoint = await listen(mustCall((serverSession) => {
2119
serverSession.onstream = mustCall(async (stream) => {
22-
// Reset the stream from the server side.
20+
stream.writer.write(new Uint8Array(1000).fill(7));
21+
while (stream.stats.maxOffsetAcknowledged < 1000n) await delay(5);
2322
stream.resetStream(42n);
24-
await assert.rejects(stream.closed, mustCall((err) => {
25-
assert.ok(err);
26-
return true;
27-
}));
28-
serverReady.resolve();
29-
await serverSession.closed;
23+
stream.closed.catch(() => {});
3024
});
31-
}), { transportParams: { maxIdleTimeout: 1 } });
25+
}));
3226

3327
const clientSession = await connect(serverEndpoint.address, {
3428
transportParams: { maxIdleTimeout: 1 },
3529
});
3630
await clientSession.opened;
3731

38-
const stream = await clientSession.createBidirectionalStream({
39-
body: encoder.encode('will be reset by server'),
40-
});
41-
42-
// Set up the closed handler before the reset to avoid unhandled rejection.
43-
const closedPromise = assert.rejects(stream.closed, mustCall((err) => {
44-
assert.ok(err);
45-
return true;
46-
}));
47-
48-
await serverReady.promise;
32+
// Keep our write side open so the stream stays alive while we read.
33+
const stream = await clientSession.createBidirectionalStream();
34+
await stream.writer.write(new Uint8Array([1]));
35+
stream.closed.catch(() => {});
4936

50-
// The async iterator should either throw or return early when the
51-
// peer resets the readable side.
37+
let received = 0;
38+
let threw;
5239
try {
53-
for await (const batch of stream) {
54-
// May receive some data before the reset arrives.
55-
assert.ok(Array.isArray(batch));
40+
for await (const chunk of stream) {
41+
for (const c of chunk) received += c.byteLength;
5642
}
57-
} catch {
58-
// The iterator may throw when the reset arrives mid-iteration.
43+
} catch (err) {
44+
threw = err;
5945
}
6046

61-
// Either way, the stream should close.
62-
await closedPromise;
63-
await clientSession.closed;
47+
// The buffered data was delivered before the error.
48+
assert.strictEqual(received, 1000);
49+
// The reset surfaced as a reset error (with its code), not a clean end.
50+
assert.strictEqual(threw?.code, 'ERR_QUIC_STREAM_RESET');
51+
52+
clientSession.close();
6453
await serverEndpoint.close();
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
// Flags: --experimental-quic --experimental-stream-iter --no-warnings
2+
3+
// Test: a readable stream truncated by the connection idle timeout delivers
4+
// the data it received, then surfaces the truncation as an error at the end of
5+
// iteration - rather than a silent clean end-of-stream that would make an
6+
// incomplete stream look complete.
7+
8+
import { hasQuic, skip, mustCall } from '../common/index.mjs';
9+
import assert from 'node:assert';
10+
11+
const { strictEqual } = assert;
12+
13+
if (!hasQuic) {
14+
skip('QUIC is not enabled');
15+
}
16+
17+
const { listen, connect } = await import('../common/quic.mjs');
18+
19+
// The body sends 1000 bytes then hangs (never a FIN), so the only thing that
20+
// ends the client's read side is the connection idle timeout.
21+
async function* stallingBody() {
22+
yield new Uint8Array(1000).fill(7);
23+
await new Promise(() => {});
24+
}
25+
26+
const serverEndpoint = await listen(mustCall((serverSession) => {
27+
serverSession.onstream = mustCall((stream) => {
28+
stream.setBody(stallingBody());
29+
stream.closed.catch(() => {});
30+
});
31+
}));
32+
33+
const clientSession = await connect(serverEndpoint.address, {
34+
// Short connection idle timeout so the truncation happens quickly.
35+
transportParams: { maxIdleTimeout: 1 },
36+
});
37+
await clientSession.opened;
38+
39+
const stream = await clientSession.createBidirectionalStream();
40+
await stream.writer.write(new Uint8Array([1]));
41+
stream.closed.catch(() => {});
42+
43+
let received = 0;
44+
let threw;
45+
try {
46+
for await (const chunk of stream) {
47+
for (const c of chunk) received += c.byteLength;
48+
}
49+
} catch (err) {
50+
threw = err;
51+
}
52+
53+
// All the buffered data was delivered before the error.
54+
strictEqual(received, 1000);
55+
// The truncation surfaced as an error at the end, not a clean end-of-stream.
56+
strictEqual(threw?.code, 'ERR_QUIC_STREAM_ABORTED');
57+
58+
await serverEndpoint.close();

‎test/parallel/test-quic-stream-setbody-errors.mjs‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,13 @@ await clientSession.opened;
5353
message: /writer already accessed/,
5454
});
5555

56-
for await (const _ of stream) { /* drain */ } // eslint-disable-line no-unused-vars
56+
// The server handles only the first stream and then closes its session, so
57+
// this stream is never answered and never receives a FIN. Reading it
58+
// therefore surfaces the truncation rather than ending cleanly.
59+
await assert.rejects((async () => {
60+
// eslint-disable-next-line no-unused-vars
61+
for await (const _ of stream) { /* drain */ }
62+
})(), { code: 'ERR_QUIC_STREAM_ABORTED' });
5763
await stream.closed;
5864
}
5965

0 commit comments

Comments
 (0)