diff --git a/.docker/observability/README.md b/.docker/observability/README.md index 3d40319b1..e24e928a3 100644 --- a/.docker/observability/README.md +++ b/.docker/observability/README.md @@ -56,7 +56,7 @@ docker compose -f docker-compose.yml -f docker-compose.override.yml up -d - docker compose v2.0.0 or before ```bash -docke-compose -f docker-compose.yml -f docker-compose.override.yml up -d +docker-compose -f docker-compose.yml -f docker-compose.override.yml up -d ``` To access the Grafana dashboard, navigate to `http://localhost:3000` in your browser. The default username and password diff --git a/Cargo.lock b/Cargo.lock index 58743d0f3..2ad5b40ff 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3661,7 +3661,7 @@ dependencies = [ "opentelemetry-semantic-conventions", "percent-encoding", "pin-project", - "prost 0.14.3", + "prost", "prost-types", "reqwest 0.13.2", "rustc_version", @@ -3772,7 +3772,7 @@ dependencies = [ "md5", "percent-encoding", "pin-project", - "prost 0.14.3", + "prost", "prost-types", "reqwest 0.13.2", "serde", @@ -4455,9 +4455,9 @@ checksum = "92ecc6618181def0457392ccd0ee51198e065e016d1d527a7ac1b6dc7c1f09d2" [[package]] name = "jemalloc_pprof" -version = "0.8.1" +version = "0.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "74ff642505c7ce8d31c0d43ec0e235c6fd4585d9b8172d8f9dd04d36590200b5" +checksum = "8a0d44c349cfe2654897fadcb9de4f0bfbf48288ec344f700b2bd59f152dd209" dependencies = [ "anyhow", "libc", @@ -4690,9 +4690,9 @@ checksum = "2c4a545a15244c7d945065b5d392b2d2d7f21526fba56ce51467b06ed445e8f7" [[package]] name = "libc" -version = "0.2.180" +version = "0.2.181" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bcc35a38544a891a5f7c865aca548a982ccb3b8650a5b06d0fd33a10283c56fc" +checksum = "459427e2af2b9c839b132acb702a1c654d95e10f8c326bfc2ad11310e458b1c5" [[package]] name = "libloading" @@ -4906,9 +4906,9 @@ dependencies = [ [[package]] name = "mappings" -version = "0.7.1" +version = "0.7.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "db4d277bb50d4508057e7bddd7fcd19ef4a4cc38051b6a5a36868d75ae2cbeb9" +checksum = "8bab1e61a4b76757edb59cd81fcaa7f3ba9018d43b527d9abfad877b4c6c60f2" dependencies = [ "anyhow", "libc", @@ -5504,7 +5504,7 @@ dependencies = [ "opentelemetry-http", "opentelemetry-proto", "opentelemetry_sdk", - "prost 0.14.3", + "prost", "reqwest 0.12.28", "thiserror 2.0.18", "tracing", @@ -5518,7 +5518,7 @@ checksum = "a7175df06de5eaee9909d4805a3d07e28bb752c34cab57fa9cff549da596b30f" dependencies = [ "opentelemetry", "opentelemetry_sdk", - "prost 0.14.3", + "prost", "tonic", "tonic-prost", ] @@ -5982,9 +5982,9 @@ dependencies = [ [[package]] name = "pprof_util" -version = "0.8.1" +version = "0.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4429d44e5e2c8a69399fc0070379201eed018e3df61e04eb7432811df073c224" +checksum = "eea0cc524de808a6d98d192a3d99fe95617031ad4a52ec0a0f987ef4432e8fe1" dependencies = [ "anyhow", "backtrace", @@ -5992,7 +5992,7 @@ dependencies = [ "inferno 0.12.4", "num", "paste", - "prost 0.13.5", + "prost", ] [[package]] @@ -6069,16 +6069,6 @@ dependencies = [ "thiserror 2.0.18", ] -[[package]] -name = "prost" -version = "0.13.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2796faa41db3ec313a31f7624d9286acf277b52de526150b7e69f3debf891ee5" -dependencies = [ - "bytes", - "prost-derive 0.13.5", -] - [[package]] name = "prost" version = "0.14.3" @@ -6086,7 +6076,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d2ea70524a2f82d518bce41317d0fae74151505651af45faf1ffbd6fd33f0568" dependencies = [ "bytes", - "prost-derive 0.14.3", + "prost-derive", ] [[package]] @@ -6101,7 +6091,7 @@ dependencies = [ "multimap", "petgraph", "prettyplease", - "prost 0.14.3", + "prost", "prost-types", "pulldown-cmark", "pulldown-cmark-to-cmark", @@ -6110,19 +6100,6 @@ dependencies = [ "tempfile", ] -[[package]] -name = "prost-derive" -version = "0.13.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" -dependencies = [ - "anyhow", - "itertools 0.14.0", - "proc-macro2", - "quote", - "syn 2.0.114", -] - [[package]] name = "prost-derive" version = "0.14.3" @@ -6142,7 +6119,7 @@ version = "0.14.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8991c4cbdb8bc5b11f0b074ffe286c30e523de90fee5ba8132f1399f23cb3dd7" dependencies = [ - "prost 0.14.3", + "prost", ] [[package]] @@ -6887,6 +6864,7 @@ dependencies = [ "rustfs-kms", "rustfs-lock", "rustfs-madmin", + "rustfs-metrics", "rustfs-notify", "rustfs-obs", "rustfs-policy", @@ -7269,6 +7247,19 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "rustfs-metrics" +version = "0.0.5" +dependencies = [ + "metrics", + "rustfs-ecstore", + "rustfs-utils", + "sysinfo", + "tokio", + "tokio-util", + "tracing", +] + [[package]] name = "rustfs-notify" version = "0.0.5" @@ -7381,7 +7372,7 @@ name = "rustfs-protos" version = "0.0.5" dependencies = [ "flatbuffers", - "prost 0.14.3", + "prost", "rustfs-common", "tonic", "tonic-prost", @@ -8786,12 +8777,12 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.24.0" +version = "3.25.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "655da9c7eb6305c55742045d5a8d2037996d61d8de95806335c7c86ce0f82e9c" +checksum = "0136791f7c95b1f6dd99f9cc786b91bb81c3800b639b3478e561ddb7be95e5f1" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.1", "once_cell", "rustix 1.1.3", "windows-sys 0.61.2", @@ -9163,7 +9154,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6c55a2d6a14174563de34409c9f92ff981d006f56da9c6ecd40d9d4a31500b0" dependencies = [ "bytes", - "prost 0.14.3", + "prost", "tonic", ] diff --git a/Cargo.toml b/Cargo.toml index f04a9a580..d118a09ef 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -31,9 +31,11 @@ members = [ "crates/lock", # Distributed locking implementation "crates/madmin", # Management dashboard and admin API interface "crates/mcp", # MCP server for S3 operations + "crates/metrics", # Metrics collection and reporting "crates/notify", # Notification system for events "crates/obs", # Observability utilities "crates/policy", # Policy management + "crates/protocols", # Protocol implementations (FTPS, SFTP, etc.) "crates/protos", # Protocol buffer definitions "crates/rio", # Rust I/O utilities and abstractions "crates/s3select-api", # S3 Select API interface @@ -45,7 +47,6 @@ members = [ "crates/utils", # Utility functions and helpers "crates/workers", # Worker thread pools and task scheduling "crates/zip", # ZIP file handling and compression - "crates/protocols", # Protocol implementations (FTPS, SFTP, etc.) ] resolver = "2" @@ -85,6 +86,7 @@ rustfs-kms = { path = "crates/kms", version = "0.0.5" } rustfs-lock = { path = "crates/lock", version = "0.0.5" } rustfs-madmin = { path = "crates/madmin", version = "0.0.5" } rustfs-mcp = { path = "crates/mcp", version = "0.0.5" } +rustfs-metrics = { path = "crates/metrics", version = "0.0.5" } rustfs-notify = { path = "crates/notify", version = "0.0.5" } rustfs-obs = { path = "crates/obs", version = "0.0.5" } rustfs-policy = { path = "crates/policy", version = "0.0.5" } @@ -205,7 +207,7 @@ hex-simd = "0.8.0" highway = { version = "1.3.0" } ipnetwork = { version = "0.21.1", features = ["serde"] } lazy_static = "1.5.0" -libc = "0.2.180" +libc = "0.2.181" libsystemd = "0.7.2" local-ip-address = "0.6.10" lz4 = "1.28.1" @@ -244,7 +246,7 @@ starshard = { version = "1.1.0", features = ["rayon", "async", "serde"] } strum = { version = "0.27.2", features = ["derive"] } sysinfo = "0.38.1" temp-env = "0.3.6" -tempfile = "3.24.0" +tempfile = "3.25.0" test-case = "3.3.1" thiserror = "2.0.18" tracing = { version = "0.1.44" } @@ -285,7 +287,7 @@ tikv-jemallocator = { version = "0.6", features = ["profiling", "stats", "unpref # Used to control and obtain statistics for jemalloc at runtime tikv-jemalloc-ctl = { version = "0.6", features = ["use_std", "stats", "profiling"] } # Used to generate pprof-compatible memory profiling data and support symbolization and flame graphs -jemalloc_pprof = { version = "0.8.1", features = ["symbolize", "flamegraph"] } +jemalloc_pprof = { version = "0.8.2", features = ["symbolize", "flamegraph"] } # Used to generate CPU performance analysis data and flame diagrams pprof = { version = "0.15.0", features = ["flamegraph", "protobuf-codec"] } diff --git a/crates/metrics/Cargo.toml b/crates/metrics/Cargo.toml new file mode 100644 index 000000000..7f3a29728 --- /dev/null +++ b/crates/metrics/Cargo.toml @@ -0,0 +1,38 @@ +# Copyright 2024 RustFS Team +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +[package] +name = "rustfs-metrics" +edition.workspace = true +license.workspace = true +repository.workspace = true +rust-version.workspace = true +version.workspace = true +homepage.workspace = true +description.workspace = true +keywords.workspace = true +categories.workspace = true +authors.workspace = true + +[dependencies] +rustfs-ecstore = { workspace = true } +rustfs-utils = { workspace = true } +metrics = { workspace = true } +sysinfo = { workspace = true } +tokio = { workspace = true, features = ["rt-multi-thread", "rt", "time", "macros"] } +tokio-util = { workspace = true } +tracing = { workspace = true } + +[lints] +workspace = true diff --git a/crates/metrics/src/collectors/bucket.rs b/crates/metrics/src/collectors/bucket.rs new file mode 100644 index 000000000..d3d3a944e --- /dev/null +++ b/crates/metrics/src/collectors/bucket.rs @@ -0,0 +1,201 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Per-bucket metrics collector. +//! +//! Collects usage metrics for each bucket in the cluster, including +//! size, object counts, and quota information. + +use crate::MetricType; +use crate::format::PrometheusMetric; +use std::borrow::Cow; + +/// Usage statistics for a single bucket. +#[derive(Debug, Clone, Default)] +pub struct BucketStats { + /// Bucket name + pub name: String, + /// Total bytes used by the bucket + pub size_bytes: u64, + /// Total number of objects in the bucket + pub objects_count: u64, + /// Quota limit in bytes (0 means no quota) + pub quota_bytes: u64, +} + +// Static metric definitions +const METRIC_SIZE: &str = "rustfs_bucket_usage_bytes"; +const METRIC_OBJECTS: &str = "rustfs_bucket_objects_total"; +const METRIC_QUOTA: &str = "rustfs_bucket_quota_bytes"; + +const HELP_SIZE: &str = "Total bytes used by the bucket"; +const HELP_OBJECTS: &str = "Total number of objects in the bucket"; +const HELP_QUOTA: &str = "Quota limit in bytes for the bucket"; + +/// Collects per-bucket usage metrics from the provided bucket statistics. +/// +/// # Metrics Produced +/// +/// For each bucket, the following metrics are produced with a `bucket` label: +/// +/// - `rustfs_bucket_usage_bytes`: Total bytes used by the bucket +/// - `rustfs_bucket_objects_total`: Total number of objects in the bucket +/// - `rustfs_bucket_quota_bytes`: Quota limit in bytes (0 if no quota configured) +/// +/// # Arguments +/// +/// * `buckets` - Slice of bucket statistics +/// +/// # Example +/// +/// ``` +/// use rustfs_metrics::collectors::{collect_bucket_metrics, BucketStats}; +/// +/// let buckets = vec![ +/// BucketStats { +/// name: "my-bucket".to_string(), +/// size_bytes: 1_000_000, +/// objects_count: 100, +/// quota_bytes: 10_000_000, +/// }, +/// ]; +/// let metrics = collect_bucket_metrics(&buckets); +/// assert_eq!(metrics.len(), 3); // size, objects, quota +/// ``` +#[must_use] +#[inline] +pub fn collect_bucket_metrics(buckets: &[BucketStats]) -> Vec { + if buckets.is_empty() { + return Vec::new(); + } + + let mut metrics = Vec::with_capacity(buckets.len() * 3); + + for bucket in buckets { + let bucket_label: Cow<'static, str> = Cow::Owned(bucket.name.clone()); + + // Bucket size in bytes + metrics.push( + PrometheusMetric::new(METRIC_SIZE, MetricType::Gauge, HELP_SIZE, bucket.size_bytes as f64) + .with_label("bucket", bucket_label.clone()), + ); + + // Object count + metrics.push( + PrometheusMetric::new(METRIC_OBJECTS, MetricType::Gauge, HELP_OBJECTS, bucket.objects_count as f64) + .with_label("bucket", bucket_label.clone()), + ); + + // Quota (always emit, 0 when no quota configured for consistent PromQL queries) + metrics.push( + PrometheusMetric::new(METRIC_QUOTA, MetricType::Gauge, HELP_QUOTA, bucket.quota_bytes as f64) + .with_label("bucket", bucket_label), + ); + } + + metrics +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::report_metrics; + + #[test] + fn test_collect_bucket_metrics() { + let buckets = vec![ + BucketStats { + name: "test-bucket".to_string(), + size_bytes: 1000, + objects_count: 50, + quota_bytes: 0, + }, + BucketStats { + name: "other-bucket".to_string(), + size_bytes: 2000, + objects_count: 100, + quota_bytes: 0, + }, + ]; + + let metrics = collect_bucket_metrics(&buckets); + report_metrics(&metrics); // This will compile and run, but we can't easily assert on the global recorder state here. + + // 2 buckets * 3 metrics each (size, objects, quota) = 6 metrics + assert_eq!(metrics.len(), 6); + + // Verify test-bucket metrics + let test_bucket_size = metrics + .iter() + .find(|m| m.name == METRIC_SIZE && m.labels.iter().any(|(k, v)| *k == "bucket" && v == "test-bucket")); + assert!(test_bucket_size.is_some()); + assert_eq!(test_bucket_size.map(|m| m.value), Some(1000.0)); + } + + #[test] + fn test_collect_bucket_metrics_with_quotas() { + let buckets = vec![BucketStats { + name: "quota-bucket".to_string(), + size_bytes: 500, + objects_count: 10, + quota_bytes: 10000, + }]; + + let metrics = collect_bucket_metrics(&buckets); + report_metrics(&metrics); + + // 1 bucket * 3 metrics (size, objects, quota) = 3 metrics + assert_eq!(metrics.len(), 3); + + // Verify quota metric exists + let quota_metric = metrics.iter().find(|m| m.name == METRIC_QUOTA); + assert!(quota_metric.is_some()); + assert_eq!(quota_metric.map(|m| m.value), Some(10000.0)); + } + + #[test] + fn test_collect_bucket_metrics_empty() { + let buckets: Vec = vec![]; + let metrics = collect_bucket_metrics(&buckets); + assert!(metrics.is_empty()); + } + + #[test] + fn test_collect_bucket_metrics_zero_quota_always_reported() { + let buckets = vec![BucketStats { + name: "no-quota-bucket".to_string(), + size_bytes: 100, + objects_count: 5, + quota_bytes: 0, + }]; + + let metrics = collect_bucket_metrics(&buckets); + report_metrics(&metrics); + + // Zero quota should still produce a quota metric with value 0 for consistent PromQL queries + assert_eq!(metrics.len(), 3); + let quota_metric = metrics.iter().find(|m| m.name == METRIC_QUOTA); + assert!(quota_metric.is_some()); + assert_eq!(quota_metric.map(|m| m.value), Some(0.0)); + } + + #[test] + fn test_bucket_stats_default() { + let stats = BucketStats::default(); + assert!(stats.name.is_empty()); + assert_eq!(stats.size_bytes, 0); + assert_eq!(stats.objects_count, 0); + assert_eq!(stats.quota_bytes, 0); + } +} diff --git a/crates/metrics/src/collectors/cluster.rs b/crates/metrics/src/collectors/cluster.rs new file mode 100644 index 000000000..09f2ead73 --- /dev/null +++ b/crates/metrics/src/collectors/cluster.rs @@ -0,0 +1,196 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Cluster-wide metrics collector. +//! +//! Collects aggregate metrics across the entire RustFS cluster including +//! total capacity, usage, and object counts. + +use crate::MetricType; +use crate::format::PrometheusMetric; + +/// Cluster capacity and usage statistics for metrics collection. +/// +/// This struct provides a decoupled interface for collecting cluster metrics +/// without depending on specific internal types. HTTP handlers should populate +/// this struct from their available data sources. +#[derive(Debug, Clone, Default)] +pub struct ClusterStats { + /// Total raw storage capacity across all disks in bytes + pub raw_capacity_bytes: u64, + /// Usable capacity after erasure coding overhead in bytes + pub usable_capacity_bytes: u64, + /// Currently used storage in bytes + pub used_bytes: u64, + /// Available free storage in bytes + pub free_bytes: u64, + /// Total number of objects in the cluster + pub objects_count: u64, + /// Total number of buckets in the cluster + pub buckets_count: u64, +} + +// Static metric definitions to avoid allocations +const METRIC_RAW_CAPACITY: &str = "rustfs_cluster_capacity_raw_total_bytes"; +const METRIC_USABLE_CAPACITY: &str = "rustfs_cluster_capacity_usable_total_bytes"; +const METRIC_USED: &str = "rustfs_cluster_capacity_used_bytes"; +const METRIC_FREE: &str = "rustfs_cluster_capacity_free_bytes"; +const METRIC_OBJECTS: &str = "rustfs_cluster_objects_total"; +const METRIC_BUCKETS: &str = "rustfs_cluster_buckets_total"; + +const HELP_RAW_CAPACITY: &str = "Total raw storage capacity in bytes across all disks"; +const HELP_USABLE_CAPACITY: &str = "Total usable storage capacity in bytes (accounting for erasure coding)"; +const HELP_USED: &str = "Total used storage capacity in bytes"; +const HELP_FREE: &str = "Total free storage capacity in bytes"; +const HELP_OBJECTS: &str = "Total number of objects in the cluster"; +const HELP_BUCKETS: &str = "Total number of buckets in the cluster"; + +/// Number of metrics produced by this collector. +const METRIC_COUNT: usize = 6; + +/// Collects cluster-wide metrics from the provided statistics. +/// +/// # Metrics Produced +/// +/// - `rustfs_cluster_capacity_raw_total_bytes`: Total raw storage capacity across all disks +/// - `rustfs_cluster_capacity_usable_total_bytes`: Usable capacity after erasure coding overhead +/// - `rustfs_cluster_capacity_used_bytes`: Currently used storage capacity +/// - `rustfs_cluster_capacity_free_bytes`: Available free storage capacity +/// - `rustfs_cluster_objects_total`: Total number of objects in the cluster +/// - `rustfs_cluster_buckets_total`: Total number of buckets in the cluster +/// +/// # Arguments +/// +/// * `stats` - Cluster statistics containing capacity and usage data +/// +/// # Example +/// +/// ``` +/// use rustfs_metrics::collectors::{collect_cluster_metrics, ClusterStats}; +/// +/// let stats = ClusterStats { +/// raw_capacity_bytes: 10_000_000_000, +/// usable_capacity_bytes: 8_000_000_000, +/// used_bytes: 2_000_000_000, +/// free_bytes: 6_000_000_000, +/// objects_count: 1000, +/// buckets_count: 10, +/// }; +/// let metrics = collect_cluster_metrics(&stats); +/// assert_eq!(metrics.len(), 6); +/// ``` +#[must_use] +#[inline] +pub fn collect_cluster_metrics(stats: &ClusterStats) -> Vec { + let mut metrics = Vec::with_capacity(METRIC_COUNT); + + metrics.push(PrometheusMetric::new( + METRIC_RAW_CAPACITY, + MetricType::Gauge, + HELP_RAW_CAPACITY, + stats.raw_capacity_bytes as f64, + )); + metrics.push(PrometheusMetric::new( + METRIC_USABLE_CAPACITY, + MetricType::Gauge, + HELP_USABLE_CAPACITY, + stats.usable_capacity_bytes as f64, + )); + metrics.push(PrometheusMetric::new(METRIC_USED, MetricType::Gauge, HELP_USED, stats.used_bytes as f64)); + metrics.push(PrometheusMetric::new(METRIC_FREE, MetricType::Gauge, HELP_FREE, stats.free_bytes as f64)); + metrics.push(PrometheusMetric::new( + METRIC_OBJECTS, + MetricType::Gauge, + HELP_OBJECTS, + stats.objects_count as f64, + )); + metrics.push(PrometheusMetric::new( + METRIC_BUCKETS, + MetricType::Gauge, + HELP_BUCKETS, + stats.buckets_count as f64, + )); + + metrics +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::report_metrics; + + #[test] + fn test_collect_cluster_metrics() { + let stats = ClusterStats { + raw_capacity_bytes: 3000, + usable_capacity_bytes: 2500, + used_bytes: 1200, + free_bytes: 1300, + objects_count: 100, + buckets_count: 5, + }; + + let metrics = collect_cluster_metrics(&stats); + report_metrics(&metrics); + + assert_eq!(metrics.len(), 6); + + // Verify raw capacity + let raw_capacity = metrics.iter().find(|m| m.name == METRIC_RAW_CAPACITY); + assert!(raw_capacity.is_some()); + assert_eq!(raw_capacity.map(|m| m.value), Some(3000.0)); + + // Verify used capacity + let used = metrics.iter().find(|m| m.name == METRIC_USED); + assert!(used.is_some()); + assert_eq!(used.map(|m| m.value), Some(1200.0)); + + // Verify object count + let objects = metrics.iter().find(|m| m.name == METRIC_OBJECTS); + assert!(objects.is_some()); + assert_eq!(objects.map(|m| m.value), Some(100.0)); + + // Verify bucket count + let buckets = metrics.iter().find(|m| m.name == METRIC_BUCKETS); + assert!(buckets.is_some()); + assert_eq!(buckets.map(|m| m.value), Some(5.0)); + } + + #[test] + fn test_collect_cluster_metrics_empty() { + let stats = ClusterStats::default(); + + let metrics = collect_cluster_metrics(&stats); + report_metrics(&metrics); + + assert_eq!(metrics.len(), 6); + + // All values should be zero + for metric in &metrics { + assert_eq!(metric.value, 0.0); + assert!(metric.labels.is_empty()); + } + } + + #[test] + fn test_cluster_stats_default() { + let stats = ClusterStats::default(); + assert_eq!(stats.raw_capacity_bytes, 0); + assert_eq!(stats.usable_capacity_bytes, 0); + assert_eq!(stats.used_bytes, 0); + assert_eq!(stats.free_bytes, 0); + assert_eq!(stats.objects_count, 0); + assert_eq!(stats.buckets_count, 0); + } +} diff --git a/crates/metrics/src/collectors/global.rs b/crates/metrics/src/collectors/global.rs new file mode 100644 index 000000000..9e1011f21 --- /dev/null +++ b/crates/metrics/src/collectors/global.rs @@ -0,0 +1,316 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::collectors::{ + BucketStats, ClusterStats, DiskStats, ResourceStats, collect_bucket_metrics, collect_cluster_metrics, collect_node_metrics, + collect_resource_metrics, +}; +use crate::constants::{ + DEFAULT_BUCKET_METRICS_INTERVAL, DEFAULT_CLUSTER_METRICS_INTERVAL, DEFAULT_NODE_METRICS_INTERVAL, + DEFAULT_RESOURCE_METRICS_INTERVAL, ENV_BUCKET_METRICS_INTERVAL, ENV_CLUSTER_METRICS_INTERVAL, ENV_DEFAULT_METRICS_INTERVAL, + ENV_NODE_METRICS_INTERVAL, ENV_RESOURCE_METRICS_INTERVAL, +}; +use crate::format::report_metrics; +use rustfs_ecstore::bucket::metadata_sys::get_quota_config; +use rustfs_ecstore::data_usage::load_data_usage_from_backend; +use rustfs_ecstore::pools::{get_total_usable_capacity, get_total_usable_capacity_free}; +use rustfs_ecstore::store_api::BucketOptions; +use rustfs_ecstore::{StorageAPI, new_object_layer_fn}; +use rustfs_utils::get_env_opt_u64; +use std::sync::OnceLock; +use std::time::{Duration, Instant}; +use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, System}; +use tokio_util::sync::CancellationToken; +use tracing::warn; + +/// Process start time for calculating uptime. +static PROCESS_START: OnceLock = OnceLock::new(); + +/// Get the process start time, initializing it on first call. +#[inline] +fn get_process_start() -> &'static Instant { + PROCESS_START.get_or_init(Instant::now) +} + +/// Collect cluster statistics from the storage layer. +async fn collect_cluster_stats() -> ClusterStats { + let Some(store) = new_object_layer_fn() else { + return ClusterStats::default(); + }; + + let storage_info = store.storage_info().await; + + let raw_capacity: u64 = storage_info.disks.iter().map(|d| d.total_space).sum(); + let used: u64 = storage_info.disks.iter().map(|d| d.used_space).sum(); + let usable_capacity = get_total_usable_capacity(&storage_info.disks, &storage_info) as u64; + let free = get_total_usable_capacity_free(&storage_info.disks, &storage_info) as u64; + + // Get bucket and object counts from data usage info + let (buckets_count, objects_count) = match load_data_usage_from_backend(store.clone()).await { + Ok(data_usage) => (data_usage.buckets_count, data_usage.objects_total_count), + Err(e) => { + warn!("Failed to load data usage from backend: {}", e); + // Fall back to bucket list for buckets_count, objects_count stays 0 + let buckets = store + .list_bucket(&BucketOptions { + cached: true, + ..Default::default() + }) + .await + .unwrap_or_else(|e| { + warn!("Failed to list buckets for cluster metrics: {}", e); + Vec::new() + }); + (buckets.len() as u64, 0) + } + }; + + ClusterStats { + raw_capacity_bytes: raw_capacity, + usable_capacity_bytes: usable_capacity, + used_bytes: used, + free_bytes: free, + objects_count, + buckets_count, + } +} + +/// Collect bucket statistics from the storage layer. +async fn collect_bucket_stats() -> Vec { + let Some(store) = new_object_layer_fn() else { + return Vec::new(); + }; + + // Load data usage info from backend to get bucket sizes and object counts + let data_usage = match load_data_usage_from_backend(store.clone()).await { + Ok(info) => Some(info), + Err(e) => { + warn!("Failed to load data usage from backend for bucket metrics: {}", e); + None + } + }; + + let buckets = match store + .list_bucket(&BucketOptions { + cached: true, + ..Default::default() + }) + .await + { + Ok(b) => b, + Err(e) => { + warn!("Failed to list buckets for metrics: {}", e); + return Vec::new(); + } + }; + + // Build bucket stats with real data from DataUsageInfo + let mut stats = Vec::with_capacity(buckets.len()); + for bucket in buckets { + if bucket.name.starts_with('.') { + continue; + } + + // Get size and objects_count from data usage info + let (size_bytes, objects_count) = data_usage + .as_ref() + .and_then(|du| du.buckets_usage.get(&bucket.name)) + .map(|bui| (bui.size, bui.objects_count)) + .unwrap_or((0, 0)); + + // Get quota from bucket metadata + let quota_bytes = match get_quota_config(&bucket.name).await { + Ok((quota, _)) => quota.get_quota_limit().unwrap_or(0), + Err(_) => 0, // No quota configured or error + }; + + stats.push(BucketStats { + name: bucket.name, + size_bytes, + objects_count, + quota_bytes, + }); + } + + stats +} + +/// Collect disk statistics from the storage layer. +async fn collect_disk_stats() -> Vec { + let Some(store) = new_object_layer_fn() else { + return Vec::new(); + }; + + let storage_info = store.storage_info().await; + + storage_info + .disks + .iter() + .map(|disk| DiskStats { + server: disk.endpoint.clone(), + drive: disk.drive_path.clone(), + total_bytes: disk.total_space, + used_bytes: disk.used_space, + free_bytes: disk.available_space, + }) + .collect() +} + +/// Collect resource statistics for the current process. +/// +/// Collects: +/// - Uptime: Calculated from process start time +/// - Memory: Process resident set size from sysinfo +/// - CPU: Process CPU usage percentage from sysinfo +#[inline] +fn collect_process_stats() -> ResourceStats { + let uptime_seconds = get_process_start().elapsed().as_secs(); + + // Use sysinfo for process metrics + let mut sys = System::new(); + let pid = Pid::from_u32(std::process::id()); + sys.refresh_processes_specifics( + ProcessesToUpdate::Some(&[pid]), + true, + ProcessRefreshKind::nothing().with_cpu().with_memory(), + ); + + if let Some(process) = sys.process(pid) { + ResourceStats { + cpu_percent: process.cpu_usage() as f64, + memory_bytes: process.memory(), + uptime_seconds, + } + } else { + // Fallback if process not found + ResourceStats { + cpu_percent: 0.0, + memory_bytes: 0, + uptime_seconds, + } + } +} + +/// Initialize the metrics collection system with periodic background tasks for cluster, bucket, node, and resource metrics. +/// +/// This function spawns background tasks that periodically collect metrics +/// and report them using the `metrics` crate. +/// +/// # Arguments +/// +/// * `token` - A cancellation token to gracefully stop the metrics collection tasks. +pub fn init_metrics_collectors(token: CancellationToken) { + // Initialize process start time + get_process_start(); + + // Helper closure to determine interval for a specific metric type + let get_interval = |env_key: &str, type_default: Duration| -> Duration { + // 1. Try specific env var + // 2. Fallback to global default env var (if set differently from hardcoded default) + // 3. Fallback to type specific default + + // Helper to check if value is valid (non-zero) + let is_valid = |v: u64| v > 0; + + if let Some(val) = get_env_opt_u64(env_key).filter(|&v| is_valid(v)) { + Duration::from_secs(val) + } else if let Some(val) = get_env_opt_u64(ENV_DEFAULT_METRICS_INTERVAL).filter(|&v| is_valid(v)) { + Duration::from_secs(val) + } else { + type_default + } + }; + + let cluster_interval = get_interval(ENV_CLUSTER_METRICS_INTERVAL, DEFAULT_CLUSTER_METRICS_INTERVAL); + let bucket_interval = get_interval(ENV_BUCKET_METRICS_INTERVAL, DEFAULT_BUCKET_METRICS_INTERVAL); + let node_interval = get_interval(ENV_NODE_METRICS_INTERVAL, DEFAULT_NODE_METRICS_INTERVAL); + let resource_interval = get_interval(ENV_RESOURCE_METRICS_INTERVAL, DEFAULT_RESOURCE_METRICS_INTERVAL); + + // Spawn task for cluster metrics + let token_clone = token.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(cluster_interval); + loop { + tokio::select! { + _ = interval.tick() => { + let stats = collect_cluster_stats().await; + let metrics = collect_cluster_metrics(&stats); + report_metrics(&metrics); + } + _ = token_clone.cancelled() => { + warn!("Metrics collection for cluster stats cancelled."); + return; + } + } + } + }); + + // Spawn task for bucket metrics + let token_clone = token.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(bucket_interval); + loop { + tokio::select! { + _ = interval.tick() => { + let stats = collect_bucket_stats().await; + let metrics = collect_bucket_metrics(&stats); + report_metrics(&metrics); + } + _ = token_clone.cancelled() => { + warn!("Metrics collection for bucket stats cancelled."); + return; + } + } + } + }); + + // Spawn task for node/disk metrics + let token_clone = token.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(node_interval); + loop { + tokio::select! { + _ = interval.tick() => { + let stats = collect_disk_stats().await; + let metrics = collect_node_metrics(&stats); + report_metrics(&metrics); + } + _ = token_clone.cancelled() => { + warn!("Metrics collection for node/disk stats cancelled."); + return; + } + } + } + }); + + // Spawn task for resource metrics + let token_clone = token.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(resource_interval); + loop { + tokio::select! { + _ = interval.tick() => { + // Resource stats collection is synchronous but fast + let stats = collect_process_stats(); + let metrics = collect_resource_metrics(&stats); + report_metrics(&metrics); + } + _ = token_clone.cancelled() => { + warn!("Metrics collection for resource stats cancelled."); + return; + } + } + } + }); +} diff --git a/crates/metrics/src/collectors/mod.rs b/crates/metrics/src/collectors/mod.rs new file mode 100644 index 000000000..f8241f9b8 --- /dev/null +++ b/crates/metrics/src/collectors/mod.rs @@ -0,0 +1,73 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Prometheus metric collectors for RustFS. +//! +//! This module provides collectors that convert RustFS data into Prometheus +//! metrics format. Each collector is responsible for a specific domain: +//! +//! - [`cluster`]: Cluster-wide capacity and object statistics +//! - [`bucket`]: Per-bucket usage and quota metrics +//! - [`node`]: Per-node disk capacity and health metrics +//! - [`resource`]: System resource metrics (CPU, memory, uptime) +//! +//! # Design Philosophy +//! +//! Collectors accept simple data structs rather than internal RustFS types. +//! This design allows HTTP handlers to populate the structs from their +//! available data sources without creating circular dependencies. +//! +//! # Example +//! +//! ``` +//! use rustfs_metrics::collectors::{ +//! collect_cluster_metrics, ClusterStats, +//! collect_bucket_metrics, BucketStats, +//! collect_node_metrics, DiskStats, +//! collect_resource_metrics, ResourceStats, +//! }; +//! use rustfs_metrics::report_metrics; +//! +//! // Collect cluster metrics +//! let cluster_stats = ClusterStats { +//! raw_capacity_bytes: 1_000_000_000, +//! used_bytes: 500_000_000, +//! ..Default::default() +//! }; +//! let mut metrics = collect_cluster_metrics(&cluster_stats); +//! +//! // Add bucket metrics +//! let bucket_stats = vec![BucketStats { +//! name: "my-bucket".to_string(), +//! size_bytes: 100_000, +//! objects_count: 50, +//! ..Default::default() +//! }]; +//! metrics.extend(collect_bucket_metrics(&bucket_stats)); +//! +//! // Report to metrics system +//! report_metrics(&metrics); +//! ``` + +mod bucket; +mod cluster; +pub(crate) mod global; +mod node; +mod resource; + +pub use bucket::{BucketStats, collect_bucket_metrics}; +pub use cluster::{ClusterStats, collect_cluster_metrics}; +pub use global::init_metrics_collectors; +pub use node::{DiskStats, collect_node_metrics}; +pub use resource::{ResourceStats, collect_resource_metrics}; diff --git a/crates/metrics/src/collectors/node.rs b/crates/metrics/src/collectors/node.rs new file mode 100644 index 000000000..ba60d19e2 --- /dev/null +++ b/crates/metrics/src/collectors/node.rs @@ -0,0 +1,196 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Per-node and per-disk metrics collector. +//! +//! Collects storage metrics for each disk/drive in the cluster, +//! including capacity, usage, and health status. + +use crate::MetricType; +use crate::format::PrometheusMetric; +use std::borrow::Cow; + +/// Statistics for a single disk/drive. +#[derive(Debug, Clone, Default)] +pub struct DiskStats { + /// Server endpoint (e.g., "node1:9000") + pub server: String, + /// Drive path (e.g., "/data/disk1") + pub drive: String, + /// Total capacity in bytes + pub total_bytes: u64, + /// Used space in bytes + pub used_bytes: u64, + /// Free space in bytes + pub free_bytes: u64, +} + +// Static metric definitions +const METRIC_TOTAL: &str = "rustfs_node_disk_total_bytes"; +const METRIC_USED: &str = "rustfs_node_disk_used_bytes"; +const METRIC_FREE: &str = "rustfs_node_disk_free_bytes"; + +const HELP_TOTAL: &str = "Total disk capacity in bytes"; +const HELP_USED: &str = "Used disk space in bytes"; +const HELP_FREE: &str = "Free disk space in bytes"; + +/// Collects per-node disk metrics from the provided disk statistics. +/// +/// # Metrics Produced +/// +/// For each disk, the following metrics are produced with `server` and `drive` labels: +/// +/// - `rustfs_node_disk_total_bytes`: Total capacity of the disk +/// - `rustfs_node_disk_used_bytes`: Used space on the disk +/// - `rustfs_node_disk_free_bytes`: Free space on the disk +/// +/// # Arguments +/// +/// * `disks` - Slice of disk statistics +/// +/// # Example +/// +/// ``` +/// use rustfs_metrics::collectors::{collect_node_metrics, DiskStats}; +/// +/// let disks = vec![ +/// DiskStats { +/// server: "node1:9000".to_string(), +/// drive: "/data/disk1".to_string(), +/// total_bytes: 1_000_000_000, +/// used_bytes: 400_000_000, +/// free_bytes: 600_000_000, +/// }, +/// ]; +/// let metrics = collect_node_metrics(&disks); +/// assert_eq!(metrics.len(), 3); +/// ``` +#[must_use] +#[inline] +pub fn collect_node_metrics(disks: &[DiskStats]) -> Vec { + if disks.is_empty() { + return Vec::new(); + } + + let mut metrics = Vec::with_capacity(disks.len() * 3); + + for disk in disks { + let server_label: Cow<'static, str> = Cow::Owned(disk.server.clone()); + let drive_label: Cow<'static, str> = Cow::Owned(disk.drive.clone()); + + metrics.push( + PrometheusMetric::new(METRIC_TOTAL, MetricType::Gauge, HELP_TOTAL, disk.total_bytes as f64) + .with_label("server", server_label.clone()) + .with_label("drive", drive_label.clone()), + ); + + metrics.push( + PrometheusMetric::new(METRIC_USED, MetricType::Gauge, HELP_USED, disk.used_bytes as f64) + .with_label("server", server_label.clone()) + .with_label("drive", drive_label.clone()), + ); + + metrics.push( + PrometheusMetric::new(METRIC_FREE, MetricType::Gauge, HELP_FREE, disk.free_bytes as f64) + .with_label("server", server_label) + .with_label("drive", drive_label), + ); + } + + metrics +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_collect_node_metrics() { + let disks = vec![ + DiskStats { + server: "node1:9000".to_string(), + drive: "/data/disk1".to_string(), + total_bytes: 1000000, + used_bytes: 400000, + free_bytes: 600000, + }, + DiskStats { + server: "node2:9000".to_string(), + drive: "/data/disk2".to_string(), + total_bytes: 2000000, + used_bytes: 800000, + free_bytes: 1200000, + }, + ]; + + let metrics = collect_node_metrics(&disks); + + // 2 disks * 3 metrics each = 6 metrics + assert_eq!(metrics.len(), 6); + + // Verify node1 disk1 total bytes + let node1_total = metrics.iter().find(|m| { + m.name == METRIC_TOTAL + && m.labels.iter().any(|(k, v)| *k == "server" && v == "node1:9000") + && m.labels.iter().any(|(k, v)| *k == "drive" && v == "/data/disk1") + }); + assert!(node1_total.is_some()); + assert_eq!(node1_total.map(|m| m.value), Some(1000000.0)); + + // Verify node2 disk2 used bytes + let node2_used = metrics.iter().find(|m| { + m.name == METRIC_USED + && m.labels.iter().any(|(k, v)| *k == "server" && v == "node2:9000") + && m.labels.iter().any(|(k, v)| *k == "drive" && v == "/data/disk2") + }); + assert!(node2_used.is_some()); + assert_eq!(node2_used.map(|m| m.value), Some(800000.0)); + } + + #[test] + fn test_collect_node_metrics_empty() { + let disks: Vec = vec![]; + let metrics = collect_node_metrics(&disks); + assert!(metrics.is_empty()); + } + + #[test] + fn test_collect_node_metrics_labels() { + let disks = vec![DiskStats { + server: "localhost:9000".to_string(), + drive: "/mnt/data".to_string(), + total_bytes: 500, + used_bytes: 200, + free_bytes: 300, + }]; + + let metrics = collect_node_metrics(&disks); + + for metric in &metrics { + assert_eq!(metric.labels.len(), 2); + assert!(metric.labels.iter().any(|(k, _)| *k == "server")); + assert!(metric.labels.iter().any(|(k, _)| *k == "drive")); + } + } + + #[test] + fn test_disk_stats_default() { + let stats = DiskStats::default(); + assert!(stats.server.is_empty()); + assert!(stats.drive.is_empty()); + assert_eq!(stats.total_bytes, 0); + assert_eq!(stats.used_bytes, 0); + assert_eq!(stats.free_bytes, 0); + } +} diff --git a/crates/metrics/src/collectors/resource.rs b/crates/metrics/src/collectors/resource.rs new file mode 100644 index 000000000..636c690e9 --- /dev/null +++ b/crates/metrics/src/collectors/resource.rs @@ -0,0 +1,170 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! System resource metrics collector. +//! +//! Collects system-level metrics for the RustFS process including +//! CPU usage, memory consumption, and process uptime. + +use crate::MetricType; +use crate::format::PrometheusMetric; + +/// Resource statistics for a RustFS process. +/// +/// This struct encapsulates the resource usage data that can be +/// collected from the operating system for the current process. +#[derive(Debug, Clone, Default)] +pub struct ResourceStats { + /// CPU usage as a percentage (0.0 to 100.0+) + pub cpu_percent: f64, + /// Resident memory usage in bytes + pub memory_bytes: u64, + /// Process uptime in seconds + pub uptime_seconds: u64, +} + +// Static metric definitions +const METRIC_CPU: &str = "rustfs_process_cpu_percent"; +const METRIC_MEMORY: &str = "rustfs_process_memory_bytes"; +const METRIC_UPTIME: &str = "rustfs_process_uptime_seconds"; + +const HELP_CPU: &str = "CPU usage of the RustFS process as a percentage"; +const HELP_MEMORY: &str = "Resident memory usage of the RustFS process in bytes"; +const HELP_UPTIME: &str = "Uptime of the RustFS process in seconds"; + +/// Number of metrics produced by this collector. +const METRIC_COUNT: usize = 3; + +/// Collects system resource metrics from the provided statistics. +/// +/// # Metrics Produced +/// +/// - `rustfs_process_cpu_percent`: CPU usage as a percentage +/// - `rustfs_process_memory_bytes`: Resident memory usage in bytes +/// - `rustfs_process_uptime_seconds`: Process uptime in seconds +/// +/// # Arguments +/// +/// * `stats` - Resource statistics for the current process +/// +/// # Example +/// +/// ``` +/// use rustfs_metrics::collectors::{collect_resource_metrics, ResourceStats}; +/// +/// let stats = ResourceStats { +/// cpu_percent: 25.5, +/// memory_bytes: 1024 * 1024 * 512, // 512 MB +/// uptime_seconds: 3600, +/// }; +/// let metrics = collect_resource_metrics(&stats); +/// assert_eq!(metrics.len(), 3); +/// ``` +#[must_use] +#[inline] +pub fn collect_resource_metrics(stats: &ResourceStats) -> Vec { + let mut metrics = Vec::with_capacity(METRIC_COUNT); + + metrics.push(PrometheusMetric::new(METRIC_CPU, MetricType::Gauge, HELP_CPU, stats.cpu_percent)); + + metrics.push(PrometheusMetric::new( + METRIC_MEMORY, + MetricType::Gauge, + HELP_MEMORY, + stats.memory_bytes as f64, + )); + + metrics.push(PrometheusMetric::new( + METRIC_UPTIME, + MetricType::Gauge, + HELP_UPTIME, + stats.uptime_seconds as f64, + )); + + metrics +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::format::report_metrics; + + #[test] + fn test_collect_resource_metrics() { + let stats = ResourceStats { + cpu_percent: 45.5, + memory_bytes: 1024 * 1024 * 256, + uptime_seconds: 7200, + }; + + let metrics = collect_resource_metrics(&stats); + report_metrics(&metrics); + + assert_eq!(metrics.len(), 3); + + // Verify CPU metric + let cpu = metrics.iter().find(|m| m.name == METRIC_CPU); + assert!(cpu.is_some()); + assert_eq!(cpu.map(|m| m.value), Some(45.5)); + + // Verify memory metric + let memory = metrics.iter().find(|m| m.name == METRIC_MEMORY); + assert!(memory.is_some()); + assert_eq!(memory.map(|m| m.value), Some((1024 * 1024 * 256) as f64)); + + // Verify uptime metric + let uptime = metrics.iter().find(|m| m.name == METRIC_UPTIME); + assert!(uptime.is_some()); + assert_eq!(uptime.map(|m| m.value), Some(7200.0)); + } + + #[test] + fn test_collect_resource_metrics_zero_values() { + let stats = ResourceStats::default(); + + let metrics = collect_resource_metrics(&stats); + report_metrics(&metrics); + + assert_eq!(metrics.len(), 3); + + for metric in &metrics { + assert_eq!(metric.value, 0.0); + assert!(metric.labels.is_empty()); + } + } + + #[test] + fn test_collect_resource_metrics_high_cpu() { + let stats = ResourceStats { + cpu_percent: 150.0, // Can exceed 100% on multi-core systems + memory_bytes: 0, + uptime_seconds: 0, + }; + + let metrics = collect_resource_metrics(&stats); + report_metrics(&metrics); + + let cpu = metrics.iter().find(|m| m.name == METRIC_CPU); + assert!(cpu.is_some()); + assert_eq!(cpu.map(|m| m.value), Some(150.0)); + } + + #[test] + fn test_resource_stats_default() { + let stats = ResourceStats::default(); + assert_eq!(stats.cpu_percent, 0.0); + assert_eq!(stats.memory_bytes, 0); + assert_eq!(stats.uptime_seconds, 0); + } +} diff --git a/crates/metrics/src/constants/mod.rs b/crates/metrics/src/constants/mod.rs new file mode 100644 index 000000000..6533997f0 --- /dev/null +++ b/crates/metrics/src/constants/mod.rs @@ -0,0 +1,42 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::time::Duration; + +/// Environment variable key for the global default metrics interval (seconds). +pub const ENV_DEFAULT_METRICS_INTERVAL: &str = "RUSTFS_METRICS_DEFAULT_INTERVAL_SEC"; +/// Default interval for metrics collection if not specified otherwise. +#[allow(dead_code)] +pub const DEFAULT_METRICS_INTERVAL: Duration = Duration::from_secs(60); + +/// Environment variable key for cluster metrics interval (seconds). +pub const ENV_CLUSTER_METRICS_INTERVAL: &str = "RUSTFS_METRICS_CLUSTER_INTERVAL_SEC"; +/// Default interval for collecting cluster-wide metrics (capacity, object counts). +pub const DEFAULT_CLUSTER_METRICS_INTERVAL: Duration = Duration::from_secs(60); + +/// Environment variable key for bucket metrics interval (seconds). +pub const ENV_BUCKET_METRICS_INTERVAL: &str = "RUSTFS_METRICS_BUCKET_INTERVAL_SEC"; +/// Default interval for collecting per-bucket metrics (usage, quotas). +/// This can be expensive if there are many buckets, so a longer interval is recommended. +pub const DEFAULT_BUCKET_METRICS_INTERVAL: Duration = Duration::from_secs(300); + +/// Environment variable key for node metrics interval (seconds). +pub const ENV_NODE_METRICS_INTERVAL: &str = "RUSTFS_METRICS_NODE_INTERVAL_SEC"; +/// Default interval for collecting node/disk metrics. +pub const DEFAULT_NODE_METRICS_INTERVAL: Duration = Duration::from_secs(60); + +/// Environment variable key for resource metrics interval (seconds). +pub const ENV_RESOURCE_METRICS_INTERVAL: &str = "RUSTFS_METRICS_RESOURCE_INTERVAL_SEC"; +/// Default interval for collecting system resource metrics (CPU, memory). +pub const DEFAULT_RESOURCE_METRICS_INTERVAL: Duration = Duration::from_secs(15); diff --git a/crates/metrics/src/format.rs b/crates/metrics/src/format.rs new file mode 100644 index 000000000..87c9eec03 --- /dev/null +++ b/crates/metrics/src/format.rs @@ -0,0 +1,128 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Prometheus text exposition format renderer. +//! +//! This module renders metrics in the standard Prometheus text format. +//! Optimized for minimal allocations and fast rendering. + +use crate::MetricType; +use metrics::{counter, describe_counter, describe_gauge, describe_histogram, gauge}; +use std::borrow::Cow; + +/// Report metrics using the `metrics` crate. +/// +/// This function iterates over the provided metrics and reports them using +/// the `metrics` crate's API. This allows integration with various metrics +/// exporters (e.g., Prometheus) that are configured globally. +pub fn report_metrics(metrics: &[PrometheusMetric]) { + for metric in metrics { + // Register metric description (help text) + // Note: In a real-world scenario, descriptions should ideally be registered once at startup. + // However, the `metrics` crate handles duplicate registrations gracefully. + match metric.metric_type { + MetricType::Counter => describe_counter!(metric.name, metric.help), + MetricType::Gauge => describe_gauge!(metric.name, metric.help), + MetricType::Histogram => describe_histogram!(metric.name, metric.help), + } + + // Convert labels to the format expected by `metrics` crate + let labels: Vec<(String, String)> = metric.labels.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect(); + + // Report the metric value + match metric.metric_type { + MetricType::Counter => { + // Use counter! macro to get a handle, then set absolute value. + // Note: `metrics` crate counters are typically monotonic and support `increment`. + // Setting an absolute value directly requires `absolute` method if supported by the backend/handle, + // or we assume the value provided is the absolute count we want to report. + // + // Since `metrics` 0.21+, `Counter` has an `absolute` method which sets the counter to a specific value. + // This is useful for mirroring an external counter. + let counter = counter!(metric.name, &labels); + counter.absolute(metric.value as u64); + } + MetricType::Gauge => { + let gauge = gauge!(metric.name, &labels); + gauge.set(metric.value); + } + MetricType::Histogram => { + let histogram = metrics::histogram!(metric.name, &labels); + histogram.record(metric.value); + } + } + } +} + +/// A single Prometheus metric with labels and value. +/// +/// This struct is optimized for performance by using `&'static str` for +/// the name and help text, which are typically compile-time constants. +/// Labels use `Cow<'static, str>` to avoid allocations when possible. +#[derive(Debug, Clone)] +pub struct PrometheusMetric { + /// The metric name (e.g., "http_requests_total"). + pub name: &'static str, + /// The type of this metric (counter, gauge, or histogram). + pub metric_type: MetricType, + /// Human-readable description shown in Prometheus UI. + pub help: &'static str, + /// Key-value label pairs for this metric instance. + /// Uses Cow to avoid allocations for static label keys. + pub labels: Vec<(&'static str, Cow<'static, str>)>, + /// The numeric value of this metric. + pub value: f64, +} + +impl PrometheusMetric { + /// Creates a new metric with the given name, type, help text, and value. + /// + /// Uses static strings to avoid heap allocations for metric metadata. + #[inline] + pub const fn new(name: &'static str, metric_type: MetricType, help: &'static str, value: f64) -> Self { + Self { + name, + metric_type, + help, + labels: Vec::new(), + value, + } + } + + /// Adds a single label with a static value to this metric. + #[inline] + #[allow(dead_code)] + pub fn with_label(mut self, key: &'static str, value: impl Into>) -> Self { + self.labels.push((key, value.into())); + self + } + + /// Adds a label with an owned string value. + /// + /// Use this when the label value is dynamically generated. + #[inline] + #[allow(dead_code)] + pub fn with_label_owned(mut self, key: &'static str, value: String) -> Self { + self.labels.push((key, Cow::Owned(value))); + self + } + + /// Sets all labels for this metric, replacing any existing labels. + #[inline] + #[allow(dead_code)] + pub fn with_labels(mut self, labels: Vec<(&'static str, Cow<'static, str>)>) -> Self { + self.labels = labels; + self + } +} diff --git a/crates/metrics/src/global.rs b/crates/metrics/src/global.rs new file mode 100644 index 000000000..c1bcce668 --- /dev/null +++ b/crates/metrics/src/global.rs @@ -0,0 +1,39 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use tokio_util::sync::CancellationToken; + +/// Initializes the global metrics system. This should be called once at the start of the application. +/// The provided `CancellationToken` will be used to gracefully shut down the metrics system when needed. +/// +/// # Arguments +/// * `token` - A `CancellationToken` that can be used to signal the metrics system to shut down gracefully. +/// +/// # Example +/// ```ignore +/// use tokio_util::sync::CancellationToken; +/// use rustfs_metrics::init_metrics_system; +/// +/// let token = CancellationToken::new(); +/// init_metrics_system(token.clone()); +/// +/// // Later, when you want to shut down the metrics system: +/// token.cancel(); +/// ``` +/// Note: This function should only be called once during the application's lifecycle. Calling it multiple times may lead to unexpected behavior. +pub fn init_metrics_system(token: CancellationToken) { + tracing::info!("init metrics system start"); + crate::collectors::init_metrics_collectors(token); + tracing::info!("init metrics system done"); +} diff --git a/crates/metrics/src/lib.rs b/crates/metrics/src/lib.rs new file mode 100644 index 000000000..529c5bc40 --- /dev/null +++ b/crates/metrics/src/lib.rs @@ -0,0 +1,23 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +pub mod collectors; +pub mod constants; +pub mod format; +mod global; +mod metrics_type; + +pub use format::report_metrics; +pub use global::init_metrics_system; +pub use metrics_type::*; diff --git a/crates/obs/src/metrics/audit.rs b/crates/metrics/src/metrics_type/audit.rs similarity index 93% rename from crates/obs/src/metrics/audit.rs rename to crates/metrics/src/metrics_type/audit.rs index 54564c6b5..2d6855842 100644 --- a/crates/obs/src/metrics/audit.rs +++ b/crates/metrics/src/metrics_type/audit.rs @@ -14,9 +14,6 @@ #![allow(dead_code)] -/// audit related metric descriptors -/// -/// This module contains the metric descriptors for the audit subsystem. use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/bucket.rs b/crates/metrics/src/metrics_type/bucket.rs similarity index 98% rename from crates/obs/src/metrics/bucket.rs rename to crates/metrics/src/metrics_type/bucket.rs index 516662e38..f007712e6 100644 --- a/crates/obs/src/metrics/bucket.rs +++ b/crates/metrics/src/metrics_type/bucket.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// bucket level s3 metric descriptor use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, new_histogram_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/bucket_replication.rs b/crates/metrics/src/metrics_type/bucket_replication.rs similarity index 99% rename from crates/obs/src/metrics/bucket_replication.rs rename to crates/metrics/src/metrics_type/bucket_replication.rs index df26b7575..2e9ec0dd2 100644 --- a/crates/obs/src/metrics/bucket_replication.rs +++ b/crates/metrics/src/metrics_type/bucket_replication.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Bucket copy metric descriptor use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/cluster_config.rs b/crates/metrics/src/metrics_type/cluster_config.rs similarity index 95% rename from crates/obs/src/metrics/cluster_config.rs rename to crates/metrics/src/metrics_type/cluster_config.rs index 143796485..2b10d3fb2 100644 --- a/crates/obs/src/metrics/cluster_config.rs +++ b/crates/metrics/src/metrics_type/cluster_config.rs @@ -14,9 +14,7 @@ #![allow(dead_code)] -/// Metric descriptors related to cluster configuration use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; - use std::sync::LazyLock; pub static CONFIG_RRS_PARITY_MD: LazyLock = LazyLock::new(|| { diff --git a/crates/obs/src/metrics/cluster_erasure_set.rs b/crates/metrics/src/metrics_type/cluster_erasure_set.rs similarity index 98% rename from crates/obs/src/metrics/cluster_erasure_set.rs rename to crates/metrics/src/metrics_type/cluster_erasure_set.rs index ecb232220..44b1ff812 100644 --- a/crates/obs/src/metrics/cluster_erasure_set.rs +++ b/crates/metrics/src/metrics_type/cluster_erasure_set.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Erasure code set related metric descriptors use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/cluster_health.rs b/crates/metrics/src/metrics_type/cluster_health.rs similarity index 96% rename from crates/obs/src/metrics/cluster_health.rs rename to crates/metrics/src/metrics_type/cluster_health.rs index b951c06e4..444bfa927 100644 --- a/crates/obs/src/metrics/cluster_health.rs +++ b/crates/metrics/src/metrics_type/cluster_health.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Cluster health-related metric descriptors use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/cluster_iam.rs b/crates/metrics/src/metrics_type/cluster_iam.rs similarity index 99% rename from crates/obs/src/metrics/cluster_iam.rs rename to crates/metrics/src/metrics_type/cluster_iam.rs index 021fd5671..9f8077259 100644 --- a/crates/obs/src/metrics/cluster_iam.rs +++ b/crates/metrics/src/metrics_type/cluster_iam.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// IAM related metric descriptors use crate::{MetricDescriptor, MetricName, new_counter_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/cluster_notification.rs b/crates/metrics/src/metrics_type/cluster_notification.rs similarity index 97% rename from crates/obs/src/metrics/cluster_notification.rs rename to crates/metrics/src/metrics_type/cluster_notification.rs index 04d46faac..58797f9a6 100644 --- a/crates/obs/src/metrics/cluster_notification.rs +++ b/crates/metrics/src/metrics_type/cluster_notification.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Notify the relevant metric descriptor use crate::{MetricDescriptor, MetricName, new_counter_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/cluster_usage.rs b/crates/metrics/src/metrics_type/cluster_usage.rs similarity index 98% rename from crates/obs/src/metrics/cluster_usage.rs rename to crates/metrics/src/metrics_type/cluster_usage.rs index 315b55ecc..1991e69c3 100644 --- a/crates/obs/src/metrics/cluster_usage.rs +++ b/crates/metrics/src/metrics_type/cluster_usage.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Descriptors of metrics related to cluster object and bucket usage use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/entry/descriptor.rs b/crates/metrics/src/metrics_type/entry/descriptor.rs similarity index 100% rename from crates/obs/src/metrics/entry/descriptor.rs rename to crates/metrics/src/metrics_type/entry/descriptor.rs diff --git a/crates/obs/src/metrics/entry/metric_name.rs b/crates/metrics/src/metrics_type/entry/metric_name.rs similarity index 100% rename from crates/obs/src/metrics/entry/metric_name.rs rename to crates/metrics/src/metrics_type/entry/metric_name.rs diff --git a/crates/obs/src/metrics/entry/metric_type.rs b/crates/metrics/src/metrics_type/entry/metric_type.rs similarity index 100% rename from crates/obs/src/metrics/entry/metric_type.rs rename to crates/metrics/src/metrics_type/entry/metric_type.rs diff --git a/crates/obs/src/metrics/entry/mod.rs b/crates/metrics/src/metrics_type/entry/mod.rs similarity index 95% rename from crates/obs/src/metrics/entry/mod.rs rename to crates/metrics/src/metrics_type/entry/mod.rs index 55ec5e5ce..9d7881e3b 100644 --- a/crates/obs/src/metrics/entry/mod.rs +++ b/crates/metrics/src/metrics_type/entry/mod.rs @@ -14,12 +14,12 @@ use crate::{MetricDescriptor, MetricName, MetricNamespace, MetricSubsystem, MetricType}; -pub(crate) mod descriptor; -pub(crate) mod metric_name; -pub(crate) mod metric_type; -pub(crate) mod namespace; +pub mod descriptor; +pub mod metric_name; +pub mod metric_type; +pub mod namespace; mod path_utils; -pub(crate) mod subsystem; +pub mod subsystem; /// Create a new counter metric descriptor pub fn new_counter_md( @@ -76,7 +76,7 @@ pub fn new_histogram_md( #[cfg(test)] mod tests { use super::*; - use crate::subsystems; + use crate::{MetricName, MetricNamespace, MetricSubsystem, subsystems}; #[test] fn test_new_histogram_md() { diff --git a/crates/obs/src/metrics/entry/namespace.rs b/crates/metrics/src/metrics_type/entry/namespace.rs similarity index 100% rename from crates/obs/src/metrics/entry/namespace.rs rename to crates/metrics/src/metrics_type/entry/namespace.rs diff --git a/crates/obs/src/metrics/entry/path_utils.rs b/crates/metrics/src/metrics_type/entry/path_utils.rs similarity index 100% rename from crates/obs/src/metrics/entry/path_utils.rs rename to crates/metrics/src/metrics_type/entry/path_utils.rs diff --git a/crates/obs/src/metrics/entry/subsystem.rs b/crates/metrics/src/metrics_type/entry/subsystem.rs similarity index 98% rename from crates/obs/src/metrics/entry/subsystem.rs rename to crates/metrics/src/metrics_type/entry/subsystem.rs index 03529800a..30f12f698 100644 --- a/crates/obs/src/metrics/entry/subsystem.rs +++ b/crates/metrics/src/metrics_type/entry/subsystem.rs @@ -204,8 +204,7 @@ pub mod subsystems { #[cfg(test)] mod tests { use super::*; - use crate::MetricType; - use crate::{MetricDescriptor, MetricName, MetricNamespace}; + use crate::{MetricDescriptor, MetricName, MetricNamespace, MetricType}; #[test] fn test_metric_subsystem_formatting() { diff --git a/crates/obs/src/metrics/ilm.rs b/crates/metrics/src/metrics_type/ilm.rs similarity index 98% rename from crates/obs/src/metrics/ilm.rs rename to crates/metrics/src/metrics_type/ilm.rs index a2db5178d..bfa4914c1 100644 --- a/crates/obs/src/metrics/ilm.rs +++ b/crates/metrics/src/metrics_type/ilm.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// ILM-related metric descriptors use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/logger_webhook.rs b/crates/metrics/src/metrics_type/logger_webhook.rs similarity index 97% rename from crates/obs/src/metrics/logger_webhook.rs rename to crates/metrics/src/metrics_type/logger_webhook.rs index 9981ea2b3..dfe41b179 100644 --- a/crates/obs/src/metrics/logger_webhook.rs +++ b/crates/metrics/src/metrics_type/logger_webhook.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// A descriptor for metrics related to webhook logs use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/mod.rs b/crates/metrics/src/metrics_type/mod.rs similarity index 61% rename from crates/obs/src/metrics/mod.rs rename to crates/metrics/src/metrics_type/mod.rs index 97c6a1583..4bc6d4c55 100644 --- a/crates/obs/src/metrics/mod.rs +++ b/crates/metrics/src/metrics_type/mod.rs @@ -12,26 +12,26 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub(crate) mod audit; -pub(crate) mod bucket; -pub(crate) mod bucket_replication; -pub(crate) mod cluster_config; -pub(crate) mod cluster_erasure_set; -pub(crate) mod cluster_health; -pub(crate) mod cluster_iam; -pub(crate) mod cluster_notification; -pub(crate) mod cluster_usage; -pub(crate) mod entry; -pub(crate) mod ilm; -pub(crate) mod logger_webhook; -pub(crate) mod replication; -pub(crate) mod request; -pub(crate) mod scanner; -pub(crate) mod system_cpu; -pub(crate) mod system_drive; -pub(crate) mod system_memory; -pub(crate) mod system_network; -pub(crate) mod system_process; +pub mod audit; +pub mod bucket; +pub mod bucket_replication; +pub mod cluster_config; +pub mod cluster_erasure_set; +pub mod cluster_health; +pub mod cluster_iam; +pub mod cluster_notification; +pub mod cluster_usage; +pub mod entry; +pub mod ilm; +pub mod logger_webhook; +pub mod replication; +pub mod request; +pub mod scanner; +pub mod system_cpu; +pub mod system_drive; +pub mod system_memory; +pub mod system_network; +pub mod system_process; pub use entry::descriptor::MetricDescriptor; pub use entry::metric_name::MetricName; diff --git a/crates/obs/src/metrics/replication.rs b/crates/metrics/src/metrics_type/replication.rs similarity index 99% rename from crates/obs/src/metrics/replication.rs rename to crates/metrics/src/metrics_type/replication.rs index 69444368b..b9ee5e4df 100644 --- a/crates/obs/src/metrics/replication.rs +++ b/crates/metrics/src/metrics_type/replication.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Metrics for replication subsystem use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/request.rs b/crates/metrics/src/metrics_type/request.rs similarity index 100% rename from crates/obs/src/metrics/request.rs rename to crates/metrics/src/metrics_type/request.rs diff --git a/crates/obs/src/metrics/scanner.rs b/crates/metrics/src/metrics_type/scanner.rs similarity index 98% rename from crates/obs/src/metrics/scanner.rs rename to crates/metrics/src/metrics_type/scanner.rs index a8ab6b591..47ab02264 100644 --- a/crates/obs/src/metrics/scanner.rs +++ b/crates/metrics/src/metrics_type/scanner.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Scanner-related metric descriptors use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/system_cpu.rs b/crates/metrics/src/metrics_type/system_cpu.rs similarity index 100% rename from crates/obs/src/metrics/system_cpu.rs rename to crates/metrics/src/metrics_type/system_cpu.rs diff --git a/crates/obs/src/metrics/system_drive.rs b/crates/metrics/src/metrics_type/system_drive.rs similarity index 99% rename from crates/obs/src/metrics/system_drive.rs rename to crates/metrics/src/metrics_type/system_drive.rs index 8d406a2e4..a9752946b 100644 --- a/crates/obs/src/metrics/system_drive.rs +++ b/crates/metrics/src/metrics_type/system_drive.rs @@ -14,7 +14,6 @@ #![allow(dead_code)] -/// Drive-related metric descriptors use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/system_memory.rs b/crates/metrics/src/metrics_type/system_memory.rs similarity index 89% rename from crates/obs/src/metrics/system_memory.rs rename to crates/metrics/src/metrics_type/system_memory.rs index f09352900..de45de6fd 100644 --- a/crates/obs/src/metrics/system_memory.rs +++ b/crates/metrics/src/metrics_type/system_memory.rs @@ -14,12 +14,6 @@ #![allow(dead_code)] -/// Memory-related metric descriptors -/// -/// This module provides a set of metric descriptors for system memory statistics. -/// These descriptors are initialized lazily using `std::sync::LazyLock` to ensure -/// they are only created when actually needed, improving performance and reducing -/// startup overhead. use crate::{MetricDescriptor, MetricName, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/system_network.rs b/crates/metrics/src/metrics_type/system_network.rs similarity index 89% rename from crates/obs/src/metrics/system_network.rs rename to crates/metrics/src/metrics_type/system_network.rs index 1b6888f05..6ac2bc75d 100644 --- a/crates/obs/src/metrics/system_network.rs +++ b/crates/metrics/src/metrics_type/system_network.rs @@ -14,12 +14,6 @@ #![allow(dead_code)] -/// Network-related metric descriptors -/// -/// These metrics capture internode network communication statistics including: -/// - Error counts for connection and general internode calls -/// - Network dial performance metrics -/// - Data transfer volume in both directions use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/metrics/system_process.rs b/crates/metrics/src/metrics_type/system_process.rs similarity index 95% rename from crates/obs/src/metrics/system_process.rs rename to crates/metrics/src/metrics_type/system_process.rs index be327e986..1bb662b4b 100644 --- a/crates/obs/src/metrics/system_process.rs +++ b/crates/metrics/src/metrics_type/system_process.rs @@ -14,11 +14,6 @@ #![allow(dead_code)] -/// Process related metric descriptors -/// -/// This module defines various system process metrics used for monitoring -/// the RustFS process performance, resource usage, and system integration. -/// Metrics are implemented using std::sync::LazyLock for thread-safe lazy initialization. use crate::{MetricDescriptor, MetricName, new_counter_md, new_gauge_md, subsystems}; use std::sync::LazyLock; diff --git a/crates/obs/src/lib.rs b/crates/obs/src/lib.rs index 3d749d785..a6f5c6713 100644 --- a/crates/obs/src/lib.rs +++ b/crates/obs/src/lib.rs @@ -56,7 +56,6 @@ mod config; mod error; mod global; -mod metrics; mod recorder; mod system; mod telemetry; @@ -64,7 +63,6 @@ mod telemetry; pub use config::*; pub use error::*; pub use global::*; -pub use metrics::*; pub use recorder::*; pub use system::SystemObserver; pub use telemetry::OtelGuard; diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index aec1ff2b0..cf7035b90 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -53,6 +53,7 @@ rustfs-iam = { workspace = true } rustfs-kms = { workspace = true } rustfs-lock.workspace = true rustfs-madmin = { workspace = true } +rustfs-metrics = { workspace = true } rustfs-notify = { workspace = true } rustfs-obs = { workspace = true } rustfs-policy = { workspace = true } diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index 2eaa84501..ca4ea198b 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -62,6 +62,7 @@ use rustfs_heal::{ create_ahm_services_cancel_token, heal::storage::ECStoreHealStorage, init_heal_manager, shutdown_ahm_services, }; use rustfs_iam::init_iam_sys; +use rustfs_metrics::init_metrics_system; use rustfs_obs::{init_obs, set_global_guard}; use rustfs_scanner::init_data_scanner; use rustfs_utils::net::parse_and_resolve_address; @@ -402,6 +403,11 @@ async fn run(opt: config::Opt) -> Result<()> { init_update_check(); + if rustfs_obs::observability_metric_enabled() { + // Initialize metrics system + init_metrics_system(ctx.clone()); + } + println!( "RustFS server version: {} started successfully at {}, current time: {}", version::get_version(),