Skip to content

Commit e40d915

Browse files
authored
fix(http): preserve error responses for retry classification (#3627)
Copilot-Session: 0bc45e17-c983-47d7-9a97-f6b1909ceefb
1 parent 42dd39a commit e40d915

3 files changed

Lines changed: 106 additions & 4 deletions

File tree

opentelemetry-http/CHANGELOG.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,11 @@
22

33
## vNext
44

5+
- Return HTTP error responses from the built-in reqwest and hyper clients instead
6+
of converting 4xx and 5xx statuses into transport errors. This preserves the
7+
response status and headers for exporter retry classification.
8+
If your code relied on `send_bytes` returning `Err` for non-success statuses,
9+
call `ResponseExt::error_for_status()` on the response instead.
510
- **Breaking** Removed `reqwest-rustls-webpki-roots` feature. The `webpki-roots` cargo feature was
611
removed from `reqwest` in v0.13.0. Use `reqwest-rustls` instead, which now correctly enables
712
`reqwest/rustls` (platform native trust roots). To use Mozilla's embedded CA bundle, construct a

opentelemetry-http/Cargo.toml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,10 @@ opentelemetry = { workspace = true, features = ["trace"] }
2929
reqwest = { workspace = true, optional = true }
3030
tokio = { workspace = true, features = ["time"], optional = true }
3131

32+
[dev-dependencies]
33+
futures-executor = { workspace = true }
34+
tokio = { workspace = true, features = ["net", "rt", "time"] }
35+
3236
[lints]
3337
workspace = true
3438

opentelemetry-http/src/lib.rs

Lines changed: 97 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ mod reqwest {
101101
async fn send_bytes(&self, request: Request<Bytes>) -> Result<Response<Bytes>, HttpError> {
102102
otel_debug!(name: "ReqwestClient.Send");
103103
let request = request.try_into()?;
104-
let mut response = self.execute(request).await?.error_for_status()?;
104+
let mut response = self.execute(request).await?;
105105
let headers = std::mem::take(response.headers_mut());
106106
let mut http_response = Response::builder()
107107
.status(response.status())
@@ -119,7 +119,7 @@ mod reqwest {
119119
async fn send_bytes(&self, request: Request<Bytes>) -> Result<Response<Bytes>, HttpError> {
120120
otel_debug!(name: "ReqwestBlockingClient.Send");
121121
let request = request.try_into()?;
122-
let mut response = self.execute(request)?.error_for_status()?;
122+
let mut response = self.execute(request)?;
123123
let headers = std::mem::take(response.headers_mut());
124124
let mut http_response = Response::builder()
125125
.status(response.status())
@@ -134,7 +134,6 @@ mod reqwest {
134134
#[cfg(feature = "hyper")]
135135
pub mod hyper {
136136
use super::{async_trait, Bytes, HttpClient, HttpError, Request, Response};
137-
use crate::ResponseExt;
138137
use http::HeaderValue;
139138
use http_body_util::{BodyExt, Full};
140139
use hyper::body::{Body as HttpBody, Frame};
@@ -207,7 +206,7 @@ pub mod hyper {
207206
.body(response.into_body().collect().await?.to_bytes())?;
208207
*http_response.headers_mut() = headers;
209208

210-
Ok(http_response.error_for_status()?)
209+
Ok(http_response)
211210
}
212211
}
213212

@@ -258,6 +257,21 @@ impl<T> ResponseExt for Response<T> {
258257
mod tests {
259258
use super::*;
260259
use http::HeaderValue;
260+
#[cfg(all(
261+
any(feature = "hyper", feature = "reqwest", feature = "reqwest-blocking"),
262+
not(target_arch = "wasm32")
263+
))]
264+
use std::io::{Read, Write};
265+
#[cfg(all(
266+
any(feature = "hyper", feature = "reqwest", feature = "reqwest-blocking"),
267+
not(target_arch = "wasm32")
268+
))]
269+
use std::net::{SocketAddr, TcpListener};
270+
#[cfg(all(
271+
any(feature = "hyper", feature = "reqwest", feature = "reqwest-blocking"),
272+
not(target_arch = "wasm32")
273+
))]
274+
use std::thread::JoinHandle;
261275

