tokio_quiche/quic/io/
connection_stage.rs1use std::fmt::Debug;
28use std::ops::ControlFlow;
29use std::sync::Arc;
30use std::time::Instant;
31
32use tokio::sync::mpsc;
33
34use crate::quic::connection::ApplicationOverQuic;
35use crate::quic::connection::HandshakeError;
36use crate::quic::connection::HandshakeInfo;
37use crate::quic::connection::Incoming;
38use crate::quic::connection::QuicConnectionStatsShared;
39use crate::quic::hooks::ConnectionHook;
40use crate::quic::QuicheConnection;
41use crate::QuicResult;
42
43pub trait ConnectionStage: Send + Debug {
56 fn on_read<A: ApplicationOverQuic>(
57 &mut self, _received_packets: bool, _qconn: &mut QuicheConnection,
58 _ctx: &mut ConnectionStageContext<A>,
59 ) -> QuicResult<()> {
60 Ok(())
61 }
62
63 fn on_flush<A: ApplicationOverQuic>(
64 &mut self, _qconn: &mut QuicheConnection,
65 _ctx: &mut ConnectionStageContext<A>,
66 ) -> ControlFlow<QuicResult<()>> {
67 ControlFlow::Continue(())
68 }
69
70 fn wait_deadline(&mut self) -> Option<Instant> {
71 None
72 }
73
74 fn post_wait(
75 &self, _qconn: &mut QuicheConnection,
76 ) -> ControlFlow<QuicResult<()>> {
77 ControlFlow::Continue(())
78 }
79}
80
81pub struct ConnectionStageContext<A> {
83 pub in_pkt: Option<Incoming>,
84 pub application: A,
85 pub incoming_pkt_receiver: mpsc::Receiver<Incoming>,
86 pub stats: QuicConnectionStatsShared,
87 pub connection_hook: Option<Arc<dyn ConnectionHook + Send + Sync + 'static>>,
88}
89
90#[derive(Debug)]
91pub struct Handshake {
92 pub handshake_info: HandshakeInfo,
93}
94
95impl Handshake {
96 fn check_handshake_timeout_expired(
97 &self, conn: &mut QuicheConnection,
98 ) -> QuicResult<()> {
99 if self.handshake_info.is_expired() {
100 let _ = conn.close(
101 false,
102 quiche::WireErrorCode::ApplicationError as u64,
103 &[],
104 );
105 return Err(HandshakeError::Timeout.into());
106 }
107
108 Ok(())
109 }
110}
111
112impl ConnectionStage for Handshake {
113 fn on_flush<A: ApplicationOverQuic>(
114 &mut self, qconn: &mut QuicheConnection,
115 _ctx: &mut ConnectionStageContext<A>,
116 ) -> ControlFlow<QuicResult<()>> {
117 if qconn.is_established() || qconn.is_in_early_data() {
120 ControlFlow::Break(Ok(()))
121 } else {
122 ControlFlow::Continue(())
123 }
124 }
125
126 fn wait_deadline(&mut self) -> Option<Instant> {
127 self.handshake_info.deadline()
128 }
129
130 fn post_wait(
131 &self, qconn: &mut QuicheConnection,
132 ) -> ControlFlow<QuicResult<()>> {
133 match self.check_handshake_timeout_expired(qconn) {
134 Ok(_) => ControlFlow::Continue(()),
135 Err(e) => ControlFlow::Break(Err(e)),
136 }
137 }
138}
139
140#[derive(Debug)]
141pub struct RunningApplication;
142
143impl ConnectionStage for RunningApplication {
144 fn on_read<A: ApplicationOverQuic>(
145 &mut self, received_packets: bool, qconn: &mut QuicheConnection,
146 ctx: &mut ConnectionStageContext<A>,
147 ) -> QuicResult<()> {
148 if ctx.application.should_act() {
149 if received_packets {
150 ctx.application.process_reads(qconn)?;
151 }
152
153 if qconn.is_established() {
154 ctx.application.process_writes(qconn)?;
155 }
156 }
157
158 Ok(())
159 }
160}
161
162#[derive(Debug)]
163pub struct Close {
164 pub work_loop_result: QuicResult<()>,
165}
166
167impl ConnectionStage for Close {}