1use 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
51const FILE_DIR: &str = "cache/object/write/";
55
56pub(crate) const DEFAULT_INDEX_CACHE_PERCENT: u8 = 20;
58
59const MIN_CACHE_CAPACITY: u64 = 512 * 1024 * 1024;
61
62const DOWNLOAD_TASK_CHANNEL_SIZE: usize = 64;
64
65struct DownloadTask {
67 index_key: IndexKey,
68 remote_path: String,
69 remote_store: ObjectStore,
70 file_size: u64,
71}
72
73#[derive(Debug)]
75struct FileCacheInner {
76 local_store: ObjectStore,
78 parquet_index: Cache<IndexKey, IndexValue>,
80 puffin_index: Cache<IndexKey, IndexValue>,
82}
83
84impl FileCacheInner {
85 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 fn cache_file_path(&self, key: IndexKey) -> String {
95 cache_file_path(FILE_DIR, key)
96 }
97
98 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 index.run_pending_tasks().await;
110 }
111
112 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 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 match key.file_type {
147 FileType::Parquet => parquet_size += size,
148 FileType::Puffin { .. } => puffin_size += size,
149 }
150 }
151 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 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 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 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 fn contains_key(&self, key: &IndexKey) -> bool {
275 self.memory_index(key.file_type).contains_key(key)
276 }
277}
278
279#[derive(Debug, Clone)]
282pub(crate) struct FileCache {
283 inner: Arc<FileCacheInner>,
285 puffin_capacity: u64,
287 download_task_tx: Option<Sender<DownloadTask>>,
289}
290
291pub(crate) type FileCacheRef = Arc<FileCache>;
292
293impl FileCache {
294 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 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 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 let inner = Arc::new(FileCacheInner {
335 local_store,
336 parquet_index,
337 puffin_index,
338 });
339
340 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 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 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 let _ = inner
375 .download(
376 task.index_key,
377 &task.remote_path,
378 &task.remote_store,
379 task.file_size,
380 1, )
382 .await;
383 }
384 info!("Background download worker stopped");
385 });
386 }
387
388 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 value.file_size
401 })
402 .max_capacity(capacity)
403 .async_eviction_listener(move |key, value, cause| {
404 let store = cache_store.clone();
405 let file_path = cache_file_path(FILE_DIR, *key);
407 async move {
408 if let RemovalCause::Replaced = cause {
409 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 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 #[allow(unused)]
445 pub(crate) async fn reader(&self, key: IndexKey) -> Option<Reader> {
446 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 index.remove(&key).await;
474 CACHE_MISS
475 .with_label_values(&[key.file_type.metric_label()])
476 .inc();
477 None
478 }
479
480 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 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 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 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 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 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 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 pub(crate) fn cache_file_path(&self, key: IndexKey) -> String {
568 self.inner.cache_file_path(key)
569 }
570
571 pub(crate) fn local_store(&self) -> ObjectStore {
573 self.inner.local_store.clone()
574 }
575
576 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 if let Some(index_value) = self.inner.parquet_index.get(&key).await {
586 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 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 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 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 pub(crate) fn contains_key(&self, key: &IndexKey) -> bool {
689 self.inner.contains_key(key)
690 }
691
692 pub(crate) fn puffin_cache_capacity(&self) -> u64 {
694 self.puffin_capacity
695 }
696
697 pub(crate) fn puffin_cache_size(&self) -> u64 {
699 self.inner.puffin_index.weighted_size()
700 }
701
702 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) .await
714 }
715
716 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 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 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#[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 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
783pub enum FileType {
784 Parquet,
786 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 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 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 fn metric_label(&self) -> &'static str {
819 match self {
820 FileType::Parquet => FILE_TYPE,
821 FileType::Puffin(_) => INDEX_TYPE,
822 }
823 }
824}
825
826#[derive(Debug, Clone)]
830pub(crate) struct IndexValue {
831 pub(crate) file_size: u32,
833}
834
835fn cache_file_path(cache_file_dir: &str, key: IndexKey) -> String {
839 join_path(cache_file_dir, &key.to_string())
840}
841
842fn 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, );
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 assert!(cache.reader(key).await.is_none());
886
887 local_store
889 .write(&file_path, b"hello".as_slice())
890 .await
891 .unwrap();
892
893 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, );
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 assert!(cache.reader(key).await.is_none());
928
929 local_store
931 .write(&file_path, b"hello".as_slice())
932 .await
933 .unwrap();
934 cache
936 .put(
937 IndexKey::new(region_id, file_id, FileType::Parquet),
938 IndexValue { file_size: 5 },
939 )
940 .await;
941
942 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 cache.inner.parquet_index.run_pending_tasks().await;
949 assert_eq!(5, cache.inner.parquet_index.weighted_size());
950
951 cache.remove(key).await;
953 assert!(cache.reader(key).await.is_none());
954
955 cache.inner.parquet_index.run_pending_tasks().await;
957
958 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, );
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 local_store
982 .write(&file_path, b"hello".as_slice())
983 .await
984 .unwrap();
985 cache
987 .put(
988 IndexKey::new(region_id, file_id, FileType::Parquet),
989 IndexValue { file_size: 5 },
990 )
991 .await;
992
993 local_store.delete(&file_path).await.unwrap();
995
996 assert!(cache.reader(key).await.is_none());
998 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, );
1013
1014 let region_id = RegionId::new(2000, 0);
1015 let file_type = FileType::Parquet;
1016 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 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 let cache = FileCache::new(
1039 local_store.clone(),
1040 ReadableSize::mb(10),
1041 None,
1042 None,
1043 true, );
1045 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 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, );
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 let data = b"hello greptime database";
1086 local_store
1087 .write(&file_path, data.as_slice())
1088 .await
1089 .unwrap();
1090 file_cache.put(key, IndexValue { file_size: 5 }).await;
1092 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}