fix(ecstore): give peer REST failures op and bucket context (#6200)

This commit is contained in:
Zhengchao An
2026-08-18 12:57:44 +08:00
committed by GitHub
parent c7a29ec0a7
commit 51497cb533
2 changed files with 166 additions and 32 deletions
@@ -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<BucketStats> {
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::<Vec<_>>();
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}"
);
}
}
}
@@ -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")),