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