Skip to main content

fxfs/lsm_tree/
types.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
5use crate::drop_event::DropEvent;
6use crate::object_handle::ReadObjectHandle;
7use crate::serialized_types::serialized_key::SerializeKey;
8use crate::serialized_types::{Version, Versioned, VersionedLatest};
9use anyhow::Error;
10use async_trait::async_trait;
11use fprint::TypeFingerprint;
12use futures::future::BoxFuture;
13use serde::{Deserialize, Serialize};
14use std::fmt::Debug;
15use std::future::Future;
16use std::hash::Hash;
17use std::marker::PhantomData;
18use std::sync::Arc;
19
20pub use fxfs_macros::impl_fuzzy_hash;
21
22// Force keys to be sorted first by a u64, so that they can be located approximately based on only
23// that integer without the whole key.
24pub trait SortByU64: Sized {
25    // Return the u64 that is used as the first value when deciding on sort order of the key.
26    fn get_leading_u64(&self) -> u64;
27}
28
29/// An extension to `std::hash::Hash` to support values which should be partitioned and hashed into
30/// buckets, where nearby keys will have the same hash value.  This is used for existence filtering
31/// in layer files (see `Layer::maybe_contains_key`).
32///
33/// For point-based keys, this can be the same as `std::hash::Hash`, but for range-based keys, the
34/// hash can collapse nearby ranges into the same hash value.  Since a range-based key may span
35/// several buckets, `FuzzyHash::fuzzy_hash` must be called to split the key up into each of the
36/// possible values that it overlaps with.
37pub trait FuzzyHash: Hash + Sized {
38    /// To support range-based keys, multiple hash values may need to be checked for a given key.
39    /// For example, an extent [0..1024) might return extents [0..512), [512..1024), each of which
40    /// will have a unique return value for `Self::hash`.  For point-based keys, a single hash
41    /// suffices, in which case None is returned and the hash value of `self` should be checked.
42    /// Note that in general only a small number of partitions (e.g. 2) should be checked at once.
43    /// Queries checking too many partitions will fall back to returning true from bloom filter
44    /// checks to avoid degenerate performance.
45    fn fuzzy_hash(&self) -> impl ExactSizeIterator<Item = u64>;
46
47    /// Returns whether the type is a range-based key. Used to prevent use of range-based keys as
48    /// a point query (see [`crate::lsm_tree::merge::Query::Point`]).
49    fn is_range_key(&self) -> bool {
50        false
51    }
52}
53
54impl_fuzzy_hash!(u8);
55impl_fuzzy_hash!(u32);
56impl_fuzzy_hash!(u64);
57impl_fuzzy_hash!(String);
58impl_fuzzy_hash!(Vec<u8>);
59
60/// Keys and values need to implement the following traits.  For merging, they need to implement
61/// MergeableKey.  TODO: Use trait_alias when available.
62pub trait Key:
63    Clone
64    + Debug
65    + Hash
66    + FuzzyHash
67    + OrdUpperBound
68    + Send
69    + SortByU64
70    + Sync
71    + Versioned
72    + VersionedLatest
73    + SerializeKey
74    + std::marker::Unpin
75    + 'static
76{
77}
78
79impl<K> Key for K where
80    K: Clone
81        + Debug
82        + Hash
83        + FuzzyHash
84        + OrdUpperBound
85        + Send
86        + SortByU64
87        + Sync
88        + Versioned
89        + VersionedLatest
90        + SerializeKey
91        + std::marker::Unpin
92        + 'static
93{
94}
95
96pub trait MergeableKey: Key + Eq + LayerKey + OrdLowerBound {}
97impl<K> MergeableKey for K where K: Key + Eq + LayerKey + OrdLowerBound {}
98
99/// Trait required for supporting Layer functionality.
100pub trait LayerValue:
101    Clone + Send + Sync + Versioned + VersionedLatest + Debug + std::marker::Unpin + 'static
102{
103}
104impl<V> LayerValue for V where
105    V: Clone + Send + Sync + Versioned + VersionedLatest + Debug + std::marker::Unpin + 'static
106{
107}
108
109/// Superset of `LayerValue` to additionally support tree searching, requires comparison and an
110/// `DELETED_MARKER` for indicating empty values used to indicate deletion in the `LSMTree`.
111pub trait Value: PartialEq + LayerValue {
112    /// Value used to represent that the entry is actually empty, and should be ignored.
113    const DELETED_MARKER: Self;
114}
115
116/// ItemRef is a struct that contains references to key and value, which is useful since in many
117/// cases since keys and values are stored separately so &Item is not possible.
118#[derive(Debug, PartialEq, Eq, Serialize)]
119pub struct ItemRef<'a, K, V> {
120    pub key: &'a K,
121    pub value: &'a V,
122}
123
124impl<K: Clone, V: Clone> ItemRef<'_, K, V> {
125    pub fn cloned(&self) -> Item<K, V> {
126        Item { key: self.key.clone(), value: self.value.clone() }
127    }
128
129    pub fn boxed(&self) -> BoxedItem<K, V> {
130        Box::new(self.cloned())
131    }
132}
133
134impl<'a, K, V> Clone for ItemRef<'a, K, V> {
135    fn clone(&self) -> Self {
136        *self
137    }
138}
139impl<'a, K, V> Copy for ItemRef<'a, K, V> {}
140
141/// Item is a struct that combines a key and a value.
142#[derive(Copy, Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
143#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
144pub struct Item<K, V> {
145    pub key: K,
146    pub value: V,
147}
148
149#[cfg_attr(fuzz, derive(arbitrary::Arbitrary))]
150#[derive(Copy, Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
151pub struct LegacyItem<K, V> {
152    pub key: K,
153    pub value: V,
154    pub sequence: u64,
155}
156
157impl<K: TypeFingerprint, V: TypeFingerprint> TypeFingerprint for LegacyItem<K, V> {
158    fn fingerprint() -> String {
159        "struct {key:".to_owned()
160            + &K::fingerprint()
161            + ",value:"
162            + &V::fingerprint()
163            + ",sequence:u64}"
164    }
165}
166
167pub type BoxedItem<K, V> = Box<Item<K, V>>;
168
169// Nb: type-fprint doesn't support generics yet.
170impl<K: TypeFingerprint, V: TypeFingerprint> TypeFingerprint for Item<K, V> {
171    fn fingerprint() -> String {
172        "struct {key:".to_owned() + &K::fingerprint() + ",value:" + &V::fingerprint() + "}"
173    }
174}
175
176impl<K, V> Item<K, V> {
177    pub fn new(key: K, value: V) -> Item<K, V> {
178        Item { key, value }
179    }
180
181    pub fn as_item_ref(&self) -> ItemRef<'_, K, V> {
182        self.into()
183    }
184
185    pub fn boxed(self) -> BoxedItem<K, V> {
186        Box::new(self)
187    }
188}
189
190impl<'a, K, V> From<&'a Item<K, V>> for ItemRef<'a, K, V> {
191    fn from(item: &'a Item<K, V>) -> ItemRef<'a, K, V> {
192        ItemRef { key: &item.key, value: &item.value }
193    }
194}
195
196/// The find functions will return items with keys that are greater-than or equal to the search key,
197/// so for keys that are like extents, the keys should sort (via OrdUpperBound) using the end
198/// of their ranges, and you should set the search key accordingly.
199///
200/// For example, let's say the tree holds extents 100..200, 200..250 and you want to perform a read
201/// for range 150..250, you should search for 0..151 which will first return the extent 100..200
202/// (and then the iterator can be advanced to 200..250 after). When merging, keys can overlap, so
203/// consider the case where we want to merge an extent with range 100..300 with an existing extent
204/// of 200..250. In that case, we want to treat the extent with range 100..300 as lower than the key
205/// 200..250 because we'll likely want to split the extents (e.g. perhaps we want 100..200,
206/// 200..250, 250..300), so for merging, we need to use a different comparison function and we deal
207/// with that using the OrdLowerBound trait.
208///
209/// If your keys don't have overlapping ranges that need to be merged, then these can be the same as
210/// std::cmp::Ord (use the DefaultOrdUpperBound and DefaultOrdLowerBound traits).
211
212/// Trait for ordering keys by their upper bound (e.g. `end` for range-based keys).
213///
214/// This ordering is used within layer files and for positioning iterators via `search_key()`.
215pub trait OrdUpperBound {
216    fn cmp_upper_bound(&self, other: &Self) -> std::cmp::Ordering;
217}
218
219pub trait DefaultOrdUpperBound: OrdUpperBound + Ord {}
220
221impl<T: DefaultOrdUpperBound> OrdUpperBound for T {
222    fn cmp_upper_bound(&self, other: &Self) -> std::cmp::Ordering {
223        // Default to using cmp.
224        self.cmp(other)
225    }
226}
227
228/// Trait for ordering keys by their lower bound (e.g. `start` for range-based keys).
229///
230/// This ordering is used exclusively by the merger's min-heap to stream keys out in
231/// left-to-right order (by `start`). It is distinct from `OrdUpperBound` and is not used
232/// with search keys.
233pub trait OrdLowerBound {
234    fn cmp_lower_bound(&self, other: &Self) -> std::cmp::Ordering;
235}
236
237pub trait DefaultOrdLowerBound: OrdLowerBound + Ord {}
238
239impl<T: DefaultOrdLowerBound> OrdLowerBound for T {
240    fn cmp_lower_bound(&self, other: &Self) -> std::cmp::Ordering {
241        // Default to using cmp.
242        self.cmp(other)
243    }
244}
245
246/// Result returned by `merge_type()` to determine how to properly merge values within a layerset.
247#[derive(Clone, PartialEq)]
248pub enum MergeType {
249    /// Always includes every layer in the merger, when seeking or advancing. Always correct, but
250    /// always as slow as possible.
251    FullMerge,
252
253    /// Stops seeking older layers when an exact key match is found in a newer one. Useful for keys
254    /// that only replace data, or with `next_key()` implementations to decide on continued merging.
255    OptimizedMerge,
256}
257
258/// Determines how to iterate forward from the current key, and how many older layers to include
259/// when merging. See the different variants of `MergeKeyType` for more details.
260pub trait LayerKey: Clone {
261    /// Called to determine how to perform merge behaviours while advancing through a layer set.
262    fn merge_type(&self) -> MergeType {
263        // Defaults to full merge. The slowest, but most predictable in behaviour.
264        MergeType::FullMerge
265    }
266
267    /// The next_key() call allows for an optimisation which allows the merger to avoid querying a
268    /// layer if it knows it has found the next possible key.  It only makes sense for this to
269    /// return Some() when merge_type() returns OptimizedMerge. Consider the following example
270    /// showing two layers with range based keys.
271    ///
272    ///      +----------+------------+
273    ///  0   |  0..100  |  100..200  |
274    ///      +----------+------------+
275    ///  1              |  100..200  |
276    ///                 +------------+
277    ///
278    /// If you search and find the 0..100 key, then only layer 0 will be touched.  If you then want
279    /// to advance to the 100..200 record, you can find it in layer 0 but unless you know that it
280    /// immediately follows the 0..100 key i.e. that there is no possible key, K, such that 0..100 <
281    /// K < 100..200, the merger has to consult all other layers to check.  `next_key` should return
282    /// a search key for an extent that immediately follows.  In practice, for extents, this should
283    /// be `end..end + 1`.
284    ///
285    /// This is purely an optimisation; the default None will be correct but not performant.
286    fn next_key(&self) -> Option<Self> {
287        None
288    }
289    /// Returns the search key (S) for this key (K), such that when searching for S in a layer
290    /// file, it returns the earliest possible key that might be relevant to K.  Searching in a
291    /// layer file is done using `cmp_upper_bound` and the iterator will be positioned on a key that
292    /// is greater than or equal to S.  Returning `None` here is the right thing to do for
293    /// non-ranged based keys, in which case K is used to search for the key.  In practice, the
294    /// implementation should return `Some(start..start + 1)` for range based keys and `None`
295    /// for everything else.  As an example, if the tree has extents 50..150 and 150..200 and we
296    /// wish to search for 100..200, search_key would return 100..101 which would position the
297    /// iterator on 50..150.  If this method is overridden, `is_search_key` below should also be
298    /// overridden.
299    fn search_key(&self) -> Option<Self> {
300        None
301    }
302
303    /// If you override `search_key` you should override `is_search_key`.
304    fn is_search_key(&self) -> bool {
305        true
306    }
307
308    /// Returns true if the keys are overlapping or equal.
309    fn overlaps(&self, other: &Self) -> bool;
310}
311
312#[derive(Debug, Eq, PartialEq, Clone, Copy)]
313pub enum Existence {
314    /// The key definitely exists.
315    Exists,
316    /// The key might exist.
317    MaybeExists,
318    /// The key definitely does not exist.
319    Missing,
320}
321
322#[derive(Debug, Eq, PartialEq, Clone, Copy)]
323pub enum MaybeContainsKey {
324    /// The layer definitely does not contain records relevant to the key.
325    False,
326    /// The layer might contain records relevant to the key.
327    Maybe,
328    /// The range key was too large to check against the existence filter, so the check was skipped.
329    RangeKeyTooLarge,
330}
331
332/// Layer is a trait that all layers need to implement (mutable and immutable).
333#[async_trait]
334pub trait Layer<K, V>: Send + Sync {
335    /// If the layer is persistent, returns the handle to its contents.  Returns None for in-memory
336    /// layers.
337    fn handle(&self) -> Option<&dyn ReadObjectHandle> {
338        None
339    }
340
341    /// Some layer implementations may choose to cache data in-memory.  Calling this function will
342    /// request that the layer purges unused cached data.  This is intended to run on a timer.
343    fn purge_cached_data(&self) {}
344
345    /// Immediately clears any in-memory cached data for this layer.
346    fn clear_cached_data(&self) {}
347
348    /// Searches for a key. Bound::Excluded is not supported. Bound::Unbounded positions the
349    /// iterator on the first item in the layer.
350    async fn seek(&self, bound: std::ops::Bound<&K>)
351    -> Result<BoxedLayerIterator<'_, K, V>, Error>;
352
353    /// Returns the number of items in the layer file.
354    fn len(&self) -> usize;
355
356    /// Returns whether the layer *might* contain records relevant to `key`.  Note that this can
357    /// return `MaybeContainsKey::Maybe` even if the layer has no records relevant to `key`, but it will
358    /// never return `MaybeContainsKey::False` if there are such records.  (As such, always
359    /// returning `MaybeContainsKey::Maybe` is a trivially correct implementation.)
360    /// If the key has too many hash partitions to check against the existence filter, returns
361    /// `MaybeContainsKey::RangeKeyTooLarge`.
362    fn maybe_contains_key(&self, _key: &K) -> MaybeContainsKey {
363        MaybeContainsKey::Maybe
364    }
365
366    /// Whether the bloom filter for the layer file is consulted or not.  If this is false, then
367    /// `maybe_contains_key` will always return `MaybeContainsKey::Maybe`.
368    /// Note that the persistent layer file may still have a bloom filter, but it might be ignored
369    /// (e.g. for a layer file on an older version).
370    fn has_bloom_filter(&self) -> bool {
371        false
372    }
373
374    /// This is similar to `maybe_contains_key` except that there *must* be a `key` and possible
375    /// state where this will indicate the key is missing i.e. *always* returning
376    /// `Existence::MaybeExists` is *not* a correct implementation.  If an implementation has a low
377    /// cost way to determine if the key *might* exist, then it may use that, but if not, it must
378    /// use a slower algorithm.  This method was introduced to allow for an efficient way of
379    /// determining if a particular key is free to be used.  `maybe_contains_key` cannot be used for
380    /// this purpose because implementations might *always* return `true` which would mean it would
381    /// be impossible to ever find a key that is free to use.  It might not be appropriate to use
382    /// this with range based keys: implementations should use `cmp_upper_bound` and test for
383    /// equality, which might not give the desired results for range based keys.
384    async fn key_exists(&self, key: &K) -> Result<Existence, Error>;
385
386    /// Locks the layer preventing it from being closed. This will never block i.e. there can be
387    /// many locks concurrently.  The lock is purely advisory: seek will still work even if lock has
388    /// not been called; it merely causes close to wait until all locks are released.  Returns None
389    /// if close has been called for the layer.
390    fn lock(&self) -> Option<Arc<DropEvent>>;
391
392    /// Waits for existing locks readers to finish and then returns.  Subsequent calls to lock will
393    /// return None.
394    async fn close(&self);
395
396    /// Returns the version number used by structs in this layer
397    fn get_version(&self) -> Version;
398
399    /// Records inspect data for the layer into `node`.  Called lazily when inspect is queried.
400    fn record_inspect_data(self: Arc<Self>, _node: &fuchsia_inspect::Node) {}
401}
402
403/// Something that implements LayerIterator is returned by the seek function.
404pub trait LayerIterator<K, V>: Send + Sync {
405    /// Advances the iterator (static dispatch, unboxed future).
406    fn advance(&mut self) -> impl Future<Output = Result<(), Error>> + Send
407    where
408        Self: Sized;
409
410    /// Advances the iterator for dynamic dispatch (trait objects).
411    /// If the iterator advances synchronously, returns `Ok(None)`.
412    fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error>;
413
414    /// Returns the current item. This will be None if called when the iterator is first crated i.e.
415    /// before either seek or advance has been called, and None if the iterator has reached the end
416    /// of the layer.
417    fn get(&self) -> Option<ItemRef<'_, K, V>>;
418
419    /// Creates an iterator that only yields items from the underlying iterator for which
420    /// `predicate` returns `true`.
421    fn filter<P>(
422        self,
423        predicate: P,
424    ) -> impl Future<Output = Result<FilterLayerIterator<Self, P, K, V>, Error>> + Send
425    where
426        P: for<'b> Fn(ItemRef<'b, K, V>) -> bool + Send + Sync,
427        Self: Sized,
428        K: Send + Sync,
429        V: Send + Sync,
430    {
431        FilterLayerIterator::new(self, predicate)
432    }
433}
434
435pub type BoxedLayerIterator<'iter, K, V> = Box<dyn LayerIterator<K, V> + 'iter>;
436
437impl<'iter, K, V> LayerIterator<K, V> for BoxedLayerIterator<'iter, K, V> {
438    async fn advance(&mut self) -> Result<(), Error> {
439        if let Some(fut) = self.as_mut().advance_dyn()? {
440            fut.await?;
441        }
442        Ok(())
443    }
444
445    fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error> {
446        self.as_mut().advance_dyn()
447    }
448
449    fn get(&self) -> Option<ItemRef<'_, K, V>> {
450        self.as_ref().get()
451    }
452}
453
454/// Mutable layers need an iterator that implements this in order to make merge_into work.
455pub(super) trait LayerIteratorMut<K, V>: Sync {
456    /// Advances the iterator.
457    fn advance(&mut self);
458
459    /// Returns the current item. This will be None if called when the iterator is first crated i.e.
460    /// before either seek or advance has been called, and None if the iterator has reached the end
461    /// of the layer.
462    fn get(&self) -> Option<ItemRef<'_, K, V>>;
463
464    /// Inserts the item before the item that the iterator is located at.  The insert won't be
465    /// visible until the changes are committed (see `commit`).
466    fn insert(&mut self, item: Item<K, V>);
467
468    /// Erases the current item and positions the iterator on the next item, if any.  The change
469    /// won't be visible until committed (see `commit`).
470    fn erase(&mut self);
471
472    /// Commits the changes.  This does not wait for existing readers to finish.
473    fn commit(&mut self);
474}
475
476/// Trait for writing new layers.
477pub trait LayerWriter<K, V>: Sized
478where
479    K: Debug + Send + Versioned + Sync,
480    V: Debug + Send + Versioned + Sync,
481{
482    /// Writes the given item to this layer.
483    fn write(&mut self, item: ItemRef<'_, K, V>) -> impl Future<Output = Result<(), Error>> + Send;
484
485    /// Flushes any buffered items to the backing storage. The total number of bytes written is
486    /// returned.
487    fn complete(self) -> impl Future<Output = Result<u64, Error>> + Send;
488}
489
490/// A `LayerIterator`` that filters the items of another `LayerIterator`.
491pub struct FilterLayerIterator<I, P, K, V> {
492    iter: I,
493    predicate: P,
494    _key: PhantomData<K>,
495    _value: PhantomData<V>,
496}
497
498impl<I, P, K, V> FilterLayerIterator<I, P, K, V>
499where
500    I: LayerIterator<K, V>,
501    P: for<'b> Fn(ItemRef<'b, K, V>) -> bool + Send + Sync,
502{
503    async fn new(iter: I, predicate: P) -> Result<Self, Error> {
504        let mut filter = Self { iter, predicate, _key: PhantomData, _value: PhantomData };
505        filter.skip_filtered().await?;
506        Ok(filter)
507    }
508
509    async fn skip_filtered(&mut self) -> Result<(), Error> {
510        loop {
511            match self.iter.get() {
512                Some(item) if !(self.predicate)(item) => {}
513                _ => return Ok(()),
514            }
515            self.iter.advance().await?;
516        }
517    }
518}
519
520impl<I, P, K, V> LayerIterator<K, V> for FilterLayerIterator<I, P, K, V>
521where
522    I: LayerIterator<K, V>,
523    P: for<'b> Fn(ItemRef<'b, K, V>) -> bool + Send + Sync,
524    K: Send + Sync,
525    V: Send + Sync,
526{
527    async fn advance(&mut self) -> Result<(), Error> {
528        self.iter.advance().await?;
529        self.skip_filtered().await
530    }
531
532    fn advance_dyn<'a>(&'a mut self) -> Result<Option<BoxFuture<'a, Result<(), Error>>>, Error> {
533        Ok(Some(Box::pin(self.advance())))
534    }
535
536    fn get(&self) -> Option<ItemRef<'_, K, V>> {
537        self.iter.get()
538    }
539}