Skip to main content

mito2/cache/
write_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 write-through cache for remote object stores.
16
17use std::sync::Arc;
18use std::time::{Duration, Instant};
19
20use common_base::readable_size::ReadableSize;
21use common_telemetry::{debug, info};
22use futures::AsyncWriteExt;
23use object_store::ObjectStore;
24use snafu::ResultExt;
25use store_api::storage::RegionId;
26use tokio::sync::mpsc::{UnboundedSender, unbounded_channel};
27
28use crate::access_layer::{
29    FilePathProvider, Metrics, RegionFilePathFactory, SstInfoArray, SstWriteRequest,
30    TempFileCleaner, WriteCachePathProvider, WriteType, new_fs_cache_store,
31};
32use crate::cache::file_cache::{FileCache, FileCacheRef, FileType, IndexKey, IndexValue};
33use crate::cache::manifest_cache::ManifestCache;
34use crate::error::{self, Result};
35use crate::metrics::UPLOAD_BYTES_TOTAL;
36use crate::region::opener::RegionLoadCacheTask;
37use crate::sst::file::RegionFileId;
38use crate::sst::index::IndexerBuilderImpl;
39use crate::sst::index::intermediate::IntermediateManager;
40use crate::sst::index::puffin_manager::{PuffinManagerFactory, SstPuffinManager};
41use crate::sst::parquet::writer::ParquetWriter;
42use crate::sst::parquet::{SstInfo, WriteOptions};
43use crate::sst::{DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY};
44
45/// A cache for uploading files to remote object stores.
46///
47/// It keeps files in local disk and then sends files to object stores.
48pub struct WriteCache {
49    /// Local file cache.
50    file_cache: FileCacheRef,
51    /// Puffin manager factory for index.
52    puffin_manager_factory: PuffinManagerFactory,
53    /// Intermediate manager for index.
54    intermediate_manager: IntermediateManager,
55    /// Sender for region load cache tasks.
56    task_sender: UnboundedSender<RegionLoadCacheTask>,
57    /// Optional cache for manifest files.
58    manifest_cache: Option<ManifestCache>,
59}
60
61pub type WriteCacheRef = Arc<WriteCache>;
62
63impl WriteCache {
64    /// Create the cache with a `local_store` to cache files and a
65    /// `object_store_manager` for all object stores.
66    #[allow(clippy::too_many_arguments)]
67    pub async fn new(
68        local_store: ObjectStore,
69        cache_capacity: ReadableSize,
70        ttl: Option<Duration>,
71        index_cache_percent: Option<u8>,
72        enable_background_worker: bool,
73        puffin_manager_factory: PuffinManagerFactory,
74        intermediate_manager: IntermediateManager,
75        manifest_cache: Option<ManifestCache>,
76    ) -> Result<Self> {
77        let (task_sender, task_receiver) = unbounded_channel();
78
79        let file_cache = Arc::new(FileCache::new(
80            local_store,
81            cache_capacity,
82            ttl,
83            index_cache_percent,
84            enable_background_worker,
85        ));
86        file_cache.recover(false, Some(task_receiver)).await;
87
88        Ok(Self {
89            file_cache,
90            puffin_manager_factory,
91            intermediate_manager,
92            task_sender,
93            manifest_cache,
94        })
95    }
96
97    /// Creates a write cache based on local fs.
98    #[allow(clippy::too_many_arguments)]
99    pub async fn new_fs(
100        cache_dir: &str,
101        cache_capacity: ReadableSize,
102        ttl: Option<Duration>,
103        index_cache_percent: Option<u8>,
104        enable_background_worker: bool,
105        puffin_manager_factory: PuffinManagerFactory,
106        intermediate_manager: IntermediateManager,
107        manifest_cache_capacity: ReadableSize,
108    ) -> Result<Self> {
109        info!("Init write cache on {cache_dir}, capacity: {cache_capacity}");
110
111        let local_store = new_fs_cache_store(cache_dir).await?;
112
113        // Create manifest cache if capacity is non-zero
114        let manifest_cache = if manifest_cache_capacity.as_bytes() > 0 {
115            Some(ManifestCache::new(local_store.clone(), manifest_cache_capacity, ttl, false).await)
116        } else {
117            None
118        };
119
120        Self::new(
121            local_store,
122            cache_capacity,
123            ttl,
124            index_cache_percent,
125            enable_background_worker,
126            puffin_manager_factory,
127            intermediate_manager,
128            manifest_cache,
129        )
130        .await
131    }
132
133    /// Returns the file cache of the write cache.
134    pub(crate) fn file_cache(&self) -> FileCacheRef {
135        self.file_cache.clone()
136    }
137
138    /// Returns the manifest cache if available.
139    pub(crate) fn manifest_cache(&self) -> Option<ManifestCache> {
140        self.manifest_cache.clone()
141    }
142
143    /// Build the puffin manager
144    pub(crate) fn build_puffin_manager(&self) -> SstPuffinManager {
145        let store = self.file_cache.local_store();
146        let path_provider = WriteCachePathProvider::new(self.file_cache.clone());
147        self.puffin_manager_factory.build(store, path_provider)
148    }
149
150    /// Put encoded SST data to the cache and upload to the remote object store.
151    pub(crate) async fn put_and_upload_sst(
152        &self,
153        data: &bytes::Bytes,
154        region_id: RegionId,
155        sst_info: &SstInfo,
156        upload_request: SstUploadRequest,
157    ) -> Result<Metrics> {
158        let file_id = sst_info.file_id;
159        let mut metrics = Metrics::new(WriteType::Flush);
160
161        // Create index key for the SST file
162        let parquet_key = IndexKey::new(region_id, file_id, FileType::Parquet);
163
164        // Write to cache first
165        let cache_start = Instant::now();
166        let cache_path = self.file_cache.cache_file_path(parquet_key);
167        let store = self.file_cache.local_store();
168        let cleaner = TempFileCleaner::new(region_id, store.clone());
169        let write_res = store
170            .write(&cache_path, data.clone())
171            .await
172            .context(crate::error::OpenDalSnafu);
173        if let Err(e) = write_res {
174            cleaner.clean_by_file_id(file_id).await;
175            return Err(e);
176        }
177
178        metrics.write_batch = cache_start.elapsed();
179
180        // Upload to remote store
181        let upload_start = Instant::now();
182        let region_file_id = RegionFileId::new(region_id, file_id);
183        let remote_path = upload_request
184            .dest_path_provider
185            .build_sst_file_path(region_file_id);
186
187        if let Err(e) = self
188            .upload(parquet_key, &remote_path, &upload_request.remote_store)
189            .await
190        {
191            // Clean up cache on failure
192            self.remove(parquet_key).await;
193            return Err(e);
194        }
195
196        metrics.upload_parquet = upload_start.elapsed();
197        Ok(metrics)
198    }
199
200    /// Returns the intermediate manager of the write cache.
201    pub(crate) fn intermediate_manager(&self) -> &IntermediateManager {
202        &self.intermediate_manager
203    }
204
205    /// Writes SST to the cache and then uploads it to the remote object store.
206    pub(crate) async fn write_and_upload_sst(
207        &self,
208        write_request: SstWriteRequest,
209        upload_request: SstUploadRequest,
210        write_opts: &WriteOptions,
211        metrics: &mut Metrics,
212    ) -> Result<SstInfoArray> {
213        let region_id = write_request.metadata.region_id;
214
215        let store = self.file_cache.local_store();
216        let path_provider = WriteCachePathProvider::new(self.file_cache.clone());
217        let indexer = IndexerBuilderImpl {
218            build_type: write_request.op_type.into(),
219            metadata: write_request.metadata.clone(),
220            puffin_manager: self
221                .puffin_manager_factory
222                .build(store.clone(), path_provider.clone()),
223            write_cache_enabled: true,
224            intermediate_manager: self.intermediate_manager.clone(),
225            index_options: write_request.index_options,
226            inverted_index_config: write_request.inverted_index_config,
227            fulltext_index_config: write_request.fulltext_index_config,
228            bloom_filter_index_config: write_request.bloom_filter_index_config,
229            #[cfg(feature = "vector_index")]
230            vector_index_config: write_request.vector_index_config,
231        };
232
233        let cleaner = TempFileCleaner::new(region_id, store.clone());
234        // Write to FileCache.
235        let mut writer = ParquetWriter::new_with_object_store(
236            store.clone(),
237            write_request.metadata,
238            write_request.index_config,
239            indexer,
240            path_provider.clone(),
241            metrics,
242        )
243        .await
244        .with_file_cleaner(cleaner);
245
246        let sst_info = match write_request.sst_write_format {
247            crate::sst::FormatType::PrimaryKey => {
248                writer
249                    .write_all_flat_as_primary_key(
250                        write_request.source,
251                        write_request.max_sequence,
252                        write_opts,
253                    )
254                    .await?
255            }
256            crate::sst::FormatType::Flat => {
257                writer
258                    .write_all_flat(write_request.source, write_request.max_sequence, write_opts)
259                    .await?
260            }
261        };
262
263        // Upload sst file to remote object store.
264        if sst_info.is_empty() {
265            return Ok(sst_info);
266        }
267
268        let mut upload_tracker = UploadTracker::new(region_id);
269        let mut err = None;
270        let remote_store = &upload_request.remote_store;
271        for sst in &sst_info {
272            let parquet_key = IndexKey::new(region_id, sst.file_id, FileType::Parquet);
273            let parquet_path = upload_request
274                .dest_path_provider
275                .build_sst_file_path(RegionFileId::new(region_id, sst.file_id));
276            let start = Instant::now();
277            if let Err(e) = self.upload(parquet_key, &parquet_path, remote_store).await {
278                err = Some(e);
279                break;
280            }
281            metrics.upload_parquet += start.elapsed();
282            upload_tracker.push_uploaded_file(parquet_path);
283
284            if sst.index_metadata.file_size > 0 {
285                let puffin_key = IndexKey::new(region_id, sst.file_id, FileType::Puffin(0));
286                let puffin_path = upload_request
287                    .dest_path_provider
288                    .build_index_file_path(RegionFileId::new(region_id, sst.file_id));
289                let start = Instant::now();
290                if let Err(e) = self.upload(puffin_key, &puffin_path, remote_store).await {
291                    err = Some(e);
292                    break;
293                }
294                metrics.upload_puffin += start.elapsed();
295                upload_tracker.push_uploaded_file(puffin_path);
296            }
297        }
298
299        if let Some(err) = err {
300            // Cleans files on failure.
301            upload_tracker
302                .clean(&sst_info, &self.file_cache, remote_store)
303                .await;
304            return Err(err);
305        }
306
307        Ok(sst_info)
308    }
309
310    /// Removes a file from the cache by `index_key`.
311    pub(crate) async fn remove(&self, index_key: IndexKey) {
312        self.file_cache.remove(index_key).await
313    }
314
315    /// Downloads a file in `remote_path` from the remote object store to the local cache
316    /// (specified by `index_key`).
317    pub(crate) async fn download(
318        &self,
319        index_key: IndexKey,
320        remote_path: &str,
321        remote_store: &ObjectStore,
322        file_size: u64,
323    ) -> Result<()> {
324        self.file_cache
325            .download(index_key, remote_path, remote_store, file_size)
326            .await
327    }
328
329    /// Downloads the target file into write cache only when it is not cached.
330    ///
331    /// Returns `Ok(true)` if this call performs a download, or `Ok(false)` if the
332    /// file is already present in write cache and download is skipped.
333    pub(crate) async fn download_if_absent(
334        &self,
335        index_key: IndexKey,
336        remote_path: &str,
337        remote_store: &ObjectStore,
338        file_size: u64,
339    ) -> Result<bool> {
340        if self.file_cache.contains_key(&index_key) {
341            debug!(
342                "Skip downloading file already in write cache, region: {}, file: {}",
343                index_key.region_id, index_key.file_id
344            );
345            return Ok(false);
346        }
347
348        self.download(index_key, remote_path, remote_store, file_size)
349            .await?;
350        Ok(true)
351    }
352
353    /// Uploads a Parquet file or a Puffin file to the remote object store.
354    pub(crate) async fn upload(
355        &self,
356        index_key: IndexKey,
357        upload_path: &str,
358        remote_store: &ObjectStore,
359    ) -> Result<()> {
360        let region_id = index_key.region_id;
361        let file_id = index_key.file_id;
362        let file_type = index_key.file_type;
363        let cache_path = self.file_cache.cache_file_path(index_key);
364
365        let start = Instant::now();
366        let cached_value = self
367            .file_cache
368            .local_store()
369            .stat(&cache_path)
370            .await
371            .context(error::OpenDalSnafu)?;
372        let reader = self
373            .file_cache
374            .local_store()
375            .reader(&cache_path)
376            .await
377            .context(error::OpenDalSnafu)?
378            .into_futures_async_read(0..cached_value.content_length())
379            .await
380            .context(error::OpenDalSnafu)?;
381
382        let mut writer = remote_store
383            .writer_with(upload_path)
384            .chunk(DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize)
385            .concurrent(DEFAULT_WRITE_CONCURRENCY)
386            .await
387            .context(error::OpenDalSnafu)?
388            .into_futures_async_write();
389
390        let bytes_written =
391            futures::io::copy(reader, &mut writer)
392                .await
393                .context(error::UploadSnafu {
394                    region_id,
395                    file_id,
396                    file_type,
397                })?;
398
399        // Must close to upload all data.
400        writer.close().await.context(error::UploadSnafu {
401            region_id,
402            file_id,
403            file_type,
404        })?;
405
406        UPLOAD_BYTES_TOTAL.inc_by(bytes_written);
407
408        debug!(
409            "Successfully upload file to remote, region: {}, file: {}, upload_path: {}, cost: {:?}",
410            region_id,
411            file_id,
412            upload_path,
413            start.elapsed(),
414        );
415
416        let index_value = IndexValue {
417            file_size: bytes_written as _,
418        };
419        // Register to file cache
420        self.file_cache.put(index_key, index_value).await;
421
422        Ok(())
423    }
424
425    /// Sends a region load cache task to the background processing queue.
426    ///
427    /// If the receiver has been dropped, the error is ignored.
428    pub(crate) fn load_region_cache(&self, task: RegionLoadCacheTask) {
429        let _ = self.task_sender.send(task);
430    }
431}
432
433/// Request to write and upload a SST.
434pub struct SstUploadRequest {
435    /// Destination path provider of which SST files in write cache should be uploaded to.
436    pub dest_path_provider: RegionFilePathFactory,
437    /// Remote object store to upload.
438    pub remote_store: ObjectStore,
439}
440
441/// A structs to track files to upload and clean them if upload failed.
442pub(crate) struct UploadTracker {
443    /// Id of the region to track.
444    region_id: RegionId,
445    /// Paths of files uploaded successfully.
446    files_uploaded: Vec<String>,
447}
448
449impl UploadTracker {
450    /// Creates a new instance of `UploadTracker` for a given region.
451    pub(crate) fn new(region_id: RegionId) -> Self {
452        Self {
453            region_id,
454            files_uploaded: Vec::new(),
455        }
456    }
457
458    /// Add a file path to the list of uploaded files.
459    pub(crate) fn push_uploaded_file(&mut self, path: String) {
460        self.files_uploaded.push(path);
461    }
462
463    /// Cleans uploaded files and files in the file cache at best effort.
464    pub(crate) async fn clean(
465        &self,
466        sst_info: &SstInfoArray,
467        file_cache: &FileCacheRef,
468        remote_store: &ObjectStore,
469    ) {
470        common_telemetry::info!(
471            "Start cleaning files on upload failure, region: {}, num_ssts: {}",
472            self.region_id,
473            sst_info.len()
474        );
475
476        // Cleans files in the file cache first.
477        for sst in sst_info {
478            let parquet_key = IndexKey::new(self.region_id, sst.file_id, FileType::Parquet);
479            file_cache.remove(parquet_key).await;
480
481            if sst.index_metadata.file_size > 0 {
482                let puffin_key = IndexKey::new(
483                    self.region_id,
484                    sst.file_id,
485                    FileType::Puffin(sst.index_metadata.version),
486                );
487                file_cache.remove(puffin_key).await;
488            }
489        }
490
491        // Cleans uploaded files.
492        for file_path in &self.files_uploaded {
493            if let Err(e) = remote_store.delete(file_path).await {
494                common_telemetry::error!(e; "Failed to delete file {}", file_path);
495            }
496        }
497    }
498}
499
500#[cfg(test)]
501mod tests {
502    use bytes::Bytes;
503    use common_test_util::temp_dir::create_temp_dir;
504    use object_store::ATOMIC_WRITE_DIR;
505    use parquet::file::metadata::PageIndexPolicy;
506    use store_api::region_request::PathType;
507    use store_api::storage::FileId;
508
509    use super::*;
510    use crate::access_layer::OperationType;
511    use crate::cache::file_cache::IndexValue;
512    use crate::cache::test_util::{assert_parquet_metadata_equal, new_fs_store};
513    use crate::cache::{CacheManager, CacheStrategy};
514    use crate::error::InvalidBatchSnafu;
515    use crate::read::FlatSource;
516    use crate::region::options::IndexOptions;
517    use crate::sst::parquet::reader::ParquetReaderBuilder;
518    use crate::test_util::TestEnv;
519    use crate::test_util::sst_util::{
520        new_flat_source_from_record_batches, new_record_batch_by_range,
521        sst_file_handle_with_file_id, sst_region_metadata,
522    };
523
524    #[tokio::test]
525    async fn test_write_and_upload_sst() {
526        // TODO(QuenKar): maybe find a way to create some object server for testing,
527        // and now just use local file system to mock.
528        let mut env = TestEnv::new().await;
529        let mock_store = env.init_object_store_manager();
530        let path_provider = RegionFilePathFactory::new("test".to_string(), PathType::Bare);
531
532        let local_dir = create_temp_dir("");
533        let local_store = new_fs_store(local_dir.path().to_str().unwrap());
534
535        let write_cache = env
536            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
537            .await;
538
539        // Create source.
540        let metadata = Arc::new(sst_region_metadata());
541        let region_id = metadata.region_id;
542        let source = new_flat_source_from_record_batches(vec![
543            new_record_batch_by_range(&["a", "d"], 0, 60),
544            new_record_batch_by_range(&["b", "f"], 0, 40),
545            new_record_batch_by_range(&["b", "h"], 100, 200),
546        ]);
547
548        let write_request = SstWriteRequest {
549            op_type: OperationType::Flush,
550            metadata,
551            source,
552            storage: None,
553            max_sequence: None,
554            sst_write_format: Default::default(),
555            cache_manager: Default::default(),
556            index_options: IndexOptions::default(),
557            index_config: Default::default(),
558            inverted_index_config: Default::default(),
559            fulltext_index_config: Default::default(),
560            bloom_filter_index_config: Default::default(),
561            #[cfg(feature = "vector_index")]
562            vector_index_config: Default::default(),
563        };
564
565        let upload_request = SstUploadRequest {
566            dest_path_provider: path_provider.clone(),
567            remote_store: mock_store.clone(),
568        };
569
570        let write_opts = WriteOptions {
571            row_group_size: 512,
572            ..Default::default()
573        };
574
575        // Write to cache and upload sst to mock remote store
576        let mut metrics = Metrics::new(WriteType::Flush);
577        let mut sst_infos = write_cache
578            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
579            .await
580            .unwrap();
581        let sst_info = sst_infos.remove(0);
582
583        let file_id = sst_info.file_id;
584        let sst_upload_path =
585            path_provider.build_sst_file_path(RegionFileId::new(region_id, file_id));
586        let index_upload_path =
587            path_provider.build_index_file_path(RegionFileId::new(region_id, file_id));
588
589        // Check write cache contains the key
590        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
591        assert!(write_cache.file_cache.contains_key(&key));
592
593        // Check file data
594        let remote_data = mock_store.read(&sst_upload_path).await.unwrap();
595        let cache_data = local_store
596            .read(&write_cache.file_cache.cache_file_path(key))
597            .await
598            .unwrap();
599        assert_eq!(remote_data.to_vec(), cache_data.to_vec());
600
601        // Check write cache contains the index key
602        let index_key = IndexKey::new(region_id, file_id, FileType::Puffin(0));
603        assert!(write_cache.file_cache.contains_key(&index_key));
604
605        let remote_index_data = mock_store.read(&index_upload_path).await.unwrap();
606        let cache_index_data = local_store
607            .read(&write_cache.file_cache.cache_file_path(index_key))
608            .await
609            .unwrap();
610        assert_eq!(remote_index_data.to_vec(), cache_index_data.to_vec());
611
612        // Removes the file from the cache.
613        let sst_index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
614        write_cache.remove(sst_index_key).await;
615        assert!(!write_cache.file_cache.contains_key(&sst_index_key));
616        write_cache.remove(index_key).await;
617        assert!(!write_cache.file_cache.contains_key(&index_key));
618    }
619
620    #[tokio::test]
621    async fn test_read_metadata_from_write_cache() {
622        common_telemetry::init_default_ut_logging();
623        let mut env = TestEnv::new().await;
624        let data_home = env.data_home().display().to_string();
625        let mock_store = env.init_object_store_manager();
626
627        let local_dir = create_temp_dir("");
628        let local_path = local_dir.path().to_str().unwrap();
629        let local_store = new_fs_store(local_path);
630
631        // Create a cache manager using only write cache
632        let write_cache = env
633            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
634            .await;
635        let cache_manager = Arc::new(
636            CacheManager::builder()
637                .write_cache(Some(write_cache.clone()))
638                .build(),
639        );
640        assert!(!cache_manager.sst_meta_cache_enabled());
641
642        // Create source
643        let metadata = Arc::new(sst_region_metadata());
644
645        let source = new_flat_source_from_record_batches(vec![
646            new_record_batch_by_range(&["a", "d"], 0, 60),
647            new_record_batch_by_range(&["b", "f"], 0, 40),
648            new_record_batch_by_range(&["b", "h"], 100, 200),
649        ]);
650
651        // Write to local cache and upload sst to mock remote store
652        let write_request = SstWriteRequest {
653            op_type: OperationType::Flush,
654            metadata,
655            source,
656            storage: None,
657            max_sequence: None,
658            sst_write_format: Default::default(),
659            cache_manager: cache_manager.clone(),
660            index_options: IndexOptions::default(),
661            index_config: Default::default(),
662            inverted_index_config: Default::default(),
663            fulltext_index_config: Default::default(),
664            bloom_filter_index_config: Default::default(),
665            #[cfg(feature = "vector_index")]
666            vector_index_config: Default::default(),
667        };
668        let write_opts = WriteOptions {
669            row_group_size: 512,
670            ..Default::default()
671        };
672        let upload_request = SstUploadRequest {
673            dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
674            remote_store: mock_store.clone(),
675        };
676
677        let mut metrics = Metrics::new(WriteType::Flush);
678        let mut sst_infos = write_cache
679            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
680            .await
681            .unwrap();
682        let sst_info = sst_infos.remove(0);
683        let write_parquet_metadata = sst_info.file_metadata.unwrap();
684
685        // Read metadata from write cache without preparing an in-memory metadata cache entry.
686        let handle = sst_file_handle_with_file_id(sst_info.file_id, 0, 1000);
687        let builder = ParquetReaderBuilder::new(
688            data_home,
689            PathType::Bare,
690            handle.clone(),
691            mock_store.clone(),
692        )
693        .cache(CacheStrategy::EnableAll(cache_manager.clone()))
694        .page_index_policy(PageIndexPolicy::Optional);
695        let reader = builder.build().await.unwrap().unwrap();
696        let cached_write_parquet_metadata = crate::cache::CachedSstMeta::try_new(
697            "test.sst",
698            Arc::unwrap_or_clone(write_parquet_metadata),
699        )
700        .unwrap()
701        .parquet_metadata();
702
703        // Check parquet metadata
704        assert_parquet_metadata_equal(cached_write_parquet_metadata, reader.parquet_metadata());
705    }
706
707    #[tokio::test]
708    async fn test_write_cache_clean_tmp_files() {
709        common_telemetry::init_default_ut_logging();
710        let mut env = TestEnv::new().await;
711        let data_home = env.data_home().display().to_string();
712        let mock_store = env.init_object_store_manager();
713
714        let write_cache_dir = create_temp_dir("");
715        let write_cache_path = write_cache_dir.path().to_str().unwrap();
716        let write_cache = env
717            .create_write_cache_from_path(write_cache_path, ReadableSize::mb(10))
718            .await;
719
720        // Create a cache manager using only write cache
721        let cache_manager = Arc::new(
722            CacheManager::builder()
723                .write_cache(Some(write_cache.clone()))
724                .build(),
725        );
726
727        // Create source
728        let metadata = Arc::new(sst_region_metadata());
729
730        // Creates a source that can return an error to abort the writer.
731        let record_batch = new_record_batch_by_range(&["a", "d"], 0, 60);
732        let schema = record_batch.schema();
733        let iter = Box::new(
734            [
735                Ok(record_batch),
736                InvalidBatchSnafu {
737                    reason: "Abort the writer",
738                }
739                .fail(),
740            ]
741            .into_iter(),
742        );
743        let source = FlatSource::new_iter(schema, iter);
744
745        // Write to local cache and upload sst to mock remote store
746        let write_request = SstWriteRequest {
747            op_type: OperationType::Flush,
748            metadata,
749            source,
750            storage: None,
751            max_sequence: None,
752            sst_write_format: Default::default(),
753            cache_manager: cache_manager.clone(),
754            index_options: IndexOptions::default(),
755            index_config: Default::default(),
756            inverted_index_config: Default::default(),
757            fulltext_index_config: Default::default(),
758            bloom_filter_index_config: Default::default(),
759            #[cfg(feature = "vector_index")]
760            vector_index_config: Default::default(),
761        };
762        let write_opts = WriteOptions {
763            row_group_size: 512,
764            ..Default::default()
765        };
766        let upload_request = SstUploadRequest {
767            dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
768            remote_store: mock_store.clone(),
769        };
770
771        let mut metrics = Metrics::new(WriteType::Flush);
772        write_cache
773            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
774            .await
775            .unwrap_err();
776        let atomic_write_dir = write_cache_dir.path().join(ATOMIC_WRITE_DIR);
777        let mut entries = tokio::fs::read_dir(&atomic_write_dir).await.unwrap();
778        let mut has_files = false;
779        while let Some(entry) = entries.next_entry().await.unwrap() {
780            if entry.file_type().await.unwrap().is_dir() {
781                continue;
782            }
783            has_files = true;
784            common_telemetry::warn!(
785                "Found remaining temporary file in atomic dir: {}",
786                entry.path().display()
787            );
788        }
789
790        assert!(!has_files);
791    }
792
793    #[tokio::test]
794    async fn test_download_if_absent_skips_when_cached() {
795        let mut env = TestEnv::new().await;
796        let remote_store = env.init_object_store_manager();
797
798        let local_dir = create_temp_dir("");
799        let local_store = new_fs_store(local_dir.path().to_str().unwrap());
800        let write_cache = env
801            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
802            .await;
803
804        let region_id = RegionId::new(1024, 1);
805        let file_id = FileId::random();
806        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
807        write_cache
808            .file_cache()
809            .put(key, IndexValue { file_size: 1 })
810            .await;
811
812        let downloaded = write_cache
813            .download_if_absent(key, "missing/path.parquet", &remote_store, 1)
814            .await
815            .unwrap();
816
817        assert!(!downloaded);
818    }
819
820    #[tokio::test]
821    async fn test_download_if_absent_downloads_when_missing() {
822        let mut env = TestEnv::new().await;
823        let remote_store = env.init_object_store_manager();
824
825        let local_dir = create_temp_dir("");
826        let local_store = new_fs_store(local_dir.path().to_str().unwrap());
827        let write_cache = env
828            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
829            .await;
830
831        let region_id = RegionId::new(1024, 2);
832        let file_id = FileId::random();
833        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
834        let remote_path = format!("download-if-absent/{file_id}.parquet");
835        let remote_data = Bytes::from_static(b"download-if-absent-test");
836        remote_store
837            .write(&remote_path, remote_data.clone())
838            .await
839            .unwrap();
840
841        let downloaded = write_cache
842            .download_if_absent(key, &remote_path, &remote_store, remote_data.len() as u64)
843            .await
844            .unwrap();
845
846        assert!(downloaded);
847        assert!(write_cache.file_cache().contains_key(&key));
848
849        let cached_data = local_store
850            .read(&write_cache.file_cache().cache_file_path(key))
851            .await
852            .unwrap();
853        assert_eq!(cached_data.to_vec(), remote_data.to_vec());
854    }
855}