Skip to content

Commit b4450fe

Browse files
committed
Make QUIC stream truncation behaviour configurable
1 parent 168285e commit b4450fe

13 files changed

Lines changed: 474 additions & 231 deletions

‎doc/api/quic.md‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3305,6 +3305,31 @@ value, PING frames will be sent automatically to keep the connection alive
33053305
before the idle timeout fires. The value should be less than the effective
33063306
idle timeout (`maxIdleTimeout` transport parameter) to be useful.
33073307

3308+
#### `sessionOptions.truncatedReads`
3309+
3310+
* Type: {string} One of `'error'` or `'allow'`.
3311+
* **Default:** `'error'`
3312+
3313+
Controls how reading a stream reports a truncated read. A stream's read side
3314+
can end without receiving a QUIC FIN, meaning the peer never signalled that
3315+
the whole stream had been sent and the data received may be incomplete. This
3316+
selects how the stream's async iterator reports this:
3317+
3318+
* `'error'` - The default. Peers are expected to always send a FIN to end
3319+
their data explicitly, and so any truncation is an error. The iterator yields
3320+
the data that did arrive and then throws, so an incomplete stream can never
3321+
be mistaken for a complete one. Incomplete streams will either throw a
3322+
`ERR_QUIC_STREAM_RESET` carrying the peer's error code, a connection error,
3323+
or `ERR_QUIC_STREAM_ABORTED` for other cases.
3324+
3325+
* `'allow'` - Truncated reads are allowed: only a stream or connection error
3326+
is reported, and any clean abort/cancellation or similar simply ends the
3327+
stream. A non-zero peer reset still throws `ERR_QUIC_STREAM_RESET` and a
3328+
connection error still throws its real error, but a truncation that carried
3329+
no error (an idle timeout, a graceful close, a local `stopSending()`) ends
3330+
the read cleanly with the data received. This matches `stream.closed`,
3331+
which rejects only on an error.
3332+
33083333
#### `sessionOptions.verifyPeer` (client only)
33093334

33103335
* Type: {string} One of `'strict'`, `'auto'`, or `'manual'`.

‎lib/internal/blob.js‎

