diff --git a/crates/obs/src/metrics/scheduler.rs b/crates/obs/src/metrics/scheduler.rs index ca487eb73..3cfdcaf6f 100644 --- a/crates/obs/src/metrics/scheduler.rs +++ b/crates/obs/src/metrics/scheduler.rs @@ -74,9 +74,20 @@ use crate::metrics::config::{ use crate::metrics::obs_is_disk_compression_enabled; use crate::metrics::report::{PrometheusMetric, report_metrics}; use crate::metrics::runtime_sources::bucket_monitor_available; +use crate::metrics::schema::audit::{AUDIT_FAILED_MESSAGES_MD, AUDIT_TARGET_QUEUE_LENGTH_MD, AUDIT_TOTAL_MESSAGES_MD}; use crate::metrics::schema::bucket_replication::{ BUCKET_L, BUCKET_REPL_BANDWIDTH_CURRENT_MD, BUCKET_REPL_BANDWIDTH_LIMIT_MD, TARGET_ARN_L, }; +use crate::metrics::schema::cluster_usage::{ + BUCKET_LABEL as USAGE_BUCKET_LABEL, RANGE_LABEL as USAGE_RANGE_LABEL, USAGE_BUCKET_DELETE_MARKERS_COUNT_MD, + USAGE_BUCKET_OBJECT_SIZE_DISTRIBUTION_MD, USAGE_BUCKET_OBJECT_VERSION_COUNT_DISTRIBUTION_MD, USAGE_BUCKET_OBJECTS_TOTAL_MD, + USAGE_BUCKET_QUOTA_TOTAL_BYTES_MD, USAGE_BUCKET_TOTAL_BYTES_MD, USAGE_BUCKET_VERSIONS_COUNT_MD, +}; +use crate::metrics::schema::node_bucket::{BUCKET_OBJECTS_TOTAL_MD, BUCKET_QUOTA_BYTES_MD, BUCKET_USAGE_BYTES_MD}; +use crate::metrics::schema::notification_target::{ + NOTIFICATION_TARGET_FAILED_MESSAGES_MD, NOTIFICATION_TARGET_QUEUE_LENGTH_MD, NOTIFICATION_TARGET_TOTAL_MESSAGES_MD, + TARGET_ID as NOTIFICATION_TARGET_ID_LABEL, TARGET_TYPE as NOTIFICATION_TARGET_TYPE_LABEL, +}; use crate::metrics::stats_collector::{ ProcessMetricBundle, collect_bucket_replication_bandwidth_stats, collect_bucket_replication_detail_stats, collect_bucket_stats, collect_cluster_and_health_stats, collect_cluster_config_stats, collect_cluster_usage_metric_stats, @@ -115,6 +126,7 @@ const LEGACY_RESOURCE_INTERVAL: &str = "RUSTFS_METRICS_RESOURCE_INTERVAL"; const LEGACY_AUDIT_INTERVAL: &str = "RUSTFS_METRICS_AUDIT_INTERVAL"; const LEGACY_NOTIFICATION_INTERVAL: &str = "RUSTFS_METRICS_NOTIFICATION_INTERVAL"; const LEGACY_DEFAULT_INTERVAL: &str = "RUSTFS_METRICS_DEFAULT_INTERVAL"; +const AUDIT_TARGET_ID_LABEL: &str = "target_id"; /// Default cycles to emit zero for removed replication bandwidth series before letting them expire. const DEFAULT_REPL_BW_ZERO_TOMBSTONE_CYCLES: u8 = 3; @@ -124,6 +136,10 @@ const METRICS_RUNTIME_SERVICE_NAME: &str = "metrics_runtime"; const METRICS_RUNTIME_COLLECTOR_TASKS: u8 = 10; type ReplBwKey = (String, String); // (bucket, target_arn) +type BucketKey = String; +type BucketRangeKey = (String, String); // (bucket, range) +type AuditTargetKey = String; +type NotificationTargetKey = (String, String); // (target_id, target_type) #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] #[serde(rename_all = "snake_case")] @@ -391,6 +407,205 @@ fn repl_bw_live_keys(stats: &[BucketReplicationBandwidthStats]) -> HashSet( + has_seen_valid_snapshot: &mut bool, + prev_live_keys: &mut HashSet, + zero_tombstones: &mut HashMap, + current_live_keys: HashSet, + tombstone_cycles: u8, +) { + if *has_seen_valid_snapshot { + for removed in prev_live_keys.difference(¤t_live_keys) { + zero_tombstones.insert(removed.clone(), tombstone_cycles); + } + } + + for key in ¤t_live_keys { + zero_tombstones.remove(key); + } + + *prev_live_keys = current_live_keys; + *has_seen_valid_snapshot = true; +} + +fn expire_series_zero_tombstones(zero_tombstones: &mut HashMap) { + if !zero_tombstones.is_empty() { + zero_tombstones.retain(|_, remaining| { + if *remaining <= 1 { + false + } else { + *remaining -= 1; + true + } + }); + } +} + +fn bucket_live_keys(stats: &[crate::metrics::collectors::BucketStats]) -> HashSet { + stats.iter().map(|stat| stat.name.clone()).collect() +} + +fn collect_bucket_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() * 3); + for bucket in zero_tombstones.keys() { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.clone()); + zero_metrics + .push(PrometheusMetric::from_descriptor(&BUCKET_USAGE_BYTES_MD, 0.0).with_label("bucket", bucket_label.clone())); + zero_metrics + .push(PrometheusMetric::from_descriptor(&BUCKET_OBJECTS_TOTAL_MD, 0.0).with_label("bucket", bucket_label.clone())); + zero_metrics.push(PrometheusMetric::from_descriptor(&BUCKET_QUOTA_BYTES_MD, 0.0).with_label("bucket", bucket_label)); + } + + zero_metrics +} + +fn bucket_usage_live_keys(stats: &[crate::metrics::collectors::BucketUsageStats]) -> HashSet { + stats.iter().map(|stat| stat.bucket.clone()).collect() +} + +fn bucket_usage_object_size_live_keys(stats: &[crate::metrics::collectors::BucketUsageStats]) -> HashSet { + stats + .iter() + .flat_map(|stat| { + stat.object_size_distribution + .iter() + .map(move |(range, _)| (stat.bucket.clone(), range.clone())) + }) + .collect() +} + +fn bucket_usage_version_live_keys(stats: &[crate::metrics::collectors::BucketUsageStats]) -> HashSet { + stats + .iter() + .flat_map(|stat| { + stat.version_count_distribution + .iter() + .map(move |(range, _)| (stat.bucket.clone(), range.clone())) + }) + .collect() +} + +fn collect_bucket_usage_zero_tombstone_metrics( + zero_bucket_tombstones: &HashMap, + zero_object_size_tombstones: &HashMap, + zero_version_tombstones: &HashMap, +) -> Vec { + let mut zero_metrics = + Vec::with_capacity(zero_bucket_tombstones.len() * 5 + zero_object_size_tombstones.len() + zero_version_tombstones.len()); + + for bucket in zero_bucket_tombstones.keys() { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.clone()); + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_TOTAL_BYTES_MD, 0.0) + .with_label(USAGE_BUCKET_LABEL, bucket_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_OBJECTS_TOTAL_MD, 0.0) + .with_label(USAGE_BUCKET_LABEL, bucket_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_VERSIONS_COUNT_MD, 0.0) + .with_label(USAGE_BUCKET_LABEL, bucket_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_DELETE_MARKERS_COUNT_MD, 0.0) + .with_label(USAGE_BUCKET_LABEL, bucket_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_QUOTA_TOTAL_BYTES_MD, 0.0) + .with_label(USAGE_BUCKET_LABEL, bucket_label), + ); + } + + for (bucket, range) in zero_object_size_tombstones.keys() { + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_OBJECT_SIZE_DISTRIBUTION_MD, 0.0) + .with_label(USAGE_RANGE_LABEL, Cow::Owned(range.clone())) + .with_label(USAGE_BUCKET_LABEL, Cow::Owned(bucket.clone())), + ); + } + + for (bucket, range) in zero_version_tombstones.keys() { + zero_metrics.push( + PrometheusMetric::from_descriptor(&USAGE_BUCKET_OBJECT_VERSION_COUNT_DISTRIBUTION_MD, 0.0) + .with_label(USAGE_RANGE_LABEL, Cow::Owned(range.clone())) + .with_label(USAGE_BUCKET_LABEL, Cow::Owned(bucket.clone())), + ); + } + + zero_metrics +} + +fn audit_target_live_keys(stats: &[AuditTargetStats]) -> HashSet { + stats.iter().map(|stat| stat.target_id.clone()).collect() +} + +fn collect_audit_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() * 3); + for target_id in zero_tombstones.keys() { + let target_id_label: Cow<'static, str> = Cow::Owned(target_id.clone()); + zero_metrics.push( + PrometheusMetric::from_descriptor(&AUDIT_FAILED_MESSAGES_MD, 0.0) + .with_label(AUDIT_TARGET_ID_LABEL, target_id_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&AUDIT_TARGET_QUEUE_LENGTH_MD, 0.0) + .with_label(AUDIT_TARGET_ID_LABEL, target_id_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&AUDIT_TOTAL_MESSAGES_MD, 0.0).with_label(AUDIT_TARGET_ID_LABEL, target_id_label), + ); + } + + zero_metrics +} + +fn notification_target_live_keys(stats: &[NotificationTargetStats]) -> HashSet { + stats + .iter() + .map(|stat| (stat.target_id.clone(), stat.target_type.clone())) + .collect() +} + +fn collect_notification_target_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() * 3); + for (target_id, target_type) in zero_tombstones.keys() { + let target_id_label: Cow<'static, str> = Cow::Owned(target_id.clone()); + let target_type_label: Cow<'static, str> = Cow::Owned(target_type.clone()); + zero_metrics.push( + PrometheusMetric::from_descriptor(&NOTIFICATION_TARGET_FAILED_MESSAGES_MD, 0.0) + .with_label(NOTIFICATION_TARGET_ID_LABEL, target_id_label.clone()) + .with_label(NOTIFICATION_TARGET_TYPE_LABEL, target_type_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&NOTIFICATION_TARGET_QUEUE_LENGTH_MD, 0.0) + .with_label(NOTIFICATION_TARGET_ID_LABEL, target_id_label.clone()) + .with_label(NOTIFICATION_TARGET_TYPE_LABEL, target_type_label.clone()), + ); + zero_metrics.push( + PrometheusMetric::from_descriptor(&NOTIFICATION_TARGET_TOTAL_MESSAGES_MD, 0.0) + .with_label(NOTIFICATION_TARGET_ID_LABEL, target_id_label) + .with_label(NOTIFICATION_TARGET_TYPE_LABEL, target_type_label), + ); + } + + zero_metrics +} + fn update_repl_bw_zero_tombstones( monitor_available: bool, has_seen_valid_snapshot: &mut bool, @@ -516,6 +731,14 @@ pub fn init_metrics_runtime(token: CancellationToken) { let token_clone = token.clone(); tokio::spawn(async move { let mut interval = tokio::time::interval(cluster_interval); + let tombstone_cycles = config.replication_bandwidth_zero_tombstone_cycles; + let mut has_seen_bucket_usage_snapshot = false; + let mut prev_bucket_usage_keys: HashSet = HashSet::new(); + let mut bucket_usage_zero_tombstones: HashMap = HashMap::new(); + let mut prev_bucket_usage_object_size_keys: HashSet = HashSet::new(); + let mut bucket_usage_object_size_zero_tombstones: HashMap = HashMap::new(); + let mut prev_bucket_usage_version_keys: HashSet = HashSet::new(); + let mut bucket_usage_version_zero_tombstones: HashMap = HashMap::new(); loop { tokio::select! { _ = interval.tick() => { @@ -536,7 +759,36 @@ pub fn init_metrics_runtime(token: CancellationToken) { if let Some((cluster_usage, bucket_usage)) = collect_cluster_usage_metric_stats().await { metrics.extend(collect_cluster_usage_metrics(&cluster_usage)); + update_series_zero_tombstones( + &mut has_seen_bucket_usage_snapshot, + &mut prev_bucket_usage_keys, + &mut bucket_usage_zero_tombstones, + bucket_usage_live_keys(&bucket_usage), + tombstone_cycles, + ); + update_series_zero_tombstones( + &mut has_seen_bucket_usage_snapshot, + &mut prev_bucket_usage_object_size_keys, + &mut bucket_usage_object_size_zero_tombstones, + bucket_usage_object_size_live_keys(&bucket_usage), + tombstone_cycles, + ); + update_series_zero_tombstones( + &mut has_seen_bucket_usage_snapshot, + &mut prev_bucket_usage_version_keys, + &mut bucket_usage_version_zero_tombstones, + bucket_usage_version_live_keys(&bucket_usage), + tombstone_cycles, + ); metrics.extend(collect_bucket_usage_metrics(&bucket_usage)); + metrics.extend(collect_bucket_usage_zero_tombstone_metrics( + &bucket_usage_zero_tombstones, + &bucket_usage_object_size_zero_tombstones, + &bucket_usage_version_zero_tombstones, + )); + expire_series_zero_tombstones(&mut bucket_usage_zero_tombstones); + expire_series_zero_tombstones(&mut bucket_usage_object_size_zero_tombstones); + expire_series_zero_tombstones(&mut bucket_usage_version_zero_tombstones); } if !metrics.is_empty() { @@ -555,12 +807,25 @@ pub fn init_metrics_runtime(token: CancellationToken) { let token_clone = token.clone(); tokio::spawn(async move { let mut interval = tokio::time::interval(bucket_interval); + let tombstone_cycles = config.replication_bandwidth_zero_tombstone_cycles; + let mut has_seen_bucket_snapshot = false; + let mut prev_bucket_keys: HashSet = HashSet::new(); + let mut bucket_zero_tombstones: HashMap = HashMap::new(); loop { tokio::select! { _ = interval.tick() => { let stats = collect_bucket_stats().await; - let metrics = collect_bucket_metrics(&stats); + update_series_zero_tombstones( + &mut has_seen_bucket_snapshot, + &mut prev_bucket_keys, + &mut bucket_zero_tombstones, + bucket_live_keys(&stats), + tombstone_cycles, + ); + let mut metrics = collect_bucket_metrics(&stats); + metrics.extend(collect_bucket_zero_tombstone_metrics(&bucket_zero_tombstones)); report_metrics(&metrics); + expire_series_zero_tombstones(&mut bucket_zero_tombstones); } _ = token_clone.cancelled() => { warn!(event = EVENT_METRICS_RUNTIME_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_METRICS_RUNTIME, collector = "bucket_stats", state = "cancelled", "metrics runtime state changed"); @@ -644,6 +909,10 @@ pub fn init_metrics_runtime(token: CancellationToken) { let token_clone = token.clone(); tokio::spawn(async move { let mut interval = tokio::time::interval(audit_interval); + let tombstone_cycles = config.replication_bandwidth_zero_tombstone_cycles; + let mut has_seen_audit_snapshot = false; + let mut prev_audit_target_keys: HashSet = HashSet::new(); + let mut audit_zero_tombstones: HashMap = HashMap::new(); loop { tokio::select! { _ = interval.tick() => { @@ -656,8 +925,17 @@ pub fn init_metrics_runtime(token: CancellationToken) { total_messages: snapshot.total_messages, }) .collect::>(); - let metrics = collect_audit_metrics(&stats); + update_series_zero_tombstones( + &mut has_seen_audit_snapshot, + &mut prev_audit_target_keys, + &mut audit_zero_tombstones, + audit_target_live_keys(&stats), + tombstone_cycles, + ); + let mut metrics = collect_audit_metrics(&stats); + metrics.extend(collect_audit_zero_tombstone_metrics(&audit_zero_tombstones)); report_metrics(&metrics); + expire_series_zero_tombstones(&mut audit_zero_tombstones); } _ = token_clone.cancelled() => { warn!(event = EVENT_METRICS_RUNTIME_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_METRICS_RUNTIME, collector = "audit_target_stats", state = "cancelled", "metrics runtime state changed"); @@ -671,6 +949,10 @@ pub fn init_metrics_runtime(token: CancellationToken) { let token_clone = token.clone(); tokio::spawn(async move { let mut interval = tokio::time::interval(notification_interval); + let tombstone_cycles = config.replication_bandwidth_zero_tombstone_cycles; + let mut has_seen_notification_target_snapshot = false; + let mut prev_notification_target_keys: HashSet = HashSet::new(); + let mut notification_target_zero_tombstones: HashMap = HashMap::new(); loop { tokio::select! { _ = interval.tick() => { @@ -692,8 +974,19 @@ pub fn init_metrics_runtime(token: CancellationToken) { total_messages: snapshot.total_messages, }) .collect::>(); + update_series_zero_tombstones( + &mut has_seen_notification_target_snapshot, + &mut prev_notification_target_keys, + &mut notification_target_zero_tombstones, + notification_target_live_keys(&target_stats), + tombstone_cycles, + ); metrics.extend(collect_notification_target_metrics(&target_stats)); + metrics.extend(collect_notification_target_zero_tombstone_metrics( + ¬ification_target_zero_tombstones, + )); report_metrics(&metrics); + expire_series_zero_tombstones(&mut notification_target_zero_tombstones); } _ = token_clone.cancelled() => { warn!(event = EVENT_METRICS_RUNTIME_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_METRICS_RUNTIME, collector = "notification_stats", state = "cancelled", "metrics runtime state changed"); @@ -1154,4 +1447,44 @@ mod tests { 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)); } + + #[test] + fn bucket_tombstones_zero_removed_buckets_then_expire() { + let mut has_seen_valid_snapshot = false; + let mut prev_live_keys = HashSet::new(); + let mut zero_tombstones = HashMap::new(); + let live_stats = vec![crate::metrics::collectors::BucketStats { + name: "tmp".to_string(), + size_bytes: 1024, + objects_count: 8, + quota_bytes: 2048, + }]; + + update_series_zero_tombstones( + &mut has_seen_valid_snapshot, + &mut prev_live_keys, + &mut zero_tombstones, + bucket_live_keys(&live_stats), + 2, + ); + assert!(zero_tombstones.is_empty()); + + update_series_zero_tombstones(&mut has_seen_valid_snapshot, &mut prev_live_keys, &mut zero_tombstones, HashSet::new(), 2); + assert_eq!(zero_tombstones.get("tmp"), Some(&2)); + + let metrics = collect_bucket_zero_tombstone_metrics(&zero_tombstones); + assert_eq!(metrics.len(), 3); + assert!(metrics.iter().all(|metric| metric.value == 0.0)); + assert!( + metrics + .iter() + .all(|metric| { metric.labels.iter().any(|(key, value)| *key == "bucket" && value == "tmp") }) + ); + + expire_series_zero_tombstones(&mut zero_tombstones); + assert_eq!(zero_tombstones.get("tmp"), Some(&1)); + + expire_series_zero_tombstones(&mut zero_tombstones); + assert!(zero_tombstones.is_empty()); + } }