Skip to main content

fxfs/lsm_tree/
merge.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::log::*;
6use crate::lsm_tree;
7use crate::lsm_tree::types::{
8    BoxedItem, Item, ItemRef, Key, Layer, LayerIterator, LayerIteratorMut, LayerKey,
9    MaybeContainsKey, MergeType, OrdLowerBound, Value,
10};
11use anyhow::Error;
12use futures::future::BoxFuture;
13use futures::try_join;
14use std::cmp::Ordering;
15use std::collections::BinaryHeap;
16use std::fmt::{Debug, Write};
17use std::ops::{Bound, Deref, DerefMut};
18use std::sync::{Arc, atomic};
19
20#[derive(Debug, Eq, PartialEq)]
21pub enum ItemOp<K, V> {
22    /// Keeps the item to be presented to the merger subsequently with a new merge pair.
23    Keep,
24
25    /// Discards the item and moves on to the next item in the respective layer.
26    Discard,
27
28    /// Replaces the item with something new which will be presented to the merger subsequently with
29    /// a new pair.
30    Replace(BoxedItem<K, V>),
31}
32
33#[derive(Debug, Eq, PartialEq)]
34pub enum MergeResult<K, V> {
35    /// Emits the left item unchanged. Keeps the right item. This is the common case. Once an item
36    /// has been emitted, it will never be seen again by the merge function.
37    EmitLeft,
38
39    /// All other merge results are covered by the following. Take care when replacing items
40    /// that you replace the correct item. The merger will never merge two items together from
41    /// the same layer. Consider the following scenario:
42    ///
43    ///        +-----------+              +-----------+
44    /// 0:     |    A      |              |    C      |
45    ///        +-----------+--------------+-----------+
46    /// 1:                 |      B       |
47    ///                    +--------------+
48    ///
49    /// Let's say that all three items can be merged together. The merge function will first be
50    /// presented with items A and B, at which point it has the option of replacing the left item
51    /// (i.e. A, in layer 0) or the right item (i.e. B in layer 1). However, if you replace the left
52    /// item, the merge function will not then be given the opportunity to merge it with C, so the
53    /// correct thing to do in this case is to replace the right item B in layer 1, and discard the
54    /// left item. A rule you can use is that you should avoid replacing an item with another item
55    /// whose upper bound exceeds that of the item you are replacing.
56    ///
57    /// There are some combinations that might lead to infinite loops (e.g. None, Keep, Keep) and
58    /// should obviously be avoided.
59    Other { emit: Option<BoxedItem<K, V>>, left: ItemOp<K, V>, right: ItemOp<K, V> },
60}
61
62/// Users must provide a merge function which will take pairs of items, left and right, and return a
63/// merge result. The left item's key will either be less than the right item's key, or if they are
64/// the same, then the left item will be in a lower layer index (lower layer indexes indicate more
65/// recent entries). The last remaining item is always emitted.
66pub type MergeFn<K, V> =
67    fn(&MergeLayerIterator<'_, K, V>, &MergeLayerIterator<'_, K, V>) -> MergeResult<K, V>;
68
69pub enum MergeItem<K, V> {
70    None,
71    Item(BoxedItem<K, V>),
72    Iter,
73}
74
75enum RawIterator<'a, K, V> {
76    None,
77    Const(Box<dyn LayerIterator<K, V> + 'a>),
78    Mut(Box<dyn LayerIteratorMut<K, V> + 'a>),
79}
80
81// SAFETY: LayerIteratorMut is not Sync or Send, but it's only used by merge_into, which doesn't
82// send across threads (and is sync, not async).
83unsafe impl<K, V> Send for RawIterator<'_, K, V> {}
84
85// An iterator that keeps track of where we are for each of the layers. We push these onto a
86// min-heap.
87pub struct MergeLayerIterator<'a, K, V> {
88    layer: Option<&'a dyn Layer<K, V>>,
89
90    // The underlying iterator.
91    iter: RawIterator<'a, K, V>,
92
93    // The index of the layer this is for.
94    pub layer_index: u16,
95
96    // The item we are currently pointing at.
97    item: MergeItem<K, V>,
98}
99
100impl<'a, K, V> MergeLayerIterator<'a, K, V> {
101    pub fn key(&self) -> &K {
102        self.item().key
103    }
104
105    pub fn value(&self) -> &V {
106        self.item().value
107    }
108
109    fn new(layer_index: u16, layer: &'a dyn Layer<K, V>) -> Self {
110        MergeLayerIterator {
111            layer: Some(layer),
112            iter: RawIterator::None,
113            layer_index,
114            item: MergeItem::None,
115        }
116    }
117
118    fn new_with_item(layer_index: u16, item: MergeItem<K, V>) -> Self {
119        MergeLayerIterator { layer: None, iter: RawIterator::None, layer_index, item }
120    }
121
122    fn item(&self) -> ItemRef<'_, K, V> {
123        match &self.item {
124            MergeItem::None => panic!("No item!"),
125            MergeItem::Item(item) => item.as_item_ref(),
126            MergeItem::Iter => self.get().unwrap(),
127        }
128    }
129
130    fn get(&self) -> Option<ItemRef<'_, K, V>> {
131        match &self.iter {
132            RawIterator::None => panic!("No iterator!"),
133            RawIterator::Const(iter) => iter.get(),
134            RawIterator::Mut(iter) => iter.get(),
135        }
136    }
137
138    fn set_item_from_iter(&mut self) {
139        self.item = if match &self.iter {
140            RawIterator::None => unreachable!(),
141            RawIterator::Const(iter) => iter.get(),
142            RawIterator::Mut(iter) => iter.get(),
143        }
144        .is_some()
145        {
146            MergeItem::Iter
147        } else {
148            MergeItem::None
149        };
150    }
151
152    fn take_item(&mut self) -> Option<BoxedItem<K, V>> {
153        if matches!(self.item, MergeItem::Item(_)) {
154            if let MergeItem::Item(item) = std::mem::replace(&mut self.item, MergeItem::None) {
155                return Some(item);
156            }
157        }
158        None
159    }
160
161    async fn advance(&mut self) -> Result<(), Error> {
162        if let MergeItem::Iter = self.item {
163            if let RawIterator::Const(iter) = &mut self.iter {
164                iter.advance().await?;
165            } else {
166                // This will never get called in the RawIterator::Mut case.
167                unreachable!();
168            }
169        }
170        self.set_item_from_iter();
171        Ok(())
172    }
173
174    fn replace(&mut self, item: BoxedItem<K, V>) {
175        self.item = MergeItem::Item(item);
176    }
177
178    fn is_some(&self) -> bool {
179        !matches!(self.item, MergeItem::None)
180    }
181
182    // This function exists so that we can advance multiple iterators concurrently using, say,
183    // try_join!.
184    async fn maybe_advance(&mut self, op: &ItemOp<K, V>) -> Result<(), Error> {
185        if let ItemOp::Keep = op { Ok(()) } else { self.advance().await }
186    }
187}
188
189// -- Ord and friends --
190impl<K: OrdLowerBound, V> Ord for MergeLayerIterator<'_, K, V> {
191    fn cmp(&self, other: &Self) -> Ordering {
192        // Reverse ordering because we want min-heap not max-heap.
193        other.key().cmp_lower_bound(self.key()).then(other.layer_index.cmp(&self.layer_index))
194    }
195}
196impl<K: OrdLowerBound, V> PartialOrd for MergeLayerIterator<'_, K, V> {
197    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
198        Some(self.cmp(other))
199    }
200}
201impl<K: OrdLowerBound, V> PartialEq for MergeLayerIterator<'_, K, V> {
202    fn eq(&self, other: &Self) -> bool {
203        self.cmp(other) == Ordering::Equal
204    }
205}
206impl<K: OrdLowerBound, V> Eq for MergeLayerIterator<'_, K, V> {}
207
208// As we merge items, the current item can be an item that has been replaced (and later emitted) by
209// the merge function, or an item referenced by an iterator, or nothing.
210enum CurrentItem<'a, 'b, K, V> {
211    None,
212    Item(BoxedItem<K, V>),
213    Iterator(&'a mut MergeLayerIterator<'b, K, V>),
214}
215
216impl<'a, 'b, K, V> CurrentItem<'a, 'b, K, V> {
217    fn take_iterator(&mut self) -> Option<&'a mut MergeLayerIterator<'b, K, V>> {
218        if matches!(self, CurrentItem::Iterator(_)) {
219            if let CurrentItem::Iterator(iter) = std::mem::replace(self, CurrentItem::None) {
220                return Some(iter);
221            }
222        }
223        None
224    }
225}
226
227impl<'a, K, V> From<&'a CurrentItem<'_, '_, K, V>> for Option<ItemRef<'a, K, V>> {
228    fn from(iter: &'a CurrentItem<'_, '_, K, V>) -> Option<ItemRef<'a, K, V>> {
229        match iter {
230            CurrentItem::None => None,
231            CurrentItem::Iterator(iterator) => Some(iterator.item()),
232            CurrentItem::Item(item) => Some(item.as_item_ref()),
233        }
234    }
235}
236
237/// Merger is the main entry point to merging.
238pub struct Merger<'a, K, V> {
239    // A buffer containing all the MergeLayerIterator objects.
240    iterators: Vec<MergeLayerIterator<'a, K, V>>,
241
242    // The function to be used for merging items.
243    merge_fn: MergeFn<K, V>,
244
245    // If true, additional logging is enabled.
246    trace: bool,
247
248    // Tracks statistics related to the LSM tree.
249    counters: Arc<lsm_tree::TreeCounters>,
250}
251
252/// Query describes the goal of a search in the LSM tree.  The caller specifies this to guide the
253/// merger on which layer files it must consult.
254/// Layers might be skipped during a query for a few reasons:
255///   * Existence filters might allow layer files to be skipped.
256///   * For bounded queries (i.e. any except FullScan), if the search key has
257///     MergeType::OptimizedMerge, we can use that to omit older layers once we find a match.
258#[derive(Debug, Clone)]
259pub enum Query<'a, K: Key + LayerKey + OrdLowerBound> {
260    /// Point queries look for a specific key in the LSM tree.  In this case, the existence filters
261    /// for each layer file can be used to decide if the layer file needs to be consulted.
262    /// Note that it is an error to use `Point` for range-like keys.  Either `LimitedRange` or
263    /// `FullRange` must be used instead.
264    Point(&'a K),
265
266    /// LimitedRange queries allow for iteration over a range of keys.  The returned iterator will
267    /// not necessarily terminate at the end of the range; the caller should stop iterating when the
268    /// iterator passes the end of the range: the records returned cannot be relied upon as
269    /// accurate.  For a LimitedRange query, the existence filters for each layer file can be used,
270    /// but we have to check for all possible keys in the range we wish to search.  Obviously, that
271    /// means that the range should be, well, limited.  Fuzzy hashes permit this for extent-like
272    /// keys, but these queries are not the right choice for things like searching a directory.
273    LimitedRange(&'a K),
274
275    /// FullRange queries position the iterator to a starting key, and scans forward to the first
276    /// key of a different type.  In this case, the existence filters are not used.  The key should
277    /// be a search key (see `LayerKey::search_key`) (which, for extent based keys, would normally
278    /// be `start..start + 1`).  This kind of query will still use optimized merges if the keys
279    /// found support it.
280    FullRange(&'a K),
281
282    /// FullScan queries are intended to yield every record in the tree.  In this case, the
283    /// existence filters are not used.
284    FullScan,
285}
286
287impl<'a, K: Key + LayerKey + OrdLowerBound> Query<'a, K> {
288    fn check_layer<V>(&self, layer: &dyn Layer<K, V>) -> MaybeContainsKey {
289        match self {
290            Self::Point(key) => layer.maybe_contains_key(key),
291            Self::LimitedRange(key) => layer.maybe_contains_key(key),
292            Self::FullRange(_) => MaybeContainsKey::Maybe,
293            Self::FullScan => MaybeContainsKey::Maybe,
294        }
295    }
296}
297
298#[fxfs_trace::trace]
299impl<'a, K: Key + LayerKey + OrdLowerBound, V: Value> Merger<'a, K, V> {
300    pub(super) fn new<I: Iterator<Item = &'a dyn Layer<K, V>>>(
301        layers: I,
302        merge_fn: MergeFn<K, V>,
303        counters: Arc<lsm_tree::TreeCounters>,
304    ) -> Merger<'a, K, V> {
305        Merger {
306            iterators: layers
307                .enumerate()
308                .map(|(index, layer)| MergeLayerIterator::new(index as u16, layer))
309                .collect(),
310            merge_fn: merge_fn,
311            trace: false,
312            counters,
313        }
314    }
315
316    /// Executes `query`, positioning the iterator to the first matching item.  See `Query` for
317    /// details.
318    #[trace]
319    pub async fn query(
320        &mut self,
321        query: Query<'_, K>,
322    ) -> Result<MergerIterator<'_, 'a, K, V>, Error> {
323        if let Query::Point(key) = query {
324            // NB: It is almost certainly an error to provide a range-like key to a Point query,
325            // hence this debug assertion.  This will most likely result in the wrong result, since
326            // the starting bound would be the original range key and therefore we might miss
327            // records that start earlier but overlap with the key.
328            // For example, if there is a layer which contains 0..100 and 100..200, and we search
329            // for 50..150, we should get both extents, which requires a different starting bound
330            // than 50..150 (particularly it would require the `K::search_key`).
331            debug_assert!(!key.is_range_key())
332        };
333        let len = self.iterators.len();
334        let mut excessive_partitions = false;
335        let pending_iterators = {
336            fxfs_trace::duration!("Merger::filter_layer_files", "len" => len);
337            self.iterators
338                .iter_mut()
339                .rev()
340                .filter(|l| match query.check_layer(l.layer.unwrap()) {
341                    MaybeContainsKey::False => false,
342                    MaybeContainsKey::Maybe => true,
343                    MaybeContainsKey::RangeKeyTooLarge => {
344                        excessive_partitions = true;
345                        true
346                    }
347                })
348                .collect::<Vec<&mut MergeLayerIterator<'a, K, V>>>()
349        };
350        let layer_count = pending_iterators.len();
351        {
352            self.counters.num_seeks.fetch_add(1, atomic::Ordering::Relaxed);
353            self.counters.layer_files_total.fetch_add(len, atomic::Ordering::Relaxed);
354            self.counters
355                .layer_files_skipped
356                .fetch_add(len - layer_count, atomic::Ordering::Relaxed);
357            if excessive_partitions {
358                self.counters.excessive_hash_partitions.fetch_add(1, atomic::Ordering::Relaxed);
359            }
360        }
361        log::debug!(query:?; "Consulting {}/{} layers", layer_count, len);
362        let mut merger_iter = MergerIterator {
363            merge_fn: self.merge_fn,
364            pending_iterators,
365            heap: BinaryHeap::with_capacity(layer_count),
366            item: CurrentItem::None,
367            trace: self.trace,
368            history: String::new(),
369        };
370        let owned_key;
371        let search_key = match query {
372            Query::Point(key) => Bound::Included(key),
373            Query::LimitedRange(key) => match key.search_key() {
374                Some(k) => {
375                    owned_key = k;
376                    Bound::Included(&owned_key)
377                }
378                None => Bound::Included(key),
379            },
380            Query::FullRange(key) => {
381                assert!(key.is_search_key());
382                Bound::Included(key)
383            }
384            Query::FullScan => Bound::Unbounded,
385        };
386        merger_iter.seek(search_key).await?;
387        Ok(merger_iter)
388    }
389
390    pub fn set_trace(&mut self, v: bool) {
391        self.trace = v;
392    }
393}
394
395/// This is an iterator that will allow iteration over merged layers.  The primary interface is via
396/// the LayerIterator trait.
397pub struct MergerIterator<'a, 'b, K, V> {
398    merge_fn: MergeFn<K, V>,
399
400    // Iterators that we have not yet pushed onto the heap.
401    pending_iterators: Vec<&'a mut MergeLayerIterator<'b, K, V>>,
402
403    // A heap with the merge iterators.
404    heap: BinaryHeap<&'a mut MergeLayerIterator<'b, K, V>>,
405
406    // The current item.
407    item: CurrentItem<'a, 'b, K, V>,
408
409    // If true, logs regarding merger behaviour are appended to history.
410    trace: bool,
411
412    // Holds trace information if trace is true.
413    history: String,
414}
415
416impl<'a, 'b, K: Key + LayerKey + OrdLowerBound, V: Value> MergerIterator<'a, 'b, K, V> {
417    pub fn pending_iterators_len(&self) -> usize {
418        self.pending_iterators.len()
419    }
420
421    /// Positions the iterator using `search_key`.
422    async fn seek(&mut self, search_key: Bound<&K>) -> Result<(), Error> {
423        self.push_iterators(search_key).await?;
424        self.advance_impl(search_key).await
425    }
426
427    // Merges items from an array of layers using the provided merge function. The merge function is
428    // repeatedly provided the lowest and the second lowest element, if one exists. In cases where
429    // the two lowest elements compare equal, the element with the lowest layer (i.e. whichever
430    // comes first in the layers array) will come first.  `search_key` is a bound for the next
431    // key we expect.
432    async fn advance_impl(&mut self, search_key: Bound<&K>) -> Result<(), Error> {
433        loop {
434            loop {
435                if self.heap.is_empty() {
436                    self.item = CurrentItem::None;
437                    return Ok(());
438                }
439                let lowest = self.heap.pop().unwrap();
440                let maybe_second_lowest = self.heap.pop();
441                if let Some(second_lowest) = maybe_second_lowest {
442                    let result = (self.merge_fn)(&lowest, &second_lowest);
443                    if self.trace {
444                        writeln!(
445                            self.history,
446                            "merge {:?}, {:?} -> {:?}",
447                            lowest.item(),
448                            second_lowest.item(),
449                            result
450                        )
451                        .unwrap();
452                    }
453                    match result {
454                        MergeResult::EmitLeft => {
455                            self.heap.push(second_lowest);
456                            self.item = CurrentItem::Iterator(lowest);
457                            break;
458                        }
459                        MergeResult::Other { emit, left, right } => {
460                            try_join!(
461                                lowest.maybe_advance(&left),
462                                second_lowest.maybe_advance(&right)
463                            )?;
464                            self.update_item(lowest, left);
465                            self.update_item(second_lowest, right);
466                            if let Some(emit) = emit {
467                                self.item = CurrentItem::Item(emit);
468                                break;
469                            }
470                        }
471                    }
472                } else {
473                    self.item = CurrentItem::Iterator(lowest);
474                    break;
475                }
476            }
477
478            // If the item we're about to yield isn't within `search_key`, ignore it and
479            // continue.  To see how this would happen, imagine the following scenario:
480            //
481            //       0          10          20            30
482            //       +----------+
483            //   0   |          |
484            //       +----------+-----------+
485            //   1   |                      |
486            //       +----------------------+-------------+
487            //   2   |                                    |
488            //       +------------------------------------+
489            //
490            // If we are to seek for 0..1 and then iterate over what we find, we expect the sequence
491            // 0..10, 10..20, 20..30.  After the first seek, we should only consult iterator 0.  For
492            // the first advance, we will consult iterator 1 but we need to merge the 0..20 element
493            // with the 0..10 element which we already emitted.  To make this work, we push the
494            // 0..10 item back on the heap (see advance) and then merge, but that will yield the
495            // 0..10 entry again, so here we need to skip over it, and then merge again, at which
496            // point we should see the 10..20 entry as expected.
497            match search_key {
498                Bound::Included(key)
499                    if self.get().unwrap().key.cmp_upper_bound(key) == Ordering::Less => {}
500                Bound::Excluded(key)
501                    if self.get().unwrap().key.cmp_upper_bound(key) != Ordering::Greater => {}
502                _ => return Ok(()),
503            }
504            if let Some(iterator) = self.item.take_iterator() {
505                iterator.advance().await?;
506                if iterator.is_some() {
507                    self.heap.push(iterator);
508                }
509            }
510        }
511    }
512
513    // Returns whether more iterators are required for the given `search_key`.  See push_iterators.
514    fn needs_more_iterators(&self, search_key: Bound<&K>) -> bool {
515        if self.pending_iterators.is_empty() {
516            return false;
517        }
518        if self.heap.is_empty() {
519            return true;
520        }
521        let Bound::Included(target) = search_key else { return true };
522        match target.merge_type() {
523            MergeType::FullMerge => true,
524            MergeType::OptimizedMerge => !self.heap.peek().unwrap().key().overlaps(target),
525        }
526    }
527
528    // Pushes additional iterators onto the heap until we are confident that the top element will
529    // yield what we are looking for.
530    async fn push_iterators(&mut self, search_key: Bound<&K>) -> Result<(), Error> {
531        while self.needs_more_iterators(search_key) {
532            let iter = self.pending_iterators.pop().unwrap();
533            let sub_iter = iter.layer.as_ref().unwrap().seek(search_key).await?;
534            if self.trace {
535                writeln!(
536                    self.history,
537                    "merger: search for {:?}, found {:?}",
538                    search_key,
539                    sub_iter.get()
540                )
541                .unwrap();
542            }
543            iter.iter = RawIterator::Const(sub_iter);
544            iter.set_item_from_iter();
545            if iter.is_some() {
546                self.heap.push(iter);
547            }
548        }
549        Ok(())
550    }
551
552    // Updates the merge iterator depending on |op|. If discarding, the iterator should have already
553    // been advanced.
554    fn update_item(&mut self, item: &'a mut MergeLayerIterator<'b, K, V>, op: ItemOp<K, V>) {
555        match op {
556            ItemOp::Keep => self.heap.push(item),
557            ItemOp::Discard => {
558                // The iterator should have already been advanced.
559                if item.is_some() {
560                    self.heap.push(item);
561                }
562            }
563            ItemOp::Replace(replacement) => {
564                item.replace(replacement);
565                self.heap.push(item);
566            }
567        }
568    }
569}
570
571impl<'a, K: Key + LayerKey + OrdLowerBound, V: Value> LayerIterator<K, V>
572    for MergerIterator<'a, '_, K, V>
573{
574    async fn advance(&mut self) -> Result<(), Error> {
575        let owned_key;
576        let mut search_key = Bound::Unbounded;
577        if !self.pending_iterators.is_empty() {
578            if let Some(ItemRef { key, .. }) = self.get() {
579                match key.next_key() {
580                    Some(k) => {
581                        assert!(k.is_search_key());
582                        owned_key = k;
583                        search_key = Bound::Included(&owned_key);
584                    }
585                    None => {
586                        owned_key = key.clone();
587                        search_key = Bound::Excluded(&owned_key);
588                    }
589                }
590            }
591        }
592
593        // Advance the iterator for the current item and push it onto the heap, and also push any
594        // additional iterators onto the heap (by calling push_iterators).
595        if let Some(iterator) = self.item.take_iterator() {
596            if self.needs_more_iterators(search_key) {
597                let existing_item = iterator.item().boxed();
598                iterator.advance().await?;
599                if let Bound::Included(s) = search_key
600                    && iterator.is_some()
601                    && iterator.key().merge_type() == MergeType::OptimizedMerge
602                    && s.overlaps(iterator.key())
603                {
604                    // In this case, the key immediately following is a good candidate, and all
605                    // we need to do is merge it with existing iterators; we shouldn't need to
606                    // consult with any more iterators.
607                } else {
608                    // We are going to need to consult more iterators so we need to go back to
609                    // the previous item so that we can merge with it.  See the comment in
610                    // advance_impl.
611                    iterator.replace(existing_item);
612
613                    // We must push other iterators here before pushing iterator onto the heap
614                    // because we know `iterator` would end up at the top of the heap.
615                    self.push_iterators(search_key).await?;
616                }
617            } else {
618                iterator.advance().await?;
619            }
620            if iterator.is_some() {
621                self.heap.push(iterator);
622            }
623        } else {
624            self.push_iterators(search_key).await?;
625        }
626
627        self.advance_impl(search_key).await
628    }
629
630    fn advance_dyn<'b>(&'b mut self) -> Result<Option<BoxFuture<'b, Result<(), Error>>>, Error> {
631        Ok(Some(Box::pin(self.advance())))
632    }
633
634    fn get(&self) -> Option<ItemRef<'_, K, V>> {
635        (&self.item).into()
636    }
637}
638
639struct MutMergeLayerIterator<'a, K, V>(MergeLayerIterator<'a, K, V>);
640
641impl<K, V> MutMergeLayerIterator<'_, K, V> {
642    fn advance(&mut self) {
643        if let MergeItem::Iter = self.item {
644            self.as_mut().advance();
645        }
646        self.set_item_from_iter();
647    }
648}
649
650impl<'a, K, V> AsMut<dyn LayerIteratorMut<K, V> + 'a> for MutMergeLayerIterator<'a, K, V> {
651    fn as_mut(&mut self) -> &mut (dyn LayerIteratorMut<K, V> + 'a) {
652        let RawIterator::Mut(iter) = &mut self.0.iter else { unreachable!() };
653        iter.as_mut()
654    }
655}
656
657impl<'a, K, V> Deref for MutMergeLayerIterator<'a, K, V> {
658    type Target = MergeLayerIterator<'a, K, V>;
659
660    fn deref(&self) -> &Self::Target {
661        &self.0
662    }
663}
664
665impl<K, V> DerefMut for MutMergeLayerIterator<'_, K, V> {
666    fn deref_mut(&mut self) -> &mut Self::Target {
667        &mut self.0
668    }
669}
670
671// Merges the given item into a mutable layer.
672pub(super) fn merge_into<K: Debug + OrdLowerBound, V: Debug>(
673    mut_iter: Box<dyn LayerIteratorMut<K, V> + '_>,
674    item: Item<K, V>,
675    merge_fn: MergeFn<K, V>,
676) -> Result<(), Error> {
677    let merge_item = if mut_iter.get().is_some() { MergeItem::Iter } else { MergeItem::None };
678    let mut mut_merge_iter = MutMergeLayerIterator(MergeLayerIterator {
679        layer: None,
680        iter: RawIterator::Mut(mut_iter),
681        layer_index: 1,
682        item: merge_item,
683    });
684    let mut item_merge_iter = MergeLayerIterator::new_with_item(0, MergeItem::Item(item.boxed()));
685    while mut_merge_iter.is_some() && item_merge_iter.is_some() {
686        if mut_merge_iter.0 > item_merge_iter {
687            // In this branch the mutable layer is left and the item we're merging-in is right.
688            let merge_result = merge_fn(&mut_merge_iter, &item_merge_iter);
689            debug!(
690                lhs:? = mut_merge_iter.key(),
691                rhs:? = item_merge_iter.key(),
692                result:? = merge_result;
693                "(1) merge");
694            match merge_result {
695                MergeResult::EmitLeft => {
696                    if let Some(item) = mut_merge_iter.take_item() {
697                        mut_merge_iter.as_mut().insert(*item);
698                        mut_merge_iter.set_item_from_iter();
699                    } else {
700                        mut_merge_iter.advance();
701                    }
702                }
703                MergeResult::Other { emit, left, right } => {
704                    if let Some(emit) = emit {
705                        mut_merge_iter.as_mut().insert(*emit);
706                    }
707                    match left {
708                        ItemOp::Keep => {}
709                        ItemOp::Discard => {
710                            if matches!(mut_merge_iter.item, MergeItem::Iter) {
711                                mut_merge_iter.as_mut().erase();
712                            }
713                            mut_merge_iter.set_item_from_iter();
714                        }
715                        ItemOp::Replace(item) => {
716                            if let MergeItem::Iter = mut_merge_iter.item {
717                                mut_merge_iter.as_mut().erase();
718                            }
719                            mut_merge_iter.item = MergeItem::Item(item)
720                        }
721                    }
722                    match right {
723                        ItemOp::Keep => {}
724                        ItemOp::Discard => item_merge_iter.item = MergeItem::None,
725                        ItemOp::Replace(item) => item_merge_iter.item = MergeItem::Item(item),
726                    }
727                }
728            }
729        } else {
730            // In this branch, the item we're merging-in is left and the mutable layer is right.
731            let merge_result = merge_fn(&item_merge_iter, &mut_merge_iter);
732            debug!(
733                lhs:? = mut_merge_iter.key(),
734                rhs:? = item_merge_iter.key(),
735                result:? = merge_result;
736                "(2) merge");
737            match merge_result {
738                MergeResult::EmitLeft => break, // Item is inserted outside the loop
739                MergeResult::Other { emit, left, right } => {
740                    if let Some(emit) = emit {
741                        mut_merge_iter.as_mut().insert(*emit);
742                    }
743                    match left {
744                        ItemOp::Keep => {}
745                        ItemOp::Discard => item_merge_iter.item = MergeItem::None,
746                        ItemOp::Replace(item) => item_merge_iter.item = MergeItem::Item(item),
747                    }
748                    match right {
749                        ItemOp::Keep => {}
750                        ItemOp::Discard => {
751                            if matches!(mut_merge_iter.item, MergeItem::Iter) {
752                                mut_merge_iter.as_mut().erase();
753                            }
754                            mut_merge_iter.set_item_from_iter();
755                        }
756                        ItemOp::Replace(item) => {
757                            if let MergeItem::Iter = mut_merge_iter.item {
758                                mut_merge_iter.as_mut().erase();
759                            }
760                            mut_merge_iter.item = MergeItem::Item(item)
761                        }
762                    }
763                }
764            }
765        }
766    } // while ...
767
768    // The only way we could get here with both items is via the break above, so we know the correct
769    // order required here.
770    if let MergeItem::Item(item) = item_merge_iter.item {
771        mut_merge_iter.as_mut().insert(*item);
772    }
773    if let Some(item) = mut_merge_iter.take_item() {
774        mut_merge_iter.as_mut().insert(*item);
775    }
776    if let RawIterator::Mut(mut iter) = mut_merge_iter.0.iter {
777        iter.commit();
778    }
779    Ok(())
780}
781
782#[cfg(test)]
783mod tests {
784    use super::ItemOp::{Discard, Keep, Replace};
785    use super::{MergeLayerIterator, MergeResult, Merger};
786    use crate::lsm_tree::persistent_layer::{PersistentLayer, PersistentLayerWriter};
787    use crate::lsm_tree::skip_list_layer::SkipListLayer;
788    use crate::lsm_tree::types::{
789        FuzzyHash, Item, ItemRef, Key, Layer, LayerIterator, LayerKey, LayerWriter,
790        MaybeContainsKey, MergeType, OrdLowerBound, OrdUpperBound, SortByU64,
791    };
792    use crate::lsm_tree::{self, Query, Value};
793    use crate::object_store::{self, AttributeId, ObjectKey, ObjectValue, VOLUME_DATA_KEY_ID};
794    use crate::serialized_types::{
795        LATEST_VERSION, Version, Versioned, VersionedLatest, versioned_type,
796    };
797    use crate::testing::fake_object::{FakeObject, FakeObjectHandle};
798    use crate::testing::writer::Writer;
799    use fprint::TypeFingerprint;
800    use fxfs_macros::{FuzzyHash, SerializeKey};
801    use rand::RngExt as _;
802    use std::hash::Hash;
803    use std::ops::{Bound, Range};
804    use std::sync::Arc;
805    use storage_units::BlockSize;
806
807    use crate::lsm_tree::testing::TestKey;
808
809    impl Value for i32 {
810        const DELETED_MARKER: Self = 0;
811    }
812
813    fn layer_ref_iter<K: Key, V: Value>(
814        layers: &[Arc<SkipListLayer<K, V>>],
815    ) -> impl Iterator<Item = &dyn Layer<K, V>> {
816        layers.iter().map(|x| x.as_ref() as &dyn Layer<K, V>)
817    }
818
819    fn dyn_layer_ref_iter<K: Key, V: Value>(
820        layers: &[Arc<dyn Layer<K, V>>],
821    ) -> impl Iterator<Item = &dyn Layer<K, V>> {
822        layers.iter().map(|x| x.as_ref())
823    }
824
825    fn counters() -> Arc<lsm_tree::TreeCounters> {
826        Arc::new(lsm_tree::TreeCounters::default())
827    }
828
829    #[fuchsia::test]
830    async fn test_emit_left() {
831        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
832        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
833        skip_lists[0].insert(items[1].clone()).expect("insert error");
834        skip_lists[1].insert(items[0].clone()).expect("insert error");
835        let mut merger = Merger::new(
836            layer_ref_iter(&skip_lists),
837            |_left, _right| MergeResult::EmitLeft,
838            counters(),
839        );
840        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
841        let ItemRef { key, value, .. } = iter.get().expect("missing item");
842        assert_eq!((key, value), (&items[0].key, &items[0].value));
843        iter.advance().await.unwrap();
844        let ItemRef { key, value, .. } = iter.get().expect("missing item");
845        assert_eq!((key, value), (&items[1].key, &items[1].value));
846        iter.advance().await.unwrap();
847        assert!(iter.get().is_none());
848    }
849
850    #[fuchsia::test]
851    async fn test_other_emit() {
852        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
853        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
854        skip_lists[0].insert(items[1].clone()).expect("insert error");
855        skip_lists[1].insert(items[0].clone()).expect("insert error");
856        let mut merger = Merger::new(
857            layer_ref_iter(&skip_lists),
858            |_left, _right| MergeResult::Other {
859                emit: Some(Item::new(TestKey(3..3), 3).boxed()),
860                left: Discard,
861                right: Discard,
862            },
863            counters(),
864        );
865        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
866
867        let ItemRef { key, value, .. } = iter.get().expect("missing item");
868        assert_eq!((key, value), (&TestKey(3..3), &3));
869        iter.advance().await.unwrap();
870        assert!(iter.get().is_none());
871    }
872
873    #[fuchsia::test]
874    async fn test_replace_left() {
875        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
876        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
877        skip_lists[0].insert(items[1].clone()).expect("insert error");
878        skip_lists[1].insert(items[0].clone()).expect("insert error");
879        let mut merger = Merger::new(
880            layer_ref_iter(&skip_lists),
881            |_left, _right| MergeResult::Other {
882                emit: None,
883                left: Replace(Item::new(TestKey(3..3), 3).boxed()),
884                right: Discard,
885            },
886            counters(),
887        );
888        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
889
890        // The merger should replace the left item and then after discarding the right item, it
891        // should emit the replacement.
892        let ItemRef { key, value, .. } = iter.get().expect("missing item");
893        assert_eq!((key, value), (&TestKey(3..3), &3));
894        iter.advance().await.unwrap();
895        assert!(iter.get().is_none());
896    }
897
898    #[fuchsia::test]
899    async fn test_replace_right() {
900        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
901        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
902        skip_lists[0].insert(items[1].clone()).expect("insert error");
903        skip_lists[1].insert(items[0].clone()).expect("insert error");
904        let mut merger = Merger::new(
905            layer_ref_iter(&skip_lists),
906            |_left, _right| MergeResult::Other {
907                emit: None,
908                left: Discard,
909                right: Replace(Item::new(TestKey(3..3), 3).boxed()),
910            },
911            counters(),
912        );
913        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
914
915        // The merger should replace the right item and then after discarding the left item, it
916        // should emit the replacement.
917        let ItemRef { key, value, .. } = iter.get().expect("missing item");
918        assert_eq!((key, value), (&TestKey(3..3), &3));
919        iter.advance().await.unwrap();
920        assert!(iter.get().is_none());
921    }
922
923    #[fuchsia::test]
924    async fn test_left_less_than_right() {
925        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
926        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
927        skip_lists[0].insert(items[1].clone()).expect("insert error");
928        skip_lists[1].insert(items[0].clone()).expect("insert error");
929        let mut merger = Merger::new(
930            layer_ref_iter(&skip_lists),
931            |left, right| {
932                assert_eq!((left.key(), left.value()), (&TestKey(1..1), &1));
933                assert_eq!((right.key(), right.value()), (&TestKey(2..2), &2));
934                MergeResult::EmitLeft
935            },
936            counters(),
937        );
938        merger.query(Query::FullScan).await.expect("seek failed");
939    }
940
941    #[fuchsia::test]
942    async fn test_left_equals_right() {
943        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
944        let item = Item::new(TestKey(1..1), 1);
945        skip_lists[0].insert(item.clone()).expect("insert error");
946        skip_lists[1].insert(item.clone()).expect("insert error");
947        let mut merger = Merger::new(
948            layer_ref_iter(&skip_lists),
949            |left, right| {
950                assert_eq!((left.key(), left.value()), (&TestKey(1..1), &1));
951                assert_eq!((right.key(), right.value()), (&TestKey(1..1), &1));
952                assert_eq!(left.layer_index, 0);
953                assert_eq!(right.layer_index, 1);
954                MergeResult::EmitLeft
955            },
956            counters(),
957        );
958        merger.query(Query::FullScan).await.expect("seek failed");
959    }
960
961    #[fuchsia::test]
962    async fn test_keep() {
963        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
964        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
965        skip_lists[0].insert(items[1].clone()).expect("insert error");
966        skip_lists[1].insert(items[0].clone()).expect("insert error");
967        let mut merger = Merger::new(
968            layer_ref_iter(&skip_lists),
969            |left, right| {
970                if left.key() == &TestKey(1..1) {
971                    MergeResult::Other {
972                        emit: None,
973                        left: Replace(Item::new(TestKey(3..3), 3).boxed()),
974                        right: Keep,
975                    }
976                } else {
977                    assert_eq!(left.key(), &TestKey(2..2));
978                    assert_eq!(right.key(), &TestKey(3..3));
979                    MergeResult::Other { emit: None, left: Discard, right: Keep }
980                }
981            },
982            counters(),
983        );
984        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
985
986        // The merger should first replace left and then it should call the merger again with 2 & 3
987        // and end up just keeping 3.
988        let ItemRef { key, value, .. } = iter.get().expect("missing item");
989        assert_eq!((key, value), (&TestKey(3..3), &3));
990        iter.advance().await.unwrap();
991        assert!(iter.get().is_none());
992    }
993
994    #[fuchsia::test]
995    async fn test_merge_10_layers() {
996        let skip_lists: Vec<_> = (0..10).map(|_| SkipListLayer::new(100)).collect();
997        let mut rng = rand::rng();
998        for i in 0..100 {
999            skip_lists[rng.random_range(0..10) as usize]
1000                .insert(Item::new(TestKey(i..i), i))
1001                .expect("insert error");
1002        }
1003        let mut merger = Merger::new(
1004            layer_ref_iter(&skip_lists),
1005            |_left, _right| MergeResult::EmitLeft,
1006            counters(),
1007        );
1008        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
1009
1010        for i in 0..100 {
1011            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1012            assert_eq!((key, value), (&TestKey(i..i), &i));
1013            iter.advance().await.unwrap();
1014        }
1015        assert!(iter.get().is_none());
1016    }
1017
1018    #[fuchsia::test]
1019    async fn test_merge_uses_cmp_lower_bound() {
1020        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1021        let items = [Item::new(TestKey(1..10), 1), Item::new(TestKey(2..3), 2)];
1022        skip_lists[0].insert(items[1].clone()).expect("insert error");
1023        skip_lists[1].insert(items[0].clone()).expect("insert error");
1024        let mut merger = Merger::new(
1025            layer_ref_iter(&skip_lists),
1026            |_left, _right| MergeResult::EmitLeft,
1027            counters(),
1028        );
1029        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
1030
1031        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1032        assert_eq!((key, value), (&items[0].key, &items[0].value));
1033        iter.advance().await.unwrap();
1034        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1035        assert_eq!((key, value), (&items[1].key, &items[1].value));
1036        iter.advance().await.unwrap();
1037        assert!(iter.get().is_none());
1038    }
1039
1040    #[fuchsia::test]
1041    async fn test_merge_into_emit_left() {
1042        let skip_list = SkipListLayer::new(100);
1043        let items =
1044            [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2), Item::new(TestKey(3..3), 3)];
1045        skip_list.insert(items[0].clone()).expect("insert error");
1046        skip_list.insert(items[2].clone()).expect("insert error");
1047        skip_list
1048            .merge_into(items[1].clone(), &items[0].key, |_left, _right| MergeResult::EmitLeft);
1049
1050        let mut iter = skip_list.seek(Bound::Unbounded);
1051        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1052        assert_eq!((key, value), (&items[0].key, &items[0].value));
1053        iter.advance().await.unwrap();
1054        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1055        assert_eq!((key, value), (&items[1].key, &items[1].value));
1056        iter.advance().await.unwrap();
1057        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1058        assert_eq!((key, value), (&items[2].key, &items[2].value));
1059        iter.advance().await.unwrap();
1060        assert!(iter.get().is_none());
1061    }
1062
1063    #[fuchsia::test]
1064    async fn test_merge_into_emit_last_after_replacing() {
1065        let skip_list = SkipListLayer::new(100);
1066        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
1067        skip_list.insert(items[0].clone()).expect("insert error");
1068
1069        skip_list.merge_into(items[1].clone(), &items[0].key, |left, right| {
1070            if left.key() == &TestKey(1..1) {
1071                assert_eq!(right.key(), &TestKey(2..2));
1072                MergeResult::Other {
1073                    emit: None,
1074                    left: Replace(Item::new(TestKey(3..3), 3).boxed()),
1075                    right: Keep,
1076                }
1077            } else {
1078                assert_eq!(left.key(), &TestKey(2..2));
1079                assert_eq!(right.key(), &TestKey(3..3));
1080                MergeResult::EmitLeft
1081            }
1082        });
1083
1084        let mut iter = skip_list.seek(Bound::Unbounded);
1085        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1086        assert_eq!((key, value), (&items[1].key, &items[1].value));
1087        iter.advance().await.unwrap();
1088        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1089        assert_eq!((key, value), (&TestKey(3..3), &3));
1090        iter.advance().await.unwrap();
1091        assert!(iter.get().is_none());
1092    }
1093
1094    #[fuchsia::test]
1095    async fn test_merge_into_emit_left_after_replacing() {
1096        let skip_list = SkipListLayer::new(100);
1097        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(3..3), 3)];
1098        skip_list.insert(items[0].clone()).expect("insert error");
1099
1100        skip_list.merge_into(items[1].clone(), &items[0].key, |left, right| {
1101            if left.key() == &TestKey(1..1) {
1102                assert_eq!(right.key(), &TestKey(3..3));
1103                MergeResult::Other {
1104                    emit: None,
1105                    left: Replace(Item::new(TestKey(2..2), 2).boxed()),
1106                    right: Keep,
1107                }
1108            } else {
1109                assert_eq!(left.key(), &TestKey(2..2));
1110                assert_eq!(right.key(), &TestKey(3..3));
1111                MergeResult::EmitLeft
1112            }
1113        });
1114
1115        let mut iter = skip_list.seek(Bound::Unbounded);
1116        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1117        assert_eq!((key, value), (&TestKey(2..2), &2));
1118        iter.advance().await.unwrap();
1119        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1120        assert_eq!((key, value), (&items[1].key, &items[1].value));
1121        iter.advance().await.unwrap();
1122        assert!(iter.get().is_none());
1123    }
1124
1125    // This tests emitting in both branches of merge_into, and most of the discard paths.
1126    #[fuchsia::test]
1127    async fn test_merge_into_emit_other_and_discard() {
1128        let skip_list = SkipListLayer::new(100);
1129        let items =
1130            [Item::new(TestKey(1..1), 1), Item::new(TestKey(3..3), 3), Item::new(TestKey(5..5), 3)];
1131        skip_list.insert(items[0].clone()).expect("insert error");
1132        skip_list.insert(items[2].clone()).expect("insert error");
1133
1134        skip_list.merge_into(items[1].clone(), &items[0].key, |left, right| {
1135            if left.key() == &TestKey(1..1) {
1136                // This tests the top branch in merge_into.
1137                assert_eq!(right.key(), &TestKey(3..3));
1138                MergeResult::Other {
1139                    emit: Some(Item::new(TestKey(2..2), 2).boxed()),
1140                    left: Discard,
1141                    right: Keep,
1142                }
1143            } else {
1144                // This tests the bottom branch in merge_into.
1145                assert_eq!(left.key(), &TestKey(3..3));
1146                assert_eq!(right.key(), &TestKey(5..5));
1147                MergeResult::Other {
1148                    emit: Some(Item::new(TestKey(4..4), 4).boxed()),
1149                    left: Discard,
1150                    right: Discard,
1151                }
1152            }
1153        });
1154
1155        let mut iter = skip_list.seek(Bound::Unbounded);
1156        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1157        assert_eq!((key, value), (&TestKey(2..2), &2));
1158        iter.advance().await.unwrap();
1159        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1160        assert_eq!((key, value), (&TestKey(4..4), &4));
1161        iter.advance().await.unwrap();
1162        assert!(iter.get().is_none());
1163    }
1164
1165    // This tests replacing the item and discarding the right item (the one remaining untested
1166    // discard path) in the top branch in merge_into.
1167    #[fuchsia::test]
1168    async fn test_merge_into_replace_and_discard() {
1169        let skip_list = SkipListLayer::new(100);
1170        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(3..3), 3)];
1171        skip_list.insert(items[0].clone()).expect("insert error");
1172
1173        skip_list.merge_into(items[1].clone(), &items[0].key, |_left, _right| MergeResult::Other {
1174            emit: Some(Item::new(TestKey(2..2), 2).boxed()),
1175            left: Replace(Item::new(TestKey(4..4), 4).boxed()),
1176            right: Discard,
1177        });
1178
1179        let mut iter = skip_list.seek(Bound::Unbounded);
1180        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1181        assert_eq!((key, value), (&TestKey(2..2), &2));
1182        iter.advance().await.unwrap();
1183        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1184        assert_eq!((key, value), (&TestKey(4..4), &4));
1185        iter.advance().await.unwrap();
1186        assert!(iter.get().is_none());
1187    }
1188
1189    // This tests replacing the right item in the top branch of merge_into and the left item in the
1190    // bottom branch of merge_into.
1191    #[fuchsia::test]
1192    async fn test_merge_into_replace_merge_item() {
1193        let skip_list = SkipListLayer::new(100);
1194        let items =
1195            [Item::new(TestKey(1..1), 1), Item::new(TestKey(3..3), 3), Item::new(TestKey(5..5), 5)];
1196        skip_list.insert(items[0].clone()).expect("insert error");
1197        skip_list.insert(items[2].clone()).expect("insert error");
1198
1199        skip_list.merge_into(items[1].clone(), &items[0].key, |_left, right| {
1200            if right.key() == &TestKey(3..3) {
1201                MergeResult::Other {
1202                    emit: None,
1203                    left: Discard,
1204                    right: Replace(Item::new(TestKey(2..2), 2).boxed()),
1205                }
1206            } else {
1207                assert_eq!(right.key(), &TestKey(5..5));
1208                MergeResult::Other {
1209                    emit: None,
1210                    left: Replace(Item::new(TestKey(4..4), 4).boxed()),
1211                    right: Discard,
1212                }
1213            }
1214        });
1215
1216        let mut iter = skip_list.seek(Bound::Unbounded);
1217        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1218        assert_eq!((key, value), (&TestKey(4..4), &4));
1219        iter.advance().await.unwrap();
1220        assert!(iter.get().is_none());
1221    }
1222
1223    // This tests replacing the right item in the bottom branch of merge_into.
1224    #[fuchsia::test]
1225    async fn test_merge_into_replace_existing() {
1226        let skip_list = SkipListLayer::new(100);
1227        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(3..3), 3)];
1228        skip_list.insert(items[1].clone()).expect("insert error");
1229
1230        skip_list.merge_into(items[0].clone(), &items[0].key, |_left, right| {
1231            if right.key() == &TestKey(3..3) {
1232                MergeResult::Other {
1233                    emit: None,
1234                    left: Keep,
1235                    right: Replace(Item::new(TestKey(2..2), 2).boxed()),
1236                }
1237            } else {
1238                MergeResult::EmitLeft
1239            }
1240        });
1241
1242        let mut iter = skip_list.seek(Bound::Unbounded);
1243        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1244        assert_eq!((key, value), (&items[0].key, &items[0].value));
1245        iter.advance().await.unwrap();
1246        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1247        assert_eq!((key, value), (&TestKey(2..2), &2));
1248        iter.advance().await.unwrap();
1249        assert!(iter.get().is_none());
1250    }
1251
1252    #[fuchsia::test]
1253    async fn test_merge_into_discard_last() {
1254        let skip_list = SkipListLayer::new(100);
1255        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
1256        skip_list.insert(items[0].clone()).expect("insert error");
1257
1258        skip_list.merge_into(items[1].clone(), &items[0].key, |_left, _right| MergeResult::Other {
1259            emit: None,
1260            left: Discard,
1261            right: Keep,
1262        });
1263
1264        let mut iter = skip_list.seek(Bound::Unbounded);
1265        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1266        assert_eq!((key, value), (&items[1].key, &items[1].value));
1267        iter.advance().await.unwrap();
1268        assert!(iter.get().is_none());
1269    }
1270
1271    #[fuchsia::test]
1272    async fn test_merge_into_empty() {
1273        let skip_list = SkipListLayer::new(100);
1274        let items = [Item::new(TestKey(1..1), 1)];
1275
1276        skip_list.merge_into(items[0].clone(), &items[0].key, |_left, _right| {
1277            panic!("Unexpected merge!");
1278        });
1279
1280        let mut iter = skip_list.seek(Bound::Unbounded);
1281        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1282        assert_eq!((key, value), (&items[0].key, &items[0].value));
1283        iter.advance().await.unwrap();
1284        assert!(iter.get().is_none());
1285    }
1286
1287    #[fuchsia::test]
1288    async fn test_seek_uses_minimum_number_of_iterators() {
1289        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1290        let items = [Item::new(TestKey(1..2), 1), Item::new(TestKey(1..2), 2)];
1291        skip_lists[0].insert(items[0].clone()).expect("insert error");
1292        skip_lists[1].insert(items[1].clone()).expect("insert error");
1293        let mut merger = Merger::new(
1294            layer_ref_iter(&skip_lists),
1295            |_left, _right| MergeResult::Other { emit: None, left: Discard, right: Keep },
1296            counters(),
1297        );
1298        let iter = merger
1299            .query(Query::FullRange(&items[0].key.search_key().unwrap()))
1300            .await
1301            .expect("seek failed");
1302
1303        // Seek should only search in the first skip list, so no merge should take place, and we'll
1304        // know if it has because we'll see a different value (2 rather than 1).
1305        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1306        assert_eq!((key, value), (&items[0].key, &items[0].value));
1307    }
1308
1309    // Checks that merging the given layers produces |expected| sequence of items starting from
1310    // |start|.
1311    async fn test_advance<K: Eq + Key + LayerKey + OrdLowerBound>(
1312        layers: &[&[(K, i32)]],
1313        query: Query<'_, K>,
1314        expected: &[(K, i32)],
1315    ) {
1316        let mut skip_lists = Vec::new();
1317        for &layer in layers {
1318            let skip_list = SkipListLayer::new(100);
1319            for (k, v) in layer {
1320                skip_list.insert(Item::new(k.clone(), *v)).expect("insert error");
1321            }
1322            skip_lists.push(skip_list);
1323        }
1324        let mut merger = Merger::new(
1325            layer_ref_iter(&skip_lists),
1326            |_left, _right| MergeResult::EmitLeft,
1327            counters(),
1328        );
1329        let mut iter = merger.query(query).await.expect("seek failed");
1330        for (k, v) in expected {
1331            let ItemRef { key, value, .. } = iter.get().expect("get failed");
1332            assert_eq!((key, value), (k, v));
1333            iter.advance().await.expect("advance failed");
1334        }
1335        assert!(iter.get().is_none());
1336    }
1337
1338    #[fuchsia::test]
1339    async fn test_seek_skips_replaced_items() {
1340        // The 1..2 and the 2..3 items are overwritten and merging them should be skipped.
1341        test_advance(
1342            &[
1343                &[(TestKey(1..2), 1), (TestKey(2..3), 2), (TestKey(4..5), 3)],
1344                &[(TestKey(1..2), 4), (TestKey(2..3), 5), (TestKey(3..4), 6)],
1345            ],
1346            Query::FullRange(&TestKey(1..2).search_key().unwrap()),
1347            &[(TestKey(1..2), 1), (TestKey(2..3), 2), (TestKey(3..4), 6), (TestKey(4..5), 3)],
1348        )
1349        .await;
1350    }
1351
1352    #[fuchsia::test]
1353    async fn test_advance_skips_replaced_items_at_end() {
1354        // Like the last test, the 1..2 item is overwritten and seeking for it should skip the merge
1355        // but this time, the keys are at the end.
1356        test_advance(
1357            &[&[(TestKey(1..2), 1)], &[(TestKey(1..2), 2)]],
1358            Query::FullRange(&TestKey(1..2).search_key().unwrap()),
1359            &[(TestKey(1..2), 1)],
1360        )
1361        .await;
1362    }
1363
1364    #[derive(
1365        Clone,
1366        Eq,
1367        Hash,
1368        FuzzyHash,
1369        PartialEq,
1370        Debug,
1371        serde::Serialize,
1372        serde::Deserialize,
1373        TypeFingerprint,
1374        Versioned,
1375        SerializeKey,
1376    )]
1377    struct TestKeyWithFullMerge(Range<u64>);
1378
1379    versioned_type! { 1.. => TestKeyWithFullMerge }
1380
1381    impl LayerKey for TestKeyWithFullMerge {
1382        fn merge_type(&self) -> MergeType {
1383            MergeType::FullMerge
1384        }
1385
1386        fn search_key(&self) -> Option<Self> {
1387            Some(Self(self.0.start..self.0.start + 1))
1388        }
1389
1390        fn is_search_key(&self) -> bool {
1391            self.0.end == self.0.start + 1
1392        }
1393
1394        fn overlaps(&self, other: &Self) -> bool {
1395            self.0.start < other.0.end && self.0.end > other.0.start
1396        }
1397    }
1398
1399    impl SortByU64 for TestKeyWithFullMerge {
1400        fn get_leading_u64(&self) -> u64 {
1401            self.0.end
1402        }
1403    }
1404
1405    impl OrdUpperBound for TestKeyWithFullMerge {
1406        fn cmp_upper_bound(&self, other: &TestKeyWithFullMerge) -> std::cmp::Ordering {
1407            self.0.end.cmp(&other.0.end).then(other.0.start.cmp(&self.0.start))
1408        }
1409    }
1410
1411    impl OrdLowerBound for TestKeyWithFullMerge {
1412        fn cmp_lower_bound(&self, other: &Self) -> std::cmp::Ordering {
1413            self.0.start.cmp(&other.0.start)
1414        }
1415    }
1416
1417    // Same set up as `test_seek_skips_replaced_items` except here nothing should skip. Any seek
1418    // should just place the iterator at a particular spot in the merged layers.
1419    #[fuchsia::test]
1420    async fn test_full_merge_consistent_advance_ordering() {
1421        let layer_set = [
1422            [
1423                (TestKeyWithFullMerge(1..2), 1i32),
1424                (TestKeyWithFullMerge(2..3), 2i32),
1425                (TestKeyWithFullMerge(4..5), 3i32),
1426            ]
1427            .as_slice(),
1428            [
1429                (TestKeyWithFullMerge(1..2), 4i32),
1430                (TestKeyWithFullMerge(2..3), 5i32),
1431                (TestKeyWithFullMerge(3..4), 6i32),
1432            ]
1433            .as_slice(),
1434        ];
1435
1436        let full_merge_result = [
1437            (TestKeyWithFullMerge(1..2), 1),
1438            (TestKeyWithFullMerge(1..2), 4),
1439            (TestKeyWithFullMerge(2..3), 2),
1440            (TestKeyWithFullMerge(2..3), 5),
1441            (TestKeyWithFullMerge(3..4), 6),
1442            (TestKeyWithFullMerge(4..5), 3),
1443        ];
1444
1445        test_advance(layer_set.as_slice(), Query::FullScan, &full_merge_result).await;
1446
1447        test_advance(
1448            layer_set.as_slice(),
1449            Query::FullRange(&TestKeyWithFullMerge(1..2).search_key().unwrap()),
1450            &full_merge_result,
1451        )
1452        .await;
1453
1454        test_advance(
1455            layer_set.as_slice(),
1456            Query::FullRange(&TestKeyWithFullMerge(2..3).search_key().unwrap()),
1457            &full_merge_result[2..],
1458        )
1459        .await;
1460
1461        test_advance(
1462            layer_set.as_slice(),
1463            Query::FullRange(&TestKeyWithFullMerge(3..4).search_key().unwrap()),
1464            &full_merge_result[4..],
1465        )
1466        .await;
1467    }
1468
1469    #[fuchsia::test]
1470    async fn test_full_merge_always_consult_all_layers() {
1471        // | 1 |   |
1472        // |   | 2 |
1473        // | 3 | 4 |
1474        let skip_lists =
1475            [SkipListLayer::new(100), SkipListLayer::new(100), SkipListLayer::new(100)];
1476        let items = [
1477            Item::new(TestKeyWithFullMerge(1..2), 1),
1478            Item::new(TestKeyWithFullMerge(2..3), 2),
1479            Item::new(TestKeyWithFullMerge(1..2), 3),
1480            Item::new(TestKeyWithFullMerge(2..3), 4),
1481        ];
1482        skip_lists[0].insert(items[0].clone()).expect("insert error");
1483        skip_lists[1].insert(items[1].clone()).expect("insert error");
1484        skip_lists[2].insert(items[2].clone()).expect("insert error");
1485        skip_lists[2].insert(items[3].clone()).expect("insert error");
1486        let mut merger = Merger::new(
1487            layer_ref_iter(&skip_lists),
1488            |left, right| {
1489                // Sum matching keys.
1490                if left.key() == right.key() {
1491                    MergeResult::Other {
1492                        emit: None,
1493                        left: Discard,
1494                        right: Replace(
1495                            Item::new(left.key().clone(), left.value() + right.value()).boxed(),
1496                        ),
1497                    }
1498                } else {
1499                    MergeResult::EmitLeft
1500                }
1501            },
1502            counters(),
1503        );
1504        let mut iter = merger
1505            .query(Query::FullRange(&items[0].key.search_key().unwrap()))
1506            .await
1507            .expect("seek failed");
1508
1509        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1510        assert_eq!((key, *value), (&items[0].key, items[0].value + items[2].value));
1511        iter.advance().await.expect("advance");
1512        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1513        assert_eq!((key, *value), (&items[1].key, items[1].value + items[3].value));
1514        iter.advance().await.expect("advance");
1515        assert!(iter.get().is_none());
1516    }
1517
1518    #[derive(
1519        Clone,
1520        Eq,
1521        Hash,
1522        FuzzyHash,
1523        PartialEq,
1524        Debug,
1525        serde::Serialize,
1526        serde::Deserialize,
1527        TypeFingerprint,
1528        Versioned,
1529        SerializeKey,
1530    )]
1531    struct TestKeyWithDefaultLayerKey(Range<u64>);
1532
1533    versioned_type! { 1.. => TestKeyWithDefaultLayerKey }
1534
1535    // Default layer key is using `MergeType::FullMerge` and returns None for `next_key()`.
1536    impl LayerKey for TestKeyWithDefaultLayerKey {
1537        fn search_key(&self) -> Option<Self> {
1538            Some(Self(self.0.start..self.0.start + 1))
1539        }
1540
1541        fn is_search_key(&self) -> bool {
1542            self.0.end == self.0.start + 1
1543        }
1544
1545        fn overlaps(&self, other: &Self) -> bool {
1546            self.0.start < other.0.end && self.0.end > other.0.start
1547        }
1548    }
1549
1550    impl SortByU64 for TestKeyWithDefaultLayerKey {
1551        fn get_leading_u64(&self) -> u64 {
1552            self.0.end
1553        }
1554    }
1555
1556    impl OrdUpperBound for TestKeyWithDefaultLayerKey {
1557        fn cmp_upper_bound(&self, other: &TestKeyWithDefaultLayerKey) -> std::cmp::Ordering {
1558            self.0.end.cmp(&other.0.end).then(other.0.start.cmp(&self.0.start))
1559        }
1560    }
1561
1562    impl OrdLowerBound for TestKeyWithDefaultLayerKey {
1563        fn cmp_lower_bound(&self, other: &Self) -> std::cmp::Ordering {
1564            self.0.start.cmp(&other.0.start)
1565        }
1566    }
1567
1568    #[fuchsia::test]
1569    async fn test_no_merge_unbounded_include_all_layers() {
1570        test_advance(
1571            &[
1572                &[
1573                    (TestKeyWithDefaultLayerKey(1..2), 1),
1574                    (TestKeyWithDefaultLayerKey(2..3), 2),
1575                    (TestKeyWithDefaultLayerKey(4..5), 3),
1576                ],
1577                &[
1578                    (TestKeyWithDefaultLayerKey(1..2), 4),
1579                    (TestKeyWithDefaultLayerKey(2..3), 5),
1580                    (TestKeyWithDefaultLayerKey(3..4), 6),
1581                ],
1582            ],
1583            Query::FullScan,
1584            &[
1585                (TestKeyWithDefaultLayerKey(1..2), 1),
1586                (TestKeyWithDefaultLayerKey(1..2), 4),
1587                (TestKeyWithDefaultLayerKey(2..3), 2),
1588                (TestKeyWithDefaultLayerKey(2..3), 5),
1589                (TestKeyWithDefaultLayerKey(3..4), 6),
1590                (TestKeyWithDefaultLayerKey(4..5), 3),
1591            ],
1592        )
1593        .await;
1594    }
1595
1596    #[fuchsia::test]
1597    async fn test_no_merge_proceeds_comprehensively_after_seek() {
1598        test_advance(
1599            &[
1600                &[
1601                    (TestKeyWithDefaultLayerKey(1..2), 1),
1602                    (TestKeyWithDefaultLayerKey(2..3), 2),
1603                    (TestKeyWithDefaultLayerKey(4..5), 3),
1604                ],
1605                &[
1606                    (TestKeyWithDefaultLayerKey(1..2), 4),
1607                    (TestKeyWithDefaultLayerKey(2..3), 5),
1608                    (TestKeyWithDefaultLayerKey(3..4), 6),
1609                ],
1610            ],
1611            Query::FullRange(&TestKeyWithDefaultLayerKey(1..2).search_key().unwrap()),
1612            &[
1613                (TestKeyWithDefaultLayerKey(1..2), 1),
1614                (TestKeyWithDefaultLayerKey(1..2), 4),
1615                (TestKeyWithDefaultLayerKey(2..3), 2),
1616                (TestKeyWithDefaultLayerKey(2..3), 5),
1617                (TestKeyWithDefaultLayerKey(3..4), 6),
1618                (TestKeyWithDefaultLayerKey(4..5), 3),
1619            ],
1620        )
1621        .await;
1622    }
1623
1624    #[fuchsia::test]
1625    async fn test_no_merge_seek_finds_lower_layer() {
1626        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1627        let items = [
1628            Item::new(TestKeyWithDefaultLayerKey(2..3), 1),
1629            Item::new(TestKeyWithDefaultLayerKey(0..1), 2),
1630        ];
1631        skip_lists[0].insert(items[0].clone()).expect("insert error");
1632        skip_lists[1].insert(items[1].clone()).expect("insert error");
1633        let mut merger = Merger::new(
1634            layer_ref_iter(&skip_lists),
1635            |_left, _right| MergeResult::EmitLeft,
1636            counters(),
1637        );
1638        let iter = merger
1639            .query(Query::FullRange(&items[1].key.search_key().unwrap()))
1640            .await
1641            .expect("seek failed");
1642
1643        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1644        assert_eq!((key, value), (&items[1].key, &items[1].value));
1645    }
1646
1647    #[fuchsia::test]
1648    async fn test_no_merge_seek_stops_at_exact_match() {
1649        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1650        let items = [
1651            Item::new(TestKeyWithDefaultLayerKey(2..3), 1),
1652            Item::new(TestKeyWithDefaultLayerKey(1..4), 2),
1653        ];
1654        skip_lists[0].insert(items[0].clone()).expect("insert error");
1655        skip_lists[1].insert(items[1].clone()).expect("insert error");
1656        let mut merger = Merger::new(
1657            layer_ref_iter(&skip_lists),
1658            |_left, _right| MergeResult::Other { emit: None, left: Discard, right: Keep },
1659            counters(),
1660        );
1661        let iter = merger
1662            .query(Query::FullRange(&items[0].key.search_key().unwrap()))
1663            .await
1664            .expect("seek failed");
1665
1666        // Seek should only search in the first skip list, so no merge should take place, and we'll
1667        // know if it has because we'll see a different value (2 rather than 1).
1668        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1669        assert_eq!((key, value), (&items[0].key, &items[0].value));
1670    }
1671
1672    #[fuchsia::test]
1673    async fn test_seek_less_than() {
1674        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1675        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
1676        skip_lists[0].insert(items[0].clone()).expect("insert error");
1677        skip_lists[1].insert(items[1].clone()).expect("insert error");
1678        // Search for a key before 1..1.
1679        let mut merger = Merger::new(
1680            layer_ref_iter(&skip_lists),
1681            |_left, _right| MergeResult::Other { emit: None, left: Discard, right: Keep },
1682            counters(),
1683        );
1684        let iter = merger
1685            .query(Query::FullRange(&TestKey(0..0).search_key().unwrap()))
1686            .await
1687            .expect("seek failed");
1688
1689        // This should find the 2..2 key because of our merge function.
1690        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1691        assert_eq!((key, value), (&items[1].key, &items[1].value));
1692    }
1693
1694    #[fuchsia::test]
1695    async fn test_seek_to_end() {
1696        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1697        let items = [Item::new(TestKey(1..1), 1), Item::new(TestKey(2..2), 2)];
1698        skip_lists[0].insert(items[0].clone()).expect("insert error");
1699        skip_lists[1].insert(items[1].clone()).expect("insert error");
1700        let mut merger = Merger::new(
1701            layer_ref_iter(&skip_lists),
1702            |_left, _right| MergeResult::Other { emit: None, left: Discard, right: Keep },
1703            counters(),
1704        );
1705        let iter = merger
1706            .query(Query::FullRange(&TestKey(3..3).search_key().unwrap()))
1707            .await
1708            .expect("seek failed");
1709
1710        assert!(iter.get().is_none());
1711    }
1712
1713    #[fuchsia::test]
1714    async fn test_merge_all_discarded() {
1715        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1716        let items = [Item::new(TestKey(1..2), 1), Item::new(TestKey(2..3), 2)];
1717        skip_lists[0].insert(items[1].clone()).expect("insert error");
1718        skip_lists[1].insert(items[0].clone()).expect("insert error");
1719        let mut merger = Merger::new(
1720            layer_ref_iter(&skip_lists),
1721            |_left, _right| MergeResult::Other { emit: None, left: Discard, right: Discard },
1722            counters(),
1723        );
1724        let iter = merger.query(Query::FullScan).await.expect("seek failed");
1725        assert!(iter.get().is_none());
1726    }
1727
1728    #[fuchsia::test]
1729    async fn test_overlapping_keys() {
1730        let skip_lists =
1731            [SkipListLayer::new(100), SkipListLayer::new(100), SkipListLayer::new(100)];
1732        let items = [
1733            Item::new(TestKey(0..10), 1),
1734            Item::new(TestKey(0..20), 2),
1735            Item::new(TestKey(0..30), 3),
1736        ];
1737        skip_lists[0].insert(items[0].clone()).expect("insert error");
1738        skip_lists[1].insert(items[1].clone()).expect("insert error");
1739        skip_lists[2].insert(items[2].clone()).expect("insert error");
1740        let mut merger = Merger::new(
1741            layer_ref_iter(&skip_lists),
1742            |left, right| {
1743                if left.key().0.end <= right.key().0.start {
1744                    MergeResult::EmitLeft
1745                } else {
1746                    // Hardcoded rule for specific test data to test containment split.
1747                    if left.key() == &TestKey(0..30) && right.key() == &TestKey(10..20) {
1748                        MergeResult::Other {
1749                            emit: Some(Item::new(TestKey(0..10), 1).boxed()),
1750                            left: Replace(Item::new(TestKey(10..30), 1).boxed()),
1751                            right: Keep,
1752                        }
1753                    } else {
1754                        MergeResult::Other {
1755                            emit: None,
1756                            left: Keep,
1757                            right: Replace(
1758                                Item::new(TestKey(left.key().0.end..right.key().0.end), 1).boxed(),
1759                            ),
1760                        }
1761                    }
1762                }
1763            },
1764            counters(),
1765        );
1766        let mut iter = merger
1767            .query(Query::FullRange(&TestKey(0..1).search_key().unwrap()))
1768            .await
1769            .expect("seek failed");
1770        assert_eq!(iter.pending_iterators_len(), 2);
1771        let ItemRef { key, .. } = iter.get().expect("missing item");
1772        assert_eq!(key, &TestKey(0..10));
1773        iter.advance().await.expect("advance failed");
1774        assert_eq!(iter.pending_iterators_len(), 1);
1775        let ItemRef { key, .. } = iter.get().expect("missing item");
1776        assert_eq!(key, &TestKey(10..20));
1777        iter.advance().await.expect("advance failed");
1778        assert_eq!(iter.pending_iterators_len(), 0);
1779        let ItemRef { key, .. } = iter.get().expect("missing item");
1780        assert_eq!(key, &TestKey(20..30));
1781        iter.advance().await.expect("advance failed");
1782        assert_eq!(iter.pending_iterators_len(), 0);
1783        assert_eq!(iter.get(), None);
1784    }
1785
1786    async fn write_layer<K: Key, V: Value>(items: Vec<Item<K, V>>) -> Arc<dyn Layer<K, V>> {
1787        let object = Arc::new(FakeObject::new());
1788        let write_handle = FakeObjectHandle::new(object.clone());
1789        let mut writer = PersistentLayerWriter::<_, K, V>::new(
1790            Writer::new(&write_handle).await,
1791            items.len(),
1792            BlockSize::SIZE_512B,
1793        )
1794        .await
1795        .expect("PersistentLayerWriter::new failed");
1796        for item in items {
1797            writer.write(item.as_item_ref()).await.expect("write failed");
1798        }
1799        writer.complete().await.expect("flush failed");
1800        PersistentLayer::open(FakeObjectHandle::new(object))
1801            .await
1802            .expect("open_persistent_layer failed")
1803    }
1804
1805    fn merge_sum(
1806        left: &MergeLayerIterator<'_, i32, i32>,
1807        right: &MergeLayerIterator<'_, i32, i32>,
1808    ) -> MergeResult<i32, i32> {
1809        // Sum matching keys.
1810        if left.key() == right.key() {
1811            MergeResult::Other {
1812                emit: None,
1813                left: Discard,
1814                right: Replace(Item::new(left.key().clone(), left.value() + right.value()).boxed()),
1815            }
1816        } else {
1817            MergeResult::EmitLeft
1818        }
1819    }
1820
1821    #[fuchsia::test]
1822    async fn test_merge_bloom_filters_point_query() {
1823        let layer_0_items = vec![Item::new(1, 1), Item::new(2, 1)];
1824        let layer_1_items = vec![Item::new(2, 1), Item::new(4, 1)];
1825        let items = [Item::new(1, 1), Item::new(2, 2), Item::new(4, 1)];
1826        let layers: [Arc<dyn Layer<i32, i32>>; 2] =
1827            [write_layer(layer_0_items).await, write_layer(layer_1_items).await];
1828        let mut merger = Merger::new(dyn_layer_ref_iter(&layers), merge_sum, counters());
1829
1830        {
1831            // Key that exists in layer 0 only
1832            let iter = merger.query(Query::Point(&1)).await.expect("seek failed");
1833            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1834            assert_eq!((key, *value), (&items[0].key, items[0].value));
1835        }
1836        {
1837            // Key that exists in layer 1 and 2
1838            let iter = merger.query(Query::Point(&2)).await.expect("seek failed");
1839            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1840            assert_eq!((key, *value), (&items[1].key, items[1].value));
1841        }
1842        {
1843            // Key that exists in layer 1
1844            let iter = merger.query(Query::Point(&4)).await.expect("seek failed");
1845            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1846            assert_eq!((key, *value), (&items[2].key, items[2].value));
1847        }
1848        {
1849            // Key that doesn't exist at all
1850            let iter = merger.query(Query::Point(&400)).await.expect("seek failed");
1851            assert!(iter.get().is_none());
1852        }
1853    }
1854
1855    #[fuchsia::test]
1856    async fn test_merge_bloom_filters_limited_range() {
1857        // NB: This test uses ObjectKey so that we don't have to reimplement the complex merging
1858        // logic for range-like keys.
1859        let layer_0_items = vec![Item::new(
1860            ObjectKey::extent(0, AttributeId::TEST_ID, 0..2048),
1861            ObjectValue::extent(0, VOLUME_DATA_KEY_ID),
1862        )];
1863        let layer_1_items = vec![
1864            Item::new(
1865                ObjectKey::extent(0, AttributeId::TEST_ID, 1024..4096),
1866                ObjectValue::extent(32768, VOLUME_DATA_KEY_ID),
1867            ),
1868            Item::new(
1869                ObjectKey::extent(0, AttributeId::TEST_ID, 16384..17408),
1870                ObjectValue::extent(65536, VOLUME_DATA_KEY_ID),
1871            ),
1872        ];
1873        let items = [
1874            Item::new(
1875                ObjectKey::extent(0, AttributeId::TEST_ID, 0..2048),
1876                ObjectValue::extent(0, VOLUME_DATA_KEY_ID),
1877            ),
1878            Item::new(
1879                ObjectKey::extent(0, AttributeId::TEST_ID, 2048..4096),
1880                ObjectValue::extent(33792, VOLUME_DATA_KEY_ID),
1881            ),
1882            Item::new(
1883                ObjectKey::extent(0, AttributeId::TEST_ID, 16384..17408),
1884                ObjectValue::extent(65536, VOLUME_DATA_KEY_ID),
1885            ),
1886        ];
1887        let layers: [Arc<dyn Layer<ObjectKey, ObjectValue>>; 2] =
1888            [write_layer(layer_0_items).await, write_layer(layer_1_items).await];
1889        let mut merger =
1890            Merger::new(dyn_layer_ref_iter(&layers), object_store::merge::merge, counters());
1891
1892        {
1893            // Range contains just keys in layer 1
1894            let mut iter = merger
1895                .query(Query::LimitedRange(&ObjectKey::extent(
1896                    0,
1897                    AttributeId::TEST_ID,
1898                    16384..16386,
1899                )))
1900                .await
1901                .expect("seek failed");
1902            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1903            assert_eq!(key, &items[2].key);
1904            assert_eq!(value, &items[2].value);
1905            iter.advance().await.expect("advance");
1906            assert!(iter.get().is_none());
1907        }
1908        {
1909            // Range contains keys in layer 0 and 1
1910            let mut iter = merger
1911                .query(Query::LimitedRange(&ObjectKey::extent(0, AttributeId::TEST_ID, 0..4096)))
1912                .await
1913                .expect("seek failed");
1914            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1915            assert_eq!(key, &items[0].key);
1916            assert_eq!(value, &items[0].value);
1917            iter.advance().await.expect("advance");
1918            let ItemRef { key, value, .. } = iter.get().expect("missing item");
1919            assert_eq!(key, &items[1].key);
1920            assert_eq!(value, &items[1].value);
1921            iter.advance().await.expect("advance");
1922        }
1923        {
1924            // Range contains no keys
1925            let mut iter = merger
1926                .query(Query::LimitedRange(&ObjectKey::extent(
1927                    0,
1928                    AttributeId::TEST_ID,
1929                    8192..12288,
1930                )))
1931                .await
1932                .expect("seek failed");
1933            let ItemRef { key, .. } = iter.get().expect("missing item");
1934            assert_eq!(key, &items[2].key);
1935            iter.advance().await.expect("advance");
1936            assert!(iter.get().is_none());
1937        }
1938    }
1939
1940    #[fuchsia::test]
1941    async fn test_merge_bloom_filters_full_range() {
1942        let layer_0_items = vec![Item::new(1, 1), Item::new(2, 1)];
1943        let layer_1_items = vec![Item::new(2, 1), Item::new(4, 1)];
1944        let items = [Item::new(1, 1), Item::new(2, 2), Item::new(4, 1)];
1945        let layers: [Arc<dyn Layer<i32, i32>>; 2] =
1946            [write_layer(layer_0_items).await, write_layer(layer_1_items).await];
1947        let mut merger = Merger::new(dyn_layer_ref_iter(&layers), merge_sum, counters());
1948
1949        let mut iter = merger.query(Query::FullRange(&0)).await.expect("seek failed");
1950        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1951        assert_eq!((key, *value), (&items[0].key, items[0].value));
1952        iter.advance().await.expect("advance");
1953        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1954        assert_eq!((key, *value), (&items[1].key, items[1].value));
1955        iter.advance().await.expect("advance");
1956        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1957        assert_eq!((key, *value), (&items[2].key, items[2].value));
1958        iter.advance().await.expect("advance");
1959        assert!(iter.get().is_none());
1960    }
1961
1962    #[fuchsia::test]
1963    async fn test_merge_bloom_filters_full_scan() {
1964        let layer_0_items = vec![Item::new(1, 1), Item::new(2, 1)];
1965        let layer_1_items = vec![Item::new(2, 1), Item::new(4, 1)];
1966        let items = [Item::new(1, 1), Item::new(2, 2), Item::new(4, 1)];
1967        let layers: [Arc<dyn Layer<i32, i32>>; 2] =
1968            [write_layer(layer_0_items).await, write_layer(layer_1_items).await];
1969        let mut merger = Merger::new(dyn_layer_ref_iter(&layers), merge_sum, counters());
1970
1971        let mut iter = merger.query(Query::FullScan).await.expect("seek failed");
1972        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1973        assert_eq!((key, *value), (&items[0].key, items[0].value));
1974        iter.advance().await.expect("advance");
1975        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1976        assert_eq!((key, *value), (&items[1].key, items[1].value));
1977        iter.advance().await.expect("advance");
1978        let ItemRef { key, value, .. } = iter.get().expect("missing item");
1979        assert_eq!((key, *value), (&items[2].key, items[2].value));
1980        iter.advance().await.expect("advance");
1981        assert!(iter.get().is_none());
1982    }
1983
1984    #[fuchsia::test]
1985    async fn test_optimized_merge_lazy_pull() {
1986        let skip_lists = [SkipListLayer::new(100), SkipListLayer::new(100)];
1987        let items_top = [
1988            Item::new(TestKey(10..20), 1),
1989            Item::new(TestKey(20..30), 2),
1990            Item::new(TestKey(30..40), 3),
1991        ];
1992        let items_bottom = [Item::new(TestKey(40..50), 4)];
1993
1994        for item in &items_top {
1995            skip_lists[0].insert(item.clone()).expect("insert error");
1996        }
1997        for item in &items_bottom {
1998            skip_lists[1].insert(item.clone()).expect("insert error");
1999        }
2000
2001        let mut merger =
2002            Merger::new(layer_ref_iter(&skip_lists), |_, _| MergeResult::EmitLeft, counters());
2003        merger.set_trace(true);
2004
2005        let mut iter =
2006            merger.query(Query::LimitedRange(&TestKey(10..20))).await.expect("seek failed");
2007        assert_eq!(iter.pending_iterators_len(), 1);
2008
2009        let ItemRef { key, .. } = iter.get().expect("missing item");
2010        assert_eq!(key, &TestKey(10..20));
2011
2012        iter.advance().await.expect("advance failed");
2013        assert_eq!(iter.pending_iterators_len(), 1);
2014
2015        let ItemRef { key, .. } = iter.get().expect("missing item");
2016        assert_eq!(key, &TestKey(20..30));
2017
2018        iter.advance().await.expect("advance failed");
2019        assert_eq!(iter.pending_iterators_len(), 1);
2020
2021        let ItemRef { key, .. } = iter.get().expect("missing item");
2022        assert_eq!(key, &TestKey(30..40));
2023
2024        iter.advance().await.expect("advance failed");
2025        assert_eq!(iter.pending_iterators_len(), 0);
2026
2027        let ItemRef { key, .. } = iter.get().expect("missing item");
2028        assert_eq!(key, &TestKey(40..50));
2029
2030        iter.advance().await.expect("advance failed");
2031        assert!(iter.get().is_none());
2032    }
2033
2034    #[fuchsia::test]
2035    async fn test_excessive_hash_partitions_counter() {
2036        use std::sync::atomic::Ordering;
2037
2038        let layer_0_items: Vec<_> = (0..100)
2039            .map(|i| {
2040                Item::new(
2041                    ObjectKey::extent(0, AttributeId::TEST_ID, i * 512..(i + 1) * 512),
2042                    ObjectValue::extent(0, VOLUME_DATA_KEY_ID),
2043                )
2044            })
2045            .collect();
2046        let layers: [Arc<dyn Layer<ObjectKey, ObjectValue>>; 1] =
2047            [write_layer(layer_0_items).await];
2048        let counters = counters();
2049        let mut merger =
2050            Merger::new(dyn_layer_ref_iter(&layers), object_store::merge::merge, counters.clone());
2051
2052        // A query with <= MAX_HASH_PARTITIONS (e.g. 1 partition: 0..4096)
2053        merger
2054            .query(Query::LimitedRange(&ObjectKey::extent(0, AttributeId::TEST_ID, 0..4096)))
2055            .await
2056            .expect("seek failed");
2057        assert_eq!(counters.excessive_hash_partitions.load(Ordering::Relaxed), 0);
2058
2059        let key_large = ObjectKey::extent(0, AttributeId::TEST_ID, 0..10485760);
2060        assert_eq!(layers[0].maybe_contains_key(&key_large), MaybeContainsKey::RangeKeyTooLarge);
2061
2062        // A query with > MAX_HASH_PARTITIONS (e.g. 10MB range: 0..10485760 = 10 partitions)
2063        merger.query(Query::LimitedRange(&key_large)).await.expect("seek failed");
2064        assert_eq!(counters.excessive_hash_partitions.load(Ordering::Relaxed), 1);
2065    }
2066}