diff --git a/crates/obs/src/metrics/scheduler.rs b/crates/obs/src/metrics/scheduler.rs index 0d0d46b5b..5a8e0c808 100644 --- a/crates/obs/src/metrics/scheduler.rs +++ b/crates/obs/src/metrics/scheduler.rs @@ -97,6 +97,7 @@ use crate::metrics::stats_collector::{ collect_process_metric_bundle_with, collect_replication_stats, collect_scanner_metric_stats, collect_system_cpu_and_memory_stats_with, }; +use crate::telemetry::retire_metric_series; use futures_util::FutureExt; use rustfs_audit::audit_target_metrics; use rustfs_io_metrics::ProcessSampler; @@ -638,10 +639,12 @@ fn update_series_zero_tombstones( *has_seen_valid_snapshot = true; } -fn expire_series_zero_tombstones(zero_tombstones: &mut HashMap) { +fn expire_series_zero_tombstones(zero_tombstones: &mut HashMap) -> Vec { + let mut expired = Vec::new(); if !zero_tombstones.is_empty() { - zero_tombstones.retain(|_, remaining| { + zero_tombstones.retain(|key, remaining| { if *remaining <= 1 { + expired.push(key.clone()); false } else { *remaining -= 1; @@ -649,6 +652,7 @@ fn expire_series_zero_tombstones(zero_tombstones: &mut HashMap) { } }); } + expired } fn bucket_live_keys(stats: &[crate::metrics::collectors::BucketStats]) -> HashSet { @@ -673,6 +677,14 @@ fn collect_bucket_zero_tombstone_metrics(zero_tombstones: &HashMap usize { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.to_string()); + let labels = [("bucket", bucket_label.clone())]; + retire_metric_series(&BUCKET_USAGE_BYTES_MD.get_full_metric_name(), &labels) + + retire_metric_series(&BUCKET_OBJECTS_TOTAL_MD.get_full_metric_name(), &labels) + + retire_metric_series(&BUCKET_QUOTA_BYTES_MD.get_full_metric_name(), &labels) +} + fn bucket_usage_live_keys(stats: &[crate::metrics::collectors::BucketUsageStats]) -> HashSet { stats.iter().map(|stat| stat.bucket.clone()).collect() } @@ -750,6 +762,24 @@ fn collect_bucket_usage_zero_tombstone_metrics( zero_metrics } +fn retire_bucket_usage_metric_series(bucket: &str) -> usize { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.to_string()); + let labels = [(USAGE_BUCKET_LABEL, bucket_label.clone())]; + retire_metric_series(&USAGE_BUCKET_TOTAL_BYTES_MD.get_full_metric_name(), &labels) + + retire_metric_series(&USAGE_BUCKET_OBJECTS_TOTAL_MD.get_full_metric_name(), &labels) + + retire_metric_series(&USAGE_BUCKET_VERSIONS_COUNT_MD.get_full_metric_name(), &labels) + + retire_metric_series(&USAGE_BUCKET_DELETE_MARKERS_COUNT_MD.get_full_metric_name(), &labels) + + retire_metric_series(&USAGE_BUCKET_QUOTA_TOTAL_BYTES_MD.get_full_metric_name(), &labels) +} + +fn retire_bucket_usage_distribution_series(metric_name: String, bucket: &str, range: &str) -> usize { + let labels = [ + (USAGE_RANGE_LABEL, Cow::Owned(range.to_string())), + (USAGE_BUCKET_LABEL, Cow::Owned(bucket.to_string())), + ]; + retire_metric_series(&metric_name, &labels) +} + fn audit_target_live_keys(stats: &[AuditTargetStats]) -> HashSet { stats.iter().map(|stat| stat.target_id.clone()).collect() } @@ -778,6 +808,13 @@ fn collect_audit_zero_tombstone_metrics(zero_tombstones: &HashMap usize { + let labels = [(AUDIT_TARGET_ID_LABEL, Cow::Owned(target_id.to_string()))]; + retire_metric_series(&AUDIT_FAILED_MESSAGES_MD.get_full_metric_name(), &labels) + + retire_metric_series(&AUDIT_TARGET_QUEUE_LENGTH_MD.get_full_metric_name(), &labels) + + retire_metric_series(&AUDIT_TOTAL_MESSAGES_MD.get_full_metric_name(), &labels) +} + fn notification_target_live_keys(stats: &[NotificationTargetStats]) -> HashSet { stats .iter() @@ -816,6 +853,16 @@ fn collect_notification_target_zero_tombstone_metrics( zero_metrics } +fn retire_notification_target_metric_series(target_id: &str, target_type: &str) -> usize { + let labels = [ + (NOTIFICATION_TARGET_ID_LABEL, Cow::Owned(target_id.to_string())), + (NOTIFICATION_TARGET_TYPE_LABEL, Cow::Owned(target_type.to_string())), + ]; + retire_metric_series(&NOTIFICATION_TARGET_FAILED_MESSAGES_MD.get_full_metric_name(), &labels) + + retire_metric_series(&NOTIFICATION_TARGET_QUEUE_LENGTH_MD.get_full_metric_name(), &labels) + + retire_metric_series(&NOTIFICATION_TARGET_TOTAL_MESSAGES_MD.get_full_metric_name(), &labels) +} + fn update_repl_bw_zero_tombstones( monitor_available: bool, has_seen_valid_snapshot: &mut bool, @@ -869,17 +916,21 @@ fn collect_repl_bw_zero_tombstone_metrics(zero_tombstones: &HashMap) { - if monitor_available && !zero_tombstones.is_empty() { - zero_tombstones.retain(|_, remaining| { - if *remaining <= 1 { - false - } else { - *remaining -= 1; - true - } - }); +fn retire_repl_bw_metric_series(bucket: &str, target_arn: &str) -> usize { + let labels = [ + (BUCKET_L, Cow::Owned(bucket.to_string())), + (TARGET_ARN_L, Cow::Owned(target_arn.to_string())), + ]; + retire_metric_series(&BUCKET_REPL_BANDWIDTH_LIMIT_MD.get_full_metric_name(), &labels) + + retire_metric_series(&BUCKET_REPL_BANDWIDTH_CURRENT_MD.get_full_metric_name(), &labels) +} + +fn expire_repl_bw_zero_tombstones(monitor_available: bool, zero_tombstones: &mut HashMap) -> Vec { + if !monitor_available { + return Vec::new(); } + + expire_series_zero_tombstones(zero_tombstones) } /// Initialize all metrics collectors. @@ -1005,9 +1056,23 @@ pub fn init_metrics_runtime(token: CancellationToken) { &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); + for bucket in expire_series_zero_tombstones(&mut bucket_usage_zero_tombstones) { + let _ = retire_bucket_usage_metric_series(&bucket); + } + for (bucket, range) in expire_series_zero_tombstones(&mut bucket_usage_object_size_zero_tombstones) { + let _ = retire_bucket_usage_distribution_series( + USAGE_BUCKET_OBJECT_SIZE_DISTRIBUTION_MD.get_full_metric_name(), + &bucket, + &range, + ); + } + for (bucket, range) in expire_series_zero_tombstones(&mut bucket_usage_version_zero_tombstones) { + let _ = retire_bucket_usage_distribution_series( + USAGE_BUCKET_OBJECT_VERSION_COUNT_DISTRIBUTION_MD.get_full_metric_name(), + &bucket, + &range, + ); + } } if !metrics.is_empty() { @@ -1047,7 +1112,9 @@ pub fn init_metrics_runtime(token: CancellationToken) { 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); + for bucket in expire_series_zero_tombstones(&mut bucket_zero_tombstones) { + let _ = retire_bucket_metric_series(&bucket); + } }).await; } _ = token_clone.cancelled() => { @@ -1125,7 +1192,9 @@ pub fn init_metrics_runtime(token: CancellationToken) { report_metrics(&metrics); // Phase-2: after N cycles, stop reporting -> series becomes absent after expiration. - expire_repl_bw_zero_tombstones(monitor_available, &mut zero_tombstones); + for (bucket, target_arn) in expire_repl_bw_zero_tombstones(monitor_available, &mut zero_tombstones) { + let _ = retire_repl_bw_metric_series(&bucket, &target_arn); + } }, ).await; } @@ -1168,7 +1237,9 @@ pub fn init_metrics_runtime(token: CancellationToken) { 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); + for target_id in expire_series_zero_tombstones(&mut audit_zero_tombstones) { + let _ = retire_audit_target_metric_series(&target_id); + } }).await; } _ = token_clone.cancelled() => { @@ -1221,7 +1292,9 @@ pub fn init_metrics_runtime(token: CancellationToken) { ¬ification_target_zero_tombstones, )); report_metrics(&metrics); - expire_series_zero_tombstones(&mut notification_target_zero_tombstones); + for (target_id, target_type) in expire_series_zero_tombstones(&mut notification_target_zero_tombstones) { + let _ = retire_notification_target_metric_series(&target_id, &target_type); + } }).await; } _ = token_clone.cancelled() => { @@ -1702,10 +1775,12 @@ mod tests { 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); + let expired = expire_repl_bw_zero_tombstones(true, &mut zero_tombstones); + assert!(expired.is_empty()); assert_eq!(zero_tombstones.get(&key), Some(&1)); - expire_repl_bw_zero_tombstones(true, &mut zero_tombstones); + let expired = expire_repl_bw_zero_tombstones(true, &mut zero_tombstones); + assert_eq!(expired, vec![key]); assert!(zero_tombstones.is_empty()); } @@ -1766,7 +1841,8 @@ mod tests { 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); + let expired = expire_repl_bw_zero_tombstones(false, &mut zero_tombstones); + assert!(expired.is_empty()); assert_eq!(zero_tombstones.get(&repl_bw_key("videos", "arn:rustfs:replication:target-b")), Some(&1)); } @@ -1803,10 +1879,12 @@ mod tests { .all(|metric| { metric.labels.iter().any(|(key, value)| *key == "bucket" && value == "tmp") }) ); - expire_series_zero_tombstones(&mut zero_tombstones); + let expired = expire_series_zero_tombstones(&mut zero_tombstones); + assert!(expired.is_empty()); assert_eq!(zero_tombstones.get("tmp"), Some(&1)); - expire_series_zero_tombstones(&mut zero_tombstones); + let expired = expire_series_zero_tombstones(&mut zero_tombstones); + assert_eq!(expired, vec!["tmp".to_string()]); assert!(zero_tombstones.is_empty()); } diff --git a/crates/obs/src/telemetry/mod.rs b/crates/obs/src/telemetry/mod.rs index 85b4573b4..cc5414a9f 100644 --- a/crates/obs/src/telemetry/mod.rs +++ b/crates/obs/src/telemetry/mod.rs @@ -52,7 +52,7 @@ mod rolling; use crate::TelemetryError; use crate::config::OtelConfig; pub use guard::OtelGuard; -pub use recorder::Recorder; +pub use recorder::{Recorder, retire_metric_series}; use rustfs_config::observability::ENV_OBS_LOG_DIRECTORY; use rustfs_config::{DEFAULT_LOG_LEVEL, ENVIRONMENT, observability::DEFAULT_OBS_ENVIRONMENT_PRODUCTION}; use rustfs_utils::get_env_opt_str; diff --git a/crates/obs/src/telemetry/otel.rs b/crates/obs/src/telemetry/otel.rs index 1e400bf7c..9316649db 100644 --- a/crates/obs/src/telemetry/otel.rs +++ b/crates/obs/src/telemetry/otel.rs @@ -42,7 +42,7 @@ use crate::global::set_observability_metric_enabled; use crate::telemetry::filter::build_env_filter; use crate::telemetry::guard::{OtelGuard, ProfilingAgent}; use crate::telemetry::local::{build_json_log_layer, spawn_cleanup_task}; -use crate::telemetry::recorder::Recorder; +use crate::telemetry::recorder::{Recorder, install_process_global_recorder}; use crate::telemetry::resource::build_resource; use crate::telemetry::rolling::{RollingAppender, Rotation}; // Import helper functions from local.rs (sibling module) @@ -432,7 +432,7 @@ fn build_meter_provider( .build(); global::set_meter_provider(provider.clone()); - metrics::set_global_recorder(recorder).map_err(|e| TelemetryError::InstallMetricsRecorder(e.to_string()))?; + install_process_global_recorder(recorder).map_err(|e| TelemetryError::InstallMetricsRecorder(e.to_string()))?; set_observability_metric_enabled(true); Ok(Some(provider)) } diff --git a/crates/obs/src/telemetry/recorder.rs b/crates/obs/src/telemetry/recorder.rs index 1b49642d2..d8b9e51fe 100644 --- a/crates/obs/src/telemetry/recorder.rs +++ b/crates/obs/src/telemetry/recorder.rs @@ -13,7 +13,7 @@ // limitations under the License. use crate::GlobalError; -use metrics::{Counter, CounterFn, Gauge, GaugeFn, Histogram, HistogramFn, Key, KeyName, Metadata, SharedString, Unit}; +use metrics::{Counter, CounterFn, Gauge, GaugeFn, Histogram, HistogramFn, Key, KeyName, Label, Metadata, SharedString, Unit}; use opentelemetry::{ InstrumentationScope, InstrumentationScopeBuilder, KeyValue, global, metrics::{Meter, MeterProvider}, @@ -24,7 +24,7 @@ use std::{ collections::HashMap, ops::Deref, sync::{ - Arc, Mutex, RwLock, + Arc, Mutex, OnceLock, RwLock, atomic::{AtomicU64, Ordering}, }, }; @@ -33,6 +33,7 @@ use tracing::error; const LOG_COMPONENT_OBS: &str = "obs"; const LOG_SUBSYSTEM_RECORDER: &str = "recorder"; const EVENT_RECORDER_STATE: &str = "recorder_state"; +static GLOBAL_RECORDER: OnceLock = OnceLock::new(); macro_rules! configure_builder { ($builder:expr, $metadata:expr) => {{ @@ -96,6 +97,7 @@ impl Builder { pub fn install(self) -> Result<(SdkMeterProvider, Recorder), GlobalError> { let (provider, recorder) = self.build(); metrics::set_global_recorder(recorder.clone())?; + remember_global_recorder(&recorder); Ok((provider, recorder)) } @@ -178,6 +180,18 @@ impl Recorder { value } + fn remove_cached_metric(lock: &RwLock>, key: &Key, metric_type: &str) -> bool { + let mut cache = match lock.write() { + Ok(g) => g, + Err(e) => { + error!(event = EVENT_RECORDER_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_RECORDER, metric_type = %metric_type, result = "cache_remove_lock_poisoned", error = %e, "recorder state changed"); + e.into_inner() + } + }; + + cache.remove(key).is_some() + } + fn with_metadata_lock(&self, f: F) -> R where F: FnOnce(&mut HashMap) -> R, @@ -198,6 +212,39 @@ impl Recorder { fn get_metadata_for_builder(&self, key_name: &str) -> Option { self.with_metadata_lock(|metadata| metadata.get(key_name).cloned()) } + + fn retire_metric_series_key(&self, key: &Key) -> usize { + let mut retired = 0usize; + retired += usize::from(Self::remove_cached_metric(&self.cached_counters, key, "counter")); + retired += usize::from(Self::remove_cached_metric(&self.cached_gauges, key, "gauge")); + retired += usize::from(Self::remove_cached_metric(&self.cached_histograms, key, "histogram")); + retired + } +} + +fn remember_global_recorder(recorder: &Recorder) { + let _ = GLOBAL_RECORDER.set(recorder.clone()); +} + +pub(crate) fn install_process_global_recorder(recorder: Recorder) -> Result<(), metrics::SetRecorderError> { + metrics::set_global_recorder(recorder.clone())?; + remember_global_recorder(&recorder); + Ok(()) +} + +pub fn retire_metric_series(name: &str, labels: &[(&'static str, Cow<'static, str>)]) -> usize { + let Some(recorder) = GLOBAL_RECORDER.get() else { + return 0; + }; + + let key = Key::from_parts( + name.to_string(), + labels + .iter() + .map(|(key, value)| Label::new((*key).to_string(), value.to_string())) + .collect::>(), + ); + recorder.retire_metric_series_key(&key) } impl Deref for Recorder { @@ -499,4 +546,48 @@ mod tests { let _ = recorder.register_counter(&second, &meta); } + + #[test] + fn retire_metric_series_key_removes_cached_counter() { + let recorder = test_recorder(); + let key = Key::from_parts("retired_counter", vec![metrics::Label::new("bucket", "tmp")]); + let meta = test_metadata(); + + let _counter = recorder.register_counter(&key, &meta); + assert_eq!(recorder.cached_counters.read().unwrap().len(), 1); + + let retired = recorder.retire_metric_series_key(&key); + assert_eq!(retired, 1); + assert!(recorder.cached_counters.read().unwrap().is_empty()); + } + + #[test] + fn retiring_churned_series_bounds_the_cache() { + // #1026: under bucket/target churn the recorder handle cache used to grow + // unbounded. After each churned series is retired the cache must fall back + // to a bounded size, and a same-named label set must re-register cleanly. + let recorder = test_recorder(); + let meta = test_metadata(); + + let mut keys = Vec::with_capacity(3000); + for i in 0..3000 { + let key = Key::from_parts("churn_counter", vec![metrics::Label::new("bucket", format!("bucket-{i}"))]); + let _ = recorder.register_counter(&key, &meta); + keys.push(key); + } + assert_eq!(recorder.cached_counters.read().unwrap().len(), 3000); + + for key in &keys { + assert_eq!(recorder.retire_metric_series_key(key), 1); + } + assert!( + recorder.cached_counters.read().unwrap().is_empty(), + "recorder cache must be bounded after churned series are retired" + ); + + // A previously retired label set can be registered and observed again. + let reused = &keys[0]; + let _ = recorder.register_counter(reused, &meta); + assert_eq!(recorder.cached_counters.read().unwrap().len(), 1); + } }