FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

quic: do not destroy incoming streams that have a consumer · nodejs/node@41c5108 · GitHub

/ node Public

Commit 41c5108

Browse files
authored andcommitted
quic: do not destroy incoming streams that have a consumer
An incoming stream was destroyed unless the session had an onstream callback, even when session-level stream callbacks (onheaders et al) were registered and the negotiated application (HTTP/3) would drive the stream through them. Users had to register stub onstream handlers just to keep their streams alive. Destroy an incoming stream only when the session has no consumer for it at all: no onstream callback, and no session-level stream callbacks runnable on the negotiated application (checked via the existing headersSupported session state, computed when the application is selected from ALPN). Sessions with no consumers keep the current destroy-and-warn behavior so unconsumed streams cannot accumulate and hold flow control credit. On HTTP/3 sessions only bidirectional request streams reach this path; control and QPACK streams are consumed internally by nghttp3 and are never exposed to JavaScript. Fixes: #64192 Signed-off-by: Naman Trivedi <trivenay@amazon.com> PR-URL: #65335 Fixes: #64192 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
1 parent 9ec9383 commit 41c5108

8 files changed

Lines changed: 331 additions & 8 deletions

File tree

‎doc/api/quic.md‎

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides
305305
separate APIs for creating each kind:
306306
[`session.createBidirectionalStream()`][] and
307307
[`session.createUnidirectionalStream()`][]. Streams initiated by a remote
308-
peer are delivered via the [`session.onstream`][] callback.
308+
peer are delivered via the [`session.onstream`][] callback. When the
309+
negotiated application protocol supports the stream-level callbacks (e.g.
310+
HTTP/3) and an `onheaders` callback is configured, incoming streams can
311+
instead be consumed entirely through it and registering `onstream` is
312+
optional.
309313

310314
There are two ways to write data to a stream:
311315

@@ -409,7 +413,9 @@ A typical client session progresses through these stages:
409413

410414
On the server side, call [`quic.listen()`][] with a callback. The callback
411415
fires for each incoming session after the TLS handshake begins. Incoming
412-
streams arrive via the [`session.onstream`][] callback.
416+
streams arrive via the [`session.onstream`][] callback, or, for HTTP/3
417+
sessions with an `onheaders` callback configured, directly through that
418+
callback (see the [minimal HTTP/3 server][] example).
413419

414420
[`session.destroy()`][] is available for immediate teardown — all open streams
415421
are destroyed and the session is closed without waiting for them to finish.
@@ -1108,6 +1114,15 @@ added: v23.8.0
11081114

11091115
The callback to invoke when a new stream is initiated by a remote peer. Read/write.
11101116

1117+
If no `onstream` callback is set and the stream has no other consumer, an
1118+
incoming stream is destroyed on arrival and a warning is emitted. An
1119+
`onheaders` callback counts as a consumer when the negotiated application
1120+
protocol supports it (e.g. HTTP/3), because it is invoked for every incoming
1121+
request stream. Other stream-level callbacks (`ontrailers`, `oninfo`,
1122+
`onwanttrailers`) do not, since they are conditional or outbound-only and
1123+
would leave the stream unobservable. An HTTP/3 server that handles requests
1124+
entirely through `onheaders` does not need to set `onstream`.
1125+
11111126
### `session.ondatagram`
11121127

11131128
<!-- YAML
@@ -3998,7 +4013,9 @@ import { listen } from 'node:quic';
39984013
const encoder = new TextEncoder();
39994014

