fix(rio): preserve bounded retries after peer EOF (#8272)

This commit is contained in:
Chris
2026-09-30 20:18:28 +08:00
committed by GitHub
parent 88d03d8199
commit 8655c38f1d
2 changed files with 103 additions and 6 deletions
+1
View File
@@ -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
+102 -6
View File
@@ -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<InternodeHttpErrorKind> {
/// 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<InternodeHttpErrorKind> {
let mut source: Option<&(dyn std::error::Error + 'static)> = Some(err);
while let Some(current) = source {
if let Some(io_err) = current.downcast_ref::<io::Error>() {
@@ -812,6 +810,12 @@ fn find_io_error_kind_in_chain(err: &(dyn std::error::Error + 'static)) -> Optio
_ => {}
}
}
if current
.downcast_ref::<hyper::Error>()
.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::<InternodeHttpError>())
.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.