tokio_quiche/http3/driver/
server.rs1use std::ops::Deref;
28use std::sync::Arc;
29
30use tokio::sync::mpsc;
31
32use super::datagram;
33use super::DriverHooks;
34use super::H3Command;
35use super::H3ConnectionError;
36use super::H3ConnectionResult;
37use super::H3Controller;
38use super::H3Driver;
39use super::H3Event;
40use super::InboundHeaders;
41use super::IncomingH3Headers;
42use super::StreamCtx;
43use super::STREAM_CAPACITY;
44use crate::http3::settings::Http3Settings;
45use crate::http3::settings::Http3SettingsEnforcer;
46use crate::http3::settings::Http3TimeoutType;
47use crate::http3::settings::TimeoutKey;
48use crate::quic::HandshakeInfo;
49use crate::quic::QuicCommand;
50use crate::quic::QuicheConnection;
51
52pub type ServerH3Driver = H3Driver<ServerHooks>;
56pub type ServerH3Controller = H3Controller<ServerHooks>;
59pub type ServerEventStream = mpsc::UnboundedReceiver<ServerH3Event>;
63
64#[derive(Clone, Debug)]
65pub struct RawPriorityValue(Vec<u8>);
66
67impl From<Vec<u8>> for RawPriorityValue {
68 fn from(value: Vec<u8>) -> Self {
69 RawPriorityValue(value)
70 }
71}
72
73impl Deref for RawPriorityValue {
74 type Target = [u8];
75
76 fn deref(&self) -> &Self::Target {
77 &self.0
78 }
79}
80
81#[derive(Clone, Debug)]
83pub struct IsInEarlyData(bool);
84
85impl IsInEarlyData {
86 fn new(is_in_early_data: bool) -> Self {
87 IsInEarlyData(is_in_early_data)
88 }
89}
90
91impl Deref for IsInEarlyData {
92 type Target = bool;
93
94 fn deref(&self) -> &Self::Target {
95 &self.0
96 }
97}
98
99#[derive(Debug)]
101pub enum ServerH3Event {
102 Core(H3Event),
103
104 Headers {
105 incoming_headers: IncomingH3Headers,
106 priority: Option<RawPriorityValue>,
108 is_in_early_data: IsInEarlyData,
109 },
110}
111
112impl From<H3Event> for ServerH3Event {
113 fn from(ev: H3Event) -> Self {
114 match ev {
115 H3Event::IncomingHeaders(incoming_headers) => {
116 Self::Headers {
122 incoming_headers,
123 priority: None,
124 is_in_early_data: IsInEarlyData::new(false),
125 }
126 },
127 _ => Self::Core(ev),
128 }
129 }
130}
131
132#[derive(Debug)]
134pub enum ServerH3Command {
135 Core(H3Command),
136}
137
138impl From<H3Command> for ServerH3Command {
139 fn from(cmd: H3Command) -> Self {
140 Self::Core(cmd)
141 }
142}
143
144impl From<QuicCommand> for ServerH3Command {
145 fn from(cmd: QuicCommand) -> Self {
146 Self::Core(H3Command::QuicCmd(cmd))
147 }
148}
149
150const PRE_HEADERS_BOOSTED_PRIORITY_URGENCY: u8 = 64;
154const PRE_HEADERS_BOOSTED_PRIORITY_INCREMENTAL: bool = false;
157
158pub struct ServerHooks {
159 settings_enforcer: Http3SettingsEnforcer,
161 requests: u64,
163
164 post_accept_timeout: Option<TimeoutKey>,
167 extended_connect_enabled: bool,
170}
171
172impl ServerHooks {
173 fn handle_request(
178 driver: &mut H3Driver<Self>, qconn: &mut QuicheConnection,
179 headers: InboundHeaders,
180 ) -> H3ConnectionResult<()> {
181 let InboundHeaders {
182 stream_id,
183 headers,
184 has_body,
185 } = headers;
186
187 if driver.stream_map.contains_key(&stream_id) {
191 return Ok(());
192 }
193
194 let (mut stream_ctx, send, recv) =
195 StreamCtx::new(stream_id, STREAM_CAPACITY);
196
197 if driver.hooks.extended_connect_enabled() {
204 if let Some(flow_id) =
205 datagram::extract_quarter_stream_id(stream_id, &headers)
206 {
207 let _ = driver.get_or_insert_flow(flow_id)?;
208 stream_ctx.associated_dgram_flow_id = Some(flow_id);
209 }
210 }
211
212 let latest_priority_update: Option<RawPriorityValue> = driver
213 .conn_mut()?
214 .take_last_priority_update(stream_id)
215 .ok()
216 .map(|v| v.into());
217
218 qconn
222 .stream_priority(
223 stream_id,
224 PRE_HEADERS_BOOSTED_PRIORITY_URGENCY,
225 PRE_HEADERS_BOOSTED_PRIORITY_INCREMENTAL,
226 )
227 .ok();
228
229 let headers = IncomingH3Headers {
230 stream_id,
231 headers,
232 send,
233 recv,
234 read_fin: !has_body,
235 h3_audit_stats: Arc::clone(&stream_ctx.audit_stats),
236 };
237
238 driver
239 .waiting_streams
240 .push(stream_ctx.wait_for_recv(stream_id));
241 driver.insert_stream(stream_id, stream_ctx);
242
243 driver.process_writable_stream(qconn, stream_id)?;
249
250 driver
251 .h3_event_sender
252 .send(ServerH3Event::Headers {
253 incoming_headers: headers,
254 priority: latest_priority_update,
255 is_in_early_data: IsInEarlyData::new(qconn.is_in_early_data()),
256 })
257 .map_err(|_| H3ConnectionError::ControllerWentAway)?;
258 driver.hooks.requests += 1;
259
260 Ok(())
261 }
262}
263
264#[allow(private_interfaces)]
265impl DriverHooks for ServerHooks {
266 type Command = ServerH3Command;
267 type Event = ServerH3Event;
268
269 fn new(settings: &Http3Settings) -> Self {
270 Self {
271 settings_enforcer: settings.into(),
272 requests: 0,
273 post_accept_timeout: None,
274 extended_connect_enabled: settings.enable_extended_connect,
275 }
276 }
277
278 fn extended_connect_enabled(&self) -> bool {
279 self.extended_connect_enabled
280 }
281
282 fn conn_established(
283 driver: &mut H3Driver<Self>, qconn: &mut QuicheConnection,
284 handshake_info: &HandshakeInfo,
285 ) -> H3ConnectionResult<()> {
286 assert!(
287 qconn.is_server(),
288 "ServerH3Driver requires a server-side QUIC connection"
289 );
290
291 if let Some(post_accept_timeout) =
292 driver.hooks.settings_enforcer.post_accept_timeout()
293 {
294 let remaining = post_accept_timeout
295 .checked_sub(handshake_info.elapsed())
296 .ok_or(H3ConnectionError::PostAcceptTimeout)?;
297
298 let key = driver
299 .hooks
300 .settings_enforcer
301 .add_timeout(Http3TimeoutType::PostAccept, remaining);
302 driver.hooks.post_accept_timeout = Some(key);
303 }
304
305 Ok(())
306 }
307
308 fn headers_received(
309 driver: &mut H3Driver<Self>, qconn: &mut QuicheConnection,
310 headers: InboundHeaders,
311 ) -> H3ConnectionResult<()> {
312 if driver
313 .hooks
314 .settings_enforcer
315 .enforce_requests_limit(driver.hooks.requests)
316 {
317 let _ =
318 qconn.close(true, quiche::h3::WireErrorCode::NoError as u64, &[]);
319 return Ok(());
320 }
321
322 if let Some(timeout) = driver.hooks.post_accept_timeout.take() {
323 driver.hooks.settings_enforcer.cancel_timeout(timeout);
326 }
327
328 Self::handle_request(driver, qconn, headers)
329 }
330
331 fn conn_command(
332 driver: &mut H3Driver<Self>, qconn: &mut QuicheConnection,
333 cmd: Self::Command,
334 ) -> H3ConnectionResult<()> {
335 let ServerH3Command::Core(cmd) = cmd;
336 driver.handle_core_command(qconn, cmd)
337 }
338
339 fn has_wait_action(driver: &mut H3Driver<Self>) -> bool {
340 driver.hooks.settings_enforcer.has_pending_timeouts()
341 }
342
343 async fn wait_for_action(
344 &mut self, qconn: &mut QuicheConnection,
345 ) -> H3ConnectionResult<()> {
346 self.settings_enforcer.enforce_timeouts(qconn).await?;
347 Err(H3ConnectionError::PostAcceptTimeout)
348 }
349}