Skip to main content

fxfs/
lsm_tree.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5//! # End-Key Indexing for File Extents
6//!
7//! Fxfs indexes file extents by their **end offset** because it avoids expensive reverse
8//! iteration (`Prev()`) during lookups in the LSM tree.
9//!
10//! If we indexed by start offset, a `Seek(>= X)` would return the extent starting *after* X,
11//! requiring a `Prev()` to find the extent actually containing X.
12//!
13//! By indexing by end offset, a `Seek(>= X + 1)` returns the first extent ending *after* X.
14//! Since extents do not overlap, this is the only extent that could possibly contain X,
15//! eliminating the need for backtracking.
16
17mod 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
56// For serialization.
57pub 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
90/// Metrics related to LSM tree churn, layer depths, and compaction performance.
91pub struct CompactionCounters {
92    /// Total number of compaction events that merged layers together.
93    pub compactions: u64,
94    /// Total bytes written during compaction, useful for measuring write amplification.
95    pub compaction_bytes_written: u64,
96    /// Total duration spent compacting layers.
97    pub compaction_time_ns: u64,
98    /// Number of mutable layers sealed. Useful for measuring churn.
99    pub total_layers_added: u64,
100    /// The maximum depth of the LSM tree observed. An indicator of worst-case read amplification.
101    pub max_layer_count: u64,
102    /// Log2 histogram of the sizes of newly compacted layers, used to verify compaction heuristics.
103    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
119/// Global counters and metrics for an LSM tree's lifetime.
120pub struct TreeCounters {
121    /// Number of individual key-lookup attempts (reads) through the tree.
122    pub num_seeks: AtomicUsize,
123    /// Tracks the number of layer files we might have looked at across all seeks.
124    /// Used alongside `layer_files_skipped` to compute the effectiveness of bloom filters.
125    pub layer_files_total: AtomicUsize,
126    /// Tracks how many layer files we skipped searching thanks to the bloom filter rejecting them.
127    pub layer_files_skipped: AtomicUsize,
128    /// Tracks the number of queries that had an excessive number of fuzzy hash partitions.
129    pub excessive_hash_partitions: AtomicUsize,
130    /// Embedded counters for mutable metrics that require locking.
131    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/// Writes the items yielded by the iterator into the supplied object.
147#[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
167/// LSMTree manages a tree of layers to provide a key/value store.  Each layer contains deltas on
168/// the preceding layer.  The top layer is an in-memory mutable layer.  Layers can be compacted to
169/// form a new combined layer.
170pub 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    /// Creates a new empty tree.
180    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    /// Opens an existing tree from the provided handles to the layer objects.
196    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    /// Replaces the immutable layers.
218    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    /// Opens layers for `object_ids` from `store` and appends them to the end of the tree (i.e. as
227    /// base layers). This is supposed to be used after replay when we are opening a tree and we
228    /// have discovered the base layers.
229    ///
230    /// Returns the sum of the sizes in bytes of all appended layers.
231    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    /// Appends already-opened layers to the end of the tree.
243    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    /// Resets the immutable layers.
252    pub fn reset_immutable_layers(&self) {
253        self.data.write().layers = Vec::new();
254    }
255
256    /// Seals the current mutable layer and creates a new one.
257    pub fn seal(&self) {
258        // We need to be sure there are no mutations currently in-progress.  This is currently
259        // guaranteed by ensuring that all mutations take a read lock on `data`.
260        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    /// Resets the tree to an empty state.
270    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    /// Clears all in-memory caches (object cache and persistent layer chunk caches) for this tree.
277    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    /// Returns an empty layer-set for this tree.
316    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    /// Adds all the layers (including the mutable layer) to `layer_set`.
321    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    /// Returns a clone of the current set of layers (including the mutable layer), after which one
333    /// can get an iterator.
334    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    /// Returns the current set of immutable layers after which one can get an iterator (for e.g.
341    /// compacting).  Since these layers are immutable, getting an iterator should not block
342    /// anything else.
343    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    /// Inserts an item into the mutable layer.
353    /// Returns error if item already exists.
354    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            // `seal` below relies on us holding a read lock whilst we do the mutation.
361            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    /// Replaces or inserts an item into the mutable layer.
376    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            // `seal` below relies on us holding a read lock whilst we do the mutation.
384            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    /// Merges the given item into the mutable layer.
398    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            // `seal` below relies on us holding a read lock whilst we do the mutation.
404            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    /// Returns true if `key` is cacheable by the underlying cache.
417    ///
418    /// Returns false if no cache is configured or if the key is not cacheable.
419    /// This allows tree mutation operations (`insert`, `replace_or_insert`, `merge_into`) to skip
420    /// cloning keys and values when cache invalidation would be a no-op.
421    #[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    /// Invalidates the cache entry for `key` if a cache is configured.
427    ///
428    /// If `value` is Some, the cache entry may be updated with the new value. If `value` is None,
429    /// the entry is removed from the cache. Does nothing if no cache is configured.
430    #[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    /// Searches for an exact match for the given key, applying `f` to the found [`ItemRef`].
438    /// If the item does not exist or has a value equal to [`Value::DELETED_MARKER`], returns
439    /// `Ok(None)`.
440    ///
441    /// Range keys (`FuzzyHash::is_range_key`) are not supported.
442    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        // It is important that the cache lookup is done prior to fetching the layer set as the
449        // placeholder returned acts as a sort of lock for the validity of the item that may be
450        // inserted later via that placeholder.
451        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    /// Searches for an exact match for the given key. If the item does not exist or has a value
479    /// equal to [`Value::DELETED_MARKER`], returns `Ok(None)`.
480    ///
481    /// Range keys (`FuzzyHash::is_range_key`) are not supported.
482    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    /// Searches for an exact match for the given key, returning only the value. If the item does
493    /// not exist or has a value equal to [`Value::DELETED_MARKER`], returns `Ok(None)`.
494    ///
495    /// Range keys (`FuzzyHash::is_range_key`) are not supported.
496    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    /// Returns true if an exact match exists for the given key and is not deleted.
507    ///
508    /// Range keys (`FuzzyHash::is_range_key`) are not supported.
509    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    /// Sets a mutation callback which is a callback that is triggered whenever any mutations are
521    /// applied to the tree.  This might be useful for tests that want to record the precise
522    /// sequence of mutations that are applied to the tree.
523    pub fn set_mutation_callback(&self, mutation_callback: MutationCallback<K, V>) {
524        self.data.write().mutation_callback = mutation_callback;
525    }
526
527    /// Returns the earliest version used by a layer in the tree.
528    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    /// Returns a new mutable layer.
540    pub fn new_mutable_layer() -> Arc<SkipListLayer<K, V>> {
541        SkipListLayer::new(SKIP_LIST_LAYER_ITEMS)
542    }
543
544    /// Replaces the mutable layer.
545    pub fn set_mutable_layer(&self, layer: Arc<SkipListLayer<K, V>>) {
546        self.data.write().mutable_layer = layer;
547    }
548
549    /// Records inspect data for the LSM tree into `node`.  Called lazily when inspect is queried.
550    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
600/// This is an RAII wrapper for a layer which holds a lock on the layer (via the Layer::lock
601/// method).
602pub 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
633/// A LayerSet provides a snapshot of the layers at a particular point in time, and allows you to
634/// get an iterator.  Iterators borrow the layers so something needs to hold reference count.
635pub 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    /// Searches for an exact match for the given key, applying `f` to the found [`ItemRef`].
659    /// If the item does not exist or has a value equal to [`Value::DELETED_MARKER`], returns
660    /// `Ok(None)`.
661    ///
662    /// Range keys (`FuzzyHash::is_range_key`) are not supported.
663    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    /// See `Layer::key_exists`.
702    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
727/// A yielder can be used during compactions which are low priority.
728pub 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        // Missing key
945        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        // Deleted key
951        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        // Base layer:
962        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        // Top layer:
969        // Key 1 has an update (+50) that should merge to 150.
970        // Key 4 is new (40).
971        tree.insert(Item::new(1, 50)).expect("insert error");
972        tree.insert(Item::new(4, 40)).expect("insert error");
973
974        // Key 1 should be merged across layers: 50 + 100 = 150.
975        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        // Key 2 only exists in the base layer: 200.
984        assert_eq!(tree.find_value(&2).await.expect("find_value failed"), Some(200));
985        assert!(tree.exists(&2).await.expect("exists failed"));
986
987        // Key 3 only exists in the base layer: 300.
988        assert_eq!(tree.find_value(&3).await.expect("find_value failed"), Some(300));
989        assert!(tree.exists(&3).await.expect("exists failed"));
990
991        // Key 4 only exists in the top layer: 40.
992        assert_eq!(tree.find_value(&4).await.expect("find_value failed"), Some(40));
993        assert!(tree.exists(&4).await.expect("exists failed"));
994
995        // Key 5 is missing entirely.
996        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        // Insert items in base layer (older).
1005        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        // In the newer layer:
1012        // - Update key 2 to 25
1013        // - Tombstone key 3
1014        // - Insert new key 4
1015        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        // Key 1: only in older layer.
1020        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        // Key 2: updated in newer layer (shadows older layer).
1024        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        // Key 3: deleted in newer layer (shadows older layer).
1028        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        // Key 4: only in newer layer.
1034        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        // Key 5: missing in all layers.
1038        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        // Searching for any key: the bloom filter on the immutable layer reports False,
1094        // so seek() is never called (which would panic) and the layer is skipped.
1095        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        // 2 queries (find_value and exists); each inspects mutable layer + 1 immutable mock layer.
1102        // The mock layer is skipped by the bloom filter on both queries.
1103        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        // Filter out odd keys (which also guarantees we skip the first key which is an edge case).
1153        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        // Zero counters.
1310        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1311
1312        // Look for an item, but don't find it. So no insertion. It is dropped.
1313        assert!(!a.exists(&item.key).await.expect("Failed find"));
1314        assert_eq!(inner.lock().stats(), (1, 0, 0, 1));
1315
1316        // Insert attempts to invalidate.
1317        let _ = a.insert(item.clone());
1318        assert_eq!(inner.lock().stats(), (1, 0, 1, 1));
1319
1320        // Look for item, find it and insert into the cache.
1321        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        // find_value on item2 also inserts into cache
1328        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        // exists on item3 also inserts into cache
1334        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        // find_map on item4 also inserts into cache
1340        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        // Insert or replace attempts to invalidate as well.
1349        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        // Zero counters.
1361        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1362
1363        // Insert attempts to invalidate.
1364        let _ = a.insert(item.clone());
1365        assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1366
1367        // Set up the item to find in the cache.
1368        inner.lock().result = Some(AuditCacheResult::Value(item.value.clone()));
1369
1370        // Look for item, find it in cache, so no insert.
1371        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        // Zero counters.
1386        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1387
1388        // Insert attempts to invalidate.
1389        let _ = a.insert(item.clone());
1390        assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1391
1392        // Set up the cache to return DELETED_MARKER.
1393        inner.lock().result = Some(AuditCacheResult::Value(u64::DELETED_MARKER));
1394
1395        // find_value should hit the cache, see DELETED_MARKER, and return None.
1396        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        // exists should also return false on cache hit DELETED_MARKER.
1400        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        // find should also return None on cache hit DELETED_MARKER.
1405        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        // find_map should also return None on cache hit DELETED_MARKER.
1410        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        // One invalidation from the insert.
1424        assert_eq!(inner.lock().stats(), (0, 0, 1, 0));
1425
1426        // Set up the NoCache response to find in the cache.
1427        inner.lock().result = Some(AuditCacheResult::NoCache);
1428
1429        // Look for item, it is uncacheable, so no insert.
1430        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        // Insert of an uncacheable key should not invalidate the cache.
1446        a.insert(item.clone()).expect("insert failed");
1447        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1448
1449        // replace_or_insert of an uncacheable key should not invalidate the cache.
1450        a.replace_or_insert(item.clone());
1451        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1452
1453        // merge_into of an uncacheable key should not invalidate the cache.
1454        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        // Empty layer set should return Missing.
1590        assert_eq!(
1591            layer_set.key_exists(&TestKey(0..1)).await.expect("key_exists failed"),
1592            Existence::Missing
1593        );
1594
1595        // Add a layer that returns Missing.
1596        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        // Add a layer that returns MaybeExists.
1605        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        // Add a layer that returns Exists.
1614        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        // Zero counters.
1634        assert_eq!(inner.lock().stats(), (0, 0, 0, 0));
1635
1636        // Lookup should fail and drop the placeholder.
1637        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    // Note: This code isn't really dead. it's used below in
1662    // `fuzz_lsm_tree_action`. However, the `#[fuzz]` proc macro attribute
1663    // obfuscates the usage enough to confuse the compiler.
1664    #[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}