Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Commit 41c5108

Browse files
trivenayaduh95
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: #65335Fixes: #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
constencoder=newTextEncoder();
39994014

40004015
constendpoint=awaitlisten((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(typeofthis.#inner.onstream==='function')returntrue;
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(typeofthis[kStreamCallbacks]?.onheaders!=='function')returnfalse;
4147+
returngetQuicSessionState(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(typeofinner.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(typeofinner.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
returnDataViewPrototypeGetUint8(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+
getstreamCallbacksSupported(){
505+
consthandle=this.#handle;
506+
if(handle===undefined)returnundefined;
507+
returnDataViewPrototypeGetUint8(
508+
handle,this.#offset +IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED);
509+
}
510+
496511
/** @type {boolean} */
497512
getisWrapped(){
498513
consthandle=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
virtualboolSupportsHeaders() const { returnfalse; }
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+
virtualboolSupportsStreamCallbacks() const { returnfalse; }
217+
213218
// Initiates application-level graceful shutdown signaling (e.g.,
214219
// HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked.
215220
virtualvoidBeginShutdown() {}

‎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+
enumclassStreamCallbacksSupportState : uint8_t {
332+
UNKNOWN,
333+
SUPPORTED,
334+
UNSUPPORTED,
335+
};
336+
331337
enumclassPathValidationResult : 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
boolSupportsHeaders() constoverride { returntrue; }
204204

205+
boolSupportsStreamCallbacks() constoverride { returntrue; }
206+
205207
boolis_started() constoverride { return started_; }
206208

207209
boolStart() 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)