From afa3935ebba9ad513e4bed060a873715b50d861d Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 8 Jul 2026 17:14:46 +0800 Subject: [PATCH] fix(object-data-cache): prune identity index on moka eviction (#4458) Moka evicts entries by TTL, time-to-idle, and capacity-LRU without going through invalidate_object, so the per-object identity index (by_object) never dropped the evicted keys and grew without bound, bypassing the memory cap. Register an async eviction listener that removes truly evicted keys from the identity index, holding only a Weak reference to avoid an Arc cycle. Refs rustfs/backlog#1031 --- crates/object-data-cache/src/moka_backend.rs | 62 +++++++++++++++++-- .../object-data-cache/src/starshard_index.rs | 6 ++ 2 files changed, 63 insertions(+), 5 deletions(-) diff --git a/crates/object-data-cache/src/moka_backend.rs b/crates/object-data-cache/src/moka_backend.rs index 3bdfaa908..0f4d81fd2 100644 --- a/crates/object-data-cache/src/moka_backend.rs +++ b/crates/object-data-cache/src/moka_backend.rs @@ -16,19 +16,19 @@ use crate::cache::{ObjectDataCacheFillResult, ObjectDataCacheGetPlan, ObjectData use crate::config::ObjectDataCacheConfig; use crate::entry::ObjectDataCacheEntry; use crate::index::{ObjectDataCacheIdentityIndex, ObjectDataCacheIndexInsertResult}; -use crate::key::ObjectDataCacheIdentity; +use crate::key::{ObjectDataCacheIdentity, ObjectDataCacheKey}; use crate::memory::ObjectDataCacheMemoryGate; use crate::singleflight::{ObjectDataCacheSingleflight, ObjectDataCacheSingleflightAcquire}; use crate::stats::ObjectDataCacheStats; use bytes::Bytes; -use moka::future::Cache; +use moka::future::{Cache, FutureExt}; use std::sync::Arc; /// Weighted Moka backend for reusable object bodies. #[derive(Debug)] pub struct MokaBackend { - cache: Cache>, - index: ObjectDataCacheIdentityIndex, + cache: Cache>, + index: Arc, singleflight: ObjectDataCacheSingleflight, memory_gate: ObjectDataCacheMemoryGate, } @@ -42,16 +42,43 @@ impl MokaBackend { let max_capacity = config.resolved_max_bytes()?; let ttl = config.ttl; let time_to_idle = config.time_to_idle; + + // The identity index tracks the keys cached per object so that + // invalidation can find them without a full-cache scan. Moka evicts + // entries on its own (TTL, time-to-idle, capacity-LRU) without going + // through `invalidate_object`, so without a listener those evicted keys + // would linger in the index forever and leak memory. Hold only a `Weak` + // reference in the listener to avoid an Arc cycle keeping the index (and + // thus the cache) alive. + let index = Arc::new(ObjectDataCacheIdentityIndex::new(usize::from(config.identity_keys_max))); + let index_for_eviction = Arc::downgrade(&index); + let cache = Cache::builder() .max_capacity(max_capacity) .weigher(|key, value: &Arc| value.estimated_weight(key)) .time_to_live(ttl) .time_to_idle(time_to_idle) + .async_eviction_listener(move |key, _value, cause| { + let index_for_eviction = index_for_eviction.clone(); + async move { + // Explicit removals and replacements are already reconciled + // with the index by the caller; only true evictions leak. + if !cause.was_evicted() { + return; + } + let Some(index) = index_for_eviction.upgrade() else { + return; + }; + let identity = ObjectDataCacheIdentity::new(Arc::clone(&key.bucket), Arc::clone(&key.object)); + index.remove_key(&identity, &key).await; + } + .boxed() + }) .build(); Ok(Self { cache, - index: ObjectDataCacheIdentityIndex::new(usize::from(config.identity_keys_max)), + index, singleflight: ObjectDataCacheSingleflight::new(Arc::clone(&stats)), memory_gate: ObjectDataCacheMemoryGate::new(config, stats), }) @@ -228,6 +255,31 @@ mod tests { assert!(matches!(lookup, ObjectDataCacheLookup::Miss)); } + #[tokio::test] + async fn moka_backend_eviction_prunes_identity_index() { + let backend = + MokaBackend::new(&enabled_config(), Arc::new(ObjectDataCacheStats::default())).expect("moka backend should build"); + + // Cache several distinct identities under the short TTL from the config. + for i in 0..8 { + let plan = cacheable_plan(&format!("object-{i}"), "etag-a"); + let _ = backend.fill_body(&plan, Bytes::from_static(b"hello")).await; + } + assert_eq!(backend.index.identity_count().await, 8); + + // Let every entry expire, then let moka process the expirations so the + // eviction listener runs and prunes the identity index. + tokio::time::sleep(Duration::from_millis(150)).await; + backend.cache.run_pending_tasks().await; + + assert_eq!(backend.cache.entry_count(), 0, "all entries should have expired"); + assert_eq!( + backend.index.identity_count().await, + 0, + "identity index must not outlive evicted cache entries" + ); + } + #[tokio::test] async fn moka_backend_expires_entries_by_tti() { let mut config = enabled_config(); diff --git a/crates/object-data-cache/src/starshard_index.rs b/crates/object-data-cache/src/starshard_index.rs index f3a01e9b4..c20dc6103 100644 --- a/crates/object-data-cache/src/starshard_index.rs +++ b/crates/object-data-cache/src/starshard_index.rs @@ -118,6 +118,12 @@ impl StarshardIdentityIndex { .is_some_and(|key_set| key_set.contains(key)) } + /// Returns the number of tracked identities. + #[cfg(test)] + pub(crate) async fn identity_count(&self) -> usize { + self.by_object.len().await + } + /// Removes index keys that no longer exist in the cache. pub async fn prune_missing(&self, identity: &ObjectDataCacheIdentity, mut key_exists: F) where