Skip to main content

aep_agent/
transport.rs

1use aep_core::{HttpRequest, HttpResponse, HttpTransport, TransportError};
2use async_trait::async_trait;
3use futures::StreamExt as _;
4
5pub struct ReqwestTransport {
6    client: reqwest::Client,
7    maximum_response_bytes: usize,
8}
9
10impl ReqwestTransport {
11    pub fn new(
12        maximum_response_bytes: usize,
13        timeout: std::time::Duration,
14    ) -> Result<Self, TransportError> {
15        let client = reqwest::Client::builder()
16            .redirect(reqwest::redirect::Policy::none())
17            .timeout(timeout)
18            .build()
19            .map_err(|error| TransportError::with_source("build HTTP client", error))?;
20        Ok(Self {
21            client,
22            maximum_response_bytes,
23        })
24    }
25}
26
27#[async_trait]
28impl HttpTransport for ReqwestTransport {
29    async fn send(&self, request: HttpRequest) -> Result<HttpResponse, TransportError> {
30        let mut builder = self
31            .client
32            .request(request.method, request.url.clone())
33            .headers(request.headers);
34        if !request.body.is_empty() {
35            builder = builder.body(request.body);
36        }
37        let response = builder
38            .send()
39            .await
40            .map_err(|error| TransportError::with_source("send HTTP request", error))?;
41        let status = response.status();
42        let final_url = response.url().clone();
43        let headers = response.headers().clone();
44        let mut stream = response.bytes_stream();
45        let mut body = Vec::new();
46        while let Some(chunk) = stream.next().await {
47            let chunk =
48                chunk.map_err(|error| TransportError::with_source("read HTTP response", error))?;
49            if body.len().saturating_add(chunk.len()) > self.maximum_response_bytes {
50                return Err(TransportError::new(
51                    "HTTP response exceeds the configured limit",
52                ));
53            }
54            body.extend_from_slice(&chunk);
55        }
56        Ok(HttpResponse {
57            status,
58            final_url,
59            headers,
60            body,
61        })
62    }
63}