Skip to main content

mito2/cache/
file_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//! A cache for files.
16
17use std::fmt;
18use std::ops::Range;
19use std::sync::Arc;
20use std::time::{Duration, Instant};
21
22use bytes::Bytes;
23use common_base::readable_size::ReadableSize;
24use common_runtime::Runtime;
25use common_telemetry::{debug, error, info, warn};
26use futures::{AsyncWriteExt, FutureExt, TryStreamExt};
27use moka::future::Cache;
28use moka::notification::RemovalCause;
29use moka::policy::EvictionPolicy;
30use object_store::util::join_path;
31use object_store::{ErrorKind, ObjectStore, Reader};
32use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData};
33use snafu::ResultExt;
34use store_api::storage::{FileId, RegionId};
35use tokio::sync::mpsc::{Sender, UnboundedReceiver};
36
37use crate::access_layer::TempFileCleaner;
38use crate::cache::{
39    CachedSstMeta, FILE_TYPE, INDEX_TYPE, SstMetaPreparation, decode_sst_meta, prepare_sst_meta,
40};
41use crate::error::{self, OpenDalSnafu, Result};
42use crate::metrics::{
43    CACHE_BYTES, CACHE_HIT, CACHE_MISS, WRITE_CACHE_DOWNLOAD_BYTES_TOTAL,
44    WRITE_CACHE_DOWNLOAD_ELAPSED,
45};
46use crate::region::opener::RegionLoadCacheTask;
47use crate::sst::parquet::helper::fetch_byte_ranges;
48use crate::sst::parquet::metadata::MetadataLoader;
49use crate::sst::parquet::reader::MetadataCacheMetrics;
50
51/// Subdirectory of cached files for write.
52///
53/// This must contain three layers, corresponding to [`build_prometheus_metrics_layer`](object_store::layers::build_prometheus_metrics_layer).
54const FILE_DIR: &str = "cache/object/write/";
55
56/// Default percentage for index (puffin) cache (20% of total capacity).
57pub(crate) const DEFAULT_INDEX_CACHE_PERCENT: u8 = 20;
58
59/// Minimum capacity for each cache (512MB).
60const MIN_CACHE_CAPACITY: u64 = 512 * 1024 * 1024;
61
62/// Channel capacity for background download tasks.
63const DOWNLOAD_TASK_CHANNEL_SIZE: usize = 64;
64
65/// A task to download a file in the background.
66struct DownloadTask {
67    index_key: IndexKey,
68    remote_path: String,
69    remote_store: ObjectStore,
70    file_size: u64,
71}
72
73/// Inner struct for FileCache that can be used in spawned tasks.
74#[derive(Debug)]
75struct FileCacheInner {
76    /// Local store to cache files.
77    local_store: ObjectStore,
78    /// Index to track cached Parquet files.
79    parquet_index: Cache<IndexKey, IndexValue>,
80    /// Index to track cached Puffin files.
81    puffin_index: Cache<IndexKey, IndexValue>,
82}
83
84impl FileCacheInner {
85    /// Returns the appropriate memory index for the given file type.
86    fn memory_index(&self, file_type: FileType) -> &Cache<IndexKey, IndexValue> {
87        match file_type {
88            FileType::Parquet => &self.parquet_index,
89            FileType::Puffin { .. } => &self.puffin_index,
90        }
91    }
92
93    /// Returns the cache file path for the key.
94    fn cache_file_path(&self, key: IndexKey) -> String {
95        cache_file_path(FILE_DIR, key)
96    }
97
98    /// Puts a file into the cache index.
99    ///
100    /// The `WriteCache` should ensure the file is in the correct path.
101    async fn put(&self, key: IndexKey, value: IndexValue) {
102        CACHE_BYTES
103            .with_label_values(&[key.file_type.metric_label()])
104            .add(value.file_size.into());
105        let index = self.memory_index(key.file_type);
106        index.insert(key, value).await;
107
108        // Since files are large items, we run the pending tasks immediately.
109        index.run_pending_tasks().await;
110    }
111
112    /// Recovers the index from local store.
113    async fn recover(&self) -> Result<()> {
114        let now = Instant::now();
115        let mut lister = self
116            .local_store
117            .lister_with(FILE_DIR)
118            .await
119            .context(OpenDalSnafu)?;
120        // Use i64 for total_size to reduce the risk of overflow.
121        // It is possible that the total size of the cache is larger than i32::MAX.
122        let (mut total_size, mut total_keys) = (0i64, 0);
123        let (mut parquet_size, mut puffin_size) = (0i64, 0i64);
124        while let Some(entry) = lister.try_next().await.context(OpenDalSnafu)? {
125            let meta = entry.metadata();
126            if !meta.is_file() {
127                continue;
128            }
129            let Some(key) = parse_index_key(entry.name()) else {
130                continue;
131            };
132
133            let meta = self
134                .local_store
135                .stat(entry.path())
136                .await
137                .context(OpenDalSnafu)?;
138            let file_size = meta.content_length() as u32;
139            let index = self.memory_index(key.file_type);
140            index.insert(key, IndexValue { file_size }).await;
141            let size = i64::from(file_size);
142            total_size += size;
143            total_keys += 1;
144
145            // Track sizes separately for each file type
146            match key.file_type {
147                FileType::Parquet => parquet_size += size,
148                FileType::Puffin { .. } => puffin_size += size,
149            }
150        }
151        // The metrics is a signed int gauge so we can updates it finally.
152        CACHE_BYTES
153            .with_label_values(&[FILE_TYPE])
154            .add(parquet_size);
155        CACHE_BYTES
156            .with_label_values(&[INDEX_TYPE])
157            .add(puffin_size);
158
159        // Run all pending tasks of the moka cache so that the cache size is updated
160        // and the eviction policy is applied.
161        self.parquet_index.run_pending_tasks().await;
162        self.puffin_index.run_pending_tasks().await;
163
164        let parquet_weight = self.parquet_index.weighted_size();
165        let parquet_count = self.parquet_index.entry_count();
166        let puffin_weight = self.puffin_index.weighted_size();
167        let puffin_count = self.puffin_index.entry_count();
168        info!(
169            "Recovered file cache, num_keys: {}, num_bytes: {}, parquet(count: {}, weight: {}), puffin(count: {}, weight: {}), cost: {:?}",
170            total_keys,
171            total_size,
172            parquet_count,
173            parquet_weight,
174            puffin_count,
175            puffin_weight,
176            now.elapsed()
177        );
178        Ok(())
179    }
180
181    /// Downloads a file without cleaning up on error.
182    async fn download_without_cleaning(
183        &self,
184        index_key: IndexKey,
185        remote_path: &str,
186        remote_store: &ObjectStore,
187        file_size: u64,
188        concurrency: usize,
189    ) -> Result<()> {
190        const DOWNLOAD_READER_CHUNK_SIZE: ReadableSize = ReadableSize::mb(8);
191
192        let file_type = index_key.file_type;
193        let timer = WRITE_CACHE_DOWNLOAD_ELAPSED
194            .with_label_values(&[match file_type {
195                FileType::Parquet => "download_parquet",
196                FileType::Puffin { .. } => "download_puffin",
197            }])
198            .start_timer();
199
200        let reader = remote_store
201            .reader_with(remote_path)
202            .concurrent(concurrency)
203            .chunk(DOWNLOAD_READER_CHUNK_SIZE.as_bytes() as usize)
204            .await
205            .context(error::OpenDalSnafu)?
206            .into_futures_async_read(0..file_size)
207            .await
208            .context(error::OpenDalSnafu)?;
209
210        let cache_path = self.cache_file_path(index_key);
211        let mut writer = self
212            .local_store
213            .writer(&cache_path)
214            .await
215            .context(error::OpenDalSnafu)?
216            .into_futures_async_write();
217
218        let region_id = index_key.region_id;
219        let file_id = index_key.file_id;
220        let bytes_written =
221            futures::io::copy(reader, &mut writer)
222                .await
223                .context(error::DownloadSnafu {
224                    region_id,
225                    file_id,
226                    file_type,
227                })?;
228        writer.close().await.context(error::DownloadSnafu {
229            region_id,
230            file_id,
231            file_type,
232        })?;
233
234        WRITE_CACHE_DOWNLOAD_BYTES_TOTAL.inc_by(bytes_written);
235
236        let elapsed = timer.stop_and_record();
237        debug!(
238            "Successfully download file '{}' to local '{}', file size: {}, region: {}, cost: {:?}s",
239            remote_path, cache_path, bytes_written, region_id, elapsed,
240        );
241
242        let index_value = IndexValue {
243            file_size: bytes_written as _,
244        };
245        self.put(index_key, index_value).await;
246        Ok(())
247    }
248
249    /// Downloads a file from remote store to local cache.
250    async fn download(
251        &self,
252        index_key: IndexKey,
253        remote_path: &str,
254        remote_store: &ObjectStore,
255        file_size: u64,
256        concurrency: usize,
257    ) -> Result<()> {
258        if let Err(e) = self
259            .download_without_cleaning(index_key, remote_path, remote_store, file_size, concurrency)
260            .await
261        {
262            error!(e; "Failed to download file '{}' for region {}", remote_path, index_key.region_id);
263
264            let filename = index_key.to_string();
265            TempFileCleaner::clean_atomic_dir_files(&self.local_store, &[&filename]).await;
266
267            return Err(e);
268        }
269
270        Ok(())
271    }
272
273    /// Checks if the key is in the file cache.
274    fn contains_key(&self, key: &IndexKey) -> bool {
275        self.memory_index(key.file_type).contains_key(key)
276    }
277}
278
279/// A file cache manages files on local store and evict files based
280/// on size.
281#[derive(Debug, Clone)]
282pub(crate) struct FileCache {
283    /// Inner cache state shared with background worker.
284    inner: Arc<FileCacheInner>,
285    /// Capacity of the puffin (index) cache in bytes.
286    puffin_capacity: u64,
287    /// Channel for background download tasks. None if background worker is disabled.
288    download_task_tx: Option<Sender<DownloadTask>>,
289}
290
291pub(crate) type FileCacheRef = Arc<FileCache>;
292
293impl FileCache {
294    /// Splits the configured total capacity between parquet and puffin caches
295    /// without exceeding the requested overall budget.
296    fn split_cache_capacities(total_capacity: u64, index_percent: u8) -> (u64, u64) {
297        let desired_puffin_capacity = total_capacity * u64::from(index_percent) / 100;
298        let min_cache_capacity = MIN_CACHE_CAPACITY.min(total_capacity / 2);
299        let puffin_capacity =
300            desired_puffin_capacity.clamp(min_cache_capacity, total_capacity - min_cache_capacity);
301        let parquet_capacity = total_capacity - puffin_capacity;
302        (parquet_capacity, puffin_capacity)
303    }
304
305    /// Creates a new file cache.
306    pub(crate) fn new(
307        local_store: ObjectStore,
308        capacity: ReadableSize,
309        ttl: Option<Duration>,
310        index_cache_percent: Option<u8>,
311        enable_background_worker: bool,
312    ) -> FileCache {
313        // Validate and use the provided percent or default
314        let index_percent = index_cache_percent
315            .filter(|&percent| percent > 0 && percent < 100)
316            .unwrap_or(DEFAULT_INDEX_CACHE_PERCENT);
317        let total_capacity = capacity.as_bytes();
318
319        let (parquet_capacity, puffin_capacity) =
320            Self::split_cache_capacities(total_capacity, index_percent);
321
322        info!(
323            "Initializing file cache with index_percent: {}%, total_capacity: {}, parquet_capacity: {}, puffin_capacity: {}",
324            index_percent,
325            ReadableSize(total_capacity),
326            ReadableSize(parquet_capacity),
327            ReadableSize(puffin_capacity)
328        );
329
330        let parquet_index = Self::build_cache(local_store.clone(), parquet_capacity, ttl, "file");
331        let puffin_index = Self::build_cache(local_store.clone(), puffin_capacity, ttl, "index");
332
333        // Create inner cache shared with background worker
334        let inner = Arc::new(FileCacheInner {
335            local_store,
336            parquet_index,
337            puffin_index,
338        });
339
340        // Only create channel and spawn worker if background download is enabled
341        let download_task_tx = if enable_background_worker {
342            let (tx, rx) = tokio::sync::mpsc::channel(DOWNLOAD_TASK_CHANNEL_SIZE);
343            Self::spawn_download_worker(inner.clone(), rx);
344            Some(tx)
345        } else {
346            None
347        };
348
349        FileCache {
350            inner,
351            puffin_capacity,
352            download_task_tx,
353        }
354    }
355
356    /// Spawns a background worker to process download tasks.
357    fn spawn_download_worker(
358        inner: Arc<FileCacheInner>,
359        mut download_task_rx: tokio::sync::mpsc::Receiver<DownloadTask>,
360    ) {
361        tokio::spawn(async move {
362            info!("Background download worker started");
363            while let Some(task) = download_task_rx.recv().await {
364                // Check if the file is already in the cache
365                if inner.contains_key(&task.index_key) {
366                    debug!(
367                        "Skipping background download for region {}, file {} - already in cache",
368                        task.index_key.region_id, task.index_key.file_id
369                    );
370                    continue;
371                }
372
373                // Ignores background download errors.
374                let _ = inner
375                    .download(
376                        task.index_key,
377                        &task.remote_path,
378                        &task.remote_store,
379                        task.file_size,
380                        1, // Background downloads use concurrency=1
381                    )
382                    .await;
383            }
384            info!("Background download worker stopped");
385        });
386    }
387
388    /// Builds a cache for a specific file type.
389    fn build_cache(
390        local_store: ObjectStore,
391        capacity: u64,
392        ttl: Option<Duration>,
393        label: &'static str,
394    ) -> Cache<IndexKey, IndexValue> {
395        let cache_store = local_store;
396        let mut builder = Cache::builder()
397            .eviction_policy(EvictionPolicy::lru())
398            .weigher(|_key, value: &IndexValue| -> u32 {
399                // We only measure space on local store.
400                value.file_size
401            })
402            .max_capacity(capacity)
403            .async_eviction_listener(move |key, value, cause| {
404                let store = cache_store.clone();
405                // Stores files under FILE_DIR.
406                let file_path = cache_file_path(FILE_DIR, *key);
407                async move {
408                    if let RemovalCause::Replaced = cause {
409                        // The cache is replaced by another file (maybe download again). We don't remove the same
410                        // file but updates the metrics as the file is already replaced by users.
411                        CACHE_BYTES.with_label_values(&[label]).sub(value.file_size.into());
412                        return;
413                    }
414
415                    match store.delete(&file_path).await {
416                        Ok(()) => {
417                            CACHE_BYTES.with_label_values(&[label]).sub(value.file_size.into());
418                        }
419                        Err(e) => {
420                            warn!(e; "Failed to delete cached file {} for region {}", file_path, key.region_id);
421                        }
422                    }
423                }
424                .boxed()
425            });
426        if let Some(ttl) = ttl {
427            builder = builder.time_to_idle(ttl);
428        }
429        builder.build()
430    }
431
432    /// Puts a file into the cache index.
433    ///
434    /// The `WriteCache` should ensure the file is in the correct path.
435    pub(crate) async fn put(&self, key: IndexKey, value: IndexValue) {
436        self.inner.put(key, value).await
437    }
438
439    pub(crate) async fn get(&self, key: IndexKey) -> Option<IndexValue> {
440        self.inner.memory_index(key.file_type).get(&key).await
441    }
442
443    /// Reads a file from the cache.
444    #[allow(unused)]
445    pub(crate) async fn reader(&self, key: IndexKey) -> Option<Reader> {
446        // We must use `get()` to update the estimator of the cache.
447        // See https://docs.rs/moka/latest/moka/future/struct.Cache.html#method.contains_key
448        let index = self.inner.memory_index(key.file_type);
449        if index.get(&key).await.is_none() {
450            CACHE_MISS
451                .with_label_values(&[key.file_type.metric_label()])
452                .inc();
453            return None;
454        }
455
456        let file_path = self.inner.cache_file_path(key);
457        match self.get_reader(&file_path).await {
458            Ok(Some(reader)) => {
459                CACHE_HIT
460                    .with_label_values(&[key.file_type.metric_label()])
461                    .inc();
462                return Some(reader);
463            }
464            Err(e) => {
465                if e.kind() != ErrorKind::NotFound {
466                    warn!(e; "Failed to get file for key {:?}", key);
467                }
468            }
469            Ok(None) => {}
470        }
471
472        // We removes the file from the index.
473        index.remove(&key).await;
474        CACHE_MISS
475            .with_label_values(&[key.file_type.metric_label()])
476            .inc();
477        None
478    }
479
480    /// Reads ranges from the cache.
481    pub(crate) async fn read_ranges(
482        &self,
483        key: IndexKey,
484        ranges: &[Range<u64>],
485    ) -> Option<Vec<Bytes>> {
486        let index = self.inner.memory_index(key.file_type);
487        if index.get(&key).await.is_none() {
488            CACHE_MISS
489                .with_label_values(&[key.file_type.metric_label()])
490                .inc();
491            return None;
492        }
493
494        let file_path = self.inner.cache_file_path(key);
495        // In most cases, it will use blocking read,
496        // because FileCache is normally based on local file system, which supports blocking read.
497        let bytes_result =
498            fetch_byte_ranges(&file_path, self.inner.local_store.clone(), ranges).await;
499        match bytes_result {
500            Ok(bytes) => {
501                CACHE_HIT
502                    .with_label_values(&[key.file_type.metric_label()])
503                    .inc();
504                Some(bytes)
505            }
506            Err(e) => {
507                if e.kind() != ErrorKind::NotFound {
508                    warn!(e; "Failed to get file for key {:?}", key);
509                }
510
511                // We removes the file from the index.
512                index.remove(&key).await;
513                CACHE_MISS
514                    .with_label_values(&[key.file_type.metric_label()])
515                    .inc();
516                None
517            }
518        }
519    }
520
521    /// Removes a file from the cache explicitly.
522    /// It always tries to remove the file from the local store because we may not have the file
523    /// in the memory index if upload is failed.
524    pub(crate) async fn remove(&self, key: IndexKey) {
525        let file_path = self.inner.cache_file_path(key);
526        self.inner.memory_index(key.file_type).remove(&key).await;
527        // Always delete the file from the local store.
528        if let Err(e) = self.inner.local_store.delete(&file_path).await {
529            warn!(e; "Failed to delete a cached file {}", file_path);
530        }
531    }
532
533    /// Recovers the index from local store.
534    ///
535    /// If `task_receiver` is provided, spawns a background task after recovery
536    /// to process `RegionLoadCacheTask` messages for loading files into the cache.
537    pub(crate) async fn recover(
538        &self,
539        sync: bool,
540        task_receiver: Option<UnboundedReceiver<RegionLoadCacheTask>>,
541    ) {
542        let moved_self = self.clone();
543        let handle = tokio::spawn(async move {
544            if let Err(err) = moved_self.inner.recover().await {
545                error!(err; "Failed to recover file cache.")
546            }
547
548            // Spawns background task to process region load cache tasks after recovery.
549            // So it won't block the recovery when `sync` is true.
550            if let Some(mut receiver) = task_receiver {
551                info!("Spawning background task for processing region load cache tasks");
552                tokio::spawn(async move {
553                    while let Some(task) = receiver.recv().await {
554                        task.fill_cache(&moved_self).await;
555                    }
556                    info!("Background task for processing region load cache tasks stopped");
557                });
558            }
559        });
560
561        if sync {
562            let _ = handle.await;
563        }
564    }
565
566    /// Returns the cache file path for the key.
567    pub(crate) fn cache_file_path(&self, key: IndexKey) -> String {
568        self.inner.cache_file_path(key)
569    }
570
571    /// Returns the local store of the file cache.
572    pub(crate) fn local_store(&self) -> ObjectStore {
573        self.inner.local_store.clone()
574    }
575
576    /// Get the parquet metadata in file cache.
577    /// If the file is not in the cache or fail to load metadata, return None.
578    pub(crate) async fn get_parquet_meta_data(
579        &self,
580        key: IndexKey,
581        cache_metrics: &mut MetadataCacheMetrics,
582        page_index_policy: PageIndexPolicy,
583    ) -> Option<ParquetMetaData> {
584        // Check if file cache contains the key
585        if let Some(index_value) = self.inner.parquet_index.get(&key).await {
586            // Load metadata from file cache
587            let local_store = self.local_store();
588            let file_path = self.inner.cache_file_path(key);
589            let file_size = index_value.file_size as u64;
590            let mut metadata_loader = MetadataLoader::new(local_store, &file_path, file_size);
591            metadata_loader.with_page_index_policy(page_index_policy);
592
593            match metadata_loader.load(cache_metrics).await {
594                Ok(metadata) => {
595                    CACHE_HIT
596                        .with_label_values(&[key.file_type.metric_label()])
597                        .inc();
598                    Some(metadata)
599                }
600                Err(e) => {
601                    if !e.is_object_not_found() {
602                        warn!(
603                            e; "Failed to get parquet metadata for key {:?}",
604                            key
605                        );
606                    }
607                    // We removes the file from the index.
608                    self.inner.parquet_index.remove(&key).await;
609                    CACHE_MISS
610                        .with_label_values(&[key.file_type.metric_label()])
611                        .inc();
612                    None
613                }
614            }
615        } else {
616            CACHE_MISS
617                .with_label_values(&[key.file_type.metric_label()])
618                .inc();
619            None
620        }
621    }
622
623    /// Get fused SST metadata from the file cache.
624    /// If the file is not in the cache, or metadata loading/decoding fails, return None.
625    /// Compact cache encoding failures return decoded-only metadata to the caller.
626    pub(crate) async fn get_sst_meta_data(
627        &self,
628        key: IndexKey,
629        cache_metrics: &mut MetadataCacheMetrics,
630        page_index_policy: PageIndexPolicy,
631        runtime: &Runtime,
632    ) -> Option<SstMetaPreparation> {
633        let file_path = self.inner.cache_file_path(key);
634        let metadata = self
635            .get_parquet_meta_data(key, cache_metrics, page_index_policy)
636            .await?;
637        match prepare_sst_meta(&file_path, metadata, None, page_index_policy, runtime).await {
638            Ok(metadata) => Some(metadata),
639            Err(err) => {
640                CACHE_MISS
641                    .with_label_values(&[key.file_type.metric_label()])
642                    .inc();
643                warn!(
644                    err; "Failed to prepare cached parquet metadata for key {:?}",
645                    key
646                );
647                None
648            }
649        }
650    }
651
652    /// Gets decoded SST metadata without preparing an in-memory cache entry.
653    pub(crate) async fn get_decoded_sst_meta_data(
654        &self,
655        key: IndexKey,
656        cache_metrics: &mut MetadataCacheMetrics,
657        page_index_policy: PageIndexPolicy,
658        runtime: &Runtime,
659    ) -> Option<Arc<CachedSstMeta>> {
660        let file_path = self.inner.cache_file_path(key);
661        let metadata = self
662            .get_parquet_meta_data(key, cache_metrics, page_index_policy)
663            .await?;
664        match decode_sst_meta(&file_path, metadata, None, page_index_policy, runtime).await {
665            Ok(metadata) => Some(metadata),
666            Err(err) => {
667                CACHE_MISS
668                    .with_label_values(&[key.file_type.metric_label()])
669                    .inc();
670                warn!(
671                    err; "Failed to decode cached parquet metadata for key {:?}",
672                    key
673                );
674                None
675            }
676        }
677    }
678
679    async fn get_reader(&self, file_path: &str) -> object_store::Result<Option<Reader>> {
680        if self.inner.local_store.exists(file_path).await? {
681            Ok(Some(self.inner.local_store.reader(file_path).await?))
682        } else {
683            Ok(None)
684        }
685    }
686
687    /// Checks if the key is in the file cache.
688    pub(crate) fn contains_key(&self, key: &IndexKey) -> bool {
689        self.inner.contains_key(key)
690    }
691
692    /// Returns the capacity of the puffin (index) cache in bytes.
693    pub(crate) fn puffin_cache_capacity(&self) -> u64 {
694        self.puffin_capacity
695    }
696
697    /// Returns the current weighted size (used bytes) of the puffin (index) cache.
698    pub(crate) fn puffin_cache_size(&self) -> u64 {
699        self.inner.puffin_index.weighted_size()
700    }
701
702    /// Downloads a file in `remote_path` from the remote object store to the local cache
703    /// (specified by `index_key`).
704    pub(crate) async fn download(
705        &self,
706        index_key: IndexKey,
707        remote_path: &str,
708        remote_store: &ObjectStore,
709        file_size: u64,
710    ) -> Result<()> {
711        self.inner
712            .download(index_key, remote_path, remote_store, file_size, 8) // Foreground uses concurrency=8
713            .await
714    }
715
716    /// Downloads a file in `remote_path` from the remote object store to the local cache
717    /// (specified by `index_key`) in the background. Errors are logged but not returned.
718    ///
719    /// This method attempts to send a download task to the background worker.
720    /// If the channel is full, the task is silently dropped.
721    pub(crate) fn maybe_download_background(
722        &self,
723        index_key: IndexKey,
724        remote_path: String,
725        remote_store: ObjectStore,
726        file_size: u64,
727    ) {
728        // Do nothing if background worker is disabled (channel is None)
729        let Some(tx) = &self.download_task_tx else {
730            return;
731        };
732
733        let task = DownloadTask {
734            index_key,
735            remote_path,
736            remote_store,
737            file_size,
738        };
739
740        // Try to send the task; if the channel is full, just drop it
741        if let Err(e) = tx.try_send(task) {
742            debug!(
743                "Failed to queue background download task for region {}, file {}: {:?}",
744                index_key.region_id, index_key.file_id, e
745            );
746        }
747    }
748}
749
750/// Key of file cache index.
751#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
752pub struct IndexKey {
753    pub region_id: RegionId,
754    pub file_id: FileId,
755    pub file_type: FileType,
756}
757
758impl IndexKey {
759    /// Creates a new index key.
760    pub fn new(region_id: RegionId, file_id: FileId, file_type: FileType) -> IndexKey {
761        IndexKey {
762            region_id,
763            file_id,
764            file_type,
765        }
766    }
767}
768
769impl fmt::Display for IndexKey {
770    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
771        write!(
772            f,
773            "{}.{}.{}",
774            self.region_id.as_u64(),
775            self.file_id,
776            self.file_type
777        )
778    }
779}
780
781/// Type of the file.
782#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
783pub enum FileType {
784    /// Parquet file.
785    Parquet,
786    /// Puffin file.
787    Puffin(u64),
788}
789
790impl fmt::Display for FileType {
791    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
792        match self {
793            FileType::Parquet => write!(f, "parquet"),
794            FileType::Puffin(version) => write!(f, "{}.puffin", version),
795        }
796    }
797}
798
799impl FileType {
800    /// Parses the file type from string.
801    pub(crate) fn parse(s: &str) -> Option<FileType> {
802        match s {
803            "parquet" => Some(FileType::Parquet),
804            "puffin" => Some(FileType::Puffin(0)),
805            _ => {
806                // if post-fix with .puffin, try to parse the version
807                if let Some(version_str) = s.strip_suffix(".puffin") {
808                    let version = version_str.parse::<u64>().ok()?;
809                    Some(FileType::Puffin(version))
810                } else {
811                    None
812                }
813            }
814        }
815    }
816
817    /// Returns the metric label for this file type.
818    fn metric_label(&self) -> &'static str {
819        match self {
820            FileType::Parquet => FILE_TYPE,
821            FileType::Puffin(_) => INDEX_TYPE,
822        }
823    }
824}
825
826/// An entity that describes the file in the file cache.
827///
828/// It should only keep minimal information needed by the cache.
829#[derive(Debug, Clone)]
830pub(crate) struct IndexValue {
831    /// Size of the file in bytes.
832    pub(crate) file_size: u32,
833}
834
835/// Generates the path to the cached file.
836///
837/// The file name format is `{region_id}.{file_id}.{file_type}`
838fn cache_file_path(cache_file_dir: &str, key: IndexKey) -> String {
839    join_path(cache_file_dir, &key.to_string())
840}
841
842/// Parse index key from the file name.
843fn parse_index_key(name: &str) -> Option<IndexKey> {
844    let mut split = name.splitn(3, '.');
845    let region_id = split.next().and_then(|s| {
846        let id = s.parse::<u64>().ok()?;
847        Some(RegionId::from_u64(id))
848    })?;
849    let file_id = split.next().and_then(|s| FileId::parse_str(s).ok())?;
850    let file_type = split.next().and_then(FileType::parse)?;
851
852    Some(IndexKey::new(region_id, file_id, file_type))
853}
854
855#[cfg(test)]
856mod tests {
857    use common_test_util::temp_dir::create_temp_dir;
858    use object_store::services::Fs;
859
860    use super::*;
861
862    fn new_fs_store(path: &str) -> ObjectStore {
863        let builder = Fs::default().root(path);
864        ObjectStore::new(builder).unwrap()
865    }
866
867    #[tokio::test]
868    async fn test_file_cache_ttl() {
869        let dir = create_temp_dir("");
870        let local_store = new_fs_store(dir.path().to_str().unwrap());
871
872        let cache = FileCache::new(
873            local_store.clone(),
874            ReadableSize::mb(10),
875            Some(Duration::from_millis(10)),
876            None,
877            true, // enable_background_worker
878        );
879        let region_id = RegionId::new(2000, 0);
880        let file_id = FileId::random();
881        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
882        let file_path = cache.cache_file_path(key);
883
884        // Get an empty file.
885        assert!(cache.reader(key).await.is_none());
886
887        // Write a file.
888        local_store
889            .write(&file_path, b"hello".as_slice())
890            .await
891            .unwrap();
892
893        // Add to the cache.
894        cache
895            .put(
896                IndexKey::new(region_id, file_id, FileType::Parquet),
897                IndexValue { file_size: 5 },
898            )
899            .await;
900
901        let exist = cache.reader(key).await;
902        assert!(exist.is_some());
903        tokio::time::sleep(Duration::from_millis(15)).await;
904        cache.inner.parquet_index.run_pending_tasks().await;
905        let non = cache.reader(key).await;
906        assert!(non.is_none());
907    }
908
909    #[tokio::test]
910    async fn test_file_cache_basic() {
911        let dir = create_temp_dir("");
912        let local_store = new_fs_store(dir.path().to_str().unwrap());
913
914        let cache = FileCache::new(
915            local_store.clone(),
916            ReadableSize::mb(10),
917            None,
918            None,
919            true, // enable_background_worker
920        );
921        let region_id = RegionId::new(2000, 0);
922        let file_id = FileId::random();
923        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
924        let file_path = cache.cache_file_path(key);
925
926        // Get an empty file.
927        assert!(cache.reader(key).await.is_none());
928
929        // Write a file.
930        local_store
931            .write(&file_path, b"hello".as_slice())
932            .await
933            .unwrap();
934        // Add to the cache.
935        cache
936            .put(
937                IndexKey::new(region_id, file_id, FileType::Parquet),
938                IndexValue { file_size: 5 },
939            )
940            .await;
941
942        // Read file content.
943        let reader = cache.reader(key).await.unwrap();
944        let buf = reader.read(..).await.unwrap().to_vec();
945        assert_eq!("hello", String::from_utf8(buf).unwrap());
946
947        // Get weighted size.
948        cache.inner.parquet_index.run_pending_tasks().await;
949        assert_eq!(5, cache.inner.parquet_index.weighted_size());
950
951        // Remove the file.
952        cache.remove(key).await;
953        assert!(cache.reader(key).await.is_none());
954
955        // Ensure all pending tasks of the moka cache is done before assertion.
956        cache.inner.parquet_index.run_pending_tasks().await;
957
958        // The file also not exists.
959        assert!(!local_store.exists(&file_path).await.unwrap());
960        assert_eq!(0, cache.inner.parquet_index.weighted_size());
961    }
962
963    #[tokio::test]
964    async fn test_file_cache_file_removed() {
965        let dir = create_temp_dir("");
966        let local_store = new_fs_store(dir.path().to_str().unwrap());
967
968        let cache = FileCache::new(
969            local_store.clone(),
970            ReadableSize::mb(10),
971            None,
972            None,
973            true, // enable_background_worker
974        );
975        let region_id = RegionId::new(2000, 0);
976        let file_id = FileId::random();
977        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
978        let file_path = cache.cache_file_path(key);
979
980        // Write a file.
981        local_store
982            .write(&file_path, b"hello".as_slice())
983            .await
984            .unwrap();
985        // Add to the cache.
986        cache
987            .put(
988                IndexKey::new(region_id, file_id, FileType::Parquet),
989                IndexValue { file_size: 5 },
990            )
991            .await;
992
993        // Remove the file but keep the index.
994        local_store.delete(&file_path).await.unwrap();
995
996        // Reader is none.
997        assert!(cache.reader(key).await.is_none());
998        // Key is removed.
999        assert!(!cache.inner.parquet_index.contains_key(&key));
1000    }
1001
1002    #[tokio::test]
1003    async fn test_file_cache_recover() {
1004        let dir = create_temp_dir("");
1005        let local_store = new_fs_store(dir.path().to_str().unwrap());
1006        let cache = FileCache::new(
1007            local_store.clone(),
1008            ReadableSize::mb(10),
1009            None,
1010            None,
1011            true, // enable_background_worker
1012        );
1013
1014        let region_id = RegionId::new(2000, 0);
1015        let file_type = FileType::Parquet;
1016        // Write N files.
1017        let file_ids: Vec<_> = (0..10).map(|_| FileId::random()).collect();
1018        let mut total_size = 0;
1019        for (i, file_id) in file_ids.iter().enumerate() {
1020            let key = IndexKey::new(region_id, *file_id, file_type);
1021            let file_path = cache.cache_file_path(key);
1022            let bytes = i.to_string().into_bytes();
1023            local_store.write(&file_path, bytes.clone()).await.unwrap();
1024
1025            // Add to the cache.
1026            cache
1027                .put(
1028                    IndexKey::new(region_id, *file_id, file_type),
1029                    IndexValue {
1030                        file_size: bytes.len() as u32,
1031                    },
1032                )
1033                .await;
1034            total_size += bytes.len();
1035        }
1036
1037        // Recover the cache.
1038        let cache = FileCache::new(
1039            local_store.clone(),
1040            ReadableSize::mb(10),
1041            None,
1042            None,
1043            true, // enable_background_worker
1044        );
1045        // No entry before recovery.
1046        assert!(
1047            cache
1048                .reader(IndexKey::new(region_id, file_ids[0], file_type))
1049                .await
1050                .is_none()
1051        );
1052        cache.recover(true, None).await;
1053
1054        // Check size.
1055        cache.inner.parquet_index.run_pending_tasks().await;
1056        assert_eq!(
1057            total_size,
1058            cache.inner.parquet_index.weighted_size() as usize
1059        );
1060
1061        for (i, file_id) in file_ids.iter().enumerate() {
1062            let key = IndexKey::new(region_id, *file_id, file_type);
1063            let reader = cache.reader(key).await.unwrap();
1064            let buf = reader.read(..).await.unwrap().to_vec();
1065            assert_eq!(i.to_string(), String::from_utf8(buf).unwrap());
1066        }
1067    }
1068
1069    #[tokio::test]
1070    async fn test_file_cache_read_ranges() {
1071        let dir = create_temp_dir("");
1072        let local_store = new_fs_store(dir.path().to_str().unwrap());
1073        let file_cache = FileCache::new(
1074            local_store.clone(),
1075            ReadableSize::mb(10),
1076            None,
1077            None,
1078            true, // enable_background_worker
1079        );
1080        let region_id = RegionId::new(2000, 0);
1081        let file_id = FileId::random();
1082        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
1083        let file_path = file_cache.cache_file_path(key);
1084        // Write a file.
1085        let data = b"hello greptime database";
1086        local_store
1087            .write(&file_path, data.as_slice())
1088            .await
1089            .unwrap();
1090        // Add to the cache.
1091        file_cache.put(key, IndexValue { file_size: 5 }).await;
1092        // Ranges
1093        let ranges = vec![0..5, 6..10, 15..19, 0..data.len() as u64];
1094        let bytes = file_cache.read_ranges(key, &ranges).await.unwrap();
1095
1096        assert_eq!(4, bytes.len());
1097        assert_eq!(b"hello", bytes[0].as_ref());
1098        assert_eq!(b"grep", bytes[1].as_ref());
1099        assert_eq!(b"data", bytes[2].as_ref());
1100        assert_eq!(data, bytes[3].as_ref());
1101    }
1102
1103    #[test]
1104    fn test_file_cache_capacity_respects_total_budget() {
1105        let total_capacity = ReadableSize::mb(256).as_bytes();
1106        let (parquet_capacity, puffin_capacity) =
1107            FileCache::split_cache_capacities(total_capacity, 20);
1108
1109        assert_eq!(total_capacity, parquet_capacity + puffin_capacity);
1110        assert_eq!(ReadableSize::mb(128).as_bytes(), parquet_capacity);
1111        assert_eq!(ReadableSize::mb(128).as_bytes(), puffin_capacity);
1112    }
1113
1114    #[test]
1115    fn test_file_cache_capacity_keeps_split_when_total_allows_it() {
1116        let total_capacity = ReadableSize::gb(5).as_bytes();
1117        let (parquet_capacity, puffin_capacity) =
1118            FileCache::split_cache_capacities(total_capacity, 20);
1119
1120        assert_eq!(total_capacity, parquet_capacity + puffin_capacity);
1121        assert_eq!(ReadableSize::gb(4).as_bytes(), parquet_capacity);
1122        assert_eq!(ReadableSize::gb(1).as_bytes(), puffin_capacity);
1123    }
1124
1125    #[test]
1126    fn test_cache_file_path() {
1127        let file_id = FileId::parse_str("3368731b-a556-42b8-a5df-9c31ce155095").unwrap();
1128        assert_eq!(
1129            "test_dir/5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.parquet",
1130            cache_file_path(
1131                "test_dir",
1132                IndexKey::new(RegionId::new(1234, 5), file_id, FileType::Parquet)
1133            )
1134        );
1135        assert_eq!(
1136            "test_dir/5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.parquet",
1137            cache_file_path(
1138                "test_dir/",
1139                IndexKey::new(RegionId::new(1234, 5), file_id, FileType::Parquet)
1140            )
1141        );
1142    }
1143
1144    #[test]
1145    fn test_parse_file_name() {
1146        let file_id = FileId::parse_str("3368731b-a556-42b8-a5df-9c31ce155095").unwrap();
1147        let region_id = RegionId::new(1234, 5);
1148        assert_eq!(
1149            IndexKey::new(region_id, file_id, FileType::Parquet),
1150            parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.parquet").unwrap()
1151        );
1152        assert_eq!(
1153            IndexKey::new(region_id, file_id, FileType::Puffin(0)),
1154            parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.puffin").unwrap()
1155        );
1156        assert_eq!(
1157            IndexKey::new(region_id, file_id, FileType::Puffin(42)),
1158            parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.42.puffin")
1159                .unwrap()
1160        );
1161        assert!(parse_index_key("").is_none());
1162        assert!(parse_index_key(".").is_none());
1163        assert!(parse_index_key("5299989643269").is_none());
1164        assert!(parse_index_key("5299989643269.").is_none());
1165        assert!(parse_index_key(".5299989643269").is_none());
1166        assert!(parse_index_key("5299989643269.").is_none());
1167        assert!(parse_index_key("5299989643269.3368731b-a556-42b8-a5df").is_none());
1168        assert!(parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095").is_none());
1169        assert!(
1170            parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.parque").is_none()
1171        );
1172        assert!(
1173            parse_index_key("5299989643269.3368731b-a556-42b8-a5df-9c31ce155095.parquet.puffin")
1174                .is_none()
1175        );
1176    }
1177}