Lines changed: 3 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -614,13 +614,7 @@ function createBlobReaderStream(reader) {
614614
// unbounded memory growth when the DataQueue has a large burst of data.
615615
const kMaxBatchChunks = 16;
616616

617-
async function* createBlobReaderIterable(reader, options = kEmptyObject) {
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;
617+
async function* createBlobReaderIterable(reader) {
624618
let wakeup = PromiseWithResolvers();
625619
let immediate;
626620
let fin = false;
@@ -654,9 +648,7 @@ async function* createBlobReaderIterable(reader, options = kEmptyObject) {
654648
break;
655649
}
656650
if (pullResult.status < 0) {
657-
error = (typeof getEndError === 'function' &&
658-
getEndError(pullResult.status)) ||
659-
new ERR_INVALID_STATE('The reader is not readable');
651+
error = new ERR_INVALID_STATE('The reader is not readable');
660652
break;
661653
}
662654
if (pullResult.status === 2) {
@@ -671,11 +663,7 @@ async function* createBlobReaderIterable(reader, options = kEmptyObject) {
671663
yield batch;
672664
}
673665

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

681669
if (blocked) {

‎lib/internal/quic/quic.js‎

Lines changed: 71 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,6 @@ const {
1313
ErrorCaptureStackTrace,
1414
FunctionPrototypeBind,
1515
FunctionPrototypeCall,
16-
Number,
1716
ObjectDefineProperties,
1817
ObjectKeys,
1918
PromisePrototypeThen,
@@ -114,8 +113,6 @@ const {
114113
ERR_QUIC_CONNECTION_FAILED,
115114
ERR_QUIC_ENDPOINT_CLOSED,
116115
ERR_QUIC_OPEN_STREAM_FAILED,
117-
ERR_QUIC_STREAM_ABORTED,
118-
ERR_QUIC_STREAM_RESET,
119116
ERR_QUIC_VERSION_NEGOTIATION_ERROR,
120117
},
121118
} = require('internal/errors');
@@ -398,6 +395,7 @@ const endpointRegistry = new SafeSet();
398395
* @property {number} [minVersion] The minimum acceptable QUIC version
399396
* @property {'use'|'ignore'|'default'} [preferredAddressPolicy] The preferred address policy
400397
* @property {'strict'|'auto'|'manual'} [verifyPeer='auto'] Peer certificate verification policy (client only)
398+
* @property {'error'|'allow'} [truncatedReads] Truncated read policy
401399
* @property {ApplicationOptions} [application] The application options
402400
* @property {TransportParams} [transportParams] The transport parameters
403401
* @property {string} [servername] The server name identifier (client only)
@@ -1580,6 +1578,8 @@ class QuicStream {
15801578
state: undefined,
15811579
stats: undefined,
15821580
pendingClose: undefined,
1581+
destroyError: undefined,
1582+
truncatedReads: undefined,
15831583
reader: undefined,
15841584
destroying: false,
15851585
iteratorLocked: false,
@@ -1627,9 +1627,10 @@ class QuicStream {
16271627
* @param {object} handle
16281628
* @param {QuicSession} session
16291629
* @param {number} direction
1630-
* @param {boolean} [isLocal]
1630+
* @param {boolean} isLocal
1631+
* @param {'error'|'allow'} truncatedReads
16311632
*/
1632-
constructor(privateSymbol, handle, session, direction, isLocal) {
1633+
constructor(privateSymbol, handle, session, direction, isLocal, truncatedReads) {
16331634
assertPrivateSymbol(privateSymbol);
16341635

16351636
this.#handle = handle;
@@ -1638,6 +1639,7 @@ class QuicStream {
16381639
inner.session = session;
16391640
inner.direction = direction;
16401641
inner.isLocal = isLocal;
1642+
inner.truncatedReads = truncatedReads;
16411643
inner.state = new QuicStreamState(
16421644
kPrivateConstructor, handle.state, handle.stateByteOffset);
16431645

@@ -1665,39 +1667,41 @@ class QuicStream {
16651667
inner.iteratorLocked = true;
16661668

16671669
inner.reader ??= this.#handle?.getReader();
1668-
// Non-readable stream (outbound-only unidirectional, or closed)
1669-
if (!inner.reader) return;
1670-
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));
1691-
}
1692-
return new ERR_QUIC_STREAM_ABORTED(
1693-
'Stream aborted before FIN was received');
1670+
// No reader means either a outbound-only unidirectional stream, or a
1671+
// stream already destroyed (data gone, but truncation must still be
1672+
// checked below).
1673+
if (inner.reader) {
1674+
yield* createBlobReaderIterable(inner.reader);
1675+
}
1676+
1677+
if (inner.state.readEnded && !inner.state.finReceived) {
1678+
// The readable has been truncated - ended with no clean FIN. We expose
1679+
// this in different ways depending on the truncatedReads option.
1680+
1681+
// Non-zero reset is always an error:
1682+
const peerResetCode = inner.state.resetCode;
1683+
if (peerResetCode > 0n) {
1684+
throw new QuicError(
1685+
`The QUIC stream was reset by the peer with error code ${peerResetCode}`,
1686+
{ __proto__: null,
1687+
code: 'ERR_QUIC_STREAM_RESET',
1688+
errorCode: peerResetCode });
16941689
}
1695-
return null;
1696-
};
16971690

1698-
yield* createBlobReaderIterable(inner.reader, {
1699-
getEndError: readTruncationError,
1700-
});
1691+
// If stream teardown has started (stats.destroyedAt is set) then a
1692+
// close event confirming a final error/clean close will settle
1693+
// imminently (might be settled already). We await here to rethrow
1694+
// any errors if the connection has failed.
1695+
if (this.destroyed || this.stats.destroyedAt !== 0n) {
1696+
await this.closed;
1697+
}
1698+
1699+
// Clean abort is truncation, but not necessarily an error:
1700+
if (inner.truncatedReads === 'error') {
1701+
throw new QuicError('Stream aborted before FIN was received',
1702+
{ __proto__: null, errorCode: peerResetCode ?? 0n });
1703+
}
1704+
}
17011705
}
17021706

17031707
/**
@@ -2068,9 +2072,10 @@ class QuicStream {
20682072
if (error !== undefined && typeof inner.onerror === 'function') {
20692073
invokeOnerror(inner.onerror, error);
20702074
}
2071-
const handle = this.#handle;
2072-
this[kFinishClose](error);
2073-
handle.destroy();
2075+
// handle.destroy() kicks off all the cleanup internals, eventually
2076+
// including [kFinishClose] which needs this original destroy error:
2077+
inner.destroyError = error;
2078+
this.#handle.destroy();
20742079
}
20752080

20762081
/**
@@ -2566,6 +2571,9 @@ class QuicStream {
25662571
if (this.destroyed) {
25672572
return inner.pendingClose.promise;
25682573
}
2574+
// Prefer an error staged by destroy() (the original object the caller
2575+
// passed) over the error delivered by the native close callback.
2576+
error = inner.destroyError ?? error;
25692577
if (error !== undefined) {
25702578
inner.pendingClose.reject(error);
25712579
} else {
@@ -2788,6 +2796,7 @@ class QuicSession {
27882796
// because server-side cert validation is handled by rejectUnauthorized
27892797
// at the C++ level.
27902798
verifyPeer: 'manual',
2799+
truncatedReads: 'error',
27912800
handshakeInfo: undefined,
27922801
/** @type {QuicSessionPath|undefined} */
27932802
path: undefined,
@@ -2821,8 +2830,9 @@ class QuicSession {
28212830
* @param {symbol} privateSymbol
28222831
* @param {object} handle
28232832
* @param {QuicEndpoint} endpoint
2833+
* @param {{ truncatedReads?: 'error'|'allow' }} [options]
28242834
*/
2825-
constructor(privateSymbol, handle, endpoint) {
2835+
constructor(privateSymbol, handle, endpoint, options = kEmptyObject) {
28262836
// Instances of QuicSession can only be created internally.
28272837
assertPrivateSymbol(privateSymbol);
28282838

@@ -2831,6 +2841,8 @@ class QuicSession {
28312841

28322842
const inner = this.#inner;
28332843
inner.endpoint = endpoint;
2844+
const { truncatedReads } = options;
2845+
if (truncatedReads !== undefined) inner.truncatedReads = truncatedReads;
28342846
// Move any qlog entries that arrived before the wrapper existed.
28352847
if (handle._pendingQlog !== undefined) {
28362848
inner.pendingQlog = handle._pendingQlog;
@@ -3376,7 +3388,8 @@ class QuicSession {
33763388
}
33773389

33783390
const stream = new QuicStream(
3379-
kPrivateConstructor, handle, this, direction, true /* isLocal */);
3391+
kPrivateConstructor, handle, this, direction, true /* isLocal */,
3392+
inner.truncatedReads);
33803393
inner.streams.add(stream);
33813394
if (typeof this.#inner.onerror === 'function') {
33823395
markPromiseAsHandled(stream.closed);
@@ -4161,7 +4174,7 @@ class QuicSession {
41614174
[kNewStream](handle, direction) {
41624175
const inner = this.#inner;
41634176
const stream = new QuicStream(kPrivateConstructor, handle, this, direction,
4164-
false /* isLocal */);
4177+
false /* isLocal */, inner.truncatedReads);
41654178

41664179
// Set the default byte budget for received streams.
41674180
stream.budget = kDefaultBudget;
@@ -4278,6 +4291,7 @@ class QuicEndpoint {
42784291
sessions: new SafeSet(),
42794292
stat: undefined,
42804293
stats: undefined,
4294+
truncatedReads: undefined,
42814295
onsession: undefined,
42824296
sessionCallbacks: undefined,
42834297
};
@@ -4410,8 +4424,8 @@ class QuicEndpoint {
44104424
};
44114425
}
44124426

4413-
#newSession(handle) {
4414-
const session = new QuicSession(kPrivateConstructor, handle, this);
4427+
#newSession(handle, options) {
4428+
const session = new QuicSession(kPrivateConstructor, handle, this, options);
44154429
this.#inner.sessions.add(session);
44164430
// Set default pending datagram queue size.
44174431
session.maxPendingDatagrams = kDefaultMaxPendingDatagrams;
@@ -4585,9 +4599,13 @@ class QuicEndpoint {
45854599
ontrailers,
45864600
oninfo,
45874601
onwanttrailers,
4602+
// Stored on the endpoint and applied to each incoming session.
4603+
truncatedReads,
45884604
...rest
45894605
} = options;
45904606

4607+
inner.truncatedReads = truncatedReads;
4608+
45914609
// Store session and stream callbacks to apply to each new incoming session.
45924610
inner.sessionCallbacks = {
45934611
__proto__: null,
@@ -4629,6 +4647,7 @@ class QuicEndpoint {
46294647
validateObject(options, 'options');
46304648
const {
46314649
sessionTicket,
4650+
truncatedReads,
46324651
...rest
46334652
} = options;
46344653

@@ -4637,7 +4656,7 @@ class QuicEndpoint {
46374656
if (handle === undefined) {
46384657
throw new ERR_QUIC_CONNECTION_FAILED();
46394658
}
4640-
const session = this.#newSession(handle);
4659+
const session = this.#newSession(handle, { __proto__: null, truncatedReads });
46414660
// Set callbacks before any async work to avoid missing events
46424661
// that fire during or immediately after the handshake.
46434662
applyCallbacks(session, options);
@@ -4871,7 +4890,8 @@ class QuicEndpoint {
48714890
const inner = this.#inner;
48724891
assert(typeof inner.onsession === 'function',
48734892
'onsession callback not specified');
4874-
const session = this.#newSession(handle);
4893+
const session = this.#newSession(handle,
4894+
{ __proto__: null, truncatedReads: inner.truncatedReads });
48754895
// Apply session callbacks stored at listen time before notifying
48764896
// the onsession callback, to avoid missing events that fire
48774897
// during or immediately after the handshake.
@@ -5357,6 +5377,7 @@ function processSessionOptions(options, config = kEmptyObject) {
53575377
maxDatagramSendAttempts = 5,
53585378
streamIdleTimeout,
53595379
verifyPeer = 'auto',
5380+
truncatedReads = 'error',
53605381
// HTTP/3 application-specific options. Nested under `application`
53615382
// to separate protocol-specific settings from transport-level ones.
53625383
application = kEmptyObject,
@@ -5408,6 +5429,9 @@ function processSessionOptions(options, config = kEmptyObject) {
54085429
validateOneOf(verifyPeer, 'options.verifyPeer',
54095430
['strict', 'auto', 'manual']);
54105431

5432+
validateOneOf(truncatedReads, 'options.truncatedReads',
5433+
['error', 'allow']);
5434+
54115435
validateInteger(drainingPeriodMultiplier, 'options.drainingPeriodMultiplier',
54125436
3, 255);
54135437

@@ -5470,6 +5494,7 @@ function processSessionOptions(options, config = kEmptyObject) {
54705494
verifyHostname: verifyPeer !== 'manual',
54715495
},
54725496
verifyPeer,
5497+
truncatedReads,
54735498
qlog,
54745499
maxPayloadSize,
54755500
unacknowledgedPacketThreshold,

0 commit comments

Comments
 (0)