serve: Extract serve_connection out of handle_connection

This commit is contained in:
Jonas Platte
2025-04-26 23:15:20 +02:00
parent a697d5badd
commit 2bb85477ff
+46 -28
View File
@@ -15,7 +15,10 @@ use hyper::body::Incoming;
use hyper_util::rt::{TokioExecutor, TokioIo}; use hyper_util::rt::{TokioExecutor, TokioIo};
#[cfg(any(feature = "http1", feature = "http2"))] #[cfg(any(feature = "http1", feature = "http2"))]
use hyper_util::{server::conn::auto::Builder, service::TowerToHyperService}; use hyper_util::{server::conn::auto::Builder, service::TowerToHyperService};
use tokio::sync::watch; use tokio::{
io::{AsyncRead, AsyncWrite},
sync::watch,
};
use tower::ServiceExt as _; use tower::ServiceExt as _;
use tower_service::Service; use tower_service::Service;
@@ -370,46 +373,61 @@ async fn handle_connection<L, M, S>(
.await .await
.unwrap_or_else(|err| match err {}); .unwrap_or_else(|err| match err {});
let tower_service = make_service let conn_service = make_service
.call(IncomingStream { .call(IncomingStream {
io: &io, io: &io,
remote_addr, remote_addr,
}) })
.await .await
.unwrap_or_else(|err| match err {}) .unwrap_or_else(|err| match err {});
.map_request(|req: Request<Incoming>| req.map(Body::new));
let hyper_service = TowerToHyperService::new(tower_service); tokio::spawn(serve_connection(
let signal_tx = signal_tx.clone(); io,
let close_rx = close_rx.clone(); conn_service,
signal_tx.clone(),
close_rx.clone(),
));
}
tokio::spawn(async move { async fn serve_connection<I, S>(
#[allow(unused_mut)] io: TokioIo<I>,
let mut builder = Builder::new(TokioExecutor::new()); conn_service: S,
// CONNECT protocol needed for HTTP/2 websockets signal_tx: watch::Sender<()>,
#[cfg(feature = "http2")] close_rx: watch::Receiver<()>,
builder.http2().enable_connect_protocol(); ) where
I: AsyncRead + AsyncWrite + Unpin + Send + 'static,
S: Service<Request, Response = Response, Error = Infallible> + Clone + Send + 'static,
S::Future: Send,
{
let hyper_service = TowerToHyperService::new(
conn_service.map_request(|req: Request<Incoming>| req.map(Body::new)),
);
let mut conn = pin!(builder.serve_connection_with_upgrades(io, hyper_service)); #[allow(unused_mut)]
let mut signal_closed = pin!(signal_tx.closed().fuse()); let mut builder = Builder::new(TokioExecutor::new());
// CONNECT protocol needed for HTTP/2 websockets
#[cfg(feature = "http2")]
builder.http2().enable_connect_protocol();
loop { let mut conn = pin!(builder.serve_connection_with_upgrades(io, hyper_service));
tokio::select! { let mut signal_closed = pin!(signal_tx.closed().fuse());
result = conn.as_mut() => {
if let Err(_err) = result { loop {
trace!("failed to serve connection: {_err:#}"); tokio::select! {
} result = conn.as_mut() => {
break; if let Err(_err) = result {
} trace!("failed to serve connection: {_err:#}");
_ = &mut signal_closed => {
trace!("signal received in task, starting graceful shutdown");
conn.as_mut().graceful_shutdown();
} }
break;
}
_ = &mut signal_closed => {
trace!("signal received in task, starting graceful shutdown");
conn.as_mut().graceful_shutdown();
} }
} }
}
drop(close_rx); drop(close_rx);
});
} }
/// An incoming stream. /// An incoming stream.