Skip to main content

lsm_tree/blob_tree/
mod.rs

1// Copyright (c) 2024-present, fjall-rs
2// This source code is licensed under both the Apache 2.0 and MIT License
3// (found in the LICENSE-* files in the repository)
4
5mod gc;
6pub mod handle;
7pub mod ingest;
8
9#[doc(hidden)]
10pub use gc::{FragmentationEntry, FragmentationMap};
11
12use crate::{
13    abstract_tree::{AbstractTree, RangeItem},
14    coding::Decode,
15    iter_guard::{IterGuard, IterGuardImpl},
16    table::Table,
17    tree::inner::MemtableId,
18    value::InternalValue,
19    version::Version,
20    vlog::{Accessor, BlobFile, BlobFileWriter},
21    Cache, Config, Memtable, SeqNo, TableId, TreeId, UserKey, UserValue,
22};
23use handle::BlobIndirection;
24use std::{
25    ops::RangeBounds,
26    path::{Path, PathBuf},
27    sync::{Arc, MutexGuard},
28};
29
30/// Iterator value guard
31pub struct Guard {
32    tree: crate::BlobTree,
33    version: Version,
34    kv: crate::Result<InternalValue>,
35}
36
37impl IterGuard for Guard {
38    fn into_inner_if(
39        self,
40        pred: impl Fn(&UserKey) -> bool,
41    ) -> crate::Result<(UserKey, Option<UserValue>)> {
42        let kv = self.kv?;
43
44        if pred(&kv.key.user_key) {
45            resolve_value_handle(
46                self.tree.id(),
47                self.tree.blobs_folder.as_path(),
48                &self.tree.index.config.cache,
49                &self.version,
50                kv,
51            )
52            .map(|(k, v)| (k, Some(v)))
53        } else {
54            Ok((kv.key.user_key, None))
55        }
56    }
57
58    fn key(self) -> crate::Result<UserKey> {
59        self.kv.map(|kv| kv.key.user_key)
60    }
61
62    fn size(self) -> crate::Result<u32> {
63        let kv = self.kv?;
64
65        if kv.key.value_type.is_indirection() {
66            let mut cursor = std::io::Cursor::new(kv.value);
67            Ok(BlobIndirection::decode_from(&mut cursor)?.size)
68        } else {
69            #[expect(clippy::cast_possible_truncation, reason = "values are u32 max length")]
70            Ok(kv.value.len() as u32)
71        }
72    }
73
74    fn into_inner(self) -> crate::Result<(UserKey, UserValue)> {
75        resolve_value_handle(
76            self.tree.id(),
77            self.tree.blobs_folder.as_path(),
78            &self.tree.index.config.cache,
79            &self.version,
80            self.kv?,
81        )
82    }
83}
84
85fn resolve_value_handle(
86    tree_id: TreeId,
87    blobs_folder: &Path,
88    cache: &Cache,
89    version: &Version,
90    item: InternalValue,
91) -> RangeItem {
92    if item.key.value_type.is_indirection() {
93        let mut cursor = std::io::Cursor::new(item.value);
94        let vptr = BlobIndirection::decode_from(&mut cursor)?;
95
96        // Resolve indirection using value log
97        match Accessor::new(&version.blob_files).get(
98            tree_id,
99            blobs_folder,
100            &item.key.user_key,
101            &vptr.vhandle,
102            cache,
103        ) {
104            Ok(Some(v)) => {
105                let k = item.key.user_key;
106                Ok((k, v))
107            }
108            Ok(None) => {
109                panic!(
110                    "value handle ({:?} => {:?}) did not match any blob - this is a bug; version={}",
111                    item.key.user_key, vptr.vhandle,
112                    version.id(),
113                );
114            }
115            Err(e) => Err(e),
116        }
117    } else {
118        let k = item.key.user_key;
119        let v = item.value;
120        Ok((k, v))
121    }
122}
123
124/// A key-value-separated log-structured merge tree
125///
126/// This tree is a composite structure, consisting of an
127/// index tree (LSM-tree) and a log-structured value log
128/// to reduce write amplification.
129#[derive(Clone)]
130pub struct BlobTree {
131    /// Index tree that holds value handles or small inline values
132    #[doc(hidden)]
133    pub index: crate::Tree,
134
135    blobs_folder: Arc<PathBuf>,
136}
137
138impl BlobTree {
139    pub(crate) fn open(config: Config) -> crate::Result<Self> {
140        use crate::file::{fsync_directory, BLOBS_FOLDER};
141
142        let index = crate::Tree::open(config)?;
143
144        let blobs_folder = index.config.path.join(BLOBS_FOLDER);
145        std::fs::create_dir_all(&blobs_folder)?;
146        fsync_directory(&blobs_folder)?;
147
148        let blob_file_id_to_continue_with = index
149            .current_version()
150            .blob_files
151            .list_ids()
152            .max()
153            .map(|x| x + 1)
154            .unwrap_or_default();
155
156        index
157            .0
158            .blob_file_id_counter
159            .set(blob_file_id_to_continue_with);
160
161        Ok(Self {
162            index,
163            blobs_folder: Arc::new(blobs_folder),
164        })
165    }
166}
167
168impl AbstractTree for BlobTree {
169    fn print_trace(&self, key: &[u8]) -> crate::Result<()> {
170        self.index.print_trace(key)
171    }
172
173    fn table_file_cache_size(&self) -> usize {
174        self.index.table_file_cache_size()
175    }
176
177    fn get_version_history_lock(
178        &self,
179    ) -> std::sync::RwLockWriteGuard<'_, crate::version::SuperVersions> {
180        self.index.get_version_history_lock()
181    }
182
183    fn next_table_id(&self) -> TableId {
184        self.index.next_table_id()
185    }
186
187    fn id(&self) -> crate::TreeId {
188        self.index.id()
189    }
190
191    fn get_internal_entry(&self, key: &[u8], seqno: SeqNo) -> crate::Result<Option<InternalValue>> {
192        self.index.get_internal_entry(key, seqno)
193    }
194
195    fn current_version(&self) -> Version {
196        self.index.current_version()
197    }
198
199    #[cfg(feature = "metrics")]
200    fn metrics(&self) -> &Arc<crate::Metrics> {
201        self.index.metrics()
202    }
203
204    fn version_free_list_len(&self) -> usize {
205        self.index.version_free_list_len()
206    }
207
208    fn prefix<K: AsRef<[u8]>>(
209        &self,
210        prefix: K,
211        seqno: SeqNo,
212        index: Option<(Arc<Memtable>, SeqNo)>,
213    ) -> Box<dyn DoubleEndedIterator<Item = IterGuardImpl> + Send + 'static> {
214        use crate::range::prefix_to_range;
215
216        let super_version = self.index.get_version_for_snapshot(seqno);
217        let tree = self.clone();
218
219        let range = prefix_to_range(prefix.as_ref());
220
221        Box::new(
222            crate::Tree::create_internal_range(super_version.clone(), &range, seqno, index).map(
223                move |kv| {
224                    IterGuardImpl::Blob(Guard {
225                        tree: tree.clone(),
226                        version: super_version.version.clone(),
227                        kv,
228                    })
229                },
230            ),
231        )
232    }
233
234    fn range<K: AsRef<[u8]>, R: RangeBounds<K>>(
235        &self,
236        range: R,
237        seqno: SeqNo,
238        index: Option<(Arc<Memtable>, SeqNo)>,
239    ) -> Box<dyn DoubleEndedIterator<Item = IterGuardImpl> + Send + 'static> {
240        let super_version = self.index.get_version_for_snapshot(seqno);
241        let tree = self.clone();
242
243        Box::new(
244            crate::Tree::create_internal_range(super_version.clone(), &range, seqno, index).map(
245                move |kv| {
246                    IterGuardImpl::Blob(Guard {
247                        tree: tree.clone(),
248                        version: super_version.version.clone(),
249                        kv,
250                    })
251                },
252            ),
253        )
254    }
255
256    fn tombstone_count(&self) -> u64 {
257        self.index.tombstone_count()
258    }
259
260    fn weak_tombstone_count(&self) -> u64 {
261        self.index.weak_tombstone_count()
262    }
263
264    fn weak_tombstone_reclaimable_count(&self) -> u64 {
265        self.index.weak_tombstone_reclaimable_count()
266    }
267
268    fn drop_range<K: AsRef<[u8]>, R: RangeBounds<K>>(&self, range: R) -> crate::Result<()> {
269        self.index.drop_range(range)
270    }
271
272    fn clear(&self) -> crate::Result<()> {
273        let config = self.tree_config();
274        let mut versions = self.get_version_history_lock();
275
276        versions.upgrade_version(
277            &config.path,
278            |v| {
279                let mut copy = v.clone();
280                copy.active_memtable =
281                    Arc::new(Memtable::new(self.index.memtable_id_counter.next()));
282                copy.sealed_memtables = Arc::default();
283                copy.version = Version::new(v.version.id() + 1, self.tree_type());
284                Ok(copy)
285            },
286            &config.seqno,
287            &config.visible_seqno,
288        )
289    }
290
291    fn major_compact(&self, target_size: u64, seqno_threshold: SeqNo) -> crate::Result<()> {
292        self.index.major_compact(target_size, seqno_threshold)
293    }
294
295    fn clear_active_memtable(&self) {
296        self.index.clear_active_memtable();
297    }
298
299    fn l0_run_count(&self) -> usize {
300        self.index.l0_run_count()
301    }
302
303    fn blob_file_count(&self) -> usize {
304        self.current_version().blob_file_count()
305    }
306
307    // NOTE: We skip reading from the value log
308    // because the vHandles already store the value size
309    fn size_of<K: AsRef<[u8]>>(&self, key: K, seqno: SeqNo) -> crate::Result<Option<u32>> {
310        let Some(item) = self.index.get_internal_entry(key.as_ref(), seqno)? else {
311            return Ok(None);
312        };
313
314        Ok(Some(if item.key.value_type.is_indirection() {
315            let mut cursor = std::io::Cursor::new(item.value);
316            let vptr = BlobIndirection::decode_from(&mut cursor)?;
317            vptr.size
318        } else {
319            #[expect(clippy::cast_possible_truncation, reason = "values are u32 length max")]
320            {
321                item.value.len() as u32
322            }
323        }))
324    }
325
326    fn stale_blob_bytes(&self) -> u64 {
327        self.current_version().gc_stats().stale_bytes()
328    }
329
330    fn filter_size(&self) -> u64 {
331        self.index.filter_size()
332    }
333
334    fn pinned_filter_size(&self) -> usize {
335        self.index.pinned_filter_size()
336    }
337
338    fn pinned_block_index_size(&self) -> usize {
339        self.index.pinned_block_index_size()
340    }
341
342    fn sealed_memtable_count(&self) -> usize {
343        self.index.sealed_memtable_count()
344    }
345
346    fn get_flush_lock(&self) -> MutexGuard<'_, ()> {
347        self.index.get_flush_lock()
348    }
349
350    fn flush_to_tables(
351        &self,
352        stream: impl Iterator<Item = crate::Result<InternalValue>>,
353    ) -> crate::Result<Option<(Vec<Table>, Option<Vec<BlobFile>>)>> {
354        use crate::{
355            coding::Encode, file::BLOBS_FOLDER, file::TABLES_FOLDER,
356            table::multi_writer::MultiWriter,
357        };
358
359        let start = std::time::Instant::now();
360
361        let table_folder = self.index.config.path.join(TABLES_FOLDER);
362
363        let data_block_size = self.index.config.data_block_size_policy.get(0);
364
365        let data_block_restart_interval =
366            self.index.config.data_block_restart_interval_policy.get(0);
367        let index_block_restart_interval =
368            self.index.config.index_block_restart_interval_policy.get(0);
369
370        let data_block_compression = self.index.config.data_block_compression_policy.get(0);
371        let index_block_compression = self.index.config.index_block_compression_policy.get(0);
372
373        let data_block_hash_ratio = self.index.config.data_block_hash_ratio_policy.get(0);
374
375        let index_partitioning = self.index.config.index_block_partitioning_policy.get(0);
376        let filter_partitioning = self.index.config.filter_block_partitioning_policy.get(0);
377
378        log::debug!("Flushing memtable(s) and performing key-value separation, data_block_restart_interval={data_block_restart_interval}, index_block_restart_interval={index_block_restart_interval}, data_block_size={data_block_size}, data_block_compression={data_block_compression:?}, index_block_compression={index_block_compression:?}");
379        log::debug!("=> to table(s) in {}", table_folder.display());
380        log::debug!("=> to blob file(s) at {}", self.blobs_folder.display());
381
382        let mut table_writer = MultiWriter::new(
383            table_folder.clone(),
384            self.index.table_id_counter.clone(),
385            64 * 1_024 * 1_024,
386            0,
387        )?
388        .use_data_block_restart_interval(data_block_restart_interval)
389        .use_index_block_restart_interval(index_block_restart_interval)
390        .use_data_block_compression(data_block_compression)
391        .use_index_block_compression(index_block_compression)
392        .use_data_block_size(data_block_size)
393        .use_data_block_hash_ratio(data_block_hash_ratio)
394        .use_bloom_policy({
395            use crate::config::FilterPolicyEntry::{Bloom, None};
396            use crate::table::filter::BloomConstructionPolicy;
397
398            match self.index.config.filter_policy.get(0) {
399                Bloom(policy) => policy,
400                None => BloomConstructionPolicy::BitsPerKey(0.0),
401            }
402        });
403
404        if index_partitioning {
405            table_writer = table_writer.use_partitioned_index();
406        }
407        if filter_partitioning {
408            table_writer = table_writer.use_partitioned_filter();
409        }
410
411        #[expect(
412            clippy::expect_used,
413            reason = "cannot create blob tree without defining kv separation options"
414        )]
415        let kv_opts = self
416            .index
417            .config
418            .kv_separation_opts
419            .as_ref()
420            .expect("kv separation options should exist");
421
422        let mut blob_writer = BlobFileWriter::new(
423            self.index.0.blob_file_id_counter.clone(),
424            self.index.config.path.join(BLOBS_FOLDER),
425            self.id(),
426            self.index.config.descriptor_table.clone(),
427        )?
428        .use_target_size(kv_opts.file_target_size)
429        .use_compression(crate::vlog::blob_file::writer::BlobCompression::Standard(
430            kv_opts.compression,
431        ));
432
433        let separation_threshold = kv_opts.separation_threshold;
434
435        for item in stream {
436            let item = item?;
437
438            if item.is_tombstone() {
439                // NOTE: Still need to add tombstone to index tree
440                // But no blob to blob writer
441                table_writer.write(InternalValue::new(item.key, UserValue::empty()))?;
442                continue;
443            }
444
445            let value = item.value;
446
447            #[expect(clippy::cast_possible_truncation, reason = "values are u32 length max")]
448            let value_size = value.len() as u32;
449
450            if value_size >= separation_threshold {
451                let vhandle = blob_writer.write(&item.key.user_key, item.key.seqno, &value)?;
452
453                let indirection = BlobIndirection {
454                    vhandle,
455                    size: value_size,
456                };
457
458                table_writer.write({
459                    let mut vptr =
460                        InternalValue::new(item.key.clone(), indirection.encode_into_vec());
461                    vptr.key.value_type = crate::ValueType::Indirection;
462                    vptr
463                })?;
464
465                table_writer.register_blob(indirection);
466            } else {
467                table_writer.write(InternalValue::new(item.key, value))?;
468            }
469        }
470
471        let blob_files = blob_writer.finish()?;
472
473        let result = table_writer.finish()?;
474
475        log::debug!("Flushed memtable(s) in {:?}", start.elapsed());
476
477        let pin_filter = self.index.config.filter_block_pinning_policy.get(0);
478        let pin_index = self.index.config.index_block_pinning_policy.get(0);
479
480        // Load tables
481        let tables = result
482            .into_iter()
483            .map(|(table_id, checksum)| -> crate::Result<Table> {
484                Table::recover(
485                    table_folder.join(table_id.to_string()),
486                    checksum,
487                    0,
488                    self.index.id,
489                    self.index.config.cache.clone(),
490                    self.index.config.descriptor_table.clone(),
491                    pin_filter,
492                    pin_index,
493                    #[cfg(feature = "metrics")]
494                    self.index.metrics.clone(),
495                )
496            })
497            .collect::<crate::Result<Vec<_>>>()?;
498
499        Ok(Some((tables, Some(blob_files))))
500    }
501
502    fn register_tables(
503        &self,
504        tables: &[Table],
505        blob_files: Option<&[BlobFile]>,
506        frag_map: Option<FragmentationMap>,
507        sealed_memtables_to_delete: &[MemtableId],
508        gc_watermark: SeqNo,
509    ) -> crate::Result<()> {
510        self.index.register_tables(
511            tables,
512            blob_files,
513            frag_map,
514            sealed_memtables_to_delete,
515            gc_watermark,
516        )
517    }
518
519    fn compact(
520        &self,
521        strategy: Arc<dyn crate::compaction::CompactionStrategy>,
522        seqno_threshold: SeqNo,
523    ) -> crate::Result<()> {
524        self.index.compact(strategy, seqno_threshold)
525    }
526
527    fn get_next_table_id(&self) -> TableId {
528        self.index.get_next_table_id()
529    }
530
531    fn tree_config(&self) -> &Config {
532        &self.index.config
533    }
534
535    fn get_highest_seqno(&self) -> Option<SeqNo> {
536        self.index.get_highest_seqno()
537    }
538
539    fn active_memtable(&self) -> Arc<Memtable> {
540        self.index.active_memtable()
541    }
542
543    fn rotate_memtable(&self) -> Option<Arc<Memtable>> {
544        self.index.rotate_memtable()
545    }
546
547    fn table_count(&self) -> usize {
548        self.index.table_count()
549    }
550
551    fn level_table_count(&self, idx: usize) -> Option<usize> {
552        self.index.level_table_count(idx)
553    }
554
555    fn approximate_len(&self) -> usize {
556        self.index.approximate_len()
557    }
558
559    // NOTE: Override the default implementation to not fetch
560    // data from the value log, so we get much faster key reads
561    fn is_empty(&self, seqno: SeqNo, index: Option<(Arc<Memtable>, SeqNo)>) -> crate::Result<bool> {
562        self.index.is_empty(seqno, index)
563    }
564
565    // NOTE: Override the default implementation to not fetch
566    // data from the value log, so we get much faster key reads
567    fn contains_key<K: AsRef<[u8]>>(&self, key: K, seqno: SeqNo) -> crate::Result<bool> {
568        self.index.contains_key(key, seqno)
569    }
570
571    // NOTE: Override the default implementation to not fetch
572    // data from the value log, so we get much faster scans
573    fn len(&self, seqno: SeqNo, index: Option<(Arc<Memtable>, SeqNo)>) -> crate::Result<usize> {
574        self.index.len(seqno, index)
575    }
576
577    fn disk_space(&self) -> u64 {
578        let version = self.current_version();
579        self.index.disk_space() + version.blob_files.on_disk_size()
580    }
581
582    fn get_highest_memtable_seqno(&self) -> Option<SeqNo> {
583        self.index.get_highest_memtable_seqno()
584    }
585
586    fn get_highest_persisted_seqno(&self) -> Option<SeqNo> {
587        self.index.get_highest_persisted_seqno()
588    }
589
590    fn insert<K: Into<UserKey>, V: Into<UserValue>>(
591        &self,
592        key: K,
593        value: V,
594        seqno: SeqNo,
595    ) -> (u64, u64) {
596        self.index.insert(key, value.into(), seqno)
597    }
598
599    fn get<K: AsRef<[u8]>>(&self, key: K, seqno: SeqNo) -> crate::Result<Option<crate::UserValue>> {
600        let key = key.as_ref();
601
602        #[expect(clippy::expect_used, reason = "lock is expected to not be poisoned")]
603        let super_version = self
604            .index
605            .version_history
606            .read()
607            .expect("lock is poisoned")
608            .get_version_for_snapshot(seqno);
609
610        let Some(item) = crate::Tree::get_internal_entry_from_version(&super_version, key, seqno)?
611        else {
612            return Ok(None);
613        };
614
615        let (_, v) = resolve_value_handle(
616            self.id(),
617            self.blobs_folder.as_path(),
618            &self.index.config.cache,
619            &super_version.version,
620            item,
621        )?;
622
623        Ok(Some(v))
624    }
625
626    fn remove<K: Into<UserKey>>(&self, key: K, seqno: SeqNo) -> (u64, u64) {
627        self.index.remove(key, seqno)
628    }
629
630    fn remove_weak<K: Into<UserKey>>(&self, key: K, seqno: SeqNo) -> (u64, u64) {
631        self.index.remove_weak(key, seqno)
632    }
633}