From c521b210d3dfde8ddfab25067850a2abc77a3feb Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 7 Sep 2026 16:15:41 +0800 Subject: [PATCH 1/2] fix(cache-moka): bound the caches by bytes instead of entry count DEFAULT_CACHE_SIZE_BYTES went straight to Cache::new, but without a weigher moka reads max_capacity as a number of entries, so the 32MiB budget admitted 33_554_432 manifests. Add the weigher iceberg::io::ObjectCache already uses. --- crates/integrations/cache-moka/src/lib.rs | 52 ++++++++++++++++++++++- 1 file changed, 50 insertions(+), 2 deletions(-) diff --git a/crates/integrations/cache-moka/src/lib.rs b/crates/integrations/cache-moka/src/lib.rs index a3314d4c67..068cd105a4 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,19 @@ 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| size_of_val(value.as_ref()) as u32) + .max_capacity(max_capacity_bytes) + .build() +} + struct MokaObjectCache(moka::sync::Cache); impl ObjectCache for MokaObjectCache @@ -54,8 +68,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 +102,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)); + } + } +} From d8bf612362c07f945cfdff2377cbef9bc668bdf5 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 7 Sep 2026 19:18:17 +0800 Subject: [PATCH 2/2] fix(cache-moka): saturate the weigher conversion to u32 Review feedback: a truncating cast could wrap for a hypothetical value larger than u32::MAX. Unreachable for Manifest (72 bytes) and ManifestList (24 bytes), but the helper is generic, so saturate instead. --- crates/integrations/cache-moka/src/lib.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/crates/integrations/cache-moka/src/lib.rs b/crates/integrations/cache-moka/src/lib.rs index 068cd105a4..deef15b252 100644 --- a/crates/integrations/cache-moka/src/lib.rs +++ b/crates/integrations/cache-moka/src/lib.rs @@ -32,7 +32,11 @@ const DEFAULT_CACHE_SIZE_BYTES: u64 = 32 * 1024 * 1024; // 32MiB fn byte_bounded_cache(max_capacity_bytes: u64) -> moka::sync::Cache> where V: Send + Sync + 'static { moka::sync::Cache::builder() - .weigher(|_, value: &Arc| size_of_val(value.as_ref()) as u32) + .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() }