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