1pub(crate) mod cache_size;
18
19pub(crate) mod file_cache;
20pub(crate) mod index;
21pub(crate) mod manifest_cache;
22#[cfg(test)]
23pub(crate) mod test_util;
24pub(crate) mod write_cache;
25
26use std::collections::{BTreeMap, HashMap};
27use std::mem;
28use std::ops::Range;
29use std::sync::{Arc, RwLock, Weak};
30
31use bytes::Bytes;
32use common_base::readable_size::ReadableSize;
33use common_datasource::compression::CompressionType;
34use common_runtime::Runtime;
35use common_runtime::runtime::RuntimeTrait;
36use common_telemetry::warn;
37use datatypes::arrow::buffer::BooleanBuffer;
38use datatypes::arrow::record_batch::RecordBatch;
39use datatypes::types::json_type::JsonNativeType;
40use datatypes::value::Value;
41use datatypes::vectors::VectorRef;
42use index::bloom_filter_index::{BloomFilterIndexCache, BloomFilterIndexCacheRef};
43use index::result_cache::IndexResultCache;
44use moka::notification::RemovalCause;
45use moka::sync::Cache;
46use object_store::ObjectStore;
47use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
48use parquet::file::metadata::{
49 FileMetaData, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, ParquetMetaDataWriter,
50};
51use puffin::puffin_manager::cache::{PuffinMetadataCache, PuffinMetadataCacheRef};
52use smallvec::SmallVec;
53use snafu::{OptionExt, ResultExt};
54use store_api::metadata::{RegionMetadata, RegionMetadataRef};
55use store_api::storage::{ColumnId, ConcreteDataType, FileId, RegionId, TimeSeriesRowSelector};
56pub use write_cache::{WriteCacheUploadStoreWrapper, WriteCacheUploadStoreWrapperRef};
57
58use crate::cache::cache_size::parquet_meta_size;
59use crate::cache::file_cache::{FileType, IndexKey};
60use crate::cache::index::inverted_index::{InvertedIndexCache, InvertedIndexCacheRef};
61#[cfg(feature = "vector_index")]
62use crate::cache::index::vector_index::{VectorIndexCache, VectorIndexCacheRef};
63use crate::cache::write_cache::WriteCacheRef;
64use crate::error::{
65 CompressObjectSnafu, DecompressObjectSnafu, InvalidMetadataSnafu, InvalidParquetSnafu,
66 JoinSnafu, ReadParquetSnafu, Result, UnexpectedSnafu, WriteParquetSnafu,
67};
68use crate::memtable::record_batch_estimated_size;
69use crate::metrics::{CACHE_BYTES, CACHE_EVICTION, CACHE_HIT, CACHE_MISS};
70use crate::read::Batch;
71use crate::read::range_cache::{RangeScanCacheKey, RangeScanCacheValue};
72use crate::read::read_columns::JsonTargetTypes;
73use crate::sst::file::{RegionFileId, RegionIndexId};
74use crate::sst::parquet::PARQUET_METADATA_KEY;
75use crate::sst::parquet::read_columns::ParquetReadColumns;
76use crate::sst::parquet::reader::MetadataCacheMetrics;
77
78const SST_META_TYPE: &str = "sst_meta";
80const SST_META_DECODED_TYPE: &str = "sst_meta_decoded";
82const VECTOR_TYPE: &str = "vector";
84const PAGE_TYPE: &str = "page";
86const FILE_TYPE: &str = "file";
88const INDEX_TYPE: &str = "index";
90const SELECTOR_RESULT_TYPE: &str = "selector_result";
92const RANGE_RESULT_TYPE: &str = "range_result";
94const PREFILTER_RESULT_TYPE: &str = "prefilter_result";
96const RANGE_RESULT_CONCAT_MEMORY_LIMIT: ReadableSize = ReadableSize::mb(512);
97const RANGE_RESULT_CONCAT_MEMORY_PERMIT: ReadableSize = ReadableSize::kb(1);
98
99#[derive(Debug)]
100pub(crate) struct RangeResultMemoryLimiter {
101 semaphore: Arc<tokio::sync::Semaphore>,
102 permit_bytes: usize,
103 total_permits: usize,
104 #[cfg(test)]
107 waited_acquires: std::sync::atomic::AtomicUsize,
108}
109
110impl Default for RangeResultMemoryLimiter {
111 fn default() -> Self {
112 Self::new(
113 RANGE_RESULT_CONCAT_MEMORY_LIMIT.as_bytes() as usize,
114 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize,
115 )
116 }
117}
118
119impl RangeResultMemoryLimiter {
120 pub(crate) fn new(limit_bytes: usize, permit_bytes: usize) -> Self {
121 let permit_bytes = permit_bytes.max(1);
122 let total_permits = limit_bytes
123 .div_ceil(permit_bytes)
124 .clamp(1, tokio::sync::Semaphore::MAX_PERMITS);
125 Self {
126 semaphore: Arc::new(tokio::sync::Semaphore::new(total_permits)),
127 permit_bytes,
128 total_permits,
129 #[cfg(test)]
130 waited_acquires: std::sync::atomic::AtomicUsize::new(0),
131 }
132 }
133
134 #[cfg(test)]
135 pub(crate) fn permit_bytes(&self) -> usize {
136 self.permit_bytes
137 }
138
139 #[cfg(test)]
140 pub(crate) fn waited_acquires(&self) -> usize {
141 self.waited_acquires
142 .load(std::sync::atomic::Ordering::Acquire)
143 }
144
145 #[cfg(test)]
146 pub(crate) fn available_permits(&self) -> usize {
147 self.semaphore.available_permits()
148 }
149
150 pub(crate) async fn acquire(&self, bytes: usize) -> Result<tokio::sync::SemaphorePermit<'_>> {
151 let permits = bytes.div_ceil(self.permit_bytes).max(1);
152 if permits > self.total_permits {
153 return UnexpectedSnafu {
154 reason: format!(
155 "range result memory request of {bytes} bytes exceeds limiter capacity of {} bytes",
156 self.total_permits.saturating_mul(self.permit_bytes)
157 ),
158 }
159 .fail();
160 }
161 #[cfg(test)]
164 if self.semaphore.available_permits() < permits {
165 self.waited_acquires
166 .fetch_add(1, std::sync::atomic::Ordering::Release);
167 }
168
169 self.semaphore
170 .acquire_many(permits as u32)
171 .await
172 .map_err(|_| {
173 UnexpectedSnafu {
174 reason: "range result memory limiter is unexpectedly closed",
175 }
176 .build()
177 })
178 }
179}
180
181#[derive(Debug)]
186pub(crate) struct CachedSstMeta {
187 parquet_metadata: Arc<ParquetMetaData>,
188 parquet_metadata_size: usize,
189 region_metadata: RegionMetadataRef,
190 region_metadata_weight: usize,
191 page_index_policy: PageIndexPolicy,
192}
193
194#[derive(Debug)]
200pub(crate) struct CompactSstMeta {
201 encoded_metadata: Bytes,
202 decoded_size: usize,
203 region_metadata: Weak<RegionMetadata>,
204 page_index_policy: PageIndexPolicy,
205 decoded: tokio::sync::Mutex<Weak<CachedSstMeta>>,
206}
207
208#[derive(Debug)]
210pub(crate) struct PreparedSstMeta {
211 compact: Arc<CompactSstMeta>,
212 decoded: Arc<CachedSstMeta>,
213}
214
215impl PreparedSstMeta {
216 pub(crate) fn decoded(&self) -> Arc<CachedSstMeta> {
217 self.decoded.clone()
218 }
219}
220
221#[derive(Debug)]
223pub(crate) enum SstMetaPreparation {
224 Prepared(PreparedSstMeta),
226 DecodedOnly {
228 decoded: Arc<CachedSstMeta>,
229 encoding_error: crate::error::Error,
230 },
231}
232
233impl SstMetaPreparation {
234 pub(crate) fn decoded(&self) -> Arc<CachedSstMeta> {
235 match self {
236 Self::Prepared(metadata) => metadata.decoded(),
237 Self::DecodedOnly { decoded, .. } => decoded.clone(),
238 }
239 }
240}
241
242impl CompactSstMeta {
243 async fn decode(&self, runtime: &Runtime) -> Result<Arc<CachedSstMeta>> {
244 let mut decoded_guard = self.decoded.lock().await;
245 if let Some(decoded) = decoded_guard.upgrade() {
246 return Ok(decoded);
247 }
248
249 let encoded_metadata = self.encoded_metadata.clone();
250 let decoded_size = self.decoded_size;
251 let region_metadata = self.region_metadata.upgrade();
252 let page_index_policy = self.page_index_policy;
253 let decoded = runtime
254 .spawn_blocking(move || {
255 let bytes = zstd::bulk::decompress(&encoded_metadata, decoded_size).context(
256 DecompressObjectSnafu {
257 compress_type: CompressionType::Zstd,
258 path: "cached SST metadata",
259 },
260 )?;
261 let bytes = Bytes::from(bytes);
262 let mut reader = ParquetMetaDataReader::new()
263 .with_column_index_policy(PageIndexPolicy::Skip)
264 .with_offset_index_policy(page_index_policy);
265 reader.try_parse(&bytes).context(ReadParquetSnafu {
266 path: "cached SST metadata",
267 })?;
268 let metadata = reader.finish().context(ReadParquetSnafu {
269 path: "cached SST metadata",
270 })?;
271 CachedSstMeta::try_new_with_page_index_policy(
272 "cached SST metadata",
273 metadata,
274 region_metadata,
275 page_index_policy,
276 )
277 .map(Arc::new)
278 })
279 .await
280 .context(JoinSnafu)??;
281
282 *decoded_guard = Arc::downgrade(&decoded);
283 Ok(decoded)
284 }
285
286 fn satisfies_page_index_policy(&self, requested: PageIndexPolicy) -> bool {
287 satisfies_page_index_policy(self.page_index_policy, requested)
288 }
289}
290
291pub(crate) async fn decode_sst_meta(
293 file_path: &str,
294 parquet_metadata: ParquetMetaData,
295 region_metadata: Option<RegionMetadataRef>,
296 page_index_policy: PageIndexPolicy,
297 runtime: &Runtime,
298) -> Result<Arc<CachedSstMeta>> {
299 let file_path = file_path.to_string();
300 runtime
301 .spawn_blocking(move || {
302 let parquet_metadata = strip_column_indexes(parquet_metadata);
303 CachedSstMeta::try_new_with_page_index_policy(
304 &file_path,
305 parquet_metadata,
306 region_metadata,
307 page_index_policy,
308 )
309 .map(Arc::new)
310 })
311 .await
312 .context(JoinSnafu)?
313}
314
315pub(crate) async fn prepare_sst_meta(
317 file_path: &str,
318 parquet_metadata: ParquetMetaData,
319 region_metadata: Option<RegionMetadataRef>,
320 page_index_policy: PageIndexPolicy,
321 runtime: &Runtime,
322) -> Result<SstMetaPreparation> {
323 let file_path = file_path.to_string();
324 runtime
325 .spawn_blocking(move || {
326 prepare_sst_meta_sync(
327 &file_path,
328 parquet_metadata,
329 region_metadata,
330 page_index_policy,
331 )
332 })
333 .await
334 .context(JoinSnafu)?
335}
336
337pub(crate) fn prepare_sst_meta_sync(
339 file_path: &str,
340 parquet_metadata: ParquetMetaData,
341 region_metadata: Option<RegionMetadataRef>,
342 page_index_policy: PageIndexPolicy,
343) -> Result<SstMetaPreparation> {
344 let parquet_metadata = strip_column_indexes(parquet_metadata);
345 let cache_encoding = encode_compact_sst_meta(file_path, &parquet_metadata);
346 finish_sst_meta_preparation(
347 file_path,
348 parquet_metadata,
349 region_metadata,
350 page_index_policy,
351 cache_encoding,
352 )
353}
354
355fn strip_column_indexes(parquet_metadata: ParquetMetaData) -> ParquetMetaData {
356 let mut builder = parquet_metadata.into_builder();
359 builder.take_column_index();
360 builder.build()
361}
362
363fn encode_compact_sst_meta(
364 file_path: &str,
365 parquet_metadata: &ParquetMetaData,
366) -> Result<(Bytes, usize)> {
367 let mut encoded = Vec::new();
368 ParquetMetaDataWriter::new(&mut encoded, parquet_metadata)
369 .finish()
370 .context(WriteParquetSnafu)?;
371 let decoded_size = encoded.len();
372 let encoded_metadata = zstd::bulk::compress(&encoded, 3).context(CompressObjectSnafu {
373 compress_type: CompressionType::Zstd,
374 path: file_path,
375 })?;
376
377 Ok((
381 Bytes::from(encoded_metadata.into_boxed_slice()),
382 decoded_size,
383 ))
384}
385
386fn finish_sst_meta_preparation(
387 file_path: &str,
388 parquet_metadata: ParquetMetaData,
389 region_metadata: Option<RegionMetadataRef>,
390 page_index_policy: PageIndexPolicy,
391 cache_encoding: Result<(Bytes, usize)>,
392) -> Result<SstMetaPreparation> {
393 let decoded = Arc::new(CachedSstMeta::try_new_with_page_index_policy(
394 file_path,
395 parquet_metadata,
396 region_metadata,
397 page_index_policy,
398 )?);
399 let (encoded_metadata, decoded_size) = match cache_encoding {
400 Ok(encoded) => encoded,
401 Err(encoding_error) => {
402 return Ok(SstMetaPreparation::DecodedOnly {
403 decoded,
404 encoding_error,
405 });
406 }
407 };
408 let compact = Arc::new(CompactSstMeta {
409 encoded_metadata,
410 decoded_size,
411 region_metadata: Arc::downgrade(&decoded.region_metadata),
412 page_index_policy,
413 decoded: tokio::sync::Mutex::new(Arc::downgrade(&decoded)),
414 });
415
416 Ok(SstMetaPreparation::Prepared(PreparedSstMeta {
417 compact,
418 decoded,
419 }))
420}
421
422impl CachedSstMeta {
423 #[cfg(test)]
424 pub(crate) fn try_new(file_path: &str, parquet_metadata: ParquetMetaData) -> Result<Self> {
425 let page_index_policy = infer_loaded_page_index_policy(&parquet_metadata);
426 Self::try_new_with_page_index_policy(file_path, parquet_metadata, None, page_index_policy)
427 }
428
429 pub(crate) fn try_new_with_region_metadata(
430 file_path: &str,
431 parquet_metadata: ParquetMetaData,
432 region_metadata: Option<RegionMetadataRef>,
433 ) -> Result<Self> {
434 let page_index_policy = infer_loaded_page_index_policy(&parquet_metadata);
435 Self::try_new_with_page_index_policy(
436 file_path,
437 parquet_metadata,
438 region_metadata,
439 page_index_policy,
440 )
441 }
442
443 pub(crate) fn try_new_with_page_index_policy(
444 file_path: &str,
445 parquet_metadata: ParquetMetaData,
446 region_metadata: Option<RegionMetadataRef>,
447 page_index_policy: PageIndexPolicy,
448 ) -> Result<Self> {
449 let file_metadata = parquet_metadata.file_metadata();
450 let key_values = file_metadata
451 .key_value_metadata()
452 .context(InvalidParquetSnafu {
453 file: file_path,
454 reason: "missing key value meta",
455 })?;
456 let meta_value = key_values
457 .iter()
458 .find(|kv| kv.key == PARQUET_METADATA_KEY)
459 .with_context(|| InvalidParquetSnafu {
460 file: file_path,
461 reason: format!("key {} not found", PARQUET_METADATA_KEY),
462 })?;
463 let json = meta_value
464 .value
465 .as_ref()
466 .with_context(|| InvalidParquetSnafu {
467 file: file_path,
468 reason: format!("No value for key {}", PARQUET_METADATA_KEY),
469 })?;
470 let region_metadata = match region_metadata {
471 Some(region_metadata) => region_metadata,
472 None => Arc::new(
473 store_api::metadata::RegionMetadata::from_json(json)
474 .context(InvalidMetadataSnafu)?,
475 ),
476 };
477 let region_metadata_weight = region_metadata.estimated_size().max(json.len());
479 let parquet_metadata = Arc::new(strip_region_metadata_from_parquet(parquet_metadata));
480 let parquet_metadata_size = parquet_meta_size(&parquet_metadata);
481
482 Ok(Self {
483 parquet_metadata,
484 parquet_metadata_size,
485 region_metadata,
486 region_metadata_weight,
487 page_index_policy,
488 })
489 }
490
491 pub(crate) fn parquet_metadata(&self) -> Arc<ParquetMetaData> {
492 self.parquet_metadata.clone()
493 }
494
495 pub(crate) fn parquet_metadata_size(&self) -> usize {
497 self.parquet_metadata_size
498 }
499
500 pub(crate) fn region_metadata(&self) -> RegionMetadataRef {
501 self.region_metadata.clone()
502 }
503
504 fn satisfies_page_index_policy(&self, requested: PageIndexPolicy) -> bool {
505 satisfies_page_index_policy(self.page_index_policy, requested)
506 }
507}
508
509fn satisfies_page_index_policy(cached: PageIndexPolicy, requested: PageIndexPolicy) -> bool {
510 match requested {
511 PageIndexPolicy::Skip => true,
512 PageIndexPolicy::Optional => cached != PageIndexPolicy::Skip,
513 PageIndexPolicy::Required => cached == PageIndexPolicy::Required,
514 }
515}
516
517fn infer_loaded_page_index_policy(parquet_metadata: &ParquetMetaData) -> PageIndexPolicy {
518 if parquet_metadata.offset_index().is_some() {
519 PageIndexPolicy::Optional
520 } else {
521 PageIndexPolicy::Skip
522 }
523}
524
525fn strip_region_metadata_from_parquet(parquet_metadata: ParquetMetaData) -> ParquetMetaData {
526 let file_metadata = parquet_metadata.file_metadata();
527 let filtered_key_values = file_metadata.key_value_metadata().and_then(|key_values| {
528 let filtered = key_values
529 .iter()
530 .filter(|kv| kv.key != PARQUET_METADATA_KEY)
531 .cloned()
532 .collect::<Vec<_>>();
533 (!filtered.is_empty()).then_some(filtered)
534 });
535 let stripped_file_metadata = FileMetaData::new(
536 file_metadata.version(),
537 file_metadata.num_rows(),
538 file_metadata.created_by().map(ToString::to_string),
539 filtered_key_values,
540 file_metadata.schema_descr_ptr(),
541 file_metadata.column_orders().cloned(),
542 );
543
544 let mut builder = parquet_metadata.into_builder();
545 let row_groups = builder.take_row_groups();
546 let offset_index = builder.take_offset_index();
547
548 parquet::file::metadata::ParquetMetaDataBuilder::new(stripped_file_metadata)
549 .set_row_groups(row_groups)
550 .set_offset_index(offset_index)
551 .build()
552}
553
554fn removal_cause_str(cause: RemovalCause) -> &'static str {
555 match cause {
556 RemovalCause::Expired => "expired",
557 RemovalCause::Explicit => "explicit",
558 RemovalCause::Replaced => "replaced",
559 RemovalCause::Size => "size",
560 }
561}
562
563#[derive(Debug, Clone, PartialEq, Eq, Hash)]
564pub(crate) struct PrefilterRowSelector {
565 row_count: usize,
566 skip: bool,
567}
568
569impl From<&RowSelector> for PrefilterRowSelector {
574 fn from(selector: &RowSelector) -> Self {
575 Self {
576 row_count: selector.row_count,
577 skip: selector.skip,
578 }
579 }
580}
581
582#[derive(Debug, Clone, PartialEq, Eq, Hash)]
584pub(crate) struct PrefilterKey {
585 file_id: FileId,
586 row_group_idx: u32,
587 row_selection: Option<Arc<Vec<PrefilterRowSelector>>>,
588 schema_version: u64,
589 filter_exprs: SmallVec<[String; 1]>,
590 mem_usage: usize,
591}
592
593impl PrefilterKey {
594 pub(crate) fn row_selection_snapshot(
595 row_selection: Option<&RowSelection>,
596 ) -> Option<Arc<Vec<PrefilterRowSelector>>> {
597 row_selection.map(|selection| {
598 Arc::new(
599 selection
600 .iter()
601 .map(PrefilterRowSelector::from)
602 .collect::<Vec<_>>(),
603 )
604 })
605 }
606
607 pub(crate) fn new(
608 file_id: FileId,
609 row_group_idx: u32,
610 row_selection: Option<Arc<Vec<PrefilterRowSelector>>>,
611 schema_version: u64,
612 filter_exprs: SmallVec<[String; 1]>,
613 ) -> Self {
614 let row_selection_bytes = row_selection
615 .as_ref()
616 .map(|selection| selection.len() * mem::size_of::<PrefilterRowSelector>())
617 .unwrap_or(0);
618 let spilled_expr_bytes = if filter_exprs.spilled() {
619 filter_exprs.capacity() * mem::size_of::<String>()
620 } else {
621 0
622 };
623 let expr_bytes = filter_exprs.iter().map(|s| s.capacity()).sum::<usize>();
624
625 Self {
626 file_id,
627 row_group_idx,
628 row_selection,
629 schema_version,
630 filter_exprs,
631 mem_usage: mem::size_of::<Self>()
632 + row_selection_bytes
633 + spilled_expr_bytes
634 + expr_bytes,
635 }
636 }
637
638 fn mem_usage(&self) -> usize {
639 self.mem_usage
640 }
641}
642
643type PrefilterResultCache = Cache<PrefilterKey, Arc<BooleanBuffer>>;
644
645fn new_prefilter_result_cache(capacity: u64) -> PrefilterResultCache {
646 Cache::builder()
647 .max_capacity(capacity)
648 .weigher(prefilter_result_cache_weight)
649 .eviction_listener(|k, v, cause| {
650 let size = prefilter_result_cache_weight(&k, &v);
651 CACHE_BYTES
652 .with_label_values(&[PREFILTER_RESULT_TYPE])
653 .sub(size.into());
654 CACHE_EVICTION
655 .with_label_values(&[PREFILTER_RESULT_TYPE, removal_cause_str(cause)])
656 .inc();
657 })
658 .build()
659}
660
661fn prefilter_result_cache_weight(k: &PrefilterKey, v: &Arc<BooleanBuffer>) -> u32 {
662 (k.mem_usage() + mem::size_of::<BooleanBuffer>() + v.values().len()) as u32
663}
664
665#[derive(Clone)]
667pub enum CacheStrategy {
668 EnableAll(CacheManagerRef),
671 Compaction(CacheManagerRef),
676 Disabled,
678}
679
680impl CacheStrategy {
681 pub(crate) fn sst_meta_runtime(&self) -> Runtime {
683 match self {
684 CacheStrategy::Compaction(_) => common_runtime::compact_runtime(),
685 CacheStrategy::EnableAll(_) | CacheStrategy::Disabled => {
686 common_runtime::global_runtime()
687 }
688 }
689 }
690
691 pub(crate) fn sst_meta_cache_enabled(&self) -> bool {
693 match self {
694 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
695 cache_manager.sst_meta_cache_enabled()
696 }
697 CacheStrategy::Disabled => false,
698 }
699 }
700
701 pub(crate) async fn get_sst_meta_data(
703 &self,
704 file_id: RegionFileId,
705 metrics: &mut MetadataCacheMetrics,
706 page_index_policy: PageIndexPolicy,
707 ) -> Option<Arc<CachedSstMeta>> {
708 match self {
709 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
710 cache_manager
711 .get_sst_meta_data(
712 file_id,
713 metrics,
714 page_index_policy,
715 &self.sst_meta_runtime(),
716 )
717 .await
718 }
719 CacheStrategy::Disabled => {
720 metrics.cache_miss += 1;
721 None
722 }
723 }
724 }
725
726 pub(crate) fn get_sst_meta_data_from_mem_cache(
728 &self,
729 file_id: RegionFileId,
730 page_index_policy: PageIndexPolicy,
731 ) -> Option<Arc<CachedSstMeta>> {
732 match self {
733 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
734 cache_manager.get_sst_meta_data_from_mem_cache(file_id, page_index_policy)
735 }
736 CacheStrategy::Disabled => None,
737 }
738 }
739
740 pub fn get_parquet_meta_data_from_mem_cache(
742 &self,
743 file_id: RegionFileId,
744 ) -> Option<Arc<ParquetMetaData>> {
745 self.get_sst_meta_data_from_mem_cache(file_id, PageIndexPolicy::Skip)
746 .map(|metadata| metadata.parquet_metadata())
747 }
748
749 pub(crate) fn put_prepared_sst_meta(
751 &self,
752 file_id: RegionFileId,
753 metadata: PreparedSstMeta,
754 retain_decoded: bool,
755 ) {
756 match self {
757 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
758 cache_manager.put_prepared_sst_meta(file_id, metadata, retain_decoded);
759 }
760 CacheStrategy::Disabled => {}
761 }
762 }
763
764 pub fn put_parquet_meta_data(
766 &self,
767 file_id: RegionFileId,
768 metadata: Arc<ParquetMetaData>,
769 region_metadata: Option<RegionMetadataRef>,
770 ) {
771 match self {
772 CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => {
773 cache_manager.put_parquet_meta_data(file_id, metadata, region_metadata);
774 }
775 CacheStrategy::Disabled => {}
776 }
777 }
778
779 pub(crate) fn get_prefilter_result(&self, key: &PrefilterKey) -> Option<Arc<BooleanBuffer>> {
782 match self {
783 CacheStrategy::EnableAll(cache_manager) => cache_manager.get_prefilter_result(key),
784 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
785 }
786 }
787
788 pub(crate) fn put_prefilter_result(&self, key: PrefilterKey, result: Arc<BooleanBuffer>) {
791 if let CacheStrategy::EnableAll(cache_manager) = self {
792 cache_manager.put_prefilter_result(key, result);
793 }
794 }
795
796 pub fn remove_parquet_meta_data(&self, file_id: RegionFileId) {
798 match self {
799 CacheStrategy::EnableAll(cache_manager) => {
800 cache_manager.remove_parquet_meta_data(file_id);
801 }
802 CacheStrategy::Compaction(cache_manager) => {
803 cache_manager.remove_parquet_meta_data(file_id);
804 }
805 CacheStrategy::Disabled => {}
806 }
807 }
808
809 pub fn get_repeated_vector(
812 &self,
813 data_type: &ConcreteDataType,
814 value: &Value,
815 ) -> Option<VectorRef> {
816 match self {
817 CacheStrategy::EnableAll(cache_manager) => {
818 cache_manager.get_repeated_vector(data_type, value)
819 }
820 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
821 }
822 }
823
824 pub fn put_repeated_vector(&self, value: Value, vector: VectorRef) {
827 if let CacheStrategy::EnableAll(cache_manager) = self {
828 cache_manager.put_repeated_vector(value, vector);
829 }
830 }
831
832 pub fn get_page_ranges(
835 &self,
836 file_id: FileId,
837 row_group_idx: usize,
838 ranges: &[Range<u64>],
839 ) -> Option<PageRangeLookup> {
840 match self {
841 CacheStrategy::EnableAll(cache_manager) => {
842 cache_manager.get_page_ranges(file_id, row_group_idx, ranges)
843 }
844 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
845 }
846 }
847
848 pub fn put_page_ranges(
851 &self,
852 file_id: FileId,
853 row_group_idx: usize,
854 ranges: &[Range<u64>],
855 pages: &[Bytes],
856 ) {
857 if let CacheStrategy::EnableAll(cache_manager) = self {
858 cache_manager.put_page_ranges(file_id, row_group_idx, ranges, pages);
859 }
860 }
861
862 pub async fn evict_puffin_cache(&self, file_id: RegionIndexId) {
864 match self {
865 CacheStrategy::EnableAll(cache_manager) => {
866 cache_manager.evict_puffin_cache(file_id).await
867 }
868 CacheStrategy::Compaction(cache_manager) => {
869 cache_manager.evict_puffin_cache(file_id).await
870 }
871 CacheStrategy::Disabled => {}
872 }
873 }
874
875 pub fn get_selector_result(
878 &self,
879 selector_key: &SelectorResultKey,
880 ) -> Option<Arc<SelectorResultValue>> {
881 match self {
882 CacheStrategy::EnableAll(cache_manager) => {
883 cache_manager.get_selector_result(selector_key)
884 }
885 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
886 }
887 }
888
889 pub fn put_selector_result(
892 &self,
893 selector_key: SelectorResultKey,
894 result: Arc<SelectorResultValue>,
895 ) {
896 if let CacheStrategy::EnableAll(cache_manager) = self {
897 cache_manager.put_selector_result(selector_key, result);
898 }
899 }
900
901 #[allow(dead_code)]
904 pub(crate) fn get_range_result(
905 &self,
906 key: &RangeScanCacheKey,
907 ) -> Option<Arc<RangeScanCacheValue>> {
908 match self {
909 CacheStrategy::EnableAll(cache_manager) => cache_manager.get_range_result(key),
910 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
911 }
912 }
913
914 pub(crate) fn put_range_result(
917 &self,
918 key: RangeScanCacheKey,
919 result: Arc<RangeScanCacheValue>,
920 ) {
921 if let CacheStrategy::EnableAll(cache_manager) = self {
922 cache_manager.put_range_result(key, result);
923 }
924 }
925
926 pub(crate) fn has_range_result_cache(&self) -> bool {
928 match self {
929 CacheStrategy::EnableAll(cache_manager) => cache_manager.has_range_result_cache(),
930 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => false,
931 }
932 }
933
934 pub(crate) fn range_result_memory_limiter(&self) -> Option<&Arc<RangeResultMemoryLimiter>> {
935 match self {
936 CacheStrategy::EnableAll(cache_manager) => {
937 Some(cache_manager.range_result_memory_limiter())
938 }
939 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
940 }
941 }
942
943 pub(crate) fn range_result_cache_size(&self) -> Option<usize> {
944 match self {
945 CacheStrategy::EnableAll(cache_manager) => {
946 Some(cache_manager.range_result_cache_size())
947 }
948 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
949 }
950 }
951
952 pub fn write_cache(&self) -> Option<&WriteCacheRef> {
955 match self {
956 CacheStrategy::EnableAll(cache_manager) => cache_manager.write_cache(),
957 CacheStrategy::Compaction(cache_manager) => cache_manager.write_cache(),
958 CacheStrategy::Disabled => None,
959 }
960 }
961
962 pub fn inverted_index_cache(&self) -> Option<&InvertedIndexCacheRef> {
965 match self {
966 CacheStrategy::EnableAll(cache_manager) => cache_manager.inverted_index_cache(),
967 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
968 }
969 }
970
971 pub fn bloom_filter_index_cache(&self) -> Option<&BloomFilterIndexCacheRef> {
974 match self {
975 CacheStrategy::EnableAll(cache_manager) => cache_manager.bloom_filter_index_cache(),
976 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
977 }
978 }
979
980 #[cfg(feature = "vector_index")]
983 pub fn vector_index_cache(&self) -> Option<&VectorIndexCacheRef> {
984 match self {
985 CacheStrategy::EnableAll(cache_manager) => cache_manager.vector_index_cache(),
986 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
987 }
988 }
989
990 pub fn puffin_metadata_cache(&self) -> Option<&PuffinMetadataCacheRef> {
993 match self {
994 CacheStrategy::EnableAll(cache_manager) => cache_manager.puffin_metadata_cache(),
995 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
996 }
997 }
998
999 pub fn index_result_cache(&self) -> Option<&IndexResultCache> {
1002 match self {
1003 CacheStrategy::EnableAll(cache_manager) => cache_manager.index_result_cache(),
1004 CacheStrategy::Compaction(_) | CacheStrategy::Disabled => None,
1005 }
1006 }
1007
1008 pub fn maybe_download_background(
1010 &self,
1011 index_key: IndexKey,
1012 remote_path: String,
1013 remote_store: ObjectStore,
1014 file_size: u64,
1015 ) {
1016 if let CacheStrategy::EnableAll(cache_manager) = self
1017 && let Some(write_cache) = cache_manager.write_cache()
1018 {
1019 write_cache.file_cache().maybe_download_background(
1020 index_key,
1021 remote_path,
1022 remote_store,
1023 file_size,
1024 );
1025 }
1026 }
1027}
1028
1029#[derive(Default)]
1033pub struct CacheManager {
1034 sst_meta_cache: Option<SstMetaCache>,
1036 sst_decoded_meta_cache: Option<SstDecodedMetaCache>,
1038 vector_cache: Option<VectorCache>,
1040 page_cache: Option<Arc<PageRangeCache>>,
1042 write_cache: Option<WriteCacheRef>,
1044 inverted_index_cache: Option<InvertedIndexCacheRef>,
1046 bloom_filter_index_cache: Option<BloomFilterIndexCacheRef>,
1048 #[cfg(feature = "vector_index")]
1050 vector_index_cache: Option<VectorIndexCacheRef>,
1051 puffin_metadata_cache: Option<PuffinMetadataCacheRef>,
1053 selector_result_cache: Option<SelectorResultCache>,
1055 range_result_cache: Option<RangeResultCache>,
1057 range_result_cache_size: u64,
1059 range_result_memory_limiter: Arc<RangeResultMemoryLimiter>,
1061 index_result_cache: Option<IndexResultCache>,
1063 prefilter_result_cache: Option<PrefilterResultCache>,
1065}
1066
1067pub type CacheManagerRef = Arc<CacheManager>;
1068
1069impl CacheManager {
1070 pub fn builder() -> CacheManagerBuilder {
1072 CacheManagerBuilder::default()
1073 }
1074
1075 pub(crate) async fn get_sst_meta_data(
1078 &self,
1079 file_id: RegionFileId,
1080 metrics: &mut MetadataCacheMetrics,
1081 page_index_policy: PageIndexPolicy,
1082 runtime: &Runtime,
1083 ) -> Option<Arc<CachedSstMeta>> {
1084 let cache_key = SstMetaKey(file_id.region_id(), file_id.file_id());
1085 let compact = self
1086 .get_compact_sst_meta(&cache_key)
1087 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1088 let decoded = self
1089 .get_decoded_sst_meta(&cache_key)
1090 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1091
1092 if let Some(compact) = compact {
1093 CACHE_HIT.with_label_values(&[SST_META_TYPE]).inc();
1094 metrics.mem_cache_hit += 1;
1095 if let Some(decoded) = decoded {
1096 CACHE_HIT.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1097 return Some(decoded);
1098 }
1099
1100 CACHE_MISS.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1101 match compact.decode(runtime).await {
1102 Ok(decoded) => {
1103 self.put_sst_meta_data(file_id, decoded.clone());
1104 return Some(decoded);
1105 }
1106 Err(err) => {
1107 warn!(err; "Failed to decode compact SST metadata, region_id: {}, file_id: {}", file_id.region_id(), file_id.file_id());
1108 self.remove_parquet_meta_data(file_id);
1109 }
1110 }
1111 } else if let Some(decoded) = decoded {
1112 CACHE_HIT.with_label_values(&[SST_META_TYPE]).inc();
1115 CACHE_HIT.with_label_values(&[SST_META_DECODED_TYPE]).inc();
1116 metrics.mem_cache_hit += 1;
1117 return Some(decoded);
1118 } else {
1119 CACHE_MISS.with_label_values(&[SST_META_TYPE]).inc();
1120 }
1121
1122 let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1123 if let Some(write_cache) = &self.write_cache {
1124 let file_cache = write_cache.file_cache();
1125 if self.sst_meta_cache_enabled() {
1126 if let Some(metadata) = file_cache
1127 .get_sst_meta_data(key, metrics, page_index_policy, runtime)
1128 .await
1129 {
1130 metrics.file_cache_hit += 1;
1131 let decoded = metadata.decoded();
1132 match metadata {
1133 SstMetaPreparation::Prepared(metadata) => {
1134 self.put_prepared_sst_meta(file_id, metadata, true);
1135 }
1136 SstMetaPreparation::DecodedOnly { encoding_error, .. } => {
1137 warn!(
1138 encoding_error;
1139 "Failed to encode file-cached SST metadata for memory cache, region_id: {}, file_id: {}",
1140 file_id.region_id(),
1141 file_id.file_id()
1142 );
1143 }
1144 }
1145 return Some(decoded);
1146 }
1147 } else if let Some(decoded) = file_cache
1148 .get_decoded_sst_meta_data(key, metrics, page_index_policy, runtime)
1149 .await
1150 {
1151 metrics.file_cache_hit += 1;
1152 return Some(decoded);
1153 }
1154 }
1155
1156 metrics.cache_miss += 1;
1157 None
1158 }
1159
1160 pub(crate) fn get_sst_meta_data_from_mem_cache(
1163 &self,
1164 file_id: RegionFileId,
1165 page_index_policy: PageIndexPolicy,
1166 ) -> Option<Arc<CachedSstMeta>> {
1167 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1168 let value = self
1169 .get_decoded_sst_meta(&key)
1170 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy));
1171 update_hit_miss(value, SST_META_DECODED_TYPE)
1172 }
1173
1174 pub fn get_parquet_meta_data_from_mem_cache(
1177 &self,
1178 file_id: RegionFileId,
1179 ) -> Option<Arc<ParquetMetaData>> {
1180 self.get_sst_meta_data_from_mem_cache(file_id, PageIndexPolicy::Skip)
1181 .map(|metadata| metadata.parquet_metadata())
1182 }
1183
1184 pub(crate) fn put_sst_meta_data(&self, file_id: RegionFileId, metadata: Arc<CachedSstMeta>) {
1186 if let Some(cache) = &self.sst_decoded_meta_cache {
1187 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1188 CACHE_BYTES
1189 .with_label_values(&[SST_META_DECODED_TYPE])
1190 .add(decoded_meta_cache_weight(&key, &metadata).into());
1191 cache.insert(key, metadata);
1192 }
1193 }
1194
1195 pub(crate) fn put_prepared_sst_meta(
1197 &self,
1198 file_id: RegionFileId,
1199 metadata: PreparedSstMeta,
1200 retain_decoded: bool,
1201 ) {
1202 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1203 if let Some(cache) = &self.sst_meta_cache {
1204 CACHE_BYTES
1205 .with_label_values(&[SST_META_TYPE])
1206 .add(meta_cache_weight(&key, &metadata.compact).into());
1207 cache.insert(key.clone(), metadata.compact);
1208 }
1209 if retain_decoded {
1210 self.put_sst_meta_data(file_id, metadata.decoded);
1211 }
1212 }
1213
1214 fn get_compact_sst_meta(&self, key: &SstMetaKey) -> Option<Arc<CompactSstMeta>> {
1215 self.sst_meta_cache.as_ref()?.get(key)
1216 }
1217
1218 fn get_decoded_sst_meta(&self, key: &SstMetaKey) -> Option<Arc<CachedSstMeta>> {
1219 self.sst_decoded_meta_cache.as_ref()?.get(key)
1220 }
1221
1222 pub(crate) async fn get_compact_sst_meta_data(
1227 &self,
1228 file_id: RegionFileId,
1229 page_index_policy: PageIndexPolicy,
1230 ) -> Option<Arc<CachedSstMeta>> {
1231 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1232 let compact = self
1233 .get_compact_sst_meta(&key)
1234 .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy))?;
1235 match compact.decode(&common_runtime::global_runtime()).await {
1236 Ok(metadata) => Some(metadata),
1237 Err(err) => {
1238 warn!(err; "Failed to decode compact SST metadata, region_id: {}, file_id: {}", file_id.region_id(), file_id.file_id());
1239 self.remove_parquet_meta_data(file_id);
1240 None
1241 }
1242 }
1243 }
1244
1245 pub fn put_parquet_meta_data(
1247 &self,
1248 file_id: RegionFileId,
1249 metadata: Arc<ParquetMetaData>,
1250 region_metadata: Option<RegionMetadataRef>,
1251 ) {
1252 if self.sst_decoded_meta_cache.is_some() {
1253 let file_path = format!(
1254 "region_id={}, file_id={}",
1255 file_id.region_id(),
1256 file_id.file_id()
1257 );
1258 match CachedSstMeta::try_new_with_region_metadata(
1259 &file_path,
1260 Arc::unwrap_or_clone(metadata),
1261 region_metadata,
1262 ) {
1263 Ok(metadata) => self.put_sst_meta_data(file_id, Arc::new(metadata)),
1264 Err(err) => warn!(
1265 err; "Failed to decode region metadata while caching parquet metadata, region_id: {}, file_id: {}",
1266 file_id.region_id(),
1267 file_id.file_id()
1268 ),
1269 }
1270 }
1271 }
1272
1273 pub fn remove_parquet_meta_data(&self, file_id: RegionFileId) {
1275 let key = SstMetaKey(file_id.region_id(), file_id.file_id());
1276 if let Some(cache) = &self.sst_meta_cache {
1277 cache.remove(&key);
1278 }
1279 if let Some(cache) = &self.sst_decoded_meta_cache {
1280 cache.remove(&key);
1281 }
1282 }
1283
1284 pub(crate) fn sst_meta_cache_is_full(&self) -> bool {
1286 let Some(cache) = &self.sst_meta_cache else {
1287 return true;
1288 };
1289 cache
1290 .policy()
1291 .max_capacity()
1292 .is_some_and(|capacity| cache.weighted_size() >= capacity)
1293 }
1294
1295 pub(crate) fn sst_meta_cache_enabled(&self) -> bool {
1297 self.sst_meta_cache.is_some()
1298 }
1299
1300 pub fn get_repeated_vector(
1302 &self,
1303 data_type: &ConcreteDataType,
1304 value: &Value,
1305 ) -> Option<VectorRef> {
1306 self.vector_cache.as_ref().and_then(|vector_cache| {
1307 let value = vector_cache.get(&(data_type.clone(), value.clone()));
1308 update_hit_miss(value, VECTOR_TYPE)
1309 })
1310 }
1311
1312 pub fn put_repeated_vector(&self, value: Value, vector: VectorRef) {
1314 if let Some(cache) = &self.vector_cache {
1315 let key = (vector.data_type(), value);
1316 CACHE_BYTES
1317 .with_label_values(&[VECTOR_TYPE])
1318 .add(vector_cache_weight(&key, &vector).into());
1319 cache.insert(key, vector);
1320 }
1321 }
1322
1323 pub fn get_page_ranges(
1325 &self,
1326 file_id: FileId,
1327 row_group_idx: usize,
1328 ranges: &[Range<u64>],
1329 ) -> Option<PageRangeLookup> {
1330 self.page_cache.as_ref().map(|page_cache| {
1331 let lookup = page_cache.lookup(file_id, row_group_idx, ranges);
1332 if lookup.cached_bytes > 0 {
1333 CACHE_HIT.with_label_values(&[PAGE_TYPE]).inc();
1334 }
1335 if !lookup.missing_ranges.is_empty() {
1336 CACHE_MISS.with_label_values(&[PAGE_TYPE]).inc();
1337 }
1338 lookup
1339 })
1340 }
1341
1342 pub fn put_page_ranges(
1344 &self,
1345 file_id: FileId,
1346 row_group_idx: usize,
1347 ranges: &[Range<u64>],
1348 pages: &[Bytes],
1349 ) {
1350 if let Some(cache) = &self.page_cache {
1351 cache.insert_ranges(file_id, row_group_idx, ranges, pages);
1352 }
1353 }
1354
1355 pub async fn evict_puffin_cache(&self, file_id: RegionIndexId) {
1357 if let Some(cache) = &self.bloom_filter_index_cache {
1358 cache.invalidate_file(file_id.file_id());
1359 }
1360
1361 if let Some(cache) = &self.inverted_index_cache {
1362 cache.invalidate_file(file_id.file_id());
1363 }
1364
1365 if let Some(cache) = &self.index_result_cache {
1366 cache.invalidate_file(file_id.file_id());
1367 }
1368
1369 #[cfg(feature = "vector_index")]
1370 if let Some(cache) = &self.vector_index_cache {
1371 cache.invalidate_file(file_id.file_id());
1372 }
1373
1374 if let Some(cache) = &self.puffin_metadata_cache {
1375 cache.remove(&file_id.to_string());
1376 }
1377
1378 if let Some(write_cache) = &self.write_cache {
1379 write_cache
1380 .remove(IndexKey::new(
1381 file_id.region_id(),
1382 file_id.file_id(),
1383 FileType::Puffin(file_id.version),
1384 ))
1385 .await;
1386 }
1387 }
1388
1389 pub fn get_selector_result(
1391 &self,
1392 selector_key: &SelectorResultKey,
1393 ) -> Option<Arc<SelectorResultValue>> {
1394 self.selector_result_cache
1395 .as_ref()
1396 .and_then(|selector_result_cache| selector_result_cache.get(selector_key))
1397 }
1398
1399 pub fn put_selector_result(
1401 &self,
1402 selector_key: SelectorResultKey,
1403 result: Arc<SelectorResultValue>,
1404 ) {
1405 if let Some(cache) = &self.selector_result_cache {
1406 CACHE_BYTES
1407 .with_label_values(&[SELECTOR_RESULT_TYPE])
1408 .add(selector_result_cache_weight(&selector_key, &result).into());
1409 cache.insert(selector_key, result);
1410 }
1411 }
1412
1413 #[allow(dead_code)]
1415 pub(crate) fn get_range_result(
1416 &self,
1417 key: &RangeScanCacheKey,
1418 ) -> Option<Arc<RangeScanCacheValue>> {
1419 self.range_result_cache
1420 .as_ref()
1421 .and_then(|cache| update_hit_miss(cache.get(key), RANGE_RESULT_TYPE))
1422 }
1423
1424 pub(crate) fn put_range_result(
1426 &self,
1427 key: RangeScanCacheKey,
1428 result: Arc<RangeScanCacheValue>,
1429 ) {
1430 if let Some(cache) = &self.range_result_cache {
1431 CACHE_BYTES
1432 .with_label_values(&[RANGE_RESULT_TYPE])
1433 .add(range_result_cache_weight(&key, &result).into());
1434 cache.insert(key, result);
1435 }
1436 }
1437
1438 pub(crate) fn has_range_result_cache(&self) -> bool {
1440 self.range_result_cache.is_some()
1441 }
1442
1443 pub(crate) fn range_result_memory_limiter(&self) -> &Arc<RangeResultMemoryLimiter> {
1444 &self.range_result_memory_limiter
1445 }
1446
1447 pub(crate) fn range_result_cache_size(&self) -> usize {
1448 self.range_result_cache_size as usize
1449 }
1450
1451 pub(crate) fn write_cache(&self) -> Option<&WriteCacheRef> {
1453 self.write_cache.as_ref()
1454 }
1455
1456 pub(crate) fn inverted_index_cache(&self) -> Option<&InvertedIndexCacheRef> {
1457 self.inverted_index_cache.as_ref()
1458 }
1459
1460 pub(crate) fn bloom_filter_index_cache(&self) -> Option<&BloomFilterIndexCacheRef> {
1461 self.bloom_filter_index_cache.as_ref()
1462 }
1463
1464 #[cfg(feature = "vector_index")]
1465 pub(crate) fn vector_index_cache(&self) -> Option<&VectorIndexCacheRef> {
1466 self.vector_index_cache.as_ref()
1467 }
1468
1469 pub(crate) fn puffin_metadata_cache(&self) -> Option<&PuffinMetadataCacheRef> {
1470 self.puffin_metadata_cache.as_ref()
1471 }
1472
1473 pub(crate) fn index_result_cache(&self) -> Option<&IndexResultCache> {
1474 self.index_result_cache.as_ref()
1475 }
1476
1477 pub(crate) fn get_prefilter_result(&self, key: &PrefilterKey) -> Option<Arc<BooleanBuffer>> {
1478 self.prefilter_result_cache
1479 .as_ref()
1480 .and_then(|cache| update_hit_miss(cache.get(key), PREFILTER_RESULT_TYPE))
1481 }
1482
1483 pub(crate) fn put_prefilter_result(&self, key: PrefilterKey, result: Arc<BooleanBuffer>) {
1484 if let Some(cache) = &self.prefilter_result_cache {
1485 CACHE_BYTES
1486 .with_label_values(&[PREFILTER_RESULT_TYPE])
1487 .add(prefilter_result_cache_weight(&key, &result).into());
1488 cache.insert(key, result);
1489 }
1490 }
1491}
1492
1493pub fn selector_result_cache_miss() {
1495 CACHE_MISS.with_label_values(&[SELECTOR_RESULT_TYPE]).inc()
1496}
1497
1498pub fn selector_result_cache_hit() {
1500 CACHE_HIT.with_label_values(&[SELECTOR_RESULT_TYPE]).inc()
1501}
1502
1503#[derive(Default)]
1505pub struct CacheManagerBuilder {
1506 sst_meta_cache_size: u64,
1507 vector_cache_size: u64,
1508 page_cache_size: u64,
1509 index_metadata_size: u64,
1510 index_content_size: u64,
1511 index_content_page_size: u64,
1512 index_result_cache_size: u64,
1513 prefilter_result_cache_size: u64,
1514 puffin_metadata_size: u64,
1515 write_cache: Option<WriteCacheRef>,
1516 selector_result_cache_size: u64,
1517 range_result_cache_size: u64,
1518}
1519
1520impl CacheManagerBuilder {
1521 pub fn sst_meta_cache_size(mut self, bytes: u64) -> Self {
1523 self.sst_meta_cache_size = bytes;
1524 self
1525 }
1526
1527 pub fn vector_cache_size(mut self, bytes: u64) -> Self {
1529 self.vector_cache_size = bytes;
1530 self
1531 }
1532
1533 pub fn page_cache_size(mut self, bytes: u64) -> Self {
1535 self.page_cache_size = bytes;
1536 self
1537 }
1538
1539 pub fn write_cache(mut self, cache: Option<WriteCacheRef>) -> Self {
1541 self.write_cache = cache;
1542 self
1543 }
1544
1545 pub fn index_metadata_size(mut self, bytes: u64) -> Self {
1547 self.index_metadata_size = bytes;
1548 self
1549 }
1550
1551 pub fn index_content_size(mut self, bytes: u64) -> Self {
1553 self.index_content_size = bytes;
1554 self
1555 }
1556
1557 pub fn index_content_page_size(mut self, bytes: u64) -> Self {
1559 self.index_content_page_size = bytes;
1560 self
1561 }
1562
1563 pub fn index_result_cache_size(mut self, bytes: u64) -> Self {
1565 self.index_result_cache_size = bytes;
1566 self
1567 }
1568
1569 pub fn prefilter_result_cache_size(mut self, bytes: u64) -> Self {
1571 self.prefilter_result_cache_size = bytes;
1572 self
1573 }
1574
1575 pub fn puffin_metadata_size(mut self, bytes: u64) -> Self {
1577 self.puffin_metadata_size = bytes;
1578 self
1579 }
1580
1581 pub fn selector_result_cache_size(mut self, bytes: u64) -> Self {
1583 self.selector_result_cache_size = bytes;
1584 self
1585 }
1586
1587 pub fn range_result_cache_size(mut self, bytes: u64) -> Self {
1589 self.range_result_cache_size = bytes;
1590 self
1591 }
1592
1593 pub fn build(self) -> CacheManager {
1595 let compact_meta_capacity = self.sst_meta_cache_size.div_ceil(2);
1599 let decoded_meta_capacity = self.sst_meta_cache_size / 2;
1600 let sst_meta_cache = (compact_meta_capacity != 0).then(|| {
1601 Cache::builder()
1602 .max_capacity(compact_meta_capacity)
1603 .weigher(meta_cache_weight)
1604 .eviction_listener(|k, v, cause| {
1605 let size = meta_cache_weight(&k, &v);
1606 CACHE_BYTES
1607 .with_label_values(&[SST_META_TYPE])
1608 .sub(size.into());
1609 CACHE_EVICTION
1610 .with_label_values(&[SST_META_TYPE, removal_cause_str(cause)])
1611 .inc();
1612 })
1613 .build()
1614 });
1615 let sst_decoded_meta_cache = (decoded_meta_capacity != 0).then(|| {
1616 Cache::builder()
1617 .max_capacity(decoded_meta_capacity)
1618 .weigher(decoded_meta_cache_weight)
1619 .eviction_listener(|k, v, cause| {
1620 let size = decoded_meta_cache_weight(&k, &v);
1621 CACHE_BYTES
1622 .with_label_values(&[SST_META_DECODED_TYPE])
1623 .sub(size.into());
1624 CACHE_EVICTION
1625 .with_label_values(&[SST_META_DECODED_TYPE, removal_cause_str(cause)])
1626 .inc();
1627 })
1628 .build()
1629 });
1630 let vector_cache = (self.vector_cache_size != 0).then(|| {
1631 Cache::builder()
1632 .max_capacity(self.vector_cache_size)
1633 .weigher(vector_cache_weight)
1634 .eviction_listener(|k, v, cause| {
1635 let size = vector_cache_weight(&k, &v);
1636 CACHE_BYTES
1637 .with_label_values(&[VECTOR_TYPE])
1638 .sub(size.into());
1639 CACHE_EVICTION
1640 .with_label_values(&[VECTOR_TYPE, removal_cause_str(cause)])
1641 .inc();
1642 })
1643 .build()
1644 });
1645 let page_cache =
1646 (self.page_cache_size != 0).then(|| PageRangeCache::new(self.page_cache_size));
1647 let inverted_index_cache = InvertedIndexCache::new(
1648 self.index_metadata_size,
1649 self.index_content_size,
1650 self.index_content_page_size,
1651 );
1652 let bloom_filter_index_cache = BloomFilterIndexCache::new(
1654 self.index_metadata_size,
1655 self.index_content_size,
1656 self.index_content_page_size,
1657 );
1658 #[cfg(feature = "vector_index")]
1659 let vector_index_cache = (self.index_content_size != 0)
1660 .then(|| Arc::new(VectorIndexCache::new(self.index_content_size)));
1661 let index_result_cache = (self.index_result_cache_size != 0)
1662 .then(|| IndexResultCache::new(self.index_result_cache_size));
1663 let prefilter_result_cache = (self.prefilter_result_cache_size != 0)
1664 .then(|| new_prefilter_result_cache(self.prefilter_result_cache_size));
1665 let puffin_metadata_cache =
1666 PuffinMetadataCache::new(self.puffin_metadata_size, &CACHE_BYTES);
1667 let selector_result_cache = (self.selector_result_cache_size != 0).then(|| {
1668 Cache::builder()
1669 .max_capacity(self.selector_result_cache_size)
1670 .weigher(selector_result_cache_weight)
1671 .eviction_listener(|k, v, cause| {
1672 let size = selector_result_cache_weight(&k, &v);
1673 CACHE_BYTES
1674 .with_label_values(&[SELECTOR_RESULT_TYPE])
1675 .sub(size.into());
1676 CACHE_EVICTION
1677 .with_label_values(&[SELECTOR_RESULT_TYPE, removal_cause_str(cause)])
1678 .inc();
1679 })
1680 .build()
1681 });
1682 let range_result_cache = (self.range_result_cache_size != 0).then(|| {
1683 Cache::builder()
1684 .max_capacity(self.range_result_cache_size)
1685 .weigher(range_result_cache_weight)
1686 .eviction_listener(move |k, v, cause| {
1687 let size = range_result_cache_weight(&k, &v);
1688 CACHE_BYTES
1689 .with_label_values(&[RANGE_RESULT_TYPE])
1690 .sub(size.into());
1691 CACHE_EVICTION
1692 .with_label_values(&[RANGE_RESULT_TYPE, removal_cause_str(cause)])
1693 .inc();
1694 })
1695 .build()
1696 });
1697 CacheManager {
1698 sst_meta_cache,
1699 sst_decoded_meta_cache,
1700 vector_cache,
1701 page_cache,
1702 write_cache: self.write_cache,
1703 inverted_index_cache: Some(Arc::new(inverted_index_cache)),
1704 bloom_filter_index_cache: Some(Arc::new(bloom_filter_index_cache)),
1705 #[cfg(feature = "vector_index")]
1706 vector_index_cache,
1707 puffin_metadata_cache: Some(Arc::new(puffin_metadata_cache)),
1708 selector_result_cache,
1709 range_result_cache,
1710 range_result_cache_size: self.range_result_cache_size,
1711 range_result_memory_limiter: Arc::new(RangeResultMemoryLimiter::new(
1712 self.range_result_cache_size as usize,
1713 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize,
1714 )),
1715 index_result_cache,
1716 prefilter_result_cache,
1717 }
1718 }
1719}
1720
1721fn meta_cache_weight(k: &SstMetaKey, v: &Arc<CompactSstMeta>) -> u32 {
1722 let size = k.estimated_size() + mem::size_of::<CompactSstMeta>() + v.encoded_metadata.len();
1725 u32::try_from(size).unwrap_or(u32::MAX)
1726}
1727
1728fn decoded_meta_cache_weight(k: &SstMetaKey, v: &Arc<CachedSstMeta>) -> u32 {
1729 let size = k.estimated_size() + v.parquet_metadata_size + v.region_metadata_weight;
1730 u32::try_from(size).unwrap_or(u32::MAX)
1731}
1732
1733fn vector_cache_weight(_k: &(ConcreteDataType, Value), v: &VectorRef) -> u32 {
1734 (mem::size_of::<ConcreteDataType>() + mem::size_of::<Value>() + v.memory_size()) as u32
1736}
1737
1738fn page_cache_weight(k: &PageFragmentKey, v: &Bytes) -> u32 {
1739 (k.estimated_size() + mem::size_of::<Bytes>() + v.len()) as u32
1740}
1741
1742fn selector_result_cache_weight(k: &SelectorResultKey, v: &Arc<SelectorResultValue>) -> u32 {
1743 (mem::size_of_val(k) + v.estimated_size()) as u32
1744}
1745
1746fn range_result_cache_weight(k: &RangeScanCacheKey, v: &Arc<RangeScanCacheValue>) -> u32 {
1747 (k.estimated_size() + v.estimated_size()) as u32
1748}
1749
1750fn update_hit_miss<T>(value: Option<T>, cache_type: &str) -> Option<T> {
1752 if value.is_some() {
1753 CACHE_HIT.with_label_values(&[cache_type]).inc();
1754 } else {
1755 CACHE_MISS.with_label_values(&[cache_type]).inc();
1756 }
1757 value
1758}
1759
1760#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1762struct SstMetaKey(RegionId, FileId);
1763
1764impl SstMetaKey {
1765 fn estimated_size(&self) -> usize {
1767 mem::size_of::<Self>()
1768 }
1769}
1770
1771#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1772struct PageFragmentGroupKey {
1773 file_id: FileId,
1774 row_group_idx: usize,
1775}
1776
1777#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1779pub struct PageFragmentKey {
1780 file_id: FileId,
1782 row_group_idx: usize,
1784 start: u64,
1786 end: u64,
1788}
1789
1790impl PageFragmentKey {
1791 fn new(file_id: FileId, row_group_idx: usize, range: &Range<u64>) -> PageFragmentKey {
1792 PageFragmentKey {
1793 file_id,
1794 row_group_idx,
1795 start: range.start,
1796 end: range.end,
1797 }
1798 }
1799
1800 fn group_key(&self) -> PageFragmentGroupKey {
1801 PageFragmentGroupKey {
1802 file_id: self.file_id,
1803 row_group_idx: self.row_group_idx,
1804 }
1805 }
1806
1807 fn estimated_size(&self) -> usize {
1809 mem::size_of::<Self>()
1810 }
1811}
1812
1813#[derive(Clone)]
1815pub struct PageRangePart {
1816 pub range: Range<u64>,
1818 pub bytes: Bytes,
1820}
1821
1822pub struct PageRangeLookup {
1824 pub cached_parts: Vec<Vec<PageRangePart>>,
1826 pub missing_ranges: Vec<Range<u64>>,
1828 pub cached_range_count: usize,
1830 pub cached_bytes: u64,
1832}
1833
1834impl PageRangeLookup {
1835 pub fn is_fully_cached(&self) -> bool {
1836 self.missing_ranges.is_empty()
1837 }
1838}
1839
1840type PageFragmentRangeIndex = BTreeMap<(u64, u64), PageFragmentKey>;
1841type PageFragmentIndex = HashMap<PageFragmentGroupKey, PageFragmentRangeIndex>;
1842
1843pub struct PageRangeCache {
1845 cache: Cache<PageFragmentKey, Bytes>,
1846 index: RwLock<PageFragmentIndex>,
1847}
1848
1849impl PageRangeCache {
1850 fn new(capacity: u64) -> Arc<PageRangeCache> {
1851 Arc::new_cyclic(|weak_cache: &std::sync::Weak<PageRangeCache>| {
1852 let cache = Cache::builder()
1853 .max_capacity(capacity)
1854 .weigher(page_cache_weight)
1855 .eviction_listener({
1856 let weak_cache = weak_cache.clone();
1857 move |k, v, cause| {
1858 let size = page_cache_weight(&k, &v);
1859 CACHE_BYTES.with_label_values(&[PAGE_TYPE]).sub(size.into());
1860 CACHE_EVICTION
1861 .with_label_values(&[PAGE_TYPE, removal_cause_str(cause)])
1862 .inc();
1863
1864 if let Some(cache) = weak_cache.upgrade()
1865 && !matches!(cause, RemovalCause::Replaced)
1866 {
1867 cache.remove_index_entry(*k);
1868 }
1869 }
1870 })
1871 .build();
1872
1873 PageRangeCache {
1874 cache,
1875 index: RwLock::new(HashMap::new()),
1876 }
1877 })
1878 }
1879
1880 fn lookup(
1881 &self,
1882 file_id: FileId,
1883 row_group_idx: usize,
1884 ranges: &[Range<u64>],
1885 ) -> PageRangeLookup {
1886 let mut cached_parts = Vec::with_capacity(ranges.len());
1887 let mut missing_ranges = Vec::new();
1888 let mut cached_range_count = 0;
1889 let mut cached_bytes = 0;
1890
1891 for range in ranges {
1892 if range.start >= range.end {
1893 cached_parts.push(Vec::new());
1894 continue;
1895 }
1896
1897 let mut parts = Vec::new();
1898 let candidates = self.find_index_candidates(file_id, row_group_idx, range);
1899 let mut stale_keys = Vec::new();
1900
1901 for fragment_key in candidates {
1902 if let Some(bytes) = self.cache.get(&fragment_key) {
1903 let part_start = range.start.max(fragment_key.start);
1904 let part_end = range.end.min(fragment_key.end);
1905 let slice_start = (part_start - fragment_key.start) as usize;
1906 let slice_end = (part_end - fragment_key.start) as usize;
1907 parts.push(PageRangePart {
1908 range: part_start..part_end,
1909 bytes: bytes.slice(slice_start..slice_end),
1910 });
1911 } else {
1912 stale_keys.push(fragment_key);
1913 }
1914 }
1915 for key in stale_keys {
1916 self.remove_uncached_index_entry(key);
1917 }
1918
1919 let mut cursor = range.start;
1920 let mut compacted_parts: Vec<PageRangePart> = Vec::with_capacity(parts.len());
1921 for part in parts {
1922 if part.range.end <= cursor {
1923 continue;
1924 }
1925
1926 let part = if part.range.start < cursor {
1927 let offset = (cursor - part.range.start) as usize;
1928 PageRangePart {
1929 range: cursor..part.range.end,
1930 bytes: part.bytes.slice(offset..),
1931 }
1932 } else {
1933 part
1934 };
1935
1936 if cursor < part.range.start {
1937 missing_ranges.push(cursor..part.range.start);
1938 }
1939 cached_bytes += part.range.end - part.range.start;
1940 cached_range_count += 1;
1941 cursor = part.range.end;
1942 compacted_parts.push(part);
1943
1944 if cursor >= range.end {
1945 break;
1946 }
1947 }
1948
1949 if cursor < range.end {
1950 missing_ranges.push(cursor..range.end);
1951 }
1952 cached_parts.push(compacted_parts);
1953 }
1954
1955 PageRangeLookup {
1956 cached_parts,
1957 missing_ranges,
1958 cached_range_count,
1959 cached_bytes,
1960 }
1961 }
1962
1963 fn insert_ranges(
1964 &self,
1965 file_id: FileId,
1966 row_group_idx: usize,
1967 ranges: &[Range<u64>],
1968 pages: &[Bytes],
1969 ) {
1970 for (range, bytes) in ranges.iter().zip(pages) {
1971 if range.start >= range.end || bytes.len() as u64 != range.end - range.start {
1972 continue;
1973 }
1974
1975 let key = PageFragmentKey::new(file_id, row_group_idx, range);
1976 let bytes = Bytes::copy_from_slice(bytes.as_ref());
1977 let size = page_cache_weight(&key, &bytes);
1978 CACHE_BYTES.with_label_values(&[PAGE_TYPE]).add(size.into());
1979 self.cache.insert(key, bytes);
1980 let mut index = self.index.write().unwrap();
1981 index
1982 .entry(key.group_key())
1983 .or_default()
1984 .insert((key.start, key.end), key);
1985 }
1986 }
1987
1988 fn find_index_candidates(
1989 &self,
1990 file_id: FileId,
1991 row_group_idx: usize,
1992 range: &Range<u64>,
1993 ) -> Vec<PageFragmentKey> {
1994 let group_key = PageFragmentGroupKey {
1995 file_id,
1996 row_group_idx,
1997 };
1998 let index = self.index.read().unwrap();
1999 index
2000 .get(&group_key)
2001 .map(|ranges| {
2002 ranges
2003 .range(..(range.end, 0))
2004 .filter_map(|(_, fragment_key)| {
2005 (fragment_key.end > range.start).then_some(*fragment_key)
2006 })
2007 .collect()
2008 })
2009 .unwrap_or_default()
2010 }
2011
2012 fn remove_uncached_index_entry(&self, key: PageFragmentKey) {
2013 let group_key = key.group_key();
2014 let mut index = self.index.write().unwrap();
2015 if self.cache.contains_key(&key) {
2016 return;
2017 }
2018
2019 Self::remove_index_entry_locked(&mut index, group_key, key);
2020 }
2021
2022 fn remove_index_entry(&self, key: PageFragmentKey) {
2023 let group_key = key.group_key();
2024 let mut index = self.index.write().unwrap();
2025 Self::remove_index_entry_locked(&mut index, group_key, key);
2026 }
2027
2028 fn remove_index_entry_locked(
2029 index: &mut PageFragmentIndex,
2030 group_key: PageFragmentGroupKey,
2031 key: PageFragmentKey,
2032 ) {
2033 let Some(ranges) = index.get_mut(&group_key) else {
2034 return;
2035 };
2036
2037 let removed = ranges
2038 .get(&(key.start, key.end))
2039 .is_some_and(|current| current == &key);
2040 if removed {
2041 ranges.remove(&(key.start, key.end));
2042 }
2043 if ranges.is_empty() {
2044 index.remove(&group_key);
2045 }
2046 }
2047}
2048
2049#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
2051pub struct SelectorResultKey {
2052 pub file_id: FileId,
2054 pub row_group_idx: usize,
2056 pub selector: TimeSeriesRowSelector,
2058}
2059
2060pub enum SelectorResult {
2062 PrimaryKey(Vec<Batch>),
2064 Flat(Vec<RecordBatch>),
2066}
2067
2068pub struct SelectorResultValue {
2070 pub result: SelectorResult,
2072 pub read_cols: ParquetReadColumns,
2074 pub json_target_types: JsonTargetTypes,
2079}
2080
2081impl SelectorResultValue {
2082 pub fn new(result: Vec<Batch>, read_cols: ParquetReadColumns) -> SelectorResultValue {
2084 SelectorResultValue {
2085 result: SelectorResult::PrimaryKey(result),
2086 read_cols,
2087 json_target_types: Arc::default(),
2088 }
2089 }
2090
2091 pub fn new_flat(
2093 result: Vec<RecordBatch>,
2094 read_cols: ParquetReadColumns,
2095 json_target_types: JsonTargetTypes,
2096 ) -> SelectorResultValue {
2097 SelectorResultValue {
2098 result: SelectorResult::Flat(result),
2099 read_cols,
2100 json_target_types,
2101 }
2102 }
2103
2104 fn estimated_size(&self) -> usize {
2106 let result_size: usize = match &self.result {
2107 SelectorResult::PrimaryKey(batches) => {
2108 batches.iter().map(|batch| batch.memory_size()).sum()
2109 }
2110 SelectorResult::Flat(batches) => batches.iter().map(record_batch_estimated_size).sum(),
2111 };
2112 result_size
2113 + self.json_target_types.len() * (size_of::<ColumnId>() + size_of::<JsonNativeType>())
2114 }
2115}
2116
2117type SstMetaCache = Cache<SstMetaKey, Arc<CompactSstMeta>>;
2119type SstDecodedMetaCache = Cache<SstMetaKey, Arc<CachedSstMeta>>;
2121type VectorCache = Cache<(ConcreteDataType, Value), VectorRef>;
2125type SelectorResultCache = Cache<SelectorResultKey, Arc<SelectorResultValue>>;
2127type RangeResultCache = Cache<RangeScanCacheKey, Arc<RangeScanCacheValue>>;
2129
2130#[cfg(test)]
2131mod tests {
2132 use std::sync::Arc;
2133
2134 use api::v1::SemanticType;
2135 use api::v1::index::{BloomFilterMeta, InvertedIndexMetas};
2136 use datatypes::schema::ColumnSchema;
2137 use datatypes::vectors::Int64Vector;
2138 use parquet::file::page_index::offset_index::{OffsetIndexMetaData, PageLocation};
2139 use puffin::file_metadata::FileMetadata;
2140 use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
2141 use store_api::storage::ColumnId;
2142
2143 use super::*;
2144 use crate::cache::index::bloom_filter_index::Tag;
2145 use crate::cache::index::result_cache::PredicateKey;
2146 use crate::cache::test_util::{
2147 parquet_meta, sst_parquet_meta, sst_parquet_meta_with_region_metadata,
2148 };
2149 use crate::read::range_cache::{
2150 RangeScanCacheKey, RangeScanCacheValue, ScanRequestFingerprintBuilder,
2151 };
2152 use crate::read::read_columns::ReadColumns;
2153 use crate::sst::parquet::row_selection::RowGroupSelection;
2154
2155 #[tokio::test]
2156 async fn test_disable_cache() {
2157 let cache = CacheManager::default();
2158 assert!(cache.sst_meta_cache.is_none());
2159 assert!(cache.sst_decoded_meta_cache.is_none());
2160 assert!(cache.vector_cache.is_none());
2161 assert!(cache.page_cache.is_none());
2162
2163 let region_id = RegionId::new(1, 1);
2164 let file_id = RegionFileId::new(region_id, FileId::random());
2165 let metadata = parquet_meta();
2166 let mut metrics = MetadataCacheMetrics::default();
2167 cache.put_parquet_meta_data(file_id, metadata, None);
2168 assert!(
2169 cache
2170 .get_sst_meta_data(
2171 file_id,
2172 &mut metrics,
2173 Default::default(),
2174 &common_runtime::global_runtime()
2175 )
2176 .await
2177 .is_none()
2178 );
2179
2180 let value = Value::Int64(10);
2181 let vector: VectorRef = Arc::new(Int64Vector::from_slice([10, 10, 10, 10]));
2182 cache.put_repeated_vector(value.clone(), vector.clone());
2183 assert!(
2184 cache
2185 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2186 .is_none()
2187 );
2188
2189 cache.put_page_ranges(
2190 file_id.file_id(),
2191 1,
2192 &[Range { start: 0, end: 5 }],
2193 &[Bytes::from_static(b"abcde")],
2194 );
2195 assert!(
2196 cache
2197 .get_page_ranges(file_id.file_id(), 1, &[Range { start: 0, end: 5 }])
2198 .is_none()
2199 );
2200
2201 assert!(cache.write_cache().is_none());
2202 }
2203
2204 #[test]
2205 fn test_sst_meta_cache_splits_capacity() {
2206 let cache = CacheManager::builder().sst_meta_cache_size(101).build();
2207
2208 assert_eq!(
2209 Some(51),
2210 cache
2211 .sst_meta_cache
2212 .as_ref()
2213 .unwrap()
2214 .policy()
2215 .max_capacity()
2216 );
2217 assert_eq!(
2218 Some(50),
2219 cache
2220 .sst_decoded_meta_cache
2221 .as_ref()
2222 .unwrap()
2223 .policy()
2224 .max_capacity()
2225 );
2226 }
2227
2228 #[tokio::test]
2229 async fn test_parquet_meta_cache() {
2230 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2231 let mut metrics = MetadataCacheMetrics::default();
2232 let region_id = RegionId::new(1, 1);
2233 let file_id = RegionFileId::new(region_id, FileId::random());
2234 assert!(
2235 cache
2236 .get_sst_meta_data(
2237 file_id,
2238 &mut metrics,
2239 Default::default(),
2240 &common_runtime::global_runtime()
2241 )
2242 .await
2243 .is_none()
2244 );
2245 let (metadata, region_metadata) = sst_parquet_meta();
2246 cache.put_parquet_meta_data(file_id, metadata, None);
2247 let cached = cache
2248 .get_sst_meta_data(
2249 file_id,
2250 &mut metrics,
2251 Default::default(),
2252 &common_runtime::global_runtime(),
2253 )
2254 .await
2255 .unwrap();
2256 assert_eq!(region_metadata, cached.region_metadata());
2257 assert!(
2258 cached
2259 .parquet_metadata()
2260 .file_metadata()
2261 .key_value_metadata()
2262 .is_none_or(|key_values| {
2263 key_values
2264 .iter()
2265 .all(|key_value| key_value.key != PARQUET_METADATA_KEY)
2266 })
2267 );
2268 cache.remove_parquet_meta_data(file_id);
2269 assert!(
2270 cache
2271 .get_sst_meta_data(
2272 file_id,
2273 &mut metrics,
2274 Default::default(),
2275 &common_runtime::global_runtime()
2276 )
2277 .await
2278 .is_none()
2279 );
2280 }
2281
2282 #[tokio::test]
2283 async fn test_parquet_meta_cache_with_provided_region_metadata() {
2284 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2285 let mut metrics = MetadataCacheMetrics::default();
2286 let region_id = RegionId::new(1, 1);
2287 let file_id = RegionFileId::new(region_id, FileId::random());
2288 let (metadata, region_metadata) = sst_parquet_meta();
2289
2290 cache.put_parquet_meta_data(file_id, metadata, Some(region_metadata.clone()));
2291
2292 let cached = cache
2293 .get_sst_meta_data(
2294 file_id,
2295 &mut metrics,
2296 Default::default(),
2297 &common_runtime::global_runtime(),
2298 )
2299 .await
2300 .unwrap();
2301 assert!(Arc::ptr_eq(®ion_metadata, &cached.region_metadata()));
2302 }
2303
2304 #[tokio::test]
2305 async fn test_compact_sst_meta_round_trip() {
2306 let cache = CacheManager::builder()
2307 .sst_meta_cache_size(1024 * 1024)
2308 .build();
2309 let region_id = RegionId::new(1, 1);
2310 let file_id = RegionFileId::new(region_id, FileId::random());
2311 let (metadata, expected_region_metadata) = sst_parquet_meta();
2312 let metadata = Arc::unwrap_or_clone(metadata);
2313 let offset_indexes = metadata
2314 .row_groups()
2315 .iter()
2316 .map(|row_group| {
2317 row_group
2318 .columns()
2319 .iter()
2320 .map(|_| OffsetIndexMetaData {
2321 page_locations: vec![PageLocation {
2322 offset: 42,
2323 compressed_page_size: 128,
2324 first_row_index: 0,
2325 }],
2326 unencoded_byte_array_data_bytes: None,
2327 })
2328 .collect()
2329 })
2330 .collect();
2331 let metadata = metadata
2332 .into_builder()
2333 .set_offset_index(Some(offset_indexes))
2334 .build();
2335 let expected_rows = metadata.file_metadata().num_rows();
2336 let prepared = prepare_sst_meta(
2337 "test.parquet",
2338 metadata,
2339 None,
2340 PageIndexPolicy::Required,
2341 &common_runtime::global_runtime(),
2342 )
2343 .await
2344 .unwrap();
2345 let SstMetaPreparation::Prepared(prepared) = prepared else {
2346 panic!("valid metadata should produce a compact cache entry");
2347 };
2348
2349 cache.put_prepared_sst_meta(file_id, prepared, false);
2351 let mut metrics = MetadataCacheMetrics::default();
2352 let cached = cache
2353 .get_sst_meta_data(
2354 file_id,
2355 &mut metrics,
2356 PageIndexPolicy::Required,
2357 &common_runtime::global_runtime(),
2358 )
2359 .await
2360 .unwrap();
2361
2362 assert_eq!(1, metrics.mem_cache_hit);
2363 assert_eq!(0, metrics.cache_miss);
2364 assert_eq!(
2365 expected_rows,
2366 cached.parquet_metadata.file_metadata().num_rows()
2367 );
2368 assert!(cached.parquet_metadata.offset_index().is_some());
2369 assert_eq!(expected_region_metadata, cached.region_metadata());
2370 assert!(
2371 cached
2372 .parquet_metadata
2373 .file_metadata()
2374 .key_value_metadata()
2375 .is_none_or(|key_values| key_values
2376 .iter()
2377 .all(|key_value| key_value.key != PARQUET_METADATA_KEY))
2378 );
2379 }
2380
2381 #[test]
2382 fn test_cache_encoding_failure_preserves_decoded_metadata() {
2383 let (metadata, expected_region_metadata) = sst_parquet_meta();
2384 let expected_rows = metadata.file_metadata().num_rows();
2385 let cache_encoding: Result<(Bytes, usize)> = Err(UnexpectedSnafu {
2386 reason: "injected compact metadata encoding failure",
2387 }
2388 .build());
2389
2390 let preparation = finish_sst_meta_preparation(
2391 "test.parquet",
2392 Arc::unwrap_or_clone(metadata),
2393 None,
2394 PageIndexPolicy::Skip,
2395 cache_encoding,
2396 )
2397 .unwrap();
2398 let SstMetaPreparation::DecodedOnly {
2399 decoded,
2400 encoding_error,
2401 } = preparation
2402 else {
2403 panic!("cache encoding failure should return decoded-only metadata");
2404 };
2405
2406 assert_eq!(
2407 expected_rows,
2408 decoded.parquet_metadata.file_metadata().num_rows()
2409 );
2410 assert_eq!(expected_region_metadata, decoded.region_metadata());
2411 assert!(
2412 encoding_error
2413 .to_string()
2414 .contains("injected compact metadata encoding failure")
2415 );
2416 }
2417
2418 #[tokio::test]
2419 async fn test_parquet_meta_cache_respects_page_index_policy() {
2420 let cache = CacheManager::builder().sst_meta_cache_size(2000).build();
2421 let region_id = RegionId::new(1, 1);
2422 let file_id = RegionFileId::new(region_id, FileId::random());
2423 let (metadata, _) = sst_parquet_meta();
2424
2425 let skip_metadata = Arc::new(
2426 CachedSstMeta::try_new_with_page_index_policy(
2427 "test.parquet",
2428 Arc::unwrap_or_clone(metadata.clone()),
2429 None,
2430 PageIndexPolicy::Skip,
2431 )
2432 .unwrap(),
2433 );
2434 cache.put_sst_meta_data(file_id, skip_metadata);
2435
2436 let mut metrics = MetadataCacheMetrics::default();
2437 assert!(
2438 cache
2439 .get_sst_meta_data(
2440 file_id,
2441 &mut metrics,
2442 PageIndexPolicy::Optional,
2443 &common_runtime::global_runtime()
2444 )
2445 .await
2446 .is_none()
2447 );
2448 assert_eq!(1, metrics.cache_miss);
2449
2450 let optional_metadata = Arc::new(
2451 CachedSstMeta::try_new_with_page_index_policy(
2452 "test.parquet",
2453 Arc::unwrap_or_clone(metadata),
2454 None,
2455 PageIndexPolicy::Optional,
2456 )
2457 .unwrap(),
2458 );
2459 cache.put_sst_meta_data(file_id, optional_metadata);
2460
2461 let mut metrics = MetadataCacheMetrics::default();
2462 assert!(
2463 cache
2464 .get_sst_meta_data(
2465 file_id,
2466 &mut metrics,
2467 PageIndexPolicy::Optional,
2468 &common_runtime::global_runtime()
2469 )
2470 .await
2471 .is_some()
2472 );
2473 assert_eq!(1, metrics.mem_cache_hit);
2474
2475 let mut metrics = MetadataCacheMetrics::default();
2476 assert!(
2477 cache
2478 .get_sst_meta_data(
2479 file_id,
2480 &mut metrics,
2481 PageIndexPolicy::Skip,
2482 &common_runtime::global_runtime()
2483 )
2484 .await
2485 .is_some()
2486 );
2487 assert_eq!(1, metrics.mem_cache_hit);
2488 }
2489
2490 #[test]
2491 fn test_decoded_meta_cache_weight_accounts_for_region_metadata() {
2492 let region_metadata = Arc::new(wide_region_metadata(128));
2493 let json_len = region_metadata.to_json().unwrap().len();
2494 let metadata = sst_parquet_meta_with_region_metadata(region_metadata.clone());
2495 let cached = Arc::new(
2496 CachedSstMeta::try_new("test.parquet", Arc::unwrap_or_clone(metadata)).unwrap(),
2497 );
2498 let key = SstMetaKey(region_metadata.region_id, FileId::random());
2499
2500 assert!(cached.region_metadata_weight > json_len);
2501 assert_eq!(
2502 decoded_meta_cache_weight(&key, &cached) as usize,
2503 key.estimated_size() + cached.parquet_metadata_size + cached.region_metadata_weight
2504 );
2505 assert_eq!(
2506 cached.parquet_metadata_size,
2507 parquet_meta_size(&cached.parquet_metadata)
2508 );
2509 }
2510
2511 #[test]
2512 fn test_decoded_meta_cache_weight_saturates_on_overflow() {
2513 let region_metadata = Arc::new(wide_region_metadata(1));
2514 let metadata = sst_parquet_meta_with_region_metadata(region_metadata.clone());
2515 let mut cached =
2516 CachedSstMeta::try_new("test.parquet", Arc::unwrap_or_clone(metadata)).unwrap();
2517 cached.region_metadata_weight = u32::MAX as usize + 1;
2518 let cached = Arc::new(cached);
2519 let key = SstMetaKey(region_metadata.region_id, FileId::random());
2520
2521 assert_eq!(u32::MAX, decoded_meta_cache_weight(&key, &cached));
2522 }
2523
2524 #[test]
2525 fn test_repeated_vector_cache() {
2526 let cache = CacheManager::builder().vector_cache_size(4096).build();
2527 let value = Value::Int64(10);
2528 assert!(
2529 cache
2530 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2531 .is_none()
2532 );
2533 let vector: VectorRef = Arc::new(Int64Vector::from_slice([10, 10, 10, 10]));
2534 cache.put_repeated_vector(value.clone(), vector.clone());
2535 let cached = cache
2536 .get_repeated_vector(&ConcreteDataType::int64_datatype(), &value)
2537 .unwrap();
2538 assert_eq!(vector, cached);
2539 }
2540
2541 #[test]
2542 fn test_page_cache() {
2543 let cache = CacheManager::builder().page_cache_size(1000).build();
2544 let file_id = FileId::random();
2545 let uncached = 0..10;
2546 assert_eq!(
2547 vec![0..10],
2548 cache
2549 .get_page_ranges(file_id, 0, std::slice::from_ref(&uncached))
2550 .unwrap()
2551 .missing_ranges
2552 );
2553
2554 let cached = 100..500;
2555 cache.put_page_ranges(
2556 file_id,
2557 0,
2558 std::slice::from_ref(&cached),
2559 &[Bytes::from(vec![7; 400])],
2560 );
2561
2562 let subrange = 200..300;
2563 let lookup = cache
2564 .get_page_ranges(file_id, 0, std::slice::from_ref(&subrange))
2565 .unwrap();
2566 assert!(lookup.is_fully_cached());
2567 assert_eq!(100, lookup.cached_bytes);
2568 assert_eq!(1, lookup.cached_parts.len());
2569 assert_eq!(200..300, lookup.cached_parts[0][0].range);
2570 assert_eq!(100, lookup.cached_parts[0][0].bytes.len());
2571
2572 let overlapping = 400..600;
2573 let lookup = cache
2574 .get_page_ranges(file_id, 0, std::slice::from_ref(&overlapping))
2575 .unwrap();
2576 assert!(!lookup.is_fully_cached());
2577 assert_eq!(100, lookup.cached_bytes);
2578 assert_eq!(vec![500..600], lookup.missing_ranges);
2579 assert_eq!(400..500, lookup.cached_parts[0][0].range);
2580 }
2581
2582 #[test]
2583 fn test_page_cache_detaches_fragment_bytes() {
2584 let cache = PageRangeCache::new(1000);
2585 let file_id = FileId::random();
2586 let backing = Bytes::from(vec![1; 1024]);
2587 let page = backing.slice(512..522);
2588 let page_ptr = page.as_ptr();
2589 let range = 0..10;
2590
2591 cache.insert_ranges(
2592 file_id,
2593 0,
2594 std::slice::from_ref(&range),
2595 std::slice::from_ref(&page),
2596 );
2597
2598 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2599 assert!(lookup.is_fully_cached());
2600 assert_eq!(1, lookup.cached_parts[0].len());
2601 assert_eq!(&page[..], &lookup.cached_parts[0][0].bytes[..]);
2602 assert_ne!(page_ptr, lookup.cached_parts[0][0].bytes.as_ptr());
2603 }
2604
2605 #[test]
2606 fn test_page_cache_replaces_fragment() {
2607 let cache = PageRangeCache::new(1000);
2608 let file_id = FileId::random();
2609 let range = 0..10;
2610
2611 cache.insert_ranges(
2612 file_id,
2613 0,
2614 std::slice::from_ref(&range),
2615 &[Bytes::from(vec![1; 10])],
2616 );
2617 cache.insert_ranges(
2618 file_id,
2619 0,
2620 std::slice::from_ref(&range),
2621 &[Bytes::from(vec![2; 10])],
2622 );
2623 cache.cache.run_pending_tasks();
2624 assert_eq!(
2625 vec![PageFragmentKey::new(file_id, 0, &range)],
2626 cache.find_index_candidates(file_id, 0, &range)
2627 );
2628
2629 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2630 assert!(lookup.is_fully_cached());
2631 assert_eq!(&vec![2; 10][..], &lookup.cached_parts[0][0].bytes[..]);
2632 }
2633
2634 #[test]
2635 fn test_page_cache_retains_disjoint_inserts_for_same_row_group() {
2636 let cache = PageRangeCache::new(1000);
2637 let file_id = FileId::random();
2638 let range1 = 0..10;
2639 let range2 = 20..30;
2640
2641 cache.insert_ranges(
2642 file_id,
2643 0,
2644 std::slice::from_ref(&range1),
2645 &[Bytes::from(vec![1; 10])],
2646 );
2647 cache.insert_ranges(
2648 file_id,
2649 0,
2650 std::slice::from_ref(&range2),
2651 &[Bytes::from(vec![2; 10])],
2652 );
2653
2654 let lookup = cache.lookup(file_id, 0, &[range1, range2]);
2655 assert!(lookup.is_fully_cached());
2656 assert_eq!(2, lookup.cached_range_count);
2657 assert_eq!(&vec![1; 10][..], &lookup.cached_parts[0][0].bytes[..]);
2658 assert_eq!(&vec![2; 10][..], &lookup.cached_parts[1][0].bytes[..]);
2659 }
2660
2661 #[test]
2662 fn test_page_cache_fragment_eviction() {
2663 let file_id = FileId::random();
2664 let range = 0..10;
2665 let key = PageFragmentKey::new(file_id, 0, &range);
2666 let page = Bytes::from(vec![1; 10]);
2667 let cache = PageRangeCache::new(page_cache_weight(&key, &page) as u64);
2668
2669 cache.insert_ranges(
2670 file_id,
2671 0,
2672 std::slice::from_ref(&range),
2673 &[Bytes::from(vec![1; 10])],
2674 );
2675 assert!(
2676 cache
2677 .lookup(file_id, 0, std::slice::from_ref(&range))
2678 .is_fully_cached()
2679 );
2680
2681 cache.cache.invalidate(&key);
2682 cache.cache.run_pending_tasks();
2683 assert!(cache.find_index_candidates(file_id, 0, &range).is_empty());
2684
2685 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2686 assert!(!lookup.is_fully_cached());
2687 assert_eq!(vec![0..10], lookup.missing_ranges);
2688 }
2689
2690 #[test]
2691 fn test_page_cache_rejects_oversized_fragment() {
2692 let cache = PageRangeCache::new(1);
2693 let file_id = FileId::random();
2694 let range = 0..10;
2695
2696 cache.insert_ranges(
2697 file_id,
2698 0,
2699 std::slice::from_ref(&range),
2700 &[Bytes::from(vec![1; 10])],
2701 );
2702 cache.cache.run_pending_tasks();
2703
2704 let lookup = cache.lookup(file_id, 0, std::slice::from_ref(&range));
2705 assert!(!lookup.is_fully_cached());
2706 assert_eq!(vec![0..10], lookup.missing_ranges);
2707 }
2708
2709 #[test]
2710 fn test_selector_result_cache() {
2711 let cache = CacheManager::builder()
2712 .selector_result_cache_size(1000)
2713 .build();
2714 let file_id = FileId::random();
2715 let key = SelectorResultKey {
2716 file_id,
2717 row_group_idx: 0,
2718 selector: TimeSeriesRowSelector::LastRow { after_merge: false },
2719 };
2720 assert!(cache.get_selector_result(&key).is_none());
2721 let result = Arc::new(SelectorResultValue::new(
2722 Vec::new(),
2723 ParquetReadColumns::from_deduped(Vec::new()),
2724 ));
2725 cache.put_selector_result(key, result);
2726 assert!(cache.get_selector_result(&key).is_some());
2727 }
2728
2729 #[test]
2730 fn test_prefilter_result_cache() {
2731 let disabled = CacheManager::builder().build();
2732 let file_id = FileId::random();
2733 let key = PrefilterKey::new(
2734 file_id,
2735 0,
2736 None,
2737 1,
2738 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()]),
2739 );
2740 let selection = Arc::new(BooleanBuffer::new_set(3));
2741
2742 disabled.put_prefilter_result(key.clone(), selection.clone());
2743 assert!(disabled.get_prefilter_result(&key).is_none());
2744
2745 let cache = Arc::new(
2746 CacheManager::builder()
2747 .prefilter_result_cache_size(1000)
2748 .build(),
2749 );
2750 assert!(cache.get_prefilter_result(&key).is_none());
2751 cache.put_prefilter_result(key.clone(), selection.clone());
2752 assert_eq!(
2753 cache.get_prefilter_result(&key).unwrap().as_ref(),
2754 selection.as_ref()
2755 );
2756
2757 let enable_all = CacheStrategy::EnableAll(cache.clone());
2758 assert!(enable_all.get_prefilter_result(&key).is_some());
2759
2760 let compaction = CacheStrategy::Compaction(cache.clone());
2761 assert!(compaction.get_prefilter_result(&key).is_none());
2762 compaction.put_prefilter_result(key.clone(), selection.clone());
2763 assert!(cache.get_prefilter_result(&key).is_some());
2764
2765 let disabled_strategy = CacheStrategy::Disabled;
2766 assert!(disabled_strategy.get_prefilter_result(&key).is_none());
2767 disabled_strategy.put_prefilter_result(key.clone(), selection);
2768 assert!(cache.get_prefilter_result(&key).is_some());
2769 }
2770
2771 #[test]
2772 fn test_prefilter_key_distinguishes_dimensions() {
2773 let file_id = FileId::random();
2774 let row_selection = RowSelection::from(vec![RowSelector::skip(1), RowSelector::select(3)]);
2775 let other_row_selection =
2776 RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(2)]);
2777 let row_selection = PrefilterKey::row_selection_snapshot(Some(&row_selection));
2778 let other_row_selection = PrefilterKey::row_selection_snapshot(Some(&other_row_selection));
2779 let base = PrefilterKey::new(
2780 file_id,
2781 0,
2782 row_selection.clone(),
2783 1,
2784 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()]),
2785 );
2786
2787 assert_ne!(
2788 base,
2789 PrefilterKey::new(
2790 FileId::random(),
2791 0,
2792 row_selection.clone(),
2793 1,
2794 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2795 )
2796 );
2797 assert_ne!(
2798 base,
2799 PrefilterKey::new(
2800 file_id,
2801 1,
2802 row_selection.clone(),
2803 1,
2804 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2805 )
2806 );
2807 assert_ne!(
2808 base,
2809 PrefilterKey::new(
2810 file_id,
2811 0,
2812 other_row_selection,
2813 1,
2814 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2815 )
2816 );
2817 assert_ne!(
2818 base,
2819 PrefilterKey::new(
2820 file_id,
2821 0,
2822 row_selection.clone(),
2823 1,
2824 SmallVec::from_vec(vec!["tag_0 IN ([b])".to_string()])
2825 )
2826 );
2827 assert_ne!(
2828 base,
2829 PrefilterKey::new(
2830 file_id,
2831 0,
2832 row_selection.clone(),
2833 2,
2834 SmallVec::from_vec(vec!["tag_0 IN ([a])".to_string()])
2835 )
2836 );
2837 let pk_group = PrefilterKey::new(
2838 file_id,
2839 0,
2840 row_selection,
2841 1,
2842 SmallVec::from_vec(vec![
2843 "tag_0 IN ([a])".to_string(),
2844 "tag_1 IN ([x])".to_string(),
2845 ]),
2846 );
2847 assert_ne!(base, pk_group);
2848 }
2849
2850 #[test]
2851 fn test_range_result_cache() {
2852 let cache = Arc::new(
2853 CacheManager::builder()
2854 .range_result_cache_size(1024 * 1024)
2855 .build(),
2856 );
2857
2858 let key = RangeScanCacheKey {
2859 region_id: RegionId::new(1, 1),
2860 row_groups: vec![(FileId::random(), 0)],
2861 scan: ScanRequestFingerprintBuilder {
2862 read_columns: ReadColumns::new(std::iter::empty()),
2863 read_column_types: vec![],
2864 filters: vec!["tag_0 = 1".to_string()],
2865 time_filters: vec![],
2866 series_row_selector: None,
2867 append_mode: false,
2868 filter_deleted: true,
2869 merge_mode: crate::region::options::MergeMode::LastRow,
2870 sequence_range: None,
2871 partition_expr_version: 0,
2872 }
2873 .build(),
2874 };
2875 let value = Arc::new(RangeScanCacheValue::new(Vec::new(), 0));
2876
2877 assert!(cache.get_range_result(&key).is_none());
2878 cache.put_range_result(key.clone(), value.clone());
2879 assert!(cache.get_range_result(&key).is_some());
2880
2881 let enable_all = CacheStrategy::EnableAll(cache.clone());
2882 assert!(enable_all.get_range_result(&key).is_some());
2883
2884 let compaction = CacheStrategy::Compaction(cache.clone());
2885 assert!(compaction.get_range_result(&key).is_none());
2886 compaction.put_range_result(key.clone(), value.clone());
2887 assert!(cache.get_range_result(&key).is_some());
2888
2889 let disabled = CacheStrategy::Disabled;
2890 assert!(disabled.get_range_result(&key).is_none());
2891 disabled.put_range_result(key.clone(), value);
2892 assert!(cache.get_range_result(&key).is_some());
2893 }
2894
2895 #[test]
2896 fn test_range_result_cache_size_configures_limiter() {
2897 let cache_size = 3 * 1024_u64;
2898 let cache = CacheManager::builder()
2899 .range_result_cache_size(cache_size)
2900 .build();
2901
2902 assert_eq!(cache.range_result_cache_size(), cache_size as usize);
2903 assert_eq!(
2904 cache.range_result_memory_limiter().permit_bytes(),
2905 RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize
2906 );
2907 assert_eq!(
2908 cache.range_result_memory_limiter().available_permits(),
2909 (cache_size as usize).div_ceil(RANGE_RESULT_CONCAT_MEMORY_PERMIT.as_bytes() as usize)
2910 );
2911 }
2912
2913 #[tokio::test]
2914 async fn range_result_memory_limiter_rejects_oversized_request() {
2915 let limiter = RangeResultMemoryLimiter::new(2 * 1024, 1024);
2916 assert_eq!(limiter.available_permits(), 2);
2917
2918 let err = limiter.acquire(10 * 1024).await.unwrap_err();
2919 assert!(
2920 err.to_string().contains("exceeds limiter capacity"),
2921 "unexpected error: {err}"
2922 );
2923 assert_eq!(limiter.available_permits(), 2);
2924 }
2925
2926 #[tokio::test]
2927 async fn range_result_memory_limiter_allows_request_up_to_capacity() {
2928 let limiter = RangeResultMemoryLimiter::new(2 * 1024, 1024);
2929 let permit = limiter.acquire(2 * 1024).await.unwrap();
2930 assert_eq!(limiter.available_permits(), 0);
2931 drop(permit);
2932 assert_eq!(limiter.available_permits(), 2);
2933 }
2934
2935 #[tokio::test]
2936 async fn test_evict_puffin_cache_clears_all_entries() {
2937 use std::collections::{BTreeMap, HashMap};
2938
2939 let cache = CacheManager::builder()
2940 .index_metadata_size(128)
2941 .index_content_size(128)
2942 .index_content_page_size(64)
2943 .index_result_cache_size(128)
2944 .puffin_metadata_size(128)
2945 .build();
2946 let cache = Arc::new(cache);
2947
2948 let region_id = RegionId::new(1, 1);
2949 let index_id = RegionIndexId::new(RegionFileId::new(region_id, FileId::random()), 0);
2950 let column_id: ColumnId = 1;
2951
2952 let bloom_cache = cache.bloom_filter_index_cache().unwrap().clone();
2953 let inverted_cache = cache.inverted_index_cache().unwrap().clone();
2954 let result_cache = cache.index_result_cache().unwrap();
2955 let puffin_metadata_cache = cache.puffin_metadata_cache().unwrap().clone();
2956
2957 let bloom_key = (
2958 index_id.file_id(),
2959 index_id.version,
2960 column_id,
2961 Tag::Skipping,
2962 );
2963 bloom_cache.put_metadata(bloom_key, Arc::new(BloomFilterMeta::default()));
2964 inverted_cache.put_metadata(
2965 (index_id.file_id(), index_id.version),
2966 Arc::new(InvertedIndexMetas::default()),
2967 );
2968 let predicate = PredicateKey::new_bloom(Arc::new(BTreeMap::new()));
2969 let selection = Arc::new(RowGroupSelection::default());
2970 result_cache.put(predicate.clone(), index_id.file_id(), selection);
2971 let file_id_str = index_id.to_string();
2972 let metadata = Arc::new(FileMetadata {
2973 blobs: Vec::new(),
2974 properties: HashMap::new(),
2975 });
2976 puffin_metadata_cache.put_metadata(file_id_str.clone(), metadata);
2977
2978 assert!(bloom_cache.get_metadata(bloom_key).is_some());
2979 assert!(
2980 inverted_cache
2981 .get_metadata((index_id.file_id(), index_id.version))
2982 .is_some()
2983 );
2984 assert!(result_cache.get(&predicate, index_id.file_id()).is_some());
2985 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_some());
2986
2987 cache.evict_puffin_cache(index_id).await;
2988
2989 assert!(bloom_cache.get_metadata(bloom_key).is_none());
2990 assert!(
2991 inverted_cache
2992 .get_metadata((index_id.file_id(), index_id.version))
2993 .is_none()
2994 );
2995 assert!(result_cache.get(&predicate, index_id.file_id()).is_none());
2996 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_none());
2997
2998 bloom_cache.put_metadata(bloom_key, Arc::new(BloomFilterMeta::default()));
3000 inverted_cache.put_metadata(
3001 (index_id.file_id(), index_id.version),
3002 Arc::new(InvertedIndexMetas::default()),
3003 );
3004 result_cache.put(
3005 predicate.clone(),
3006 index_id.file_id(),
3007 Arc::new(RowGroupSelection::default()),
3008 );
3009 puffin_metadata_cache.put_metadata(
3010 file_id_str.clone(),
3011 Arc::new(FileMetadata {
3012 blobs: Vec::new(),
3013 properties: HashMap::new(),
3014 }),
3015 );
3016
3017 let strategy = CacheStrategy::EnableAll(cache.clone());
3018 strategy.evict_puffin_cache(index_id).await;
3019
3020 assert!(bloom_cache.get_metadata(bloom_key).is_none());
3021 assert!(
3022 inverted_cache
3023 .get_metadata((index_id.file_id(), index_id.version))
3024 .is_none()
3025 );
3026 assert!(result_cache.get(&predicate, index_id.file_id()).is_none());
3027 assert!(puffin_metadata_cache.get_metadata(&file_id_str).is_none());
3028 }
3029
3030 fn wide_region_metadata(column_count: u32) -> RegionMetadata {
3031 let region_id = RegionId::new(1024, 7);
3032 let mut builder = RegionMetadataBuilder::new(region_id);
3033 let mut primary_key = Vec::new();
3034
3035 for column_id in 0..column_count {
3036 let semantic_type = if column_id < 32 {
3037 primary_key.push(column_id);
3038 SemanticType::Tag
3039 } else {
3040 SemanticType::Field
3041 };
3042 let mut column_schema = ColumnSchema::new(
3043 format!("wide_column_{column_id}"),
3044 ConcreteDataType::string_datatype(),
3045 true,
3046 );
3047 column_schema
3048 .mut_metadata()
3049 .insert(format!("cache_key_{column_id}"), "cache_value".repeat(4));
3050 builder.push_column_metadata(ColumnMetadata {
3051 column_schema,
3052 semantic_type,
3053 column_id,
3054 });
3055 }
3056
3057 builder.push_column_metadata(ColumnMetadata {
3058 column_schema: ColumnSchema::new(
3059 "ts",
3060 ConcreteDataType::timestamp_millisecond_datatype(),
3061 false,
3062 ),
3063 semantic_type: SemanticType::Timestamp,
3064 column_id: column_count,
3065 });
3066 builder.primary_key(primary_key);
3067
3068 builder.build().unwrap()
3069 }
3070}