From abffa5cf1b2d1e51f48aa6f59cec5b6e77ef9413 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Tue, 18 Aug 2026 12:35:36 +0800 Subject: [PATCH] chore(storage): drop dead io-schedule metrics and helpers (#6199) --- docs/architecture/global-state-inventory.md | 2 +- rustfs/src/storage/concurrency/io_schedule.rs | 212 +----------------- rustfs/src/storage/concurrency/mod.rs | 70 +----- 3 files changed, 19 insertions(+), 265 deletions(-) diff --git a/docs/architecture/global-state-inventory.md b/docs/architecture/global-state-inventory.md index 6c06b10c8..8407ad864 100644 --- a/docs/architecture/global-state-inventory.md +++ b/docs/architecture/global-state-inventory.md @@ -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. | diff --git a/rustfs/src/storage/concurrency/io_schedule.rs b/rustfs/src/storage/concurrency/io_schedule.rs index d611dfeb2..8d202708a 100644 --- a/rustfs/src/storage/concurrency/io_schedule.rs +++ b/rustfs/src/storage/concurrency/io_schedule.rs @@ -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 IoPriorityQueue { } } -// ============================================ -// 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 // ============================================ diff --git a/rustfs/src/storage/concurrency/mod.rs b/rustfs/src/storage/concurrency/mod.rs index 9b14eebda..eb47aa4ed 100644 --- a/rustfs/src/storage/concurrency/mod.rs +++ b/rustfs/src/storage/concurrency/mod.rs @@ -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() -}