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