Skip to main content

mito2/
cache.rs

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