diff --git a/rustfs/src/storage/concurrency.rs b/rustfs/src/storage/concurrency.rs index 331bd8901..5481c4c27 100644 --- a/rustfs/src/storage/concurrency.rs +++ b/rustfs/src/storage/concurrency.rs @@ -186,26 +186,26 @@ pub fn get_concurrency_aware_buffer_size(file_size: i64, base_buffer_size: usize /// ``` pub fn get_advanced_buffer_size(file_size: i64, base_buffer_size: usize, is_sequential: bool) -> usize { let concurrent_requests = ACTIVE_GET_REQUESTS.load(Ordering::Relaxed); - + // For very small files, use smaller buffers regardless of concurrency if file_size > 0 && file_size < 256 * KI_B as i64 { return (file_size as usize / 4).max(16 * KI_B).min(64 * KI_B); } - + // Base calculation from standard function let standard_size = get_concurrency_aware_buffer_size(file_size, base_buffer_size); - + // For sequential reads, we can be more aggressive with buffer sizes if is_sequential && concurrent_requests <= MEDIUM_CONCURRENCY_THRESHOLD { return ((standard_size as f64 * 1.5) as usize).min(2 * MI_B); } - + // For high concurrency with large files, optimize for parallel processing if concurrent_requests > HIGH_CONCURRENCY_THRESHOLD && file_size > 10 * MI_B as i64 { // Use smaller, more numerous buffers for better parallelism return (standard_size as f64 * 0.8) as usize; } - + standard_size } @@ -261,18 +261,18 @@ impl HotObjectCache { // Clone the data reference while holding read lock let data = Arc::clone(&cached.data); drop(cache); - + // Now acquire write lock to promote in LRU and update stats let mut cache_write = self.cache.write().await; if let Some(cached) = cache_write.get(key) { cached.hit_count.fetch_add(1, Ordering::Relaxed); - + #[cfg(feature = "metrics")] { use metrics::counter; counter!("rustfs_object_cache_hits").increment(1); } - + return Some(data); } } @@ -346,13 +346,13 @@ impl HotObjectCache { max_object_size: self.max_object_size, } } - + /// Check if a key exists in cache without promoting it in LRU async fn contains(&self, key: &str) -> bool { let cache = self.cache.read().await; cache.peek(key).is_some() } - + /// Get multiple objects from cache in a single operation /// /// This is more efficient than calling get() multiple times as it acquires @@ -360,12 +360,12 @@ impl HotObjectCache { async fn get_batch(&self, keys: &[String]) -> Vec>>> { let mut cache = self.cache.write().await; let mut results = Vec::with_capacity(keys.len()); - + for key in keys { if let Some(cached) = cache.get(key) { cached.hit_count.fetch_add(1, Ordering::Relaxed); results.push(Some(Arc::clone(&cached.data))); - + #[cfg(feature = "metrics")] { use metrics::counter; @@ -373,7 +373,7 @@ impl HotObjectCache { } } else { results.push(None); - + #[cfg(feature = "metrics")] { use metrics::counter; @@ -381,24 +381,25 @@ impl HotObjectCache { } } } - + results } - + /// Remove a specific key from cache /// /// Returns true if the key was found and removed, false otherwise. async fn remove(&self, key: &str) -> bool { let mut cache = self.cache.write().await; - if let Some((_, cached)) = cache.pop(key) { + if let Some(cached) = cache.pop(key) { let current = self.current_size.load(Ordering::Relaxed); - self.current_size.store(current.saturating_sub(cached.size), Ordering::Relaxed); + self.current_size + .store(current.saturating_sub(cached.size), Ordering::Relaxed); true } else { false } } - + /// Get the most frequently accessed keys /// /// Returns up to `limit` keys sorted by hit count in descending order. @@ -408,7 +409,7 @@ impl HotObjectCache { .iter() .map(|(k, v)| (k.clone(), v.hit_count.load(Ordering::Relaxed))) .collect(); - + keys_with_hits.sort_by(|a, b| b.1.cmp(&a.1)); keys_with_hits.truncate(limit); keys_with_hits @@ -491,35 +492,35 @@ impl ConcurrencyManager { pub async fn clear_cache(&self) { self.cache.clear().await } - + /// Check if a key exists in cache /// /// This is a lightweight check that doesn't promote the key in LRU order. pub async fn is_cached(&self, key: &str) -> bool { self.cache.contains(key).await } - + /// Get multiple objects from cache in batch /// /// More efficient than individual get_cached() calls when fetching multiple objects. pub async fn get_cached_batch(&self, keys: &[String]) -> Vec>>> { self.cache.get_batch(keys).await } - + /// Remove a specific object from cache /// /// Returns true if the object was cached and removed. pub async fn remove_cached(&self, key: &str) -> bool { self.cache.remove(key).await } - + /// Get the most frequently accessed keys /// /// Useful for identifying hot objects and optimizing cache strategies. pub async fn get_hot_keys(&self, limit: usize) -> Vec<(String, usize)> { self.cache.get_hot_keys(limit).await } - + /// Warm up cache with frequently accessed objects /// /// This can be called during server startup or maintenance windows diff --git a/rustfs/src/storage/concurrent_get_object_test.rs b/rustfs/src/storage/concurrent_get_object_test.rs index 511709cd9..120641a2f 100644 --- a/rustfs/src/storage/concurrent_get_object_test.rs +++ b/rustfs/src/storage/concurrent_get_object_test.rs @@ -24,8 +24,7 @@ #[cfg(test)] mod tests { use crate::storage::concurrency::{ - ConcurrencyManager, GetObjectGuard, - get_concurrency_aware_buffer_size, get_advanced_buffer_size + ConcurrencyManager, GetObjectGuard, get_advanced_buffer_size, get_concurrency_aware_buffer_size, }; use rustfs_config::{KI_B, MI_B}; use std::sync::Arc; @@ -289,24 +288,24 @@ mod tests { #[tokio::test] async fn test_cache_batch_operations() { let manager = ConcurrencyManager::new(); - + // Cache multiple objects for i in 0..10 { let key = format!("batch/object{}", i); let data = vec![i as u8; 100 * KI_B]; // 100KB each manager.cache_object(key, data).await; } - + // Test batch get let keys: Vec = (0..10).map(|i| format!("batch/object{}", i)).collect(); let results = manager.get_cached_batch(&keys).await; - + assert_eq!(results.len(), 10, "Should return result for each key"); for (i, result) in results.iter().enumerate() { assert!(result.is_some(), "Object {} should be in cache", i); assert_eq!(result.as_ref().unwrap().len(), 100 * KI_B); } - + // Test batch get with missing keys let mixed_keys = vec![ "batch/object0".to_string(), @@ -315,7 +314,7 @@ mod tests { "missing/key2".to_string(), ]; let mixed_results = manager.get_cached_batch(&mixed_keys).await; - + assert_eq!(mixed_results.len(), 4); assert!(mixed_results[0].is_some(), "First key should exist"); assert!(mixed_results[1].is_none(), "Second key should be missing"); @@ -327,7 +326,7 @@ mod tests { #[tokio::test] async fn test_cache_warming() { let manager = ConcurrencyManager::new(); - + // Prepare objects for warming let mut warm_objects = Vec::new(); for i in 0..5 { @@ -335,14 +334,14 @@ mod tests { let data = vec![i as u8; 512 * KI_B]; // 512KB each warm_objects.push((key, data)); } - + // Warm cache manager.warm_cache(warm_objects).await; - + // Verify all objects are cached let stats = manager.cache_stats().await; assert_eq!(stats.entries, 5, "All objects should be cached"); - + for i in 0..5 { let key = format!("warm/object{}", i); assert!(manager.is_cached(&key).await, "Object {} should be cached", i); @@ -353,14 +352,14 @@ mod tests { #[tokio::test] async fn test_hot_keys_tracking() { let manager = ConcurrencyManager::new(); - + // Cache objects with different access patterns for i in 0..5 { let key = format!("hot/object{}", i); let data = vec![i as u8; 100 * KI_B]; manager.cache_object(key, data).await; } - + // Simulate access patterns (object 0 and 1 are hot) for _ in 0..10 { let _ = manager.get_cached("hot/object0").await; @@ -371,10 +370,10 @@ mod tests { for _ in 0..2 { let _ = manager.get_cached("hot/object2").await; } - + // Get hot keys let hot_keys = manager.get_hot_keys(3).await; - + assert_eq!(hot_keys.len(), 3, "Should return top 3 keys"); assert_eq!(hot_keys[0].0, "hot/object0", "Most accessed should be first"); assert!(hot_keys[0].1 >= 10, "Object 0 should have at least 10 hits"); @@ -386,22 +385,22 @@ mod tests { #[tokio::test] async fn test_cache_removal() { let manager = ConcurrencyManager::new(); - + // Cache an object let key = "remove/test".to_string(); let data = vec![1u8; 100 * KI_B]; manager.cache_object(key.clone(), data).await; - + // Verify it's cached assert!(manager.is_cached(&key).await, "Object should be cached"); - + // Remove it let removed = manager.remove_cached(&key).await; assert!(removed, "Should successfully remove cached object"); - + // Verify it's gone assert!(!manager.is_cached(&key).await, "Object should no longer be cached"); - + // Try to remove non-existent key let not_removed = manager.remove_cached("nonexistent").await; assert!(!not_removed, "Should return false for non-existent key"); @@ -411,37 +410,34 @@ mod tests { #[tokio::test] async fn test_advanced_buffer_sizing() { let base_buffer = 256 * KI_B; // 256KB base - + // Test small file optimization let small_size = get_advanced_buffer_size(128 * KI_B as i64, base_buffer, false); assert!(small_size < base_buffer, "Small files should use smaller buffers"); assert!(small_size >= 16 * KI_B, "Should not go below minimum"); - + // Test sequential read optimization let _guard = ConcurrencyManager::track_request(); let sequential_size = get_advanced_buffer_size(10 * MI_B as i64, base_buffer, true); let random_size = get_advanced_buffer_size(10 * MI_B as i64, base_buffer, false); - + // Sequential reads should get larger buffers at low concurrency assert!( - sequential_size >= random_size, + sequential_size >= random_size, "Sequential reads should have equal or larger buffers at low concurrency" ); - + drop(_guard); - + // Test high concurrency with large files let mut guards = Vec::new(); for _ in 0..10 { guards.push(ConcurrencyManager::track_request()); } - + let high_concurrency_size = get_advanced_buffer_size(50 * MI_B as i64, base_buffer, false); - assert!( - high_concurrency_size <= base_buffer, - "High concurrency should reduce buffer size" - ); - + assert!(high_concurrency_size <= base_buffer, "High concurrency should reduce buffer size"); + drop(guards); } @@ -449,14 +445,14 @@ mod tests { #[tokio::test] async fn test_concurrent_cache_access() { let manager = Arc::new(ConcurrencyManager::new()); - + // Pre-populate cache for i in 0..50 { let key = format!("concurrent/object{}", i); let data = vec![i as u8; 100 * KI_B]; manager.cache_object(key, data).await; } - + // Spawn multiple concurrent readers let mut handles = Vec::new(); for worker_id in 0..10 { @@ -473,7 +469,7 @@ mod tests { }); handles.push(handle); } - + // Wait for all workers to complete let mut total_hits = 0; for handle in handles { @@ -481,7 +477,7 @@ mod tests { println!("Worker {} got {} cache hits", worker_id, hits); total_hits += hits; } - + // All workers should get hits for all cached objects assert_eq!(total_hits, 500, "Should have 500 total hits (10 workers * 50 objects)"); } @@ -490,23 +486,23 @@ mod tests { #[tokio::test] async fn test_is_cached_no_promotion() { let manager = ConcurrencyManager::new(); - + // Cache two objects manager.cache_object("first".to_string(), vec![1u8; 100 * KI_B]).await; manager.cache_object("second".to_string(), vec![2u8; 100 * KI_B]).await; - + // Check first without accessing it assert!(manager.is_cached("first").await); - + // Access second multiple times for _ in 0..5 { let _ = manager.get_cached("second").await; } - + // Both should still be cached assert!(manager.is_cached("first").await); assert!(manager.is_cached("second").await); - + // Get hot keys - second should be hotter let hot_keys = manager.get_hot_keys(2).await; assert_eq!(hot_keys[0].0, "second", "Second should be hottest"); diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 96f2678de..a56c042bf 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -124,6 +124,7 @@ use rustfs_utils::{ use rustfs_zip::CompressionFormat; use s3s::header::{X_AMZ_RESTORE, X_AMZ_RESTORE_OUTPUT_PATH}; use s3s::{S3, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, dto::*, s3_error}; +use std::convert::Infallible; use std::ops::Add; use std::{ collections::HashMap, @@ -132,7 +133,6 @@ use std::{ str::FromStr, sync::{Arc, LazyLock}, }; -use std::convert::Infallible; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tokio::{io::AsyncRead, sync::mpsc}; use tokio_stream::wrappers::ReceiverStream; @@ -1645,9 +1645,10 @@ impl S3 for FS { if let Some(cached_data) = manager.get_cached(&cache_key).await { debug!("Serving object from cache: {}", cache_key); + let value = cached_data.clone(); // Build response from cached data let body = Some(StreamingBlob::wrap::<_, Infallible>(futures::stream::once(async move { - Ok(bytes::Bytes::from((*cached_data).clone())) + Ok(bytes::Bytes::from((*value).clone())) }))); let output = GetObjectOutput {