| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 9ec9383 commit 41c5108
8 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -305,7 +305,11 @@ unidirectional (data flows in only one direction). The `quic` module provides | |||
| 305 | 305 | separate APIs for creating each kind: | |
| 306 | 306 | [`session.createBidirectionalStream()`][] and | |
| 307 | 307 | [`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. | ||
| 309 | 313 | ||
| 310 | 314 | There are two ways to write data to a stream: | |
| 311 | 315 | ||
@@ -409,7 +413,9 @@ A typical client session progresses through these stages: | |||
| 409 | 413 | ||
| 410 | 414 | On the server side, call [`quic.listen()`][] with a callback. The callback | |
| 411 | 415 | 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). | ||
| 413 | 419 | ||
| 414 | 420 | [`session.destroy()`][] is available for immediate teardown — all open streams | |
| 415 | 421 | are destroyed and the session is closed without waiting for them to finish. | |
@@ -1108,6 +1114,15 @@ added: v23.8.0 | |||
| 1108 | 1114 | ||
| 1109 | 1115 | The callback to invoke when a new stream is initiated by a remote peer. Read/write. | |
| 1110 | 1116 | ||
| 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 | + | ||
| 1111 | 1126 | ### `session.ondatagram` | |
| 1112 | 1127 | ||
| 1113 | 1128 | <!-- YAML | |
@@ -3998,7 +4013,9 @@ import { listen } from 'node:quic'; | |||
| 3998 | 4013 | const encoder = new TextEncoder(); | |
| 3999 | 4014 | ||
| 4000 | 4015 | const endpoint = await listen((session) => { | |
| 4001 | - // The session.onstream callback fires for each new client-initiated stream. | ||
| 4016 | + // The session.onstream callback fires for each new client-initiated | ||
| 4017 | + // stream. It is optional here: with `onheaders` configured below, | ||
| 4018 | + // request streams are consumed through that callback. | ||
| 4002 | 4019 | }, { | |
| 4003 | 4020 | sni: { '*': { keys: [defaultKey], certs: [defaultCert] } }, | |
| 4004 | 4021 | // ALPN defaults to 'h3'. | |
@@ -4632,5 +4649,6 @@ throughput issues caused by flow control. | |||
| 4632 | 4649 | [`stream.writer`]: #streamwriter | |
| 4633 | 4650 | [`writer.fail()`]: #streamwriter | |
| 4634 | 4651 | [`writer.fail(reason)`]: #streamwriter | |
| 4652 | + [minimal HTTP/3 server]: #minimal-http3-server | ||
| 4635 | 4653 | [qlog]: https://datatracker.ietf.org/doc/draft-ietf-quic-qlog-main-schema/ | |
| 4636 | 4654 | [qvis]: https://qvis.quictools.info/ | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -4129,6 +4129,24 @@ class QuicSession { | |||
| 4129 | 4129 | this.#inner.verifyPeer = value; | |
| 4130 | 4130 | } | |
| 4131 | 4131 | ||
| 4132 | + /** | ||
| 4133 | + * True if an incoming stream has a consumer registered on this session: | ||
| 4134 | + * either an onstream callback, or - when the negotiated application | ||
| 4135 | + * supports headers (e.g. HTTP/3) - session-level stream callbacks that | ||
| 4136 | + * the application layer will invoke (onheaders et al). | ||
| 4137 | + * @returns {boolean} | ||
| 4138 | + */ | ||
| 4139 | + #hasStreamConsumer() { | ||
| 4140 | + if (typeof this.#inner.onstream === 'function') return true; | ||
| 4141 | + // Only onheaders is guaranteed to fire for every incoming stream when | ||
| 4142 | + // the negotiated application supports stream callbacks (e.g. HTTP/3). | ||
| 4143 | + // Other stream callbacks are conditional (ontrailers, oninfo) or | ||
| 4144 | + // outbound-only (onwanttrailers) and do not expose the stream, so they | ||
| 4145 | + // do not count as a consumer. | ||
| 4146 | + if (typeof this[kStreamCallbacks]?.onheaders !== 'function') return false; | ||
| 4147 | + return getQuicSessionState(this).streamCallbacksSupported === 1; | ||
| 4148 | + } | ||
| 4149 | + | ||
| 4132 | 4150 | /** | |
| 4133 | 4151 | * @param {object} handle | |
| 4134 | 4152 | * @param {number} direction | |
@@ -4141,10 +4159,13 @@ class QuicSession { | |||
| 4141 | 4159 | // Set the default byte budget for received streams. | |
| 4142 | 4160 | stream.budget = kDefaultBudget; | |
| 4143 | 4161 | ||
| 4144 | - // A new stream was received. If we don't have an onstream callback, then | ||
| 4145 | - // there's nothing we can do about it. Destroy the stream in this case. | ||
| 4146 | - if (typeof inner.onstream !== 'function') { | ||
| 4147 | - process.emitWarning('A new stream was received but no onstream callback was provided'); | ||
| 4162 | + // A new stream was received. If the session has no consumer for it - | ||
| 4163 | + // neither an onstream callback nor, on a session whose application | ||
| 4164 | + // supports headers (e.g. HTTP/3), any session-level stream callbacks - | ||
| 4165 | + // there's nothing that could ever read it. Destroy the stream in this | ||
| 4166 | + // case rather than letting it hold flow control credit. | ||
| 4167 | + if (!this.#hasStreamConsumer()) { | ||
| 4168 | + process.emitWarning('A new stream was received but no stream consumer callback was provided'); | ||
| 4148 | 4169 | stream.destroy(); | |
| 4149 | 4170 | return; | |
| 4150 | 4171 | } | |
@@ -4175,7 +4196,14 @@ class QuicSession { | |||
| 4175 | 4196 | }); | |
| 4176 | 4197 | } | |
| 4177 | 4198 | ||
| 4178 | - safeCallbackInvoke(inner.onstream, this, stream); | ||
| 4199 | + // Deliver the stream to the onstream consumer if one is registered. | ||
| 4200 | + // Reaching this point without one means #hasStreamConsumer accepted | ||
| 4201 | + // the stream on behalf of the application layer: the session-level | ||
| 4202 | + // stream callbacks were applied above and the application (e.g. | ||
| 4203 | + // HTTP/3) drives the stream, so there is nothing to invoke here. | ||
| 4204 | + if (typeof inner.onstream === 'function') { | ||
| 4205 | + safeCallbackInvoke(inner.onstream, this, stream); | ||
| 4206 | + } | ||
| 4179 | 4207 | } | |
| 4180 | 4208 | ||
| 4181 | 4209 | [kRemoveStream](stream) { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -72,6 +72,7 @@ const { | |||
| 72 | 72 | IDX_STATE_SESSION_STREAM_OPEN_ALLOWED, | |
| 73 | 73 | IDX_STATE_SESSION_PRIORITY_SUPPORTED, | |
| 74 | 74 | IDX_STATE_SESSION_HEADERS_SUPPORTED, | |
| 75 | + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED, | ||
| 75 | 76 | IDX_STATE_SESSION_WRAPPED, | |
| 76 | 77 | IDX_STATE_SESSION_APPLICATION_TYPE, | |
| 77 | 78 | IDX_STATE_SESSION_NO_ERROR_CODE, | |
@@ -119,6 +120,7 @@ assert(IDX_STATE_SESSION_HANDSHAKE_CONFIRMED !== undefined); | |||
| 119 | 120 | assert(IDX_STATE_SESSION_STREAM_OPEN_ALLOWED !== undefined); | |
| 120 | 121 | assert(IDX_STATE_SESSION_PRIORITY_SUPPORTED !== undefined); | |
| 121 | 122 | assert(IDX_STATE_SESSION_HEADERS_SUPPORTED !== undefined); | |
| 123 | + assert(IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED !== undefined); | ||
| 122 | 124 | assert(IDX_STATE_SESSION_WRAPPED !== undefined); | |
| 123 | 125 | assert(IDX_STATE_SESSION_APPLICATION_TYPE !== undefined); | |
| 124 | 126 | assert(IDX_STATE_SESSION_NO_ERROR_CODE !== undefined); | |
@@ -493,6 +495,19 @@ class QuicSessionState { | |||
| 493 | 495 | return DataViewPrototypeGetUint8(handle, this.#offset + IDX_STATE_SESSION_HEADERS_SUPPORTED); | |
| 494 | 496 | } | |
| 495 | 497 | ||
| 498 | + /** | ||
| 499 | + * Whether the negotiated application dispatches the session-level | ||
| 500 | + * stream callbacks (onheaders et al) for incoming streams. | ||
| 501 | + * Returns 0 (unknown), 1 (supported), or 2 (not supported). | ||
| 502 | + * @type {number} | ||
| 503 | + */ | ||
| 504 | + get streamCallbacksSupported() { | ||
| 505 | + const handle = this.#handle; | ||
| 506 | + if (handle === undefined) return undefined; | ||
| 507 | + return DataViewPrototypeGetUint8( | ||
| 508 | + handle, this.#offset + IDX_STATE_SESSION_STREAM_CALLBACKS_SUPPORTED); | ||
| 509 | + } | ||
| 510 | + | ||
| 496 | 511 | /** @type {boolean} */ | |
| 497 | 512 | get isWrapped() { | |
| 498 | 513 | const handle = this.#handle; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -210,6 +210,11 @@ class Session::Application : public MemoryRetainer { | |||
| 210 | 210 | // do not support headers should return false (the default). | |
| 211 | 211 | virtual bool SupportsHeaders() const { return false; } | |
| 212 | 212 | ||
| 213 | + // True if this application dispatches the session-level stream | ||
| 214 | + // callbacks (onheaders et al) for incoming streams when they are | ||
| 215 | + // registered on the session. | ||
| 216 | + virtual bool SupportsStreamCallbacks() const { return false; } | ||
| 217 | + | ||
| 213 | 218 | // Initiates application-level graceful shutdown signaling (e.g., | |
| 214 | 219 | // HTTP/3 GOAWAY). Called when Session::Close(GRACEFUL) is invoked. | |
| 215 | 220 | virtual void BeginShutdown() {} | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -328,6 +328,12 @@ enum class HeadersSupportState : uint8_t { | |||
| 328 | 328 | UNSUPPORTED, | |
| 329 | 329 | }; | |
| 330 | 330 | ||
| 331 | + enum class StreamCallbacksSupportState : uint8_t { | ||
| 332 | + UNKNOWN, | ||
| 333 | + SUPPORTED, | ||
| 334 | + UNSUPPORTED, | ||
| 335 | + }; | ||
| 336 | + | ||
| 331 | 337 | enum class PathValidationResult : uint8_t { | |
| 332 | 338 | SUCCESS = NGTCP2_PATH_VALIDATION_RESULT_SUCCESS, | |
| 333 | 339 | FAILURE = NGTCP2_PATH_VALIDATION_RESULT_FAILURE, | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -202,6 +202,8 @@ class Http3ApplicationImpl final : public Session::Application { | |||
| 202 | 202 | ||
| 203 | 203 | bool SupportsHeaders() const override { return true; } | |
| 204 | 204 | ||
| 205 | + bool SupportsStreamCallbacks() const override { return true; } | ||
| 206 | + | ||
| 205 | 207 | bool is_started() const override { return started_; } | |
| 206 | 208 | ||
| 207 | 209 | bool Start() override { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -136,6 +136,7 @@ uint64_t MaxDatagramPayload(uint64_t max_frame_size) { | |||
| 136 | 136 | V(STREAM_OPEN_ALLOWED, stream_open_allowed, uint8_t) \ | |
| 137 | 137 | V(PRIORITY_SUPPORTED, priority_supported, uint8_t) \ | |
| 138 | 138 | V(HEADERS_SUPPORTED, headers_supported, uint8_t) \ | |
| 139 | + V(STREAM_CALLBACKS_SUPPORTED, stream_callbacks_supported, uint8_t) \ | ||
| 139 | 140 | V(WRAPPED, wrapped, uint8_t) \ | |
| 140 | 141 | V(APPLICATION_TYPE, application_type, uint8_t) \ | |
| 141 | 142 | V(NO_ERROR_CODE, no_error_code, error_code) \ | |
@@ -2649,6 +2650,10 @@ void Session::SetApplication(std::unique_ptr<Application> app) { | |||
| 2649 | 2650 | impl_->state()->headers_supported = static_cast<uint8_t>( | |
| 2650 | 2651 | app->SupportsHeaders() ? HeadersSupportState::SUPPORTED | |
| 2651 | 2652 | : HeadersSupportState::UNSUPPORTED); | |
| 2653 | + impl_->state()->stream_callbacks_supported = | ||
| 2654 | + static_cast<uint8_t>(app->SupportsStreamCallbacks() | ||
| 2655 | + ? StreamCallbacksSupportState::SUPPORTED | ||
| 2656 | + : StreamCallbacksSupportState::UNSUPPORTED); | ||
| 2652 | 2657 | // Surface the application's "no error" and "internal error" codes via | |
| 2653 | 2658 | // session state so that JS-side code (e.g. the stream writer's fail() | |
| 2654 | 2659 | // path) can resolve the right wire code for the negotiated ALPN | |
| Back | FazBrowse Home | New Git URL |
0 commit comments