mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
perf(io-metrics): add a global switch to gate general metric emission (#4733)
The io-metrics crate already had per-stage GET/PUT switches, but the request headers/summaries and ~40 general recorders (I/O scheduler, bytes-pool, zero-copy, bandwidth, system-resource, error/timeout/retry) emitted unconditionally, paying their label allocations and arithmetic even when no metric recorder is installed. Add `METRICS_ENABLED: AtomicBool` (default false) with `set_metrics_enabled()` / `metrics_enabled()`, isomorphic to the existing stage switches, and gate: - GET headers/summaries via `get_stage_metrics_enabled()`: record_get_object_request_start / _request_started / _request_result / _timeout / _completion / _total_duration_with_path, plus legacy record_get_object. - PUT headers via `put_stage_metrics_enabled()`: record_put_object_request_start / _request_result, plus legacy record_put_object. - 41 general recorders via `metrics_enabled()`. Startup enables all three switches together under `observability_metric_enabled()` (startup_observability.rs), with a passthrough in startup_runtime_sources.rs. The global recorder is itself only installed when observability metric export is on, so gating off is pure cost savings when export is disabled and a no-op when it is enabled. Deliberately NOT gated (would break function, not just observation): the EC encode in-flight and GET buffered-bytes accounting guards (add/remove_ec_encode_inflight_bytes, track_get_object_buffered_bytes) whose values are read back by current_*; and the stateful sibling modules (lock/deadlock/backpressure/autotuner/adaptive_ttl/collector/internode). The console realtime metrics endpoint (/admin/v3/metrics) samples system state and internode metrics directly, so it does not depend on the gated recorders. Adds a toggle test exercising both the enabled and disabled paths, serialized on the existing METRICS_FLAG_LOCK to stay robust under nextest's shared-process model. Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -59,6 +59,20 @@ use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
static PUT_STAGE_METRICS_ENABLED: AtomicBool = AtomicBool::new(false);
|
||||
static GET_STAGE_METRICS_ENABLED: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
/// Global switch for all remaining (non GET/PUT-stage) metric emission in this
|
||||
/// crate's free `record_*` functions — I/O scheduler, bytes-pool, zero-copy,
|
||||
/// bandwidth, system-resource, and error/timeout/retry counters.
|
||||
///
|
||||
/// When `false`, those recorders become no-ops so callers skip the label
|
||||
/// allocations and arithmetic they would otherwise perform before the metric
|
||||
/// macro. Set to `true` during startup when OTEL metric export is enabled,
|
||||
/// alongside the GET/PUT stage switches.
|
||||
///
|
||||
/// This switch intentionally does NOT gate functions that maintain functional
|
||||
/// internal state read back by the system (EC encode in-flight accounting,
|
||||
/// GET whole-object buffered-bytes tracking); those must always run.
|
||||
static METRICS_ENABLED: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
/// Enable or disable detailed per-stage PUT metrics.
|
||||
///
|
||||
/// Called once during startup, typically gated by `rustfs_obs::observability_metric_enabled()`.
|
||||
@@ -70,6 +84,13 @@ pub fn set_get_stage_metrics_enabled(enabled: bool) {
|
||||
GET_STAGE_METRICS_ENABLED.store(enabled, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Enable or disable general (non GET/PUT-stage) metric emission.
|
||||
///
|
||||
/// Called once during startup, typically gated by `rustfs_obs::observability_metric_enabled()`.
|
||||
pub fn set_metrics_enabled(enabled: bool) {
|
||||
METRICS_ENABLED.store(enabled, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Returns `true` if detailed per-stage PUT metrics are enabled.
|
||||
///
|
||||
/// Callers should check this before calling `Instant::now()` for stage timing
|
||||
@@ -84,6 +105,12 @@ pub fn get_stage_metrics_enabled() -> bool {
|
||||
GET_STAGE_METRICS_ENABLED.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Returns `true` if general (non GET/PUT-stage) metric emission is enabled.
|
||||
#[inline(always)]
|
||||
pub fn metrics_enabled() -> bool {
|
||||
METRICS_ENABLED.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
// Public modules
|
||||
pub mod adaptive_ttl;
|
||||
pub mod autotuner;
|
||||
@@ -268,6 +295,9 @@ impl Drop for MemoryGaugeGuard {
|
||||
/// Record GetObject request start.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_request_start(concurrent_requests: usize) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_get_object_requests_total").increment(1);
|
||||
gauge!("rustfs_io_get_object_concurrent_requests").set(concurrent_requests as f64);
|
||||
}
|
||||
@@ -275,12 +305,18 @@ pub fn record_get_object_request_start(concurrent_requests: usize) {
|
||||
/// Record GetObject request start without concurrency context.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_request_started() {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_get_object_requests_total").increment(1);
|
||||
}
|
||||
|
||||
/// Record GetObject request result.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_request_result(status: &str, duration_secs: f64) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_get_object_request_results_total", "status" => status.to_string()).increment(1);
|
||||
histogram!("rustfs_io_get_object_request_duration_seconds", "status" => status.to_string()).record(duration_secs);
|
||||
}
|
||||
@@ -288,6 +324,9 @@ pub fn record_get_object_request_result(status: &str, duration_secs: f64) {
|
||||
/// Record PutObject request start.
|
||||
#[inline(always)]
|
||||
pub fn record_put_object_request_start(concurrent_requests: usize) {
|
||||
if !put_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_put_object_requests_total").increment(1);
|
||||
gauge!("rustfs_io_put_object_concurrent_requests").set(concurrent_requests as f64);
|
||||
}
|
||||
@@ -295,6 +334,9 @@ pub fn record_put_object_request_start(concurrent_requests: usize) {
|
||||
/// Record PutObject request result.
|
||||
#[inline(always)]
|
||||
pub fn record_put_object_request_result(status: &str, duration_secs: f64) {
|
||||
if !put_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_put_object_request_results_total", "status" => status.to_string()).increment(1);
|
||||
histogram!("rustfs_io_put_object_request_duration_seconds", "status" => status.to_string()).record(duration_secs);
|
||||
}
|
||||
@@ -302,6 +344,9 @@ pub fn record_put_object_request_result(status: &str, duration_secs: f64) {
|
||||
/// Record GetObject timeout for a specific stage.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_timeout(stage: Option<&str>, elapsed_secs: Option<f64>) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
match stage {
|
||||
Some(stage) => counter!("rustfs_io_get_object_timeout_total", "stage" => stage.to_string()).increment(1),
|
||||
None => counter!("rustfs_io_get_object_timeout_total").increment(1),
|
||||
@@ -315,6 +360,9 @@ pub fn record_get_object_timeout(stage: Option<&str>, elapsed_secs: Option<f64>)
|
||||
/// Record GetObject completion.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_completion(total_duration_secs: f64, response_size_bytes: i64, buffer_size_bytes: usize) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_get_object_completed_total").increment(1);
|
||||
histogram!("rustfs_io_get_object_total_duration_seconds").record(total_duration_secs);
|
||||
histogram!("rustfs_io_get_object_response_size_bytes").record(response_size_bytes as f64);
|
||||
@@ -474,12 +522,18 @@ pub fn record_get_object_memory_body_stream_poll(source: &'static str, outcome:
|
||||
/// Record I/O queue congestion observation.
|
||||
#[inline(always)]
|
||||
pub fn record_io_queue_congestion() {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_queue_congestion_total").increment(1);
|
||||
}
|
||||
|
||||
/// Record I/O priority assignment.
|
||||
#[inline(always)]
|
||||
pub fn record_io_priority_assignment(priority: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_priority_assigned_total", "priority" => priority.to_string()).increment(1);
|
||||
}
|
||||
|
||||
@@ -1402,6 +1456,9 @@ pub fn record_get_object_metadata_phase_duration_with_early_stop(duration_secs:
|
||||
/// Record GET object total duration with reader path label.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object_total_duration_with_path(duration_secs: f64, reader_path: &'static str) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
histogram!(
|
||||
"rustfs_io_get_object_total_duration_seconds_with_path",
|
||||
"reader_path" => reader_path
|
||||
@@ -1454,6 +1511,9 @@ pub fn record_get_object_pipeline_failure_for_path(path: &'static str, stage: &'
|
||||
/// * `duration_ms` - Time taken for the read operation in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_zero_copy_read(size_bytes: usize, duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_reads_total").increment(1);
|
||||
histogram!("rustfs_zero_copy_read_size_bytes").record(size_bytes as f64);
|
||||
histogram!("rustfs_zero_copy_read_duration_ms").record(duration_ms);
|
||||
@@ -1471,6 +1531,9 @@ pub fn record_zero_copy_read(size_bytes: usize, duration_ms: f64) {
|
||||
/// * `bytes_saved` - Number of bytes that would have been copied without zero-copy
|
||||
#[inline(always)]
|
||||
pub fn record_memory_copy_saved(bytes_saved: usize) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_memory_saved_bytes_total").increment(bytes_saved as u64);
|
||||
}
|
||||
|
||||
@@ -1484,6 +1547,9 @@ pub fn record_memory_copy_saved(bytes_saved: usize) {
|
||||
/// * `reason` - Reason for the fallback (e.g., "mmap_unavailable", "file_too_large")
|
||||
#[inline(always)]
|
||||
pub fn record_zero_copy_fallback(reason: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_fallback_total", "reason" => reason.to_string()).increment(1);
|
||||
counter!(mmap_copy::FALLBACK_TOTAL, "reason" => reason.to_string()).increment(1);
|
||||
}
|
||||
@@ -1501,6 +1567,9 @@ pub fn record_zero_copy_fallback(reason: &str) {
|
||||
/// * `from_pool` - Whether buffer was reused from pool
|
||||
#[inline(always)]
|
||||
pub fn record_bytes_pool_acquire(tier: &str, size: usize, from_pool: bool) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_bytes_pool_acquisitions_total", "tier" => tier.to_string()).increment(1);
|
||||
gauge!("rustfs_bytes_pool_size_bytes", "tier" => tier.to_string()).set(size as f64);
|
||||
|
||||
@@ -1518,6 +1587,9 @@ pub fn record_bytes_pool_acquire(tier: &str, size: usize, from_pool: bool) {
|
||||
/// * `tier` - Pool tier ("small", "medium", "large", "xlarge")
|
||||
#[inline(always)]
|
||||
pub fn record_bytes_pool_return(tier: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_bytes_pool_returns_total", "tier" => tier.to_string()).increment(1);
|
||||
}
|
||||
|
||||
@@ -1529,6 +1601,9 @@ pub fn record_bytes_pool_return(tier: &str) {
|
||||
/// * `bytes` - Currently allocated bytes
|
||||
#[inline(always)]
|
||||
pub fn record_bytes_pool_allocated(tier: &str, bytes: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_bytes_pool_allocated_bytes", "tier" => tier.to_string()).set(bytes as f64);
|
||||
}
|
||||
|
||||
@@ -1540,6 +1615,9 @@ pub fn record_bytes_pool_allocated(tier: &str, bytes: u64) {
|
||||
/// * `hit_rate` - Hit rate (0.0 - 1.0)
|
||||
#[inline(always)]
|
||||
pub fn record_bytes_pool_hit_rate(tier: &str, hit_rate: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_bytes_pool_hit_rate", "tier" => tier.to_string()).set(hit_rate * 100.0);
|
||||
}
|
||||
|
||||
@@ -1554,6 +1632,9 @@ pub fn record_bytes_pool_hit_rate(tier: &str, hit_rate: f64) {
|
||||
/// * `outcome` - Acquisition outcome ("hit" or "miss")
|
||||
#[inline(always)]
|
||||
pub fn record_bytespool_acquisition(tier: &'static str, outcome: &'static str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_bytespool_acquisition_total", "tier" => tier, "outcome" => outcome).increment(1);
|
||||
}
|
||||
|
||||
@@ -1568,6 +1649,9 @@ pub fn record_bytespool_acquisition(tier: &'static str, outcome: &'static str) {
|
||||
/// * `outcome` - Return outcome ("recycled" or "dropped")
|
||||
#[inline(always)]
|
||||
pub fn record_bytespool_return(tier: &'static str, outcome: &'static str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_bytespool_return_total", "tier" => tier, "outcome" => outcome).increment(1);
|
||||
}
|
||||
|
||||
@@ -1579,6 +1663,9 @@ pub fn record_bytespool_return(tier: &'static str, outcome: &'static str) {
|
||||
/// * `duration_ms` - Time taken for the write operation in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_zero_copy_write(size_bytes: usize, duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_write_total").increment(1);
|
||||
histogram!("rustfs_zero_copy_write_size_bytes").record(size_bytes as f64);
|
||||
histogram!("rustfs_zero_copy_write_duration_ms").record(duration_ms);
|
||||
@@ -1598,6 +1685,9 @@ pub fn record_zero_copy_write(size_bytes: usize, duration_ms: f64) {
|
||||
/// * `reason` - Reason for the fallback
|
||||
#[inline(always)]
|
||||
pub fn record_zero_copy_write_fallback(reason: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_write_fallback_total", "reason" => reason.to_string()).increment(1);
|
||||
counter!(buffered_write::FALLBACK_TOTAL, "reason" => reason.to_string()).increment(1);
|
||||
}
|
||||
@@ -1609,6 +1699,9 @@ pub fn record_zero_copy_write_fallback(reason: &str) {
|
||||
/// * `size_bytes` - Number of bytes saved from zero-copy
|
||||
#[inline(always)]
|
||||
pub fn record_bytes_saved(size_bytes: usize) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_zero_copy_bytes_saved_total").increment(size_bytes as u64);
|
||||
}
|
||||
|
||||
@@ -1627,6 +1720,9 @@ pub fn record_bytes_saved(size_bytes: usize) {
|
||||
/// interpreted as the definitive source of truth for data-plane copy mode.
|
||||
#[inline(always)]
|
||||
pub fn record_get_object(duration_ms: f64, size_bytes: i64) {
|
||||
if !get_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_s3_get_object_total").increment(1);
|
||||
histogram!("rustfs_s3_get_object_duration_ms").record(duration_ms);
|
||||
|
||||
@@ -1644,6 +1740,9 @@ pub fn record_get_object(duration_ms: f64, size_bytes: i64) {
|
||||
/// * `zero_copy_eligible` - Whether the request was eligible for a zero-copy path
|
||||
#[inline(always)]
|
||||
pub fn record_put_object(duration_ms: f64, size_bytes: i64, zero_copy_eligible: bool) {
|
||||
if !put_stage_metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_s3_put_object_total").increment(1);
|
||||
histogram!("rustfs_s3_put_object_duration_ms").record(duration_ms);
|
||||
|
||||
@@ -1745,6 +1844,9 @@ pub fn record_put_object_stage_duration(stage: &'static str, duration_ms: f64) {
|
||||
/// operations that are NOT part of the PUT object hot path.
|
||||
#[inline(always)]
|
||||
pub fn record_stage_duration(stage: &'static str, duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
histogram!("rustfs_internal_stage_duration_ms", "stage" => stage).record(duration_ms);
|
||||
}
|
||||
|
||||
@@ -1757,6 +1859,9 @@ pub fn record_stage_duration(stage: &'static str, duration_ms: f64) {
|
||||
/// * `is_truncated` - Whether the response was truncated
|
||||
#[inline(always)]
|
||||
pub fn record_list_objects(duration_ms: f64, objects_count: u64, is_truncated: bool) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_s3_list_objects_total").increment(1);
|
||||
histogram!("rustfs_s3_list_objects_duration_ms").record(duration_ms);
|
||||
histogram!("rustfs_s3_list_objects_count").record(objects_count as f64);
|
||||
@@ -1774,6 +1879,9 @@ pub fn record_list_objects(duration_ms: f64, objects_count: u64, is_truncated: b
|
||||
/// * `version_deleted` - Whether a specific version was deleted
|
||||
#[inline(always)]
|
||||
pub fn record_delete_object(duration_ms: f64, version_deleted: bool) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_s3_delete_object_total").increment(1);
|
||||
histogram!("rustfs_s3_delete_object_duration_ms").record(duration_ms);
|
||||
|
||||
@@ -1796,6 +1904,9 @@ pub fn record_delete_object(duration_ms: f64, version_deleted: bool) {
|
||||
/// * `concurrent_requests` - Number of concurrent requests
|
||||
#[inline(always)]
|
||||
pub fn record_io_strategy(storage_media: &str, access_pattern: &str, buffer_size: usize, concurrent_requests: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_strategy_total",
|
||||
"storage_media" => storage_media.to_string(),
|
||||
"access_pattern" => access_pattern.to_string(),
|
||||
@@ -1817,6 +1928,9 @@ pub fn record_io_strategy(storage_media: &str, access_pattern: &str, buffer_size
|
||||
/// * `duration_ms` - Time spent waiting for disk permit
|
||||
#[inline(always)]
|
||||
pub fn record_permit_wait(duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
histogram!("rustfs_io_permit_wait_duration_ms").record(duration_ms);
|
||||
}
|
||||
|
||||
@@ -1828,6 +1942,9 @@ pub fn record_permit_wait(duration_ms: f64) {
|
||||
/// * `concurrent_requests` - Number of concurrent requests
|
||||
#[inline(always)]
|
||||
pub fn record_io_load_level(load_level: &str, concurrent_requests: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_load_level",
|
||||
"level" => load_level.to_string(),
|
||||
)
|
||||
@@ -1845,6 +1962,9 @@ pub fn record_io_load_level(load_level: &str, concurrent_requests: u64) {
|
||||
/// * `entries` - Number of entries in the cache
|
||||
#[inline(always)]
|
||||
pub fn record_cache_size(tier: &str, size_bytes: usize, entries: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_cache_size_bytes",
|
||||
"tier" => tier.to_string(),
|
||||
)
|
||||
@@ -1868,6 +1988,9 @@ pub fn record_cache_size(tier: &str, size_bytes: usize, entries: u64) {
|
||||
/// * `tier` - Bandwidth tier ("low", "medium", "high", "unknown")
|
||||
#[inline(always)]
|
||||
pub fn record_bandwidth(bytes_per_second: u64, tier: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
let tier_label = if tier.is_empty() { "unknown" } else { tier };
|
||||
gauge!("rustfs_bandwidth_current_bps", "tier" => "all").set(bytes_per_second as f64);
|
||||
gauge!("rustfs_bandwidth_current_bps", "tier" => tier_label.to_string()).set(bytes_per_second as f64);
|
||||
@@ -1883,6 +2006,9 @@ pub fn record_bandwidth(bytes_per_second: u64, tier: &str) {
|
||||
/// * `duration_ms` - Duration of the transfer in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_data_transfer(bytes: u64, duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_io_transfer_bytes_total").increment(bytes);
|
||||
histogram!("rustfs_io_transfer_duration_ms").record(duration_ms);
|
||||
|
||||
@@ -1904,6 +2030,9 @@ pub fn record_data_transfer(bytes: u64, duration_ms: f64) {
|
||||
/// * `total_bytes` - Total memory in bytes
|
||||
#[inline(always)]
|
||||
pub fn record_memory_usage(used_bytes: u64, total_bytes: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_memory_used_bytes").set(used_bytes as f64);
|
||||
gauge!("rustfs_memory_total_bytes").set(total_bytes as f64);
|
||||
|
||||
@@ -1916,6 +2045,9 @@ pub fn record_memory_usage(used_bytes: u64, total_bytes: u64) {
|
||||
/// Record process-level memory split metrics.
|
||||
#[inline(always)]
|
||||
pub fn record_process_memory_split(resident_bytes: u64, virtual_bytes: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_memory_process_resident_bytes").set(resident_bytes as f64);
|
||||
gauge!("rustfs_memory_process_virtual_bytes").set(virtual_bytes as f64);
|
||||
}
|
||||
@@ -1930,6 +2062,9 @@ pub fn record_cgroup_memory_split(
|
||||
active_file_bytes: Option<u64>,
|
||||
inactive_file_bytes: Option<u64>,
|
||||
) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
if let Some(current_bytes) = current_bytes {
|
||||
gauge!("rustfs_memory_cgroup_current_bytes").set(current_bytes as f64);
|
||||
}
|
||||
@@ -1965,6 +2100,9 @@ pub struct AllocatorMemoryObservation {
|
||||
/// Record allocator-level memory attribution when supported by the allocator.
|
||||
#[inline(always)]
|
||||
pub fn record_allocator_memory_observation(backend: &'static str, observation: AllocatorMemoryObservation) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
if let Some(reserved_bytes) = observation.reserved_bytes {
|
||||
gauge!("rustfs_memory_allocator_reserved_bytes", "backend" => backend).set(reserved_bytes as f64);
|
||||
}
|
||||
@@ -2039,6 +2177,9 @@ pub fn current_get_object_buffered_bytes() -> u64 {
|
||||
/// * `percent` - CPU usage percentage (0.0 - 100.0)
|
||||
#[inline(always)]
|
||||
pub fn record_cpu_usage(percent: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_cpu_usage_percent").set(percent);
|
||||
}
|
||||
|
||||
@@ -2052,6 +2193,9 @@ pub fn record_cpu_usage(percent: f64) {
|
||||
/// * `write_ops` - Number of write operations
|
||||
#[inline(always)]
|
||||
pub fn record_disk_io(read_bytes: u64, write_bytes: u64, read_ops: u64, write_ops: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_disk_read_bytes_total").increment(read_bytes);
|
||||
counter!("rustfs_disk_write_bytes_total").increment(write_bytes);
|
||||
counter!("rustfs_disk_read_ops_total").increment(read_ops);
|
||||
@@ -2070,6 +2214,9 @@ pub fn record_disk_io(read_bytes: u64, write_bytes: u64, read_ops: u64, write_op
|
||||
/// * `error_type` - Error type (e.g., "timeout", "disk_error", "network")
|
||||
#[inline(always)]
|
||||
pub fn record_error(operation: &str, error_type: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_errors_total",
|
||||
"operation" => operation.to_string(),
|
||||
"type" => error_type.to_string(),
|
||||
@@ -2085,6 +2232,9 @@ pub fn record_error(operation: &str, error_type: &str) {
|
||||
/// * `duration_ms` - Duration before timeout
|
||||
#[inline(always)]
|
||||
pub fn record_timeout(operation: &str, duration_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_timeouts_total",
|
||||
"operation" => operation.to_string(),
|
||||
)
|
||||
@@ -2104,6 +2254,9 @@ pub fn record_timeout(operation: &str, duration_ms: f64) {
|
||||
/// * `attempt_number` - Attempt number (1-based)
|
||||
#[inline(always)]
|
||||
pub fn record_retry(operation: &str, attempt_number: u32) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!("rustfs_retries_total",
|
||||
"operation" => operation.to_string(),
|
||||
)
|
||||
@@ -2126,6 +2279,9 @@ pub fn record_retry(operation: &str, attempt_number: u32) {
|
||||
/// * `latency_ms` - I/O latency in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_io_latency(latency_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
histogram!("rustfs_io_latency_ms").record(latency_ms);
|
||||
}
|
||||
|
||||
@@ -2136,6 +2292,9 @@ pub fn record_io_latency(latency_ms: f64) {
|
||||
/// * `latency_ms` - P95 I/O latency in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_io_latency_p95(latency_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_io_latency_p95_ms").set(latency_ms);
|
||||
}
|
||||
|
||||
@@ -2146,6 +2305,9 @@ pub fn record_io_latency_p95(latency_ms: f64) {
|
||||
/// * `latency_ms` - P99 I/O latency in milliseconds
|
||||
#[inline(always)]
|
||||
pub fn record_io_latency_p99(latency_ms: f64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!("rustfs_io_latency_p99_ms").set(latency_ms);
|
||||
}
|
||||
|
||||
@@ -2500,10 +2662,37 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn test_record_stage_duration_generic() {
|
||||
// Generic stage duration should always record (no gating flag)
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
// Generic stage duration is gated by the common `metrics_enabled()` switch.
|
||||
set_metrics_enabled(true);
|
||||
record_stage_duration("metacache_walk_dir_primary", 15.0);
|
||||
record_stage_duration("store_list_objects_walk_internal", 8.5);
|
||||
record_stage_duration("lifecycle_free_version_recovery_failed", 120.0);
|
||||
assert!(metrics_enabled());
|
||||
set_metrics_enabled(false);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_metrics_enabled_toggle_gates_common_recorders() {
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
|
||||
// Disabled: the common recorders must be no-ops (no panic, no recording).
|
||||
set_metrics_enabled(false);
|
||||
assert!(!metrics_enabled());
|
||||
record_stage_duration("metacache_walk_dir_primary", 15.0);
|
||||
record_list_objects(50.0, 100, false);
|
||||
record_error("get_object", "timeout");
|
||||
record_cpu_usage(25.5);
|
||||
|
||||
// Enabled: the same recorders run their emission bodies without panicking.
|
||||
set_metrics_enabled(true);
|
||||
assert!(metrics_enabled());
|
||||
record_stage_duration("metacache_walk_dir_primary", 15.0);
|
||||
record_list_objects(50.0, 100, false);
|
||||
record_error("get_object", "timeout");
|
||||
record_cpu_usage(25.5);
|
||||
|
||||
set_metrics_enabled(false);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -2655,6 +2844,9 @@ pub use metric_names::{aligned_pread, buffered_write, mmap_copy, zero_copy};
|
||||
/// including the operation type and size.
|
||||
#[inline(always)]
|
||||
pub fn record_zero_copy_buffer_operation(operation: &str, size: usize) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!(
|
||||
zero_copy::BUFFER_OPERATIONS_TOTAL,
|
||||
"operation" => operation.to_string()
|
||||
@@ -2674,6 +2866,9 @@ pub fn record_zero_copy_buffer_operation(operation: &str, size: usize) {
|
||||
/// which should be minimized in zero-copy paths.
|
||||
#[inline(always)]
|
||||
pub fn record_memory_copy(count: u32, size: usize) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!(zero_copy::MEMORY_COPY_TOTAL).increment(count as u64);
|
||||
|
||||
counter!(zero_copy::MEMORY_COPY_BYTES_TOTAL).increment(size as u64);
|
||||
@@ -2687,6 +2882,9 @@ pub fn record_memory_copy(count: u32, size: usize) {
|
||||
/// for zero-copy data sharing.
|
||||
#[inline(always)]
|
||||
pub fn record_shared_ref_operation(operation: &str) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!(
|
||||
zero_copy::SHARED_REF_OPERATIONS_TOTAL,
|
||||
"operation" => operation.to_string()
|
||||
@@ -2700,6 +2898,9 @@ pub fn record_shared_ref_operation(operation: &str) {
|
||||
/// adjustments.
|
||||
#[inline(always)]
|
||||
pub fn record_bufreader_optimization(layers_eliminated: u32, buffer_size: usize) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
counter!(zero_copy::BUFREADER_LAYERS_ELIMINATED_TOTAL).increment(layers_eliminated as u64);
|
||||
|
||||
histogram!(zero_copy::BUFREADER_BUFFER_SIZE_BYTES).record(buffer_size as f64);
|
||||
@@ -2711,6 +2912,9 @@ pub fn record_bufreader_optimization(layers_eliminated: u32, buffer_size: usize)
|
||||
/// status.
|
||||
#[inline(always)]
|
||||
pub fn record_direct_io_operation(operation: &str, size: usize, success: bool) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
let status = if success { "success" } else { "fallback" };
|
||||
|
||||
counter!(
|
||||
@@ -2747,6 +2951,9 @@ pub fn record_direct_io_operation(operation: &str, size: usize, success: bool) {
|
||||
/// This function updates gauge metrics for overall zero-copy performance.
|
||||
#[inline(always)]
|
||||
pub fn update_zero_copy_performance_metrics(copy_count: u32, throughput_mbps: f64, memory_saved: u64) {
|
||||
if !metrics_enabled() {
|
||||
return;
|
||||
}
|
||||
gauge!(zero_copy::AVG_COPY_COUNT).set(copy_count as f64);
|
||||
|
||||
gauge!(zero_copy::THROUGHPUT_MBPS).set(throughput_mbps);
|
||||
|
||||
@@ -29,6 +29,7 @@ pub(crate) async fn init_observability_runtime(store: Arc<ECStore>, ctx: Cancell
|
||||
init_compression_total_memory_from_backend(store).await;
|
||||
startup_runtime_sources::set_put_stage_metrics_enabled(true);
|
||||
startup_runtime_sources::set_get_stage_metrics_enabled(true);
|
||||
startup_runtime_sources::set_metrics_enabled(true);
|
||||
startup_runtime_sources::init_metrics_runtime(ctx.clone());
|
||||
crate::memory_observability::init_memory_observability(ctx.clone());
|
||||
init_auto_tuner(ctx).await;
|
||||
|
||||
@@ -104,6 +104,10 @@ pub(crate) fn set_get_stage_metrics_enabled(enabled: bool) {
|
||||
rustfs_io_metrics::set_get_stage_metrics_enabled(enabled);
|
||||
}
|
||||
|
||||
pub(crate) fn set_metrics_enabled(enabled: bool) {
|
||||
rustfs_io_metrics::set_metrics_enabled(enabled);
|
||||
}
|
||||
|
||||
pub(crate) fn init_tls_metrics() {
|
||||
rustfs_tls_runtime::init_tls_metrics();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user