aster_forge_http/
reqwest_client.rs

1use super::read_reqwest_body_limited;
2
3/// Errors produced while adapting a buffered `http` request to reqwest.
4#[derive(Debug, thiserror::Error)]
5pub enum BufferedHttpError {
6    /// The request could not be converted or sent.
7    #[error("outbound HTTP transport failed")]
8    Transport(#[source] reqwest::Error),
9    /// The response body could not be read within the configured bound.
10    #[error("outbound HTTP response body failed: {0}")]
11    ResponseBody(String),
12    /// The buffered response could not be assembled.
13    #[error("outbound HTTP response could not be assembled")]
14    ResponseBuild(#[source] http::Error),
15}
16
17/// Executes an in-memory HTTP request with reqwest and buffers a strictly bounded response body.
18///
19/// Redirect and timeout behavior come from the supplied reqwest client. Callers retain ownership of
20/// endpoint policy, user agent, status handling, and product error mapping.
21///
22/// # Errors
23///
24/// Returns a transport error when the request cannot be converted or sent, a response-body error
25/// when the configured limit is exceeded or streaming fails, or a response-build error when the
26/// buffered response cannot be reconstructed.
27pub async fn execute_reqwest_buffered_limited(
28    client: &reqwest::Client,
29    request: http::Request<Vec<u8>>,
30    max_response_bytes: usize,
31) -> Result<http::Response<Vec<u8>>, BufferedHttpError> {
32    let request = reqwest::Request::try_from(request).map_err(BufferedHttpError::Transport)?;
33    let response = client
34        .execute(request)
35        .await
36        .map_err(BufferedHttpError::Transport)?;
37    let status = response.status();
38    let version = response.version();
39    let headers = response.headers().clone();
40    let body = read_reqwest_body_limited(
41        response,
42        "outbound HTTP response body",
43        max_response_bytes,
44        BufferedHttpError::ResponseBody,
45    )
46    .await?;
47
48    let mut builder = http::Response::builder().status(status).version(version);
49    for (name, value) in &headers {
50        builder = builder.header(name, value);
51    }
52    builder.body(body).map_err(BufferedHttpError::ResponseBuild)
53}
54
55#[cfg(test)]
56mod tests {
57    use tokio::io::{AsyncReadExt, AsyncWriteExt};
58
59    use super::{BufferedHttpError, execute_reqwest_buffered_limited};
60
61    async fn spawn_response(response: &'static [u8]) -> std::net::SocketAddr {
62        let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0))
63            .await
64            .expect("test listener should bind");
65        let address = listener
66            .local_addr()
67            .expect("test listener should expose address");
68        tokio::spawn(async move {
69            let (mut socket, _) = listener
70                .accept()
71                .await
72                .expect("test server should accept request");
73            let mut buffer = [0_u8; 1024];
74            let _ = socket
75                .read(&mut buffer)
76                .await
77                .expect("test server should read request");
78            socket
79                .write_all(response)
80                .await
81                .expect("test server should write response");
82        });
83        address
84    }
85
86    #[tokio::test]
87    async fn buffered_client_preserves_status_headers_and_body() {
88        let address = spawn_response(
89            b"HTTP/1.1 202 Accepted\r\nContent-Length: 4\r\nX-Test: yes\r\nConnection: close\r\n\r\nbody",
90        )
91        .await;
92        let request = http::Request::get(format!("http://{address}/token"))
93            .body(Vec::new())
94            .expect("test request should build");
95
96        let response = execute_reqwest_buffered_limited(&reqwest::Client::new(), request, 4)
97            .await
98            .expect("bounded request should succeed");
99
100        assert_eq!(response.status(), http::StatusCode::ACCEPTED);
101        assert_eq!(response.headers()["x-test"], "yes");
102        assert_eq!(response.body(), b"body");
103    }
104
105    #[tokio::test]
106    async fn buffered_client_rejects_chunked_body_over_limit() {
107        let address = spawn_response(
108            b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n14\r\nsecret-response-body\r\n0\r\n\r\n",
109        )
110        .await;
111        let request = http::Request::get(format!("http://{address}/discovery"))
112            .body(Vec::new())
113            .expect("test request should build");
114
115        let error = execute_reqwest_buffered_limited(&reqwest::Client::new(), request, 4)
116            .await
117            .expect_err("oversized body should fail");
118
119        assert!(matches!(error, BufferedHttpError::ResponseBody(_)));
120        assert!(error.to_string().contains("exceeds 4 bytes limit"));
121        assert!(!error.to_string().contains("secret-response-body"));
122        assert!(!format!("{error:?}").contains("secret-response-body"));
123    }
124}