postrust_proxy/vendored/
forwarder.rs1use crate::config::Backend;
6use crate::vendored::hyper_ext::ProxyBody;
7use crate::vendored::types::ProxyError;
8use hyper::body::Incoming;
9use hyper::{Request, Response};
10use hyper_util::client::legacy::Client;
11use hyper_util::rt::TokioExecutor as HyperTokioExecutor;
12use std::time::Duration;
13
14type HttpConnector = hyper_util::client::legacy::connect::HttpConnector;
16
17pub struct ForwarderClient {
19 http_client: Client<HttpConnector, ProxyBody>,
21 timeout: Duration,
23}
24
25impl ForwarderClient {
26 pub fn new(timeout: Duration) -> Self {
28 let http_connector = HttpConnector::new();
29 let http_client = Client::builder(HyperTokioExecutor::new())
30 .pool_idle_timeout(Duration::from_secs(30))
31 .pool_max_idle_per_host(32)
32 .build(http_connector);
33
34 Self {
35 http_client,
36 timeout,
37 }
38 }
39
40 pub async fn forward(
42 &self,
43 backend: &Backend,
44 mut request: Request<ProxyBody>,
45 ) -> Result<Response<Incoming>, ProxyError> {
46 let uri = request.uri();
48 let path_and_query = uri.path_and_query().map(|pq| pq.as_str()).unwrap_or("/");
49
50 let upstream_uri = format!("{}://{}{}", backend.scheme, backend.address, path_and_query);
51
52 *request.uri_mut() = upstream_uri
54 .parse()
55 .map_err(|e| ProxyError::Request(format!("Invalid upstream URI: {}", e)))?;
56
57 let response = tokio::time::timeout(self.timeout, self.http_client.request(request))
59 .await
60 .map_err(|_| ProxyError::Timeout)?
61 .map_err(|e| ProxyError::Connection(e.to_string()))?;
62
63 Ok(response)
64 }
65}
66
67impl Default for ForwarderClient {
68 fn default() -> Self {
69 Self::new(Duration::from_secs(30))
70 }
71}