(
&self,
conn: Connection,
_tokio_metrics_intervals: impl Iterator<Item = TaskMetrics> + Send + 'static,
)
| 671 | const NAME: &'static str = "http"; |
| 672 | |
| 673 | fn handle_connection( |
| 674 | &self, |
| 675 | conn: Connection, |
| 676 | _tokio_metrics_intervals: impl Iterator<Item = TaskMetrics> + Send + 'static, |
| 677 | ) -> ConnectionHandler { |
| 678 | let router = self.router.clone(); |
| 679 | let tls_context = self.tls.clone(); |
| 680 | let mut conn = TokioIo::new(conn); |
| 681 | |
| 682 | Box::pin(async { |
| 683 | let direct_peer_addr = conn.inner().peer_addr().context("fetching peer addr")?; |
| 684 | let peer_addr = conn |
| 685 | .inner_mut() |
| 686 | .take_proxy_header_address() |
| 687 | .await |
| 688 | .map(|a| a.source) |
| 689 | .unwrap_or(direct_peer_addr); |
| 690 | |
| 691 | let (conn, conn_protocol) = match tls_context { |
| 692 | Some(tls_context) => { |
| 693 | let mut ssl_stream = SslStream::new(Ssl::new(&tls_context.get())?, conn)?; |
| 694 | if let Err(e) = Pin::new(&mut ssl_stream).accept().await { |
| 695 | let _ = ssl_stream.get_mut().inner_mut().shutdown().await; |
| 696 | return Err(e.into()); |
| 697 | } |
| 698 | (MaybeHttpsStream::Https(ssl_stream), ConnProtocol::Https) |
| 699 | } |
| 700 | _ => (MaybeHttpsStream::Http(conn), ConnProtocol::Http), |
| 701 | }; |
| 702 | let mut make_tower_svc = router |
| 703 | .layer(Extension(conn_protocol)) |
| 704 | .into_make_service_with_connect_info::<SocketAddr>(); |
| 705 | let tower_svc = make_tower_svc.call(peer_addr).await.unwrap(); |
| 706 | let hyper_svc = hyper::service::service_fn(|req| tower_svc.clone().call(req)); |
| 707 | let http = hyper::server::conn::http1::Builder::new(); |
| 708 | http.serve_connection(conn, hyper_svc) |
| 709 | .with_upgrades() |
| 710 | .err_into() |
| 711 | .await |
| 712 | }) |
| 713 | } |
| 714 | } |
| 715 | |
| 716 | pub async fn handle_leader_status( |
nothing calls this directly
no test coverage detected