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).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            row_group_size: write_opts.row_group_size,
221            puffin_manager: self
222                .puffin_manager_factory
223                .build(store.clone(), path_provider.clone()),
224            write_cache_enabled: true,
225            intermediate_manager: self.intermediate_manager.clone(),
226            index_options: write_request.index_options,
227            inverted_index_config: write_request.inverted_index_config,
228            fulltext_index_config: write_request.fulltext_index_config,
229            bloom_filter_index_config: write_request.bloom_filter_index_config,
230            #[cfg(feature = "vector_index")]
231            vector_index_config: write_request.vector_index_config,
232        };
233
234        let cleaner = TempFileCleaner::new(region_id, store.clone());
235        // Write to FileCache.
236        let mut writer = ParquetWriter::new_with_object_store(
237            store.clone(),
238            write_request.metadata,
239            write_request.index_config,
240            indexer,
241            path_provider.clone(),
242            metrics,
243        )
244        .await
245        .with_file_cleaner(cleaner);
246
247        let sst_info = match write_request.source {
248            either::Left(source) => {
249                writer
250                    .write_all(source, write_request.max_sequence, write_opts)
251                    .await?
252            }
253            either::Right(flat_source) => writer.write_all_flat(flat_source, write_opts).await?,
254        };
255
256        // Upload sst file to remote object store.
257        if sst_info.is_empty() {
258            return Ok(sst_info);
259        }
260
261        let mut upload_tracker = UploadTracker::new(region_id);
262        let mut err = None;
263        let remote_store = &upload_request.remote_store;
264        for sst in &sst_info {
265            let parquet_key = IndexKey::new(region_id, sst.file_id, FileType::Parquet);
266            let parquet_path = upload_request
267                .dest_path_provider
268                .build_sst_file_path(RegionFileId::new(region_id, sst.file_id));
269            let start = Instant::now();
270            if let Err(e) = self.upload(parquet_key, &parquet_path, remote_store).await {
271                err = Some(e);
272                break;
273            }
274            metrics.upload_parquet += start.elapsed();
275            upload_tracker.push_uploaded_file(parquet_path);
276
277            if sst.index_metadata.file_size > 0 {
278                let puffin_key = IndexKey::new(region_id, sst.file_id, FileType::Puffin(0));
279                let puffin_path = upload_request
280                    .dest_path_provider
281                    .build_index_file_path(RegionFileId::new(region_id, sst.file_id));
282                let start = Instant::now();
283                if let Err(e) = self.upload(puffin_key, &puffin_path, remote_store).await {
284                    err = Some(e);
285                    break;
286                }
287                metrics.upload_puffin += start.elapsed();
288                upload_tracker.push_uploaded_file(puffin_path);
289            }
290        }
291
292        if let Some(err) = err {
293            // Cleans files on failure.
294            upload_tracker
295                .clean(&sst_info, &self.file_cache, remote_store)
296                .await;
297            return Err(err);
298        }
299
300        Ok(sst_info)
301    }
302
303    /// Removes a file from the cache by `index_key`.
304    pub(crate) async fn remove(&self, index_key: IndexKey) {
305        self.file_cache.remove(index_key).await
306    }
307
308    /// Downloads a file in `remote_path` from the remote object store to the local cache
309    /// (specified by `index_key`).
310    pub(crate) async fn download(
311        &self,
312        index_key: IndexKey,
313        remote_path: &str,
314        remote_store: &ObjectStore,
315        file_size: u64,
316    ) -> Result<()> {
317        self.file_cache
318            .download(index_key, remote_path, remote_store, file_size)
319            .await
320    }
321
322    /// Uploads a Parquet file or a Puffin file to the remote object store.
323    pub(crate) async fn upload(
324        &self,
325        index_key: IndexKey,
326        upload_path: &str,
327        remote_store: &ObjectStore,
328    ) -> Result<()> {
329        let region_id = index_key.region_id;
330        let file_id = index_key.file_id;
331        let file_type = index_key.file_type;
332        let cache_path = self.file_cache.cache_file_path(index_key);
333
334        let start = Instant::now();
335        let cached_value = self
336            .file_cache
337            .local_store()
338            .stat(&cache_path)
339            .await
340            .context(error::OpenDalSnafu)?;
341        let reader = self
342            .file_cache
343            .local_store()
344            .reader(&cache_path)
345            .await
346            .context(error::OpenDalSnafu)?
347            .into_futures_async_read(0..cached_value.content_length())
348            .await
349            .context(error::OpenDalSnafu)?;
350
351        let mut writer = remote_store
352            .writer_with(upload_path)
353            .chunk(DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize)
354            .concurrent(DEFAULT_WRITE_CONCURRENCY)
355            .await
356            .context(error::OpenDalSnafu)?
357            .into_futures_async_write();
358
359        let bytes_written =
360            futures::io::copy(reader, &mut writer)
361                .await
362                .context(error::UploadSnafu {
363                    region_id,
364                    file_id,
365                    file_type,
366                })?;
367
368        // Must close to upload all data.
369        writer.close().await.context(error::UploadSnafu {
370            region_id,
371            file_id,
372            file_type,
373        })?;
374
375        UPLOAD_BYTES_TOTAL.inc_by(bytes_written);
376
377        debug!(
378            "Successfully upload file to remote, region: {}, file: {}, upload_path: {}, cost: {:?}",
379            region_id,
380            file_id,
381            upload_path,
382            start.elapsed(),
383        );
384
385        let index_value = IndexValue {
386            file_size: bytes_written as _,
387        };
388        // Register to file cache
389        self.file_cache.put(index_key, index_value).await;
390
391        Ok(())
392    }
393
394    /// Sends a region load cache task to the background processing queue.
395    ///
396    /// If the receiver has been dropped, the error is ignored.
397    pub(crate) fn load_region_cache(&self, task: RegionLoadCacheTask) {
398        let _ = self.task_sender.send(task);
399    }
400}
401
402/// Request to write and upload a SST.
403pub struct SstUploadRequest {
404    /// Destination path provider of which SST files in write cache should be uploaded to.
405    pub dest_path_provider: RegionFilePathFactory,
406    /// Remote object store to upload.
407    pub remote_store: ObjectStore,
408}
409
410/// A structs to track files to upload and clean them if upload failed.
411pub(crate) struct UploadTracker {
412    /// Id of the region to track.
413    region_id: RegionId,
414    /// Paths of files uploaded successfully.
415    files_uploaded: Vec<String>,
416}
417
418impl UploadTracker {
419    /// Creates a new instance of `UploadTracker` for a given region.
420    pub(crate) fn new(region_id: RegionId) -> Self {
421        Self {
422            region_id,
423            files_uploaded: Vec::new(),
424        }
425    }
426
427    /// Add a file path to the list of uploaded files.
428    pub(crate) fn push_uploaded_file(&mut self, path: String) {
429        self.files_uploaded.push(path);
430    }
431
432    /// Cleans uploaded files and files in the file cache at best effort.
433    pub(crate) async fn clean(
434        &self,
435        sst_info: &SstInfoArray,
436        file_cache: &FileCacheRef,
437        remote_store: &ObjectStore,
438    ) {
439        common_telemetry::info!(
440            "Start cleaning files on upload failure, region: {}, num_ssts: {}",
441            self.region_id,
442            sst_info.len()
443        );
444
445        // Cleans files in the file cache first.
446        for sst in sst_info {
447            let parquet_key = IndexKey::new(self.region_id, sst.file_id, FileType::Parquet);
448            file_cache.remove(parquet_key).await;
449
450            if sst.index_metadata.file_size > 0 {
451                let puffin_key = IndexKey::new(
452                    self.region_id,
453                    sst.file_id,
454                    FileType::Puffin(sst.index_metadata.version),
455                );
456                file_cache.remove(puffin_key).await;
457            }
458        }
459
460        // Cleans uploaded files.
461        for file_path in &self.files_uploaded {
462            if let Err(e) = remote_store.delete(file_path).await {
463                common_telemetry::error!(e; "Failed to delete file {}", file_path);
464            }
465        }
466    }
467}
468
469#[cfg(test)]
470mod tests {
471    use common_test_util::temp_dir::create_temp_dir;
472    use object_store::ATOMIC_WRITE_DIR;
473    use store_api::region_request::PathType;
474
475    use super::*;
476    use crate::access_layer::OperationType;
477    use crate::cache::test_util::new_fs_store;
478    use crate::cache::{CacheManager, CacheStrategy};
479    use crate::error::InvalidBatchSnafu;
480    use crate::read::Source;
481    use crate::region::options::IndexOptions;
482    use crate::sst::parquet::reader::ParquetReaderBuilder;
483    use crate::test_util::TestEnv;
484    use crate::test_util::sst_util::{
485        assert_parquet_metadata_eq, new_batch_by_range, new_source, sst_file_handle_with_file_id,
486        sst_region_metadata,
487    };
488
489    #[tokio::test]
490    async fn test_write_and_upload_sst() {
491        // TODO(QuenKar): maybe find a way to create some object server for testing,
492        // and now just use local file system to mock.
493        let mut env = TestEnv::new().await;
494        let mock_store = env.init_object_store_manager();
495        let path_provider = RegionFilePathFactory::new("test".to_string(), PathType::Bare);
496
497        let local_dir = create_temp_dir("");
498        let local_store = new_fs_store(local_dir.path().to_str().unwrap());
499
500        let write_cache = env
501            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
502            .await;
503
504        // Create Source
505        let metadata = Arc::new(sst_region_metadata());
506        let region_id = metadata.region_id;
507        let source = new_source(&[
508            new_batch_by_range(&["a", "d"], 0, 60),
509            new_batch_by_range(&["b", "f"], 0, 40),
510            new_batch_by_range(&["b", "h"], 100, 200),
511        ]);
512
513        let write_request = SstWriteRequest {
514            op_type: OperationType::Flush,
515            metadata,
516            source: either::Left(source),
517            storage: None,
518            max_sequence: None,
519            cache_manager: Default::default(),
520            index_options: IndexOptions::default(),
521            index_config: Default::default(),
522            inverted_index_config: Default::default(),
523            fulltext_index_config: Default::default(),
524            bloom_filter_index_config: Default::default(),
525            #[cfg(feature = "vector_index")]
526            vector_index_config: Default::default(),
527        };
528
529        let upload_request = SstUploadRequest {
530            dest_path_provider: path_provider.clone(),
531            remote_store: mock_store.clone(),
532        };
533
534        let write_opts = WriteOptions {
535            row_group_size: 512,
536            ..Default::default()
537        };
538
539        // Write to cache and upload sst to mock remote store
540        let mut metrics = Metrics::new(WriteType::Flush);
541        let mut sst_infos = write_cache
542            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
543            .await
544            .unwrap();
545        let sst_info = sst_infos.remove(0);
546
547        let file_id = sst_info.file_id;
548        let sst_upload_path =
549            path_provider.build_sst_file_path(RegionFileId::new(region_id, file_id));
550        let index_upload_path =
551            path_provider.build_index_file_path(RegionFileId::new(region_id, file_id));
552
553        // Check write cache contains the key
554        let key = IndexKey::new(region_id, file_id, FileType::Parquet);
555        assert!(write_cache.file_cache.contains_key(&key));
556
557        // Check file data
558        let remote_data = mock_store.read(&sst_upload_path).await.unwrap();
559        let cache_data = local_store
560            .read(&write_cache.file_cache.cache_file_path(key))
561            .await
562            .unwrap();
563        assert_eq!(remote_data.to_vec(), cache_data.to_vec());
564
565        // Check write cache contains the index key
566        let index_key = IndexKey::new(region_id, file_id, FileType::Puffin(0));
567        assert!(write_cache.file_cache.contains_key(&index_key));
568
569        let remote_index_data = mock_store.read(&index_upload_path).await.unwrap();
570        let cache_index_data = local_store
571            .read(&write_cache.file_cache.cache_file_path(index_key))
572            .await
573            .unwrap();
574        assert_eq!(remote_index_data.to_vec(), cache_index_data.to_vec());
575
576        // Removes the file from the cache.
577        let sst_index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
578        write_cache.remove(sst_index_key).await;
579        assert!(!write_cache.file_cache.contains_key(&sst_index_key));
580        write_cache.remove(index_key).await;
581        assert!(!write_cache.file_cache.contains_key(&index_key));
582    }
583
584    #[tokio::test]
585    async fn test_read_metadata_from_write_cache() {
586        common_telemetry::init_default_ut_logging();
587        let mut env = TestEnv::new().await;
588        let data_home = env.data_home().display().to_string();
589        let mock_store = env.init_object_store_manager();
590
591        let local_dir = create_temp_dir("");
592        let local_path = local_dir.path().to_str().unwrap();
593        let local_store = new_fs_store(local_path);
594
595        // Create a cache manager using only write cache
596        let write_cache = env
597            .create_write_cache(local_store.clone(), ReadableSize::mb(10))
598            .await;
599        let cache_manager = Arc::new(
600            CacheManager::builder()
601                .write_cache(Some(write_cache.clone()))
602                .build(),
603        );
604
605        // Create source
606        let metadata = Arc::new(sst_region_metadata());
607
608        let source = new_source(&[
609            new_batch_by_range(&["a", "d"], 0, 60),
610            new_batch_by_range(&["b", "f"], 0, 40),
611            new_batch_by_range(&["b", "h"], 100, 200),
612        ]);
613
614        // Write to local cache and upload sst to mock remote store
615        let write_request = SstWriteRequest {
616            op_type: OperationType::Flush,
617            metadata,
618            source: either::Left(source),
619            storage: None,
620            max_sequence: None,
621            cache_manager: cache_manager.clone(),
622            index_options: IndexOptions::default(),
623            index_config: Default::default(),
624            inverted_index_config: Default::default(),
625            fulltext_index_config: Default::default(),
626            bloom_filter_index_config: Default::default(),
627            #[cfg(feature = "vector_index")]
628            vector_index_config: Default::default(),
629        };
630        let write_opts = WriteOptions {
631            row_group_size: 512,
632            ..Default::default()
633        };
634        let upload_request = SstUploadRequest {
635            dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
636            remote_store: mock_store.clone(),
637        };
638
639        let mut metrics = Metrics::new(WriteType::Flush);
640        let mut sst_infos = write_cache
641            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
642            .await
643            .unwrap();
644        let sst_info = sst_infos.remove(0);
645        let write_parquet_metadata = sst_info.file_metadata.unwrap();
646
647        // Read metadata from write cache
648        let handle = sst_file_handle_with_file_id(sst_info.file_id, 0, 1000);
649        let builder = ParquetReaderBuilder::new(
650            data_home,
651            PathType::Bare,
652            handle.clone(),
653            mock_store.clone(),
654        )
655        .cache(CacheStrategy::EnableAll(cache_manager.clone()));
656        let reader = builder.build().await.unwrap();
657
658        // Check parquet metadata
659        assert_parquet_metadata_eq(write_parquet_metadata, reader.parquet_metadata());
660    }
661
662    #[tokio::test]
663    async fn test_write_cache_clean_tmp_files() {
664        common_telemetry::init_default_ut_logging();
665        let mut env = TestEnv::new().await;
666        let data_home = env.data_home().display().to_string();
667        let mock_store = env.init_object_store_manager();
668
669        let write_cache_dir = create_temp_dir("");
670        let write_cache_path = write_cache_dir.path().to_str().unwrap();
671        let write_cache = env
672            .create_write_cache_from_path(write_cache_path, ReadableSize::mb(10))
673            .await;
674
675        // Create a cache manager using only write cache
676        let cache_manager = Arc::new(
677            CacheManager::builder()
678                .write_cache(Some(write_cache.clone()))
679                .build(),
680        );
681
682        // Create source
683        let metadata = Arc::new(sst_region_metadata());
684
685        // Creates a source that can return an error to abort the writer.
686        let source = Source::Iter(Box::new(
687            [
688                Ok(new_batch_by_range(&["a", "d"], 0, 60)),
689                InvalidBatchSnafu {
690                    reason: "Abort the writer",
691                }
692                .fail(),
693            ]
694            .into_iter(),
695        ));
696
697        // Write to local cache and upload sst to mock remote store
698        let write_request = SstWriteRequest {
699            op_type: OperationType::Flush,
700            metadata,
701            source: either::Left(source),
702            storage: None,
703            max_sequence: None,
704            cache_manager: cache_manager.clone(),
705            index_options: IndexOptions::default(),
706            index_config: Default::default(),
707            inverted_index_config: Default::default(),
708            fulltext_index_config: Default::default(),
709            bloom_filter_index_config: Default::default(),
710            #[cfg(feature = "vector_index")]
711            vector_index_config: Default::default(),
712        };
713        let write_opts = WriteOptions {
714            row_group_size: 512,
715            ..Default::default()
716        };
717        let upload_request = SstUploadRequest {
718            dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
719            remote_store: mock_store.clone(),
720        };
721
722        let mut metrics = Metrics::new(WriteType::Flush);
723        write_cache
724            .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
725            .await
726            .unwrap_err();
727        let atomic_write_dir = write_cache_dir.path().join(ATOMIC_WRITE_DIR);
728        let mut entries = tokio::fs::read_dir(&atomic_write_dir).await.unwrap();
729        let mut has_files = false;
730        while let Some(entry) = entries.next_entry().await.unwrap() {
731            if entry.file_type().await.unwrap().is_dir() {
732                continue;
733            }
734            has_files = true;
735            common_telemetry::warn!(
736                "Found remaining temporary file in atomic dir: {}",
737                entry.path().display()
738            );
739        }
740
741        assert!(!has_files);
742    }
743}