use crate::config::Backend;
use crate::vendored::hyper_ext::ProxyBody;
use crate::vendored::types::ProxyError;
use hyper::body::Incoming;
use hyper::{Request, Response};
use hyper_util::client::legacy::Client;
use hyper_util::rt::TokioExecutor as HyperTokioExecutor;
use std::time::Duration;
type HttpConnector = hyper_util::client::legacy::connect::HttpConnector;
pub struct ForwarderClient {
http_client: Client<HttpConnector, ProxyBody>,
timeout: Duration,
}
impl ForwarderClient {
pub fn new(timeout: Duration) -> Self {
let http_connector = HttpConnector::new();
let http_client = Client::builder(HyperTokioExecutor::new())
.pool_idle_timeout(Duration::from_secs(30))
.pool_max_idle_per_host(32)
.build(http_connector);
Self {
http_client,
timeout,
}
}
pub async fn forward(
&self,
backend: &Backend,
mut request: Request<ProxyBody>,
) -> Result<Response<Incoming>, ProxyError> {
let uri = request.uri();
let path_and_query = uri.path_and_query().map(|pq| pq.as_str()).unwrap_or("/");
let upstream_uri = format!("{}://{}{}", backend.scheme, backend.address, path_and_query);
*request.uri_mut() = upstream_uri
.parse()
.map_err(|e| ProxyError::Request(format!("Invalid upstream URI: {}", e)))?;
let response = tokio::time::timeout(self.timeout, self.http_client.request(request))
.await
.map_err(|_| ProxyError::Timeout)?
.map_err(|e| ProxyError::Connection(e.to_string()))?;
Ok(response)
}
}
impl Default for ForwarderClient {
fn default() -> Self {
Self::new(Duration::from_secs(30))
}
}