mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-18 10:43:15 +00:00
chore(storage): drop dead io-schedule metrics and helpers (#6199)
This commit is contained in:
@@ -107,7 +107,7 @@ inventory. Generic function-local names such as `CACHE`, `LOCK`, `INIT`, and
|
||||
| `AUTH_FS` | `rustfs/src/storage/access.rs` | Cache or constant / owner-local cache | Authorization tag-condition lookup keeps its filesystem helper private to the access owner. |
|
||||
| `LOCK_STATS` | `rustfs/src/storage/lock_optimizer.rs` | Process-global owner-local metrics | Lock optimization statistics stay private behind lock optimizer helper APIs. |
|
||||
| `DEADLOCK_DETECTOR` | `rustfs/src/storage/deadlock_detector.rs` | Process-global owner-local state | Deadlock detector lifecycle state stays private to the storage deadlock detector owner. |
|
||||
| `CONCURRENCY_MANAGER`, `ACTIVE_GET_REQUESTS`, `ACTIVE_PUT_REQUESTS`, `IO_PRIORITY_METRICS` | `rustfs/src/storage/concurrency/*` | Process-global owner-local scheduler state | Storage concurrency manager, counters, and metrics remain inside the storage concurrency owner boundary. |
|
||||
| `CONCURRENCY_MANAGER`, `ACTIVE_GET_REQUESTS`, `ACTIVE_PUT_REQUESTS` | `rustfs/src/storage/concurrency/*` | Process-global owner-local scheduler state | Storage concurrency manager and request counters remain inside the storage concurrency owner boundary. |
|
||||
| `GET_OBJECT_BUFFER_THRESHOLD_WARNED`, `GET_READER_STREAM_BUFFER_SIZE_OVERRIDE`, function-local `ENABLED`, `OBJECT_SEEK_SUPPORT_THRESHOLD`, `OBJECT_SEEK_SUPPORT_CONCURRENCY_THRESHOLDS` | `rustfs/src/app/object_usecase.rs` | Cache or constant / owner-local cache | Object GET/seek tuning caches and warning guards stay private to object usecase helpers. |
|
||||
| `SUPPORTED_HEADERS` | `rustfs/src/storage/options.rs` | Cache or constant / owner-local constant | Supported-header lookup state stays private to storage option parsing. |
|
||||
| `AUDIT_TARGET_SPECS`, `NOTIFICATION_TARGET_SPECS` | `rustfs/src/admin/handlers/audit.rs`, `rustfs/src/admin/handlers/event.rs`, `rustfs/src/admin/handlers/plugins_instances.rs` | Cache or constant / owner-local constant | Admin target descriptor tables stay private to their handler owners. |
|
||||
|
||||
@@ -14,21 +14,12 @@
|
||||
|
||||
//! I/O scheduling types for adaptive buffer sizing and load management.
|
||||
//!
|
||||
//! # Migration Note
|
||||
//!
|
||||
//! This module contains types that are also available in `rustfs_io_core`.
|
||||
//! For new code, prefer using types from `rustfs_io_core` directly:
|
||||
//!
|
||||
//! ```ignore
|
||||
//! // Recommended: Use io-core types
|
||||
//! use rustfs_io_core::{
|
||||
//! IoLoadLevel, IoPriority, IoSchedulerConfig,
|
||||
//! calculate_optimal_buffer_size, get_buffer_size_for_media,
|
||||
//! };
|
||||
//! ```
|
||||
//!
|
||||
//! This module remains for backward compatibility and provides additional
|
||||
//! runtime monitoring features (`IoPriorityMetrics`, `IoStrategyDebugInfo`).
|
||||
//! This is the live scheduling implementation. `rustfs_io_core` supplies the
|
||||
//! shared config shapes (`IoSchedulerConfig`, `IoPriorityQueueConfig`) that the
|
||||
//! types here project into through `to_core_config`, plus the `io_profile`
|
||||
//! storage-media model; bandwidth samples come from `rustfs_io_metrics`.
|
||||
//! Same-named io-core types are those config shapes, not a backing
|
||||
//! implementation this module delegates to.
|
||||
|
||||
use rustfs_config::{KI_B, MI_B};
|
||||
use rustfs_io_core::io_profile::{AccessPattern, StorageMedia, StorageProfile};
|
||||
@@ -1762,169 +1753,6 @@ impl<T> IoPriorityQueue<T> {
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// I/O Priority Queue Metrics
|
||||
// ============================================
|
||||
|
||||
/// Global metrics for I/O priority queue monitoring.
|
||||
///
|
||||
/// These metrics are emitted through the shared metrics pipeline and provide
|
||||
/// visibility into the priority queue behavior.
|
||||
#[allow(dead_code)]
|
||||
pub struct IoPriorityMetrics {
|
||||
/// High priority queue depth.
|
||||
pub high_queue_depth: AtomicU64,
|
||||
/// Normal priority queue depth.
|
||||
pub normal_queue_depth: AtomicU64,
|
||||
/// Low priority queue depth.
|
||||
pub low_queue_depth: AtomicU64,
|
||||
/// High priority total wait time in nanoseconds.
|
||||
pub high_wait_time_ns: AtomicU64,
|
||||
/// Normal priority total wait time in nanoseconds.
|
||||
pub normal_wait_time_ns: AtomicU64,
|
||||
/// Low priority total wait time in nanoseconds.
|
||||
pub low_wait_time_ns: AtomicU64,
|
||||
/// Total starvation events count.
|
||||
pub starvation_events: AtomicU64,
|
||||
/// High priority requests processed.
|
||||
pub high_processed: AtomicU64,
|
||||
/// Normal priority requests processed.
|
||||
pub normal_processed: AtomicU64,
|
||||
/// Low priority requests processed.
|
||||
pub low_processed: AtomicU64,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl Default for IoPriorityMetrics {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl IoPriorityMetrics {
|
||||
/// Create a new metrics instance.
|
||||
pub const fn new() -> Self {
|
||||
Self {
|
||||
high_queue_depth: AtomicU64::new(0),
|
||||
normal_queue_depth: AtomicU64::new(0),
|
||||
low_queue_depth: AtomicU64::new(0),
|
||||
high_wait_time_ns: AtomicU64::new(0),
|
||||
normal_wait_time_ns: AtomicU64::new(0),
|
||||
low_wait_time_ns: AtomicU64::new(0),
|
||||
starvation_events: AtomicU64::new(0),
|
||||
high_processed: AtomicU64::new(0),
|
||||
normal_processed: AtomicU64::new(0),
|
||||
low_processed: AtomicU64::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
/// Update queue depths from status.
|
||||
#[allow(dead_code)]
|
||||
pub fn update_queue_depths(&self, status: &IoQueueStatus) {
|
||||
self.high_queue_depth
|
||||
.store(status.high_priority_waiting as u64, Ordering::Relaxed);
|
||||
self.normal_queue_depth
|
||||
.store(status.normal_priority_waiting as u64, Ordering::Relaxed);
|
||||
self.low_queue_depth
|
||||
.store(status.low_priority_waiting as u64, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Record a starvation event.
|
||||
#[allow(dead_code)]
|
||||
pub fn record_starvation(&self) {
|
||||
self.starvation_events.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Record a processed request.
|
||||
#[allow(dead_code)]
|
||||
pub fn record_processed(&self, priority: IoPriority) {
|
||||
match priority {
|
||||
IoPriority::High => self.high_processed.fetch_add(1, Ordering::Relaxed),
|
||||
IoPriority::Normal => self.normal_processed.fetch_add(1, Ordering::Relaxed),
|
||||
IoPriority::Low => self.low_processed.fetch_add(1, Ordering::Relaxed),
|
||||
};
|
||||
}
|
||||
|
||||
/// Record wait time for a priority level.
|
||||
pub fn record_wait_time(&self, priority: IoPriority, wait_ns: u64) {
|
||||
match priority {
|
||||
IoPriority::High => self.high_wait_time_ns.fetch_add(wait_ns, Ordering::Relaxed),
|
||||
IoPriority::Normal => self.normal_wait_time_ns.fetch_add(wait_ns, Ordering::Relaxed),
|
||||
IoPriority::Low => self.low_wait_time_ns.fetch_add(wait_ns, Ordering::Relaxed),
|
||||
};
|
||||
}
|
||||
|
||||
/// Get high priority queue depth.
|
||||
pub fn get_high_queue_depth(&self) -> u64 {
|
||||
self.high_queue_depth.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Get normal priority queue depth.
|
||||
pub fn get_normal_queue_depth(&self) -> u64 {
|
||||
self.normal_queue_depth.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Get low priority queue depth.
|
||||
pub fn get_low_queue_depth(&self) -> u64 {
|
||||
self.low_queue_depth.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Get total starvation events.
|
||||
pub fn get_starvation_events(&self) -> u64 {
|
||||
self.starvation_events.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Get metrics summary for logging/debugging.
|
||||
pub fn summary(&self) -> String {
|
||||
format!(
|
||||
"high_queue={}, normal_queue={}, low_queue={}, starvation={}, high_proc={}, normal_proc={}, low_proc={}",
|
||||
self.get_high_queue_depth(),
|
||||
self.get_normal_queue_depth(),
|
||||
self.get_low_queue_depth(),
|
||||
self.get_starvation_events(),
|
||||
self.high_processed.load(Ordering::Relaxed),
|
||||
self.normal_processed.load(Ordering::Relaxed),
|
||||
self.low_processed.load(Ordering::Relaxed)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// Global I/O priority metrics instance.
|
||||
#[allow(dead_code)]
|
||||
pub static IO_PRIORITY_METRICS: IoPriorityMetrics = IoPriorityMetrics::new();
|
||||
|
||||
/// Get optimized buffer size for I/O operations.
|
||||
///
|
||||
/// This function provides adaptive buffer sizing based on:
|
||||
/// - File size (small files get smaller buffers)
|
||||
/// - Concurrent request count (high concurrency gets smaller buffers)
|
||||
/// - Base buffer size from configuration
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `file_size` - Size of the file being read/written (-1 for unknown)
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// Optimal buffer size in bytes
|
||||
///
|
||||
/// # Example
|
||||
///
|
||||
/// ```ignore
|
||||
/// let buffer_size = get_buffer_size_opt_in(1024 * 1024); // 1MB file
|
||||
/// assert!(buffer_size >= 64 * 1024); // At least 64KB
|
||||
/// ```
|
||||
#[allow(dead_code)]
|
||||
pub fn get_buffer_size_opt_in(file_size: i64) -> usize {
|
||||
// Get base buffer size from configuration
|
||||
let base_buffer_size =
|
||||
rustfs_utils::get_env_usize(rustfs_config::ENV_OBJECT_IO_BUFFER_SIZE, rustfs_config::DEFAULT_OBJECT_IO_BUFFER_SIZE);
|
||||
|
||||
// Apply concurrency-aware adjustments
|
||||
get_concurrency_aware_buffer_size(file_size, base_buffer_size)
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// Unit Tests
|
||||
// ============================================
|
||||
@@ -1933,13 +1761,12 @@ pub fn get_buffer_size_opt_in(file_size: i64) -> usize {
|
||||
#[allow(unused_imports)]
|
||||
mod tests {
|
||||
use super::{
|
||||
IoLoadLevel, IoPriority, IoPriorityMetrics, IoPriorityQueue, IoPriorityQueueConfig, IoSchedulerConfig,
|
||||
IoSchedulingContext, IoStrategy, get_advanced_buffer_size, get_buffer_size_opt_in, get_concurrency_aware_buffer_size,
|
||||
IoLoadLevel, IoPriority, IoPriorityQueue, IoPriorityQueueConfig, IoSchedulerConfig, IoSchedulingContext, IoStrategy,
|
||||
get_advanced_buffer_size, get_concurrency_aware_buffer_size,
|
||||
};
|
||||
use rustfs_io_core::io_profile::{AccessPattern, StorageMedia};
|
||||
use rustfs_io_metrics::bandwidth::{BandwidthSnapshot, BandwidthTier};
|
||||
use serial_test::serial;
|
||||
use std::sync::atomic::Ordering;
|
||||
use std::time::Duration;
|
||||
|
||||
#[tokio::test]
|
||||
@@ -2126,29 +1953,6 @@ mod tests {
|
||||
assert_eq!(config.starvation_threshold_secs, 120);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_io_priority_metrics() {
|
||||
let metrics = IoPriorityMetrics::new();
|
||||
|
||||
// Test initial state
|
||||
assert_eq!(metrics.get_high_queue_depth(), 0);
|
||||
assert_eq!(metrics.get_normal_queue_depth(), 0);
|
||||
assert_eq!(metrics.get_low_queue_depth(), 0);
|
||||
assert_eq!(metrics.get_starvation_events(), 0);
|
||||
|
||||
// Test recording
|
||||
metrics.record_starvation();
|
||||
assert_eq!(metrics.get_starvation_events(), 1);
|
||||
|
||||
metrics.record_processed(IoPriority::High);
|
||||
metrics.record_processed(IoPriority::High);
|
||||
metrics.record_processed(IoPriority::Normal);
|
||||
|
||||
assert_eq!(metrics.high_processed.load(Ordering::Relaxed), 2);
|
||||
assert_eq!(metrics.normal_processed.load(Ordering::Relaxed), 1);
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// Multi-Factor Strategy Tests
|
||||
// ============================================
|
||||
|
||||
@@ -24,16 +24,14 @@
|
||||
//! - **Concurrency Management**: Coordination of concurrent GetObject requests
|
||||
//! - **Request Tracking**: RAII guards for request lifecycle management
|
||||
//!
|
||||
//! # Migration Note
|
||||
//! # Relationship to the shared crates
|
||||
//!
|
||||
//! Core algorithms have been migrated to `rustfs-io-core` and metrics to
|
||||
//! `rustfs-io-metrics`. This module maintains API compatibility while
|
||||
//! delegating to the new implementations.
|
||||
//! The scheduling algorithm lives in [`io_schedule`], not in `rustfs-io-core`:
|
||||
//! this module does not delegate to it. `rustfs-io-core` owns the shared
|
||||
//! config shapes and the `io_profile` storage-media model that [`io_schedule`]
|
||||
//! consumes, and `rustfs-io-metrics` owns bandwidth sampling and metric
|
||||
//! recording.
|
||||
|
||||
// Sub-modules
|
||||
// pub mod bandwidth_monitor; // Migrated to rustfs-io-metrics
|
||||
// pub mod global_metrics; // Migrated to rustfs-io-metrics
|
||||
// pub mod io_profile; // Migrated to rustfs-io-core
|
||||
pub mod io_schedule;
|
||||
pub mod manager;
|
||||
pub mod request_guard;
|
||||
@@ -45,9 +43,8 @@ pub mod request_guard;
|
||||
// I/O scheduling types (from io_schedule.rs for backward compatibility)
|
||||
#[allow(unused_imports)]
|
||||
pub use io_schedule::{
|
||||
IO_PRIORITY_METRICS, IoLoadLevel, IoPriority, IoPriorityMetrics, IoPriorityQueue, IoPriorityQueueConfig, IoQueueStatus,
|
||||
IoSchedulerConfig, IoStrategy, get_advanced_buffer_size, get_buffer_size_opt_in, get_concurrency_aware_buffer_size,
|
||||
get_put_concurrency_aware_buffer_size,
|
||||
IoLoadLevel, IoPriority, IoPriorityQueue, IoPriorityQueueConfig, IoQueueStatus, IoSchedulerConfig, IoStrategy,
|
||||
get_advanced_buffer_size, get_concurrency_aware_buffer_size, get_put_concurrency_aware_buffer_size,
|
||||
};
|
||||
|
||||
// Request tracking
|
||||
@@ -56,24 +53,6 @@ pub use request_guard::{GetObjectGuard, PutObjectGuard};
|
||||
// Concurrency manager
|
||||
pub use manager::{ConcurrencyManager, DiskReadAdmission, PutObjectAdmission};
|
||||
|
||||
// ============================================
|
||||
// New Module Re-exports (for gradual migration)
|
||||
// ============================================
|
||||
|
||||
// Re-export types from rustfs-io-core for convenience
|
||||
pub use rustfs_io_core::{
|
||||
// Backpressure types
|
||||
BackpressureMonitor,
|
||||
// Deadlock detection types
|
||||
DeadlockDetector,
|
||||
// Scheduler types
|
||||
IoScheduler,
|
||||
// Lock optimization types
|
||||
LockOptimizer,
|
||||
};
|
||||
|
||||
// Re-export types from rustfs-io-metrics for convenience
|
||||
|
||||
// ============================================
|
||||
// Helper Functions
|
||||
// ============================================
|
||||
@@ -83,37 +62,8 @@ pub fn get_concurrency_manager() -> &'static ConcurrencyManager {
|
||||
ConcurrencyManager::global()
|
||||
}
|
||||
|
||||
/// Reset the active get requests counter (for testing).
|
||||
#[allow(dead_code)]
|
||||
pub fn reset_active_get_requests() {
|
||||
io_schedule::ACTIVE_GET_REQUESTS.store(0, std::sync::atomic::Ordering::Relaxed);
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
/// Reset the active put requests counter (for testing).
|
||||
#[cfg(test)]
|
||||
pub fn reset_active_put_requests() {
|
||||
io_schedule::ACTIVE_PUT_REQUESTS.store(0, std::sync::atomic::Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Create a new I/O scheduler with default configuration.
|
||||
#[allow(dead_code)]
|
||||
pub fn create_io_scheduler() -> IoScheduler {
|
||||
IoScheduler::with_defaults()
|
||||
}
|
||||
|
||||
/// Create a new backpressure monitor with default configuration.
|
||||
#[allow(dead_code)]
|
||||
pub fn create_backpressure_monitor() -> BackpressureMonitor {
|
||||
BackpressureMonitor::with_defaults()
|
||||
}
|
||||
|
||||
/// Create a new deadlock detector with default configuration.
|
||||
#[allow(dead_code)]
|
||||
pub fn create_deadlock_detector() -> DeadlockDetector {
|
||||
DeadlockDetector::with_defaults()
|
||||
}
|
||||
|
||||
/// Create a new lock optimizer with default configuration.
|
||||
#[allow(dead_code)]
|
||||
pub fn create_lock_optimizer() -> LockOptimizer {
|
||||
LockOptimizer::with_defaults()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user