From 8655c38f1dfd21ae2580edb260641df1dc596e5c Mon Sep 17 00:00:00 2001 From: Chris Date: Wed, 30 Sep 2026 20:18:28 +0800 Subject: [PATCH] fix(rio): preserve bounded retries after peer EOF (#8272) --- crates/rio/Cargo.toml | 1 + crates/rio/src/http_reader.rs | 108 ++++++++++++++++++++++++++++++++-- 2 files changed, 103 insertions(+), 6 deletions(-) diff --git a/crates/rio/Cargo.toml b/crates/rio/Cargo.toml index c30b83a86..28b23e1b1 100644 --- a/crates/rio/Cargo.toml +++ b/crates/rio/Cargo.toml @@ -63,6 +63,7 @@ arc-swap.workspace = true tokio = { workspace = true, features = ["io-util", "macros", "net", "rt-multi-thread", "sync", "time"] } rand = { workspace = true, features = ["serde"] } http.workspace = true +hyper.workspace = true aes-gcm = { workspace = true, features = ["rand_core"] } crc-fast = { workspace = true } pin-project-lite.workspace = true diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index 3a1e58fe8..8d454b99e 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -758,7 +758,7 @@ fn internode_request_context(method: &Method, url: &str, operation: Option<&'sta /// /// Ordering matters: /// 1. A caller-reported timeout wins outright (`ConnectTimeout`). -/// 2. A typed `io::ErrorKind` anywhere in the source chain wins over any +/// 2. A typed I/O error or HTTP EOF anywhere in the source chain wins over any /// string or body signal — a real `ConnectionRefused` must never be /// mislabeled `DnsResolutionFailed` just because "dns" appears in the text. /// 3. Only for connect-phase failures with no typed kind do we consult the @@ -775,7 +775,7 @@ fn classify_transport_error( return InternodeHttpErrorKind::ConnectTimeout; } - if let Some(kind) = find_io_error_kind_in_chain(err) { + if let Some(kind) = find_transport_error_kind_in_chain(err) { return kind; } @@ -793,10 +793,8 @@ fn classify_transport_error( InternodeHttpErrorKind::Unknown } -/// Walk the error source chain looking for a `std::io::Error` and map its -/// [`io::ErrorKind`] onto our typed classification. Returns `None` when no -/// `io::Error` is present or its kind carries no actionable signal. -fn find_io_error_kind_in_chain(err: &(dyn std::error::Error + 'static)) -> Option { +/// Map typed I/O failures and HTTP EOFs from the error source chain. +fn find_transport_error_kind_in_chain(err: &(dyn std::error::Error + 'static)) -> Option { let mut source: Option<&(dyn std::error::Error + 'static)> = Some(err); while let Some(current) = source { if let Some(io_err) = current.downcast_ref::() { @@ -812,6 +810,12 @@ fn find_io_error_kind_in_chain(err: &(dyn std::error::Error + 'static)) -> Optio _ => {} } } + if current + .downcast_ref::() + .is_some_and(hyper::Error::is_incomplete_message) + { + return Some(InternodeHttpErrorKind::ConnectionReset); + } source = current.source(); } None @@ -2038,6 +2042,98 @@ mod tests { assert_eq!(classify_transport_error(&err, false, false, false), InternodeHttpErrorKind::Unknown); } + #[tokio::test] + async fn peer_close_before_upload_response_remains_retryable() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("peer listener should bind"); + let address = listener.local_addr().expect("peer address should be available"); + let peer = tokio::spawn(async move { + for _ in 0..2 { + let (mut socket, _) = listener.accept().await.expect("peer should receive the upload"); + let mut received = Vec::new(); + let mut buffer = [0_u8; 4096]; + loop { + let count = socket.read(&mut buffer).await.expect("peer should read the request"); + assert_ne!(count, 0, "upload should reach the peer before it closes"); + received.extend_from_slice(&buffer[..count]); + if received + .windows(b"restart-mid-upload".len()) + .any(|part| part == b"restart-mid-upload") + { + break; + } + } + socket + .shutdown() + .await + .expect("peer should close before returning response headers"); + } + }); + let error = reqwest::Client::builder() + .no_proxy() + .build() + .expect("upload client should build") + .put(format!("http://{address}/rustfs/rpc/put_file_stream_v1")) + .body("restart-mid-upload") + .send() + .await + .expect_err("closing the peer before its response should fail the upload"); + let kind = classify_reqwest_error(&error); + assert_eq!(kind, InternodeHttpErrorKind::ConnectionReset); + assert!(kind.is_retryable(), "a peer close during an upload must retain bounded retry eligibility"); + + let mut writer = + HttpWriter::new(format!("http://{address}/rustfs/rpc/put_file_stream_v1"), Method::PUT, HeaderMap::new()) + .await + .expect("internode writer should be created"); + writer.write_all(b"restart-mid-upload").await.expect("upload should start"); + let error = writer.shutdown().await.expect_err("peer close should fail the shard writer"); + let source = error + .get_ref() + .and_then(|source| source.downcast_ref::()) + .expect("shard failure should preserve the internode error classification"); + assert_eq!(source.kind(), InternodeHttpErrorKind::ConnectionReset); + assert!(source.kind().is_retryable()); + assert_eq!(source.context().operation(), Some(INTERNODE_OPERATION_PUT_FILE_STREAM)); + peer.await.expect("peer should finish normally"); + } + + #[tokio::test] + async fn malformed_upload_response_is_not_an_eof_retry() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("peer listener should bind"); + let address = listener.local_addr().expect("peer address should be available"); + let peer = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("peer should receive the upload"); + let mut buffer = [0_u8; 4096]; + assert!(socket.read(&mut buffer).await.expect("peer should read the request") > 0); + socket + .write_all(b"invalid HTTP response\r\n\r\n") + .await + .expect("peer should send the invalid response"); + socket.shutdown().await.expect("peer should close"); + }); + let error = reqwest::Client::builder() + .no_proxy() + .build() + .expect("upload client should build") + .put(format!("http://{address}/rustfs/rpc/put_file_stream_v1")) + .body("payload") + .send() + .await + .expect_err("invalid response headers should fail the request"); + let kind = classify_reqwest_error(&error); + assert_eq!(kind, InternodeHttpErrorKind::Unknown); + assert!(!kind.is_retryable(), "a malformed response must not become a transport retry"); + peer.await.expect("peer should finish normally"); + } + #[test] fn classify_typed_wins_over_body() { // io ConnectionReset present AND is_body: typed classification wins.