diff --git a/crates/integrations/cache-moka/src/lib.rs b/crates/integrations/cache-moka/src/lib.rs index a3314d4c67..deef15b252 100644 --- a/crates/integrations/cache-moka/src/lib.rs +++ b/crates/integrations/cache-moka/src/lib.rs @@ -16,6 +16,7 @@ // under the License. use std::hash::Hash; +use std::mem::size_of_val; use std::sync::Arc; use iceberg::cache::{ObjectCache, ObjectCacheProvide}; @@ -23,6 +24,23 @@ use iceberg::spec::{Manifest, ManifestList}; const DEFAULT_CACHE_SIZE_BYTES: u64 = 32 * 1024 * 1024; // 32MiB +/// A cache whose `max_capacity` is a byte budget rather than an entry count. +/// +/// Without a weigher `moka` treats `max_capacity` as a number of entries, so passing a byte +/// figure to `Cache::new` leaves the cache effectively unbounded. This mirrors the weigher +/// `iceberg::io::ObjectCache` uses. +fn byte_bounded_cache(max_capacity_bytes: u64) -> moka::sync::Cache> +where V: Send + Sync + 'static { + moka::sync::Cache::builder() + .weigher(|_, value: &Arc| { + // Saturate rather than truncate: `moka` weights are `u32`, and a silently + // wrapped weight would let the cache grow past `max_capacity_bytes`. + u32::try_from(size_of_val(value.as_ref())).unwrap_or(u32::MAX) + }) + .max_capacity(max_capacity_bytes) + .build() +} + struct MokaObjectCache(moka::sync::Cache); impl ObjectCache for MokaObjectCache @@ -54,8 +72,8 @@ impl Default for MokaObjectCacheProvider { impl MokaObjectCacheProvider { /// Creates a new `MokaObjectCacheProvider` with default cache sizes. pub fn new() -> Self { - let manifest_cache = MokaObjectCache(moka::sync::Cache::new(DEFAULT_CACHE_SIZE_BYTES)); - let manifest_list_cache = MokaObjectCache(moka::sync::Cache::new(DEFAULT_CACHE_SIZE_BYTES)); + let manifest_cache = MokaObjectCache(byte_bounded_cache(DEFAULT_CACHE_SIZE_BYTES)); + let manifest_list_cache = MokaObjectCache(byte_bounded_cache(DEFAULT_CACHE_SIZE_BYTES)); Self { manifest_cache, @@ -88,3 +106,37 @@ impl ObjectCacheProvide for MokaObjectCacheProvider { &self.manifest_list_cache } } + +#[cfg(test)] +mod tests { + use super::*; + + /// A cache with no weigher gives every entry weight 1, so `weighted_size` would be the + /// entry count and a byte figure for `max_capacity` would admit 33_554_432 manifests. + #[test] + fn test_cache_weighs_entries_by_size_not_count() { + let cache = byte_bounded_cache::<[u8; 512]>(DEFAULT_CACHE_SIZE_BYTES); + cache.insert("a".to_string(), Arc::new([0u8; 512])); + cache.insert("b".to_string(), Arc::new([0u8; 512])); + cache.run_pending_tasks(); + + assert_eq!(cache.entry_count(), 2); + assert_eq!(cache.weighted_size(), 2 * 512); + assert_eq!( + cache.policy().max_capacity(), + Some(DEFAULT_CACHE_SIZE_BYTES) + ); + } + + #[test] + fn test_default_provider_caches_are_byte_bounded() { + let provider = MokaObjectCacheProvider::new(); + + for max in [ + provider.manifest_cache.0.policy().max_capacity(), + provider.manifest_list_cache.0.policy().max_capacity(), + ] { + assert_eq!(max, Some(DEFAULT_CACHE_SIZE_BYTES)); + } + } +}