1mod bloom_filter;
18pub mod cache;
19pub mod merge;
20pub mod persistent_layer;
21pub mod skip_list_layer;
22pub mod types;
23
24#[cfg(any(test, fuzz))]
25pub mod testing;
26
27use crate::drop_event::DropEvent;
28use crate::log::*;
29use crate::metrics::DurationMeasureScope;
30use crate::object_handle::{LayerObject, WriteBytes};
31use crate::object_store::ObjectStore;
32use crate::serialized_types::{LATEST_VERSION, Version};
33
34use anyhow::Error;
35use cache::{ObjectCache, ObjectCacheResult};
36
37use fuchsia_inspect::HistogramProperty;
38use fuchsia_sync::RwLock;
39use fxfs_crypto::Crypt;
40use persistent_layer::{PersistentLayer, PersistentLayerWriter};
41use skip_list_layer::SkipListLayer;
42use std::fmt;
43use std::ops::Bound;
44use std::sync::atomic::{AtomicUsize, Ordering};
45use std::sync::{Arc, Mutex};
46use storage_units::BlockSize;
47use types::{
48 Existence, Item, ItemRef, Key, Layer, LayerIterator, LayerKey, LayerWriter, MaybeContainsKey,
49 MergeType, MergeableKey, OrdLowerBound, Value,
50};
51
52pub use merge::Query;
53
54const SKIP_LIST_LAYER_ITEMS: usize = 512;
55
56pub use persistent_layer::{
58 LayerHeader as PersistentLayerHeader, LayerHeaderV39 as PersistentLayerHeaderV39,
59 LayerInfo as PersistentLayerInfo, LayerInfoV39 as PersistentLayerInfoV39, SyncPersistentLayer,
60 layer_from_handle, open_layers,
61};
62
63pub async fn layers_from_handles<K: Key, V: Value>(
64 handles: impl IntoIterator<Item = impl LayerObject + 'static>,
65) -> Result<Vec<Arc<dyn Layer<K, V>>>, Error> {
66 let mut layers = Vec::new();
67 for handle in handles {
68 layers.push(PersistentLayer::open(handle).await?);
69 }
70 Ok(layers)
71}
72
73#[derive(Eq, PartialEq, Debug)]
74pub enum Operation {
75 Insert,
76 ReplaceOrInsert,
77 MergeInto,
78}
79
80pub type MutationCallback<K, V> = Option<Box<dyn Fn(Operation, &Item<K, V>) + Send + Sync>>;
81
82struct Inner<K, V> {
83 mutable_layer: Arc<SkipListLayer<K, V>>,
84 layers: Vec<Arc<dyn Layer<K, V>>>,
85 mutation_callback: MutationCallback<K, V>,
86}
87
88pub const LOG2_HISTOGRAM_BUCKETS: usize = 32;
89
90pub struct CompactionCounters {
92 pub compactions: u64,
94 pub compaction_bytes_written: u64,
96 pub compaction_time_ns: u64,
98 pub total_layers_added: u64,
100 pub max_layer_count: u64,
102 pub layer_size_histogram: [u64; LOG2_HISTOGRAM_BUCKETS],
104}
105
106impl Default for CompactionCounters {
107 fn default() -> Self {
108 Self {
109 compactions: 0,
110 compaction_bytes_written: 0,
111 compaction_time_ns: 0,
112 total_layers_added: 0,
113 max_layer_count: 0,
114 layer_size_histogram: [0; LOG2_HISTOGRAM_BUCKETS],
115 }
116 }
117}
118
119pub struct TreeCounters {
121 pub num_seeks: AtomicUsize,
123 pub layer_files_total: AtomicUsize,
126 pub layer_files_skipped: AtomicUsize,
128 pub excessive_hash_partitions: AtomicUsize,
130 pub compaction: Mutex<CompactionCounters>,
132}
133
134impl Default for TreeCounters {
135 fn default() -> Self {
136 Self {
137 num_seeks: AtomicUsize::new(0),
138 layer_files_total: AtomicUsize::new(0),
139 layer_files_skipped: AtomicUsize::new(0),
140 excessive_hash_partitions: AtomicUsize::new(0),
141 compaction: Mutex::new(CompactionCounters::default()),
142 }
143 }
144}
145
146#[fxfs_trace::trace]
148pub async fn compact_with_iterator<K: Key, V: Value, W: WriteBytes + Send>(
149 mut iterator: impl LayerIterator<K, V>,
150 num_items: usize,
151 writer: W,
152 block_size: BlockSize,
153 mut yielder: Option<impl Yielder>,
154) -> Result<u64, Error> {
155 let mut writer = PersistentLayerWriter::<W, K, V>::new(writer, num_items, block_size).await?;
156 while let Some(item_ref) = iterator.get() {
157 debug!(item_ref:?; "compact: writing");
158 writer.write(item_ref).await?;
159 iterator.advance().await?;
160 if let Some(y) = yielder.as_mut() {
161 y.yield_now().await;
162 }
163 }
164 writer.complete().await
165}
166
167pub struct LSMTree<K, V> {
171 data: RwLock<Inner<K, V>>,
172 merge_fn: merge::MergeFn<K, V>,
173 cache: Option<Box<dyn ObjectCache<K, V>>>,
174 counters: Arc<TreeCounters>,
175}
176
177#[fxfs_trace::trace]
178impl<'tree, K: MergeableKey, V: Value> LSMTree<K, V> {
179 pub fn new(merge_fn: merge::MergeFn<K, V>, cache: Option<Box<dyn ObjectCache<K, V>>>) -> Self {
181 let counters = TreeCounters::default();
182 counters.compaction.lock().unwrap().max_layer_count = 1;
183 LSMTree {
184 data: RwLock::new(Inner {
185 mutable_layer: Self::new_mutable_layer(),
186 layers: Vec::new(),
187 mutation_callback: None,
188 }),
189 merge_fn,
190 cache,
191 counters: Arc::new(counters),
192 }
193 }
194
195 pub async fn open(
197 merge_fn: merge::MergeFn<K, V>,
198 handles: impl IntoIterator<Item = impl LayerObject + 'static>,
199 cache: Option<Box<dyn ObjectCache<K, V>>>,
200 ) -> Result<Self, Error> {
201 let layers = layers_from_handles(handles).await?;
202 let max_layer_count = layers.len() as u64 + 1;
203 let counters = TreeCounters::default();
204 counters.compaction.lock().unwrap().max_layer_count = max_layer_count;
205 Ok(LSMTree {
206 data: RwLock::new(Inner {
207 mutable_layer: Self::new_mutable_layer(),
208 layers,
209 mutation_callback: None,
210 }),
211 merge_fn,
212 cache,
213 counters: Arc::new(counters),
214 })
215 }
216
217 pub fn set_layers(&self, layers: Vec<Arc<dyn Layer<K, V>>>) {
219 let mut data = self.data.write();
220 data.layers = layers;
221 let layer_count = data.layers.len() + 1;
222 let mut counters = self.counters.compaction.lock().unwrap();
223 counters.max_layer_count = std::cmp::max(counters.max_layer_count, layer_count as u64);
224 }
225
226 pub async fn append_layers(
232 &self,
233 store: &Arc<ObjectStore>,
234 object_ids: impl IntoIterator<Item = u64>,
235 crypt: Option<Arc<dyn Crypt>>,
236 ) -> Result<u64, Error> {
237 let (layers, total_size) = open_layers(store, object_ids, crypt).await?;
238 self.append_open_layers(layers);
239 Ok(total_size)
240 }
241
242 pub fn append_open_layers(&self, mut layers: Vec<Arc<dyn Layer<K, V>>>) {
244 let mut data = self.data.write();
245 data.layers.append(&mut layers);
246 let layer_count = data.layers.len() + 1;
247 let mut counters = self.counters.compaction.lock().unwrap();
248 counters.max_layer_count = std::cmp::max(counters.max_layer_count, layer_count as u64);
249 }
250
251 pub fn reset_immutable_layers(&self) {
253 self.data.write().layers = Vec::new();
254 }
255
256 pub fn seal(&self) {
258 let mut data = self.data.write();
261 let layer = std::mem::replace(&mut data.mutable_layer, Self::new_mutable_layer());
262 data.layers.insert(0, layer);
263 let layer_count = data.layers.len() + 1;
264 let mut counters = self.counters.compaction.lock().unwrap();
265 counters.max_layer_count = std::cmp::max(counters.max_layer_count, layer_count as u64);
266 counters.total_layers_added += 1;
267 }
268
269 pub fn reset(&self) {
271 let mut data = self.data.write();
272 data.layers = Vec::new();
273 data.mutable_layer = Self::new_mutable_layer();
274 }
275
276 pub fn clear_cache(&self) {
278 if let Some(cache) = &self.cache {
279 cache.clear();
280 }
281 for layer in self.immutable_layer_set().layers {
282 layer.clear_cached_data();
283 }
284 }
285
286 pub fn cache_len(&self) -> usize {
287 self.cache.as_ref().map_or(0, |c| c.len())
288 }
289
290 pub fn report_compaction_metrics(
291 &self,
292 bytes_written: u64,
293 duration: std::time::Duration,
294 layer_count: usize,
295 ) {
296 let mut counters = self.counters.compaction.lock().unwrap();
297 counters.compactions += 1;
298 counters.compaction_bytes_written += bytes_written;
299 counters.compaction_time_ns += duration.as_nanos() as u64;
300
301 let bucket = if bytes_written == 0 {
302 0
303 } else {
304 std::cmp::min(LOG2_HISTOGRAM_BUCKETS - 1, 63 - bytes_written.leading_zeros() as usize)
305 };
306 counters.layer_size_histogram[bucket] += 1;
307
308 crate::metrics::lsm_tree_metrics().compaction_layer_stack_depth.insert(layer_count as u64);
309 }
310
311 pub fn compaction_bytes_written(&self) -> u64 {
312 self.counters.compaction.lock().unwrap().compaction_bytes_written
313 }
314
315 pub fn empty_layer_set(&self) -> LayerSet<K, V> {
317 LayerSet { layers: Vec::new(), merge_fn: self.merge_fn, counters: self.counters.clone() }
318 }
319
320 pub fn add_all_layers_to_layer_set(&self, layer_set: &mut LayerSet<K, V>) {
322 let data = self.data.read();
323 layer_set.layers.reserve_exact(data.layers.len() + 1);
324 layer_set
325 .layers
326 .push(LockedLayer::from(data.mutable_layer.clone() as Arc<dyn Layer<K, V>>));
327 for layer in &data.layers {
328 layer_set.layers.push(layer.clone().into());
329 }
330 }
331
332 pub fn layer_set(&self) -> LayerSet<K, V> {
335 let mut layer_set = self.empty_layer_set();
336 self.add_all_layers_to_layer_set(&mut layer_set);
337 layer_set
338 }
339
340 pub fn immutable_layer_set(&self) -> LayerSet<K, V> {
344 let data = self.data.read();
345 let mut layers = Vec::with_capacity(data.layers.len());
346 for layer in &data.layers {
347 layers.push(layer.clone().into());
348 }
349 LayerSet { layers, merge_fn: self.merge_fn, counters: self.counters.clone() }
350 }
351
352 pub fn insert(&self, item: Item<K, V>) -> Result<(), Error> {
355 let _measure = DurationMeasureScope::new(&crate::metrics::lsm_tree_metrics().insert);
356
357 let item_dup = if self.is_cacheable(&item.key) { Some(item.clone()) } else { None };
358
359 {
360 let data = self.data.read();
362 if let Some(mutation_callback) = data.mutation_callback.as_ref() {
363 mutation_callback(Operation::Insert, &item);
364 }
365 data.mutable_layer.insert(item)?;
366 }
367
368 if let Some(item) = item_dup {
369 let val = if item.value == V::DELETED_MARKER { None } else { Some(item.value) };
370 self.invalidate_cache(&item.key, val);
371 }
372 Ok(())
373 }
374
375 pub fn replace_or_insert(&self, item: Item<K, V>) {
377 let _measure =
378 DurationMeasureScope::new(&crate::metrics::lsm_tree_metrics().replace_or_insert);
379
380 let item_dup = if self.is_cacheable(&item.key) { Some(item.clone()) } else { None };
381
382 {
383 let data = self.data.read();
385 if let Some(mutation_callback) = data.mutation_callback.as_ref() {
386 mutation_callback(Operation::ReplaceOrInsert, &item);
387 }
388 data.mutable_layer.replace_or_insert(item);
389 }
390
391 if let Some(item) = item_dup {
392 let val = if item.value == V::DELETED_MARKER { None } else { Some(item.value) };
393 self.invalidate_cache(&item.key, val);
394 }
395 }
396
397 pub fn merge_into(&self, item: Item<K, V>, lower_bound: &K) {
399 let _measure = DurationMeasureScope::new(&crate::metrics::lsm_tree_metrics().merge_into);
400
401 let key = if self.is_cacheable(&item.key) { Some(item.key.clone()) } else { None };
402 {
403 let data = self.data.read();
405 if let Some(mutation_callback) = data.mutation_callback.as_ref() {
406 mutation_callback(Operation::MergeInto, &item);
407 }
408 data.mutable_layer.merge_into(item, lower_bound, self.merge_fn);
409 }
410
411 if let Some(key) = key {
412 self.invalidate_cache(&key, None);
413 }
414 }
415
416 #[inline(always)]
422 fn is_cacheable(&self, key: &K) -> bool {
423 if let Some(cache) = &self.cache { cache.is_cacheable(key) } else { false }
424 }
425
426 #[inline(always)]
431 fn invalidate_cache(&self, key: &K, value: Option<V>) {
432 if let Some(cache) = &self.cache {
433 cache.invalidate(key, value);
434 }
435 }
436
437 pub async fn find_map<R, F>(&self, search_key: &K, f: F) -> Result<Option<R>, Error>
443 where
444 K: Eq,
445 F: FnOnce(ItemRef<'_, K, V>) -> R,
446 {
447 let _measure = DurationMeasureScope::new(&crate::metrics::lsm_tree_metrics().find);
448 let mut token = if let Some(cache) = &self.cache {
452 match cache.lookup_or_reserve(search_key) {
453 ObjectCacheResult::Value(value) => {
454 if *value == V::DELETED_MARKER {
455 return Ok(None);
456 } else {
457 return Ok(Some(f(ItemRef { key: search_key, value: &value })));
458 }
459 }
460 ObjectCacheResult::Placeholder(token) => Some(token),
461 ObjectCacheResult::NoCache => None,
462 }
463 } else {
464 None
465 };
466 let layer_set = self.layer_set();
467 let result = layer_set
468 .find_map(search_key, |item_ref| {
469 if let Some(token) = token.take() {
470 token.complete(Some(item_ref.value));
471 }
472 f(item_ref)
473 })
474 .await?;
475 Ok(result)
476 }
477
478 pub fn find(
483 &self,
484 search_key: &K,
485 ) -> impl Future<Output = Result<Option<Item<K, V>>, Error>> + Send
486 where
487 K: Eq,
488 {
489 self.find_map(search_key, |item| Item { key: item.key.clone(), value: item.value.clone() })
490 }
491
492 pub fn find_value(
497 &self,
498 search_key: &K,
499 ) -> impl Future<Output = Result<Option<V>, Error>> + Send
500 where
501 K: Eq,
502 {
503 self.find_map(search_key, |item| item.value.clone())
504 }
505
506 pub async fn exists(&self, search_key: &K) -> Result<bool, Error>
510 where
511 K: Eq,
512 {
513 Ok(self.find_map(search_key, |_| ()).await?.is_some())
514 }
515
516 pub fn mutable_layer(&self) -> Arc<SkipListLayer<K, V>> {
517 self.data.read().mutable_layer.clone()
518 }
519
520 pub fn set_mutation_callback(&self, mutation_callback: MutationCallback<K, V>) {
524 self.data.write().mutation_callback = mutation_callback;
525 }
526
527 pub fn get_earliest_version(&self) -> Version {
529 let mut earliest_version = LATEST_VERSION;
530 for layer in self.layer_set().layers {
531 let layer_version = layer.get_version();
532 if layer_version < earliest_version {
533 earliest_version = layer_version;
534 }
535 }
536 return earliest_version;
537 }
538
539 pub fn new_mutable_layer() -> Arc<SkipListLayer<K, V>> {
541 SkipListLayer::new(SKIP_LIST_LAYER_ITEMS)
542 }
543
544 pub fn set_mutable_layer(&self, layer: Arc<SkipListLayer<K, V>>) {
546 self.data.write().mutable_layer = layer;
547 }
548
549 pub fn record_inspect_data(&self, root: &fuchsia_inspect::Node) {
551 let layer_set = self.layer_set();
552 root.record_child("layers", move |node| {
553 let mut index = 0;
554 for layer in layer_set.layers {
555 node.record_child(format!("{index}"), move |node| {
556 layer.1.record_inspect_data(node)
557 });
558 index += 1;
559 }
560 });
561 {
562 let counters = self.counters.compaction.lock().unwrap();
563 root.record_uint("num_seeks", self.counters.num_seeks.load(Ordering::Relaxed) as u64);
564 root.record_uint("bloom_filter_success_percent", {
565 let layer_files_total = self.counters.layer_files_total.load(Ordering::Relaxed);
566 let layer_files_skipped = self.counters.layer_files_skipped.load(Ordering::Relaxed);
567 if layer_files_total == 0 {
568 0
569 } else {
570 (layer_files_skipped * 100).div_ceil(layer_files_total) as u64
571 }
572 });
573 root.record_uint(
574 "excessive_hash_partitions",
575 self.counters.excessive_hash_partitions.load(Ordering::Relaxed) as u64,
576 );
577 root.record_uint("compactions", counters.compactions);
578 root.record_uint("compaction_bytes_written", counters.compaction_bytes_written);
579 root.record_uint("compaction_time_ns", counters.compaction_time_ns);
580 root.record_uint("total_layers_added", counters.total_layers_added);
581 root.record_uint("max_layer_count", counters.max_layer_count);
582
583 let layer_sizes = root.create_uint_exponential_histogram(
584 "layer_size_histogram_log2",
585 fuchsia_inspect::ExponentialHistogramParams {
586 floor: 1,
587 initial_step: 1,
588 step_multiplier: 2,
589 buckets: LOG2_HISTOGRAM_BUCKETS,
590 },
591 );
592 for (i, count) in counters.layer_size_histogram.iter().enumerate() {
593 layer_sizes.insert_multiple(1u64 << i, *count as usize);
594 }
595 root.record(layer_sizes);
596 }
597 }
598}
599
600pub struct LockedLayer<K, V>(Arc<DropEvent>, Arc<dyn Layer<K, V>>);
603
604impl<K, V> LockedLayer<K, V> {
605 pub async fn close_layer(self) {
606 let layer = self.1;
607 std::mem::drop(self.0);
608 layer.close().await;
609 }
610}
611
612impl<K, V> From<Arc<dyn Layer<K, V>>> for LockedLayer<K, V> {
613 fn from(layer: Arc<dyn Layer<K, V>>) -> Self {
614 let event = layer.lock().unwrap();
615 Self(event, layer)
616 }
617}
618
619impl<K, V> std::ops::Deref for LockedLayer<K, V> {
620 type Target = Arc<dyn Layer<K, V>>;
621
622 fn deref(&self) -> &Self::Target {
623 &self.1
624 }
625}
626
627impl<K, V> AsRef<dyn Layer<K, V>> for LockedLayer<K, V> {
628 fn as_ref(&self) -> &(dyn Layer<K, V> + 'static) {
629 self.1.as_ref()
630 }
631}
632
633pub struct LayerSet<K, V> {
636 pub layers: Vec<LockedLayer<K, V>>,
637 merge_fn: merge::MergeFn<K, V>,
638 counters: Arc<TreeCounters>,
639}
640
641impl<K: Key + LayerKey + OrdLowerBound, V: Value> LayerSet<K, V> {
642 pub fn sum_len(&self) -> usize {
643 let mut size = 0;
644 for layer in &self.layers {
645 size += layer.len()
646 }
647 size
648 }
649
650 pub fn merger(&self) -> merge::Merger<'_, K, V> {
651 merge::Merger::new(
652 self.layers.iter().map(|x| x.as_ref()),
653 self.merge_fn,
654 self.counters.clone(),
655 )
656 }
657
658 pub async fn find_map<R, F>(&self, search_key: &K, f: F) -> Result<Option<R>, Error>
664 where
665 K: Eq,
666 F: FnOnce(ItemRef<'_, K, V>) -> R,
667 {
668 if search_key.merge_type() == MergeType::OptimizedMerge {
669 self.counters.num_seeks.fetch_add(1, Ordering::Relaxed);
670 for layer in &self.layers {
671 self.counters.layer_files_total.fetch_add(1, Ordering::Relaxed);
672 if layer.maybe_contains_key(search_key) == MaybeContainsKey::False {
673 self.counters.layer_files_skipped.fetch_add(1, Ordering::Relaxed);
674 continue;
675 }
676 let iter = layer.seek(Bound::Included(search_key)).await?;
677 if let Some(item_ref) = iter.get() {
678 if item_ref.key == search_key {
679 if *item_ref.value == V::DELETED_MARKER {
680 return Ok(None);
681 } else {
682 return Ok(Some(f(item_ref)));
683 }
684 }
685 }
686 }
687 return Ok(None);
688 }
689
690 let mut merger = self.merger();
691 Ok(match merger.query(Query::Point(search_key)).await?.get() {
692 Some(item_ref)
693 if item_ref.key == search_key && *item_ref.value != V::DELETED_MARKER =>
694 {
695 Some(f(item_ref))
696 }
697 _ => None,
698 })
699 }
700
701 pub async fn key_exists(&self, key: &K) -> Result<Existence, Error> {
703 for l in &self.layers {
704 match l.key_exists(key).await? {
705 e @ (Existence::Exists | Existence::MaybeExists) => return Ok(e),
706 _ => {}
707 }
708 }
709 Ok(Existence::Missing)
710 }
711}
712
713impl<K, V> fmt::Debug for LayerSet<K, V> {
714 fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
715 fmt.debug_list()
716 .entries(self.layers.iter().map(|l| {
717 if let Some(handle) = l.handle() {
718 format!("{}", handle.object_id())
719 } else {
720 format!("{:?}", Arc::as_ptr(l))
721 }
722 }))
723 .finish()
724 }
725}
726
727pub trait Yielder: Send {
729 fn yield_now(&mut self) -> impl Future<Output = ()> + Send;
730}
731
732#[cfg(test)]
733mod tests {
734 use super::{LSMTree, Yielder, compact_with_iterator};
735 use crate::drop_event::DropEvent;
736 use crate::lsm_tree::cache::{ObjectCache, ObjectCachePlaceholder, ObjectCacheResult};
737 use crate::lsm_tree::merge::{ItemOp, MergeLayerIterator, MergeResult};
738 use crate::lsm_tree::types::{
739 BoxedLayerIterator, Existence, Item, ItemRef, Key, Layer, LayerIterator, MaybeContainsKey,
740 Value,
741 };
742 use crate::lsm_tree::{Query, layers_from_handles};
743 use crate::object_handle::ObjectHandle;
744 use crate::serialized_types::{LATEST_VERSION, Version};
745 use crate::testing::fake_object::{FakeObject, FakeObjectHandle};
746 use crate::testing::writer::Writer;
747 use anyhow::{Error, anyhow};
748 use async_trait::async_trait;
749
750 use fuchsia_sync::{Mutex, MutexGuard};
751
752 use rand::rng;
753 use rand::seq::SliceRandom as _;
754
755 use std::sync::Arc;
756 use std::sync::atomic::Ordering;
757
758 use super::testing::TestKey;
759
760 fn emit_left_merge_fn(
761 _left: &MergeLayerIterator<'_, TestKey, u64>,
762 _right: &MergeLayerIterator<'_, TestKey, u64>,
763 ) -> MergeResult<TestKey, u64> {
764 MergeResult::EmitLeft
765 }
766
767 fn emit_left_i32_merge_fn(
768 _left: &MergeLayerIterator<'_, i32, u64>,
769 _right: &MergeLayerIterator<'_, i32, u64>,
770 ) -> MergeResult<i32, u64> {
771 MergeResult::EmitLeft
772 }
773
774 fn merge_sum_i32(
775 left: &MergeLayerIterator<'_, i32, u64>,
776 right: &MergeLayerIterator<'_, i32, u64>,
777 ) -> MergeResult<i32, u64> {
778 if left.key() == right.key() {
779 MergeResult::Other {
780 emit: None,
781 left: ItemOp::Discard,
782 right: ItemOp::Replace(
783 Item::new(*left.key(), *left.value() + *right.value()).boxed(),
784 ),
785 }
786 } else {
787 MergeResult::EmitLeft
788 }
789 }
790
791 impl Value for u64 {
792 const DELETED_MARKER: Self = 0;
793 }
794
795 struct NoOpYielder;
796 impl Yielder for NoOpYielder {
797 async fn yield_now(&mut self) {}
798 }
799
800 #[fuchsia::test]
801 async fn test_iteration() {
802 let tree = LSMTree::new(emit_left_merge_fn, None);
803 let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
804 tree.insert(items[0].clone()).expect("insert error");
805 tree.insert(items[1].clone()).expect("insert error");
806 let layers = tree.layer_set();
807 let mut merger = layers.merger();
808 let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
809 let ItemRef { key, value, .. } = iter.get().expect("missing item");
810 assert_eq!((key, value), (&items[0].key, &items[0].value));
811 iter.advance().await.expect("advance failed");
812 let ItemRef { key, value, .. } = iter.get().expect("missing item");
813 assert_eq!((key, value), (&items[1].key, &items[1].value));
814 iter.advance().await.expect("advance failed");
815 assert!(iter.get().is_none());
816 }
817
818 #[fuchsia::test]
819 async fn test_compact() {
820 let tree = LSMTree::new(emit_left_merge_fn, None);
821 let items = [
822 Item::new(TestKey(1..1), 1),
823 Item::new(TestKey(2..2), 2),
824 Item::new(TestKey(3..3), 3),
825 Item::new(TestKey(4..4), 4),
826 ];
827 tree.insert(items[0].clone()).expect("insert error");
828 tree.insert(items[1].clone()).expect("insert error");
829 tree.seal();
830 tree.insert(items[2].clone()).expect("insert error");
831 tree.insert(items[3].clone()).expect("insert error");
832 tree.seal();
833 let object = Arc::new(FakeObject::new());
834 let handle = FakeObjectHandle::new(object.clone());
835 {
836 let layer_set = tree.immutable_layer_set();
837 let mut merger = layer_set.merger();
838 let iter = merger.query(Query::FullScan).await.expect("create merger");
839 compact_with_iterator(
840 iter,
841 items.len(),
842 Writer::new(&handle).await,
843 handle.block_size(),
844 Option::<NoOpYielder>::None,
845 )
846 .await
847 .expect("compact failed");
848 }
849 tree.set_layers(layers_from_handles([handle]).await.expect("layers_from_handles failed"));
850 let handle = FakeObjectHandle::new(object.clone());
851 let tree = LSMTree::open(emit_left_merge_fn, [handle], None).await.expect("open failed");
852
853 let layers = tree.layer_set();
854 let mut merger = layers.merger();
855 let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
856 for i in 1..5 {
857 let ItemRef { key, value, .. } = iter.get().expect("missing item");
858 assert_eq!((key, value), (&TestKey(i..i), &i));
859 iter.advance().await.expect("advance failed");
860 }
861 assert!(iter.get().is_none());
862 }
863
864 #[fuchsia::test]
865 async fn test_find() {
866 let items = [
867 Item::new(TestKey(1..1), 1),
868 Item::new(TestKey(2..2), 2),
869 Item::new(TestKey(3..3), 3),
870 Item::new(TestKey(4..4), 4),
871 ];
872 let tree = LSMTree::new(emit_left_merge_fn, None);
873 tree.insert(items[0].clone()).expect("insert error");
874 tree.insert(items[1].clone()).expect("insert error");
875 tree.seal();
876 tree.insert(items[2].clone()).expect("insert error");
877 tree.insert(items[3].clone()).expect("insert error");
878
879 let item = tree.find(&items[1].key).await.expect("find failed").expect("not found");
880 assert_eq!(item, items[1]);
881 assert!(tree.find(&TestKey(100..100)).await.expect("find failed").is_none());
882
883 assert_eq!(tree.find_value(&items[2].key).await.expect("find_value failed"), Some(3));
884 assert_eq!(tree.find_value(&TestKey(100..100)).await.expect("find_value failed"), None);
885
886 assert!(tree.exists(&items[0].key).await.expect("exists failed"));
887 assert!(!tree.exists(&TestKey(100..100)).await.expect("exists failed"));
888
889 assert_eq!(
890 tree.find_map(&items[3].key, |item| item.value * 10).await.expect("find_map failed"),
891 Some(40)
892 );
893 assert_eq!(
894 tree.find_map(&TestKey(100..100), |item| item.value * 10)
895 .await
896 .expect("find_map failed"),
897 None
898 );
899 }
900
901 #[fuchsia::test]
902 async fn test_find_no_return_deleted_values() {
903 let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), u64::DELETED_MARKER)];
904 let tree = LSMTree::new(emit_left_merge_fn, None);
905 tree.insert(items[0].clone()).expect("insert error");
906 tree.insert(items[1].clone()).expect("insert error");
907
908 let item = tree.find(&items[0].key).await.expect("find failed").expect("not found");
909 assert_eq!(item, items[0]);
910 assert_eq!(tree.find_value(&items[0].key).await.expect("find_value failed"), Some(1));
911 assert_eq!(
912 tree.find_map(&items[0].key, |item| *item.value).await.expect("find_map failed"),
913 Some(1)
914 );
915 assert!(tree.exists(&items[0].key).await.expect("exists failed"));
916
917 assert!(tree.find(&items[1].key).await.expect("find failed").is_none());
918 assert!(tree.find_value(&items[1].key).await.expect("find_value failed").is_none());
919 assert!(tree.find_map(&items[1].key, |_| ()).await.expect("find_map failed").is_none());
920 assert!(!tree.exists(&items[1].key).await.expect("exists failed"));
921 }
922
923 #[fuchsia::test]
924 async fn test_find_full_merge() {
925 let tree = LSMTree::new(emit_left_i32_merge_fn, None);
926 let items = [
927 Item::new(1, 10),
928 Item::new(2, 20),
929 Item::new(3, 30),
930 Item::new(4, u64::DELETED_MARKER),
931 ];
932 for item in &items {
933 tree.insert(item.clone()).expect("insert error");
934 }
935
936 assert_eq!(tree.find(&1).await.expect("find failed"), Some(items[0].clone()));
937 assert_eq!(tree.find_value(&2).await.expect("find_value failed"), Some(20));
938 assert!(tree.exists(&3).await.expect("exists failed"));
939 assert_eq!(
940 tree.find_map(&3, |item| *item.value * 2).await.expect("find_map failed"),
941 Some(60)
942 );
943
944 assert_eq!(tree.find(&99).await.expect("find failed"), None);
946 assert_eq!(tree.find_value(&99).await.expect("find_value failed"), None);
947 assert!(!tree.exists(&99).await.expect("exists failed"));
948 assert_eq!(tree.find_map(&99, |_| ()).await.expect("find_map failed"), None);
949
950 assert_eq!(tree.find(&4).await.expect("find failed"), None);
952 assert_eq!(tree.find_value(&4).await.expect("find_value failed"), None);
953 assert!(!tree.exists(&4).await.expect("exists failed"));
954 assert_eq!(tree.find_map(&4, |_| ()).await.expect("find_map failed"), None);
955 }
956
957 #[fuchsia::test]
958 async fn test_find_full_merge_multi_layer() {
959 let tree = LSMTree::new(merge_sum_i32, None);
960
961 tree.insert(Item::new(1, 100)).expect("insert error");
963 tree.insert(Item::new(2, 200)).expect("insert error");
964 tree.insert(Item::new(3, 300)).expect("insert error");
965
966 tree.seal();
967
968 tree.insert(Item::new(1, 50)).expect("insert error");
972 tree.insert(Item::new(4, 40)).expect("insert error");
973
974 assert_eq!(tree.find_value(&1).await.expect("find_value failed"), Some(150));
976 assert_eq!(tree.find(&1).await.expect("find failed"), Some(Item::new(1, 150)));
977 assert!(tree.exists(&1).await.expect("exists failed"));
978 assert_eq!(
979 tree.find_map(&1, |item| *item.value * 2).await.expect("find_map failed"),
980 Some(300)
981 );
982
983 assert_eq!(tree.find_value(&2).await.expect("find_value failed"), Some(200));
985 assert!(tree.exists(&2).await.expect("exists failed"));
986
987 assert_eq!(tree.find_value(&3).await.expect("find_value failed"), Some(300));
989 assert!(tree.exists(&3).await.expect("exists failed"));
990
991 assert_eq!(tree.find_value(&4).await.expect("find_value failed"), Some(40));
993 assert!(tree.exists(&4).await.expect("exists failed"));
994
995 assert_eq!(tree.find_value(&5).await.expect("find_value failed"), None);
997 assert!(!tree.exists(&5).await.expect("exists failed"));
998 }
999
1000 #[fuchsia::test]
1001 async fn test_find_optimized_merge_multi_layer() {
1002 let tree = LSMTree::new(emit_left_merge_fn, None);
1003
1004 tree.insert(Item::new(TestKey(1..1), 10)).expect("insert error");
1006 tree.insert(Item::new(TestKey(2..2), 20)).expect("insert error");
1007 tree.insert(Item::new(TestKey(3..3), 30)).expect("insert error");
1008
1009 tree.seal();
1010
1011 tree.insert(Item::new(TestKey(2..2), 25)).expect("insert error");
1016 tree.insert(Item::new(TestKey(3..3), u64::DELETED_MARKER)).expect("insert error");
1017 tree.insert(Item::new(TestKey(4..4), 40)).expect("insert error");
1018
1019 assert_eq!(tree.find_value(&TestKey(1..1)).await.expect("find_value failed"), Some(10));
1021 assert!(tree.exists(&TestKey(1..1)).await.expect("exists failed"));
1022
1023 assert_eq!(tree.find_value(&TestKey(2..2)).await.expect("find_value failed"), Some(25));
1025 assert!(tree.exists(&TestKey(2..2)).await.expect("exists failed"));
1026
1027 assert_eq!(tree.find(&TestKey(3..3)).await.expect("find failed"), None);
1029 assert_eq!(tree.find_value(&TestKey(3..3)).await.expect("find_value failed"), None);
1030 assert!(!tree.exists(&TestKey(3..3)).await.expect("exists failed"));
1031 assert_eq!(tree.find_map(&TestKey(3..3), |_| ()).await.expect("find_map failed"), None);
1032
1033 assert_eq!(tree.find_value(&TestKey(4..4)).await.expect("find_value failed"), Some(40));
1035 assert!(tree.exists(&TestKey(4..4)).await.expect("exists failed"));
1036
1037 assert_eq!(tree.find_value(&TestKey(5..5)).await.expect("find_value failed"), None);
1039 assert!(!tree.exists(&TestKey(5..5)).await.expect("exists failed"));
1040 }
1041
1042 #[fuchsia::test]
1043 async fn test_find_bloom_filter_rejection() {
1044 struct BloomFilterMockLayer {
1045 drop_event: Mutex<Option<Arc<DropEvent>>>,
1046 }
1047
1048 impl BloomFilterMockLayer {
1049 fn new() -> Self {
1050 Self { drop_event: Mutex::new(Some(Arc::new(DropEvent::new()))) }
1051 }
1052 }
1053
1054 #[async_trait]
1055 impl<K: Key, V: Value> Layer<K, V> for BloomFilterMockLayer {
1056 async fn seek(
1057 &self,
1058 _bound: std::ops::Bound<&K>,
1059 ) -> Result<BoxedLayerIterator<'_, K, V>, Error> {
1060 panic!("seek should not be called on layer skipped by bloom filter");
1061 }
1062
1063 fn lock(&self) -> Option<Arc<DropEvent>> {
1064 self.drop_event.lock().clone()
1065 }
1066
1067 fn len(&self) -> usize {
1068 0
1069 }
1070
1071 async fn close(&self) {}
1072
1073 fn get_version(&self) -> Version {
1074 LATEST_VERSION
1075 }
1076
1077 fn maybe_contains_key(&self, _key: &K) -> MaybeContainsKey {
1078 MaybeContainsKey::False
1079 }
1080
1081 async fn key_exists(&self, _key: &K) -> Result<Existence, Error> {
1082 unimplemented!()
1083 }
1084 }
1085
1086 let tree = LSMTree::new(emit_left_merge_fn, None);
1087 let layer: Arc<dyn Layer<TestKey, u64>> = Arc::new(BloomFilterMockLayer::new());
1088 tree.set_layers(vec![layer]);
1089
1090 let skipped_before = tree.counters.layer_files_skipped.load(Ordering::Relaxed);
1091 let total_before = tree.counters.layer_files_total.load(Ordering::Relaxed);
1092
1093 assert_eq!(tree.find_value(&TestKey(1..1)).await.expect("find failed"), None);
1096 assert!(!tree.exists(&TestKey(1..1)).await.expect("exists failed"));
1097
1098 let skipped_after = tree.counters.layer_files_skipped.load(Ordering::Relaxed);
1099 let total_after = tree.counters.layer_files_total.load(Ordering::Relaxed);
1100
1101 assert_eq!(skipped_after - skipped_before, 2);
1104 assert_eq!(total_after - total_before, 4);
1105 }
1106
1107 #[fuchsia::test]
1108 async fn test_empty_seal() {
1109 let tree = LSMTree::new(emit_left_merge_fn, None);
1110 tree.seal();
1111 let item = Item::new(TestKey(1..1), 1);
1112 tree.insert(item.clone()).expect("insert error");
1113 let object = Arc::new(FakeObject::new());
1114 let handle = FakeObjectHandle::new(object.clone());
1115 {
1116 let layer_set = tree.immutable_layer_set();
1117 let mut merger = layer_set.merger();
1118 let iter = merger.query(Query::FullScan).await.expect("create merger");
1119 compact_with_iterator(
1120 iter,
1121 0,
1122 Writer::new(&handle).await,
1123 handle.block_size(),
1124 Option::<NoOpYielder>::None,
1125 )
1126 .await
1127 .expect("compact failed");
1128 }
1129 tree.set_layers(layers_from_handles([handle]).await.expect("layers_from_handles failed"));
1130 let found_item = tree.find(&item.key).await.expect("find failed").expect("not found");
1131 assert_eq!(found_item, item);
1132 assert!(!tree.exists(&TestKey(2..2)).await.expect("find failed"));
1133 }
1134
1135 #[fuchsia::test]
1136 async fn test_filter() {
1137 let items = [
1138 Item::new(TestKey(1..1), 1),
1139 Item::new(TestKey(2..2), 2),
1140 Item::new(TestKey(3..3), 3),
1141 Item::new(TestKey(4..4), 4),
1142 ];
1143 let tree = LSMTree::new(emit_left_merge_fn, None);
1144 tree.insert(items[0].clone()).expect("insert error");
1145 tree.insert(items[1].clone()).expect("insert error");
1146 tree.insert(items[2].clone()).expect("insert error");
1147 tree.insert(items[3].clone()).expect("insert error");
1148
1149 let layers = tree.layer_set();
1150 let mut merger = layers.merger();
1151
1152 let mut iter = merger
1154 .query(Query::FullScan)
1155 .await
1156 .expect("seek failed")
1157 .filter(|item: ItemRef<'_, TestKey, u64>| item.key.0.start % 2 == 0)
1158 .await
1159 .expect("filter failed");
1160
1161 assert_eq!(iter.get(), Some(items[1].as_item_ref()));
1162 iter.advance().await.expect("advance failed");
1163 assert_eq!(iter.get(), Some(items[3].as_item_ref()));
1164 iter.advance().await.expect("advance failed");
1165 assert!(iter.get().is_none());
1166 }
1167
1168 #[fuchsia::test]
1169 async fn test_insert_order_agnostic() {
1170 let items = [
1171 Item::new(TestKey(1..1), 1),
1172 Item::new(TestKey(2..2), 2),
1173 Item::new(TestKey(3..3), 3),
1174 Item::new(TestKey(4..4), 4),
1175 Item::new(TestKey(5..5), 5),
1176 Item::new(TestKey(6..6), 6),
1177 ];
1178 let a = LSMTree::new(emit_left_merge_fn, None);
1179 for item in &items {
1180 a.insert(item.clone()).expect("insert error");
1181 }
1182 let b = LSMTree::new(emit_left_merge_fn, None);
1183 let mut shuffled = items.clone();
1184 shuffled.shuffle(&mut rng());
1185 for item in &shuffled {
1186 b.insert(item.clone()).expect("insert error");
1187 }
1188 let layers = a.layer_set();
1189 let mut merger = layers.merger();
1190 let mut iter_a = merger.query(Query::FullScan).await.expect("seek failed");
1191 let layers = b.layer_set();
1192 let mut merger = layers.merger();
1193 let mut iter_b = merger.query(Query::FullScan).await.expect("seek failed");
1194
1195 for item in items {
1196 assert_eq!(Some(item.as_item_ref()), iter_a.get());
1197 assert_eq!(Some(item.as_item_ref()), iter_b.get());
1198 iter_a.advance().await.expect("advance failed");
1199 iter_b.advance().await.expect("advance failed");
1200 }
1201 assert!(iter_a.get().is_none());
1202 assert!(iter_b.get().is_none());
1203 }
1204
1205 enum AuditCacheResult<V> {
1206 Value(V),
1207 NoCache,
1208 }
1209
1210 struct AuditCacheInner<V: Value> {
1211 lookups: u64,
1212 completions: u64,
1213 invalidations: u64,
1214 drops: u64,
1215 result: Option<AuditCacheResult<V>>,
1216 is_cacheable: bool,
1217 }
1218
1219 impl<V: Value> AuditCacheInner<V> {
1220 fn stats(&self) -> (u64, u64, u64, u64) {
1221 (self.lookups, self.completions, self.invalidations, self.drops)
1222 }
1223 }
1224
1225 struct AuditCache<V: Value> {
1226 inner: Arc<Mutex<AuditCacheInner<V>>>,
1227 }
1228
1229 impl<V: Value> AuditCache<V> {
1230 fn new() -> Self {
1231 Self {
1232 inner: Arc::new(Mutex::new(AuditCacheInner {
1233 lookups: 0,
1234 completions: 0,
1235 invalidations: 0,
1236 drops: 0,
1237 result: None,
1238 is_cacheable: true,
1239 })),
1240 }
1241 }
1242 }
1243
1244 struct AuditPlaceholder<V: Value> {
1245 inner: Arc<Mutex<AuditCacheInner<V>>>,
1246 completed: Mutex<bool>,
1247 }
1248
1249 impl<V: Value> ObjectCachePlaceholder<V> for AuditPlaceholder<V> {
1250 fn complete(self: Box<Self>, _: Option<&V>) {
1251 self.inner.lock().completions += 1;
1252 *self.completed.lock() = true;
1253 }
1254 }
1255
1256 impl<V: Value> Drop for AuditPlaceholder<V> {
1257 fn drop(&mut self) {
1258 if !*self.completed.lock() {
1259 self.inner.lock().drops += 1;
1260 }
1261 }
1262 }
1263
1264 impl<K: Key + std::cmp::PartialEq, V: Value> ObjectCache<K, V> for AuditCache<V> {
1265 fn lookup_or_reserve(&self, _key: &K) -> ObjectCacheResult<'_, V> {
1266 {
1267 let mut inner = self.inner.lock();
1268 inner.lookups += 1;
1269 match MutexGuard::try_map_or_err(inner, |inner| match &mut inner.result {
1270 Some(AuditCacheResult::Value(value)) => Ok(value),
1271 Some(AuditCacheResult::NoCache) => Err(false),
1272 None => Err(true),
1273 }) {
1274 Ok(guard) => {
1275 return ObjectCacheResult::Value(guard);
1276 }
1277 Err((_, false)) => {
1278 return ObjectCacheResult::NoCache;
1279 }
1280 Err((_, true)) => {}
1281 }
1282 }
1283 ObjectCacheResult::Placeholder(Box::new(AuditPlaceholder {
1284 inner: self.inner.clone(),
1285 completed: Mutex::new(false),
1286 }))
1287 }
1288
1289 fn is_cacheable(&self, _key: &K) -> bool {
1290 self.inner.lock().is_cacheable
1291 }
1292
1293 fn invalidate(&self, _key: &K, _value: Option<V>) {
1294 self.inner.lock().invalidations += 1;
1295 }
1296
1297 fn clear(&self) {
1298 self.inner.lock().result = None;
1299 }
1300 }
1301
1302 #[fuchsia::test]
1303 async fn test_cache_handling() {
1304 let item = Item::new(TestKey(1..1), 1);
1305 let cache = Box::new(AuditCache::new());
1306 let inner = cache.inner.clone();
1307 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1308
1309 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1311
1312 assert!(!a.exists(&item.key).await.expect("Failed find"));
1314 assert_eq!(inner.lock().stats(), (1, 0, 0, 1));
1315
1316 let _ = a.insert(item.clone());
1318 assert_eq!(inner.lock().stats(), (1, 0, 1, 1));
1319
1320 assert_eq!(
1322 a.find(&item.key).await.expect("Failed find").expect("Item should be found.").value,
1323 item.value
1324 );
1325 assert_eq!(inner.lock().stats(), (2, 1, 1, 1));
1326
1327 let item2 = Item::new(TestKey(2..2), 2);
1329 let _ = a.insert(item2.clone());
1330 assert_eq!(a.find_value(&item2.key).await.expect("Failed find_value"), Some(item2.value));
1331 assert_eq!(inner.lock().stats(), (3, 2, 2, 1));
1332
1333 let item3 = Item::new(TestKey(3..3), 3);
1335 let _ = a.insert(item3.clone());
1336 assert!(a.exists(&item3.key).await.expect("Failed exists"));
1337 assert_eq!(inner.lock().stats(), (4, 3, 3, 1));
1338
1339 let item4 = Item::new(TestKey(4..4), 4);
1341 let _ = a.insert(item4.clone());
1342 assert_eq!(
1343 a.find_map(&item4.key, |item| item.value * 2).await.expect("Failed find_map"),
1344 Some(8)
1345 );
1346 assert_eq!(inner.lock().stats(), (5, 4, 4, 1));
1347
1348 a.replace_or_insert(item.clone());
1350 assert_eq!(inner.lock().stats(), (5, 4, 5, 1));
1351 }
1352
1353 #[fuchsia::test]
1354 async fn test_cache_hit() {
1355 let item = Item::new(TestKey(1..1), 1);
1356 let cache = Box::new(AuditCache::new());
1357 let inner = cache.inner.clone();
1358 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1359
1360 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1362
1363 let _ = a.insert(item.clone());
1365 assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1366
1367 inner.lock().result = Some(AuditCacheResult::Value(item.value.clone()));
1369
1370 assert_eq!(
1372 a.find_value(&item.key).await.expect("Failed find").expect("Item should be found."),
1373 item.value
1374 );
1375 assert_eq!(inner.lock().stats(), (1, 0, 1, 0));
1376 }
1377
1378 #[fuchsia::test]
1379 async fn test_cache_hit_deleted_marker() {
1380 let item = Item::new(TestKey(1..1), 1);
1381 let cache = Box::new(AuditCache::new());
1382 let inner = cache.inner.clone();
1383 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1384
1385 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1387
1388 let _ = a.insert(item.clone());
1390 assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1391
1392 inner.lock().result = Some(AuditCacheResult::Value(u64::DELETED_MARKER));
1394
1395 assert_eq!(a.find_value(&item.key).await.expect("Failed find_value"), None);
1397 assert_eq!(inner.lock().stats(), (1, 0, 1, 0));
1398
1399 inner.lock().result = Some(AuditCacheResult::Value(u64::DELETED_MARKER));
1401 assert!(!a.exists(&item.key).await.expect("Failed exists"));
1402 assert_eq!(inner.lock().stats(), (2, 0, 1, 0));
1403
1404 inner.lock().result = Some(AuditCacheResult::Value(u64::DELETED_MARKER));
1406 assert!(a.find(&item.key).await.expect("Failed find").is_none());
1407 assert_eq!(inner.lock().stats(), (3, 0, 1, 0));
1408
1409 inner.lock().result = Some(AuditCacheResult::Value(u64::DELETED_MARKER));
1411 assert!(a.find_map(&item.key, |_| ()).await.expect("Failed find_map").is_none());
1412 assert_eq!(inner.lock().stats(), (4, 0, 1, 0));
1413 }
1414
1415 #[fuchsia::test]
1416 async fn test_cache_says_uncacheable() {
1417 let item = Item::new(TestKey(1..1), 1);
1418 let cache = Box::new(AuditCache::new());
1419 let inner = cache.inner.clone();
1420 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1421 let _ = a.insert(item.clone());
1422
1423 assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1425
1426 inner.lock().result = Some(AuditCacheResult::NoCache);
1428
1429 assert_eq!(
1431 a.find_value(&item.key).await.expect("Failed find").expect("Should find item"),
1432 item.value
1433 );
1434 assert_eq!(inner.lock().stats(), (1, 0, 1, 0));
1435 }
1436
1437 #[fuchsia::test]
1438 async fn test_uncacheable_key_skips_invalidation() {
1439 let item = Item::new(TestKey(1..1), 1);
1440 let cache = Box::new(AuditCache::new());
1441 cache.inner.lock().is_cacheable = false;
1442 let inner = cache.inner.clone();
1443 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1444
1445 a.insert(item.clone()).expect("insert failed");
1447 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1448
1449 a.replace_or_insert(item.clone());
1451 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1452
1453 a.merge_into(item.clone(), &TestKey(1..1));
1455 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1456 }
1457
1458 struct FailLayer {
1459 drop_event: Mutex<Option<Arc<DropEvent>>>,
1460 }
1461
1462 impl FailLayer {
1463 fn new() -> Self {
1464 Self { drop_event: Mutex::new(Some(Arc::new(DropEvent::new()))) }
1465 }
1466 }
1467
1468 #[async_trait]
1469 impl<K: Key, V: Value> Layer<K, V> for FailLayer {
1470 async fn seek(
1471 &self,
1472 _bound: std::ops::Bound<&K>,
1473 ) -> Result<BoxedLayerIterator<'_, K, V>, Error> {
1474 Err(anyhow!("Purposely failed seek"))
1475 }
1476
1477 fn lock(&self) -> Option<Arc<DropEvent>> {
1478 self.drop_event.lock().clone()
1479 }
1480
1481 fn len(&self) -> usize {
1482 0
1483 }
1484
1485 async fn close(&self) {
1486 let listener = match std::mem::replace(&mut (*self.drop_event.lock()), None) {
1487 Some(drop_event) => drop_event.listen(),
1488 None => return,
1489 };
1490 listener.await;
1491 }
1492
1493 fn get_version(&self) -> Version {
1494 LATEST_VERSION
1495 }
1496
1497 async fn key_exists(&self, _key: &K) -> Result<Existence, Error> {
1498 unimplemented!();
1499 }
1500 }
1501
1502 struct MockLayer {
1503 exists_result: Existence,
1504 drop_event: Mutex<Option<Arc<DropEvent>>>,
1505 cleared: std::sync::atomic::AtomicBool,
1506 }
1507
1508 impl MockLayer {
1509 fn new(exists_result: Existence) -> Self {
1510 Self {
1511 exists_result,
1512 drop_event: Mutex::new(Some(Arc::new(DropEvent::new()))),
1513 cleared: std::sync::atomic::AtomicBool::new(false),
1514 }
1515 }
1516 }
1517
1518 #[async_trait]
1519 impl<K: Key, V: Value> Layer<K, V> for MockLayer {
1520 async fn seek(
1521 &self,
1522 _bound: std::ops::Bound<&K>,
1523 ) -> Result<BoxedLayerIterator<'_, K, V>, Error> {
1524 unimplemented!()
1525 }
1526
1527 fn clear_cached_data(&self) {
1528 self.cleared.store(true, Ordering::Relaxed);
1529 }
1530
1531 fn lock(&self) -> Option<Arc<DropEvent>> {
1532 self.drop_event.lock().clone()
1533 }
1534
1535 fn len(&self) -> usize {
1536 0
1537 }
1538
1539 async fn close(&self) {
1540 let listener = match std::mem::replace(&mut (*self.drop_event.lock()), None) {
1541 Some(drop_event) => drop_event.listen(),
1542 None => return,
1543 };
1544 listener.await;
1545 }
1546
1547 fn get_version(&self) -> Version {
1548 LATEST_VERSION
1549 }
1550
1551 async fn key_exists(&self, _key: &K) -> Result<Existence, Error> {
1552 Ok(self.exists_result)
1553 }
1554 }
1555
1556 #[fuchsia::test]
1557 async fn test_clear_cache() {
1558 let item = Item::new(TestKey(1..1), 1);
1559 let cache = Box::new(AuditCache::new());
1560 let inner = cache.inner.clone();
1561 let tree = LSMTree::new(emit_left_merge_fn, Some(cache));
1562
1563 let layer1 = Arc::new(MockLayer::new(Existence::MaybeExists));
1564 let layer2 = Arc::new(MockLayer::new(Existence::MaybeExists));
1565 tree.set_layers(vec![
1566 layer1.clone() as Arc<dyn Layer<TestKey, u64>>,
1567 layer2.clone() as Arc<dyn Layer<TestKey, u64>>,
1568 ]);
1569
1570 inner.lock().result = Some(AuditCacheResult::Value(item.value));
1571 assert_eq!(tree.find_value(&item.key).await.expect("find_value failed"), Some(item.value));
1572 assert!(!layer1.cleared.load(Ordering::Relaxed));
1573 assert!(!layer2.cleared.load(Ordering::Relaxed));
1574
1575 tree.clear_cache();
1576
1577 assert!(inner.lock().result.is_none());
1578 assert!(layer1.cleared.load(Ordering::Relaxed));
1579 assert!(layer2.cleared.load(Ordering::Relaxed));
1580 }
1581
1582 #[fuchsia::test]
1583 async fn test_layer_set_key_exists() {
1584 use super::LockedLayer;
1585
1586 let tree = LSMTree::new(emit_left_merge_fn, None);
1587 let mut layer_set = tree.empty_layer_set();
1588
1589 assert_eq!(
1591 layer_set.key_exists(&TestKey(0..1)).await.expect("key_exists failed"),
1592 Existence::Missing
1593 );
1594
1595 layer_set.layers.push(LockedLayer::from(
1597 Arc::new(MockLayer::new(Existence::Missing)) as Arc<dyn Layer<TestKey, u64>>
1598 ));
1599 assert_eq!(
1600 layer_set.key_exists(&TestKey(0..1)).await.expect("key_exists failed"),
1601 Existence::Missing
1602 );
1603
1604 layer_set.layers.push(LockedLayer::from(
1606 Arc::new(MockLayer::new(Existence::MaybeExists)) as Arc<dyn Layer<TestKey, u64>>
1607 ));
1608 assert_eq!(
1609 layer_set.key_exists(&TestKey(0..1)).await.expect("key_exists failed"),
1610 Existence::MaybeExists
1611 );
1612
1613 layer_set.layers.insert(
1615 0,
1616 LockedLayer::from(
1617 Arc::new(MockLayer::new(Existence::Exists)) as Arc<dyn Layer<TestKey, u64>>
1618 ),
1619 );
1620 assert_eq!(
1621 layer_set.key_exists(&TestKey(0..1)).await.expect("key_exists failed"),
1622 Existence::Exists
1623 );
1624 }
1625
1626 #[fuchsia::test]
1627 async fn test_failed_lookup() {
1628 let cache = Box::new(AuditCache::new());
1629 let inner = cache.inner.clone();
1630 let a = LSMTree::new(emit_left_merge_fn, Some(cache));
1631 a.set_layers(vec![Arc::new(FailLayer::new())]);
1632
1633 assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1635
1636 assert!(a.find(&TestKey(1..1)).await.is_err());
1638 assert_eq!(inner.lock().stats(), (1, 0, 0, 1));
1639 }
1640}
1641
1642#[cfg(fuzz)]
1643mod fuzz {
1644 use crate::lsm_tree::types::{Item, Value};
1645 use crate::serialized_types::{
1646 LATEST_VERSION, Version, Versioned, VersionedLatest, versioned_type,
1647 };
1648 use arbitrary::Arbitrary;
1649
1650 use fuzz::fuzz;
1651
1652 use super::testing::TestKey;
1653
1654 impl Versioned for u64 {}
1655 versioned_type! { 1.. => u64 }
1656
1657 impl Value for u64 {
1658 const DELETED_MARKER: Self = 0;
1659 }
1660
1661 #[allow(dead_code)]
1665 #[derive(Arbitrary)]
1666 enum FuzzAction {
1667 Insert(Item<TestKey, u64>),
1668 ReplaceOrInsert(Item<TestKey, u64>),
1669 MergeInto(Item<TestKey, u64>, TestKey),
1670 Find(TestKey),
1671 Seal,
1672 }
1673
1674 #[fuzz]
1675 fn fuzz_lsm_tree_actions(actions: Vec<FuzzAction>) {
1676 use super::LSMTree;
1677 use crate::lsm_tree::merge::{MergeLayerIterator, MergeResult};
1678 use futures::executor::block_on;
1679
1680 fn emit_left_merge_fn(
1681 _left: &MergeLayerIterator<'_, TestKey, u64>,
1682 _right: &MergeLayerIterator<'_, TestKey, u64>,
1683 ) -> MergeResult<TestKey, u64> {
1684 MergeResult::EmitLeft
1685 }
1686
1687 let tree = LSMTree::new(emit_left_merge_fn, None);
1688 for action in actions {
1689 match action {
1690 FuzzAction::Insert(item) => {
1691 let _ = tree.insert(item);
1692 }
1693 FuzzAction::ReplaceOrInsert(item) => {
1694 tree.replace_or_insert(item);
1695 }
1696 FuzzAction::Find(key) => {
1697 block_on(tree.find(&key)).expect("find failed");
1698 }
1699 FuzzAction::MergeInto(item, bound) => tree.merge_into(item, &bound),
1700 FuzzAction::Seal => tree.seal(),
1701 };
1702 }
1703 }
1704}