Skip to main content

mito2/
access_layer.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
15use std::path::Path;
16use std::sync::Arc;
17use std::time::{Duration, Instant};
18
19use async_stream::try_stream;
20use common_base::readable_size::ReadableSize;
21use common_runtime::runtime::RuntimeTrait;
22use common_telemetry::warn;
23use common_time::Timestamp;
24use futures::{Stream, TryStreamExt};
25use object_store::services::Fs;
26use object_store::util::{join_dir, with_instrument_layers};
27use object_store::{ATOMIC_WRITE_DIR, ErrorKind, OLD_ATOMIC_WRITE_DIR, ObjectStore};
28use parquet::file::metadata::PageIndexPolicy;
29use smallvec::SmallVec;
30use snafu::ResultExt;
31use store_api::metadata::RegionMetadataRef;
32use store_api::region_request::PathType;
33use store_api::sst_entry::StorageSstEntry;
34use store_api::storage::{FileId, RegionId, SequenceNumber};
35
36use crate::cache::file_cache::{FileCacheRef, FileType, IndexKey};
37use crate::cache::write_cache::SstUploadRequest;
38use crate::cache::{CacheManagerRef, SstMetaPreparation, prepare_sst_meta_sync};
39use crate::config::{BloomFilterConfig, FulltextIndexConfig, IndexConfig, InvertedIndexConfig};
40use crate::error::{
41    CleanDirSnafu, DeleteIndexSnafu, DeleteIndexesSnafu, DeleteSstsSnafu, OpenDalSnafu, Result,
42};
43use crate::metrics::{COMPACTION_STAGE_ELAPSED, FLUSH_ELAPSED};
44use crate::read::FlatSource;
45use crate::region::options::IndexOptions;
46use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId};
47use crate::sst::index::IndexerBuilderImpl;
48use crate::sst::index::intermediate::IntermediateManager;
49use crate::sst::index::puffin_manager::{PuffinManagerFactory, SstPuffinManager};
50use crate::sst::location::{self, region_dir_from_table_dir};
51use crate::sst::parquet::reader::ParquetReaderBuilder;
52use crate::sst::parquet::writer::ParquetWriter;
53use crate::sst::parquet::{SstInfo, WriteOptions};
54use crate::sst::{DEFAULT_WRITE_CONCURRENCY, FormatType};
55
56pub type AccessLayerRef = Arc<AccessLayer>;
57/// SST write results.
58pub type SstInfoArray = SmallVec<[SstInfo; 2]>;
59
60/// Write operation type.
61#[derive(Eq, PartialEq, Debug)]
62pub enum WriteType {
63    /// Writes from flush
64    Flush,
65    /// Writes from compaction.
66    Compaction,
67}
68
69#[derive(Debug)]
70pub struct Metrics {
71    pub(crate) write_type: WriteType,
72    pub(crate) iter_source: Duration,
73    pub(crate) write_batch: Duration,
74    pub(crate) update_index: Duration,
75    pub(crate) upload_parquet: Duration,
76    pub(crate) upload_puffin: Duration,
77    pub(crate) compact_memtable: Duration,
78}
79
80impl Metrics {
81    pub fn new(write_type: WriteType) -> Self {
82        Self {
83            write_type,
84            iter_source: Default::default(),
85            write_batch: Default::default(),
86            update_index: Default::default(),
87            upload_parquet: Default::default(),
88            upload_puffin: Default::default(),
89            compact_memtable: Default::default(),
90        }
91    }
92
93    pub(crate) fn merge(mut self, other: Self) -> Self {
94        assert_eq!(self.write_type, other.write_type);
95        self.iter_source += other.iter_source;
96        self.write_batch += other.write_batch;
97        self.update_index += other.update_index;
98        self.upload_parquet += other.upload_parquet;
99        self.upload_puffin += other.upload_puffin;
100        self.compact_memtable += other.compact_memtable;
101        self
102    }
103
104    pub(crate) fn observe(self) {
105        match self.write_type {
106            WriteType::Flush => {
107                FLUSH_ELAPSED
108                    .with_label_values(&["iter_source"])
109                    .observe(self.iter_source.as_secs_f64());
110                FLUSH_ELAPSED
111                    .with_label_values(&["write_batch"])
112                    .observe(self.write_batch.as_secs_f64());
113                FLUSH_ELAPSED
114                    .with_label_values(&["update_index"])
115                    .observe(self.update_index.as_secs_f64());
116                FLUSH_ELAPSED
117                    .with_label_values(&["upload_parquet"])
118                    .observe(self.upload_parquet.as_secs_f64());
119                FLUSH_ELAPSED
120                    .with_label_values(&["upload_puffin"])
121                    .observe(self.upload_puffin.as_secs_f64());
122                if !self.compact_memtable.is_zero() {
123                    FLUSH_ELAPSED
124                        .with_label_values(&["compact_memtable"])
125                        .observe(self.upload_puffin.as_secs_f64());
126                }
127            }
128            WriteType::Compaction => {
129                COMPACTION_STAGE_ELAPSED
130                    .with_label_values(&["iter_source"])
131                    .observe(self.iter_source.as_secs_f64());
132                COMPACTION_STAGE_ELAPSED
133                    .with_label_values(&["write_batch"])
134                    .observe(self.write_batch.as_secs_f64());
135                COMPACTION_STAGE_ELAPSED
136                    .with_label_values(&["update_index"])
137                    .observe(self.update_index.as_secs_f64());
138                COMPACTION_STAGE_ELAPSED
139                    .with_label_values(&["upload_parquet"])
140                    .observe(self.upload_parquet.as_secs_f64());
141                COMPACTION_STAGE_ELAPSED
142                    .with_label_values(&["upload_puffin"])
143                    .observe(self.upload_puffin.as_secs_f64());
144            }
145        };
146    }
147}
148
149/// A layer to access SST files under the same directory.
150pub struct AccessLayer {
151    table_dir: String,
152    /// Path type for generating file paths.
153    path_type: PathType,
154    /// Target object store.
155    object_store: ObjectStore,
156    /// Puffin manager factory for index.
157    puffin_manager_factory: PuffinManagerFactory,
158    /// Intermediate manager for inverted index.
159    intermediate_manager: IntermediateManager,
160}
161
162impl std::fmt::Debug for AccessLayer {
163    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
164        f.debug_struct("AccessLayer")
165            .field("table_dir", &self.table_dir)
166            .finish()
167    }
168}
169
170impl AccessLayer {
171    /// Returns a new [AccessLayer] for specific `table_dir`.
172    pub fn new(
173        table_dir: impl Into<String>,
174        path_type: PathType,
175        object_store: ObjectStore,
176        puffin_manager_factory: PuffinManagerFactory,
177        intermediate_manager: IntermediateManager,
178    ) -> AccessLayer {
179        AccessLayer {
180            table_dir: table_dir.into(),
181            path_type,
182            object_store,
183            puffin_manager_factory,
184            intermediate_manager,
185        }
186    }
187
188    /// Returns the directory of the table.
189    pub fn table_dir(&self) -> &str {
190        &self.table_dir
191    }
192
193    /// Returns the object store of the layer.
194    pub fn object_store(&self) -> &ObjectStore {
195        &self.object_store
196    }
197
198    /// Returns the path type of the layer.
199    pub fn path_type(&self) -> PathType {
200        self.path_type
201    }
202
203    /// Returns the puffin manager factory.
204    pub fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
205        &self.puffin_manager_factory
206    }
207
208    /// Returns the intermediate manager.
209    pub fn intermediate_manager(&self) -> &IntermediateManager {
210        &self.intermediate_manager
211    }
212
213    /// Build the puffin manager.
214    pub(crate) fn build_puffin_manager(&self) -> SstPuffinManager {
215        let store = self.object_store.clone();
216        let path_provider =
217            RegionFilePathFactory::new(self.table_dir().to_string(), self.path_type());
218        self.puffin_manager_factory.build(store, path_provider)
219    }
220
221    pub(crate) async fn delete_index(
222        &self,
223        index_file_id: RegionIndexId,
224    ) -> Result<(), crate::error::Error> {
225        let path = location::index_file_path(
226            &self.table_dir,
227            RegionIndexId::new(index_file_id.file_id, index_file_id.version),
228            self.path_type,
229        );
230        self.object_store
231            .delete(&path)
232            .await
233            .context(DeleteIndexSnafu {
234                file_id: index_file_id.file_id(),
235            })?;
236        Ok(())
237    }
238
239    pub(crate) async fn delete_ssts(
240        &self,
241        region_id: RegionId,
242        file_ids: &[FileId],
243    ) -> Result<(), crate::error::Error> {
244        if file_ids.is_empty() {
245            return Ok(());
246        }
247
248        let attempted_files = file_ids.to_vec();
249        // Deleter does not normalize leading slashes like Operator::delete does.
250        let paths: Vec<_> = file_ids
251            .iter()
252            .map(|file_id| {
253                location::sst_file_path(
254                    &self.table_dir,
255                    RegionFileId::new(region_id, *file_id),
256                    self.path_type,
257                )
258                .trim_start_matches('/')
259                .to_string()
260            })
261            .collect();
262
263        let mut deleter = self
264            .object_store
265            .deleter()
266            .await
267            .with_context(|_| DeleteSstsSnafu {
268                region_id,
269                file_ids: attempted_files.clone(),
270            })?;
271        deleter
272            .delete_iter(paths)
273            .await
274            .with_context(|_| DeleteSstsSnafu {
275                region_id,
276                file_ids: attempted_files.clone(),
277            })?;
278        deleter.close().await.with_context(|_| DeleteSstsSnafu {
279            region_id,
280            file_ids: attempted_files,
281        })?;
282
283        Ok(())
284    }
285
286    pub(crate) async fn delete_indexes(
287        &self,
288        index_ids: &[RegionIndexId],
289    ) -> Result<(), crate::error::Error> {
290        if index_ids.is_empty() {
291            return Ok(());
292        }
293
294        let file_ids: Vec<_> = index_ids
295            .iter()
296            .map(|index_id| index_id.file_id())
297            .collect();
298        let paths: Vec<_> = index_ids
299            .iter()
300            .map(|index_id| {
301                location::index_file_path(&self.table_dir, *index_id, self.path_type)
302                    .trim_start_matches('/')
303                    .to_string()
304            })
305            .collect();
306
307        let mut deleter = self
308            .object_store
309            .deleter()
310            .await
311            .context(DeleteIndexesSnafu {
312                file_ids: file_ids.clone(),
313            })?;
314        deleter
315            .delete_iter(paths)
316            .await
317            .context(DeleteIndexesSnafu {
318                file_ids: file_ids.clone(),
319            })?;
320        deleter
321            .close()
322            .await
323            .context(DeleteIndexesSnafu { file_ids })?;
324
325        Ok(())
326    }
327
328    /// Returns the directory of the region in the table.
329    pub fn build_region_dir(&self, region_id: RegionId) -> String {
330        region_dir_from_table_dir(&self.table_dir, region_id, self.path_type)
331    }
332
333    /// Returns a reader builder for specific `file`.
334    pub(crate) fn read_sst(&self, file: FileHandle) -> ParquetReaderBuilder {
335        ParquetReaderBuilder::new(
336            self.table_dir.clone(),
337            self.path_type,
338            file,
339            self.object_store.clone(),
340        )
341    }
342
343    /// Writes a SST with specific `file_id` and `metadata` to the layer.
344    ///
345    /// Returns the info of the SST. If no data written, returns None.
346    pub async fn write_sst(
347        &self,
348        request: SstWriteRequest,
349        write_opts: &WriteOptions,
350        metrics: &mut Metrics,
351    ) -> Result<SstInfoArray> {
352        let op_type = request.op_type;
353        let region_id = request.metadata.region_id;
354        let region_metadata = request.metadata.clone();
355        let cache_manager = request.cache_manager.clone();
356        let override_sequence = if request.preserve_row_sequence {
357            None
358        } else {
359            request.max_sequence
360        };
361
362        let sst_info = if let Some(write_cache) = cache_manager.write_cache() {
363            // Write to the write cache.
364            write_cache
365                .write_and_upload_sst(
366                    request,
367                    SstUploadRequest {
368                        dest_path_provider: RegionFilePathFactory::new(
369                            self.table_dir.clone(),
370                            self.path_type,
371                        ),
372                        remote_store: self.object_store.clone(),
373                    },
374                    write_opts,
375                    metrics,
376                )
377                .await?
378        } else {
379            // Write cache is disabled.
380            let store = self.object_store.clone();
381            let path_provider = RegionFilePathFactory::new(self.table_dir.clone(), self.path_type);
382            let indexer_builder = IndexerBuilderImpl {
383                build_type: request.op_type.into(),
384                metadata: request.metadata.clone(),
385                puffin_manager: self
386                    .puffin_manager_factory
387                    .build(store, path_provider.clone()),
388                write_cache_enabled: false,
389                intermediate_manager: self.intermediate_manager.clone(),
390                index_options: request.index_options,
391                inverted_index_config: request.inverted_index_config,
392                fulltext_index_config: request.fulltext_index_config,
393                bloom_filter_index_config: request.bloom_filter_index_config,
394            };
395            // We disable write cache on file system but we still use atomic write.
396            // TODO(yingwen): If we support other non-fs stores without the write cache, then
397            // we may have find a way to check whether we need the cleaner.
398            let cleaner = TempFileCleaner::new(region_id, self.object_store.clone());
399            let mut writer = ParquetWriter::new_with_object_store(
400                self.object_store.clone(),
401                request.metadata,
402                request.index_config,
403                indexer_builder,
404                path_provider,
405                metrics,
406            )
407            .await
408            .with_file_cleaner(cleaner);
409            match request.sst_write_format {
410                FormatType::PrimaryKey => {
411                    writer
412                        .write_all_flat_as_primary_key(
413                            request.source,
414                            override_sequence,
415                            write_opts,
416                        )
417                        .await?
418                }
419                FormatType::Flat => {
420                    writer
421                        .write_all_flat(request.source, override_sequence, write_opts)
422                        .await?
423                }
424            }
425        };
426
427        // Put parquet metadata to cache manager.
428        if !sst_info.is_empty() && cache_manager.sst_meta_cache_enabled() {
429            let runtime = match op_type {
430                OperationType::Compact => common_runtime::compact_runtime(),
431                OperationType::Flush => common_runtime::global_runtime(),
432            };
433            for sst in &sst_info {
434                if let Some(parquet_metadata) = &sst.file_metadata {
435                    let file_id = RegionFileId::new(region_id, sst.file_id);
436                    let file_path = format!(
437                        "region_id={}, file_id={}",
438                        file_id.region_id(),
439                        file_id.file_id()
440                    );
441                    let page_index_policy = if parquet_metadata.offset_index().is_some() {
442                        PageIndexPolicy::Optional
443                    } else {
444                        PageIndexPolicy::Skip
445                    };
446                    let parquet_metadata = parquet_metadata.clone();
447                    let region_metadata = region_metadata.clone();
448                    let cache_manager = cache_manager.clone();
449                    // Compact cache preparation is best-effort. Run the entire operation in one
450                    // detached blocking task so it neither blocks an async worker nor delays the
451                    // SST write.
452                    runtime.spawn_blocking(move || {
453                        match prepare_sst_meta_sync(
454                            &file_path,
455                            Arc::unwrap_or_clone(parquet_metadata),
456                            Some(region_metadata),
457                            page_index_policy,
458                        ) {
459                            Ok(SstMetaPreparation::Prepared(metadata)) => {
460                                cache_manager.put_prepared_sst_meta(file_id, metadata, true);
461                            }
462                            Ok(SstMetaPreparation::DecodedOnly { encoding_error, .. }) => warn!(
463                                encoding_error;
464                                "Failed to encode parquet metadata for cache, file: {}",
465                                file_path
466                            ),
467                            Err(err) => {
468                                warn!(err; "Failed to cache parquet metadata for {}", file_path);
469                            }
470                        }
471                    });
472                }
473            }
474        }
475
476        Ok(sst_info)
477    }
478
479    /// Puts encoded SST bytes to the write cache (if enabled) and uploads it to the object store.
480    pub(crate) async fn put_sst(
481        &self,
482        data: &bytes::Bytes,
483        region_file_id: RegionFileId,
484        cache_manager: &CacheManagerRef,
485        write_buffer_size: ReadableSize,
486    ) -> Result<Metrics> {
487        if let Some(write_cache) = cache_manager.write_cache() {
488            // Write to cache and upload to remote store
489            let upload_request = SstUploadRequest {
490                dest_path_provider: RegionFilePathFactory::new(
491                    self.table_dir.clone(),
492                    self.path_type,
493                ),
494                remote_store: self.object_store.clone(),
495            };
496            write_cache
497                .put_and_upload_sst(data, region_file_id, upload_request, write_buffer_size)
498                .await
499        } else {
500            let start = Instant::now();
501            let cleaner =
502                TempFileCleaner::new(region_file_id.region_id(), self.object_store.clone());
503            let path_provider = RegionFilePathFactory::new(self.table_dir.clone(), self.path_type);
504            let sst_file_path = path_provider.build_sst_file_path(region_file_id);
505            let mut writer = self
506                .object_store
507                .writer_with(&sst_file_path)
508                .chunk(write_buffer_size.as_bytes() as usize)
509                .concurrent(DEFAULT_WRITE_CONCURRENCY)
510                .await
511                .context(OpenDalSnafu)?;
512            if let Err(err) = writer.write(data.clone()).await.context(OpenDalSnafu) {
513                cleaner.clean_by_file_id(region_file_id.file_id()).await;
514                return Err(err);
515            }
516            if let Err(err) = writer.close().await.context(OpenDalSnafu) {
517                cleaner.clean_by_file_id(region_file_id.file_id()).await;
518                return Err(err);
519            }
520            let mut metrics = Metrics::new(WriteType::Flush);
521            metrics.write_batch = start.elapsed();
522            Ok(metrics)
523        }
524    }
525
526    /// Lists the SST entries from the storage layer in the table directory.
527    pub fn storage_sst_entries(&self) -> impl Stream<Item = Result<StorageSstEntry>> + use<> {
528        let object_store = self.object_store.clone();
529        let table_dir = self.table_dir.clone();
530
531        try_stream! {
532            let mut lister = object_store
533                .lister_with(table_dir.as_str())
534                .recursive(true)
535                .await
536                .context(OpenDalSnafu)?;
537
538            while let Some(entry) = lister.try_next().await.context(OpenDalSnafu)? {
539                let metadata = entry.metadata();
540                if metadata.is_dir() {
541                    continue;
542                }
543
544                let path = entry.path();
545                if !path.ends_with(".parquet") && !path.ends_with(".puffin") {
546                    continue;
547                }
548
549                let file_size = metadata.content_length();
550                let file_size = if file_size == 0 { None } else { Some(file_size) };
551                let last_modified_ms = metadata
552                    .last_modified()
553                    .map(|ts| Timestamp::new_millisecond(ts.into_inner().as_millisecond()));
554
555                let entry = StorageSstEntry {
556                    file_path: path.to_string(),
557                    file_size,
558                    last_modified_ms,
559                    node_id: None,
560                };
561
562                yield entry;
563            }
564        }
565    }
566}
567
568/// `OperationType` represents the origin of the `SstWriteRequest`.
569#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
570pub enum OperationType {
571    Flush,
572    Compact,
573}
574
575/// Contents to build a SST.
576pub struct SstWriteRequest {
577    pub op_type: OperationType,
578    pub metadata: RegionMetadataRef,
579    pub source: FlatSource,
580    pub cache_manager: CacheManagerRef,
581    /// Optional uniform row sequence for writes that do not preserve sequences.
582    /// Compaction passes `None` to retain the reader's effective input sequences.
583    pub max_sequence: Option<SequenceNumber>,
584    pub sst_write_format: FormatType,
585
586    pub preserve_row_sequence: bool,
587
588    /// Configs for index
589    pub index_options: IndexOptions,
590    pub index_config: IndexConfig,
591    pub inverted_index_config: InvertedIndexConfig,
592    pub fulltext_index_config: FulltextIndexConfig,
593    pub bloom_filter_index_config: BloomFilterConfig,
594}
595
596/// Cleaner to remove temp files on the atomic write dir.
597pub(crate) struct TempFileCleaner {
598    region_id: RegionId,
599    object_store: ObjectStore,
600}
601
602impl TempFileCleaner {
603    /// Constructs the cleaner for the region and store.
604    pub(crate) fn new(region_id: RegionId, object_store: ObjectStore) -> Self {
605        Self {
606            region_id,
607            object_store,
608        }
609    }
610
611    /// Removes the SST and index file from the local atomic dir by the file id.
612    /// This only removes the initial index, since the index version is always 0 for a new SST, this method should be safe to pass 0.
613    pub(crate) async fn clean_by_file_id(&self, file_id: FileId) {
614        let sst_key = IndexKey::new(self.region_id, file_id, FileType::Parquet).to_string();
615        let index_key = IndexKey::new(self.region_id, file_id, FileType::Puffin(0)).to_string();
616
617        Self::clean_atomic_dir_files(&self.object_store, &[&sst_key, &index_key]).await;
618    }
619
620    /// Removes the files from the local atomic dir by their names.
621    pub(crate) async fn clean_atomic_dir_files(
622        local_store: &ObjectStore,
623        names_to_remove: &[&str],
624    ) {
625        // We don't know the actual suffix of the file under atomic dir, so we have
626        // to list the dir. The cost should be acceptable as there won't be to many files.
627        let Ok(entries) = local_store.list(ATOMIC_WRITE_DIR).await.inspect_err(|e| {
628            if e.kind() != ErrorKind::NotFound {
629                common_telemetry::error!(e; "Failed to list tmp files for {:?}", names_to_remove)
630            }
631        }) else {
632            return;
633        };
634
635        // In our case, we can ensure the file id is unique so it is safe to remove all files
636        // with the same file id under the atomic write dir.
637        let actual_files: Vec<_> = entries
638            .into_iter()
639            .filter_map(|entry| {
640                if entry.metadata().is_dir() {
641                    return None;
642                }
643
644                // Remove name that matches files_to_remove.
645                let should_remove = names_to_remove
646                    .iter()
647                    .any(|file| entry.name().starts_with(file));
648                if should_remove {
649                    Some(entry.path().to_string())
650                } else {
651                    None
652                }
653            })
654            .collect();
655
656        common_telemetry::warn!(
657            "Clean files {:?} under atomic write dir for {:?}",
658            actual_files,
659            names_to_remove
660        );
661
662        if let Err(e) = local_store.delete_iter(actual_files).await {
663            common_telemetry::error!(e; "Failed to delete tmp file for {:?}", names_to_remove);
664        }
665    }
666}
667
668pub(crate) async fn new_fs_cache_store(root: &str) -> Result<ObjectStore> {
669    // Preserve native filesystem prefixes such as Windows UNC shares.
670    let atomic_write_dir = Path::new(root).join(ATOMIC_WRITE_DIR);
671    clean_dir(&atomic_write_dir.to_string_lossy()).await?;
672
673    // Compatible code. Remove this after a major release.
674    let old_atomic_temp_dir = join_dir(root, OLD_ATOMIC_WRITE_DIR);
675    clean_dir(&old_atomic_temp_dir).await?;
676
677    let builder = Fs::default()
678        .root(root)
679        .atomic_write_dir(&atomic_write_dir.to_string_lossy());
680    let store = ObjectStore::new(builder).context(OpenDalSnafu)?;
681
682    Ok(with_instrument_layers(store, false))
683}
684
685/// Clean the directory.
686async fn clean_dir(dir: &str) -> Result<()> {
687    if tokio::fs::try_exists(dir)
688        .await
689        .context(CleanDirSnafu { dir })?
690    {
691        tokio::fs::remove_dir_all(dir)
692            .await
693            .context(CleanDirSnafu { dir })?;
694    }
695
696    Ok(())
697}
698
699/// Path provider for SST file and index file.
700pub trait FilePathProvider: Send + Sync {
701    /// Creates index file path of given file id. Version default to 0, and not shown in the path.
702    fn build_index_file_path(&self, file_id: RegionFileId) -> String;
703
704    /// Creates index file path of given index id (with version support).
705    fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String;
706
707    /// Creates SST file path of given file id.
708    fn build_sst_file_path(&self, file_id: RegionFileId) -> String;
709}
710
711/// Path provider that builds paths in local write cache.
712#[derive(Clone)]
713pub(crate) struct WriteCachePathProvider {
714    file_cache: FileCacheRef,
715}
716
717impl WriteCachePathProvider {
718    /// Creates a new `WriteCachePathProvider` instance.
719    pub fn new(file_cache: FileCacheRef) -> Self {
720        Self { file_cache }
721    }
722}
723
724impl FilePathProvider for WriteCachePathProvider {
725    fn build_index_file_path(&self, file_id: RegionFileId) -> String {
726        let puffin_key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Puffin(0));
727        self.file_cache.cache_file_path(puffin_key)
728    }
729
730    fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
731        let puffin_key = IndexKey::new(
732            index_id.region_id(),
733            index_id.file_id(),
734            FileType::Puffin(index_id.version),
735        );
736        self.file_cache.cache_file_path(puffin_key)
737    }
738
739    fn build_sst_file_path(&self, file_id: RegionFileId) -> String {
740        let parquet_file_key =
741            IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
742        self.file_cache.cache_file_path(parquet_file_key)
743    }
744}
745
746/// Path provider that builds paths in region storage path.
747#[derive(Clone, Debug)]
748pub(crate) struct RegionFilePathFactory {
749    pub(crate) table_dir: String,
750    pub(crate) path_type: PathType,
751}
752
753impl RegionFilePathFactory {
754    /// Creates a new `RegionFilePathFactory` instance.
755    pub fn new(table_dir: String, path_type: PathType) -> Self {
756        Self {
757            table_dir,
758            path_type,
759        }
760    }
761}
762
763impl FilePathProvider for RegionFilePathFactory {
764    fn build_index_file_path(&self, file_id: RegionFileId) -> String {
765        location::index_file_path_legacy(&self.table_dir, file_id, self.path_type)
766    }
767
768    fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
769        location::index_file_path(&self.table_dir, index_id, self.path_type)
770    }
771
772    fn build_sst_file_path(&self, file_id: RegionFileId) -> String {
773        location::sst_file_path(&self.table_dir, file_id, self.path_type)
774    }
775}
776
777#[cfg(test)]
778mod tests {
779    use std::path::PathBuf;
780
781    use bytes::Bytes;
782    use common_test_util::temp_dir::create_temp_dir;
783
784    use super::*;
785    use crate::cache::CacheManager;
786    use crate::cache::test_util::new_fs_store;
787    use crate::test_util::TestEnv;
788    use crate::test_util::sst_util::WriteChunkRecorder;
789
790    #[rstest::rstest]
791    #[tokio::test]
792    async fn test_put_sst_uses_configured_buffer_size(
793        #[values(false, true)] enable_write_cache: bool,
794        #[values(1024, 4096)] chunk_size: usize,
795    ) {
796        let mut env = TestEnv::new().await;
797        let remote_chunks = WriteChunkRecorder::default();
798        let remote_store = env.init_object_store_manager().layer(remote_chunks.layer());
799        let local_dir = create_temp_dir("encoded-sst-cache");
800        let local_chunks = WriteChunkRecorder::default();
801        let local_store =
802            new_fs_store(local_dir.path().to_str().unwrap()).layer(local_chunks.layer());
803        let write_cache = if enable_write_cache {
804            Some(
805                env.create_write_cache(local_store.clone(), ReadableSize::mb(10))
806                    .await,
807            )
808        } else {
809            None
810        };
811        let cache_manager = Arc::new(
812            CacheManager::builder()
813                .write_cache(write_cache.clone())
814                .build(),
815        );
816        let access_layer = AccessLayer::new(
817            "test",
818            PathType::Bare,
819            remote_store.clone(),
820            env.get_puffin_manager(),
821            env.get_intermediate_manager(),
822        );
823        let file_id = RegionFileId::new(RegionId::new(1024, 1), FileId::random());
824        // Two full chunks and a partial tail distinguish the configured size from the default.
825        let encoded = Bytes::from(
826            (0..2 * chunk_size + 17)
827                .map(|i| (i % 251) as u8)
828                .collect::<Vec<_>>(),
829        );
830        access_layer
831            .put_sst(
832                &encoded,
833                file_id,
834                &cache_manager,
835                ReadableSize(chunk_size as u64),
836            )
837            .await
838            .unwrap();
839
840        let remote_path = location::sst_file_path("test", file_id, PathType::Bare);
841        remote_chunks.assert_chunks(&remote_path, chunk_size, encoded.len());
842        assert_eq!(
843            remote_store.read(&remote_path).await.unwrap().to_bytes(),
844            encoded
845        );
846        if let Some(write_cache) = write_cache {
847            let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
848            let cache_path = write_cache.file_cache().cache_file_path(key);
849            local_chunks.assert_chunks(&cache_path, chunk_size, encoded.len());
850            assert_eq!(
851                local_store.read(&cache_path).await.unwrap().to_bytes(),
852                encoded
853            );
854        }
855    }
856
857    #[tokio::test]
858    async fn test_new_fs_cache_store() {
859        let root = common_test_util::temp_dir::create_temp_dir("fs-cache-store");
860        let dirs = [
861            root.path().join(ATOMIC_WRITE_DIR),
862            PathBuf::from(join_dir(
863                root.path().to_str().unwrap(),
864                OLD_ATOMIC_WRITE_DIR,
865            )),
866        ];
867        for dir in &dirs {
868            tokio::fs::create_dir_all(dir).await.unwrap();
869            tokio::fs::write(dir.join("stale"), b"stale").await.unwrap();
870        }
871        let store = new_fs_cache_store(root.path().to_str().unwrap())
872            .await
873            .unwrap();
874        for dir in &dirs {
875            assert!(!dir.join("stale").exists());
876        }
877        store.write("index", "contents").await.unwrap();
878        assert_eq!(
879            tokio::fs::read(root.path().join("index")).await.unwrap(),
880            b"contents"
881        );
882    }
883}