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}