1mod 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
30pub 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 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#[derive(Clone)]
130pub struct BlobTree {
131 #[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 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 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 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 fn is_empty(&self, seqno: SeqNo, index: Option<(Arc<Memtable>, SeqNo)>) -> crate::Result<bool> {
562 self.index.is_empty(seqno, index)
563 }
564
565 fn contains_key<K: AsRef<[u8]>>(&self, key: K, seqno: SeqNo) -> crate::Result<bool> {
568 self.index.contains_key(key, seqno)
569 }
570
571 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}