From 1c660362d9901d0e930a5244ee3430dc0d171653 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Sat, 15 Aug 2026 20:09:54 +0800 Subject: [PATCH] fix(replication): carry failure rolling windows through cluster aggregation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review: both metrics endpoints aggregate first, and FailStats::merge dropped the process-local samples (which also never cross the peer-RPC wire — serde-skipped), so lastMinute/lastHour serialized as zero right after a failure while totals was nonzero. - FailStats gains serializable last_minute/last_hour window snapshots (serde default: old nodes read zeros, new fields are ignored by old decoders), recomputed on every add_size and re-stamped at the per-node collection point (get_latest_replication_stats), and summed by merge. - The wire DTO takes the component-wise max of the live samples and the snapshot, so both the single-node and the aggregated path report the window. - Regression test drives a stat through rmp round trip + merge before serialization, as requested. Also restore the #[allow(dead_code)] attribute to route_policy — the new module declaration had been inserted between the attribute and its item, which broke the -D warnings CI lanes. --- .../bucket/replication/replication_state.rs | 6 ++ crates/replication/src/stats.rs | 28 ++++++++ rustfs/src/admin/mod.rs | 2 +- rustfs/src/admin/replication_metrics_wire.rs | 72 ++++++++++++++++--- 4 files changed, 97 insertions(+), 11 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_state.rs b/crates/ecstore/src/bucket/replication/replication_state.rs index d16a507e5..cd0e303f3 100644 --- a/crates/ecstore/src/bucket/replication/replication_state.rs +++ b/crates/ecstore/src/bucket/replication/replication_state.rs @@ -704,6 +704,12 @@ impl ReplicationStats { } else { BucketReplicationStats::new() }; + // Stamp the serializable failure windows from the live samples: the + // samples themselves do not cross the peer-RPC wire, so this snapshot + // is what cluster aggregation and the metrics endpoints see. + for stat in replication_stats.stats.values_mut() { + stat.fail_stats.refresh_windows(); + } let uptime = if cache.contains_key(bucket) { SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) diff --git a/crates/replication/src/stats.rs b/crates/replication/src/stats.rs index a6c000d43..621aca473 100644 --- a/crates/replication/src/stats.rs +++ b/crates/replication/src/stats.rs @@ -520,6 +520,14 @@ struct FailureSample { pub struct FailStats { pub count: i64, pub size: i64, + /// Rolling-window snapshots refreshed at collection time + /// ([`Self::refresh_windows`]). The raw samples (`recent`) are process + /// local (serde-skipped), so these fields are what survives the peer-RPC + /// wire and [`Self::merge`]-based cluster aggregation. + #[serde(default)] + pub last_minute: FailedMetric, + #[serde(default)] + pub last_hour: FailedMetric, #[serde(skip)] recent: VecDeque, } @@ -535,6 +543,16 @@ impl FailStats { self.size = self.size.saturating_add(size); self.recent.push_back(FailureSample { observed_at, size }); self.prune(observed_at); + self.refresh_windows(); + } + + /// Recompute the serializable rolling-window snapshots from the local + /// samples. Only meaningful on the live per-node struct: a deserialized + /// or merged struct has no samples, and refreshing it would wipe the + /// aggregated windows. + pub fn refresh_windows(&mut self) { + self.last_minute = self.recent_since(Duration::from_secs(60)); + self.last_hour = self.recent_since(Duration::from_secs(3600)); } fn prune(&mut self, observed_at: Instant) { @@ -565,6 +583,16 @@ impl FailStats { Self { count: self.count.saturating_add(other.count), size: self.size.saturating_add(other.size), + // The window snapshots sum across nodes; the raw samples do not + // travel and stay empty on aggregated structs. + last_minute: FailedMetric { + count: self.last_minute.count.saturating_add(other.last_minute.count), + size: self.last_minute.size.saturating_add(other.last_minute.size), + }, + last_hour: FailedMetric { + count: self.last_hour.count.saturating_add(other.last_hour.count), + size: self.last_hour.size.saturating_add(other.last_hour.size), + }, recent: VecDeque::new(), } } diff --git a/rustfs/src/admin/mod.rs b/rustfs/src/admin/mod.rs index 4afb9a5f4..d61abf932 100644 --- a/rustfs/src/admin/mod.rs +++ b/rustfs/src/admin/mod.rs @@ -17,9 +17,9 @@ mod auth; pub mod console; pub mod handlers; mod plugin_contract; +pub(crate) mod replication_metrics_wire; // Contract inventory is validated by tests before later runtime integration. #[allow(dead_code)] -pub(crate) mod replication_metrics_wire; pub(crate) mod route_policy; pub mod router; pub(crate) mod runtime_sources; diff --git a/rustfs/src/admin/replication_metrics_wire.rs b/rustfs/src/admin/replication_metrics_wire.rs index 713f8eac6..4d64e915c 100644 --- a/rustfs/src/admin/replication_metrics_wire.rs +++ b/rustfs/src/admin/replication_metrics_wire.rs @@ -207,17 +207,29 @@ pub(crate) struct TargetMetricsWire { } fn target_timed_err_stats(stat: &InternalReplicationStat) -> TimedErrStatsWire { - let last_minute = stat.fail_stats.recent_since(Duration::from_secs(60)); - let last_hour = stat.fail_stats.recent_since(Duration::from_secs(3600)); + // Cluster aggregation merges FailStats without the process-local samples, + // so the serializable window snapshots (refreshed at each node's + // collection point, summed by merge) are authoritative here; the live + // samples only ever agree with or lag them, so take the larger. + let sampled_minute = stat.fail_stats.recent_since(Duration::from_secs(60)); + let sampled_hour = stat.fail_stats.recent_since(Duration::from_secs(3600)); + let window = |sampled_count: i64, sampled_size: i64, snapshot_count: i64, snapshot_size: i64| RStatWire { + count: sampled_count.max(snapshot_count) as f64, + bytes: sampled_size.max(snapshot_size), + }; TimedErrStatsWire { - last_minute: RStatWire { - count: last_minute.count as f64, - bytes: last_minute.size, - }, - last_hour: RStatWire { - count: last_hour.count as f64, - bytes: last_hour.size, - }, + last_minute: window( + sampled_minute.count, + sampled_minute.size, + stat.fail_stats.last_minute.count, + stat.fail_stats.last_minute.size, + ), + last_hour: window( + sampled_hour.count, + sampled_hour.size, + stat.fail_stats.last_hour.count, + stat.fail_stats.last_hour.size, + ), totals: RStatWire { count: stat.failed.count as f64, bytes: stat.failed.size, @@ -486,6 +498,46 @@ mod tests { assert_eq!(json["downtimeInfo"], serde_json::json!({})); } + /// Review regression: both metrics endpoints aggregate first, and the + /// FailStats merge drops the process-local samples — the rolling windows + /// must survive a peer-RPC round trip plus aggregation and still reach + /// the wire body. + #[test] + fn failure_windows_survive_aggregation_before_serialization() { + // Node A: live failure, windows stamped at the collection point. + let mut node_a = crate::admin::storage_api::replication::BucketReplicationStat::default(); + node_a.fail_stats.add_size(512, None::<&std::io::Error>); + node_a.failed = node_a.fail_stats.to_metric(); + + // Node A's stats cross the peer RPC wire: the samples are dropped, + // the window snapshots travel. + let encoded = rmp_serde::to_vec_named(&node_a).expect("stat should encode"); + let remote: crate::admin::storage_api::replication::BucketReplicationStat = + rmp_serde::from_slice(&encoded).expect("stat should decode"); + + // Aggregation merges the remote stat with an empty local one. + let merged_fail = remote.fail_stats.merge(&Default::default()); + let mut aggregated = crate::admin::storage_api::replication::BucketReplicationStat::default(); + aggregated.failed = merged_fail.to_metric(); + aggregated.fail_stats = merged_fail; + + let mut stats = BucketStats::default(); + stats + .replication_stats + .stats + .insert("arn:minio:replication::t:b".to_string(), aggregated); + + let json = serde_json::to_value(MetricsWire::from(&stats.replication_stats)).expect("wire should serialize"); + let failed = &json["Stats"]["arn:minio:replication::t:b"]["failed"]; + assert_eq!(failed["totals"]["count"], 1.0); + assert_eq!( + failed["lastMinute"]["count"], 1.0, + "the rolling minute window must survive RPC + aggregation" + ); + assert_eq!(failed["lastMinute"]["bytes"], 512); + assert_eq!(failed["lastHour"]["count"], 1.0); + } + /// Pin the intra-cluster peer-RPC wire format of the internal stats: it /// is msgpack with the Rust field names as map keys /// (`rmp_serde::to_vec_named` in node_service.rs). If someone "fixes"