mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
fix(http): classify peer closures from error sources (#8307)
fix(http): classify peer transport closure by error source Co-authored-by: Codex <codex@bednarowski.ca>
This commit is contained in:
+105
-11
@@ -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::<hyper::Error>() {
|
||||
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::<Error>() {
|
||||
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::<Error>() {
|
||||
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<Request<()>, 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<metrics::Unit>,
|
||||
|
||||
Reference in New Issue
Block a user