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