1use 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
28pub trait TimeProvider: Send + Sync {
30 fn now(&self) -> Instant;
32}
33
34#[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 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 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
161pub 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 pub fn new(memory_limit: usize) -> Self {
178 Self::new_with_ttl(memory_limit, None)
179 }
180
181 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 pub fn with_name(mut self, name: impl Into<String>) -> Self {
193 self.name = name.into();
194 self
195 }
196
197 pub fn with_time_provider(mut self, provider: Arc<dyn TimeProvider>) -> Self {
199 self.time_provider = provider;
200 self
201 }
202
203 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 assert!(cache.get(&object_meta.location).is_none());
381
382 let cached_entry =
384 CachedFileMetadataEntry::new(object_meta.clone(), Arc::clone(&metadata));
385 cache.put(&object_meta.location, cached_entry);
386
387 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 let result2 = cache.get(&object_meta.location).unwrap();
396 assert!(result2.is_valid_for(&object_meta));
397
398 let object_meta2 = create_test_object_meta("test", 2048);
400 let result3 = cache.get(&object_meta2.location).unwrap();
401 assert!(!result3.is_valid_for(&object_meta2));
403
404 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 cache.remove(&object_meta.location);
414 assert!(!cache.contains_key(&object_meta.location));
415
416 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 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 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 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 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 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 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 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 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 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 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 cache.clear();
583 assert_eq!(cache.len(), 0);
584 assert_eq!(cache.memory_used(), 0);
585
586 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 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 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 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 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 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 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 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 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 assert!(cache.get(&path).is_none());
838
839 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 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 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 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 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); }
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 let cached_value = CachedFileMetadata::new(
960 meta.clone(),
961 Arc::clone(&schema_fingerprint),
962 Arc::new(Statistics::new_unknown(&schema)),
963 None, );
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 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 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 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 let meta_v2 = create_test_meta("test.parquet", 200);
1015
1016 let cached = cache.get(&path).unwrap();
1017 assert!(!cached.is_valid_for(&meta_v2, &schema_fingerprint));
1019
1020 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 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 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 let cached = cache.get(&path).unwrap();
1065 assert!(cached.is_valid_for(&meta_v1, &schema_fingerprint));
1066 assert!(cached.ordering.is_some());
1067
1068 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, e_tag: None,
1076 version: None,
1077 };
1078
1079 let cached = cache.get(&path).unwrap();
1081 assert!(!cached.is_valid_for(&meta_v2, &schema_fingerprint));
1082
1083 let ordering_v2 = ordering(); 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 let cached = cache.get(&path).unwrap();
1095 assert!(!cached.is_valid_for(&meta_v1, &schema_fingerprint));
1096
1097 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 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 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 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 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 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 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 fn create_object_meta(path: &str, location_size: usize) -> ObjectMeta {
1340 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 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 assert!(!cache.contains_key(&key));
1388 assert_eq!(cache.len(), 0);
1389
1390 assert!(cache.get(&key).is_none());
1392
1393 let meta = create_test_object_meta("file1", 50);
1395 cache.put(&key, CachedFileList::new(vec![meta]));
1396
1397 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 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 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 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 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 let cache = DefaultCache::new(entry_size * 3);
1463
1464 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 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)); 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 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 let _ = cache.get(&key1);
1505
1506 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)); assert!(!cache.contains_key(&key2)); 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 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 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 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 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 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 assert_eq!(cache.len(), 2);
1572 assert!(!cache.contains_key(&key1)); assert!(!cache.contains_key(&key2)); 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 cache.put(&key1, value1);
1593 cache.put(&key2, value2);
1594 cache.put(&key3, value3);
1595 assert_eq!(cache.len(), 3);
1596
1597 cache.update_cache_limit(entry_size);
1599
1600 assert_eq!(cache.len(), 1);
1602 assert!(cache.contains_key(&key3));
1603 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 cache.put(&key1, value1);
1624 cache.put(&key2, value2.clone());
1625 cache.put(&key3, value3_v1);
1626 assert_eq!(cache.len(), 3);
1627
1628 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 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)); assert!(cache.contains_key(&key2));
1645 assert!(cache.contains_key(&key3));
1646
1647 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 assert!(cache.get(&key1).is_some());
1690 assert!(cache.get(&key2).is_some());
1691 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 mock_time.inc(Duration::from_millis(150));
1717
1718 assert!(!cache.contains_key(&key1));
1720 assert_eq!(cache.len(), 1); assert!(!cache.contains_key(&key2));
1722 assert_eq!(cache.len(), 0); }
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 cache.put(&key3, value3);
1747 assert!(!cache.contains_key(&key1)); 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)); assert!(cache.contains_key(&key3)); }
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.put(&key, value.clone());
1767
1768 let result = cache.get(&key);
1770 assert!(result.is_some());
1771 assert_eq!(result.unwrap().files.len(), 2);
1772
1773 thread::sleep(Duration::from_millis(150));
1775
1776 let result2 = cache.get(&key);
1778 assert!(result2.is_none());
1779 }
1780
1781 #[test]
1782 fn test_meta_heap_bytes_calculation() {
1783 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); 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); 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); 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); }
1823
1824 #[test]
1825 fn test_memory_tracking() {
1826 let cache = DefaultCache::new(1000);
1827
1828 {
1830 assert_eq!(cache.memory_used(), 0);
1831 }
1832
1833 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 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 cache.remove(&key1);
1855 {
1856 assert_eq!(cache.memory_used(), entry_size_2);
1857 }
1858
1859 cache.clear();
1861 {
1862 assert_eq!(cache.memory_used(), 0);
1863 }
1864 }
1865
1866 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 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 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 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 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 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 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 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 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}