Skip to content

Commit 18720ff

Browse files
committed
tokio-quiche: add ServerH3Event::Headers
Priority signals take the form of HTTP headers and/or frames. Some clients only send frames, while some clients can both. At the point in time a server reads request HEADERS, a client may sent a PRIORITY_UPDATE frame that should take precedence. With this change, we read take the last PRIORITY_UPDATE processed by quiche, and pass it alongside IncomingH3Headers in a new ServerH3Event:Headers event so that Tokio-Quiche-based servers can use it if they wish.
1 parent a67d4da commit 18720ff

5 files changed

Lines changed: 98 additions & 38 deletions

File tree

‎tokio-quiche/examples/async_http3_server/server.rs‎

Lines changed: 29 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ use tokio_quiche::http3::driver::H3Event;
4444
use tokio_quiche::http3::driver::IncomingH3Headers;
4545
use tokio_quiche::http3::driver::OutboundFrame;
4646
use tokio_quiche::http3::driver::OutboundFrameSender;
47+
use tokio_quiche::http3::driver::RawPriorityValue;
4748
use tokio_quiche::http3::driver::ServerEventStream;
4849
use tokio_quiche::http3::driver::ServerH3Event;
4950
use tokio_quiche::BoxError;
@@ -102,20 +103,14 @@ where
102103
) -> QuicResult<()> {
103104
loop {
104105
match h3_event_receiver.recv().await {
105-
Some(event) => self.handle_h3_event(event).await?,
106+
Some(event) => self.handle_server_h3_event(event).await?,
106107
None => return Ok(()), /* The sender was dropped, implying
107108
* connection was terminated */
108109
}
109110
}
110111
}
111112