262276
#[test]
263277
fn http_headers_get() {
@@ -309,6 +323,85 @@ mod tests {
309323
assert!(got.contains(&"headername2"));
310324
}
311325

326+
#[cfg(all(
327+
any(feature = "hyper", feature = "reqwest", feature = "reqwest-blocking"),
328+
not(target_arch = "wasm32")
329+
))]
330+
fn spawn_error_response_server() -> (SocketAddr, JoinHandle<()>) {
331+
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
332+
let address = listener.local_addr().unwrap();
333+
let server = std::thread::spawn(move || {
334+
let (mut stream, _) = listener.accept().unwrap();
335+
let mut request = [0; 1024];
336+
let _ = stream.read(&mut request).unwrap();
337+
stream
338+
.write_all(
339+
b"HTTP/1.1 429 Too Many Requests\r\n\
340+
Retry-After: 7\r\n\
341+
Content-Length: 0\r\n\
342+
Connection: close\r\n\r\n",
343+
)
344+
.unwrap();
345+
});
346+
(address, server)
347+
}
348+
349+
#[cfg(all(feature = "reqwest-blocking", not(target_arch = "wasm32")))]
350+
#[test]
351+
fn reqwest_blocking_preserves_error_response_status_and_headers() {
352+
let (address, server) = spawn_error_response_server();
353+
let client = ::reqwest::blocking::Client::new();
354+
let request = Request::post(format!("http://{address}/v1/traces"))
355+
.body(Bytes::new())
356+
.unwrap();
357+
let response = futures_executor::block_on(client.send_bytes(request)).unwrap();
358+
359+
server.join().unwrap();
360+
assert_eq!(response.status(), http::StatusCode::TOO_MANY_REQUESTS);
361+
assert_eq!(response.headers().get("retry-after").unwrap(), "7");
362+
}
363+
364+
#[cfg(all(feature = "reqwest", not(target_arch = "wasm32")))]
365+
#[test]
366+
fn reqwest_async_preserves_error_response_status_and_headers() {
367+
let (address, server) = spawn_error_response_server();
368+
let client = ::reqwest::Client::new();
369+
let request = Request::post(format!("http://{address}/v1/traces"))
370+
.body(Bytes::new())
371+
.unwrap();
372+
let runtime = tokio::runtime::Builder::new_current_thread()
373+
.enable_all()
374+
.build()
375+
.unwrap();
376+
let response = runtime.block_on(client.send_bytes(request)).unwrap();
377+
378+
server.join().unwrap();
379+
assert_eq!(response.status(), http::StatusCode::TOO_MANY_REQUESTS);
380+
assert_eq!(response.headers().get("retry-after").unwrap(), "7");
381+
}
382+
383+
#[cfg(all(feature = "hyper", not(target_arch = "wasm32")))]
384+
#[test]
385+
fn hyper_preserves_error_response_status_and_headers() {
386+
let (address, server) = spawn_error_response_server();
387+
let client = crate::hyper::HyperClient::with_default_connector(
388+
std::time::Duration::from_secs(2),
389+
None,
390+
);
391+
let request = Request::post(format!("http://{address}/v1/traces"))
392+
.body(Bytes::new())
393+
.unwrap();
394+
let runtime = tokio::runtime::Builder::new_current_thread()
395+
.enable_all()
396+
.build()
397+
.unwrap();
398+
let response = runtime.block_on(client.send_bytes(request)).unwrap();
399+
400+
server.join().unwrap();
401+
assert_eq!(response.status(), http::StatusCode::TOO_MANY_REQUESTS);
402+
assert_eq!(response.headers().get("retry-after").unwrap(), "7");
403+
}
404+
312405
#[test]
313406
fn http_headers_reserve() {
314407
let mut carrier = http::HeaderMap::new();

0 commit comments

Comments
 (0)