tokio_quiche/quic/connection/mod.rs
1// Copyright (C) 2025, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8// * Redistributions of source code must retain the above copyright notice,
9// this list of conditions and the following disclaimer.
10//
11// * Redistributions in binary form must reproduce the above copyright
12// notice, this list of conditions and the following disclaimer in the
13// documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27mod error;
28mod id;
29mod map;
30
31pub use self::error::HandshakeError;
32pub use self::id::ConnectionIdGenerator;
33pub use self::id::SharedConnectionIdGenerator;
34pub use self::id::SimpleConnectionIdGenerator;
35pub(crate) use self::map::ConnectionMap;
36
37use boring::ssl::SslRef;
38use datagram_socket::AsSocketStats;
39use datagram_socket::DatagramSocketSend;
40use datagram_socket::MaybeConnectedSocket;
41use datagram_socket::QuicAuditStats;
42use datagram_socket::ShutdownConnection;
43use datagram_socket::SocketStats;
44use foundations::telemetry::log;
45use futures::Future;
46use quiche::ConnectionId;
47use std::fmt;
48use std::io;
49use std::net::SocketAddr;
50use std::sync::Arc;
51use std::sync::Mutex;
52use std::task::Poll;
53use std::time::Duration;
54use std::time::Instant;
55use std::time::SystemTime;
56use tokio::sync::mpsc;
57use tokio_util::task::AbortOnDropHandle;
58
59use self::error::make_handshake_result;
60use super::hooks::ConnectionHook;
61use super::io::connection_stage::Close;
62use super::io::connection_stage::ConnectionStageContext;
63use super::io::connection_stage::Handshake;
64use super::io::connection_stage::RunningApplication;
65use super::io::worker::Closing;
66use super::io::worker::IoWorkerParams;
67use super::io::worker::Running;
68use super::io::worker::RunningOrClosing;
69use super::io::worker::WriteState;
70use super::QuicheConnection;
71use crate::metrics::Metrics;
72use crate::quic::io::worker::IoWorker;
73use crate::quic::io::worker::WriterConfig;
74use crate::quic::io::worker::INCOMING_QUEUE_SIZE;
75use crate::quic::router::ConnectionMapCommand;
76use crate::QuicResult;
77
78/// Wrapper for connection statistics recorded by [quiche].
79#[derive(Debug)]
80pub struct QuicConnectionStats {
81 /// Aggregate connection statistics across all paths.
82 pub stats: quiche::Stats,
83 /// Specific statistics about the connection's active path.
84 pub path_stats: Option<quiche::PathStats>,
85}
86pub(crate) type QuicConnectionStatsShared = Arc<Mutex<QuicConnectionStats>>;
87
88impl QuicConnectionStats {
89 pub(crate) fn from_conn(qconn: &QuicheConnection) -> Self {
90 Self {
91 stats: qconn.stats(),
92 path_stats: qconn.path_stats().find(|stats| stats.active),
93 }
94 }
95
96 fn startup_exit_to_socket_stats(
97 value: quiche::StartupExit,
98 ) -> datagram_socket::StartupExit {
99 let reason = match value.reason {
100 quiche::StartupExitReason::Loss =>
101 datagram_socket::StartupExitReason::Loss,
102 quiche::StartupExitReason::BandwidthPlateau =>
103 datagram_socket::StartupExitReason::BandwidthPlateau,
104 quiche::StartupExitReason::PersistentQueue =>
105 datagram_socket::StartupExitReason::PersistentQueue,
106 quiche::StartupExitReason::ConservativeSlowStartRounds =>
107 datagram_socket::StartupExitReason::ConservativeSlowStartRounds,
108 };
109
110 datagram_socket::StartupExit {
111 cwnd: value.cwnd,
112 bandwidth: value.bandwidth,
113 reason,
114 }
115 }
116}
117
118impl AsSocketStats for QuicConnectionStats {
119 fn as_socket_stats(&self) -> SocketStats {
120 SocketStats {
121 pmtu: self
122 .path_stats
123 .as_ref()
124 .map(|p| p.pmtu as u16)
125 .unwrap_or_default(),
126 rtt_us: self
127 .path_stats
128 .as_ref()
129 .map(|p| p.rtt.as_micros() as i64)
130 .unwrap_or_default(),
131 min_rtt_us: self
132 .path_stats
133 .as_ref()
134 .and_then(|p| p.min_rtt.map(|x| x.as_micros() as i64))
135 .unwrap_or_default(),
136 max_rtt_us: self
137 .path_stats
138 .as_ref()
139 .and_then(|p| p.max_rtt.map(|x| x.as_micros() as i64))
140 .unwrap_or_default(),
141 rtt_var_us: self
142 .path_stats
143 .as_ref()
144 .map(|p| p.rttvar.as_micros() as i64)
145 .unwrap_or_default(),
146 cwnd: self
147 .path_stats
148 .as_ref()
149 .map(|p| p.cwnd as u64)
150 .unwrap_or_default(),
151 total_pto_count: self
152 .path_stats
153 .as_ref()
154 .map(|p| p.total_pto_count as u64)
155 .unwrap_or_default(),
156 packets_sent: self.stats.sent as u64,
157 packets_recvd: self.stats.recv as u64,
158 packets_lost: self.stats.lost as u64,
159 packets_lost_spurious: self.stats.spurious_lost as u64,
160 packets_retrans: self.stats.retrans as u64,
161 bytes_sent: self.stats.sent_bytes,
162 bytes_recvd: self.stats.recv_bytes,
163 bytes_lost: self.stats.lost_bytes,
164 bytes_retrans: self.stats.stream_retrans_bytes,
165 bytes_unsent: 0, /* not implemented yet, kept for compatibility
166 * with TCP */
167 delivery_rate: self
168 .path_stats
169 .as_ref()
170 .map(|p| p.delivery_rate)
171 .unwrap_or_default(),
172 max_bandwidth: self.path_stats.as_ref().and_then(|p| p.max_bandwidth),
173 startup_exit: self
174 .path_stats
175 .as_ref()
176 .and_then(|p| p.startup_exit)
177 .map(QuicConnectionStats::startup_exit_to_socket_stats),
178 bytes_in_flight_duration_us: self
179 .stats
180 .bytes_in_flight_duration
181 .as_micros() as u64,
182 }
183 }
184}
185
186/// A received network packet with additional metadata.
187#[derive(Debug)]
188pub struct Incoming {
189 /// The address that sent the inbound packet.
190 pub peer_addr: SocketAddr,
191 /// The address on which we received the inbound packet.
192 pub local_addr: SocketAddr,
193 /// The receive timestamp of the packet.
194 ///
195 /// Used for the `perf-quic-listener-metrics` feature.
196 pub rx_time: Option<SystemTime>,
197 /// The packet's contents.
198 pub buf: Vec<u8>,
199 /// If set, then `buf` is a GRO buffer containing multiple packets.
200 /// Each individual packet has a size of `gso` (except for the last one).
201 pub gro: Option<i32>,
202 /// [SO_MARK] control message value received from the socket.
203 ///
204 /// This will always be `None` after the connection has been spawned as
205 /// the message is `take()`d before spawning.
206 ///
207 /// [SO_MARK]: https://man7.org/linux/man-pages/man7/socket.7.html
208 #[cfg(target_os = "linux")]
209 pub so_mark_data: Option<[u8; 4]>,
210}
211
212/// A QUIC connection that has not performed a handshake yet.
213///
214/// This type is currently only used for server-side connections. It is created
215/// and added to the listener's connection stream after an initial packet from
216/// a client has been received and (optionally) the client's IP address has been
217/// validated.
218///
219/// To turn the initial connection into a fully established one, a QUIC
220/// handshake must be performed. Users have multiple options to facilitate this:
221/// - `start` is a simple entrypoint which spawns a task to handle the entire
222/// lifetime of the QUIC connection. The caller can then only communicate with
223/// the connection via their [`ApplicationOverQuic`].
224/// - `handshake` spawns a task for the handshake and awaits its completion.
225/// Afterwards, it pauses the connection and allows the caller to resume it
226/// later via an opaque struct. We spawn a separate task to allow the tokio
227/// scheduler free choice in where to run the handshake.
228/// - `handshake_fut` returns a future to drive the handshake for maximum
229/// flexibility.
230#[must_use = "call InitialQuicConnection::start to establish the connection"]
231pub struct InitialQuicConnection<Tx, M>
232where
233 Tx: DatagramSocketSend + Send + 'static + ?Sized,
234 M: Metrics,
235{
236 params: QuicConnectionParams<Tx, M>,
237 pub(crate) audit_log_stats: Arc<QuicAuditStats>,
238 stats: QuicConnectionStatsShared,
239 pub(crate) incoming_ev_sender: mpsc::Sender<Incoming>,
240 incoming_ev_receiver: mpsc::Receiver<Incoming>,
241}
242
243impl<Tx, M> InitialQuicConnection<Tx, M>
244where
245 Tx: DatagramSocketSend + Send + 'static + ?Sized,
246 M: Metrics,
247{
248 #[inline]
249 pub(crate) fn new(params: QuicConnectionParams<Tx, M>) -> Self {
250 let (incoming_ev_sender, incoming_ev_receiver) =
251 mpsc::channel(INCOMING_QUEUE_SIZE);
252 let audit_log_stats = Arc::new(QuicAuditStats::new(params.scid.to_vec()));
253
254 let stats = Arc::new(Mutex::new(QuicConnectionStats::from_conn(
255 ¶ms.quiche_conn,
256 )));
257
258 Self {
259 params,
260 audit_log_stats,
261 stats,
262 incoming_ev_sender,
263 incoming_ev_receiver,
264 }
265 }
266
267 /// The local address this connection listens on.
268 pub fn local_addr(&self) -> SocketAddr {
269 self.params.local_addr
270 }
271
272 /// The remote address for this connection.
273 pub fn peer_addr(&self) -> SocketAddr {
274 self.params.peer_addr
275 }
276
277 /// [boring]'s SSL object for this connection.
278 #[doc(hidden)]
279 pub fn ssl_mut(&mut self) -> &mut SslRef {
280 // Deref to pick `Connection::as_mut` over `Box::as_mut`.
281 (*self.params.quiche_conn).as_mut()
282 }
283
284 /// A handle to the [`QuicAuditStats`] for this connection.
285 ///
286 /// # Note
287 /// These stats are updated during the lifetime of the connection.
288 /// The getter exists to grab a handle early on, which can then
289 /// be stowed away and read out after the connection has closed.
290 #[inline]
291 pub fn audit_log_stats(&self) -> Arc<QuicAuditStats> {
292 Arc::clone(&self.audit_log_stats)
293 }
294
295 /// A handle to the [`QuicConnectionStats`] for this connection.
296 ///
297 /// # Note
298 /// Initially, these stats represent the state when the [quiche::Connection]
299 /// was created. They are updated when the connection is closed, so this
300 /// getter exists primarily to grab a handle early on.
301 #[inline]
302 pub fn stats(&self) -> &QuicConnectionStatsShared {
303 &self.stats
304 }
305
306 /// Creates a future to drive the connection's handshake.
307 ///
308 /// This is a lower-level alternative to the `handshake` function which
309 /// gives the caller more control over execution of the future. See
310 /// `handshake` for details on the return values.
311 pub fn handshake_fut<A: ApplicationOverQuic>(
312 self, app: A,
313 ) -> (
314 QuicConnection,
315 impl Future<Output = io::Result<Running<Arc<Tx>, M, A>>> + Send + 'static,
316 ) {
317 self.params.metrics.connections_in_memory().inc();
318
319 let conn = QuicConnection {
320 local_addr: self.params.local_addr,
321 peer_addr: self.params.peer_addr,
322 audit_log_stats: Arc::clone(&self.audit_log_stats),
323 stats: Arc::clone(&self.stats),
324 scid: self.params.scid,
325 };
326 let context = ConnectionStageContext {
327 in_pkt: self.params.initial_pkt,
328 incoming_pkt_receiver: self.incoming_ev_receiver,
329 application: app,
330 stats: Arc::clone(&self.stats),
331 connection_hook: self.params.connection_hook,
332 };
333 let conn_stage = Handshake {
334 handshake_info: self.params.handshake_info,
335 };
336 let params = IoWorkerParams {
337 socket: MaybeConnectedSocket::new(self.params.socket),
338 shutdown_tx: self.params.shutdown_tx,
339 cfg: self.params.writer_cfg,
340 audit_log_stats: self.audit_log_stats,
341 write_state: WriteState::default(),
342 conn_map_cmd_tx: self.params.conn_map_cmd_tx,
343 cid_generator: self.params.cid_generator,
344 #[cfg(feature = "perf-quic-listener-metrics")]
345 init_rx_time: self.params.init_rx_time,
346 metrics: self.params.metrics.clone(),
347 };
348
349 let handshake_fut = async move {
350 let qconn = self.params.quiche_conn;
351 let handshake_done =
352 IoWorker::new(params, conn_stage).run(qconn, context).await;
353
354 match handshake_done {
355 RunningOrClosing::Running(r) => Ok(r),
356 RunningOrClosing::Closing(Closing {
357 params,
358 work_loop_result,
359 mut context,
360 mut qconn,
361 }) => {
362 let hs_result = make_handshake_result(&work_loop_result);
363 IoWorker::new(params, Close { work_loop_result })
364 .close(&mut qconn, &mut context)
365 .await;
366 hs_result
367 },
368 }
369 };
370
371 (conn, handshake_fut)
372 }
373
374 /// Performs the QUIC handshake in a separate tokio task and awaits its
375 /// completion.
376 ///
377 /// The returned [`QuicConnection`] holds metadata about the established
378 /// connection. The connection itself is paused after `handshake`
379 /// returns and must be resumed by passing the opaque `Running` value to
380 /// [`InitialQuicConnection::resume`]. This two-step process
381 /// allows callers to collect telemetry and run code before serving their
382 /// [`ApplicationOverQuic`].
383 pub async fn handshake<A: ApplicationOverQuic>(
384 self, app: A,
385 ) -> io::Result<(QuicConnection, Running<Arc<Tx>, M, A>)> {
386 let task_metrics = self.params.metrics.clone();
387 let (conn, handshake_fut) = Self::handshake_fut(self, app);
388
389 let handshake_handle = crate::metrics::tokio_task::spawn(
390 "quic_handshake_worker",
391 task_metrics,
392 handshake_fut,
393 );
394
395 // `AbortOnDropHandle` simulates task-killswitch behavior without
396 // needing to give up ownership of the `JoinHandle`.
397 let handshake_abort_handle = AbortOnDropHandle::new(handshake_handle);
398
399 let worker = handshake_abort_handle.await??;
400
401 Ok((conn, worker))
402 }
403
404 /// Resumes a QUIC connection which was paused after a successful handshake.
405 pub fn resume<A: ApplicationOverQuic>(pre_running: Running<Arc<Tx>, M, A>) {
406 let task_metrics = pre_running.params.metrics.clone();
407 let fut = async move {
408 let Running {
409 params,
410 context,
411 qconn,
412 } = pre_running;
413 let running_worker = IoWorker::new(params, RunningApplication);
414
415 let Closing {
416 params,
417 mut context,
418 work_loop_result,
419 mut qconn,
420 } = running_worker.run(qconn, context).await;
421
422 IoWorker::new(params, Close { work_loop_result })
423 .close(&mut qconn, &mut context)
424 .await;
425 };
426
427 crate::metrics::tokio_task::spawn_with_killswitch(
428 "quic_io_worker",
429 task_metrics,
430 fut,
431 );
432 }
433
434 /// Drives a QUIC connection from handshake to close in separate tokio
435 /// tasks.
436 ///
437 /// It combines [`InitialQuicConnection::handshake`] and
438 /// [`InitialQuicConnection::resume`] into a single call.
439 pub fn start<A: ApplicationOverQuic>(self, app: A) -> QuicConnection {
440 let task_metrics = self.params.metrics.clone();
441 let (conn, handshake_fut) = Self::handshake_fut(self, app);
442 // Pin to the heap so the spawned task only carries a pointer to it
443 // instead of inlining the full future state across the await.
444 let handshake_fut = Box::pin(handshake_fut);
445
446 let fut = async move {
447 match handshake_fut.await {
448 Ok(running) => Self::resume(running),
449 Err(e) => {
450 log::error!("QUIC handshake failed in IQC::start"; "error" => e);
451 },
452 }
453 };
454
455 crate::metrics::tokio_task::spawn_with_killswitch(
456 "quic_handshake_worker",
457 task_metrics,
458 fut,
459 );
460
461 conn
462 }
463}
464
465pub(crate) struct QuicConnectionParams<Tx, M>
466where
467 Tx: DatagramSocketSend + Send + 'static + ?Sized,
468 M: Metrics,
469{
470 pub writer_cfg: WriterConfig,
471 pub initial_pkt: Option<Incoming>,
472 pub shutdown_tx: mpsc::Sender<()>,
473 pub conn_map_cmd_tx: mpsc::UnboundedSender<ConnectionMapCommand>, /* channel that signals connection map changes */
474 pub scid: ConnectionId<'static>,
475 pub cid_generator: Option<SharedConnectionIdGenerator>,
476 pub metrics: M,
477 pub connection_hook: Option<Arc<dyn ConnectionHook + Send + Sync + 'static>>,
478 #[cfg(feature = "perf-quic-listener-metrics")]
479 pub init_rx_time: Option<SystemTime>,
480 pub handshake_info: HandshakeInfo,
481 /// Boxed because this value is moved by-value through several nested
482 /// async state machines. Inlining a [`QuicheConnection`] here would
483 /// duplicate its payload across the future state slots that hold it
484 /// across an `.await`.
485 pub quiche_conn: Box<QuicheConnection>,
486 pub socket: Arc<Tx>,
487 pub local_addr: SocketAddr,
488 pub peer_addr: SocketAddr,
489}
490
491/// Metadata about an established QUIC connection.
492///
493/// While this struct allows access to some facets of a QUIC connection, it
494/// notably does not represent the [quiche::Connection] itself. The crate
495/// handles most interactions with [quiche] internally in a worker task. Users
496/// can only access the connection directly via their [`ApplicationOverQuic`]
497/// implementation.
498///
499/// See the [module-level docs](crate::quic) for an overview of how a QUIC
500/// connection is handled internally.
501pub struct QuicConnection {
502 local_addr: SocketAddr,
503 peer_addr: SocketAddr,
504 audit_log_stats: Arc<QuicAuditStats>,
505 stats: QuicConnectionStatsShared,
506 scid: ConnectionId<'static>,
507}
508
509impl QuicConnection {
510 /// The local address this connection listens on.
511 #[inline]
512 pub fn local_addr(&self) -> SocketAddr {
513 self.local_addr
514 }
515
516 /// The remote address for this connection.
517 #[inline]
518 pub fn peer_addr(&self) -> SocketAddr {
519 self.peer_addr
520 }
521
522 /// A handle to the [`QuicAuditStats`] for this connection.
523 ///
524 /// # Note
525 /// These stats are updated during the lifetime of the connection.
526 /// The getter exists to grab a handle early on, which can then
527 /// be stowed away and read out after the connection has closed.
528 #[inline]
529 pub fn audit_log_stats(&self) -> &Arc<QuicAuditStats> {
530 &self.audit_log_stats
531 }
532
533 /// A handle to the [`QuicConnectionStats`] for this connection.
534 ///
535 /// # Note
536 /// Initially, these stats represent the state when the [quiche::Connection]
537 /// was created. They are updated when the connection is closed, so this
538 /// getter exists primarily to grab a handle early on.
539 #[inline]
540 pub fn stats(&self) -> &QuicConnectionStatsShared {
541 &self.stats
542 }
543
544 /// The QUIC source connection ID used by this connection.
545 #[inline]
546 pub fn scid(&self) -> &ConnectionId<'static> {
547 &self.scid
548 }
549}
550
551impl AsSocketStats for QuicConnection {
552 #[inline]
553 fn as_socket_stats(&self) -> SocketStats {
554 // It is important to note that those stats are only updated when
555 // the connection stops, which is fine, since this is only used to
556 // log after the connection is finished.
557 self.stats.lock().unwrap().as_socket_stats()
558 }
559
560 #[inline]
561 fn as_quic_stats(&self) -> Option<&Arc<QuicAuditStats>> {
562 Some(&self.audit_log_stats)
563 }
564}
565
566impl<Tx, M> AsSocketStats for InitialQuicConnection<Tx, M>
567where
568 Tx: DatagramSocketSend + Send + 'static + ?Sized,
569 M: Metrics,
570{
571 #[inline]
572 fn as_socket_stats(&self) -> SocketStats {
573 // It is important to note that those stats are only updated when
574 // the connection stops, which is fine, since this is only used to
575 // log after the connection is finished.
576 self.stats.lock().unwrap().as_socket_stats()
577 }
578
579 #[inline]
580 fn as_quic_stats(&self) -> Option<&Arc<QuicAuditStats>> {
581 Some(&self.audit_log_stats)
582 }
583}
584
585impl<Tx, M> ShutdownConnection for InitialQuicConnection<Tx, M>
586where
587 Tx: DatagramSocketSend + Send + 'static + ?Sized,
588 M: Metrics,
589{
590 #[inline]
591 fn poll_shutdown(
592 &mut self, _cx: &mut std::task::Context,
593 ) -> std::task::Poll<io::Result<()>> {
594 // TODO: Does nothing at the moment. We always call Self::start
595 // anyway so it's not really important at this moment.
596 Poll::Ready(Ok(()))
597 }
598}
599
600impl ShutdownConnection for QuicConnection {
601 #[inline]
602 fn poll_shutdown(
603 &mut self, _cx: &mut std::task::Context,
604 ) -> std::task::Poll<io::Result<()>> {
605 // TODO: does nothing at the moment
606 Poll::Ready(Ok(()))
607 }
608}
609
610/// Details about a connection's QUIC handshake.
611#[derive(Debug, Clone)]
612pub struct HandshakeInfo {
613 /// The time at which the connection was created.
614 start_time: Instant,
615 /// The timeout before which the handshake must complete.
616 timeout: Option<Duration>,
617 /// The real duration that the handshake took to complete.
618 time_handshake: Option<Duration>,
619}
620
621impl HandshakeInfo {
622 pub(crate) fn new(start_time: Instant, timeout: Option<Duration>) -> Self {
623 Self {
624 start_time,
625 timeout,
626 time_handshake: None,
627 }
628 }
629
630 /// The time at which the connection was created.
631 #[inline]
632 pub fn start_time(&self) -> Instant {
633 self.start_time
634 }
635
636 /// How long the handshake took to complete.
637 #[inline]
638 pub fn elapsed(&self) -> Duration {
639 self.time_handshake.unwrap_or_default()
640 }
641
642 pub(crate) fn set_elapsed(&mut self) {
643 let elapsed = self.start_time.elapsed();
644 self.time_handshake = Some(elapsed)
645 }
646
647 pub(crate) fn deadline(&self) -> Option<Instant> {
648 self.timeout.map(|timeout| self.start_time + timeout)
649 }
650
651 pub(crate) fn is_expired(&self) -> bool {
652 self.timeout
653 .is_some_and(|timeout| self.start_time.elapsed() >= timeout)
654 }
655}
656
657/// A trait to implement an application served over QUIC.
658///
659/// The application is driven by an internal worker task, which also handles I/O
660/// for the connection. The worker feeds inbound packets into the
661/// [quiche::Connection], calls [`ApplicationOverQuic::process_reads`] followed
662/// by [`ApplicationOverQuic::process_writes`], and then flushes any pending
663/// outbound packets to the network. This repeats in a loop until either the
664/// connection is closed or the [`ApplicationOverQuic`] returns an error.
665///
666/// In between loop iterations, the worker yields until a new packet arrives, a
667/// timer expires, or [`ApplicationOverQuic::wait_for_data`] resolves.
668/// Implementors can interact with the underlying connection via the mutable
669/// reference passed to trait methods.
670#[allow(unused_variables)] // for default functions
671pub trait ApplicationOverQuic: Send + 'static {
672 /// Callback to customize the [`ApplicationOverQuic`] after the QUIC
673 /// handshake completed successfully.
674 ///
675 /// # Errors
676 /// Returning an error from this method immediately stops the worker loop
677 /// and transitions to the connection closing stage.
678 fn on_conn_established(
679 &mut self, qconn: &mut QuicheConnection, handshake_info: &HandshakeInfo,
680 ) -> QuicResult<()>;
681
682 /// Determines whether the application's methods will be called by the
683 /// worker.
684 ///
685 /// The function is checked in each iteration of the worker loop. Only
686 /// `on_conn_established()` bypasses this check.
687 fn should_act(&self) -> bool;
688
689 /// Waits for an event to trigger the next iteration of the worker loop.
690 ///
691 /// The returned future is awaited in parallel to inbound packets and the
692 /// connection's timers. Any one of those futures resolving triggers the
693 /// next loop iteration, so implementations should not rely on
694 /// `wait_for_data` for the bulk of their processing. Instead, after
695 /// `wait_for_data` resolves, `process_writes` should be used to pull all
696 /// available data out of the event source (for example, a channel).
697 ///
698 /// As for any future, it is **very important** that this method does not
699 /// block the runtime. If it does, the other concurrent futures will be
700 /// starved.
701 ///
702 /// # Cancel safety
703 /// This method MUST be cancel safe.
704 /// It gets called inside select! and could be (repeatedly) cancelled
705 ///
706 /// # Errors
707 /// Returning an error from this method immediately stops the worker loop
708 /// and transitions to the connection closing stage.
709 fn wait_for_data(
710 &mut self, qconn: &mut QuicheConnection,
711 ) -> impl Future<Output = QuicResult<()>> + Send;
712
713 /// Processes data received on the connection.
714 ///
715 /// This method is only called if `should_act()` returns `true` and any
716 /// packets were received since the last worker loop iteration. It
717 /// should be used to read from the connection's open streams.
718 ///
719 /// # Errors
720 /// Returning an error from this method immediately stops the worker loop
721 /// and transitions to the connection closing stage.
722 fn process_reads(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()>;
723
724 /// Adds data to be sent on the connection.
725 ///
726 /// Unlike `process_reads`, this method is called on every iteration of the
727 /// worker loop (provided `should_act()` returns true). It is called
728 /// after `process_reads` and immediately before packets are pushed to
729 /// the socket. The main use case is providing already-buffered data to
730 /// the [quiche::Connection].
731 ///
732 /// # Errors
733 /// Returning an error from this method immediately stops the worker loop
734 /// and transitions to the connection closing stage.
735 fn process_writes(&mut self, qconn: &mut QuicheConnection) -> QuicResult<()>;
736
737 /// Callback to inspect the result of the worker task, before a final packet
738 /// with a `CONNECTION_CLOSE` frame is flushed to the network.
739 ///
740 /// `connection_result` is [`Ok`] only if the connection was closed without
741 /// any local error. Otherwise, the state of `qconn` depends on the
742 /// error type and application behavior.
743 fn on_conn_close<M: Metrics>(
744 &mut self, qconn: &mut QuicheConnection, metrics: &M,
745 connection_result: &QuicResult<()>,
746 ) {
747 }
748}
749
750/// A command to execute on a [quiche::Connection] in the context of an
751/// [`ApplicationOverQuic`].
752///
753/// We expect most [`ApplicationOverQuic`] implementations (such as
754/// [H3Driver](crate::http3::driver::H3Driver)) will provide some way to submit
755/// actions for them to take, for example via a channel. This enum may be
756/// accepted as part of those actions to inspect or alter the state of the
757/// underlying connection.
758pub enum QuicCommand {
759 /// Close the connection with the given parameters.
760 ///
761 /// Some packets may still be sent after this command has been executed, so
762 /// the worker task may continue running for a bit. See
763 /// [`quiche::Connection::close`] for details.
764 ConnectionClose(ConnectionShutdownBehaviour),
765 /// Execute a custom callback on the connection.
766 Custom(Box<dyn FnOnce(&mut QuicheConnection) + Send + 'static>),
767 /// Collect the current [`SocketStats`] from the connection.
768 ///
769 /// Unlike [`QuicConnection::stats()`], these statistics are not cached and
770 /// instead are retrieved right before the command is executed.
771 Stats(Box<dyn FnOnce(datagram_socket::SocketStats) + Send + 'static>),
772 /// Collect the current [`QuicConnectionStats`] from the connection.
773 ///
774 /// These statistics are not cached and instead are retrieved right before
775 /// the command is executed.
776 ConnectionStats(Box<dyn FnOnce(QuicConnectionStats) + Send + 'static>),
777}
778
779impl QuicCommand {
780 /// Consume the command and perform its operation on `qconn`.
781 ///
782 /// This method should be called by [`ApplicationOverQuic`] implementations
783 /// when they receive a [`QuicCommand`] to execute.
784 pub fn execute(self, qconn: &mut QuicheConnection) {
785 match self {
786 Self::ConnectionClose(behavior) => {
787 let ConnectionShutdownBehaviour {
788 send_application_close,
789 error_code,
790 reason,
791 } = behavior;
792
793 let _ = qconn.close(send_application_close, error_code, &reason);
794 },
795 Self::Custom(f) => {
796 (f)(qconn);
797 },
798 Self::Stats(callback) => {
799 let stats_pair = QuicConnectionStats::from_conn(qconn);
800 (callback)(stats_pair.as_socket_stats());
801 },
802 Self::ConnectionStats(callback) => {
803 let stats_pair = QuicConnectionStats::from_conn(qconn);
804 (callback)(stats_pair);
805 },
806 }
807 }
808}
809
810impl fmt::Debug for QuicCommand {
811 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
812 match self {
813 Self::ConnectionClose(b) =>
814 f.debug_tuple("ConnectionClose").field(b).finish(),
815 Self::Custom(_) => f.debug_tuple("Custom").finish_non_exhaustive(),
816 Self::Stats(_) => f.debug_tuple("Stats").finish_non_exhaustive(),
817 Self::ConnectionStats(_) =>
818 f.debug_tuple("ConnectionStats").finish_non_exhaustive(),
819 }
820 }
821}
822
823/// Parameters to close a [quiche::Connection].
824///
825/// The connection will use these parameters for the `CONNECTION_CLOSE` frame
826/// it sends to its peer.
827#[derive(Debug, Clone)]
828pub struct ConnectionShutdownBehaviour {
829 /// Whether to send an application close or a regular close to the peer.
830 ///
831 /// If this is true but the connection is not in a state where it is safe to
832 /// send an application error (not established nor in early data), in
833 /// accordance with [RFC 9000](https://www.rfc-editor.org/rfc/rfc9000.html#section-10.2.3-3), the
834 /// error code is changed to `APPLICATION_ERROR` and the reason phrase is
835 /// cleared.
836 pub send_application_close: bool,
837 /// The [QUIC][proto-err] or [application-level][app-err] error code to send
838 /// to the peer.
839 ///
840 /// [proto-err]: https://www.rfc-editor.org/rfc/rfc9000.html#section-20.1
841 /// [app-err]: https://www.rfc-editor.org/rfc/rfc9000.html#section-20.2
842 pub error_code: u64,
843 /// The reason phrase to send to the peer.
844 pub reason: Vec<u8>,
845}