diff --git a/crates/obs/src/metrics/scheduler.rs b/crates/obs/src/metrics/scheduler.rs index 3217db20a..e02ab8fda 100644 --- a/crates/obs/src/metrics/scheduler.rs +++ b/crates/obs/src/metrics/scheduler.rs @@ -25,6 +25,7 @@ use crate::metrics::collectors::{ AuditTargetStats, + BucketReplicationBandwidthStats, NotificationStats, NotificationTargetStats, // System monitoring collectors (migrated from rustfs-obs::system) @@ -104,6 +105,78 @@ const DEFAULT_REPL_BW_ZERO_TOMBSTONE_CYCLES: u8 = 3; /// Env var that overrides the zero-emission tombstone cycles for removed replication bandwidth series. const ENV_REPL_BW_ZERO_TOMBSTONE_CYCLES: &str = "RUSTFS_METRICS_REPL_BW_ZERO_TOMBSTONE_CYCLES"; +type ReplBwKey = (String, String); // (bucket, target_arn) + +fn repl_bw_live_keys(stats: &[BucketReplicationBandwidthStats]) -> HashSet { + stats.iter().map(|s| (s.bucket.clone(), s.target_arn.clone())).collect() +} + +fn update_repl_bw_zero_tombstones( + monitor_available: bool, + has_seen_valid_snapshot: &mut bool, + prev_live_keys: &mut HashSet, + zero_tombstones: &mut HashMap, + current_live_keys: HashSet, + tombstone_cycles: u8, +) { + if !monitor_available { + return; + } + + if *has_seen_valid_snapshot { + for removed in prev_live_keys.difference(¤t_live_keys) { + zero_tombstones.insert(removed.clone(), tombstone_cycles); + } + } + + // Key becomes live again: stop zeroing immediately. + for key in ¤t_live_keys { + zero_tombstones.remove(key); + } + + *prev_live_keys = current_live_keys; + *has_seen_valid_snapshot = true; +} + +fn collect_repl_bw_zero_tombstone_metrics(zero_tombstones: &HashMap) -> Vec { + if zero_tombstones.is_empty() { + return Vec::new(); + } + + let mut zero_metrics = Vec::with_capacity(zero_tombstones.len() * 2); + for (bucket, target_arn) in zero_tombstones.keys() { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.clone()); + let target_arn_label: Cow<'static, str> = Cow::Owned(target_arn.clone()); + + zero_metrics.push( + PrometheusMetric::from_descriptor(&BUCKET_REPL_BANDWIDTH_LIMIT_MD, 0.0) + .with_label(BUCKET_L, bucket_label.clone()) + .with_label(TARGET_ARN_L, target_arn_label.clone()), + ); + + zero_metrics.push( + PrometheusMetric::from_descriptor(&BUCKET_REPL_BANDWIDTH_CURRENT_MD, 0.0) + .with_label(BUCKET_L, bucket_label) + .with_label(TARGET_ARN_L, target_arn_label), + ); + } + + zero_metrics +} + +fn expire_repl_bw_zero_tombstones(monitor_available: bool, zero_tombstones: &mut HashMap) { + if monitor_available && !zero_tombstones.is_empty() { + zero_tombstones.retain(|_, remaining| { + if *remaining <= 1 { + false + } else { + *remaining -= 1; + true + } + }); + } +} + /// Initialize all metrics collectors. /// /// This function spawns background tasks that periodically collect metrics @@ -288,7 +361,6 @@ pub fn init_metrics_runtime(token: CancellationToken) { tokio::spawn(async move { let mut interval = tokio::time::interval(bucket_replication_bandwidth_interval); let repl_bw_zero_tombstone_cycles = parse_repl_bw_zero_tombstone_cycles(); - type ReplBwKey = (String, String); // (bucket, target_arn) let mut prev_live_keys: HashSet = HashSet::new(); let mut zero_tombstones: HashMap = HashMap::new(); let mut has_seen_valid_snapshot = false; @@ -298,50 +370,23 @@ pub fn init_metrics_runtime(token: CancellationToken) { let monitor_available = get_global_bucket_monitor().is_some(); let stats = collect_bucket_replication_bandwidth_stats(); - let current_live_keys: HashSet = stats.iter().map(|s| (s.bucket.clone(), s.target_arn.clone())) - .collect(); + let current_live_keys = repl_bw_live_keys(&stats); - if monitor_available { - if has_seen_valid_snapshot { - for removed in prev_live_keys.difference(¤t_live_keys) { - zero_tombstones.insert(removed.clone(), repl_bw_zero_tombstone_cycles); - } - } - - // Key becomes live again: stop zeroing immediately. - for key in ¤t_live_keys { - zero_tombstones.remove(key); - } - - prev_live_keys = current_live_keys; - has_seen_valid_snapshot = true; - } else { + if !monitor_available { warn!("Bucket monitor unavailable; skip replication bandwidth key-state transition this cycle."); } + update_repl_bw_zero_tombstones( + monitor_available, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + current_live_keys, + repl_bw_zero_tombstone_cycles, + ); let mut metrics = collect_bucket_replication_bandwidth_metrics(&stats); // Phase-1 action: force zero for removed keys during tombstone cycles. - if !zero_tombstones.is_empty() { - let mut zero_metrics = Vec::with_capacity(zero_tombstones.len() * 2); - for (bucket, target_arn) in zero_tombstones.keys() { - let bucket_label: Cow<'static, str> = Cow::Owned(bucket.clone()); - let target_arn_label: Cow<'static, str> = Cow::Owned(target_arn.clone()); - - zero_metrics.push( - PrometheusMetric::from_descriptor(&BUCKET_REPL_BANDWIDTH_LIMIT_MD, 0.0) - .with_label(BUCKET_L, bucket_label.clone()) - .with_label(TARGET_ARN_L, target_arn_label.clone()), - ); - - zero_metrics.push( - PrometheusMetric::from_descriptor(&BUCKET_REPL_BANDWIDTH_CURRENT_MD, 0.0) - .with_label(BUCKET_L, bucket_label) - .with_label(TARGET_ARN_L, target_arn_label), - ); - } - - metrics.extend(zero_metrics); - } + metrics.extend(collect_repl_bw_zero_tombstone_metrics(&zero_tombstones)); let bucket_replication = collect_bucket_replication_detail_stats().await; metrics.extend(collect_bucket_replication_metrics(&bucket_replication)); @@ -350,16 +395,7 @@ pub fn init_metrics_runtime(token: CancellationToken) { report_metrics(&metrics); // Phase-2: after N cycles, stop reporting -> series becomes absent after expiration. - if monitor_available && !zero_tombstones.is_empty() { - zero_tombstones.retain(|_, remaining| { - if *remaining <= 1 { - false - } else { - *remaining -= 1; - true - } - }); - } + expire_repl_bw_zero_tombstones(monitor_available, &mut zero_tombstones); } _ = token_clone.cancelled() => { warn!("Metrics collection for bucket replication bandwidth stats cancelled."); @@ -646,10 +682,21 @@ fn collect_system_monitoring_metrics( #[cfg(test)] mod tests { - use super::advance_deadline; + use super::*; + use std::collections::{HashMap, HashSet}; use std::time::Duration; use tokio::time::Instant; + fn repl_bw_key(bucket: &str, target_arn: &str) -> ReplBwKey { + (bucket.to_string(), target_arn.to_string()) + } + + fn repl_bw_keys(keys: &[(&str, &str)]) -> HashSet { + keys.iter() + .map(|(bucket, target_arn)| repl_bw_key(bucket, target_arn)) + .collect() + } + #[test] fn advance_deadline_keeps_future_deadline_unchanged() { let base = Instant::now(); @@ -665,4 +712,119 @@ mod tests { advance_deadline(&mut deadline, Duration::from_secs(5), base + Duration::from_secs(12)); assert_eq!(deadline, base + Duration::from_secs(15)); } + + #[test] + fn repl_bw_tombstones_zero_removed_keys_then_expire() { + let mut has_seen_valid_snapshot = false; + let mut prev_live_keys = HashSet::new(); + let mut zero_tombstones = HashMap::new(); + let key = repl_bw_key("photos", "arn:rustfs:replication:target-a"); + + update_repl_bw_zero_tombstones( + true, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + repl_bw_keys(&[("photos", "arn:rustfs:replication:target-a")]), + 2, + ); + assert!(has_seen_valid_snapshot); + assert_eq!(prev_live_keys, repl_bw_keys(&[("photos", "arn:rustfs:replication:target-a")])); + assert!(zero_tombstones.is_empty()); + + update_repl_bw_zero_tombstones( + true, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + HashSet::new(), + 2, + ); + assert_eq!(zero_tombstones.get(&key), Some(&2)); + + let metrics = collect_repl_bw_zero_tombstone_metrics(&zero_tombstones); + assert_eq!(metrics.len(), 2); + assert!(metrics.iter().all(|metric| metric.value == 0.0)); + + let names = metrics.iter().map(|metric| metric.name.to_string()).collect::>(); + assert!(names.contains(&BUCKET_REPL_BANDWIDTH_LIMIT_MD.get_full_metric_name())); + assert!(names.contains(&BUCKET_REPL_BANDWIDTH_CURRENT_MD.get_full_metric_name())); + + for metric in metrics { + let labels = metric + .labels + .into_iter() + .map(|(key, value)| (key, value.to_string())) + .collect::>(); + assert_eq!(labels.get(BUCKET_L).map(String::as_str), Some("photos")); + assert_eq!(labels.get(TARGET_ARN_L).map(String::as_str), Some("arn:rustfs:replication:target-a")); + } + + expire_repl_bw_zero_tombstones(true, &mut zero_tombstones); + assert_eq!(zero_tombstones.get(&key), Some(&1)); + + expire_repl_bw_zero_tombstones(true, &mut zero_tombstones); + assert!(zero_tombstones.is_empty()); + } + + #[test] + fn repl_bw_tombstones_stop_zeroing_when_key_becomes_live_again() { + let mut has_seen_valid_snapshot = false; + let mut prev_live_keys = HashSet::new(); + let mut zero_tombstones = HashMap::new(); + let live_keys = repl_bw_keys(&[("photos", "arn:rustfs:replication:target-a")]); + + update_repl_bw_zero_tombstones( + true, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + live_keys.clone(), + 3, + ); + update_repl_bw_zero_tombstones( + true, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + HashSet::new(), + 3, + ); + assert_eq!(zero_tombstones.get(&repl_bw_key("photos", "arn:rustfs:replication:target-a")), Some(&3)); + + update_repl_bw_zero_tombstones( + true, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + live_keys.clone(), + 3, + ); + + assert!(zero_tombstones.is_empty()); + assert_eq!(prev_live_keys, live_keys); + } + + #[test] + fn repl_bw_tombstones_do_not_advance_when_monitor_unavailable() { + let mut has_seen_valid_snapshot = true; + let mut prev_live_keys = repl_bw_keys(&[("photos", "arn:rustfs:replication:target-a")]); + let mut zero_tombstones = HashMap::from([(repl_bw_key("videos", "arn:rustfs:replication:target-b"), 1)]); + + update_repl_bw_zero_tombstones( + false, + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + HashSet::new(), + 3, + ); + + assert!(has_seen_valid_snapshot); + assert_eq!(prev_live_keys, repl_bw_keys(&[("photos", "arn:rustfs:replication:target-a")])); + assert_eq!(zero_tombstones.get(&repl_bw_key("videos", "arn:rustfs:replication:target-b")), Some(&1)); + + expire_repl_bw_zero_tombstones(false, &mut zero_tombstones); + assert_eq!(zero_tombstones.get(&repl_bw_key("videos", "arn:rustfs:replication:target-b")), Some(&1)); + } }