Skip to main content

datafusion_execution/cache/
default_cache.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use std::sync::{Arc, Mutex};
19use std::time::Duration;
20
21use datafusion_common::TableReference;
22use datafusion_common::instant::Instant;
23use datafusion_common::{HashMap, Result};
24
25use crate::cache::lru_queue::LruQueue;
26use crate::cache::{Cache, CacheEntryInfo, CacheKey, CacheValue};
27
28/// Source of the current time used by a [`DefaultCache`] when applying TTLs.
29pub trait TimeProvider: Send + Sync {
30    /// Return the current instant.
31    fn now(&self) -> Instant;
32}
33
34/// [`TimeProvider`] backed by [`Instant::now`].
35///
36/// This is the default time source used by [`DefaultCache`]
37#[derive(Debug, Default)]
38pub struct SystemTimeProvider;
39
40impl TimeProvider for SystemTimeProvider {
41    fn now(&self) -> Instant {
42        Instant::now()
43    }
44}
45
46#[derive(Clone)]
47struct ValueEntry<V: CacheValue> {
48    value: V,
49    expires: Option<Instant>,
50}
51
52struct DefaultCacheState<K: CacheKey, V: CacheValue> {
53    lru_queue: LruQueue<K, ValueEntry<V>>,
54    hits: HashMap<K, usize>,
55    memory_limit: usize,
56    memory_used: usize,
57    ttl: Option<Duration>,
58}
59
60impl<K: CacheKey, V: CacheValue> DefaultCacheState<K, V> {
61    fn new(memory_limit: usize, ttl: Option<Duration>) -> Self {
62        Self {
63            lru_queue: LruQueue::new(),
64            hits: HashMap::new(),
65            memory_limit,
66            memory_used: 0,
67            ttl,
68        }
69    }
70
71    fn get(&mut self, key: &K, now: Instant) -> Option<V> {
72        let entry = self.lru_queue.get(key)?;
73        if let Some(exp) = entry.expires
74            && now > exp
75        {
76            self.remove(key);
77            return None;
78        }
79        let value = entry.value.clone();
80        *self.hits.entry(key.clone()).or_insert(0) += 1;
81        Some(value)
82    }
83
84    fn contains_key(&mut self, key: &K, now: Instant) -> bool {
85        let Some(entry) = self.lru_queue.peek(key) else {
86            return false;
87        };
88        match entry.expires {
89            Some(exp) if now > exp => {
90                self.remove(key);
91                false
92            }
93            _ => true,
94        }
95    }
96
97    fn put(&mut self, key: &K, value: V, now: Instant) -> Option<V> {
98        let value_size = value.size();
99
100        if value_size == 0 {
101            return None;
102        }
103
104        let key_size = key.size();
105        let total_size = key_size + value_size;
106
107        if total_size > self.memory_limit {
108            // Remove potential stale entry
109            return self.remove(key);
110        }
111
112        let expires = self.ttl.map(|ttl| now + ttl);
113        let entry = ValueEntry { value, expires };
114
115        self.memory_used += total_size;
116        self.hits.insert(key.clone(), 0);
117        let old = self.lru_queue.put(key.clone(), entry);
118        if let Some(old_entry) = &old {
119            self.memory_used -= key_size;
120            self.memory_used -= old_entry.value.size();
121        }
122
123        self.evict_entries();
124
125        old.map(|v| v.value)
126    }
127
128    fn remove(&mut self, key: &K) -> Option<V> {
129        let entry = self.lru_queue.remove(key)?;
130        self.memory_used -= key.size();
131        self.memory_used -= entry.value.size();
132        self.hits.remove(key);
133        Some(entry.value)
134    }
135
136    fn evict_entries(&mut self) {
137        while self.memory_used > self.memory_limit {
138            let Some((evicted_key, evicted)) = self.lru_queue.pop() else {
139                // cache is empty while memory_used > memory_limit, cannot happen
140                log::error!(
141                    "DefaultCache memory accounting bug: memory_used={} but cache is empty",
142                    self.memory_used
143                );
144                debug_assert!(false, "memory_used > limit with empty cache");
145                self.memory_used = 0;
146                return;
147            };
148            self.memory_used -= evicted_key.size();
149            self.memory_used -= evicted.value.size();
150            self.hits.remove(&evicted_key);
151        }
152    }
153
154    fn clear(&mut self) {
155        self.lru_queue.clear();
156        self.hits.clear();
157        self.memory_used = 0;
158    }
159}
160
161/// In-memory [`Cache`] with an LRU eviction policy, byte-based memory limit,
162/// and optional per-entry TTL.
163///
164/// Entries are evicted in least-recently-used order whenever an insert would
165/// push `memory_used` above `memory_limit`. Inserts whose own size exceeds the
166/// limit are rejected (and any prior entry under the same key is removed).
167/// When a TTL is configured, the expiration is stamped onto each entry at
168/// insertion time and checked lazily on access. Entries with size 0 are rejected.
169pub struct DefaultCache<K: CacheKey, V: CacheValue> {
170    state: Mutex<DefaultCacheState<K, V>>,
171    time_provider: Arc<dyn TimeProvider>,
172    name: String,
173}
174
175impl<K: CacheKey, V: CacheValue> DefaultCache<K, V> {
176    /// Create a cache with the given memory budget in bytes and no TTL.
177    pub fn new(memory_limit: usize) -> Self {
178        Self::new_with_ttl(memory_limit, None)
179    }
180
181    /// Create a cache with the given memory budget in bytes and an optional
182    /// TTL applied to every newly inserted entry.
183    pub fn new_with_ttl(memory_limit: usize, ttl: Option<Duration>) -> Self {
184        Self {
185            state: Mutex::new(DefaultCacheState::new(memory_limit, ttl)),
186            time_provider: Arc::new(SystemTimeProvider),
187            name: "DefaultCache".to_string(),
188        }
189    }
190
191    /// Override the cache name.
192    pub fn with_name(mut self, name: impl Into<String>) -> Self {
193        self.name = name.into();
194        self
195    }
196
197    /// Override the time source used to stamp and check TTLs.
198    pub fn with_time_provider(mut self, provider: Arc<dyn TimeProvider>) -> Self {
199        self.time_provider = provider;
200        self
201    }
202
203    /// Number of bytes currently accounted for by live entries.
204    pub fn memory_used(&self) -> usize {
205        self.state.lock().unwrap().memory_used
206    }
207}
208
209impl<K: CacheKey, V: CacheValue> Cache<K, V> for DefaultCache<K, V> {
210    fn get(&self, key: &K) -> Option<V> {
211        let now = self.time_provider.now();
212        let mut state = self.state.lock().unwrap();
213        state.get(key, now)
214    }
215
216    fn put(&self, key: &K, value: V) -> Option<V> {
217        let now = self.time_provider.now();
218        let mut state = self.state.lock().unwrap();
219        state.put(key, value, now)
220    }
221
222    fn remove(&self, k: &K) -> Option<V> {
223        let mut state = self.state.lock().unwrap();
224        state.remove(k)
225    }
226
227    fn contains_key(&self, k: &K) -> bool {
228        let now = self.time_provider.now();
229        let mut state = self.state.lock().unwrap();
230        state.contains_key(k, now)
231    }
232
233    fn len(&self) -> usize {
234        self.state.lock().unwrap().lru_queue.len()
235    }
236
237    fn clear(&self) {
238        let mut state = self.state.lock().unwrap();
239        state.clear();
240    }
241
242    fn name(&self) -> String {
243        self.name.clone()
244    }
245    fn cache_limit(&self) -> usize {
246        self.state.lock().unwrap().memory_limit
247    }
248
249    fn update_cache_limit(&self, limit: usize) {
250        let mut state = self.state.lock().unwrap();
251        state.memory_limit = limit;
252        state.evict_entries();
253    }
254
255    fn cache_ttl(&self) -> Option<Duration> {
256        self.state.lock().unwrap().ttl
257    }
258
259    fn update_cache_ttl(&self, ttl: Option<Duration>) {
260        let mut state = self.state.lock().unwrap();
261        state.ttl = ttl;
262    }
263
264    fn drop_table_entries(&self, table_ref: &TableReference) -> Result<()> {
265        let mut state = self.state.lock().unwrap();
266        let to_remove: Vec<K> = state
267            .lru_queue
268            .keys()
269            .filter(|k| k.table_ref() == Some(table_ref))
270            .cloned()
271            .collect();
272        for k in &to_remove {
273            state.remove(k);
274        }
275        Ok(())
276    }
277
278    fn list_entries(&self) -> HashMap<K, CacheEntryInfo<V>> {
279        let state = self.state.lock().unwrap();
280        state
281            .lru_queue
282            .list_entries()
283            .into_iter()
284            .map(|(k, entry)| {
285                let hits = state.hits.get(k).copied().unwrap_or(0);
286                let info = CacheEntryInfo {
287                    value: entry.value.clone(),
288                    size_bytes: entry.value.size(),
289                    hits,
290                    expires: entry.expires,
291                };
292                (k.clone(), info)
293            })
294            .collect()
295    }
296}
297
298#[cfg(test)]
299mod tests {
300    use std::sync::Arc;
301
302    use crate::cache::cache_manager::{
303        CachedFileList, DEFAULT_LIST_FILES_CACHE_MEMORY_LIMIT, meta_heap_bytes,
304    };
305    use crate::cache::cache_manager::{
306        CachedFileMetadata, DEFAULT_FILE_STATISTICS_MEMORY_LIMIT,
307    };
308    use crate::cache::cache_manager::{CachedFileMetadataEntry, FileMetadata};
309    use crate::cache::default_cache::DefaultCache;
310    use crate::cache::default_cache::TimeProvider;
311    use crate::cache::{Cache, CacheEntryInfo};
312    use crate::cache::{CacheKey, CacheValue};
313    use crate::cache::{SchemaFingerprint, TableScopedPath};
314    use arrow::array::{Int32Array, ListArray, RecordBatch};
315    use arrow::buffer::{OffsetBuffer, ScalarBuffer};
316    use arrow::datatypes::{DataType, Field, Schema, TimeUnit};
317    use chrono::DateTime;
318    use datafusion_common::HashMap;
319    use datafusion_common::TableReference;
320    use datafusion_common::heap_size::{DFHeapSize, DFHeapSizeCtx};
321    use datafusion_common::instant::Instant;
322    use datafusion_common::stats::Precision;
323    use datafusion_common::{ColumnStatistics, ScalarValue, Statistics};
324    use datafusion_expr::ColumnarValue;
325    use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
326    use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr};
327    use object_store::ObjectMeta;
328    use object_store::path::Path;
329    use std::sync::Mutex;
330    use std::thread;
331    use std::time::Duration;
332
333    pub struct TestFileMetadata {
334        metadata: String,
335    }
336
337    impl FileMetadata for TestFileMetadata {
338        fn as_any(&self) -> &dyn std::any::Any {
339            self
340        }
341
342        fn memory_size(&self) -> usize {
343            self.metadata.len()
344        }
345
346        fn extra_info(&self) -> HashMap<String, String> {
347            HashMap::from([("extra_info".to_owned(), "abc".to_owned())])
348        }
349    }
350
351    impl PartialEq for CachedFileMetadataEntry {
352        fn eq(&self, other: &Self) -> bool {
353            self.meta == other.meta
354        }
355    }
356
357    fn create_test_object_meta(path: &str, size: usize) -> ObjectMeta {
358        ObjectMeta {
359            location: Path::from(path),
360            last_modified: DateTime::parse_from_rfc3339("2025-07-29T12:12:12+00:00")
361                .unwrap()
362                .into(),
363            size: size as u64,
364            e_tag: None,
365            version: None,
366        }
367    }
368
369    #[test]
370    fn test_default_file_metadata_cache() {
371        let object_meta = create_test_object_meta("test", 1024);
372
373        let metadata: Arc<dyn FileMetadata> = Arc::new(TestFileMetadata {
374            metadata: "retrieved_metadata".to_owned(),
375        });
376
377        let cache = DefaultCache::new(1024 * 1024);
378
379        // Cache miss
380        assert!(cache.get(&object_meta.location).is_none());
381
382        // Put a value
383        let cached_entry =
384            CachedFileMetadataEntry::new(object_meta.clone(), Arc::clone(&metadata));
385        cache.put(&object_meta.location, cached_entry);
386
387        // Verify the cached value
388        assert!(cache.contains_key(&object_meta.location));
389        let result = cache.get(&object_meta.location).unwrap();
390        let test_file_metadata = Arc::downcast::<TestFileMetadata>(result.file_metadata);
391        assert!(test_file_metadata.is_ok());
392        assert_eq!(test_file_metadata.unwrap().metadata, "retrieved_metadata");
393
394        // Cache hit - check validation
395        let result2 = cache.get(&object_meta.location).unwrap();
396        assert!(result2.is_valid_for(&object_meta));
397
398        // File size changed - closure should detect invalidity
399        let object_meta2 = create_test_object_meta("test", 2048);
400        let result3 = cache.get(&object_meta2.location).unwrap();
401        // Cached entry should NOT be valid for new meta
402        assert!(!result3.is_valid_for(&object_meta2));
403
404        // Return new entry
405        let new_entry =
406            CachedFileMetadataEntry::new(object_meta2.clone(), Arc::clone(&metadata));
407        cache.put(&object_meta2.location, new_entry);
408
409        let result4 = cache.get(&object_meta2.location).unwrap();
410        assert_eq!(result4.meta.size, 2048);
411
412        // remove
413        cache.remove(&object_meta.location);
414        assert!(!cache.contains_key(&object_meta.location));
415
416        // len and clear
417        let object_meta3 = create_test_object_meta("test3", 100);
418        cache.put(
419            &object_meta.location,
420            CachedFileMetadataEntry::new(object_meta.clone(), Arc::clone(&metadata)),
421        );
422        cache.put(
423            &object_meta3.location,
424            CachedFileMetadataEntry::new(object_meta3.clone(), Arc::clone(&metadata)),
425        );
426        assert_eq!(cache.len(), 2);
427        cache.clear();
428        assert_eq!(cache.len(), 0);
429    }
430
431    fn generate_test_metadata_with_size(
432        path: &str,
433        size: usize,
434    ) -> (ObjectMeta, Arc<dyn FileMetadata>) {
435        let object_meta = ObjectMeta {
436            location: Path::from(path),
437            last_modified: chrono::Utc::now(),
438            size: size as u64,
439            e_tag: None,
440            version: None,
441        };
442        let metadata = "a".repeat(size);
443        let metadata: Arc<dyn FileMetadata> = Arc::new(TestFileMetadata { metadata });
444
445        (object_meta, metadata)
446    }
447
448    #[test]
449    fn test_default_file_metadata_cache_with_limit() {
450        // Create a cache with 1000 bytes capacity + 4 keys each key 2 bytes
451        let cache = DefaultCache::new(1000 + 4 * 2);
452
453        let (object_meta1, metadata1) = generate_test_metadata_with_size("01", 100);
454        let (object_meta2, metadata2) = generate_test_metadata_with_size("02", 500);
455        let (object_meta3, metadata3) = generate_test_metadata_with_size("03", 300);
456
457        cache.put(
458            &object_meta1.location,
459            CachedFileMetadataEntry::new(object_meta1.clone(), metadata1),
460        );
461        cache.put(
462            &object_meta2.location,
463            CachedFileMetadataEntry::new(object_meta2.clone(), metadata2),
464        );
465        cache.put(
466            &object_meta3.location,
467            CachedFileMetadataEntry::new(object_meta3.clone(), metadata3),
468        );
469
470        // all entries will fit
471        assert_eq!(cache.len(), 3);
472        assert_eq!(cache.memory_used(), 906);
473        assert!(cache.contains_key(&object_meta1.location));
474        assert!(cache.contains_key(&object_meta2.location));
475        assert!(cache.contains_key(&object_meta3.location));
476
477        // add a new entry which will remove the least recently used ("1")
478        let (object_meta4, metadata4) = generate_test_metadata_with_size("04", 200);
479        cache.put(
480            &object_meta4.location,
481            CachedFileMetadataEntry::new(object_meta4.clone(), metadata4),
482        );
483        assert_eq!(cache.len(), 3);
484        assert_eq!(cache.memory_used(), 1006);
485        assert!(!cache.contains_key(&object_meta1.location));
486        assert!(cache.contains_key(&object_meta4.location));
487
488        // get entry "2", which will move it to the top of the queue, and add a new one which will
489        // remove the new least recently used ("3")
490        let _ = cache.get(&object_meta2.location);
491        let (object_meta5, metadata5) = generate_test_metadata_with_size("05", 100);
492        cache.put(
493            &object_meta5.location,
494            CachedFileMetadataEntry::new(object_meta5.clone(), metadata5),
495        );
496        assert_eq!(cache.len(), 3);
497        assert_eq!(cache.memory_used(), 806);
498        assert!(!cache.contains_key(&object_meta3.location));
499        assert!(cache.contains_key(&object_meta5.location));
500
501        // new entry which will not be able to fit in the 1000 bytes allocated
502        let (object_meta6, metadata6) = generate_test_metadata_with_size("06", 1200);
503        cache.put(
504            &object_meta6.location,
505            CachedFileMetadataEntry::new(object_meta6.clone(), metadata6),
506        );
507        assert_eq!(cache.len(), 3);
508        assert_eq!(cache.memory_used(), 806);
509        assert!(!cache.contains_key(&object_meta6.location));
510
511        // new entry which is able to fit without removing any entry
512        let (object_meta7, metadata7) = generate_test_metadata_with_size("07", 200);
513        cache.put(
514            &object_meta7.location,
515            CachedFileMetadataEntry::new(object_meta7.clone(), metadata7),
516        );
517        assert_eq!(cache.len(), 4);
518        assert_eq!(cache.memory_used(), 1008);
519        assert!(cache.contains_key(&object_meta7.location));
520
521        // new entry which will remove all other entries
522        let (object_meta8, metadata8) = generate_test_metadata_with_size("08", 999);
523        cache.put(
524            &object_meta8.location,
525            CachedFileMetadataEntry::new(object_meta8.clone(), metadata8),
526        );
527        assert_eq!(cache.len(), 1);
528        assert_eq!(cache.memory_used(), 1001);
529        assert!(cache.contains_key(&object_meta8.location));
530
531        // when updating an entry, the previous ones are not unnecessarily removed
532        let (object_meta9, metadata9) = generate_test_metadata_with_size("09", 300);
533        let (object_meta10, metadata10) = generate_test_metadata_with_size("10", 200);
534        let (object_meta11_v1, metadata11_v1) =
535            generate_test_metadata_with_size("11", 400);
536        cache.put(
537            &object_meta9.location,
538            CachedFileMetadataEntry::new(object_meta9.clone(), metadata9),
539        );
540        cache.put(
541            &object_meta10.location,
542            CachedFileMetadataEntry::new(object_meta10.clone(), metadata10),
543        );
544        cache.put(
545            &object_meta11_v1.location,
546            CachedFileMetadataEntry::new(object_meta11_v1.clone(), metadata11_v1),
547        );
548        assert_eq!(cache.memory_used(), 906);
549        assert_eq!(cache.len(), 3);
550        let (object_meta11_v2, metadata11_v2) =
551            generate_test_metadata_with_size("11", 500);
552        cache.put(
553            &object_meta11_v2.location,
554            CachedFileMetadataEntry::new(object_meta11_v2.clone(), metadata11_v2),
555        );
556        assert_eq!(cache.memory_used(), 1006);
557        assert_eq!(cache.len(), 3);
558        assert!(cache.contains_key(&object_meta9.location));
559        assert!(cache.contains_key(&object_meta10.location));
560        assert!(cache.contains_key(&object_meta11_v2.location));
561
562        // when updating an entry that now exceeds the limit, the LRU ("09") needs to be removed
563        let (object_meta11_v3, metadata11_v3) =
564            generate_test_metadata_with_size("11", 510);
565        cache.put(
566            &object_meta11_v3.location,
567            CachedFileMetadataEntry::new(object_meta11_v3.clone(), metadata11_v3),
568        );
569        assert_eq!(cache.memory_used(), 714);
570        assert_eq!(cache.len(), 2);
571        assert!(cache.contains_key(&object_meta10.location));
572        assert!(cache.contains_key(&object_meta11_v3.location));
573
574        // manually removing an entry that is not the LRU
575        cache.remove(&object_meta11_v3.location);
576        assert_eq!(cache.len(), 1);
577        assert_eq!(cache.memory_used(), 202);
578        assert!(cache.contains_key(&object_meta10.location));
579        assert!(!cache.contains_key(&object_meta11_v3.location));
580
581        // clear
582        cache.clear();
583        assert_eq!(cache.len(), 0);
584        assert_eq!(cache.memory_used(), 0);
585
586        // resizing the cache should clear the extra entries
587        let (object_meta12, metadata12) = generate_test_metadata_with_size("12", 300);
588        let (object_meta13, metadata13) = generate_test_metadata_with_size("13", 200);
589        let (object_meta14, metadata14) = generate_test_metadata_with_size("14", 500);
590        cache.put(
591            &object_meta12.location,
592            CachedFileMetadataEntry::new(object_meta12.clone(), metadata12),
593        );
594        cache.put(
595            &object_meta13.location,
596            CachedFileMetadataEntry::new(object_meta13.clone(), metadata13),
597        );
598        cache.put(
599            &object_meta14.location,
600            CachedFileMetadataEntry::new(object_meta14.clone(), metadata14),
601        );
602        assert_eq!(cache.len(), 3);
603        assert_eq!(cache.memory_used(), 1006);
604        cache.update_cache_limit(600);
605        assert_eq!(cache.len(), 1);
606        assert_eq!(cache.memory_used(), 502);
607        assert!(!cache.contains_key(&object_meta12.location));
608        assert!(!cache.contains_key(&object_meta13.location));
609        assert!(cache.contains_key(&object_meta14.location));
610    }
611
612    #[test]
613    fn test_default_file_metadata_cache_entries_info() {
614        // Create a cache with 1000 bytes + 4 bytes for 4 keys each key 1 byte
615        let cache = DefaultCache::new(1000 + 4);
616
617        let (object_meta1, metadata1) = generate_test_metadata_with_size("1", 100);
618        let (object_meta2, metadata2) = generate_test_metadata_with_size("2", 200);
619        let (object_meta3, metadata3) = generate_test_metadata_with_size("3", 300);
620
621        // initial entries, all will have hits = 0
622        let entry_1 = CachedFileMetadataEntry::new(object_meta1.clone(), metadata1);
623        let entry_2 = CachedFileMetadataEntry::new(object_meta2.clone(), metadata2);
624        let entry_3 = CachedFileMetadataEntry::new(object_meta3.clone(), metadata3);
625
626        // Build a cache which fits exactly these 3 entries
627
628        cache.put(&object_meta1.location, entry_1.clone());
629        cache.put(&object_meta2.location, entry_2.clone());
630        cache.put(&object_meta3.location, entry_3.clone());
631        let entries = cache.list_entries();
632
633        assert_eq!(
634            entries,
635            HashMap::from([
636                (
637                    Path::from("1"),
638                    CacheEntryInfo {
639                        value: entry_1.clone(),
640                        size_bytes: 100,
641                        hits: 0,
642                        expires: None,
643                    }
644                ),
645                (
646                    Path::from("2"),
647                    CacheEntryInfo {
648                        value: entry_2.clone(),
649                        size_bytes: 200,
650                        hits: 0,
651                        expires: None,
652                    }
653                ),
654                (
655                    Path::from("3"),
656                    CacheEntryInfo {
657                        value: entry_3.clone(),
658                        size_bytes: 300,
659                        hits: 0,
660                        expires: None,
661                    }
662                )
663            ])
664        );
665
666        // new hit on "1"
667        let _ = cache.get(&object_meta1.location);
668        assert_eq!(
669            cache.list_entries(),
670            HashMap::from([
671                (
672                    Path::from("1"),
673                    CacheEntryInfo {
674                        value: entry_1.clone(),
675                        size_bytes: 100,
676                        hits: 1,
677                        expires: None,
678                    }
679                ),
680                (
681                    Path::from("2"),
682                    CacheEntryInfo {
683                        value: entry_2.clone(),
684                        size_bytes: 200,
685                        hits: 0,
686                        expires: None,
687                    }
688                ),
689                (
690                    Path::from("3"),
691                    CacheEntryInfo {
692                        value: entry_3.clone(),
693                        size_bytes: 300,
694                        hits: 0,
695                        expires: None,
696                    }
697                )
698            ])
699        );
700
701        // new entry, will evict "2"
702        let (object_meta4, metadata4) = generate_test_metadata_with_size("4", 600);
703        let entry_4 = CachedFileMetadataEntry::new(object_meta4.clone(), metadata4);
704        cache.put(&object_meta4.location, entry_4.clone());
705        assert_eq!(
706            cache.list_entries(),
707            HashMap::from([
708                (
709                    Path::from("1"),
710                    CacheEntryInfo {
711                        value: entry_1.clone(),
712                        size_bytes: 100,
713                        hits: 1,
714                        expires: None,
715                    }
716                ),
717                (
718                    Path::from("3"),
719                    CacheEntryInfo {
720                        value: entry_3.clone(),
721                        size_bytes: 300,
722                        hits: 0,
723                        expires: None,
724                    }
725                ),
726                (
727                    Path::from("4"),
728                    CacheEntryInfo {
729                        value: entry_4.clone(),
730                        size_bytes: 600,
731                        hits: 0,
732                        expires: None,
733                    }
734                )
735            ])
736        );
737
738        // replace entry "1"
739        let (object_meta1_new, metadata1_new) = generate_test_metadata_with_size("1", 50);
740        let entry_1 =
741            CachedFileMetadataEntry::new(object_meta1_new.clone(), metadata1_new);
742        cache.put(&object_meta1_new.location, entry_1.clone());
743        assert_eq!(
744            cache.list_entries(),
745            HashMap::from([
746                (
747                    Path::from("1"),
748                    CacheEntryInfo {
749                        value: entry_1.clone(),
750                        size_bytes: 50,
751                        hits: 0,
752                        expires: None,
753                    }
754                ),
755                (
756                    Path::from("3"),
757                    CacheEntryInfo {
758                        value: entry_3.clone(),
759                        size_bytes: 300,
760                        hits: 0,
761                        expires: None,
762                    }
763                ),
764                (
765                    Path::from("4"),
766                    CacheEntryInfo {
767                        value: entry_4.clone(),
768                        size_bytes: 600,
769                        hits: 0,
770                        expires: None,
771                    }
772                )
773            ])
774        );
775
776        // remove entry "4"
777        cache.remove(&object_meta4.location);
778        assert_eq!(
779            cache.list_entries(),
780            HashMap::from([
781                (
782                    Path::from("1"),
783                    CacheEntryInfo {
784                        value: entry_1.clone(),
785                        size_bytes: 50,
786                        hits: 0,
787                        expires: None,
788                    }
789                ),
790                (
791                    Path::from("3"),
792                    CacheEntryInfo {
793                        value: entry_3.clone(),
794                        size_bytes: 300,
795                        hits: 0,
796                        expires: None,
797                    }
798                )
799            ])
800        );
801
802        // clear
803        cache.clear();
804        assert_eq!(cache.list_entries(), HashMap::from([]));
805    }
806
807    fn create_test_meta(path: &str, size: u64) -> ObjectMeta {
808        ObjectMeta {
809            location: Path::from(path),
810            last_modified: DateTime::parse_from_rfc3339("2022-09-27T22:36:00+02:00")
811                .unwrap()
812                .into(),
813            size,
814            e_tag: None,
815            version: None,
816        }
817    }
818
819    #[test]
820    fn test_statistics_cache() {
821        let meta = create_test_meta("test", 1024);
822        let cache = DefaultCache::new(DEFAULT_FILE_STATISTICS_MEMORY_LIMIT);
823
824        let schema = Schema::new(vec![Field::new(
825            "test_column",
826            DataType::Timestamp(TimeUnit::Second, None),
827            false,
828        )]);
829
830        let path = TableScopedPath {
831            path: meta.location.clone(),
832            table: None,
833        };
834        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
835
836        // Cache miss
837        assert!(cache.get(&path).is_none());
838
839        // Put a value
840        let cached_value = CachedFileMetadata::new(
841            meta.clone(),
842            Arc::clone(&schema_fingerprint),
843            Arc::new(Statistics::new_unknown(&schema)),
844            None,
845        );
846        cache.put(&path, cached_value);
847
848        // Cache hit
849        let result = cache.get(&path);
850        assert!(result.is_some());
851
852        let cached = result.unwrap();
853        assert!(cached.is_valid_for(&meta, &schema_fingerprint));
854
855        let equivalent_schema_fingerprint =
856            Arc::new(SchemaFingerprint::from_schema(&schema));
857        assert!(!Arc::ptr_eq(
858            &schema_fingerprint,
859            &equivalent_schema_fingerprint
860        ));
861        assert!(cached.is_valid_for(&meta, &equivalent_schema_fingerprint));
862
863        let different_schema = Schema::new(vec![Field::new(
864            "different_column",
865            DataType::Timestamp(TimeUnit::Second, None),
866            false,
867        )]);
868        let different_schema_fingerprint =
869            Arc::new(SchemaFingerprint::from_schema(&different_schema));
870        assert!(!cached.is_valid_for(&meta, &different_schema_fingerprint));
871
872        // File size changed - validation should fail
873        let meta2 = create_test_meta("test", 2048);
874
875        let cached = cache.get(&path).unwrap();
876        assert!(!cached.is_valid_for(&meta2, &schema_fingerprint));
877
878        // Update with new value
879        let cached_value2 = CachedFileMetadata::new(
880            meta2.clone(),
881            Arc::clone(&schema_fingerprint),
882            Arc::new(Statistics::new_unknown(&schema)),
883            None,
884        );
885        cache.put(&path, cached_value2);
886
887        // Test list_entries
888        let entries = cache.list_entries();
889        assert_eq!(entries.len(), 1);
890
891        let path_3 = TableScopedPath {
892            path: Path::from("test"),
893            table: None,
894        };
895
896        let entry = entries.get(&path_3).unwrap();
897        assert_eq!(entry.value.meta.size, 2048); // Should be updated value
898    }
899
900    #[derive(Clone, Debug, PartialEq, Eq, Hash)]
901    struct MockExpr {}
902
903    impl std::fmt::Display for MockExpr {
904        fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
905            write!(f, "MockExpr")
906        }
907    }
908
909    impl PhysicalExpr for MockExpr {
910        fn data_type(
911            &self,
912            _input_schema: &Schema,
913        ) -> datafusion_common::Result<DataType> {
914            Ok(DataType::Int32)
915        }
916
917        fn nullable(&self, _input_schema: &Schema) -> datafusion_common::Result<bool> {
918            Ok(false)
919        }
920
921        fn evaluate(
922            &self,
923            _batch: &RecordBatch,
924        ) -> datafusion_common::Result<ColumnarValue> {
925            unimplemented!()
926        }
927
928        fn children(&self) -> Vec<&Arc<dyn PhysicalExpr>> {
929            vec![]
930        }
931
932        fn with_new_children(
933            self: Arc<Self>,
934            children: Vec<Arc<dyn PhysicalExpr>>,
935        ) -> datafusion_common::Result<Arc<dyn PhysicalExpr>> {
936            assert!(children.is_empty());
937            Ok(self)
938        }
939
940        fn fmt_sql(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
941            write!(f, "MockExpr")
942        }
943    }
944
945    fn ordering() -> LexOrdering {
946        let expr = Arc::new(MockExpr {}) as Arc<dyn PhysicalExpr>;
947        LexOrdering::new(vec![PhysicalSortExpr::new_default(expr)]).unwrap()
948    }
949
950    #[test]
951    fn test_ordering_cache() {
952        let meta = create_test_meta("test.parquet", 100);
953        let cache = DefaultCache::new(DEFAULT_FILE_STATISTICS_MEMORY_LIMIT);
954
955        let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]);
956        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
957
958        // Cache statistics with no ordering
959        let cached_value = CachedFileMetadata::new(
960            meta.clone(),
961            Arc::clone(&schema_fingerprint),
962            Arc::new(Statistics::new_unknown(&schema)),
963            None, // No ordering yet
964        );
965
966        let path = TableScopedPath {
967            path: meta.location.clone(),
968            table: None,
969        };
970
971        cache.put(&path, cached_value);
972
973        let result = cache.get(&path).unwrap();
974        assert!(result.ordering.is_none());
975
976        // Update to add ordering
977        let mut cached = cache.get(&path).unwrap();
978        if cached.is_valid_for(&meta, &schema_fingerprint) && cached.ordering.is_none() {
979            cached.ordering = Some(ordering());
980        }
981        cache.put(&path, cached);
982
983        let result2 = cache.get(&path).unwrap();
984        assert!(result2.ordering.is_some());
985
986        // Verify list_entries shows has_ordering = true
987        let entries = cache.list_entries();
988        assert_eq!(entries.len(), 1);
989        assert!(entries.get(&path).unwrap().value.ordering.is_some());
990    }
991
992    #[test]
993    fn test_cache_invalidation_on_file_modification() {
994        let cache = DefaultCache::new(DEFAULT_FILE_STATISTICS_MEMORY_LIMIT);
995        let path = TableScopedPath {
996            path: Path::from("test.parquet"),
997            table: None,
998        };
999        let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]);
1000        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
1001
1002        let meta_v1 = create_test_meta("test.parquet", 100);
1003
1004        // Cache initial value
1005        let cached_value = CachedFileMetadata::new(
1006            meta_v1.clone(),
1007            Arc::clone(&schema_fingerprint),
1008            Arc::new(Statistics::new_unknown(&schema)),
1009            None,
1010        );
1011        cache.put(&path, cached_value);
1012
1013        // File modified (size changed)
1014        let meta_v2 = create_test_meta("test.parquet", 200);
1015
1016        let cached = cache.get(&path).unwrap();
1017        // Should not be valid for new meta
1018        assert!(!cached.is_valid_for(&meta_v2, &schema_fingerprint));
1019
1020        // Compute new value and update
1021        let new_cached = CachedFileMetadata::new(
1022            meta_v2.clone(),
1023            Arc::clone(&schema_fingerprint),
1024            Arc::new(Statistics::new_unknown(&schema)),
1025            None,
1026        );
1027        cache.put(&path, new_cached);
1028
1029        // Should have new metadata
1030        let result = cache.get(&path).unwrap();
1031        assert_eq!(result.meta.size, 200);
1032    }
1033
1034    #[test]
1035    fn test_ordering_cache_invalidation_on_file_modification() {
1036        let cache = DefaultCache::new(DEFAULT_FILE_STATISTICS_MEMORY_LIMIT);
1037        let path = TableScopedPath {
1038            path: Path::from("test.parquet"),
1039            table: None,
1040        };
1041        let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]);
1042        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
1043
1044        // Cache with original metadata and ordering
1045        let meta_v1 = ObjectMeta {
1046            location: path.path.clone(),
1047            last_modified: DateTime::parse_from_rfc3339("2022-09-27T22:36:00+02:00")
1048                .unwrap()
1049                .into(),
1050            size: 100,
1051            e_tag: None,
1052            version: None,
1053        };
1054        let ordering_v1 = ordering();
1055        let cached_v1 = CachedFileMetadata::new(
1056            meta_v1.clone(),
1057            Arc::clone(&schema_fingerprint),
1058            Arc::new(Statistics::new_unknown(&schema)),
1059            Some(ordering_v1),
1060        );
1061        cache.put(&path, cached_v1);
1062
1063        // Verify cached ordering is valid
1064        let cached = cache.get(&path).unwrap();
1065        assert!(cached.is_valid_for(&meta_v1, &schema_fingerprint));
1066        assert!(cached.ordering.is_some());
1067
1068        // File modified (size changed)
1069        let meta_v2 = ObjectMeta {
1070            location: path.path.clone(),
1071            last_modified: DateTime::parse_from_rfc3339("2022-09-28T10:00:00+02:00")
1072                .unwrap()
1073                .into(),
1074            size: 200, // Changed
1075            e_tag: None,
1076            version: None,
1077        };
1078
1079        // Cache entry exists but should be invalid for new metadata
1080        let cached = cache.get(&path).unwrap();
1081        assert!(!cached.is_valid_for(&meta_v2, &schema_fingerprint));
1082
1083        // Cache new version with different ordering
1084        let ordering_v2 = ordering(); // New ordering instance
1085        let cached_v2 = CachedFileMetadata::new(
1086            meta_v2.clone(),
1087            Arc::clone(&schema_fingerprint),
1088            Arc::new(Statistics::new_unknown(&schema)),
1089            Some(ordering_v2),
1090        );
1091        cache.put(&path, cached_v2);
1092
1093        // Old metadata should be invalid
1094        let cached = cache.get(&path).unwrap();
1095        assert!(!cached.is_valid_for(&meta_v1, &schema_fingerprint));
1096
1097        // New metadata should be valid
1098        assert!(cached.is_valid_for(&meta_v2, &schema_fingerprint));
1099        assert!(cached.ordering.is_some());
1100    }
1101
1102    #[test]
1103    fn test_list_entries() {
1104        let cache = DefaultCache::new(DEFAULT_FILE_STATISTICS_MEMORY_LIMIT);
1105        let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]);
1106        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
1107
1108        let meta1 = create_test_meta("test1.parquet", 100);
1109
1110        let cached_value_1 = CachedFileMetadata::new(
1111            meta1.clone(),
1112            Arc::clone(&schema_fingerprint),
1113            Arc::new(Statistics::new_unknown(&schema)),
1114            None,
1115        );
1116
1117        let path_1 = TableScopedPath {
1118            path: meta1.location.clone(),
1119            table: None,
1120        };
1121
1122        cache.put(&path_1, cached_value_1.clone());
1123        let meta2 = create_test_meta("test2.parquet", 200);
1124        let cached_value_2 = CachedFileMetadata::new(
1125            meta2.clone(),
1126            Arc::clone(&schema_fingerprint),
1127            Arc::new(Statistics::new_unknown(&schema)),
1128            Some(ordering()),
1129        );
1130
1131        let path_2 = TableScopedPath {
1132            path: meta2.location.clone(),
1133            table: None,
1134        };
1135
1136        cache.put(&path_2, cached_value_2.clone());
1137
1138        let entries = cache.list_entries();
1139        assert_eq!(
1140            entries,
1141            HashMap::from([
1142                (
1143                    path_1,
1144                    CacheEntryInfo {
1145                        value: cached_value_1,
1146                        hits: 0,
1147                        size_bytes: 373,
1148                        expires: None,
1149                    }
1150                ),
1151                (
1152                    path_2,
1153                    CacheEntryInfo {
1154                        value: cached_value_2,
1155                        hits: 0,
1156                        size_bytes: 373,
1157                        expires: None,
1158                    }
1159                ),
1160            ])
1161        );
1162    }
1163
1164    #[test]
1165    fn test_cache_entry_added_when_entries_are_within_cache_limit() {
1166        let (meta_1, value_1) =
1167            create_cached_file_metadata_with_stats("test1.parquet", 10);
1168        let (meta_2, value_2) =
1169            create_cached_file_metadata_with_stats("test2.parquet", 10);
1170        let (meta_3, value_3) =
1171            create_cached_file_metadata_with_stats("test3.parquet", 10);
1172
1173        let mut ctx = DFHeapSizeCtx::default();
1174
1175        let limit_for_2_entries = meta_1.location.as_ref().heap_size(&mut ctx)
1176            + value_1.heap_size(&mut ctx)
1177            + meta_2.location.as_ref().heap_size(&mut ctx)
1178            + value_2.heap_size(&mut ctx);
1179
1180        // create a cache with a limit which fits exactly 2 entries
1181        let cache = DefaultCache::new(limit_for_2_entries);
1182        let path_1 = TableScopedPath {
1183            path: meta_1.location.clone(),
1184            table: None,
1185        };
1186
1187        let path_2 = TableScopedPath {
1188            path: meta_2.location.clone(),
1189            table: None,
1190        };
1191
1192        cache.put(&path_1, value_1.clone());
1193        cache.put(&path_2, value_2.clone());
1194
1195        assert_eq!(cache.len(), 2);
1196        assert_eq!(cache.memory_used(), limit_for_2_entries);
1197
1198        let result_1 = cache.get(&path_1);
1199        let result_2 = cache.get(&path_2);
1200        assert_eq!(result_1.unwrap(), value_1);
1201        assert_eq!(result_2.unwrap(), value_2);
1202
1203        let path_3 = TableScopedPath {
1204            path: meta_3.location.clone(),
1205            table: None,
1206        };
1207
1208        // adding the third entry evicts the first entry
1209        cache.put(&path_3, value_3.clone());
1210        assert_eq!(cache.len(), 2);
1211        assert_eq!(cache.memory_used(), limit_for_2_entries);
1212
1213        let result_1 = cache.get(&path_1);
1214        assert!(result_1.is_none());
1215
1216        let result_2 = cache.get(&path_2);
1217        let result_3 = cache.get(&path_3);
1218
1219        assert_eq!(result_2.unwrap(), value_2);
1220        assert_eq!(result_3.unwrap(), value_3);
1221
1222        // add the third entry again, making sure memory usage remains the same
1223        cache.put(&path_3, value_3.clone());
1224        assert_eq!(cache.memory_used(), limit_for_2_entries);
1225        cache.put(&path_3, value_3.clone());
1226        assert_eq!(cache.memory_used(), limit_for_2_entries);
1227
1228        let mut ctx = DFHeapSizeCtx::default();
1229        cache.remove(&path_2);
1230        assert_eq!(cache.len(), 1);
1231        assert_eq!(
1232            cache.memory_used(),
1233            meta_3.location.as_ref().heap_size(&mut ctx) + value_3.heap_size(&mut ctx)
1234        );
1235
1236        cache.clear();
1237        assert_eq!(cache.len(), 0);
1238        assert_eq!(cache.memory_used(), 0);
1239    }
1240
1241    #[test]
1242    fn test_cache_rejects_entry_which_is_too_large() {
1243        let (meta, value_too_large) =
1244            create_cached_file_metadata_with_stats("test1.parquet", 10);
1245        let mut ctx = DFHeapSizeCtx::default();
1246        let limit_less_than_the_entry = value_too_large.clone().heap_size(&mut ctx) - 1;
1247
1248        // create a cache with a size less than the entry
1249        let cache = DefaultCache::new(limit_less_than_the_entry);
1250
1251        let path_1 = TableScopedPath {
1252            path: meta.location.clone(),
1253            table: None,
1254        };
1255
1256        cache.put(&path_1, value_too_large.clone());
1257
1258        assert_eq!(cache.len(), 0);
1259        assert_eq!(cache.memory_used(), 0);
1260
1261        // Test stale entry is removed when oversized entry is added
1262        let (_, value_fits) = create_cached_file_metadata_with_stats("test1.parquet", 7);
1263        cache.put(&path_1, value_fits.clone());
1264
1265        assert_eq!(cache.len(), 1);
1266        assert_eq!(cache.memory_used(), 1514);
1267
1268        // now add an entry which is over the limit and make sure the old stale entry is removed
1269        let stale_entry = cache.put(&path_1, value_too_large.clone());
1270        assert_eq!(stale_entry, Some(value_fits));
1271
1272        assert_eq!(cache.len(), 0);
1273        assert_eq!(cache.memory_used(), 0);
1274    }
1275
1276    fn create_cached_file_metadata_with_stats(
1277        file_name: &str,
1278        series_size: i32,
1279    ) -> (ObjectMeta, CachedFileMetadata) {
1280        let series: Vec<i32> = (0..=series_size).collect();
1281        let values = Int32Array::from(series);
1282        let offsets = OffsetBuffer::new(ScalarBuffer::from(vec![0, series_size + 1]));
1283        let field = Arc::new(Field::new_list_field(DataType::Int32, false));
1284        let list_array = ListArray::new(field, offsets, Arc::new(values), None);
1285
1286        let column_statistics = ColumnStatistics {
1287            null_count: Precision::Exact(1),
1288            max_value: Precision::Exact(ScalarValue::List(Arc::new(list_array.clone()))),
1289            min_value: Precision::Exact(ScalarValue::List(Arc::new(list_array.clone()))),
1290            sum_value: Precision::Exact(ScalarValue::List(Arc::new(list_array.clone()))),
1291            distinct_count: Precision::Exact(10),
1292            byte_size: Precision::Absent,
1293        };
1294
1295        let stats = Statistics {
1296            num_rows: Precision::Exact(100),
1297            total_byte_size: Precision::Exact(100),
1298            column_statistics: vec![column_statistics.clone()],
1299        };
1300        let mut ctx = DFHeapSizeCtx::default();
1301        let object_meta = create_test_meta(file_name, stats.heap_size(&mut ctx) as u64);
1302        let schema = Schema::new(vec![Field::new("list", DataType::Int32, true)]);
1303        let schema_fingerprint = Arc::new(SchemaFingerprint::from_schema(&schema));
1304        let value = CachedFileMetadata::new(
1305            object_meta.clone(),
1306            schema_fingerprint,
1307            Arc::new(stats.clone()),
1308            None,
1309        );
1310        (object_meta, value)
1311    }
1312
1313    struct MockTimeProvider {
1314        base: Instant,
1315        offset: Mutex<Duration>,
1316    }
1317
1318    impl MockTimeProvider {
1319        fn new() -> Self {
1320            Self {
1321                base: Instant::now(),
1322                offset: Mutex::new(Duration::ZERO),
1323            }
1324        }
1325
1326        fn inc(&self, duration: Duration) {
1327            let mut offset = self.offset.lock().unwrap();
1328            *offset += duration;
1329        }
1330    }
1331
1332    impl TimeProvider for MockTimeProvider {
1333        fn now(&self) -> Instant {
1334            self.base + *self.offset.lock().unwrap()
1335        }
1336    }
1337
1338    /// Helper function to create a test ObjectMeta with a specific path and location string size
1339    fn create_object_meta(path: &str, location_size: usize) -> ObjectMeta {
1340        // Create a location string of the desired size by padding with zeros
1341        let location_str = if location_size > path.len() {
1342            format!("{}{}", path, "0".repeat(location_size - path.len()))
1343        } else {
1344            path.to_string()
1345        };
1346
1347        ObjectMeta {
1348            location: Path::from(location_str),
1349            last_modified: DateTime::parse_from_rfc3339("2022-09-27T22:36:00+02:00")
1350                .unwrap()
1351                .into(),
1352            size: 1024,
1353            e_tag: None,
1354            version: None,
1355        }
1356    }
1357
1358    /// Helper function to create a TableScopedPath and a CachedFileList with at least meta_size bytes
1359    fn create_test_list_files_entry(
1360        path: &str,
1361        count: usize,
1362        meta_size: usize,
1363        table: Option<TableReference>,
1364    ) -> (TableScopedPath, CachedFileList) {
1365        let key = TableScopedPath {
1366            table,
1367            path: Path::from(path),
1368        };
1369        let metas: Vec<ObjectMeta> = (0..count)
1370            .map(|i| create_object_meta(&format!("file{i}"), meta_size))
1371            .collect();
1372        let value = CachedFileList::new(metas);
1373        (key, value)
1374    }
1375
1376    #[test]
1377    fn test_basic_operations() {
1378        let cache = DefaultCache::new(DEFAULT_LIST_FILES_CACHE_MEMORY_LIMIT);
1379        let table_ref = Some(TableReference::from("table"));
1380        let path = Path::from("test_path");
1381        let key = TableScopedPath {
1382            table: table_ref.clone(),
1383            path,
1384        };
1385
1386        // Initially cache is empty
1387        assert!(!cache.contains_key(&key));
1388        assert_eq!(cache.len(), 0);
1389
1390        // Cache miss - get returns None
1391        assert!(cache.get(&key).is_none());
1392
1393        // Put a value
1394        let meta = create_test_object_meta("file1", 50);
1395        cache.put(&key, CachedFileList::new(vec![meta]));
1396
1397        // Entry should be cached
1398        assert!(cache.contains_key(&key));
1399        assert_eq!(cache.len(), 1);
1400        let result = cache.get(&key).unwrap();
1401        assert_eq!(result.files.len(), 1);
1402
1403        // Remove the entry
1404        let removed = cache.remove(&key).unwrap();
1405        assert_eq!(removed.files.len(), 1);
1406        assert!(!cache.contains_key(&key));
1407        assert_eq!(cache.len(), 0);
1408
1409        // Put multiple entries
1410        let (key1, value1) =
1411            create_test_list_files_entry("path1", 2, 50, table_ref.clone());
1412        let (key2, value2) = create_test_list_files_entry("path2", 3, 50, table_ref);
1413        cache.put(&key1, value1.clone());
1414        cache.put(&key2, value2.clone());
1415        assert_eq!(cache.len(), 2);
1416
1417        // List cache entries
1418        assert_eq!(
1419            cache.list_entries(),
1420            HashMap::from([
1421                (
1422                    key1.clone(),
1423                    CacheEntryInfo {
1424                        value: value1.clone(),
1425                        size_bytes: value1.size(),
1426                        hits: 0,
1427                        expires: None,
1428                    }
1429                ),
1430                (
1431                    key2.clone(),
1432                    CacheEntryInfo {
1433                        value: value2.clone(),
1434                        size_bytes: value2.size(),
1435                        hits: 0,
1436                        expires: None,
1437                    }
1438                )
1439            ])
1440        );
1441
1442        // Clear all entries
1443        cache.clear();
1444        assert_eq!(cache.len(), 0);
1445        assert!(!cache.contains_key(&key1));
1446        assert!(!cache.contains_key(&key2));
1447    }
1448
1449    #[test]
1450    fn test_lru_eviction_basic() {
1451        let table_ref = Some(TableReference::from("table"));
1452        let (key1, value1) =
1453            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1454        let (key2, value2) =
1455            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1456        let (key3, value3) =
1457            create_test_list_files_entry("path3", 1, 100, table_ref.clone());
1458
1459        let entry_size = key1.size() + value1.size();
1460
1461        // Set cache limit to exactly fit all 3 entries
1462        let cache = DefaultCache::new(entry_size * 3);
1463
1464        // All three entries should fit
1465        cache.put(&key1, value1);
1466        cache.put(&key2, value2);
1467        cache.put(&key3, value3);
1468        assert_eq!(cache.len(), 3);
1469        assert!(cache.contains_key(&key1));
1470        assert!(cache.contains_key(&key2));
1471        assert!(cache.contains_key(&key3));
1472
1473        // Adding a new entry should evict path1 (LRU)
1474        let (key4, value4) = create_test_list_files_entry("path4", 1, 100, table_ref);
1475        cache.put(&key4, value4);
1476
1477        assert_eq!(cache.len(), 3);
1478        assert!(!cache.contains_key(&key1)); // Evicted
1479        assert!(cache.contains_key(&key2));
1480        assert!(cache.contains_key(&key3));
1481        assert!(cache.contains_key(&key4));
1482    }
1483
1484    #[test]
1485    fn test_lru_ordering_after_access() {
1486        let table_ref = Some(TableReference::from("table"));
1487        let (key1, value1) =
1488            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1489        let (key2, value2) =
1490            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1491        let (key3, value3) =
1492            create_test_list_files_entry("path3", 1, 100, table_ref.clone());
1493
1494        // Set cache limit to fit exactly three entries
1495        let cache = DefaultCache::new((key1.size() + value1.size()) * 3);
1496
1497        cache.put(&key1, value1);
1498        cache.put(&key2, value2);
1499        cache.put(&key3, value3);
1500        assert_eq!(cache.len(), 3);
1501
1502        // Access path1 to move it to front (MRU)
1503        // Order is now: path2 (LRU), path3, path1 (MRU)
1504        let _ = cache.get(&key1);
1505
1506        // Adding a new entry should evict path2 (the LRU)
1507        let (key4, value4) = create_test_list_files_entry("path4", 1, 100, table_ref);
1508        cache.put(&key4, value4);
1509
1510        assert_eq!(cache.len(), 3);
1511        assert!(cache.contains_key(&key1)); // Still present (recently accessed)
1512        assert!(!cache.contains_key(&key2)); // Evicted (was LRU)
1513        assert!(cache.contains_key(&key3));
1514        assert!(cache.contains_key(&key4));
1515    }
1516
1517    #[test]
1518    fn test_reject_too_large() {
1519        let table_ref = Some(TableReference::from("table"));
1520        let (key1, value1) =
1521            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1522        let (key2, value2) =
1523            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1524
1525        // Set cache limit to fit both entries
1526        let cache = DefaultCache::new((key1.size() + value1.size()) * 2);
1527
1528        cache.put(&key1, value1);
1529        cache.put(&key2, value2);
1530        assert_eq!(cache.len(), 2);
1531
1532        // Try to add an entry that's too large to fit in the cache
1533        // The entry is not stored (too large)
1534        let (key_large, value_large) =
1535            create_test_list_files_entry("large", 1, 1000, table_ref);
1536        cache.put(&key_large, value_large);
1537
1538        // Large entry should not be added
1539        assert!(!cache.contains_key(&key_large));
1540        assert_eq!(cache.len(), 2);
1541        assert!(cache.contains_key(&key1));
1542        assert!(cache.contains_key(&key2));
1543    }
1544
1545    #[test]
1546    fn test_multiple_evictions() {
1547        let table_ref = Some(TableReference::from("table"));
1548        let (key1, value1) =
1549            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1550        let (key2, value2) =
1551            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1552        let (key3, value3) =
1553            create_test_list_files_entry("path3", 1, 100, table_ref.clone());
1554
1555        let entry_size = key1.size() + value1.size();
1556
1557        // Set cache limit for exactly 3 entries
1558        let cache = DefaultCache::new(entry_size * 3);
1559
1560        cache.put(&key1, value1);
1561        cache.put(&key2, value2);
1562        cache.put(&key3, value3);
1563        assert_eq!(cache.len(), 3);
1564
1565        // Add a large entry that requires evicting 2 entries
1566        let (key_large, value_large) =
1567            create_test_list_files_entry("large", 1, 200, table_ref);
1568        cache.put(&key_large, value_large);
1569
1570        // path1 and path2 should be evicted (both LRU), path3 and path_large remain
1571        assert_eq!(cache.len(), 2);
1572        assert!(!cache.contains_key(&key1)); // Evicted
1573        assert!(!cache.contains_key(&key2)); // Evicted
1574        assert!(cache.contains_key(&key3));
1575        assert!(cache.contains_key(&key_large));
1576    }
1577
1578    #[test]
1579    fn test_cache_limit_resize() {
1580        let table_ref = Some(TableReference::from("table"));
1581        let (key1, value1) =
1582            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1583        let (key2, value2) =
1584            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1585        let (key3, value3) = create_test_list_files_entry("path3", 1, 100, table_ref);
1586
1587        let entry_size = key1.size() + value1.size();
1588
1589        let cache = DefaultCache::new(entry_size * 3);
1590
1591        // Add three entries
1592        cache.put(&key1, value1);
1593        cache.put(&key2, value2);
1594        cache.put(&key3, value3);
1595        assert_eq!(cache.len(), 3);
1596
1597        // Resize cache to only fit one entry
1598        cache.update_cache_limit(entry_size);
1599
1600        // Should keep only the most recent entry (path3, the MRU)
1601        assert_eq!(cache.len(), 1);
1602        assert!(cache.contains_key(&key3));
1603        // Earlier entries (LRU) should be evicted
1604        assert!(!cache.contains_key(&key1));
1605        assert!(!cache.contains_key(&key2));
1606    }
1607
1608    #[test]
1609    fn test_entry_update_with_size_change() {
1610        let table_ref = Some(TableReference::from("table"));
1611        let (key1, value1) =
1612            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1613        let (key2, value2) =
1614            create_test_list_files_entry("path2", 1, 100, table_ref.clone());
1615        let (key3, value3_v1) =
1616            create_test_list_files_entry("path3", 1, 100, table_ref.clone());
1617
1618        let entry_size = key1.size() + value1.size();
1619
1620        let cache = DefaultCache::new(entry_size * 3);
1621
1622        // Add three entries
1623        cache.put(&key1, value1);
1624        cache.put(&key2, value2.clone());
1625        cache.put(&key3, value3_v1);
1626        assert_eq!(cache.len(), 3);
1627
1628        // Update path3 with same size - should not cause eviction
1629        let (_, value3_v2) =
1630            create_test_list_files_entry("path3", 1, 100, table_ref.clone());
1631        cache.put(&key3, value3_v2);
1632
1633        assert_eq!(cache.len(), 3);
1634        assert!(cache.contains_key(&key1));
1635        assert!(cache.contains_key(&key2));
1636        assert!(cache.contains_key(&key3));
1637
1638        // Update path3 with larger size that requires evicting path1 (LRU)
1639        let (_, value3_v3) = create_test_list_files_entry("path3", 1, 200, table_ref);
1640        cache.put(&key3, value3_v3.clone());
1641
1642        assert_eq!(cache.len(), 2);
1643        assert!(!cache.contains_key(&key1)); // Evicted (was LRU)
1644        assert!(cache.contains_key(&key2));
1645        assert!(cache.contains_key(&key3));
1646
1647        // List cache entries
1648        assert_eq!(
1649            cache.list_entries(),
1650            HashMap::from([
1651                (
1652                    key2,
1653                    CacheEntryInfo {
1654                        value: value2.clone(),
1655                        size_bytes: value2.size(),
1656                        hits: 0,
1657                        expires: None,
1658                    }
1659                ),
1660                (
1661                    key3,
1662                    CacheEntryInfo {
1663                        value: value3_v3.clone(),
1664                        size_bytes: value3_v3.size(),
1665                        hits: 0,
1666                        expires: None,
1667                    }
1668                )
1669            ])
1670        );
1671    }
1672
1673    #[test]
1674    fn test_cache_with_ttl() {
1675        let ttl = Duration::from_millis(100);
1676
1677        let mock_time = Arc::new(MockTimeProvider::new());
1678        let cache = DefaultCache::new_with_ttl(10000, Some(ttl))
1679            .with_time_provider(Arc::clone(&mock_time) as Arc<dyn TimeProvider>);
1680
1681        let table_ref = Some(TableReference::from("table"));
1682        let (key1, value1) =
1683            create_test_list_files_entry("path1", 2, 50, table_ref.clone());
1684        let (key2, value2) = create_test_list_files_entry("path2", 2, 50, table_ref);
1685        cache.put(&key1, value1.clone());
1686        cache.put(&key2, value2.clone());
1687
1688        // Entries should be accessible immediately
1689        assert!(cache.get(&key1).is_some());
1690        assert!(cache.get(&key2).is_some());
1691        // List cache entries
1692        assert_eq!(
1693            cache.list_entries(),
1694            HashMap::from([
1695                (
1696                    key1.clone(),
1697                    CacheEntryInfo {
1698                        value: value1.clone(),
1699                        size_bytes: value1.size(),
1700                        hits: 1,
1701                        expires: mock_time.now().checked_add(ttl),
1702                    }
1703                ),
1704                (
1705                    key2.clone(),
1706                    CacheEntryInfo {
1707                        value: value2.clone(),
1708                        size_bytes: value2.size(),
1709                        hits: 1,
1710                        expires: mock_time.now().checked_add(ttl),
1711                    }
1712                )
1713            ])
1714        );
1715        // Wait for TTL to expire
1716        mock_time.inc(Duration::from_millis(150));
1717
1718        // Entries should now return None when observed through contains_key
1719        assert!(!cache.contains_key(&key1));
1720        assert_eq!(cache.len(), 1); // key1 was removed by contains_key()
1721        assert!(!cache.contains_key(&key2));
1722        assert_eq!(cache.len(), 0); // key2 was removed by contains_key()
1723    }
1724
1725    #[test]
1726    fn test_cache_with_ttl_and_lru() {
1727        let ttl = Duration::from_millis(200);
1728
1729        let mock_time = Arc::new(MockTimeProvider::new());
1730        let cache = DefaultCache::new_with_ttl(1100, Some(ttl))
1731            .with_time_provider(Arc::clone(&mock_time) as Arc<dyn TimeProvider>);
1732
1733        let table_ref = Some(TableReference::from("table"));
1734        let (key1, value1) =
1735            create_test_list_files_entry("path1", 1, 400, table_ref.clone());
1736        let (key2, value2) =
1737            create_test_list_files_entry("path2", 1, 400, table_ref.clone());
1738
1739        let (key3, value3) = create_test_list_files_entry("path3", 1, 400, table_ref);
1740        cache.put(&key1, value1);
1741        mock_time.inc(Duration::from_millis(50));
1742        cache.put(&key2, value2);
1743        mock_time.inc(Duration::from_millis(50));
1744
1745        // path3 should evict path1 due to size limit
1746        cache.put(&key3, value3);
1747        assert!(!cache.contains_key(&key1)); // Evicted by LRU
1748        assert!(cache.contains_key(&key2));
1749        assert!(cache.contains_key(&key3));
1750
1751        mock_time.inc(Duration::from_millis(151));
1752
1753        assert!(!cache.contains_key(&key2)); // Expired
1754        assert!(cache.contains_key(&key3)); // Still valid
1755    }
1756
1757    #[test]
1758    fn test_ttl_expiration_in_get() {
1759        let ttl = Duration::from_millis(100);
1760        let cache = DefaultCache::new_with_ttl(10000, Some(ttl));
1761
1762        let table_ref = Some(TableReference::from("table"));
1763        let (key, value) = create_test_list_files_entry("path", 2, 50, table_ref);
1764
1765        // Cache the entry
1766        cache.put(&key, value.clone());
1767
1768        // Entry should be accessible immediately
1769        let result = cache.get(&key);
1770        assert!(result.is_some());
1771        assert_eq!(result.unwrap().files.len(), 2);
1772
1773        // Wait for TTL to expire
1774        thread::sleep(Duration::from_millis(150));
1775
1776        // Get should return None because entry expired
1777        let result2 = cache.get(&key);
1778        assert!(result2.is_none());
1779    }
1780
1781    #[test]
1782    fn test_meta_heap_bytes_calculation() {
1783        // Test with minimal ObjectMeta (no e_tag, no version)
1784        let meta1 = ObjectMeta {
1785            location: Path::from("test"),
1786            last_modified: chrono::Utc::now(),
1787            size: 100,
1788            e_tag: None,
1789            version: None,
1790        };
1791        assert_eq!(meta_heap_bytes(&meta1), 4); // Just the location string "test"
1792
1793        // Test with e_tag
1794        let meta2 = ObjectMeta {
1795            location: Path::from("test"),
1796            last_modified: chrono::Utc::now(),
1797            size: 100,
1798            e_tag: Some("etag123".to_string()),
1799            version: None,
1800        };
1801        assert_eq!(meta_heap_bytes(&meta2), 4 + 7); // location (4) + e_tag (7)
1802
1803        // Test with version
1804        let meta3 = ObjectMeta {
1805            location: Path::from("test"),
1806            last_modified: chrono::Utc::now(),
1807            size: 100,
1808            e_tag: None,
1809            version: Some("v1.0".to_string()),
1810        };
1811        assert_eq!(meta_heap_bytes(&meta3), 4 + 4); // location (4) + version (4)
1812
1813        // Test with both e_tag and version
1814        let meta4 = ObjectMeta {
1815            location: Path::from("test"),
1816            last_modified: chrono::Utc::now(),
1817            size: 100,
1818            e_tag: Some("tag".to_string()),
1819            version: Some("ver".to_string()),
1820        };
1821        assert_eq!(meta_heap_bytes(&meta4), 4 + 3 + 3); // location (4) + e_tag (3) + version (3)
1822    }
1823
1824    #[test]
1825    fn test_memory_tracking() {
1826        let cache = DefaultCache::new(1000);
1827
1828        // Verify cache starts with 0 memory used
1829        {
1830            assert_eq!(cache.memory_used(), 0);
1831        }
1832
1833        // Add entry and verify memory tracking
1834        let table_ref = Some(TableReference::from("table"));
1835        let (key1, value1) =
1836            create_test_list_files_entry("path1", 1, 100, table_ref.clone());
1837        cache.put(&key1, value1.clone());
1838        let entry_size_1 = key1.size() + value1.size();
1839        {
1840            assert_eq!(cache.memory_used(), entry_size_1);
1841        }
1842
1843        // Add another entry
1844        let (key2, value2) =
1845            create_test_list_files_entry("path2", 1, 200, table_ref.clone());
1846        cache.put(&key2, value2.clone());
1847        let entry_size_2 = key2.size() + value2.size();
1848
1849        {
1850            assert_eq!(cache.memory_used(), entry_size_1 + entry_size_2);
1851        }
1852
1853        // Remove first entry and verify memory decreases
1854        cache.remove(&key1);
1855        {
1856            assert_eq!(cache.memory_used(), entry_size_2);
1857        }
1858
1859        // Clear and verify memory is 0
1860        cache.clear();
1861        {
1862            assert_eq!(cache.memory_used(), 0);
1863        }
1864    }
1865
1866    // Prefix filtering tests using CachedFileList::filter_by_prefix
1867
1868    /// Helper function to create ObjectMeta with a specific location path
1869    fn create_object_meta_with_path(location: &str) -> ObjectMeta {
1870        ObjectMeta {
1871            location: Path::from(location),
1872            last_modified: DateTime::parse_from_rfc3339("2022-09-27T22:36:00+02:00")
1873                .unwrap()
1874                .into(),
1875            size: 1024,
1876            e_tag: None,
1877            version: None,
1878        }
1879    }
1880
1881    #[test]
1882    fn test_prefix_filtering() {
1883        let cache = DefaultCache::new(100000);
1884
1885        // Create files for a partitioned table
1886        let table_base = Path::from("my_table");
1887        let files = vec![
1888            create_object_meta_with_path("my_table/a=1/file1.parquet"),
1889            create_object_meta_with_path("my_table/a=1/file2.parquet"),
1890            create_object_meta_with_path("my_table/a=2/file3.parquet"),
1891            create_object_meta_with_path("my_table/a=2/file4.parquet"),
1892        ];
1893
1894        // Cache the full table listing
1895        let table_ref = Some(TableReference::from("table"));
1896        let key = TableScopedPath {
1897            table: table_ref,
1898            path: table_base,
1899        };
1900        cache.put(&key, CachedFileList::new(files));
1901
1902        let result = cache.get(&key).unwrap();
1903
1904        // Filter for partition a=1
1905        let prefix_a1 = Some(Path::from("my_table/a=1"));
1906        let filtered = result.files_matching_prefix(&prefix_a1);
1907        assert_eq!(filtered.len(), 2);
1908        assert!(
1909            filtered
1910                .iter()
1911                .all(|m| m.location.as_ref().starts_with("my_table/a=1"))
1912        );
1913
1914        // Filter for partition a=2
1915        let prefix_a2 = Some(Path::from("my_table/a=2"));
1916        let filtered_2 = result.files_matching_prefix(&prefix_a2);
1917        assert_eq!(filtered_2.len(), 2);
1918        assert!(
1919            filtered_2
1920                .iter()
1921                .all(|m| m.location.as_ref().starts_with("my_table/a=2"))
1922        );
1923
1924        // No filter returns all
1925        let all = result.files_matching_prefix(&None);
1926        assert_eq!(all.len(), 4);
1927    }
1928
1929    #[test]
1930    fn test_prefix_no_matching_files() {
1931        let cache = DefaultCache::new(100000);
1932
1933        let table_base = Path::from("my_table");
1934        let files = vec![
1935            create_object_meta_with_path("my_table/a=1/file1.parquet"),
1936            create_object_meta_with_path("my_table/a=2/file2.parquet"),
1937        ];
1938
1939        let table_ref = Some(TableReference::from("table"));
1940        let key = TableScopedPath {
1941            table: table_ref,
1942            path: table_base,
1943        };
1944        cache.put(&key, CachedFileList::new(files));
1945        let result = cache.get(&key).unwrap();
1946
1947        // Query for partition a=3 which doesn't exist
1948        let prefix_a3 = Some(Path::from("my_table/a=3"));
1949        let filtered = result.files_matching_prefix(&prefix_a3);
1950        assert!(filtered.is_empty());
1951    }
1952
1953    #[test]
1954    fn test_nested_partitions() {
1955        let cache = DefaultCache::new(100000);
1956
1957        let table_base = Path::from("events");
1958        let files = vec![
1959            create_object_meta_with_path(
1960                "events/year=2024/month=01/day=01/file1.parquet",
1961            ),
1962            create_object_meta_with_path(
1963                "events/year=2024/month=01/day=02/file2.parquet",
1964            ),
1965            create_object_meta_with_path(
1966                "events/year=2024/month=02/day=01/file3.parquet",
1967            ),
1968            create_object_meta_with_path(
1969                "events/year=2025/month=01/day=01/file4.parquet",
1970            ),
1971        ];
1972
1973        let table_ref = Some(TableReference::from("table"));
1974        let key = TableScopedPath {
1975            table: table_ref,
1976            path: table_base,
1977        };
1978        cache.put(&key, CachedFileList::new(files));
1979        let result = cache.get(&key).unwrap();
1980
1981        // Filter for year=2024/month=01
1982        let prefix_month = Some(Path::from("events/year=2024/month=01"));
1983        let filtered = result.files_matching_prefix(&prefix_month);
1984        assert_eq!(filtered.len(), 2);
1985
1986        // Filter for year=2024
1987        let prefix_year = Some(Path::from("events/year=2024"));
1988        let filtered_year = result.files_matching_prefix(&prefix_year);
1989        assert_eq!(filtered_year.len(), 3);
1990    }
1991
1992    #[test]
1993    fn test_drop_table_entries() {
1994        let cache = DefaultCache::new(DEFAULT_LIST_FILES_CACHE_MEMORY_LIMIT);
1995
1996        let table_ref1 = TableReference::from("table1");
1997        let table_ref2 = TableReference::from("table2");
1998        let (key1, value1) =
1999            create_test_list_files_entry("path1", 1, 100, Some(table_ref1.clone()));
2000        let (key2, value2) =
2001            create_test_list_files_entry("path2", 1, 100, Some(table_ref1.clone()));
2002        let (key3, value3) =
2003            create_test_list_files_entry("path3", 1, 100, Some(table_ref2.clone()));
2004
2005        cache.put(&key1, value1);
2006        cache.put(&key2, value2);
2007        cache.put(&key3, value3);
2008
2009        cache.drop_table_entries(&table_ref1).unwrap();
2010
2011        assert!(!cache.contains_key(&key1));
2012        assert!(!cache.contains_key(&key2));
2013        assert!(cache.contains_key(&key3));
2014    }
2015}