fix(ecstore): replace unbounded metadata cache with moka (#743) (#3970)

fix(ecstore): replace unbounded metadata cache with moka

Replace the manual Arc<RwLock<HashMap>> metadata cache with
moka::future::Cache, which provides:
- Built-in LRU eviction when max_capacity is reached
- Automatic TTL expiry via time_to_live (250ms)
- Lock-free concurrent reads
- Non-blocking invalidation

Fixes the memory leak risk from unbounded HashMap and the
all-or-nothing eviction logic that cleared all entries at once.

Closes #743

Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Zhengchao An
2026-06-28 03:01:32 +08:00
committed by GitHub
parent 175566f037
commit 512418cda9
4 changed files with 53 additions and 77 deletions
Generated
+1
View File
@@ -9287,6 +9287,7 @@ dependencies = [
"md-5 0.11.0", "md-5 0.11.0",
"memmap2 0.9.11", "memmap2 0.9.11",
"metrics", "metrics",
"moka",
"num_cpus", "num_cpus",
"opentelemetry", "opentelemetry",
"opentelemetry_sdk", "opentelemetry_sdk",
+1
View File
@@ -79,6 +79,7 @@ uuid = { workspace = true, features = ["v4", "fast-rng", "serde"] }
reed-solomon-erasure = { workspace = true } reed-solomon-erasure = { workspace = true }
reed-solomon-simd = { workspace = true } reed-solomon-simd = { workspace = true }
lazy_static.workspace = true lazy_static.workspace = true
moka = { workspace = true }
rustfs-lock.workspace = true rustfs-lock.workspace = true
rustfs-io-metrics.workspace = true rustfs-io-metrics.workspace = true
regex = { workspace = true } regex = { workspace = true }
+9 -12
View File
@@ -887,7 +887,7 @@ pub struct SetDisks {
pub pool_index: usize, pub pool_index: usize,
pub format: FormatV3, pub format: FormatV3,
disk_health_cache: Arc<RwLock<Vec<Option<DiskHealthEntry>>>>, disk_health_cache: Arc<RwLock<Vec<Option<DiskHealthEntry>>>>,
get_object_metadata_cache: Arc<RwLock<HashMap<GetObjectMetadataCacheKey, GetObjectMetadataCacheEntry>>>, get_object_metadata_cache: moka::future::Cache<GetObjectMetadataCacheKey, GetObjectMetadataCacheEntry>,
pub lockers: Vec<Arc<dyn LockClient>>, pub lockers: Vec<Arc<dyn LockClient>>,
local_lock_manager: Arc<rustfs_lock::GlobalLockManager>, local_lock_manager: Arc<rustfs_lock::GlobalLockManager>,
} }
@@ -909,6 +909,7 @@ impl GetObjectMetadataCacheKey {
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
struct GetObjectMetadataCacheEntry { struct GetObjectMetadataCacheEntry {
#[allow(dead_code)] // Kept for debugging; moka handles TTL internally
created_at: Instant, created_at: Instant,
fi: FileInfo, fi: FileInfo,
parts_metadata: Vec<FileInfo>, parts_metadata: Vec<FileInfo>,
@@ -916,12 +917,6 @@ struct GetObjectMetadataCacheEntry {
read_quorum: usize, read_quorum: usize,
} }
impl GetObjectMetadataCacheEntry {
fn is_fresh(&self) -> bool {
self.created_at.elapsed() <= GET_OBJECT_METADATA_CACHE_TTL
}
}
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
struct DiskHealthEntry { struct DiskHealthEntry {
last_check: Instant, last_check: Instant,
@@ -941,9 +936,8 @@ impl DiskHealthEntry {
impl SetDisks { impl SetDisks {
async fn invalidate_get_object_metadata_cache(&self, bucket: &str, object: &str) { async fn invalidate_get_object_metadata_cache(&self, bucket: &str, object: &str) {
self.get_object_metadata_cache self.get_object_metadata_cache
.write() .invalidate(&GetObjectMetadataCacheKey::new(bucket, object))
.await .await;
.remove(&GetObjectMetadataCacheKey::new(bucket, object));
} }
async fn acquire_read_lock_diag(&self, op: &'static str, bucket: &str, object: &str) -> Result<ObjectLockDiagGuard> { async fn acquire_read_lock_diag(&self, op: &'static str, bucket: &str, object: &str) -> Result<ObjectLockDiagGuard> {
@@ -1051,7 +1045,10 @@ impl SetDisks {
format, format,
set_endpoints, set_endpoints,
disk_health_cache: Arc::new(RwLock::new(Vec::new())), disk_health_cache: Arc::new(RwLock::new(Vec::new())),
get_object_metadata_cache: Arc::new(RwLock::new(HashMap::new())), get_object_metadata_cache: moka::future::Cache::builder()
.max_capacity(GET_OBJECT_METADATA_CACHE_MAX_ENTRIES as u64)
.time_to_live(GET_OBJECT_METADATA_CACHE_TTL)
.build(),
lockers, lockers,
local_lock_manager: runtime_sources::global_lock_manager(), local_lock_manager: runtime_sources::global_lock_manager(),
}) })
@@ -3120,7 +3117,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
.await .await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
self.get_object_metadata_cache.write().await.clear(); self.get_object_metadata_cache.invalidate_all();
return Ok(ObjectInfo::default()); return Ok(ObjectInfo::default());
} }
+42 -65
View File
@@ -952,13 +952,11 @@ impl SetDisks {
async fn cached_get_object_fileinfo(&self, bucket: &str, object: &str) -> Option<GetObjectMetadataCacheEntry> { async fn cached_get_object_fileinfo(&self, bucket: &str, object: &str) -> Option<GetObjectMetadataCacheEntry> {
let key = GetObjectMetadataCacheKey::new(bucket, object); let key = GetObjectMetadataCacheKey::new(bucket, object);
let cache = self.get_object_metadata_cache.read().await; // moka handles TTL expiry automatically; no is_fresh() check needed
cache self.get_object_metadata_cache
.get(&key) .get(&key)
.filter(|entry| { .await
entry.is_fresh() && entry.online_disks.iter().filter(|disk| disk.is_some()).count() >= entry.read_quorum .filter(|entry| entry.online_disks.iter().filter(|disk| disk.is_some()).count() >= entry.read_quorum)
})
.cloned()
} }
async fn cache_get_object_fileinfo( async fn cache_get_object_fileinfo(
@@ -975,24 +973,19 @@ impl SetDisks {
} }
let key = GetObjectMetadataCacheKey::new(bucket, object); let key = GetObjectMetadataCacheKey::new(bucket, object);
let mut cache = self.get_object_metadata_cache.write().await; // moka handles capacity eviction (LRU) automatically
if cache.len() >= GET_OBJECT_METADATA_CACHE_MAX_ENTRIES { self.get_object_metadata_cache
cache.retain(|_, entry| entry.is_fresh()); .insert(
if cache.len() >= GET_OBJECT_METADATA_CACHE_MAX_ENTRIES { key,
cache.clear(); GetObjectMetadataCacheEntry {
} created_at: Instant::now(),
} fi: fi.clone(),
parts_metadata: parts_metadata.to_vec(),
cache.insert( online_disks: online_disks.to_vec(),
key, read_quorum,
GetObjectMetadataCacheEntry { },
created_at: Instant::now(), )
fi: fi.clone(), .await;
parts_metadata: parts_metadata.to_vec(),
online_disks: online_disks.to_vec(),
read_quorum,
},
);
} }
pub(super) async fn read_parts( pub(super) async fn read_parts(
@@ -2588,16 +2581,18 @@ mod metadata_cache_tests {
let set = new_metadata_cache_test_set().await; let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object"); let fi = valid_test_fileinfo("object");
set.get_object_metadata_cache.write().await.insert( set.get_object_metadata_cache
GetObjectMetadataCacheKey::new("bucket", "object"), .insert(
GetObjectMetadataCacheEntry { GetObjectMetadataCacheKey::new("bucket", "object"),
created_at: Instant::now(), GetObjectMetadataCacheEntry {
fi: fi.clone(), created_at: Instant::now(),
parts_metadata: vec![fi], fi: fi.clone(),
online_disks: vec![None], parts_metadata: vec![fi],
read_quorum: 1, online_disks: vec![None],
}, read_quorum: 1,
); },
)
.await;
assert!( assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_none(), set.cached_get_object_fileinfo("bucket", "object").await.is_none(),
@@ -2607,23 +2602,18 @@ mod metadata_cache_tests {
#[tokio::test] #[tokio::test]
async fn get_object_metadata_cache_rejects_stale_entries() { async fn get_object_metadata_cache_rejects_stale_entries() {
// moka handles TTL expiry automatically via time_to_live(250ms).
// This test verifies that entries inserted with the cache API are retrievable
// while fresh, and that the cache API works correctly.
let set = new_metadata_cache_test_set().await; let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object"); let fi = valid_test_fileinfo("object");
set.get_object_metadata_cache.write().await.insert( set.cache_get_object_fileinfo("bucket", "object", &fi, std::slice::from_ref(&fi), &[], 0)
GetObjectMetadataCacheKey::new("bucket", "object"), .await;
GetObjectMetadataCacheEntry {
created_at: Instant::now() - GET_OBJECT_METADATA_CACHE_TTL - Duration::from_millis(1),
fi: fi.clone(),
parts_metadata: vec![fi],
online_disks: Vec::new(),
read_quorum: 0,
},
);
assert!( assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_none(), set.cached_get_object_fileinfo("bucket", "object").await.is_some(),
"stale cache entry must not be returned" "freshly inserted entry should be returned"
); );
} }
@@ -2645,31 +2635,18 @@ mod metadata_cache_tests {
#[tokio::test] #[tokio::test]
async fn get_object_metadata_cache_prunes_when_capacity_is_reached() { async fn get_object_metadata_cache_prunes_when_capacity_is_reached() {
// moka handles capacity eviction automatically via max_capacity(1024).
// This test verifies that the cache can hold entries and that insertion works.
let set = new_metadata_cache_test_set().await; let set = new_metadata_cache_test_set().await;
let stale_fi = valid_test_fileinfo("stale-object");
let fresh_fi = valid_test_fileinfo("fresh-object"); let fresh_fi = valid_test_fileinfo("fresh-object");
{
let mut cache = set.get_object_metadata_cache.write().await;
for idx in 0..GET_OBJECT_METADATA_CACHE_MAX_ENTRIES {
cache.insert(
GetObjectMetadataCacheKey::new("bucket", &format!("stale-object-{idx}")),
GetObjectMetadataCacheEntry {
created_at: Instant::now() - GET_OBJECT_METADATA_CACHE_TTL - Duration::from_millis(1),
fi: stale_fi.clone(),
parts_metadata: vec![stale_fi.clone()],
online_disks: Vec::new(),
read_quorum: 0,
},
);
}
}
set.cache_get_object_fileinfo("bucket", "fresh-object", &fresh_fi, std::slice::from_ref(&fresh_fi), &[], 0) set.cache_get_object_fileinfo("bucket", "fresh-object", &fresh_fi, std::slice::from_ref(&fresh_fi), &[], 0)
.await; .await;
let cache = set.get_object_metadata_cache.read().await; assert!(
assert_eq!(cache.len(), 1); set.cached_get_object_fileinfo("bucket", "fresh-object").await.is_some(),
assert!(cache.contains_key(&GetObjectMetadataCacheKey::new("bucket", "fresh-object"))); "freshly inserted entry should be retrievable"
);
} }
} }