From 51497cb5335b5f02bdf16355bb413bad182da757 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Tue, 18 Aug 2026 12:57:44 +0800 Subject: [PATCH] fix(ecstore): give peer REST failures op and bucket context (#6200) --- .../src/cluster/rpc/peer_rest_client.rs | 196 +++++++++++++++--- .../ecstore/src/cluster/rpc/peer_s3_client.rs | 2 + 2 files changed, 166 insertions(+), 32 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 1426fc91f..3bac43266 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -86,6 +86,25 @@ const PEER_REST_RECOVERY_MAX_BACKOFF: Duration = Duration::from_secs(30); const SCANNER_ACTIVITY_MAX_MESSAGE_SIZE: usize = 1024; const REPLICATION_STATS_MAX_MESSAGE_SIZE: usize = 8 * 1024 * 1024; +/// Error for a peer that reported `success = false` without an `error_info` payload. +/// +/// Same shape as `peer_s3_client::peer_failure_without_details`, over `StorageError` +/// instead of `DiskError`. The message names the operation (and the bucket, where the +/// operation has one) and nothing else, for two reasons: +/// +/// - `finalize_result` classifies failures by message substring, so any text matching +/// `message_has_network_needle` would take an answering peer offline and evict its +/// connection over a plain application-level rejection. +/// - Quorum aggregation (`reduce_errs`) buckets `Io` errors by kind plus rendered +/// message, so a per-peer detail such as the peer address would split one shared +/// failure into single-count buckets and downgrade the dominant error. +fn peer_failure_without_details(op: &str, bucket: Option<&str>) -> Error { + match bucket { + Some(bucket) => Error::other(format!("{op}({bucket}): peer returned failure without error details")), + None => Error::other(format!("{op}: peer returned failure without error details")), + } +} + fn decode_bucket_stats_response(response: GetBucketStatsDataResponse) -> Result { if !response.success { return Err(Error::other( @@ -696,7 +715,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("local_storage_info", None)); } let data = response.storage_info; @@ -719,7 +738,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("server_info", None)); } let data = response.server_properties; @@ -742,7 +761,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_cpus", None)); } let data = response.cpus; @@ -765,7 +784,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_net_info", None)); } let data = response.net_info; @@ -788,7 +807,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_partitions", None)); } let data = response.partitions; @@ -811,7 +830,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_os_info", None)); } let data = response.os_info; @@ -832,7 +851,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_se_linux_info", None)); } let data = response.sys_services; @@ -857,7 +876,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_sys_config", None)); } let data = response.sys_config; @@ -882,7 +901,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_sys_errors", None)); } let data = response.sys_errors; @@ -907,7 +926,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_mem_info", None)); } let data = response.mem_info; @@ -939,7 +958,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_metrics", None)); } let data = response.realtime_metrics; @@ -964,7 +983,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_live_events", None)); } Ok(PeerLiveEventsBatch { @@ -989,7 +1008,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("get_proc_info", None)); } let data = response.proc_info; @@ -1016,7 +1035,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("start_profiling", None)); } Ok(()) } @@ -1323,7 +1342,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_bucket_metadata", Some(bucket))); } Ok(()) } @@ -1346,7 +1365,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("delete_bucket_metadata", Some(bucket))); } Ok(()) } @@ -1369,7 +1388,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("delete_policy", None)); } Ok(()) } @@ -1392,7 +1411,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_policy", None)); } Ok(()) } @@ -1417,7 +1436,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_policy_mapping", None)); } Ok(()) } @@ -1440,7 +1459,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("delete_user", None)); } Ok(()) } @@ -1463,7 +1482,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("delete_service_account", None)); } Ok(()) } @@ -1487,7 +1506,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_user", None)); } Ok(()) } @@ -1510,7 +1529,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_service_account", None)); } Ok(()) } @@ -1533,7 +1552,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_group", None)); } Ok(()) } @@ -1554,7 +1573,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("reload_site_replication_config", None)); } Ok(()) } @@ -1597,7 +1616,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("signal_service", None)); } validate_signal_service_protocol(sig, sub_sys, response.protocol_version)?; Ok(response) @@ -1667,7 +1686,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("reload_pool_meta", None)); } Ok(()) @@ -1691,7 +1710,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("stop_rebalance", None)); } Ok(()) @@ -1725,7 +1744,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("load_rebalance_meta", None)); } Ok(()) @@ -1753,7 +1772,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("start_decommission", None)); } Ok(()) @@ -1777,7 +1796,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("decommission_cancel", None)); } Ok(()) @@ -1801,7 +1820,7 @@ impl PeerRestClient { if let Some(msg) = response.error_info { return Err(Error::other(msg)); } - return Err(Error::other("")); + return Err(peer_failure_without_details("clear_decommission", None)); } Ok(()) @@ -1947,6 +1966,8 @@ fn tier_config_reload_status_outcome(status: tonic::Status) -> TierConfigReloadO mod tests { use super::*; use crate::config::com::STORAGE_CLASS_SUB_SYS; + use crate::disk::error::DiskError; + use crate::disk::error_reduce::reduce_errs; use crate::layout::{disks_layout::DisksLayout, endpoints::SetupType}; use rustfs_config::{ENV_KUBERNETES_SERVICE_HOST, ENV_LOCAL_ENDPOINT_HOST, ENV_STARTUP_TOPOLOGY_WAIT_MODE}; use serde_json::Value; @@ -3098,4 +3119,115 @@ mod tests { && span.get("request_id").and_then(Value::as_str) == Some("req-peer-rest") })); } + + /// Every operation name passed to `peer_failure_without_details` in this file. + const PEER_FAILURE_OPS: &[&str] = &[ + "local_storage_info", + "server_info", + "get_cpus", + "get_net_info", + "get_partitions", + "get_os_info", + "get_se_linux_info", + "get_sys_config", + "get_sys_errors", + "get_mem_info", + "get_metrics", + "get_live_events", + "get_proc_info", + "start_profiling", + "load_bucket_metadata", + "delete_bucket_metadata", + "delete_policy", + "load_policy", + "load_policy_mapping", + "delete_user", + "delete_service_account", + "load_user", + "load_service_account", + "load_group", + "reload_site_replication_config", + "signal_service", + "reload_pool_meta", + "stop_rebalance", + "load_rebalance_meta", + "start_decommission", + "decommission_cancel", + "clear_decommission", + ]; + + #[test] + fn peer_failure_without_details_names_operation_and_bucket() { + for op in PEER_FAILURE_OPS { + let message = peer_failure_without_details(op, None).to_string(); + assert!(message.contains(op), "{op} message must name the operation: {message}"); + } + + for op in ["load_bucket_metadata", "delete_bucket_metadata"] { + let message = peer_failure_without_details(op, Some("ops-bucket")).to_string(); + assert!(message.contains(op), "{op} message must name the operation: {message}"); + assert!(message.contains("ops-bucket"), "{op} message must name the bucket: {message}"); + } + } + + #[test] + fn peer_failure_without_details_keeps_one_reduce_errs_bucket_per_operation() { + // reduce_errs groups Io errors by kind plus rendered message: peers failing the + // same operation must stay a single dominant error instead of one bucket per peer. + let per_peer_errs = (0..4) + .map(|_| Some(DiskError::from(peer_failure_without_details("load_bucket_metadata", Some("shared"))))) + .collect::>(); + let (count, dominant) = reduce_errs(&per_peer_errs, &[]); + assert_eq!(count, 4, "one shared failure must not split into per-peer buckets"); + assert_eq!( + dominant, + Some(DiskError::from(peer_failure_without_details("load_bucket_metadata", Some("shared")))) + ); + + assert_ne!( + peer_failure_without_details("load_bucket_metadata", Some("shared")).to_string(), + peer_failure_without_details("delete_bucket_metadata", Some("shared")).to_string() + ); + assert_ne!( + peer_failure_without_details("load_bucket_metadata", Some("bucket-a")).to_string(), + peer_failure_without_details("load_bucket_metadata", Some("bucket-b")).to_string() + ); + } + + #[test] + fn peer_failure_without_details_never_reads_as_a_network_failure() { + // `finalize_result` marks the peer offline and evicts its connection whenever the + // message matches a network needle. A peer that answered `success = false` is alive, + // so no operation or bucket name may push this text over that classifier. + for op in PEER_FAILURE_OPS { + let err = peer_failure_without_details(op, None); + assert!( + !PeerRestClient::is_network_like_error(&err), + "{op} must not read as a transport failure: {err}" + ); + + let scoped = peer_failure_without_details(op, Some("bucket-name")); + assert!( + !PeerRestClient::is_network_like_error(&scoped), + "{op} must not read as a transport failure: {scoped}" + ); + } + + // The bucket name is caller-supplied. Every needle carries a space, which S3 bucket + // names cannot, and the name is closed by `)` before the literal text resumes, so no + // needle can straddle the boundary either. + for bucket in [ + "timed-out", + "connection-reset", + "transport-error", + "broken-pipe", + "unavailable-logs", + ] { + let err = peer_failure_without_details("load_bucket_metadata", Some(bucket)); + assert!( + !PeerRestClient::is_network_like_error(&err), + "bucket {bucket} must not push the message over the network classifier: {err}" + ); + } + } } diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index 02c9f75f1..d44651dce 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -220,6 +220,8 @@ fn pool_write_quorum(participant_count: usize) -> usize { /// buckets `Error::Io` by kind plus rendered message, so any per-peer detail (address, /// timing) would split one shared failure into single-count buckets and downgrade a real /// dominant error into `ErasureWriteQuorum`. +/// +/// `peer_rest_client` carries the same helper over `StorageError` for the same response shape. fn peer_failure_without_details(op: &str, bucket: Option<&str>) -> Error { match bucket { Some(bucket) => Error::other(format!("{op}({bucket}): peer returned failure without error details")),