40004015
const endpoint = await listen((session) => {
4001-
// The session.onstream callback fires for each new client-initiated stream.
4016+
// The session.onstream callback fires for each new client-initiated
4017+
// stream. It is optional here: with `onheaders` configured below,
4018+
// request streams are consumed through that callback.
40024019
}, {
40034020
sni: { '*': { keys: [defaultKey], certs: [defaultCert] } },
40044021
// ALPN defaults to 'h3'.
@@ -4632,5 +4649,6 @@ throughput issues caused by flow control.
46324649
[`stream.writer`]: #streamwriter
46334650
[`writer.fail()`]: #streamwriter
46344651
[`writer.fail(reason)`]: #streamwriter
4652+
[minimal HTTP/3 server]: #minimal-http3-server
46354653
[qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/
46364654
[qvis]: https://qvis.quictools.info/

‎lib/internal/quic/quic.js‎

Lines changed: 33 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4129,6 +4129,24 @@ class QuicSession {
41294129
this.#inner.verifyPeer = value;
41304130
}
41314131

4132+
/**
4133+
* True if an incoming stream has a consumer registered on this session:
4134+
* either an onstream callback, or - when the negotiated application
4135+
* supports headers (e.g. HTTP/3) - session-level stream callbacks that
4136+
* the application layer will invoke (onheaders et al).
4137+
* @returns {boolean}
4138+
*/
4139+
#hasStreamConsumer() {
4140+
if (typeof this.#inner.onstream === 'function') return true;
4141+
// Only onheaders is guaranteed to fire for every incoming stream when
4142+
// the negotiated application supports stream callbacks (e.g. HTTP/3).
4143+
// Other stream callbacks are conditional (ontrailers, oninfo) or
4144+
// outbound-only (onwanttrailers) and do not expose the stream, so they
4145+
// do not count as a consumer.
4146+
if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false;
4147+
return getQuicSessionState(this).streamCallbacksSupported === 1;
4148+
}
4149+
41324150
/**
41334151
* @param {object} handle
41344152
* @param {number} direction
@@ -4141,10 +4159,13 @@ class QuicSession {
41414159
// Set the default byte budget for received streams.
41424160
stream.budget = kDefaultBudget;
41434161

4144-
// A new stream was received. If we don't have an onstream callback, then
4145-
// there's nothing we can do about it. Destroy the stream in this case.
4146-
if (typeof inner.onstream !== 'function') {
4147-
process.emitWarning('A new stream was received but no onstream callback was provided');
4162+
// A new stream was received. If the session has no consumer for it -
4163+
// neither an onstream callback nor, on a session whose application
4164+
// supports headers (e.g. HTTP/3), any session-level stream callbacks -
4165+
// there's nothing that could ever read it. Destroy the stream in this
4166+
// case rather than letting it hold flow control credit.
4167+
if (!this.#hasStreamConsumer()) {
4168+
process.emitWarning('A new stream was received but no stream consumer callback was provided');
41484169
stream.destroy();
41494170
return;
41504171
}
@@ -4175,7 +4196,14 @@ class QuicSession {
41754196
});
41764197
}
41774198

4178-
safeCallbackInvoke(inner.onstream, this, stream);
4199+
// Deliver the stream to the onstream consumer if one is registered.
4200+
// Reaching this point without one means #hasStreamConsumer accepted
4201+
// the stream on behalf of the application layer: the session-level
4202+
// stream callbacks were applied above and the application (e.g.
4203+
// HTTP/3) drives the stream, so there is nothing to invoke here.
4204+
if (typeof inner.onstream === 'function') {
4205+
safeCallbackInvoke(inner.onstream, this, stream);
4206+
}
41794207
}
41804208

41814209
[kRemoveStream](stream) {

‎lib/internal/quic/state.js‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ const {
7272
IDX_STATE_SESSION_STREAM_OPEN_ALLOWED,
7373
IDX_STATE_SESSION_PRIORITY_SUPPORTED,
7474
IDX_STATE_SESSION_HEADERS_SUPPORTED,
75+
IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED,
7576
IDX_STATE_SESSION_WRAPPED,
7677
IDX_STATE_SESSION_APPLICATION_TYPE,
7778
IDX_STATE_SESSION_NO_ERROR_CODE,
@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined);
119120
assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined);
120121
assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined);
121122
assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined);
123+
assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined);
122124
assert(IDX_STATE_SESSION_WRAPPED !== undefined);
123125
assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined);
124126
assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined);
@@ -493,6 +495,19 @@ class QuicSessionState {
493495
return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED);
494496
}
495497

498+
/**
499+
* Whether the negotiated application dispatches the session-level
500+
* stream callbacks (onheaders et al) for incoming streams.
501+
* Returns 0 (unknown), 1 (supported), or 2 (not supported).
502+
* @type {number}
503+
*/
504+
get streamCallbacksSupported() {
505+
const handle = this.#handle;
506+
if (handle === undefined) return undefined;
507+
return DataViewPrototypeGetUint8(
508+
handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
get isWrapped() {
498513
const handle = this.#handle;

‎src/quic/application.h‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer {
210210
// do not support headers should return false (the default).
211211
virtual bool SupportsHeaders() const { return false; }
212212

213+
// True if this application dispatches the session-level stream
214+
// callbacks (onheaders et al) for incoming streams when they are
215+
// registered on the session.
216+
virtual bool SupportsStreamCallbacks() const { return false; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtual void BeginShutdown() {}

‎src/quic/defs.h‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t {
328328
UNSUPPORTED,
329329
};
330330

331+
enum class StreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enum class PathValidationResult : uint8_t {
332338
SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS,
333339
FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE,

‎src/quic/http3.cc‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application {
202202

203203
bool SupportsHeaders() const override { return true; }
204204

205+
bool SupportsStreamCallbacks() const override { return true; }
206+
205207
bool is_started() const override { return started_; }
206208

207209
bool Start() override {

‎src/quic/session.cc‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) {
136136
V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \
137137
V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \
138138
V(HEADERS_SUPPORTED, headers_supported, uint8_t) \
139+
V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \
139140
V(WRAPPED, wrapped, uint8_t) \
140141
V(APPLICATION_TYPE, application_type, uint8_t) \
141142
V(NO_ERROR_CODE, no_error_code, error_code) \
@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) {
26492650
impl_->state()->headers_supported = static_cast<uint8_t>(
26502651
app->SupportsHeaders() ? HeadersSupportState::SUPPORTED
26512652
: HeadersSupportState::UNSUPPORTED);
2653+
impl_->state()->stream_callbacks_supported =
2654+
static_cast<uint8_t>(app->SupportsStreamCallbacks()
2655+
? StreamCallbacksSupportState::SUPPORTED
2656+
: StreamCallbacksSupportState::UNSUPPORTED);
26522657
// Surface the application's "no error" and "internal error" codes via
26532658
// session state so that JS-side code (e.g. the stream writer's fail()
26542659
// path) can resolve the right wire code for the negotiated ALPN

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL