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