From afa367d220fdf2b40de1d78a38b10e52e715eba5 Mon Sep 17 00:00:00 2001 From: AL Date: Sat, 3 Oct 2026 05:18:41 -0400 Subject: [PATCH] fix(http): classify peer closures from error sources (#8307) fix(http): classify peer transport closure by error source Co-authored-by: Codex --- rustfs/src/server/http.rs | 116 ++++++++++++++++++++++++++++++++++---- 1 file changed, 105 insertions(+), 11 deletions(-) diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 47003af4d..111c2c98c 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -2334,18 +2334,8 @@ fn process_connection( /// Handles connection errors by logging them with appropriate severity fn handle_connection_error(peer_addr: Option<&str>, err: &(dyn std::error::Error + 'static)) { let peer_addr = peer_addr.unwrap_or("unknown"); - let s = err.to_string(); - if s.contains("connection reset") || s.contains("broken pipe") { - log_transport_closed(peer_addr, "connection_reset", &s, "client_disconnect"); - return; - } - if let Some(hyper_err) = err.downcast_ref::() { - if hyper_err.is_incomplete_message() { - log_transport_closed(peer_addr, "incomplete_message", &hyper_err.to_string(), "client_disconnect"); - } else if hyper_err.is_closed() { - log_transport_closed(peer_addr, "connection_closed", &hyper_err.to_string(), "client_disconnect"); - } else if hyper_err.is_parse() { + if hyper_err.is_parse() { log_transport_failed(peer_addr, "parse_failure", &hyper_err.to_string()); } else if hyper_err.is_user() { // is_user() = "error from user's Body stream": the application @@ -2366,6 +2356,14 @@ fn handle_connection_error(peer_addr: Option<&str>, err: &(dyn std::error::Error result = "transport_error", "HTTP transport failed" ); + } else if let Some(kind) = + peer_transport_close_kind(err).or_else(|| hyper_body_write_close_kind(&format!("{hyper_err:?}"))) + { + log_transport_closed(peer_addr, kind, &hyper_err.to_string(), "client_disconnect"); + } else if hyper_err.is_incomplete_message() { + log_transport_closed(peer_addr, "incomplete_message", &hyper_err.to_string(), "client_disconnect"); + } else if hyper_err.is_closed() { + log_transport_closed(peer_addr, "connection_closed", &hyper_err.to_string(), "client_disconnect"); } else if hyper_err.is_canceled() { log_transport_closed(peer_addr, "canceled", &hyper_err.to_string(), "client_disconnect"); } else if format!("{:?}", hyper_err).contains("HeaderTimeout") { @@ -2383,6 +2381,10 @@ fn handle_connection_error(peer_addr: Option<&str>, err: &(dyn std::error::Error ); } } else if let Some(io_err) = err.downcast_ref::() { + if let Some(kind) = peer_transport_close_kind(err) { + log_transport_closed(peer_addr, kind, &io_err.to_string(), "client_disconnect"); + return; + } error!( event = EVENT_HTTP_TRANSPORT_FAILED, component = LOG_COMPONENT_SERVER, @@ -2407,6 +2409,37 @@ fn handle_connection_error(peer_addr: Option<&str>, err: &(dyn std::error::Error } } +fn peer_transport_close_kind(err: &(dyn std::error::Error + 'static)) -> Option<&'static str> { + let mut source = Some(err); + while let Some(error) = source { + if let Some(io_error) = error.downcast_ref::() { + let kind = match io_error.kind() { + std::io::ErrorKind::BrokenPipe => Some("broken_pipe"), + std::io::ErrorKind::ConnectionReset => Some("connection_reset"), + std::io::ErrorKind::ConnectionAborted => Some("connection_aborted"), + _ => None, + }; + if kind.is_some() { + return kind; + } + } + source = error.source(); + } + None +} + +fn hyper_body_write_close_kind(debug: &str) -> Option<&'static str> { + let debug = debug.to_ascii_lowercase(); + let fields = debug.strip_prefix("hyper::error(bodywrite, os {")?.split("message:").next()?; + let kind = fields.split(',').find_map(|field| field.trim().strip_prefix("kind: "))?; + match kind.trim_end_matches([' ', '}']) { + "brokenpipe" => Some("broken_pipe"), + "connectionreset" => Some("connection_reset"), + "connectionaborted" => Some("connection_aborted"), + _ => None, + } +} + #[allow(clippy::result_large_err)] fn check_auth(req: Request<()>) -> std::result::Result, Status> { let local_node_name = @@ -2612,6 +2645,67 @@ mod tests { use tokio::sync::{Notify, mpsc}; use tower::{Layer, Service, ServiceBuilder}; + #[test] + fn peer_transport_close_kind_uses_io_error_sources() { + #[derive(Debug)] + struct SourceError(Error); + + impl std::fmt::Display for SourceError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fmt(f) + } + } + + impl std::error::Error for SourceError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + Some(&self.0) + } + } + + for (io_kind, expected) in [ + (std::io::ErrorKind::BrokenPipe, "broken_pipe"), + (std::io::ErrorKind::ConnectionReset, "connection_reset"), + (std::io::ErrorKind::ConnectionAborted, "connection_aborted"), + ] { + let nested = SourceError(Error::new(io_kind, "peer closed")); + assert_eq!(peer_transport_close_kind(&nested), Some(expected)); + } + assert_eq!(peer_transport_close_kind(&Error::other("Broken pipe")), None); + assert_eq!( + peer_transport_close_kind(&Error::new(std::io::ErrorKind::UnexpectedEof, "stream failed")), + None + ); + } + + #[test] + fn capitalized_hyper_body_write_peer_close_uses_narrow_fallback() { + let observed = "hyper::Error(BodyWrite, Os { code: 32, kind: BrokenPipe, message: \"Broken pipe\" })"; + assert_eq!(hyper_body_write_close_kind(observed), Some("broken_pipe")); + assert_eq!(hyper_body_write_close_kind(&observed.to_ascii_lowercase()), Some("broken_pipe")); + assert_eq!( + hyper_body_write_close_kind( + "hyper::Error(BodyWrite, Os { code: 104, kind: ConnectionReset, message: \"Connection reset by peer\" })" + ), + Some("connection_reset") + ); + assert_eq!( + hyper_body_write_close_kind( + "hyper::Error(BodyWrite, Os { code: 103, kind: ConnectionAborted, message: \"Software caused connection abort\" })" + ), + Some("connection_aborted") + ); + assert_eq!(hyper_body_write_close_kind("hyper::Error(User, Broken pipe)"), None); + assert_eq!(hyper_body_write_close_kind("hyper::Error(Parse, Os { kind: BrokenPipe })"), None); + assert_eq!( + hyper_body_write_close_kind("hyper::Error(BodyWrite, Os { kind: Other, message: \"Broken pipe\" })"), + None + ); + assert_eq!( + hyper_body_write_close_kind("hyper::Error(BodyWrite, Os { kind: Other, message: \"kind: BrokenPipe\" })"), + None + ); + } + type MetricRow = ( metrics_util::CompositeKey, Option,