From c1162d31919a20ce693fecbc0abbf1f3c833bf4a Mon Sep 17 00:00:00 2001 From: Jonas Platte Date: Sat, 26 Apr 2025 21:37:14 +0200 Subject: [PATCH] Extract handle_connection out of do_serve --- axum/src/serve/mod.rs | 115 ++++++++++++++++++++++++------------------ 1 file changed, 66 insertions(+), 49 deletions(-) diff --git a/axum/src/serve/mod.rs b/axum/src/serve/mod.rs index 9a0d96a4..d003bdc5 100644 --- a/axum/src/serve/mod.rs +++ b/axum/src/serve/mod.rs @@ -304,55 +304,7 @@ where } }; - let io = TokioIo::new(io); - - trace!("connection {remote_addr:?} accepted"); - - make_service - .ready() - .await - .unwrap_or_else(|err| match err {}); - - let tower_service = make_service - .call(IncomingStream { - io: &io, - remote_addr, - }) - .await - .unwrap_or_else(|err| match err {}) - .map_request(|req: Request| req.map(Body::new)); - - let hyper_service = TowerToHyperService::new(tower_service); - let signal_tx = signal_tx.clone(); - let close_rx = close_rx.clone(); - - tokio::spawn(async move { - #[allow(unused_mut)] - let mut builder = Builder::new(TokioExecutor::new()); - // CONNECT protocol needed for HTTP/2 websockets - #[cfg(feature = "http2")] - builder.http2().enable_connect_protocol(); - - let mut conn = pin!(builder.serve_connection_with_upgrades(io, hyper_service)); - let mut signal_closed = pin!(signal_tx.closed().fuse()); - - loop { - tokio::select! { - result = conn.as_mut() => { - if let Err(_err) = result { - trace!("failed to serve connection: {_err:#}"); - } - break; - } - _ = &mut signal_closed => { - trace!("signal received in task, starting graceful shutdown"); - conn.as_mut().graceful_shutdown(); - } - } - } - - drop(close_rx); - }); + handle_connection(&mut make_service, &signal_tx, &close_rx, io, remote_addr).await; } drop(close_rx); @@ -365,6 +317,71 @@ where close_tx.closed().await; } +async fn handle_connection( + make_service: &mut M, + signal_tx: &watch::Sender<()>, + close_rx: &watch::Receiver<()>, + io: ::Io, + remote_addr: ::Addr, +) where + L: Listener, + L::Addr: Debug, + M: for<'a> Service, Error = Infallible, Response = S> + Send + 'static, + for<'a> >>::Future: Send, + S: Service + Clone + Send + 'static, + S::Future: Send, +{ + let io = TokioIo::new(io); + + trace!("connection {remote_addr:?} accepted"); + + make_service + .ready() + .await + .unwrap_or_else(|err| match err {}); + + let tower_service = make_service + .call(IncomingStream { + io: &io, + remote_addr, + }) + .await + .unwrap_or_else(|err| match err {}) + .map_request(|req: Request| req.map(Body::new)); + + let hyper_service = TowerToHyperService::new(tower_service); + let signal_tx = signal_tx.clone(); + let close_rx = close_rx.clone(); + + tokio::spawn(async move { + #[allow(unused_mut)] + let mut builder = Builder::new(TokioExecutor::new()); + // CONNECT protocol needed for HTTP/2 websockets + #[cfg(feature = "http2")] + builder.http2().enable_connect_protocol(); + + let mut conn = pin!(builder.serve_connection_with_upgrades(io, hyper_service)); + let mut signal_closed = pin!(signal_tx.closed().fuse()); + + loop { + tokio::select! { + result = conn.as_mut() => { + if let Err(_err) = result { + trace!("failed to serve connection: {_err:#}"); + } + break; + } + _ = &mut signal_closed => { + trace!("signal received in task, starting graceful shutdown"); + conn.as_mut().graceful_shutdown(); + } + } + } + + drop(close_rx); + }); +} + /// An incoming stream. /// /// Used with [`serve`] and [`IntoMakeServiceWithConnectInfo`].