fix(replication): carry failure rolling windows through cluster aggregation

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.
This commit is contained in:
唐小鸭
2026-08-15 20:09:54 +08:00
parent 111d10027a
commit 1c660362d9
4 changed files with 97 additions and 11 deletions
@@ -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)
+28
View File
@@ -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<FailureSample>,
}
@@ -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(),
}
}
+1 -1
View File
@@ -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;
+62 -10
View File
@@ -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"