aster_forge_http/
reqwest_client.rs1use super::read_reqwest_body_limited;
2
3#[derive(Debug, thiserror::Error)]
5pub enum BufferedHttpError {
6 #[error("outbound HTTP transport failed")]
8 Transport(#[source] reqwest::Error),
9 #[error("outbound HTTP response body failed: {0}")]
11 ResponseBody(String),
12 #[error("outbound HTTP response could not be assembled")]
14 ResponseBuild(#[source] http::Error),
15}
16
17pub 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}