1use 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>;
57pub type SstInfoArray = SmallVec<[SstInfo; 2]>;
59
60#[derive(Eq, PartialEq, Debug)]
62pub enum WriteType {
63 Flush,
65 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
149pub struct AccessLayer {
151 table_dir: String,
152 path_type: PathType,
154 object_store: ObjectStore,
156 puffin_manager_factory: PuffinManagerFactory,
158 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 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 pub fn table_dir(&self) -> &str {
190 &self.table_dir
191 }
192
193 pub fn object_store(&self) -> &ObjectStore {
195 &self.object_store
196 }
197
198 pub fn path_type(&self) -> PathType {
200 self.path_type
201 }
202
203 pub fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
205 &self.puffin_manager_factory
206 }
207
208 pub fn intermediate_manager(&self) -> &IntermediateManager {
210 &self.intermediate_manager
211 }
212
213 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 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 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 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 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_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 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 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 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 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 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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
570pub enum OperationType {
571 Flush,
572 Compact,
573}
574
575pub struct SstWriteRequest {
577 pub op_type: OperationType,
578 pub metadata: RegionMetadataRef,
579 pub source: FlatSource,
580 pub cache_manager: CacheManagerRef,
581 pub max_sequence: Option<SequenceNumber>,
584 pub sst_write_format: FormatType,
585
586 pub preserve_row_sequence: bool,
587
588 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
596pub(crate) struct TempFileCleaner {
598 region_id: RegionId,
599 object_store: ObjectStore,
600}
601
602impl TempFileCleaner {
603 pub(crate) fn new(region_id: RegionId, object_store: ObjectStore) -> Self {
605 Self {
606 region_id,
607 object_store,
608 }
609 }
610
611 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 pub(crate) async fn clean_atomic_dir_files(
622 local_store: &ObjectStore,
623 names_to_remove: &[&str],
624 ) {
625 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 let actual_files: Vec<_> = entries
638 .into_iter()
639 .filter_map(|entry| {
640 if entry.metadata().is_dir() {
641 return None;
642 }
643
644 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 let atomic_write_dir = Path::new(root).join(ATOMIC_WRITE_DIR);
671 clean_dir(&atomic_write_dir.to_string_lossy()).await?;
672
673 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
685async 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
699pub trait FilePathProvider: Send + Sync {
701 fn build_index_file_path(&self, file_id: RegionFileId) -> String;
703
704 fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String;
706
707 fn build_sst_file_path(&self, file_id: RegionFileId) -> String;
709}
710
711#[derive(Clone)]
713pub(crate) struct WriteCachePathProvider {
714 file_cache: FileCacheRef,
715}
716
717impl WriteCachePathProvider {
718 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#[derive(Clone, Debug)]
748pub(crate) struct RegionFilePathFactory {
749 pub(crate) table_dir: String,
750 pub(crate) path_type: PathType,
751}
752
753impl RegionFilePathFactory {
754 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 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}