1use 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 Keep,
24
25 Discard,
27
28 Replace(BoxedItem<K, V>),
31}
32
33#[derive(Debug, Eq, PartialEq)]
34pub enum MergeResult<K, V> {
35 EmitLeft,
38
39 Other { emit: Option<BoxedItem<K, V>>, left: ItemOp<K, V>, right: ItemOp<K, V> },
60}
61
62pub 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
81unsafe impl<K, V> Send for RawIterator<'_, K, V> {}
84
85pub struct MergeLayerIterator<'a, K, V> {
88 layer: Option<&'a dyn Layer<K, V>>,
89
90 iter: RawIterator<'a, K, V>,
92
93 pub layer_index: u16,
95
96 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 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 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
189impl<K: OrdLowerBound, V> Ord for MergeLayerIterator<'_, K, V> {
191 fn cmp(&self, other: &Self) -> Ordering {
192 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
208enum 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
237pub struct Merger<'a, K, V> {
239 iterators: Vec<MergeLayerIterator<'a, K, V>>,
241
242 merge_fn: MergeFn<K, V>,
244
245 trace: bool,
247
248 counters: Arc<lsm_tree::TreeCounters>,
250}
251
252#[derive(Debug, Clone)]
259pub enum Query<'a, K: Key + LayerKey + OrdLowerBound> {
260 Point(&'a K),
265
266 LimitedRange(&'a K),
274
275 FullRange(&'a K),
281
282 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 #[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 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
395pub struct MergerIterator<'a, 'b, K, V> {
398 merge_fn: MergeFn<K, V>,
399
400 pending_iterators: Vec<&'a mut MergeLayerIterator<'b, K, V>>,
402
403 heap: BinaryHeap<&'a mut MergeLayerIterator<'b, K, V>>,
405
406 item: CurrentItem<'a, 'b, K, V>,
408
409 trace: bool,
411
412 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 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 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 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 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 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 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 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 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 } else {
608 iterator.replace(existing_item);
612
613 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
671pub(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 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 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, 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 } 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 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 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 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 #[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 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 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 #[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 #[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 #[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 let ItemRef { key, value, .. } = iter.get().expect("missing item");
1306 assert_eq!((key, value), (&items[0].key, &items[0].value));
1307 }
1308
1309 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 merger.query(Query::LimitedRange(&key_large)).await.expect("seek failed");
2064 assert_eq!(counters.excessive_hash_partitions.load(Ordering::Relaxed), 1);
2065 }
2066}