Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 54 additions & 2 deletions crates/integrations/cache-moka/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,31 @@
// under the License.

use std::hash::Hash;
use std::mem::size_of_val;
use std::sync::Arc;

use iceberg::cache::{ObjectCache, ObjectCacheProvide};
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<V>(max_capacity_bytes: u64) -> moka::sync::Cache<String, Arc<V>>
where V: Send + Sync + 'static {
moka::sync::Cache::builder()
.weigher(|_, value: &Arc<V>| {
// 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<K, V>(moka::sync::Cache<K, V>);

impl<K, V> ObjectCache<K, V> for MokaObjectCache<K, V>
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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));
}
}
}
Loading