1use std::{
2 borrow::Cow,
3 collections::HashSet,
4 fmt::Display,
5 hash::BuildHasherDefault,
6 io::{BufWriter, ErrorKind, Write},
7 mem::take,
8 ops::RangeInclusive,
9 path::{Path, PathBuf},
10 sync::{
11 OnceLock,
12 atomic::{AtomicBool, AtomicU32, Ordering},
13 },
14};
15
16use anyhow::{Context, Result, bail};
17use auto_hash_map::AutoSet;
18use byteorder::{BE, ReadBytesExt, WriteBytesExt};
19use dashmap::DashSet;
20#[cfg(feature = "mmap")]
21use either::Either;
22use fs_err::{self as fs, File, OpenOptions, ReadDir};
23use jiff::Timestamp;
24#[cfg(feature = "mmap")]
25use memmap2::Mmap;
26use nohash_hasher::BuildNoHashHasher;
27use parking_lot::{Mutex, RwLock};
28use rustc_hash::FxHasher;
29use serde::{Deserialize, Serialize};
30use smallvec::SmallVec;
31use tracing::span::EnteredSpan;
32
33pub use crate::compaction::selector::CompactConfig;
34#[cfg(feature = "mmap")]
35use crate::{AccessMode, mmap_helper::advise_mmap_for_persistence};
36use crate::{
37 DbConfig, FamilyKind, QueryKey,
38 arc_bytes::ArcBytes,
39 compaction::selector::{Compactable, get_merge_segments},
40 compression::{Compression, checksum_block, decompress_into_arc},
41 constants::{
42 DATA_THRESHOLD_PER_COMPACTED_FILE, KEY_BLOCK_AVG_SIZE, KEY_BLOCK_CACHE_SIZE,
43 MAX_ENTRIES_PER_COMPACTED_FILE, VALUE_BLOCK_AVG_SIZE, VALUE_BLOCK_CACHE_SIZE,
44 },
45 key::{StoreKey, hash_key},
46 lookup_entry::{IterValue, LookupEntry, LookupValue},
47 merge_iter::MergeIter,
48 meta_file::{MetaEntryFlags, MetaFile, MetaLookupResult, StaticSortedFileRange},
49 meta_file_builder::MetaFileBuilder,
50 parallel_scheduler::ParallelScheduler,
51 rc_bytes::RcBytes,
52 sst_filter::SstFilter,
53 static_sorted_file::{BlockCache, SstLookupResult, StaticSortedFileIter},
54 static_sorted_file_builder::{StaticSortedFileBuilderMeta, StreamingSstWriter},
55 write_batch::{FinishResult, NewFile, WriteBatch},
56};
57
58#[cfg(feature = "stats")]
59#[derive(Debug)]
60pub struct CacheStatistics {
61 pub hit_rate: f32,
62 pub fill: f32,
63 pub items: usize,
64 pub size: u64,
65 pub hits: u64,
66 pub misses: u64,
67}
68
69#[cfg(feature = "stats")]
70impl CacheStatistics {
71 fn new<Key, Val, We, B, L>(cache: &quick_cache::sync::Cache<Key, Val, We, B, L>) -> Self
72 where
73 Key: Eq + std::hash::Hash,
74 Val: Clone,
75 We: quick_cache::Weighter<Key, Val> + Clone,
76 B: std::hash::BuildHasher + Clone,
77 L: quick_cache::Lifecycle<Key, Val> + Clone,
78 {
79 let size = cache.weight();
80 let hits = cache.hits();
81 let misses = cache.misses();
82 Self {
83 hit_rate: hits as f32 / (hits + misses) as f32,
84 fill: size as f32 / cache.capacity() as f32,
85 items: cache.len(),
86 size,
87 hits,
88 misses,
89 }
90 }
91}
92
93#[cfg(feature = "stats")]
94#[derive(Debug)]
95pub struct Statistics {
96 pub meta_files: usize,
97 pub sst_files: usize,
98 pub key_block_cache: CacheStatistics,
99 pub value_block_cache: CacheStatistics,
100 pub hits: u64,
101 pub misses: u64,
102 pub miss_family: u64,
103 pub miss_range: u64,
104 pub miss_amqf: u64,
105 pub miss_key: u64,
106}
107
108#[cfg(feature = "stats")]
109#[derive(Default)]
110struct TrackedStats {
111 hits_deleted: std::sync::atomic::AtomicU64,
112 hits_small: std::sync::atomic::AtomicU64,
113 hits_blob: std::sync::atomic::AtomicU64,
114 miss_family: std::sync::atomic::AtomicU64,
115 miss_range: std::sync::atomic::AtomicU64,
116 miss_amqf: std::sync::atomic::AtomicU64,
117 miss_key: std::sync::atomic::AtomicU64,
118 miss_global: std::sync::atomic::AtomicU64,
119}
120
121enum ActiveWriteState {
123 Active(&'static str),
126 Error,
129}
130
131enum DeferredDeletion {
137 Sst(u32),
138 Meta(u32),
139 Blob(u32),
140}
141
142pub(crate) struct WriteOperationGuard<'a> {
149 active: &'a Mutex<Option<ActiveWriteState>>,
151 path: &'a Path,
153 seq_before: u32,
157 succeeded: bool,
159}
160
161impl WriteOperationGuard<'_> {
162 pub(crate) fn success(&mut self) {
166 self.succeeded = true;
167 }
168}
169
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
188pub struct CurrentDbVersion {
189 pub max_sequence_number: u32,
191 pub commit_time: Timestamp,
193}
194
195pub fn read_current_version(path: &Path) -> Result<Option<CurrentDbVersion>> {
200 let current_path = path.join("CURRENT");
201 let content = match fs::read(¤t_path) {
202 Ok(content) => content,
203 Err(e) if e.kind() == ErrorKind::NotFound => return Ok(None),
204 Err(e) => return Err(e).context("Failed to read CURRENT file"),
205 };
206
207 serde_json::from_slice::<CurrentDbVersion>(&content)
208 .with_context(|| {
209 format!(
210 "CURRENT file at {} is corrupt ({} bytes)",
211 current_path.display(),
212 content.len()
213 )
214 })
215 .map(Some)
216}
217
218fn commit_current(path: &Path, seq: u32) -> Result<()> {
227 let version: &CurrentDbVersion = &CurrentDbVersion {
228 max_sequence_number: seq,
229 commit_time: Timestamp::now(),
230 };
231 let mut contents =
232 serde_json::to_string(version).context("Failed to serialize the CURRENT file")?;
233 contents.push('\n');
234 let next_path = path.join("CURRENT.next");
235 let mut next_file = File::create(&next_path)?;
236 next_file.write_all(contents.as_bytes())?;
237 next_file.sync_data()?;
238 drop(next_file);
239 fs::rename(&next_path, path.join("CURRENT"))?;
240 #[cfg(not(windows))]
252 File::open(path)
253 .and_then(|dir| dir.sync_data())
254 .context("Failed to sync database directory after updating CURRENT")?;
255 Ok(())
256}
257
258fn delete_orphan_files(path: &Path, seq_before: u32) -> Result<()> {
263 commit_current(path, seq_before).context("Unable to restore CURRENT file")?;
266
267 for entry in fs::read_dir(path)? {
268 let entry = entry?;
269 let path = entry.path();
270 if let Some(ext) = path.extension().and_then(|s| s.to_str())
271 && let Some(seq) = path
272 .file_stem()
273 .and_then(|s| s.to_str())
274 .and_then(|s| s.parse::<u32>().ok())
275 && seq > seq_before
276 {
277 match ext {
278 "sst" | "meta" | "blob" | "del" => fs::remove_file(&path)?,
279 _ => {}
280 }
281 }
282 }
283 Ok(())
284}
285
286impl Drop for WriteOperationGuard<'_> {
287 fn drop(&mut self) {
288 if self.succeeded {
289 *self.active.lock() = None;
291 return;
292 }
293
294 match delete_orphan_files(self.path, self.seq_before) {
297 Ok(()) => *self.active.lock() = None,
298 Err(_) => *self.active.lock() = Some(ActiveWriteState::Error),
299 }
300 }
301}
302
303pub struct TurboPersistence<S: ParallelScheduler, const FAMILIES: usize> {
306 parallel_scheduler: S,
307 path: PathBuf,
309 read_only: bool,
312 inner: RwLock<Inner<FAMILIES>>,
314 is_empty: AtomicBool,
317 active_write_operation: Mutex<Option<ActiveWriteState>>,
320 deferred_deletions: Mutex<Vec<DeferredDeletion>>,
324 key_block_cache: OnceLock<BlockCache>,
328 value_block_cache: OnceLock<BlockCache>,
331 config: DbConfig<FAMILIES>,
333 #[cfg(feature = "stats")]
335 stats: TrackedStats,
336}
337
338struct Inner<const FAMILIES: usize> {
340 meta_files_by_family: [Vec<MetaFile>; FAMILIES],
344 current_sequence_number: u32,
346 accessed_key_hashes: [DashSet<u64, BuildNoHashHasher<u64>>; FAMILIES],
350}
351
352impl<const FAMILIES: usize> Inner<FAMILIES> {
353 fn is_empty(&self) -> bool {
354 self.meta_files_by_family.iter().all(Vec::is_empty)
355 }
356
357 fn push_meta_file(&mut self, meta_file: MetaFile) {
358 let family = meta_file.family() as usize;
359 debug_assert!(family < FAMILIES, "meta file family is out of bounds");
360 let shard = &mut self.meta_files_by_family[family];
361 debug_assert!(
362 shard.last().is_none_or(|previous| {
363 previous.sequence_number() < meta_file.sequence_number()
364 }),
365 "meta file appended out of sequence order for family {family}"
366 );
367 shard.push(meta_file);
368 }
369
370 #[cfg(debug_assertions)]
371 fn debug_assert_meta_invariants(&self) {
372 for (family, meta_files) in self.meta_files_by_family.iter().enumerate() {
373 debug_assert!(
374 meta_files
375 .iter()
376 .all(|meta| meta.family() as usize == family),
377 "meta file stored in the wrong family shard"
378 );
379 debug_assert!(
380 meta_files
381 .windows(2)
382 .all(|pair| pair[0].sequence_number() < pair[1].sequence_number()),
383 "meta files in family {family} are not in ascending sequence order"
384 );
385 }
386 }
387}
388
389pub struct CommitOptions {
390 new_meta_files: Vec<NewFile>,
391 new_sst_files: Vec<NewFile>,
392 new_blob_files: Vec<NewFile>,
393 sst_files_to_delete: Vec<DeletedFile>,
394 blob_seq_numbers_to_delete: Vec<u32>,
395 sequence_number: u32,
396 keys_written: u64,
397}
398
399#[derive(Clone, Copy)]
402struct DeletedFile {
403 seq: u32,
404 size: u64,
406}
407
408#[derive(Clone, Copy, Debug, Default)]
411pub struct CommitStats {
412 pub bytes_written: u64,
414 pub bytes_deleted: u64,
416}
417
418impl Display for CommitStats {
419 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
420 let CommitStats {
421 bytes_written,
422 bytes_deleted,
423 } = self;
424 write!(
425 f,
426 "bytes_written={bytes_written} bytes_deleted={bytes_deleted}"
427 )
428 }
429}
430
431struct OpenOpts<S: ParallelScheduler, const FAMILIES: usize> {
432 path: PathBuf,
433 read_only: bool,
434 parallel_scheduler: S,
435 config: DbConfig<FAMILIES>,
436}
437
438impl<S: ParallelScheduler + Default, const FAMILIES: usize> TurboPersistence<S, FAMILIES> {
439 pub fn open(path: PathBuf) -> Result<Self> {
444 Self::open_with_parallel_scheduler(path, Default::default())
445 }
446
447 pub fn open_with_config(path: PathBuf, config: DbConfig<FAMILIES>) -> Result<Self> {
449 Self::open_with_config_and_parallel_scheduler(path, config, Default::default())
450 }
451
452 pub fn open_read_only_with_config(path: PathBuf, config: DbConfig<FAMILIES>) -> Result<Self> {
455 Self::open_read_only_with_parallel_scheduler(path, config, Default::default())
456 }
457
458 pub fn empty_in_memory_with_config(config: DbConfig<FAMILIES>) -> Self {
462 Self::new(OpenOpts {
465 path: PathBuf::new(),
466 read_only: true,
467 parallel_scheduler: Default::default(),
468 config,
469 })
470 }
471}
472
473impl<S: ParallelScheduler, const FAMILIES: usize> TurboPersistence<S, FAMILIES> {
474 fn new(
475 OpenOpts {
476 path,
477 read_only,
478 parallel_scheduler,
479 config,
480 }: OpenOpts<S, FAMILIES>,
481 ) -> Self {
482 Self {
483 parallel_scheduler,
484 path,
485 read_only,
486 inner: RwLock::new(Inner {
487 meta_files_by_family: [(); FAMILIES].map(|_| Vec::new()),
488 current_sequence_number: 0,
489 accessed_key_hashes: [(); FAMILIES]
490 .map(|_| DashSet::with_hasher(BuildNoHashHasher::default())),
491 }),
492 is_empty: AtomicBool::new(true),
493 active_write_operation: Mutex::new(None),
494 deferred_deletions: Mutex::new(Vec::new()),
495 key_block_cache: OnceLock::new(),
496 value_block_cache: OnceLock::new(),
497 config,
498 #[cfg(feature = "stats")]
499 stats: TrackedStats::default(),
500 }
501 }
502
503 pub fn open_with_parallel_scheduler(path: PathBuf, parallel_scheduler: S) -> Result<Self> {
508 Self::open_with_config_and_parallel_scheduler(path, DbConfig::default(), parallel_scheduler)
509 }
510
511 pub fn open_with_config_and_parallel_scheduler(
513 path: PathBuf,
514 config: DbConfig<FAMILIES>,
515 parallel_scheduler: S,
516 ) -> Result<Self> {
517 let mut db = Self::new(OpenOpts {
518 path,
519 read_only: false,
520 parallel_scheduler,
521 config,
522 });
523 db.open_directory(false)?;
524 Ok(db)
525 }
526
527 fn open_read_only_with_parallel_scheduler(
530 path: PathBuf,
531 config: DbConfig<FAMILIES>,
532 parallel_scheduler: S,
533 ) -> Result<Self> {
534 let mut db = Self::new(OpenOpts {
535 path,
536 read_only: true,
537 parallel_scheduler,
538 config,
539 });
540 db.open_directory(false)?;
541 Ok(db)
542 }
543
544 fn open_directory(&mut self, read_only: bool) -> Result<()> {
546 match fs::read_dir(&self.path) {
547 Ok(entries) => {
548 if !self
549 .load_directory(entries, read_only)
550 .context("Loading persistence directory failed")?
551 {
552 if read_only {
553 bail!("Failed to open database");
554 }
555 commit_current(&self.path, 0)
556 .context("Initializing persistence directory failed")?;
557 }
558 Ok(())
559 }
560 Err(e) => {
561 if !read_only && e.kind() == ErrorKind::NotFound {
562 self.create_and_init_directory()
563 .context("Creating and initializing persistence directory failed")?;
564 Ok(())
565 } else {
566 Err(e).context("Failed to open database")
567 }
568 }
569 }
570 }
571
572 fn create_and_init_directory(&mut self) -> Result<()> {
574 fs::create_dir_all(&self.path)?;
575 commit_current(&self.path, 0)
576 }
577
578 fn load_directory(&mut self, entries: ReadDir, read_only: bool) -> Result<bool> {
580 let mut meta_files = Vec::new();
581 let current = match read_current_version(&self.path)? {
582 Some(version) => version.max_sequence_number,
583 None if !read_only => return Ok(false),
584 None => bail!("Failed to open database: CURRENT file is missing"),
585 };
586
587 let mut deleted_files = HashSet::new();
588 for entry in entries {
589 let entry = entry?;
590 let path = entry.path();
591 if let Some(ext) = path.extension().and_then(|s| s.to_str()) {
592 if path.file_stem().and_then(|s| s.to_str()) == Some("CURRENT") {
596 if !read_only {
597 fs::remove_file(&path)?;
598 }
599 continue;
600 }
601 let seq: u32 = path
602 .file_stem()
603 .context("File has no file stem")?
604 .to_str()
605 .context("File stem is not valid utf-8")?
606 .parse()?;
607 if deleted_files.contains(&seq) {
608 continue;
609 }
610 if seq > current {
611 if !read_only {
612 fs::remove_file(&path)?;
613 }
614 } else {
615 match ext {
616 "meta" => {
617 meta_files.push(seq);
618 }
619 "del" => {
620 let mut content = &*fs::read(&path)?;
621 let mut no_existing_files = true;
622 while !content.is_empty() {
623 let seq = content.read_u32::<BE>()?;
624 deleted_files.insert(seq);
625 if !read_only {
626 let sst_file = self.path.join(format!("{seq:08}.sst"));
628 let meta_file = self.path.join(format!("{seq:08}.meta"));
629 let blob_file = self.path.join(format!("{seq:08}.blob"));
630 for path in [sst_file, meta_file, blob_file] {
631 if fs::exists(&path)? {
632 fs::remove_file(path)?;
633 no_existing_files = false;
634 }
635 }
636 }
637 }
638 if !read_only && no_existing_files {
639 fs::remove_file(&path)?;
640 }
641 }
642 "blob" | "sst" => {
643 }
645 _ => {
646 if !path
647 .file_name()
648 .is_some_and(|s| s.as_encoded_bytes().starts_with(b"."))
649 {
650 bail!("Unexpected file in persistence directory: {:?}", path);
651 }
652 }
653 }
654 }
655 } else {
656 match path.file_stem().and_then(|s| s.to_str()) {
657 Some("CURRENT") => {
658 }
660 Some("LOG") => {
661 }
663 _ => {
664 if !path
665 .file_name()
666 .is_some_and(|s| s.as_encoded_bytes().starts_with(b"."))
667 {
668 bail!("Unexpected file in persistence directory: {:?}", path);
669 }
670 }
671 }
672 }
673 }
674
675 meta_files.retain(|seq| !deleted_files.contains(seq));
676 meta_files.sort_unstable();
677 let mut meta_files = self
678 .parallel_scheduler
679 .parallel_map_collect::<_, _, Result<Vec<MetaFile>>>(&meta_files, |&seq| {
680 let meta_file = MetaFile::open(
681 &self.path,
682 seq,
683 Some(&self.config.family_configs),
684 self.config.access_mode,
685 )?;
686 Ok(meta_file)
687 })?;
688
689 let mut sst_filter = SstFilter::new();
690 for meta_file in meta_files.iter_mut().rev() {
691 sst_filter.apply_filter(meta_file);
692 }
693
694 let inner = self.inner.get_mut();
695 for meta_file in meta_files {
696 inner.push_meta_file(meta_file);
697 }
698 #[cfg(debug_assertions)]
699 inner.debug_assert_meta_invariants();
700 inner.current_sequence_number = current;
701 self.is_empty.store(inner.is_empty(), Ordering::Relaxed);
702 Ok(true)
703 }
704
705 #[tracing::instrument(level = "info", name = "reading database blob", skip_all)]
707 fn read_blob(&self, seq: u32, compression: Compression) -> Result<ArcBytes> {
708 let path = self.path.join(format!("{seq:08}.blob"));
709 #[cfg(feature = "mmap")]
710 let file = File::open(&path)?;
711 #[cfg(feature = "mmap")]
712 let data: Either<Mmap, Vec<u8>> = match self.config.access_mode {
713 AccessMode::Mmap => {
714 let mmap = unsafe { Mmap::map(file.file()) }.with_context(|| {
715 format!(
716 "Failed to mmap blob file {} ({} bytes)",
717 path.display(),
718 file.metadata().map(|m| m.len()).unwrap_or(0)
719 )
720 })?;
721 #[cfg(unix)]
722 mmap.advise(memmap2::Advice::Sequential)?;
723 #[cfg(unix)]
724 mmap.advise(memmap2::Advice::WillNeed)?;
725 advise_mmap_for_persistence(&mmap)?;
726 Either::Left(mmap)
727 }
728 AccessMode::File => Either::Right(fs::read(&path)?),
729 };
730 #[cfg(feature = "mmap")]
731 let mut reader: &[u8] = match &data {
732 Either::Left(mmap) => mmap,
733 Either::Right(bytes) => bytes,
734 };
735 #[cfg(not(feature = "mmap"))]
737 let data = fs::read(&path)?;
738 #[cfg(not(feature = "mmap"))]
739 let mut reader: &[u8] = &data;
740 let uncompressed_length = reader
741 .read_u32::<BE>()
742 .context("Failed to read uncompressed length from blob file")?;
743 let expected_checksum = reader.read_u32::<BE>()?;
744
745 let actual_checksum = checksum_block(reader);
747 if actual_checksum != expected_checksum {
748 bail!(
749 "Cache corruption detected: checksum mismatch in blob file {:08}.blob (expected \
750 {:08x}, got {:08x})",
751 seq,
752 expected_checksum,
753 actual_checksum
754 );
755 }
756
757 let buffer = decompress_into_arc(compression, uncompressed_length, reader)?;
758 Ok(ArcBytes::from(buffer))
759 }
760
761 pub fn is_empty(&self) -> bool {
763 self.is_empty.load(Ordering::Relaxed)
764 }
765
766 pub fn has_unrecoverable_write_error(&self) -> bool {
769 matches!(
770 *self.active_write_operation.lock(),
771 Some(ActiveWriteState::Error)
772 )
773 }
774
775 fn acquire_write_operation(&self, name: &'static str) -> Result<WriteOperationGuard<'_>> {
779 if self.read_only {
780 bail!("Cannot perform write operations on a read-only database");
781 }
782 let mut slot = self.active_write_operation.lock();
783 match &*slot {
784 Some(ActiveWriteState::Active(active_name)) => {
785 bail!(
786 "Another {active_name} is already active (only a single write operation is \
787 allowed at a time)"
788 );
789 }
790 Some(ActiveWriteState::Error) => {
791 bail!(
792 "A previous write operation failed with an unrecoverable error; no further \
793 writes are possible"
794 );
795 }
796 None => {}
797 }
798 *slot = Some(ActiveWriteState::Active(name));
799 drop(slot); let seq_before = self.inner.read().current_sequence_number;
801 Ok(WriteOperationGuard {
802 active: &self.active_write_operation,
803 path: &self.path,
804 seq_before,
805 succeeded: false,
806 })
807 }
808
809 pub fn write_batch<K: StoreKey + Send + Sync>(&self) -> Result<WriteBatch<'_, K, S, FAMILIES>> {
814 let guard = self.acquire_write_operation("write batch")?;
815 let current = guard.seq_before;
817 Ok(WriteBatch::new(
818 guard,
819 self.path.clone(),
820 current,
821 self.parallel_scheduler.clone(),
822 self.config.family_configs,
823 ))
824 }
825
826 fn key_block_cache(&self) -> &BlockCache {
827 self.key_block_cache.get_or_init(|| {
828 BlockCache::with(
829 KEY_BLOCK_CACHE_SIZE as usize / KEY_BLOCK_AVG_SIZE,
830 KEY_BLOCK_CACHE_SIZE,
831 Default::default(),
832 Default::default(),
833 Default::default(),
834 )
835 })
836 }
837
838 fn value_block_cache(&self) -> &BlockCache {
839 self.value_block_cache.get_or_init(|| {
840 BlockCache::with(
841 VALUE_BLOCK_CACHE_SIZE as usize / VALUE_BLOCK_AVG_SIZE,
842 VALUE_BLOCK_CACHE_SIZE,
843 Default::default(),
844 Default::default(),
845 Default::default(),
846 )
847 })
848 }
849
850 pub fn clear_cache(&self) {
852 self.clear_block_caches();
853 for meta in self.inner.write().meta_files_by_family.iter_mut().flatten() {
854 meta.clear_cache();
855 }
856 }
857
858 pub fn clear_block_caches(&self) {
861 if let Some(cache) = self.key_block_cache.get() {
862 cache.clear();
863 }
864 if let Some(cache) = self.value_block_cache.get() {
865 cache.clear();
866 }
867 }
868
869 pub fn prepare_all_sst_caches(&self) {
872 for meta in self.inner.write().meta_files_by_family.iter_mut().flatten() {
873 meta.prepare_sst_cache();
874 }
875 }
876
877 fn open_log(&self) -> Result<BufWriter<File>> {
878 if self.read_only {
879 unreachable!("Only write operations can open the log file");
880 }
881 let log_path = self.path.join("LOG");
882 let log_file = OpenOptions::new()
883 .create(true)
884 .append(true)
885 .open(log_path)?;
886 Ok(BufWriter::new(log_file))
887 }
888
889 pub fn commit_write_batch<K: StoreKey + Send + Sync>(
892 &self,
893 mut write_batch: WriteBatch<'_, K, S, FAMILIES>,
894 ) -> Result<CommitStats> {
895 if self.read_only {
896 unreachable!("It's not possible to create a write batch for a read-only database");
897 }
898 let FinishResult {
899 sequence_number,
900 new_meta_files,
901 new_sst_files,
902 new_blob_files,
903 keys_written,
904 } = write_batch.finish(|family| {
905 let inner = self.inner.read();
906 let set = &inner.accessed_key_hashes[family as usize];
907 let initial_capacity = set.len() * 20 / 19;
910 let mut amqf =
916 qfilter::Filter::with_fingerprint_size(initial_capacity as u64, u64::BITS as u8)
917 .unwrap();
918 set.retain(|hash| {
921 amqf.insert_fingerprint(false, *hash)
925 .expect("Failed to insert fingerprint");
926 false
927 });
928 amqf
929 })?;
930 let stats = self.commit(CommitOptions {
931 new_meta_files,
932 new_sst_files,
933 new_blob_files,
934 sst_files_to_delete: vec![],
935 blob_seq_numbers_to_delete: vec![],
936 sequence_number,
937 keys_written,
938 })?;
939 write_batch.mark_succeeded();
941 Ok(stats)
942 }
943
944 fn commit(
947 &self,
948 CommitOptions {
949 mut new_meta_files,
950 new_sst_files,
951 new_blob_files,
952 sst_files_to_delete,
953 mut blob_seq_numbers_to_delete,
954 sequence_number: mut seq,
955 keys_written,
956 }: CommitOptions,
957 ) -> Result<CommitStats, anyhow::Error> {
958 let time = Timestamp::now();
959
960 new_meta_files.sort_unstable_by_key(|f| f.seq);
961
962 let mut stats = CommitStats::default();
963
964 let sync_span = tracing::trace_span!("sync new files").entered();
965
966 enum SyncItem {
967 Meta(u32, File),
968 Sst(File),
969 Blob(u32, File),
970 }
971 enum SyncResult {
972 Meta(MetaFile),
973 Sst,
974 Blob(u32, File),
975 }
976
977 let mut sync_items: Vec<SyncItem> =
978 Vec::with_capacity(new_meta_files.len() + new_sst_files.len() + new_blob_files.len());
979 for NewFile { seq, file, size } in new_meta_files {
980 stats.bytes_written += size;
981 sync_items.push(SyncItem::Meta(seq, file));
982 }
983 for NewFile { file, size, .. } in new_sst_files {
984 stats.bytes_written += size;
985 sync_items.push(SyncItem::Sst(file));
986 }
987 for NewFile { seq, file, size } in new_blob_files {
988 stats.bytes_written += size;
989 sync_items.push(SyncItem::Blob(seq, file));
990 }
991
992 let results: Vec<SyncResult> = self
993 .parallel_scheduler
994 .parallel_map_collect_owned::<_, _, Result<Vec<_>>>(sync_items, |item| match item {
995 SyncItem::Meta(seq, file) => {
996 file.sync_data()?;
997 let meta_file = MetaFile::open(
998 &self.path,
999 seq,
1000 Some(&self.config.family_configs),
1001 self.config.access_mode,
1002 )?;
1003 Ok(SyncResult::Meta(meta_file))
1004 }
1005 SyncItem::Sst(file) => {
1006 file.sync_data()?;
1007 Ok(SyncResult::Sst)
1008 }
1009 SyncItem::Blob(seq, file) => {
1010 file.sync_data()?;
1011 Ok(SyncResult::Blob(seq, file))
1012 }
1013 })?;
1014
1015 let mut new_meta_files: Vec<MetaFile> = Vec::new();
1016 let mut new_blob_files: Vec<(u32, File)> = Vec::new();
1017 for result in results {
1018 match result {
1019 SyncResult::Meta(mf) => new_meta_files.push(mf),
1020 SyncResult::Sst => {}
1021 SyncResult::Blob(seq, file) => new_blob_files.push((seq, file)),
1022 }
1023 }
1024
1025 let mut sst_filter = SstFilter::new();
1026 for meta_file in new_meta_files.iter_mut().rev() {
1027 sst_filter.apply_filter(meta_file);
1028 }
1029
1030 drop(sync_span);
1035
1036 let new_meta_info = new_meta_files
1037 .iter()
1038 .map(|meta| {
1039 let ssts = meta
1040 .entries()
1041 .iter()
1042 .zip(meta.hash_ranges())
1043 .map(|(entry, range)| {
1044 let seq = entry.sequence_number();
1045 let size = entry.size();
1046 let flags = entry.flags();
1047 (seq, range.min_hash, range.max_hash, size, flags)
1048 })
1049 .collect::<Vec<_>>();
1050 (
1051 meta.sequence_number(),
1052 meta.family(),
1053 ssts,
1054 meta.obsolete_sst_files().to_vec(),
1055 )
1056 })
1057 .collect::<Vec<_>>();
1058
1059 let has_delete_file;
1068 let mut meta_seq_numbers_to_delete = [(); FAMILIES].map(|_| Vec::new());
1069 let mut entries_to_remove = [(); FAMILIES].map(|_| Vec::new());
1070 stats.bytes_deleted += sst_files_to_delete.iter().map(|f| f.size).sum::<u64>();
1073 let mut sst_seq_numbers_to_delete = sst_files_to_delete
1075 .iter()
1076 .map(|f| f.seq)
1077 .collect::<Vec<_>>();
1078
1079 {
1080 let inner = self.inner.read();
1081
1082 for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() {
1086 entries_to_remove[family].extend(
1087 meta_files
1088 .iter()
1089 .rev()
1090 .map(|meta_file| sst_filter.apply_filter_collect(meta_file)),
1091 );
1092 }
1093
1094 for meta_file in new_meta_files.iter().rev() {
1101 let should_remove = sst_filter.apply_and_get_remove(meta_file);
1102 debug_assert!(
1103 !should_remove,
1104 "newly created meta file should never be a candidate for removal"
1105 );
1106 }
1107 for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() {
1108 for i in (0..meta_files.len()).rev() {
1109 let to_remove = &entries_to_remove[family][meta_files.len() - 1 - i];
1111 if sst_filter.apply_and_get_remove_after_removing(&meta_files[i], to_remove) {
1112 meta_seq_numbers_to_delete[family].push(meta_files[i].sequence_number());
1113 stats.bytes_deleted += meta_files[i].byte_size();
1115 }
1116 }
1117 }
1118
1119 has_delete_file = !sst_files_to_delete.is_empty()
1123 || !blob_seq_numbers_to_delete.is_empty()
1124 || meta_seq_numbers_to_delete
1125 .iter()
1126 .any(|seqs| !seqs.is_empty());
1127 }
1128
1129 stats.bytes_deleted += blob_seq_numbers_to_delete
1135 .iter()
1136 .map(|seq| {
1137 fs::metadata(self.path.join(format!("{seq:08}.blob")))
1138 .map(|m| m.len())
1139 .unwrap_or(0)
1140 })
1141 .sum::<u64>();
1142
1143 if has_delete_file {
1144 seq += 1;
1145 }
1146
1147 self.parallel_scheduler.block_in_place(|| {
1148 if has_delete_file {
1149 sst_seq_numbers_to_delete.sort_unstable();
1150 for seqs in &mut meta_seq_numbers_to_delete {
1151 seqs.sort_unstable();
1152 }
1153 blob_seq_numbers_to_delete.sort_unstable();
1154 let mut buf = Vec::with_capacity(
1156 (sst_seq_numbers_to_delete.len()
1157 + meta_seq_numbers_to_delete
1158 .iter()
1159 .map(Vec::len)
1160 .sum::<usize>()
1161 + blob_seq_numbers_to_delete.len())
1162 * size_of::<u32>(),
1163 );
1164 for seq in sst_seq_numbers_to_delete.iter() {
1165 buf.write_u32::<BE>(*seq)?;
1166 }
1167 for seq in meta_seq_numbers_to_delete.iter().flatten() {
1168 buf.write_u32::<BE>(*seq)?;
1169 }
1170 for seq in blob_seq_numbers_to_delete.iter() {
1171 buf.write_u32::<BE>(*seq)?;
1172 }
1173 let del_path = self.path.join(format!("{seq:08}.del"));
1174 let mut file = File::create(&del_path)?;
1175 file.write_all(&buf)?;
1176 file.sync_data()?;
1177 }
1178
1179 commit_current(&self.path, seq).context("Committing CURRENT file failed")?;
1180
1181 if let Err(e) = (|| {
1199 let mut log = self.open_log()?;
1200 writeln!(log, "Time {time}")?;
1201 let span = time.until(Timestamp::now())?;
1202 writeln!(log, "Commit {seq:08} {keys_written} keys in {span:#}")?;
1203 writeln!(log, "FAM | META SEQ | SST SEQ | RANGE")?;
1204 for (meta_seq, family, ssts, obsolete) in new_meta_info {
1205 for (seq, min, max, size, flags) in ssts {
1206 writeln!(
1207 log,
1208 "{family:3} | {meta_seq:08} | {seq:08} SST | {} ({} MiB, {})",
1209 range_to_str(min, max),
1210 size / 1024 / 1024,
1211 flags
1212 )?;
1213 }
1214 for obsolete in obsolete.chunks(15) {
1215 write!(log, "{family:3} | {meta_seq:08} |")?;
1216 for seq in obsolete {
1217 write!(log, " {seq:08}")?;
1218 }
1219 writeln!(log, " OBSOLETE SST")?;
1220 }
1221 }
1222
1223 fn write_seq_numbers<W: std::io::Write, T>(
1224 log: &mut W,
1225 items: &[T],
1226 label: &str,
1227 extract_seq: fn(&T) -> u32,
1228 ) -> std::io::Result<()> {
1229 for chunk in items.chunks(15) {
1230 write!(log, " | |")?;
1231 for item in chunk {
1232 write!(log, " {:08}", extract_seq(item))?;
1233 }
1234 writeln!(log, " {}", label)?;
1235 }
1236 Ok(())
1237 }
1238
1239 new_blob_files.sort_unstable_by_key(|(seq, _)| *seq);
1240 write_seq_numbers(&mut log, &new_blob_files, "NEW BLOB", |&(seq, _)| seq)?;
1241 write_seq_numbers(
1242 &mut log,
1243 &blob_seq_numbers_to_delete,
1244 "BLOB DELETED",
1245 |&seq| seq,
1246 )?;
1247 write_seq_numbers(
1248 &mut log,
1249 &sst_seq_numbers_to_delete,
1250 "SST DELETED",
1251 |&seq| seq,
1252 )?;
1253 for seqs in &meta_seq_numbers_to_delete {
1254 write_seq_numbers(&mut log, seqs, "META DELETED", |&seq| seq)?;
1255 }
1256 anyhow::Ok(())
1257 })() {
1258 eprintln!("turbo-persistence: failed to write LOG after commit {seq:08}: {e:#}");
1259 }
1260
1261 anyhow::Ok(())
1262 })?;
1263
1264 {
1270 let mut inner = self.inner.write();
1271
1272 for (meta_files, family_removals) in
1274 inner.meta_files_by_family.iter_mut().zip(entries_to_remove)
1275 {
1276 for (meta_file, to_remove) in
1277 meta_files.iter_mut().zip(family_removals.into_iter().rev())
1278 {
1279 if !to_remove.is_empty() {
1280 meta_file.retain_entries(|seq| !to_remove.contains(&seq));
1281 }
1282 }
1283 }
1284
1285 for meta_file in new_meta_files.drain(..) {
1286 inner.push_meta_file(meta_file);
1287 }
1288 for (meta_files, seqs_to_delete) in inner
1289 .meta_files_by_family
1290 .iter_mut()
1291 .zip(&meta_seq_numbers_to_delete)
1292 {
1293 if !seqs_to_delete.is_empty() {
1294 let to_delete: HashSet<u32> = seqs_to_delete.iter().copied().collect();
1295 meta_files.retain(|meta| !to_delete.contains(&meta.sequence_number()));
1296 }
1297 }
1298 #[cfg(debug_assertions)]
1299 inner.debug_assert_meta_invariants();
1300 inner.current_sequence_number = seq;
1301 self.is_empty.store(inner.is_empty(), Ordering::Relaxed);
1302 drop(inner);
1304 }
1305
1306 self.deferred_deletions.lock().extend(
1311 Self::try_delete_files(&self.path, &sst_seq_numbers_to_delete, "sst")
1312 .map(DeferredDeletion::Sst)
1313 .chain(
1314 meta_seq_numbers_to_delete
1315 .iter()
1316 .flat_map(|seqs| Self::try_delete_files(&self.path, seqs, "meta"))
1317 .map(DeferredDeletion::Meta),
1318 )
1319 .chain(
1320 Self::try_delete_files(&self.path, &blob_seq_numbers_to_delete, "blob")
1321 .map(DeferredDeletion::Blob),
1322 ),
1323 );
1324
1325 self.retry_deferred_deletions();
1327
1328 #[cfg(feature = "verbose_log")]
1330 {
1331 let _: Result<(), _> = (|| -> anyhow::Result<()> {
1332 let mut log = self.open_log()?;
1333 writeln!(log, "New database state:")?;
1334 writeln!(log, "FAM | META SEQ | SST SEQ FLAGS | RANGE")?;
1335 let inner = self.inner.read();
1336 for (family, meta_files) in inner.meta_files_by_family.iter().enumerate() {
1337 for meta in meta_files {
1338 let meta_seq = meta.sequence_number();
1339 for (entry, range) in meta.entries().iter().zip(meta.hash_ranges()) {
1340 let seq = entry.sequence_number();
1341 writeln!(
1342 log,
1343 "{family:3} | {meta_seq:08} | {seq:08} {:>6} | {}",
1344 entry.flags(),
1345 range_to_str(range.min_hash, range.max_hash)
1346 )?;
1347 }
1348 }
1349 }
1350 Ok(())
1351 })();
1352 }
1353
1354 Ok(stats)
1355 }
1356
1357 pub fn full_compact(&self) -> Result<()> {
1360 self.compact(&CompactConfig {
1361 min_merge_count: 2,
1362 optimal_merge_count: usize::MAX,
1363 max_merge_count: usize::MAX,
1364 max_merge_bytes: u64::MAX,
1365 min_merge_duplication_bytes: 0,
1366 optimal_merge_duplication_bytes: u64::MAX,
1367 max_merge_segment_count: usize::MAX,
1368 })?;
1369 Ok(())
1370 }
1371
1372 pub fn compact(&self, compact_config: &CompactConfig) -> Result<Option<CommitStats>> {
1380 let mut guard = self.acquire_write_operation("compaction")?;
1381
1382 self.clear_cache();
1387
1388 let mut sequence_number;
1389 let mut new_meta_files = Vec::new();
1390 let mut new_sst_files = Vec::new();
1391 let mut sst_files_to_delete = Vec::new();
1392 let mut blob_seq_numbers_to_delete = Vec::new();
1393 let mut keys_written = 0;
1394
1395 {
1396 let inner = self.inner.read();
1397 sequence_number = AtomicU32::new(inner.current_sequence_number);
1398 self.compact_internal(
1399 &inner.meta_files_by_family,
1400 &sequence_number,
1401 &mut new_meta_files,
1402 &mut new_sst_files,
1403 &mut sst_files_to_delete,
1404 &mut blob_seq_numbers_to_delete,
1405 &mut keys_written,
1406 compact_config,
1407 )
1408 .context("Failed to compact database")?;
1409 }
1410
1411 let has_changes = !new_meta_files.is_empty();
1412 let stats = if has_changes {
1413 let stats = self
1414 .commit(CommitOptions {
1415 new_meta_files,
1416 new_sst_files,
1417 new_blob_files: Vec::new(),
1418 sst_files_to_delete,
1419 blob_seq_numbers_to_delete,
1420 sequence_number: *sequence_number.get_mut(),
1421 keys_written,
1422 })
1423 .context("Failed to commit the database compaction")?;
1424 Some(stats)
1425 } else {
1426 None
1427 };
1428
1429 guard.success();
1430 Ok(stats)
1431 }
1432
1433 fn compact_internal(
1435 &self,
1436 meta_files_by_family: &[Vec<MetaFile>; FAMILIES],
1437 sequence_number: &AtomicU32,
1438 new_meta_files: &mut Vec<NewFile>,
1439 new_sst_files: &mut Vec<NewFile>,
1440 sst_files_to_delete: &mut Vec<DeletedFile>,
1441 blob_seq_numbers_to_delete: &mut Vec<u32>,
1442 keys_written: &mut u64,
1443 compact_config: &CompactConfig,
1444 ) -> Result<()> {
1445 if meta_files_by_family.iter().all(Vec::is_empty) {
1446 return Ok(());
1447 }
1448
1449 struct SstWithRange {
1450 meta_index: usize,
1452 index_in_meta: u32,
1453 seq: u32,
1454 range: StaticSortedFileRange,
1455 size: u64,
1456 flags: MetaEntryFlags,
1457 }
1458
1459 impl Compactable for SstWithRange {
1460 fn range(&self) -> RangeInclusive<u64> {
1461 self.range.min_hash..=self.range.max_hash
1462 }
1463
1464 fn size(&self) -> u64 {
1465 self.size
1466 }
1467
1468 fn category(&self) -> u8 {
1469 if self.flags.cold() { 1 } else { 0 }
1472 }
1473 }
1474
1475 let sst_by_family = meta_files_by_family
1476 .iter()
1477 .enumerate()
1478 .map(|(family, meta_files)| {
1479 meta_files
1480 .iter()
1481 .enumerate()
1482 .flat_map(|(meta_index, meta)| {
1483 debug_assert_eq!(
1484 meta.family() as usize,
1485 family,
1486 "meta file stored in the wrong family shard during compaction"
1487 );
1488 meta.entries()
1489 .iter()
1490 .enumerate()
1491 .map(move |(index_in_meta, entry)| SstWithRange {
1492 meta_index,
1493 index_in_meta: index_in_meta as u32,
1494 seq: entry.sequence_number(),
1495 range: meta.range(index_in_meta as u32),
1496 size: entry.size(),
1497 flags: entry.flags(),
1498 })
1499 })
1500 .collect::<Vec<_>>()
1501 })
1502 .collect::<Vec<_>>();
1503
1504 let path = &self.path;
1505
1506 let log_mutex = Mutex::new(());
1507
1508 struct PartialResultPerFamily {
1509 new_meta_file: Option<NewFile>,
1510 new_sst_files: Vec<NewFile>,
1511 sst_files_to_delete: Vec<DeletedFile>,
1512 blob_seq_numbers_to_delete: Vec<u32>,
1513 keys_written: u64,
1514 }
1515
1516 let mut compact_config = compact_config.clone();
1517 let merge_jobs = sst_by_family
1518 .into_iter()
1519 .enumerate()
1520 .filter_map(|(family, ssts_with_ranges)| {
1521 if compact_config.max_merge_segment_count == 0 {
1522 return None;
1523 }
1524 let (merge_jobs, real_merge_job_size) =
1525 get_merge_segments(&ssts_with_ranges, &compact_config);
1526 compact_config.max_merge_segment_count -= real_merge_job_size;
1527 Some((family, ssts_with_ranges, merge_jobs))
1528 })
1529 .collect::<Vec<_>>();
1530
1531 let result = self
1532 .parallel_scheduler
1533 .parallel_map_collect_owned::<_, _, Result<Vec<_>>>(
1534 merge_jobs,
1535 |(family, ssts_with_ranges, merge_jobs)| {
1536 let meta_files = &meta_files_by_family[family];
1537 let family = family as u32;
1538 debug_assert!(
1539 meta_files.iter().all(|meta| meta.family() == family),
1540 "compaction received a meta file from the wrong family shard"
1541 );
1542
1543 if merge_jobs.is_empty() {
1544 return Ok(PartialResultPerFamily {
1545 new_meta_file: None,
1546 new_sst_files: Vec::new(),
1547 sst_files_to_delete: Vec::new(),
1548 blob_seq_numbers_to_delete: Vec::new(),
1549 keys_written: 0,
1550 });
1551 }
1552
1553 let used_key_hashes: Option<qfilter::Filter> = {
1558 let filters: Vec<qfilter::FilterRef<'_>> = meta_files
1559 .iter()
1560 .filter_map(|meta_file| {
1561 meta_file.deserialize_used_key_hashes_amqf().transpose()
1562 })
1563 .collect::<Result<Vec<_>>>()?
1564 .into_iter()
1565 .filter(|amqf| !amqf.is_empty())
1566 .collect();
1567 if filters.is_empty() {
1568 None
1569 } else if filters.len() == 1 {
1570 Some(filters[0].to_owned())
1572 } else {
1573 let total_len: u64 = filters.iter().map(|f| f.len()).sum();
1574 let mut merged =
1577 qfilter::Filter::with_fingerprint_size(total_len, u64::BITS as u8)
1578 .expect("Failed to create merged AMQF filter");
1579 for filter in &filters {
1580 merged
1581 .merge(false, filter)
1582 .expect("Failed to merge AMQF filters");
1583 }
1584 merged.shrink_to_fit();
1585 Some(merged)
1586 }
1587 };
1588
1589 let sst_files_to_delete = merge_jobs
1592 .iter()
1593 .filter(|l| l.len() > 1)
1594 .flat_map(|l| l.iter().copied())
1595 .map(|index| DeletedFile {
1596 seq: ssts_with_ranges[index].seq,
1597 size: ssts_with_ranges[index].size,
1598 })
1599 .collect::<Vec<_>>();
1600
1601 let span = tracing::trace_span!(
1603 "merge files",
1604 family = self.config.family_configs[family as usize].name
1605 );
1606 enum PartialMergeResult<'l> {
1607 Merged {
1608 new_sst_files: Vec<(u32, File, StaticSortedFileBuilderMeta<'static>)>,
1609 blob_seq_numbers_to_delete: Vec<u32>,
1610 keys_written: u64,
1611 indices: SmallVec<[usize; 1]>,
1612 },
1613 Move {
1614 seq: u32,
1615 meta: StaticSortedFileBuilderMeta<'l>,
1616 },
1617 }
1618 let merge_result = self
1619 .parallel_scheduler
1620 .parallel_map_collect_owned::<_, _, Result<Vec<_>>>(merge_jobs, |indices| {
1621 let _span = span.clone().entered();
1622
1623 if indices.len() == 1 {
1624 let index = indices[0];
1626 let meta_index = ssts_with_ranges[index].meta_index;
1627 let index_in_meta = ssts_with_ranges[index].index_in_meta;
1628 let meta_file = &meta_files[meta_index];
1629 let entry = meta_file.entry(index_in_meta);
1630 let amqf = Cow::Borrowed(entry.raw_amqf(meta_file.amqf_data()));
1631 let hash_range = meta_file.hash_range(index_in_meta);
1632 let meta = StaticSortedFileBuilderMeta {
1633 min_hash: hash_range.min_hash,
1634 max_hash: hash_range.max_hash,
1635 amqf,
1636 block_count: entry.block_count(),
1637 size: entry.size(),
1638 flags: entry.flags(),
1639 entries: 0,
1640 };
1641 return Ok(PartialMergeResult::Move {
1642 seq: entry.sequence_number(),
1643 meta,
1644 });
1645 }
1646
1647 let tombstone_is_dead = {
1652 let oldest_index_in_job = indices
1660 .iter()
1661 .copied()
1662 .min()
1663 .expect("merge jobs are not empty");
1664 let older_filters = ssts_with_ranges[..oldest_index_in_job]
1665 .iter()
1666 .map(|sst| {
1667 let meta_file = &meta_files[sst.meta_index];
1668 let entry = meta_file.entry(sst.index_in_meta);
1669 let range = meta_file.hash_range(sst.index_in_meta);
1670 (range.min_hash, range.max_hash, entry.amqf())
1671 })
1672 .collect::<Vec<_>>();
1673 move |hash: u64| {
1674 !older_filters.iter().any(|(min, max, amqf)| {
1675 hash >= *min
1676 && hash <= *max
1677 && amqf.contains_fingerprint(hash)
1678 })
1679 }
1680 };
1681 let iters = indices
1685 .iter()
1686 .map(|&index| {
1687 let meta_index = ssts_with_ranges[index].meta_index;
1688 let index_in_meta = ssts_with_ranges[index].index_in_meta;
1689 let meta_file = &meta_files[meta_index];
1690 let entry = meta_file.entry(index_in_meta);
1691 StaticSortedFileIter::open(
1692 path,
1693 entry.sst_metadata(),
1694 meta_file.compression(),
1695 self.config.access_mode,
1696 )
1697 })
1698 .collect::<Result<Vec<_>>>()?;
1699
1700 let iter = MergeIter::new(iters.into_iter())?;
1701
1702 let mut blob_seq_numbers_to_delete: Vec<u32> = Vec::new();
1703
1704 struct Collector {
1705 writer: Option<(u32, StreamingSstWriter<LookupEntry>)>,
1712 flags: MetaEntryFlags,
1713 compression: Compression,
1714 new_sst_files:
1715 Vec<(u32, File, StaticSortedFileBuilderMeta<'static>)>,
1716 last_hash: Option<u64>,
1719 }
1720 impl Collector {
1721 fn new(flags: MetaEntryFlags, compression: Compression) -> Self {
1722 Self {
1723 writer: None,
1724 flags,
1725 compression,
1726 new_sst_files: Vec::new(),
1727 last_hash: None,
1728 }
1729 }
1730
1731 fn ensure_writer(
1733 &mut self,
1734 path: &Path,
1735 sequence_number: &AtomicU32,
1736 ) -> Result<&mut StreamingSstWriter<LookupEntry>>
1737 {
1738 if self.writer.is_none() {
1739 let seq =
1740 sequence_number.fetch_add(1, Ordering::SeqCst) + 1;
1741 let sst_path = path.join(format!("{seq:08}.sst"));
1742 let writer = StreamingSstWriter::new(
1743 &sst_path,
1744 self.flags,
1745 MAX_ENTRIES_PER_COMPACTED_FILE as u64,
1746 self.compression,
1747 )?;
1748 self.writer = Some((seq, writer));
1749 }
1750 Ok(&mut self.writer.as_mut().unwrap().1)
1751 }
1752
1753 fn close_sst_file(&mut self, keys_written: &mut u64) -> Result<()> {
1757 if let Some((seq, writer)) = self.writer.take() {
1758 let _span =
1759 tracing::trace_span!("close merged sst file").entered();
1760 let (meta, file) = writer.close()?;
1761 *keys_written += meta.entries;
1762 self.new_sst_files.push((seq, file, meta));
1763 }
1764 Ok(())
1765 }
1766
1767 fn cancel(&mut self) {
1769 if let Some((_, writer)) = self.writer.take() {
1770 writer.cancel();
1771 }
1772 }
1773
1774 fn add_entry(
1778 &mut self,
1779 entry: LookupEntry,
1780 path: &Path,
1781 sequence_number: &AtomicU32,
1782 keys_written: &mut u64,
1783 ) -> Result<()> {
1784 let key_changed = self.last_hash != Some(entry.hash);
1785 if key_changed
1788 && let Some((_, ref writer)) = self.writer
1789 && writer.is_full(
1790 MAX_ENTRIES_PER_COMPACTED_FILE,
1791 DATA_THRESHOLD_PER_COMPACTED_FILE,
1792 )
1793 {
1794 self.close_sst_file(keys_written)?;
1795 }
1796 self.last_hash = Some(entry.hash);
1797 let writer = self.ensure_writer(path, sequence_number)?;
1798 if let Err(err) = writer.add(entry) {
1799 self.cancel();
1800 return Err(err);
1801 }
1802 Ok(())
1803 }
1804 }
1805 #[cfg(debug_assertions)]
1806 impl Drop for Collector {
1807 fn drop(&mut self) {
1808 if !std::thread::panicking() {
1809 assert!(
1810 self.writer.is_none(),
1811 "Collector dropped with an open writer"
1812 );
1813 }
1814 }
1815 }
1816 let compression =
1817 self.config.family_configs[family as usize].compression;
1818 let mut used_collector =
1819 Collector::new(MetaEntryFlags::WARM, compression);
1820 let mut unused_collector =
1821 Collector::new(MetaEntryFlags::COLD, compression);
1822 let mut current_key: Option<RcBytes> = None;
1823 let mut keys_written = 0;
1824
1825 let mut skip_remaining_for_this_key = false;
1832 let mut deleted_values_for_this_key: AutoSet<
1835 RcBytes,
1836 BuildHasherDefault<FxHasher>,
1837 1,
1838 > = AutoSet::default();
1839 let family_config = &self.config.family_configs[family as usize];
1840
1841 let result: Result<_> = (|| {
1842 for entry in iter {
1843 let entry = entry?;
1844 if current_key.as_ref() != Some(&entry.key) {
1845 skip_remaining_for_this_key = false;
1847 deleted_values_for_this_key.clear();
1848 current_key = Some(entry.key.clone());
1849 }
1850 if let IterValue::KeyValueDeleted { value } = &entry.value {
1854 deleted_values_for_this_key.insert(value.clone());
1855 if tombstone_is_dead(entry.hash) {
1859 continue;
1860 }
1861 } else if !deleted_values_for_this_key.is_empty()
1862 && let IterValue::Slice { value } = &entry.value
1864 && deleted_values_for_this_key.contains(value)
1865 {
1866 continue;
1869 }
1870 if !skip_remaining_for_this_key {
1871 let is_used =
1872 used_key_hashes.as_ref().is_some_and(|amqf| {
1873 amqf.contains_fingerprint(entry.hash)
1874 });
1875 let collector = if is_used {
1876 &mut used_collector
1877 } else {
1878 &mut unused_collector
1879 };
1880 match family_config.kind {
1881 FamilyKind::MultiValue => {
1882 if matches!(entry.value, IterValue::KeyDeleted) {
1886 skip_remaining_for_this_key = true;
1887 }
1888 }
1889 FamilyKind::SingleValue => {
1890 skip_remaining_for_this_key = true;
1893 }
1894 }
1895 if matches!(entry.value, IterValue::KeyDeleted)
1898 && tombstone_is_dead(entry.hash)
1899 {
1900 continue;
1901 }
1902 collector.add_entry(
1903 entry,
1904 path,
1905 sequence_number,
1906 &mut keys_written,
1907 )?;
1908 } else {
1909 if let IterValue::Blob { sequence_number } = &entry.value {
1913 blob_seq_numbers_to_delete.push(*sequence_number);
1914 }
1915 }
1916 }
1917
1918 used_collector.close_sst_file(&mut keys_written)?;
1920 unused_collector.close_sst_file(&mut keys_written)?;
1921
1922 let mut new_sst_files = take(&mut unused_collector.new_sst_files);
1923 new_sst_files.append(&mut used_collector.new_sst_files);
1924 Ok(PartialMergeResult::Merged {
1925 new_sst_files,
1926 blob_seq_numbers_to_delete,
1927 keys_written,
1928 indices,
1929 })
1930 })();
1931 if result.is_err() {
1932 used_collector.cancel();
1933 unused_collector.cancel();
1934 }
1935 result
1936 })
1937 .with_context(|| {
1938 format!("Failed to merge database files for family {family}")
1939 })?;
1940
1941 let Some((sst_files_len, blob_delete_len)) = merge_result
1942 .iter()
1943 .map(|r| {
1944 if let PartialMergeResult::Merged {
1945 new_sst_files,
1946 blob_seq_numbers_to_delete,
1947 indices: _,
1948 keys_written: _,
1949 } = r
1950 {
1951 (new_sst_files.len(), blob_seq_numbers_to_delete.len())
1952 } else {
1953 (0, 0)
1954 }
1955 })
1956 .reduce(|(a1, a2), (b1, b2)| (a1 + b1, a2 + b2))
1957 else {
1958 unreachable!()
1959 };
1960
1961 let mut new_sst_files = Vec::with_capacity(sst_files_len);
1962 let mut blob_seq_numbers_to_delete = Vec::with_capacity(blob_delete_len);
1963
1964 let meta_seq = sequence_number.fetch_add(1, Ordering::SeqCst) + 1;
1965 let mut meta_file_builder = MetaFileBuilder::new(
1966 family,
1967 self.config.family_configs[family as usize].compression,
1968 );
1969
1970 let mut keys_written = 0;
1971 self.parallel_scheduler.block_in_place(|| {
1972 let guard = log_mutex.lock();
1973 let mut log = self.open_log()?;
1974 writeln!(log, "{family:3} | {meta_seq:08} | Compaction:",)?;
1975
1976 for result in merge_result {
1977 match result {
1978 PartialMergeResult::Merged {
1979 new_sst_files: merged_new_sst_files,
1980 blob_seq_numbers_to_delete: merged_blob_seq_numbers_to_delete,
1981 keys_written: merged_keys_written,
1982 indices,
1983 } => {
1984 writeln!(
1985 log,
1986 "{family:3} | {meta_seq:08} | MERGE \
1987 ({merged_keys_written} keys):"
1988 )?;
1989 for i in indices.iter() {
1990 let seq = ssts_with_ranges[*i].seq;
1991 let (min, max) = ssts_with_ranges[*i].range().into_inner();
1992 writeln!(
1993 log,
1994 "{family:3} | {meta_seq:08} | {seq:08} INPUT | {}",
1995 range_to_str(min, max)
1996 )?;
1997 }
1998 for (seq, file, meta) in merged_new_sst_files {
1999 let min = meta.min_hash;
2000 let max = meta.max_hash;
2001 writeln!(
2002 log,
2003 "{family:3} | {meta_seq:08} | {seq:08} OUTPUT | {} \
2004 ({})",
2005 range_to_str(min, max),
2006 meta.flags
2007 )?;
2008
2009 let size = meta.size;
2010 meta_file_builder.add(seq, meta);
2011 new_sst_files.push(NewFile { seq, file, size });
2012 }
2013 blob_seq_numbers_to_delete
2014 .extend(merged_blob_seq_numbers_to_delete);
2015 keys_written += merged_keys_written;
2016 }
2017 PartialMergeResult::Move { seq, meta } => {
2018 let min = meta.min_hash;
2019 let max = meta.max_hash;
2020 writeln!(
2021 log,
2022 "{family:3} | {meta_seq:08} | {seq:08} MOVED | {}",
2023 range_to_str(min, max)
2024 )?;
2025
2026 meta_file_builder.add(seq, meta);
2027 }
2028 }
2029 }
2030 drop(log);
2031 drop(guard);
2032
2033 anyhow::Ok(())
2034 })?;
2035
2036 for f in sst_files_to_delete.iter() {
2037 meta_file_builder.add_obsolete_sst_file(f.seq);
2038 }
2039 let new_meta_file = {
2044 let _span = tracing::trace_span!("write meta file").entered();
2045 let (file, size) = self
2046 .parallel_scheduler
2047 .block_in_place(|| meta_file_builder.write(&self.path, meta_seq))?;
2048 NewFile {
2049 seq: meta_seq,
2050 file,
2051 size,
2052 }
2053 };
2054
2055 Ok(PartialResultPerFamily {
2056 new_meta_file: Some(new_meta_file),
2057 new_sst_files,
2058 sst_files_to_delete,
2059 blob_seq_numbers_to_delete,
2060 keys_written,
2061 })
2062 },
2063 )?;
2064
2065 for PartialResultPerFamily {
2066 new_meta_file: inner_new_meta_file,
2067 new_sst_files: mut inner_new_sst_files,
2068 sst_files_to_delete: mut inner_sst_files_to_delete,
2069 blob_seq_numbers_to_delete: mut inner_blob_seq_numbers_to_delete,
2070 keys_written: inner_keys_written,
2071 } in result
2072 {
2073 new_meta_files.extend(inner_new_meta_file);
2074 new_sst_files.append(&mut inner_new_sst_files);
2075 sst_files_to_delete.append(&mut inner_sst_files_to_delete);
2076 blob_seq_numbers_to_delete.append(&mut inner_blob_seq_numbers_to_delete);
2077 *keys_written += inner_keys_written;
2078 }
2079
2080 Ok(())
2081 }
2082
2083 pub fn get<K: QueryKey>(&self, family: usize, key: &K) -> Result<Option<ArcBytes>> {
2086 debug_assert!(family < FAMILIES, "Family index out of bounds");
2087 if self.config.family_configs[family].kind != FamilyKind::SingleValue {
2088 panic!(
2090 "only single valued tables can be queried with `get', call `get_multiple` instead"
2091 )
2092 }
2093 let span = tracing::trace_span!(
2094 "database read",
2095 name = self.config.family_configs[family].name,
2096 result_size = tracing::field::Empty
2097 )
2098 .entered();
2099 let results = self.get_impl::<K, false>(family, key, &span)?;
2100 debug_assert!(results.len() <= 1, "get() should return at most one result");
2101 Ok(results.into_iter().next())
2102 }
2103
2104 pub fn get_multiple<K: QueryKey>(
2114 &self,
2115 family: usize,
2116 key: &K,
2117 ) -> Result<SmallVec<[ArcBytes; 1]>> {
2118 debug_assert!(family < FAMILIES, "Family index out of bounds");
2119 if self.config.family_configs[family].kind != FamilyKind::MultiValue {
2120 panic!("only multi-valued tables can be queried with `get_multiple`")
2122 }
2123 let span = tracing::trace_span!(
2124 "database read multiple",
2125 name = self.config.family_configs[family].name,
2126 result_count = tracing::field::Empty,
2127 result_size = tracing::field::Empty
2128 )
2129 .entered();
2130 let results = self.get_impl::<K, true>(family, key, &span)?;
2131 Ok(results)
2132 }
2133
2134 fn get_impl<K: QueryKey, const FIND_ALL: bool>(
2139 &self,
2140 family: usize,
2141 key: &K,
2142 span: &EnteredSpan,
2143 ) -> Result<SmallVec<[ArcBytes; 1]>> {
2144 let hash = hash_key(key);
2145 let inner = self.inner.read();
2146 let mut output: SmallVec<[ArcBytes; 1]> = SmallVec::new();
2147 #[cfg(feature = "stats")]
2150 let mut found_in_sst = false;
2151
2152 let mut deleted_values: AutoSet<ArcBytes, BuildHasherDefault<FxHasher>, 1> =
2156 AutoSet::default();
2157
2158 let mut size = 0;
2159
2160 let key_block_cache = self.key_block_cache();
2161 let value_block_cache = self.value_block_cache();
2162 debug_assert!(
2163 inner.meta_files_by_family[family]
2164 .iter()
2165 .all(|meta| meta.family() as usize == family),
2166 "meta file stored in the wrong family shard while querying family {family}"
2167 );
2168 for meta in inner.meta_files_by_family[family].iter().rev() {
2169 match meta.lookup::<K, FIND_ALL>(
2170 family as u32,
2171 hash,
2172 key,
2173 key_block_cache,
2174 value_block_cache,
2175 )? {
2176 MetaLookupResult::FamilyMiss => {
2177 #[cfg(feature = "stats")]
2178 self.stats.miss_family.fetch_add(1, Ordering::Relaxed);
2179 }
2180 MetaLookupResult::RangeMiss => {
2181 #[cfg(feature = "stats")]
2182 self.stats.miss_range.fetch_add(1, Ordering::Relaxed);
2183 }
2184 MetaLookupResult::QuickFilterMiss => {
2185 #[cfg(feature = "stats")]
2186 self.stats.miss_amqf.fetch_add(1, Ordering::Relaxed);
2187 }
2188 MetaLookupResult::SstLookup(result) => match result {
2189 SstLookupResult::Found(values) => {
2190 #[cfg(feature = "stats")]
2191 {
2192 found_in_sst = true;
2193 }
2194 inner.accessed_key_hashes[family].insert(hash);
2195 for value in values {
2196 match value {
2197 LookupValue::KeyDeleted => {
2198 #[cfg(feature = "stats")]
2199 self.stats.hits_deleted.fetch_add(1, Ordering::Relaxed);
2200 if !FIND_ALL {
2201 span.record("result_size", "deleted");
2202 return Ok(SmallVec::new());
2203 }
2204 if output.is_empty() {
2208 span.record("result_size", "deleted");
2209 } else {
2210 span.record("result_size", size);
2211 }
2212 return Ok(output);
2213 }
2214 LookupValue::KeyValueDeleted { value } => {
2215 #[cfg(feature = "stats")]
2216 self.stats.hits_deleted.fetch_add(1, Ordering::Relaxed);
2217 deleted_values.insert(value);
2220 }
2221 LookupValue::Slice { value } => {
2222 #[cfg(feature = "stats")]
2223 self.stats.hits_small.fetch_add(1, Ordering::Relaxed);
2224 if deleted_values.contains(&value) {
2225 continue;
2226 }
2227 if !FIND_ALL {
2228 span.record("result_size", value.len());
2229 return Ok(SmallVec::from_buf([value]));
2230 }
2231 size += value.len();
2232 output.push(value);
2233 }
2234 LookupValue::Blob { sequence_number } => {
2235 #[cfg(feature = "stats")]
2236 self.stats.hits_blob.fetch_add(1, Ordering::Relaxed);
2237 let blob = self.read_blob(
2238 sequence_number,
2239 self.config.family_configs[family].compression,
2240 )?;
2241 if deleted_values.iter().any(|d| **d == *blob) {
2242 continue;
2243 }
2244 if !FIND_ALL {
2245 span.record("result_size", blob.len());
2246 return Ok(SmallVec::from_buf([blob]));
2247 }
2248 size += blob.len();
2249 output.push(blob);
2250 }
2251 }
2252 }
2253 }
2254 SstLookupResult::NotFound => {
2255 #[cfg(feature = "stats")]
2256 self.stats.miss_key.fetch_add(1, Ordering::Relaxed);
2257 }
2258 },
2259 }
2260 }
2261
2262 #[cfg(feature = "stats")]
2263 if !found_in_sst {
2264 self.stats.miss_global.fetch_add(1, Ordering::Relaxed);
2265 }
2266
2267 if FIND_ALL {
2268 span.record("result_count", output.len());
2269 }
2270 if output.is_empty() {
2271 span.record("result_size", "not_found");
2272 } else {
2273 span.record("result_size", size);
2274 }
2275 Ok(output)
2276 }
2277
2278 pub fn batch_get<K: QueryKey>(
2279 &self,
2280 family: usize,
2281 keys: &[K],
2282 ) -> Result<Vec<Option<ArcBytes>>> {
2283 debug_assert!(family < FAMILIES, "Family index out of bounds");
2284 if self.config.family_configs[family].kind != FamilyKind::SingleValue {
2285 panic!("only single valued tables can be queried with `batch_get'")
2287 }
2288 let span = tracing::trace_span!(
2289 "database batch read",
2290 name = self.config.family_configs[family].name,
2291 keys = keys.len(),
2292 not_found = tracing::field::Empty,
2293 deleted = tracing::field::Empty,
2294 result_size = tracing::field::Empty
2295 )
2296 .entered();
2297 let mut cells: Vec<(u64, usize, Option<LookupValue>)> = Vec::with_capacity(keys.len());
2298 let mut empty_cells = keys.len();
2299 for (index, key) in keys.iter().enumerate() {
2300 let hash = hash_key(key);
2301 cells.push((hash, index, None));
2302 }
2303 cells.sort_by_key(|(hash, _, _)| *hash);
2304 let inner = self.inner.read();
2305 let key_block_cache = self.key_block_cache();
2306 let value_block_cache = self.value_block_cache();
2307 debug_assert!(
2308 inner.meta_files_by_family[family]
2309 .iter()
2310 .all(|meta| meta.family() as usize == family),
2311 "meta file stored in the wrong family shard while querying family {family}"
2312 );
2313 for meta in inner.meta_files_by_family[family].iter().rev() {
2314 let _result = meta.batch_lookup(
2315 family as u32,
2316 keys,
2317 &mut cells,
2318 &mut empty_cells,
2319 key_block_cache,
2320 value_block_cache,
2321 )?;
2322
2323 #[cfg(feature = "stats")]
2324 {
2325 let crate::meta_file::MetaBatchLookupResult {
2326 family_miss,
2327 range_misses,
2328 quick_filter_misses,
2329 sst_misses,
2330 hits: _,
2331 } = _result;
2332 if family_miss {
2333 self.stats.miss_family.fetch_add(1, Ordering::Relaxed);
2334 }
2335 if range_misses > 0 {
2336 self.stats
2337 .miss_range
2338 .fetch_add(range_misses as u64, Ordering::Relaxed);
2339 }
2340 if quick_filter_misses > 0 {
2341 self.stats
2342 .miss_amqf
2343 .fetch_add(quick_filter_misses as u64, Ordering::Relaxed);
2344 }
2345 if sst_misses > 0 {
2346 self.stats
2347 .miss_key
2348 .fetch_add(sst_misses as u64, Ordering::Relaxed);
2349 }
2350 }
2351
2352 if empty_cells == 0 {
2353 break;
2354 }
2355 }
2356 let mut deleted = 0;
2357 let mut not_found = 0;
2358 let mut result_size = 0;
2359 let mut results = vec![None; keys.len()];
2360 for (hash, index, result) in cells {
2361 if let Some(result) = result {
2362 inner.accessed_key_hashes[family].insert(hash);
2363 let result = match result {
2364 LookupValue::KeyDeleted => {
2365 #[cfg(feature = "stats")]
2366 self.stats.hits_deleted.fetch_add(1, Ordering::Relaxed);
2367 deleted += 1;
2368 None
2369 }
2370 LookupValue::KeyValueDeleted { .. } => {
2371 bail!(
2374 "unexpected key-value tombstone in SingleValue family {}",
2375 self.config.family_configs[family].name
2376 )
2377 }
2378 LookupValue::Slice { value } => {
2379 #[cfg(feature = "stats")]
2380 self.stats.hits_small.fetch_add(1, Ordering::Relaxed);
2381 result_size += value.len();
2382 Some(value)
2383 }
2384 LookupValue::Blob { sequence_number } => {
2385 #[cfg(feature = "stats")]
2386 self.stats.hits_blob.fetch_add(1, Ordering::Relaxed);
2387 let blob = self.read_blob(
2388 sequence_number,
2389 self.config.family_configs[family].compression,
2390 )?;
2391 result_size += blob.len();
2392 Some(blob)
2393 }
2394 };
2395 results[index] = result;
2396 } else {
2397 #[cfg(feature = "stats")]
2398 self.stats.miss_global.fetch_add(1, Ordering::Relaxed);
2399 not_found += 1;
2400 }
2401 }
2402 span.record("not_found", not_found);
2403 span.record("deleted", deleted);
2404 span.record("result_size", result_size);
2405 Ok(results)
2406 }
2407
2408 #[cfg(feature = "stats")]
2410 pub fn statistics(&self) -> Statistics {
2411 let inner = self.inner.read();
2412 Statistics {
2413 meta_files: inner.meta_files_by_family.iter().map(Vec::len).sum(),
2414 sst_files: inner
2415 .meta_files_by_family
2416 .iter()
2417 .flatten()
2418 .map(|meta| meta.entries().len())
2419 .sum(),
2420 key_block_cache: CacheStatistics::new(self.key_block_cache()),
2421 value_block_cache: CacheStatistics::new(self.value_block_cache()),
2422 hits: self.stats.hits_deleted.load(Ordering::Relaxed)
2423 + self.stats.hits_small.load(Ordering::Relaxed)
2424 + self.stats.hits_blob.load(Ordering::Relaxed),
2425 misses: self.stats.miss_global.load(Ordering::Relaxed),
2426 miss_family: self.stats.miss_family.load(Ordering::Relaxed),
2427 miss_range: self.stats.miss_range.load(Ordering::Relaxed),
2428 miss_amqf: self.stats.miss_amqf.load(Ordering::Relaxed),
2429 miss_key: self.stats.miss_key.load(Ordering::Relaxed),
2430 }
2431 }
2432
2433 pub fn meta_info(&self) -> Result<Vec<MetaFileInfo>> {
2434 Ok(self
2435 .inner
2436 .read()
2437 .meta_files_by_family
2438 .iter()
2439 .flat_map(|meta_files| meta_files.iter().rev())
2440 .map(|meta_file| {
2441 let entries = meta_file
2442 .entries()
2443 .iter()
2444 .zip(meta_file.hash_ranges())
2445 .map(|(entry, range)| {
2446 let amqf = entry.raw_amqf(meta_file.amqf_data());
2447 MetaFileEntryInfo {
2448 sequence_number: entry.sequence_number(),
2449 min_hash: range.min_hash,
2450 max_hash: range.max_hash,
2451 sst_size: entry.size(),
2452 flags: entry.flags(),
2453 amqf_size: entry.amqf_size(),
2454 amqf_entries: amqf.len(),
2455 block_count: entry.block_count(),
2456 }
2457 })
2458 .collect();
2459 MetaFileInfo {
2460 sequence_number: meta_file.sequence_number(),
2461 family: meta_file.family(),
2462 obsolete_sst_files: meta_file.obsolete_sst_files().to_vec(),
2463 entries,
2464 }
2465 })
2466 .collect())
2467 }
2468
2469 pub fn shutdown(&self) -> Result<()> {
2472 #[cfg(feature = "print_stats")]
2473 println!("{:#?}", self.statistics());
2474 self.retry_deferred_deletions();
2475 Ok(())
2476 }
2477
2478 fn try_delete_files<'a>(
2481 dir: &'a Path,
2482 seqs: &'a [u32],
2483 ext: &'a str,
2484 ) -> impl Iterator<Item = u32> + 'a {
2485 seqs.iter()
2486 .copied()
2487 .filter(move |&seq| fs::remove_file(dir.join(format!("{seq:08}.{ext}"))).is_err())
2488 }
2489
2490 fn retry_deferred_deletions(&self) {
2495 let mut deferred = self.deferred_deletions.lock();
2496 deferred.retain(|entry| {
2497 let (seq, ext) = match *entry {
2498 DeferredDeletion::Sst(seq) => (seq, "sst"),
2499 DeferredDeletion::Meta(seq) => (seq, "meta"),
2500 DeferredDeletion::Blob(seq) => (seq, "blob"),
2501 };
2502 fs::remove_file(self.path.join(format!("{seq:08}.{ext}"))).is_err()
2504 });
2505 }
2506}
2507
2508fn range_to_str(min: u64, max: u64) -> String {
2509 use std::fmt::Write;
2510 const DISPLAY_SIZE: usize = 100;
2511 const TOTAL_SIZE: u64 = u64::MAX;
2512 let start_pos = (min as u128 * DISPLAY_SIZE as u128 / TOTAL_SIZE as u128) as usize;
2513 let end_pos = (max as u128 * DISPLAY_SIZE as u128 / TOTAL_SIZE as u128) as usize;
2514 let mut range_str = String::new();
2515 for i in 0..DISPLAY_SIZE {
2516 if i == start_pos && i == end_pos {
2517 range_str.push('O');
2518 } else if i == start_pos {
2519 range_str.push('[');
2520 } else if i == end_pos {
2521 range_str.push(']');
2522 } else if i > start_pos && i < end_pos {
2523 range_str.push('=');
2524 } else {
2525 range_str.push(' ');
2526 }
2527 }
2528 write!(range_str, " | {min:016x}-{max:016x}").unwrap();
2529 range_str
2530}
2531
2532pub struct MetaFileInfo {
2533 pub sequence_number: u32,
2534 pub family: u32,
2535 pub obsolete_sst_files: Vec<u32>,
2536 pub entries: Vec<MetaFileEntryInfo>,
2537}
2538
2539pub struct MetaFileEntryInfo {
2540 pub sequence_number: u32,
2541 pub min_hash: u64,
2542 pub max_hash: u64,
2543 pub amqf_size: u32,
2544 pub amqf_entries: usize,
2545 pub sst_size: u64,
2546 pub flags: MetaEntryFlags,
2547 pub block_count: u16,
2548}