112-
/// Handle a [`ServerH3Event`].
113-
///
114-
/// For simplicity's sake, we only handle a couple of events here.
115-
// TODO(evanrittenhouse): support POST requests
116-
async fn handle_h3_event(&mut self, event: ServerH3Event) -> QuicResult<()> {
117-
let ServerH3Event::Core(event) = event;
118-
113+
fn handle_h3_event(event: H3Event) -> QuicResult<()> {
119114
match event {
120115
// Received an explicit connection level error. Not much to do here.
121116
H3Event::ConnectionError(err) => QuicResult::Err(Box::new(err)),
@@ -130,15 +125,30 @@ where
130125
QuicResult::Err(err)
131126
},
132127

133-
// Received headers for a new stream from the H3Driver.
134-
H3Event::IncomingHeaders(headers) => {
135-
self.handle_incoming_headers(headers).await;
136-
128+
_ => {
129+
log::info!("received unhandled event: {event:?}");
137130
Ok(())
138131
},
132+
}
133+
}
139134

140-
_ => {
141-
log::info!("received unhandled event: {event:?}");
135+
/// Handle a [`ServerH3Event`].
136+
///
137+
/// For simplicity's sake, we only handle a couple of events here.
138+
// TODO(evanrittenhouse): support POST requests
139+
async fn handle_server_h3_event(
140+
&mut self, event: ServerH3Event,
141+
) -> QuicResult<()> {
142+
match event {
143+
ServerH3Event::Core(event) => Self::handle_h3_event(event),
144+
145+
ServerH3Event::Headers {
146+
incoming_headers,
147+
priority,
148+
} => {
149+
// Received headers for a new stream from the H3Driver.
150+
self.handle_incoming_headers(incoming_headers, priority)
151+
.await;
142152
Ok(())
143153
},
144154
}
@@ -149,7 +159,10 @@ where
149159
/// This function transforms the incoming headers into a [`Request`],
150160
/// creating the proper response body if requested. It then spawns a
151161
/// Tokio task which calls the `service_fn` on the [`Request`].
152-
async fn handle_incoming_headers(&mut self, headers: IncomingH3Headers) {
162+
async fn handle_incoming_headers(
163+
&mut self, headers: IncomingH3Headers,
164+
_priority: Option<RawPriorityValue>,
165+
) {
153166
log::info!("received headers: {:?}", &headers);
154167

155168
let IncomingH3Headers {
@@ -168,6 +181,7 @@ where
168181

169182
let service_fn = Arc::clone(&self.service_fn);
170183

184+
// TODO: use the _priority input parameter in request handling
171185
tokio::spawn(async move {
172186
Self::handle_request(service_fn, req, frame_sender).await;
173187
});

‎tokio-quiche/src/http3/driver/mod.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@ pub use self::client::ClientH3Driver;
8484
pub use self::client::ClientH3Event;
8585
pub use self::client::ClientRequestSender;
8686
pub use self::client::NewClientRequest;
87+
pub use self::server::RawPriorityValue;
8788
pub use self::server::ServerEventStream;
8889
pub use self::server::ServerH3Command;
8990
pub use self::server::ServerH3Controller;

‎tokio-quiche/src/http3/driver/server.rs‎

Lines changed: 41 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
2525
// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
2626

27+
use std::ops::Deref;
2728
use std::sync::Arc;
2829

2930
use tokio::sync::mpsc;
@@ -60,15 +61,44 @@ pub type ServerH3Controller = H3Controller<ServerHooks>;
6061
/// transfer data.
6162
pub type ServerEventStream = mpsc::UnboundedReceiver<ServerH3Event>;
6263

64+
#[derive(Clone, Debug)]
65+
pub struct RawPriorityValue(Vec<u8>);
66+
67+
impl From<Vec<u8>> for RawPriorityValue {
68+
fn from(value: Vec<u8>) -> Self {
69+
RawPriorityValue(value)
70+
}
71+
}
72+
73+
impl Deref for RawPriorityValue {
74+
type Target = [u8];
75+
76+
fn deref(&self) -> &Self::Target {
77+
&self.0
78+
}
79+
}
80+
6381
/// Events produced by [ServerH3Driver].
6482
#[derive(Debug)]
6583
pub enum ServerH3Event {
6684
Core(H3Event),
85+
86+
Headers {
87+
incoming_headers: IncomingH3Headers,
88+
/// The latest PRIORITY_UPDATE frame value, if any.
89+
priority: Option<RawPriorityValue>,
90+
},
6791
}
6892

6993
impl From<H3Event> for ServerH3Event {
7094
fn from(ev: H3Event) -> Self {
71-
Self::Core(ev)
95+
match ev {
96+
H3Event::IncomingHeaders(incoming_headers) => Self::Headers {
97+
incoming_headers,
98+
priority: None,
99+
},
100+
_ => Self::Core(ev),
101+
}
72102
}
73103
}
74104

@@ -130,6 +160,12 @@ impl ServerHooks {
130160
stream_ctx.associated_dgram_flow_id = Some(flow_id);
131161
}
132162

163+
let latest_priority_update: Option<RawPriorityValue> = driver
164+
.conn_mut()?
165+
.take_last_priority_update(stream_id)
166+
.ok()
167+
.map(|v| v.into());
168+
133169
let headers = IncomingH3Headers {
134170
stream_id,
135171
headers,
@@ -146,7 +182,10 @@ impl ServerHooks {
146182

147183
driver
148184
.h3_event_sender
149-
.send(H3Event::IncomingHeaders(headers).into())
185+
.send(ServerH3Event::Headers {
186+
incoming_headers: headers,
187+
priority: latest_priority_update,
188+
})
150189
.map_err(|_| H3ConnectionError::ControllerWentAway)?;
151190
driver.hooks.requests += 1;
152191

‎tokio-quiche/tests/fixtures/mod.rs‎

Lines changed: 16 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -133,21 +133,25 @@ pub async fn serve_connection_details(
133133

134134
loop {
135135
select! {
136-
Some(frame) = event_rx.recv() => {
137-
let ServerH3Event::Core(frame) = frame;
138-
match frame {
139-
H3Event::IncomingSettings {..} | H3Event::BodyBytesReceived { .. } | H3Event::StreamClosed { .. } => {},
140-
H3Event::IncomingHeaders(headers) => {
136+
Some(event) = event_rx.recv() => {
137+
match event {
138+
ServerH3Event::Core(event) => {
139+
match event {
140+
H3Event::IncomingSettings {..} | H3Event::BodyBytesReceived { .. } | H3Event::StreamClosed { .. } | H3Event::IncomingHeaders(..) => {},
141+
H3Event::ConnectionError(err) => { break Err(err.into()); }
142+
H3Event::ConnectionShutdown(Some(err)) => { break Err(err.into()); }
143+
_ => unreachable!()
144+
}
145+
}
146+
147+
ServerH3Event::Headers{ incoming_headers, ..} => {
141148
let IncomingH3Headers {
142-
stream_id, headers, send, recv, ..
143-
} = headers;
149+
stream_id, headers, send, recv, ..
150+
} = incoming_headers;
144151

145-
request_counter.fetch_add(1, Ordering::SeqCst);
146-
request_futs.push(handle_forwarded_headers_frame(stream_id, headers, send, recv));
152+
request_counter.fetch_add(1, Ordering::SeqCst);
153+
request_futs.push(handle_forwarded_headers_frame(stream_id, headers, send, recv));
147154
}
148-
H3Event::ConnectionError(err) => { break Err(err.into()); }
149-
H3Event::ConnectionShutdown(Some(err)) => { break Err(err.into()); }
150-
_ => unreachable!()
151155
}
152156
}
153157
Some(_) = request_futs.next() => {}

‎tokio-quiche/tests/integration_tests/headers.rs‎

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -45,12 +45,18 @@ async fn test_additional_headers() {
4545
move |mut h3_conn| async move {
4646
let event_rx = h3_conn.h3_controller.event_receiver_mut();
4747

48-
while let Some(frame) = event_rx.recv().await {
49-
let ServerH3Event::Core(frame) = frame;
48+
while let Some(event) = event_rx.recv().await {
49+
match event {
50+
ServerH3Event::Core(event) => match event {
51+
H3Event::ConnectionShutdown(_) => break,
5052

51-
match frame {
52-
H3Event::IncomingHeaders(headers) => {
53-
let IncomingH3Headers { mut send, .. } = headers;
53+
_ => (),
54+
},
55+
56+
ServerH3Event::Headers {
57+
incoming_headers, ..
58+
} => {
59+
let IncomingH3Headers { mut send, .. } = incoming_headers;
5460

5561
// Send initial headers.
5662
send.send(OutboundFrame::Headers(
@@ -73,10 +79,6 @@ async fn test_additional_headers() {
7379
.await
7480
.unwrap();
7581
},
76-
77-
H3Event::ConnectionShutdown(_) => break,
78-
79-
_ => (),
8082
}
8183
}
8284
},

0 commit comments

Comments
 (0)