1use std::ops::Range;
18use std::pin::Pin;
19use std::sync::Arc;
20use std::sync::atomic::AtomicU64;
21use std::sync::atomic::Ordering;
22
23use futures::Stream;
24use futures::StreamExt;
25use futures::stream::BoxStream;
26use futures::stream::Peekable;
27use mysten_common::zip_debug_eq::ZipDebugEqIteratorExt;
28use roaring::RoaringBitmap;
29
30use super::BitmapBucketSource;
31use super::BitmapQuery;
32use super::BucketItem;
33use super::BucketStream;
34use super::DedupedQuery;
35use super::LeafHead;
36use super::LeafStop;
37use super::ScanDirection;
38use super::ScanStop;
39use super::SkipPolicy;
40use super::Watermarked;
41use super::WatermarkedBucketStream;
42use super::advance_in_direction;
43use super::bound_in_direction;
44use super::bucket_edges;
45use super::build_term_specs;
46use super::collapse;
47use super::count_on_floor_refs;
48use super::eval_term_at_bucket;
49use super::frontier_advanced;
50use super::leaf_skip_targets;
51use super::recompute_unreferenced;
52use super::strictly_before;
53use super::take_snapshot_bitmap;
54
55#[derive(Clone, Copy, Debug, Eq, PartialEq)]
61pub struct BitmapScanMetrics {
62 pub buckets_evaluated: u64,
65 pub buckets_discarded: u64,
67 pub leaf_seeks: u64,
69}
70
71#[derive(Clone)]
75pub(crate) struct BitmapScanBudget {
76 initial: u64,
77 remaining: Arc<AtomicU64>,
78 discarded: Arc<AtomicU64>,
79 seeks: Arc<AtomicU64>,
80}
81
82impl BitmapScanBudget {
83 pub(crate) fn new(initial: u64) -> Self {
84 Self {
85 initial,
86 remaining: Arc::new(AtomicU64::new(initial)),
87 discarded: Arc::new(AtomicU64::new(0)),
88 seeks: Arc::new(AtomicU64::new(0)),
89 }
90 }
91
92 fn try_take(&self) -> bool {
94 self.remaining
95 .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |b| {
96 if b == 0 { None } else { Some(b - 1) }
97 })
98 .is_ok()
99 }
100
101 fn take_first(&self) {
113 let _ = self
114 .remaining
115 .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |b| {
116 Some(b.saturating_sub(1))
117 });
118 }
119
120 fn buckets_evaluated(&self) -> u64 {
121 self.initial
122 .saturating_sub(self.remaining.load(Ordering::SeqCst))
123 }
124 fn note_discarded(&self) {
125 self.discarded.fetch_add(1, Ordering::Relaxed);
126 }
127
128 fn note_seek(&self) {
129 self.seeks.fetch_add(1, Ordering::Relaxed);
130 }
131
132 fn discarded(&self) -> u64 {
133 self.discarded.load(Ordering::Relaxed)
134 }
135
136 fn seeks(&self) -> u64 {
137 self.seeks.load(Ordering::Relaxed)
138 }
139}
140
141struct ObserveOnDrop<F: FnOnce(BitmapScanMetrics) + Send + 'static> {
145 budget: BitmapScanBudget,
146 callback: Option<F>,
147}
148
149impl<F: FnOnce(BitmapScanMetrics) + Send + 'static> Drop for ObserveOnDrop<F> {
150 fn drop(&mut self) {
151 if let Some(cb) = self.callback.take() {
152 cb(BitmapScanMetrics {
153 buckets_evaluated: self.budget.buckets_evaluated(),
154 buckets_discarded: self.budget.discarded(),
155 leaf_seeks: self.budget.seeks(),
156 });
157 }
158 }
159}
160
161pub fn eval_bitmap_query_stream<S, F>(
171 source: S,
172 query: BitmapQuery,
173 range: Range<u64>,
174 bucket_size: u64,
175 direction: ScanDirection,
176 budget: u64,
177 policy: SkipPolicy,
178 on_metrics: F,
179) -> BoxStream<'static, Result<Watermarked<u64>, ScanStop>>
180where
181 S: BitmapBucketSource,
182 F: FnOnce(BitmapScanMetrics) + Send + 'static,
183{
184 let leaves = query.unique_leaf_count();
185 if (budget as usize) < leaves {
186 return async_stream::stream! {
190 yield Err(ScanStop::Fault(anyhow::anyhow!(
192 "bitmap scan budget {budget} is insufficient for {leaves} leaf streams; \
193 server is misconfigured"
194 )));
195 }
196 .boxed();
197 }
198 let budget = BitmapScanBudget::new(budget);
199 let bucket_stream = eval_bitmap_query_bucket_stream(
200 source,
201 query,
202 range.clone(),
203 bucket_size,
204 direction,
205 budget.clone(),
206 policy,
207 );
208 let inner = flatten_watermarked_buckets(bucket_stream, range.clone(), bucket_size, direction);
209 let guard = ObserveOnDrop {
213 budget,
214 callback: Some(on_metrics),
215 };
216 async_stream::stream! {
217 let _guard = guard;
218 futures::pin_mut!(inner);
219 while let Some(item) = inner.next().await {
220 yield item;
221 }
222 }
223 .boxed()
224}
225
226async fn peek_leaf<S>(
229 mut leaf: Pin<&mut Peekable<S>>,
230 idx: usize,
231 bucket_size: u64,
232 range: &Range<u64>,
233 direction: ScanDirection,
234 terminus: u64,
235) -> (usize, LeafHead, Option<u64>)
236where
237 S: Stream<Item = BucketItem>,
238{
239 match leaf.as_mut().peek().await {
240 Some(Ok((bucket, _))) => {
241 let (pre, _post) = bucket_edges(*bucket, bucket_size, range, direction);
242 (idx, LeafHead::Bucket(*bucket), Some(pre))
243 }
244 None => (idx, LeafHead::Eof, Some(terminus)),
245 Some(Err(_)) => (idx, LeafHead::Error, None),
246 }
247}
248
249pub(crate) fn eval_bitmap_query_bucket_stream<S>(
256 source: S,
257 query: BitmapQuery,
258 range: Range<u64>,
259 bucket_size: u64,
260 direction: ScanDirection,
261 budget: BitmapScanBudget,
262 policy: SkipPolicy,
263) -> WatermarkedBucketStream
264where
265 S: BitmapBucketSource,
266{
267 let DedupedQuery {
268 keys: unique_keys,
269 mut terms,
270 } = build_term_specs(query.terms);
271 let leaf_count = unique_keys.len();
272 let terminus = if direction.is_ascending() {
273 range.end
274 } else {
275 range.start
276 };
277 let request_floor = if direction.is_ascending() {
278 range.start
279 } else {
280 range.end
281 };
282
283 async_stream::try_stream! {
284 let mut leaves: Vec<Peekable<BucketStream>> = Vec::with_capacity(leaf_count);
285 for key in &unique_keys {
286 let raw = source.scan_bucket_stream(key.clone(), range.clone(), direction);
287 leaves.push(
288 budgeted_bucket_stream(raw, budget.clone(), true)
289 .boxed()
290 .peekable(),
291 );
292 }
293 let mut unreferenced = vec![false; leaf_count];
294 let mut front = vec![request_floor; leaf_count];
295 let mut drained = vec![0u64; leaf_count];
296 let mut progress_frontier = Some(request_floor);
297
298 loop {
299 let mut peeks = Vec::new();
300 for (i, leaf) in leaves.iter_mut().enumerate() {
301 if !unreferenced[i] {
302 peeks.push(peek_leaf(
303 Pin::new(leaf),
304 i,
305 bucket_size,
306 &range,
307 direction,
308 terminus,
309 ));
310 }
311 }
312 let results = futures::future::join_all(peeks).await;
313 let mut class: Vec<Option<LeafHead>> = (0..leaf_count).map(|_| None).collect();
314 for (i, head, scanned_to) in results {
315 if let Some(position) = scanned_to {
316 front[i] = advance_in_direction(front[i], position, direction);
317 }
318 class[i] = Some(head);
319 }
320
321 for term in terms.iter_mut() {
322 if !term.unsatisfiable
323 && term
324 .includes
325 .iter()
326 .any(|&i| matches!(class[i], Some(LeafHead::Eof)))
327 {
328 term.unsatisfiable = true;
329 }
330 }
331 recompute_unreferenced(&terms, &class, &mut unreferenced);
332
333 let targets = leaf_skip_targets(&terms, &class, &unreferenced, direction);
334 for (i, target) in targets.iter().enumerate() {
335 if let Some(target) = target {
336 let (pre, _) = bucket_edges(*target, bucket_size, &range, direction);
337 front[i] = advance_in_direction(front[i], pre, direction);
338 }
339 }
340
341 let mut errors: Vec<LeafStop> = Vec::new();
342 for i in 0..leaf_count {
343 if !unreferenced[i] && matches!(class[i], Some(LeafHead::Error)) {
344 match Pin::new(&mut leaves[i]).next().await {
345 Some(Err(error)) => errors.push(error),
346 _ => unreachable!("peek classified Error"),
347 }
348 }
349 }
350
351 let active: Vec<usize> = (0..leaf_count).filter(|&i| !unreferenced[i]).collect();
352 if active.is_empty() {
353 return;
354 }
355
356 let floor_pos = active
357 .iter()
358 .map(|&i| front[i])
359 .reduce(|a, b| bound_in_direction(a, b, direction))
360 .expect("active non-empty");
361 let collapsed = (!errors.is_empty()).then(|| collapse(errors, floor_pos));
362 let scan_limited = matches!(collapsed, Some(ScanStop::ScanLimit { .. }));
363 if !scan_limited && frontier_advanced(progress_frontier, floor_pos, direction) {
364 yield Watermarked::Watermark(floor_pos);
365 progress_frontier = Some(floor_pos);
366 }
367 if let Some(stop) = collapsed {
368 Err(stop)?;
369 }
370
371 let lagging: Vec<usize> = active
375 .iter()
376 .copied()
377 .filter(|&i| targets[i].is_some())
378 .collect();
379 for i in 0..leaf_count {
383 if targets[i].is_none() {
384 drained[i] = 0;
385 }
386 }
387
388 let eval_bucket = active
392 .iter()
393 .filter(|&&i| targets[i].is_none())
394 .filter_map(|&i| match class[i] {
395 Some(LeafHead::Bucket(bucket)) => Some(bucket),
396 _ => None,
397 })
398 .reduce(|a, b| bound_in_direction(a, b, direction))
399 .expect("at least one active leaf has a ready bucket");
400 let lagging_target = lagging
405 .iter()
406 .filter_map(|&i| targets[i])
407 .reduce(|a, b| bound_in_direction(a, b, direction));
408 let evaluate = lagging_target
409 .is_none_or(|target| strictly_before(eval_bucket, target, direction));
410
411 if evaluate {
412 let (_, post) = bucket_edges(eval_bucket, bucket_size, &range, direction);
413 let mut snapshot: Vec<Option<RoaringBitmap>> =
417 (0..leaf_count).map(|_| None).collect();
418 let mut on_floor = vec![false; leaf_count];
419 for i in 0..leaf_count {
420 if !unreferenced[i]
421 && targets[i].is_none()
422 && matches!(class[i], Some(LeafHead::Bucket(b)) if b == eval_bucket)
423 {
424 on_floor[i] = true;
425 front[i] = advance_in_direction(front[i], post, direction);
426 snapshot[i] = match Pin::new(&mut leaves[i]).next().await {
427 Some(Ok((_, bitmap))) => Some(bitmap),
428 _ => None,
429 };
430 }
431 }
432 let mut remaining_refs = count_on_floor_refs(&terms, &on_floor);
435 let mut result: Option<RoaringBitmap> = None;
436 for term in &terms {
437 if term.unsatisfiable {
438 continue;
439 }
440 let includes = term
441 .includes
442 .iter()
443 .map(|&i| {
444 take_snapshot_bitmap(
445 &mut snapshot,
446 &mut remaining_refs,
447 &on_floor,
448 i,
449 )
450 })
451 .collect();
452 let excludes = term
453 .excludes
454 .iter()
455 .map(|&i| {
456 take_snapshot_bitmap(
457 &mut snapshot,
458 &mut remaining_refs,
459 &on_floor,
460 i,
461 )
462 })
463 .collect();
464 if let Some(bitmap) = eval_term_at_bucket(includes, excludes) {
465 result = Some(match result {
466 None => bitmap,
467 Some(acc) => acc | bitmap,
468 });
469 }
470 }
471 if let Some(bitmap) = result {
472 yield Watermarked::Item((eval_bucket, bitmap));
473 }
474 if frontier_advanced(progress_frontier, post, direction) {
475 yield Watermarked::Watermark(post);
476 progress_frontier = Some(post);
477 }
478 }
479
480 let catch_up = leaves.iter_mut().zip_debug_eq(drained.iter_mut())
484 .enumerate()
485 .filter_map(|(i, (leaf, drained))| {
486 let target = targets[i]?;
487 let source = source.clone();
488 let key = unique_keys[i].clone();
489 let range = range.clone();
490 let budget = budget.clone();
491 Some(async move {
492 loop {
496 let dead_row = matches!(
497 Pin::new(&mut *leaf).peek().await,
498 Some(Ok((bucket, _)))
499 if strictly_before(*bucket, target, direction)
500 );
501 if !dead_row {
502 break;
503 }
504 if policy
508 .drain_probe_rows
509 .is_some_and(|probe| *drained >= u64::from(probe.get()))
510 {
511 budget.note_discarded();
515 let narrowed = match direction {
518 ScanDirection::Ascending => {
519 target.saturating_mul(bucket_size).max(range.start)
520 ..range.end
521 }
522 ScanDirection::Descending => {
523 range.start
524 ..target
525 .saturating_add(1)
526 .saturating_mul(bucket_size)
527 .min(range.end)
528 }
529 };
530 let raw = source.scan_bucket_stream(key, narrowed, direction);
531 *leaf = budgeted_bucket_stream(raw, budget.clone(), false)
535 .boxed()
536 .peekable();
537 budget.note_seek();
538 break;
539 }
540 let discarded = Pin::new(&mut *leaf).next().await;
543 debug_assert!(matches!(discarded, Some(Ok(_))));
544 *drained += 1;
545 budget.note_discarded();
546 }
547 })
548 });
549 futures::future::join_all(catch_up).await;
550 }
551 }
552 .boxed()
553}
554
555fn budgeted_bucket_stream<S>(
560 inner: S,
561 budget: BitmapScanBudget,
562 reserve_first: bool,
563) -> impl Stream<Item = BucketItem> + Send + 'static
564where
565 S: Stream<Item = BucketItem> + Send + 'static,
566{
567 async_stream::try_stream! {
568 futures::pin_mut!(inner);
569 let mut first = reserve_first;
570 while let Some(item) = inner.next().await {
571 let item = item?;
572 if first {
573 budget.take_first();
574 first = false;
575 } else if !budget.try_take() {
576 Err(LeafStop::BudgetExhausted)?;
577 }
578 yield item;
579 }
580 }
581}
582
583pub fn flatten_watermarked_buckets<S, E>(
587 stream: S,
588 range: Range<u64>,
589 bucket_size: u64,
590 direction: ScanDirection,
591) -> impl Stream<Item = Result<Watermarked<u64>, E>> + Send + 'static
592where
593 S: Stream<Item = Result<Watermarked<(u64, RoaringBitmap)>, E>> + Send + 'static,
594 E: Send + 'static,
595{
596 async_stream::try_stream! {
597 if range.is_empty() {
598 return;
599 }
600 let start_bucket = range.start / bucket_size;
601 let end_bucket = (range.end - 1) / bucket_size;
602 futures::pin_mut!(stream);
603 while let Some(item) = stream.next().await {
604 match item? {
605 Watermarked::Watermark(p) => yield Watermarked::Watermark(p),
606 Watermarked::Item((bucket_id, bitmap)) => {
607 let bucket_start = bucket_id * bucket_size;
608 let is_first = bucket_id == start_bucket;
609 let is_last = bucket_id == end_bucket;
610 let lo = if is_first {
611 (range.start - bucket_start) as u32
612 } else {
613 0
614 };
615 let hi = if is_last {
616 ((range.end - bucket_start).min(bucket_size)) as u32
617 } else {
618 bucket_size as u32
619 };
620
621 if direction.is_ascending() {
622 for bit in bitmap.iter() {
623 if bit >= lo && bit < hi {
624 yield Watermarked::Item(bucket_start + bit as u64);
625 }
626 }
627 } else {
628 for bit in bitmap.iter().rev() {
629 if bit >= lo && bit < hi {
630 yield Watermarked::Item(bucket_start + bit as u64);
631 }
632 }
633 }
634 }
635 }
636 }
637 }
638}
639
640#[cfg(test)]
641mod tests {
642 use std::collections::BTreeMap;
643 use std::sync::Arc;
644
645 use futures::TryStreamExt;
646 use futures::stream;
647
648 use super::*;
649 use crate::bitmap_query::BitmapLiteral;
650 use crate::bitmap_query::BitmapTerm;
651 use crate::bitmap_query::BucketStream;
652
653 const BUCKET_SIZE: u64 = 100_000;
654 type TestBuckets = BTreeMap<Vec<u8>, Vec<(u64, Vec<u32>)>>;
655
656 #[derive(Clone)]
657 struct TestBucketSource {
658 buckets: Arc<TestBuckets>,
659 }
660
661 impl BitmapBucketSource for TestBucketSource {
662 fn scan_bucket_stream(
663 &self,
664 dimension_key: Vec<u8>,
665 range: Range<u64>,
666 direction: ScanDirection,
667 ) -> BucketStream {
668 let mut buckets = self
669 .buckets
670 .get(&dimension_key)
671 .cloned()
672 .unwrap_or_default();
673 if range.is_empty() {
674 buckets.clear();
675 } else {
676 let first_bucket = range.start / BUCKET_SIZE;
677 let last_bucket = (range.end - 1) / BUCKET_SIZE;
678 buckets.retain(|(bucket, _)| first_bucket <= *bucket && *bucket <= last_bucket);
679 }
680 if matches!(direction, ScanDirection::Descending) {
681 buckets.reverse();
682 }
683 make_bucket_stream(
684 buckets
685 .iter()
686 .map(|(bucket_id, bits)| (*bucket_id, bits.as_slice()))
687 .collect(),
688 )
689 }
690 }
691
692 fn make_bitmap(bits: &[u32]) -> RoaringBitmap {
693 let mut bm = RoaringBitmap::new();
694 for &b in bits {
695 bm.insert(b);
696 }
697 bm
698 }
699
700 fn make_bucket_stream(items: Vec<(u64, &[u32])>) -> BucketStream {
701 let items: Vec<BucketItem> = items
702 .into_iter()
703 .map(|(bid, bits)| Ok((bid, make_bitmap(bits))))
704 .collect();
705 stream::iter(items).boxed()
706 }
707
708 async fn drain_marked(
711 stream: BoxStream<'static, Result<Watermarked<u64>, ScanStop>>,
712 ) -> Result<(Vec<u64>, Vec<u64>), ScanStop> {
713 let all: Vec<Watermarked<u64>> = stream.try_collect().await?;
714 let mut items = Vec::new();
715 let mut watermarks = Vec::new();
716 for m in all {
717 match m {
718 Watermarked::Item(v) => items.push(v),
719 Watermarked::Watermark(f) => watermarks.push(f),
720 }
721 }
722 Ok((items, watermarks))
723 }
724
725 fn test_key(value: &[u8]) -> Vec<u8> {
726 crate::dimensions::encode_dimension_key(crate::dimensions::IndexDimension::Sender, value)
727 }
728
729 fn include(value: &[u8]) -> BitmapLiteral {
730 BitmapLiteral::include(test_key(value)).unwrap()
731 }
732
733 fn exclude(value: &[u8]) -> BitmapLiteral {
734 BitmapLiteral::exclude(test_key(value)).unwrap()
735 }
736
737 #[tokio::test]
743 async fn nested_term_starvation_bundles_frontier_in_scan_limit() {
744 let source = TestBucketSource {
745 buckets: Arc::new(BTreeMap::from([
746 (
747 test_key(b"a"),
748 vec![
749 (0, vec![1]),
750 (1, vec![1]),
751 (2, vec![1]),
752 (3, vec![1]),
753 (4, vec![1]),
754 ],
755 ),
756 (test_key(b"b"), vec![(50, vec![1])]),
757 (test_key(b"c"), vec![(40, vec![7])]),
758 ])),
759 };
760 let query = BitmapQuery::new(vec![
761 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
762 BitmapTerm::new(vec![include(b"c")]).unwrap(),
763 ])
764 .unwrap();
765
766 let stream = eval_bitmap_query_stream(
768 source,
769 query,
770 0..(60 * BUCKET_SIZE),
771 BUCKET_SIZE,
772 ScanDirection::Ascending,
773 3,
774 SkipPolicy::DRAIN_ONLY,
775 |_| {},
776 );
777 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
778
779 let items: Vec<u64> = all
780 .iter()
781 .filter_map(|result| match result {
782 Ok(Watermarked::Item(item)) => Some(*item),
783 _ => None,
784 })
785 .collect();
786 assert_eq!(items, vec![40 * BUCKET_SIZE + 7]);
787
788 let watermarks: Vec<u64> = all
789 .iter()
790 .filter_map(|r| match r {
791 Ok(Watermarked::Watermark(position)) => Some(*position),
792 _ => None,
793 })
794 .collect();
795 assert_eq!(watermarks, vec![40 * BUCKET_SIZE, 41 * BUCKET_SIZE]);
796 assert_eq!(
797 all.len(),
798 items.len() + watermarks.len() + 1,
799 "the stopping round must add only its terminal, not another beacon"
800 );
801 let err = all
802 .last()
803 .expect("non-empty")
804 .as_ref()
805 .expect_err("scan must terminate with an error");
806 match err {
807 ScanStop::ScanLimit { scan_frontier } => assert_eq!(
808 *scan_frontier,
809 50 * BUCKET_SIZE,
810 "terminal must carry the logically advanced merged floor"
811 ),
812 other => panic!("expected ScanLimit, got {other:?}"),
813 }
814 }
815
816 #[tokio::test]
827 async fn flush_on_error_delivers_below_floor_sibling_result() {
828 let source = TestBucketSource {
829 buckets: Arc::new(BTreeMap::from([
830 (
831 test_key(b"a"),
832 vec![
833 (0, vec![1]),
834 (1, vec![1]),
835 (2, vec![1]),
836 (3, vec![1]),
837 (4, vec![1]),
838 ],
839 ),
840 (test_key(b"b"), vec![(50, vec![1])]),
841 (test_key(b"c"), vec![(0, vec![7])]),
842 ])),
843 };
844 let query = BitmapQuery::new(vec![
845 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
846 BitmapTerm::new(vec![include(b"c")]).unwrap(),
847 ])
848 .unwrap();
849
850 let stream = eval_bitmap_query_stream(
851 source,
852 query,
853 0..(60 * BUCKET_SIZE),
854 BUCKET_SIZE,
855 ScanDirection::Ascending,
856 3,
857 SkipPolicy::DRAIN_ONLY,
858 |_| {},
859 );
860 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
861
862 let items: Vec<u64> = all
864 .iter()
865 .filter_map(|r| match r {
866 Ok(Watermarked::Item(v)) => Some(*v),
867 _ => None,
868 })
869 .collect();
870 assert_eq!(items, vec![7], "c's below-floor match must be delivered");
871
872 let watermarks: Vec<u64> = all
873 .iter()
874 .filter_map(|r| match r {
875 Ok(Watermarked::Watermark(position)) => Some(*position),
876 _ => None,
877 })
878 .collect();
879 assert_eq!(watermarks, vec![BUCKET_SIZE]);
880 assert_eq!(
881 all.len(),
882 items.len() + watermarks.len() + 1,
883 "the stopping round must add only its terminal, not another beacon"
884 );
885
886 let err = all
887 .last()
888 .expect("non-empty")
889 .as_ref()
890 .expect_err("scan must terminate with an error");
891 let scan_frontier = match err {
892 ScanStop::ScanLimit { scan_frontier } => *scan_frontier,
893 other => panic!("expected ScanLimit, got {other:?}"),
894 };
895 assert_eq!(
896 scan_frontier,
897 50 * BUCKET_SIZE,
898 "terminal frontier must include the logical gap jump"
899 );
900 assert!(
901 items.iter().all(|&i| i < scan_frontier),
902 "delivered items must be below the resume frontier"
903 );
904 }
905
906 #[test]
907 fn bitmap_query_validation_rejects_empty_shapes() {
908 assert!(BitmapQuery::new(Vec::new()).is_err());
909 assert!(BitmapLiteral::include(Vec::new()).is_err());
910 assert!(
911 BitmapLiteral::include(vec![crate::dimensions::IndexDimension::Sender.tag_byte()])
912 .is_err()
913 );
914 assert!(BitmapLiteral::include(vec![0xff, 0x00]).is_err());
915 assert!(BitmapTerm::new(vec![exclude(b"neg")]).is_err());
916 }
917
918 #[tokio::test]
923 async fn shared_include_across_terms_scans_dimension_once() {
924 use crate::bitmap_query::test_utils::CountingBucketSource;
925
926 let source = CountingBucketSource::new(BTreeMap::from([
927 (test_key(b"a"), vec![(0, vec![1, 2, 3])]),
928 (test_key(b"b"), vec![(0, vec![1])]),
929 (test_key(b"c"), vec![(0, vec![2])]),
930 ]));
931 let query = BitmapQuery::new(vec![
932 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
933 BitmapTerm::new(vec![include(b"a"), include(b"c")]).unwrap(),
934 ])
935 .unwrap();
936
937 let stream = eval_bitmap_query_stream(
938 source.clone(),
939 query,
940 0..200_000,
941 BUCKET_SIZE,
942 ScanDirection::Ascending,
943 u64::MAX,
944 SkipPolicy::DRAIN_ONLY,
945 |_| {},
946 );
947 let (items, _watermarks) = drain_marked(stream).await.unwrap();
948
949 assert_eq!(items, vec![1, 2]);
952 assert_eq!(source.scan_count(&test_key(b"a")), 1);
954 assert_eq!(source.scan_count(&test_key(b"b")), 1);
955 assert_eq!(source.scan_count(&test_key(b"c")), 1);
956 }
957
958 #[tokio::test]
962 async fn shared_key_across_include_and_exclude_terms_scans_once() {
963 use crate::bitmap_query::test_utils::CountingBucketSource;
964
965 let source = CountingBucketSource::new(BTreeMap::from([
966 (test_key(b"a"), vec![(0, vec![1, 2])]),
967 (test_key(b"b"), vec![(0, vec![1, 2, 3])]),
968 ]));
969 let query = BitmapQuery::new(vec![
973 BitmapTerm::new(vec![include(b"b"), exclude(b"a")]).unwrap(),
974 BitmapTerm::new(vec![include(b"b"), include(b"a")]).unwrap(),
975 ])
976 .unwrap();
977
978 let stream = eval_bitmap_query_stream(
979 source.clone(),
980 query,
981 0..100_000,
982 BUCKET_SIZE,
983 ScanDirection::Ascending,
984 u64::MAX,
985 SkipPolicy::DRAIN_ONLY,
986 |_| {},
987 );
988 let (items, _watermarks) = drain_marked(stream).await.unwrap();
989
990 assert_eq!(items, vec![1, 2, 3]);
991 assert_eq!(source.scan_count(&test_key(b"a")), 1);
992 assert_eq!(source.scan_count(&test_key(b"b")), 1);
993 }
994
995 #[tokio::test]
996 async fn eval_bitmap_query_stream_uses_backend_source() {
997 let source = TestBucketSource {
998 buckets: Arc::new(BTreeMap::from([
999 (test_key(b"a"), vec![(0, vec![1, 2, 3]), (1, vec![5])]),
1000 (test_key(b"b"), vec![(0, vec![2, 3]), (1, vec![5])]),
1001 (test_key(b"c"), vec![(0, vec![3])]),
1002 ])),
1003 };
1004 let query = BitmapQuery::new(vec![
1005 BitmapTerm::new(vec![include(b"a"), include(b"b"), exclude(b"c")]).unwrap(),
1006 ])
1007 .unwrap();
1008
1009 let stream = eval_bitmap_query_stream(
1010 source,
1011 query,
1012 0..200_000,
1013 BUCKET_SIZE,
1014 ScanDirection::Ascending,
1015 u64::MAX,
1016 SkipPolicy::DRAIN_ONLY,
1017 |_| {},
1018 );
1019 let (items, _watermarks) = drain_marked(stream).await.unwrap();
1020
1021 assert_eq!(items, vec![2, BUCKET_SIZE + 5]);
1022 }
1023
1024 #[tokio::test]
1025 async fn eval_bitmap_query_stream_descending() {
1026 let source = TestBucketSource {
1027 buckets: Arc::new(BTreeMap::from([
1028 (test_key(b"a"), vec![(0, vec![1, 2, 3]), (1, vec![5])]),
1029 (test_key(b"b"), vec![(0, vec![2, 3]), (1, vec![5])]),
1030 (test_key(b"c"), vec![(0, vec![3])]),
1031 ])),
1032 };
1033 let query = BitmapQuery::new(vec![
1034 BitmapTerm::new(vec![include(b"a"), include(b"b"), exclude(b"c")]).unwrap(),
1035 ])
1036 .unwrap();
1037
1038 let stream = eval_bitmap_query_stream(
1039 source,
1040 query,
1041 0..200_000,
1042 BUCKET_SIZE,
1043 ScanDirection::Descending,
1044 u64::MAX,
1045 SkipPolicy::DRAIN_ONLY,
1046 |_| {},
1047 );
1048 let (items, _watermarks) = drain_marked(stream).await.unwrap();
1049
1050 assert_eq!(items, vec![BUCKET_SIZE + 5, 2]);
1051 }
1052
1053 #[tokio::test]
1056 async fn flatten_watermarked_buckets_ascending() {
1057 let range = 50u64..(2 * BUCKET_SIZE + 50_001);
1058 let marked_buckets = stream::iter(vec![
1059 Ok::<_, LeafStop>(Watermarked::Item((
1060 0,
1061 make_bitmap(&[10, 50, (BUCKET_SIZE - 1) as u32]),
1062 ))),
1063 Ok(Watermarked::Watermark(BUCKET_SIZE)),
1064 Ok(Watermarked::Item((
1065 1,
1066 make_bitmap(&[0, (BUCKET_SIZE - 1) as u32]),
1067 ))),
1068 Ok(Watermarked::Watermark(2 * BUCKET_SIZE)),
1069 Ok(Watermarked::Item((2, make_bitmap(&[0, 50_000, 50_001])))),
1070 Ok(Watermarked::Watermark(2 * BUCKET_SIZE + 50_001)),
1071 ]);
1072 let out: Vec<Watermarked<u64>> = flatten_watermarked_buckets(
1073 marked_buckets,
1074 range,
1075 BUCKET_SIZE,
1076 ScanDirection::Ascending,
1077 )
1078 .try_collect()
1079 .await
1080 .unwrap();
1081
1082 assert_eq!(
1083 out,
1084 vec![
1085 Watermarked::Item(50),
1086 Watermarked::Item(BUCKET_SIZE - 1),
1087 Watermarked::Watermark(BUCKET_SIZE),
1088 Watermarked::Item(BUCKET_SIZE),
1089 Watermarked::Item(2 * BUCKET_SIZE - 1),
1090 Watermarked::Watermark(2 * BUCKET_SIZE),
1091 Watermarked::Item(2 * BUCKET_SIZE),
1092 Watermarked::Item(2 * BUCKET_SIZE + 50_000),
1093 Watermarked::Watermark(2 * BUCKET_SIZE + 50_001),
1094 ],
1095 );
1096 }
1097
1098 #[tokio::test]
1099 async fn scan_budget_below_unique_leaf_count_yields_misconfig_error() {
1100 let source = TestBucketSource {
1107 buckets: Arc::new(BTreeMap::from([(
1108 test_key(b"a"),
1109 vec![(0, vec![1, 2]), (1, vec![3])],
1110 )])),
1111 };
1112 let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
1113
1114 let metrics = std::sync::Arc::new(std::sync::Mutex::new(None));
1115 let metrics_sink = metrics.clone();
1116 let stream = eval_bitmap_query_stream(
1117 source,
1118 query,
1119 0..200_000,
1120 BUCKET_SIZE,
1121 ScanDirection::Ascending,
1122 0,
1123 SkipPolicy::DRAIN_ONLY,
1124 move |m| *metrics_sink.lock().unwrap() = Some(m),
1125 );
1126 let err = drain_marked(stream).await.unwrap_err();
1127
1128 assert!(
1129 matches!(err, ScanStop::Fault(_)),
1130 "must be a Fault (cursorless), never a clean ScanLimit end"
1131 );
1132 assert!(
1133 err.to_string().contains("insufficient for"),
1134 "expected misconfig error, got {err:?}"
1135 );
1136 assert!(
1139 metrics.lock().unwrap().is_none(),
1140 "misconfig early-out should not emit scan metrics"
1141 );
1142 }
1143
1144 #[tokio::test]
1145 async fn scan_budget_shared_across_dimensions() {
1146 let source = TestBucketSource {
1150 buckets: Arc::new(BTreeMap::from([
1151 (
1152 test_key(b"a"),
1153 vec![(0, vec![1]), (1, vec![2]), (2, vec![3])],
1154 ),
1155 (
1156 test_key(b"b"),
1157 vec![(0, vec![1]), (1, vec![2]), (2, vec![3])],
1158 ),
1159 (
1160 test_key(b"c"),
1161 vec![(0, vec![1]), (1, vec![2]), (2, vec![3])],
1162 ),
1163 ])),
1164 };
1165 let query = BitmapQuery::new(vec![
1166 BitmapTerm::new(vec![include(b"a"), include(b"b"), include(b"c")]).unwrap(),
1167 ])
1168 .unwrap();
1169
1170 let metrics = std::sync::Arc::new(std::sync::Mutex::new(None));
1171 let metrics_sink = metrics.clone();
1172 let stream = eval_bitmap_query_stream(
1173 source,
1174 query,
1175 0..300_000,
1176 BUCKET_SIZE,
1177 ScanDirection::Ascending,
1178 4,
1179 SkipPolicy::DRAIN_ONLY,
1180 move |m| *metrics_sink.lock().unwrap() = Some(m),
1181 );
1182 let err = drain_marked(stream).await.unwrap_err();
1183
1184 assert!(
1185 matches!(err, ScanStop::ScanLimit { .. }),
1186 "expected ScanLimit, got {err:?}"
1187 );
1188 assert_eq!(metrics.lock().unwrap().unwrap().buckets_evaluated, 4);
1191 }
1192
1193 #[tokio::test]
1198 async fn scan_budget_exclude_side_exhaustion_does_not_leak_includes() {
1199 let source = TestBucketSource {
1200 buckets: Arc::new(BTreeMap::from([
1201 (
1202 test_key(b"inc"),
1203 vec![(0, vec![1]), (1, vec![2]), (2, vec![3])],
1204 ),
1205 (
1206 test_key(b"exc"),
1207 vec![(0, vec![1]), (1, vec![2]), (2, vec![3])],
1208 ),
1209 ])),
1210 };
1211 let query = BitmapQuery::new(vec![
1212 BitmapTerm::new(vec![include(b"inc"), exclude(b"exc")]).unwrap(),
1213 ])
1214 .unwrap();
1215
1216 let stream = eval_bitmap_query_stream(
1223 source,
1224 query,
1225 0..300_000,
1226 BUCKET_SIZE,
1227 ScanDirection::Ascending,
1228 2,
1229 SkipPolicy::DRAIN_ONLY,
1230 |_| {},
1231 );
1232 let result = drain_marked(stream).await;
1233
1234 let err = result.expect_err("must surface scan-limit, not silently emit includes");
1236 assert!(
1237 matches!(err, ScanStop::ScanLimit { .. }),
1238 "expected ScanLimit, got {err:?}"
1239 );
1240 }
1241
1242 #[tokio::test]
1247 async fn sparse_intersect_bundles_frontier_in_scan_limit() {
1248 let source = TestBucketSource {
1249 buckets: Arc::new(BTreeMap::from([
1250 (
1251 test_key(b"a"),
1252 vec![
1253 (0, vec![1]),
1254 (1, vec![1]),
1255 (2, vec![1]),
1256 (3, vec![1]),
1257 (4, vec![1]),
1258 ],
1259 ),
1260 (test_key(b"b"), vec![(100, vec![1])]),
1261 ])),
1262 };
1263 let query = BitmapQuery::new(vec![
1264 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
1265 ])
1266 .unwrap();
1267
1268 let stream = eval_bitmap_query_stream(
1269 source,
1270 query,
1271 0..(110 * BUCKET_SIZE),
1272 BUCKET_SIZE,
1273 ScanDirection::Ascending,
1274 4,
1275 SkipPolicy::DRAIN_ONLY,
1276 |_| {},
1277 );
1278 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
1279
1280 assert!(
1281 all.iter().all(|r| !matches!(r, Ok(Watermarked::Item(_)))),
1282 "disjoint intersect must not emit items"
1283 );
1284 let err = all
1285 .last()
1286 .expect("non-empty")
1287 .as_ref()
1288 .expect_err("scan must terminate with an error");
1289 let scan_frontier = match err {
1290 ScanStop::ScanLimit { scan_frontier } => *scan_frontier,
1291 other => panic!("expected ScanLimit, got {other:?}"),
1292 };
1293 assert_eq!(
1294 scan_frontier,
1295 100 * BUCKET_SIZE,
1296 "terminal must carry the logically advanced merged floor"
1297 );
1298 let watermarks: Vec<u64> = all
1299 .iter()
1300 .filter_map(|r| match r {
1301 Ok(Watermarked::Watermark(position)) => Some(*position),
1302 _ => None,
1303 })
1304 .collect();
1305 assert_eq!(
1306 watermarks,
1307 vec![100 * BUCKET_SIZE],
1308 "the stopping round must not append another frontier beacon"
1309 );
1310 assert_eq!(all.len(), watermarks.len() + 1);
1311 }
1312
1313 #[tokio::test]
1314 async fn eval_emits_watermarks_at_bucket_boundaries_ascending() {
1315 let source = TestBucketSource {
1316 buckets: Arc::new(BTreeMap::from([(
1317 test_key(b"a"),
1318 vec![(0, vec![1]), (3, vec![2]), (7, vec![3])],
1319 )])),
1320 };
1321 let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
1322
1323 let stream = eval_bitmap_query_stream(
1324 source,
1325 query,
1326 0..(8 * BUCKET_SIZE),
1327 BUCKET_SIZE,
1328 ScanDirection::Ascending,
1329 u64::MAX,
1330 SkipPolicy::DRAIN_ONLY,
1331 |_| {},
1332 );
1333 let (items, watermarks) = drain_marked(stream).await.unwrap();
1334
1335 assert_eq!(items, vec![1, 3 * BUCKET_SIZE + 2, 7 * BUCKET_SIZE + 3]);
1337 assert_eq!(
1341 watermarks,
1342 vec![
1343 BUCKET_SIZE,
1344 3 * BUCKET_SIZE,
1345 4 * BUCKET_SIZE,
1346 7 * BUCKET_SIZE,
1347 8 * BUCKET_SIZE,
1348 ]
1349 );
1350 }
1351
1352 #[tokio::test]
1353 async fn descending_exclusive_upper_bound_is_not_progress() {
1354 let source = TestBucketSource {
1355 buckets: Arc::new(BTreeMap::from([(
1356 test_key(b"a"),
1357 vec![(0, vec![1]), (3, vec![2]), (7, vec![3])],
1358 )])),
1359 };
1360 let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
1361
1362 let stream = eval_bitmap_query_stream(
1363 source,
1364 query,
1365 0..(8 * BUCKET_SIZE),
1366 BUCKET_SIZE,
1367 ScanDirection::Descending,
1368 u64::MAX,
1369 SkipPolicy::DRAIN_ONLY,
1370 |_| {},
1371 );
1372 let (items, watermarks) = drain_marked(stream).await.unwrap();
1373
1374 assert_eq!(items, vec![7 * BUCKET_SIZE + 3, 3 * BUCKET_SIZE + 2, 1]);
1375 assert_eq!(
1378 watermarks,
1379 vec![
1380 7 * BUCKET_SIZE,
1381 4 * BUCKET_SIZE,
1382 3 * BUCKET_SIZE,
1383 BUCKET_SIZE,
1384 0,
1385 ]
1386 );
1387 }
1388
1389 #[tokio::test]
1390 async fn natural_completion_omits_terminus_but_retains_earned_progress() {
1391 let source = TestBucketSource {
1396 buckets: Arc::new(BTreeMap::from([
1397 (
1398 test_key(b"a"),
1399 vec![(0, vec![1]), (2, vec![3]), (4, vec![5])],
1400 ),
1401 (
1402 test_key(b"b"),
1403 vec![(1, vec![1]), (3, vec![3]), (5, vec![5])],
1404 ),
1405 ])),
1406 };
1407 let query = BitmapQuery::new(vec![
1408 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
1409 ])
1410 .unwrap();
1411
1412 let stream = eval_bitmap_query_stream(
1413 source,
1414 query,
1415 0..(6 * BUCKET_SIZE),
1416 BUCKET_SIZE,
1417 ScanDirection::Ascending,
1418 u64::MAX,
1419 SkipPolicy::DRAIN_ONLY,
1420 |_| {},
1421 );
1422 let (items, watermarks) = drain_marked(stream).await.unwrap();
1423
1424 assert!(items.is_empty(), "disjoint intersect must not emit items");
1425 assert_eq!(
1426 watermarks,
1427 vec![
1428 BUCKET_SIZE,
1429 2 * BUCKET_SIZE,
1430 3 * BUCKET_SIZE,
1431 4 * BUCKET_SIZE,
1432 5 * BUCKET_SIZE,
1433 ],
1434 );
1435 }
1436
1437 #[tokio::test]
1438 async fn eval_watermark_ordering_invariant_item_then_watermark() {
1439 let source = TestBucketSource {
1444 buckets: Arc::new(BTreeMap::from([(
1445 test_key(b"a"),
1446 vec![(0, vec![10, 20, 30]), (1, vec![40, 50])],
1447 )])),
1448 };
1449 let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
1450
1451 let stream = eval_bitmap_query_stream(
1452 source,
1453 query,
1454 0..(2 * BUCKET_SIZE),
1455 BUCKET_SIZE,
1456 ScanDirection::Ascending,
1457 u64::MAX,
1458 SkipPolicy::DRAIN_ONLY,
1459 |_| {},
1460 );
1461 let all: Vec<Watermarked<u64>> = stream.try_collect().await.unwrap();
1462
1463 assert_eq!(
1467 all,
1468 vec![
1469 Watermarked::Item(10),
1470 Watermarked::Item(20),
1471 Watermarked::Item(30),
1472 Watermarked::Watermark(BUCKET_SIZE),
1473 Watermarked::Item(BUCKET_SIZE + 40),
1474 Watermarked::Item(BUCKET_SIZE + 50),
1475 Watermarked::Watermark(2 * BUCKET_SIZE),
1476 ],
1477 );
1478 }
1479
1480 #[tokio::test]
1481 async fn flatten_watermarked_buckets_descending() {
1482 let range = 50u64..(2 * BUCKET_SIZE + 50_001);
1483 let marked_buckets = stream::iter(vec![
1484 Ok::<_, LeafStop>(Watermarked::Item((2, make_bitmap(&[0, 50_000, 50_001])))),
1485 Ok(Watermarked::Watermark(2 * BUCKET_SIZE)),
1486 Ok(Watermarked::Item((
1487 1,
1488 make_bitmap(&[0, (BUCKET_SIZE - 1) as u32]),
1489 ))),
1490 Ok(Watermarked::Watermark(BUCKET_SIZE)),
1491 Ok(Watermarked::Item((
1492 0,
1493 make_bitmap(&[10, 50, (BUCKET_SIZE - 1) as u32]),
1494 ))),
1495 Ok(Watermarked::Watermark(50)),
1496 ]);
1497 let out: Vec<Watermarked<u64>> = flatten_watermarked_buckets(
1498 marked_buckets,
1499 range,
1500 BUCKET_SIZE,
1501 ScanDirection::Descending,
1502 )
1503 .try_collect()
1504 .await
1505 .unwrap();
1506
1507 assert_eq!(
1508 out,
1509 vec![
1510 Watermarked::Item(2 * BUCKET_SIZE + 50_000),
1511 Watermarked::Item(2 * BUCKET_SIZE),
1512 Watermarked::Watermark(2 * BUCKET_SIZE),
1513 Watermarked::Item(2 * BUCKET_SIZE - 1),
1514 Watermarked::Item(BUCKET_SIZE),
1515 Watermarked::Watermark(BUCKET_SIZE),
1516 Watermarked::Item(BUCKET_SIZE - 1),
1517 Watermarked::Item(50),
1518 Watermarked::Watermark(50),
1519 ],
1520 );
1521 }
1522
1523 #[test]
1526 fn collapse_single_budget_stop_binds_frontier() {
1527 assert!(matches!(
1528 collapse(vec![LeafStop::BudgetExhausted], 17 * BUCKET_SIZE),
1529 ScanStop::ScanLimit { scan_frontier } if scan_frontier == 17 * BUCKET_SIZE
1530 ));
1531 }
1532
1533 #[test]
1536 fn collapse_all_budget_stops_bind_frontier() {
1537 assert!(matches!(
1538 collapse(
1539 vec![LeafStop::BudgetExhausted, LeafStop::BudgetExhausted],
1540 23 * BUCKET_SIZE
1541 ),
1542 ScanStop::ScanLimit { scan_frontier } if scan_frontier == 23 * BUCKET_SIZE
1543 ));
1544 }
1545
1546 fn gap_probe_policy() -> SkipPolicy {
1547 SkipPolicy {
1548 drain_probe_rows: std::num::NonZeroU32::new(2),
1549 }
1550 }
1551
1552 fn skewed_and_source() -> TestBucketSource {
1553 TestBucketSource {
1554 buckets: Arc::new(BTreeMap::from([
1555 (test_key(b"a"), vec![(0, vec![1]), (50, vec![1])]),
1556 (
1557 test_key(b"b"),
1558 (0..=50).map(|bucket| (bucket, vec![1])).collect(),
1559 ),
1560 ])),
1561 }
1562 }
1563
1564 #[tokio::test]
1565 async fn skewed_and_seeks_past_dead_gap() {
1566 let query = BitmapQuery::new(vec![
1567 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
1568 ])
1569 .unwrap();
1570 let (metrics_tx, metrics_rx) = std::sync::mpsc::channel();
1571 let stream = eval_bitmap_query_stream(
1572 skewed_and_source(),
1573 query,
1574 0..(51 * BUCKET_SIZE),
1575 BUCKET_SIZE,
1576 ScanDirection::Ascending,
1577 1_000,
1578 gap_probe_policy(),
1579 move |observed| metrics_tx.send(observed).unwrap(),
1580 );
1581 let (items, watermarks) = drain_marked(stream).await.unwrap();
1582
1583 assert_eq!(items, vec![1, 50 * BUCKET_SIZE + 1]);
1584 assert_eq!(
1585 watermarks,
1586 vec![BUCKET_SIZE, 50 * BUCKET_SIZE, 51 * BUCKET_SIZE]
1587 );
1588 assert_eq!(
1589 metrics_rx.recv().unwrap(),
1590 BitmapScanMetrics {
1591 buckets_evaluated: 7,
1592 buckets_discarded: 3,
1593 leaf_seeks: 1,
1594 }
1595 );
1596 }
1597
1598 #[tokio::test]
1599 async fn skewed_and_seeks_past_dead_gap_descending() {
1600 let query = BitmapQuery::new(vec![
1601 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
1602 ])
1603 .unwrap();
1604 let (metrics_tx, metrics_rx) = std::sync::mpsc::channel();
1605 let stream = eval_bitmap_query_stream(
1606 skewed_and_source(),
1607 query,
1608 0..(51 * BUCKET_SIZE),
1609 BUCKET_SIZE,
1610 ScanDirection::Descending,
1611 1_000,
1612 gap_probe_policy(),
1613 move |observed| metrics_tx.send(observed).unwrap(),
1614 );
1615 let (items, watermarks) = drain_marked(stream).await.unwrap();
1616
1617 assert_eq!(items, vec![50 * BUCKET_SIZE + 1, 1]);
1618 assert_eq!(watermarks, vec![50 * BUCKET_SIZE, BUCKET_SIZE, 0]);
1619 assert_eq!(
1620 metrics_rx.recv().unwrap(),
1621 BitmapScanMetrics {
1622 buckets_evaluated: 7,
1623 buckets_discarded: 3,
1624 leaf_seeks: 1,
1625 }
1626 );
1627 }
1628
1629 #[tokio::test]
1630 async fn budget_death_mid_gap_carries_jumped_frontier() {
1631 let query = BitmapQuery::new(vec![
1632 BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
1633 ])
1634 .unwrap();
1635 let first: Vec<_> = eval_bitmap_query_stream(
1636 skewed_and_source(),
1637 query.clone(),
1638 0..(51 * BUCKET_SIZE),
1639 BUCKET_SIZE,
1640 ScanDirection::Ascending,
1641 5,
1642 gap_probe_policy(),
1643 |_| {},
1644 )
1645 .collect()
1646 .await;
1647 let first_items: Vec<_> = first
1648 .iter()
1649 .filter_map(|result| match result {
1650 Ok(Watermarked::Item(item)) => Some(*item),
1651 _ => None,
1652 })
1653 .collect();
1654 assert_eq!(first_items, vec![1]);
1655 assert!(matches!(
1656 first.last(),
1657 Some(Err(ScanStop::ScanLimit { scan_frontier }))
1658 if *scan_frontier == 50 * BUCKET_SIZE
1659 ));
1660
1661 let resumed = eval_bitmap_query_stream(
1662 skewed_and_source(),
1663 query,
1664 (50 * BUCKET_SIZE)..(51 * BUCKET_SIZE),
1665 BUCKET_SIZE,
1666 ScanDirection::Ascending,
1667 100,
1668 gap_probe_policy(),
1669 |_| {},
1670 );
1671 let (resumed_items, _) = drain_marked(resumed).await.unwrap();
1672 assert_eq!(resumed_items, vec![50 * BUCKET_SIZE + 1]);
1673 }
1674
1675 #[tokio::test]
1676 async fn shared_leaf_never_skips_past_its_own_terms_candidate() {
1677 let source = TestBucketSource {
1678 buckets: Arc::new(BTreeMap::from([
1679 (test_key(b"a"), vec![(0, vec![1]), (50, vec![1])]),
1680 (test_key(b"e"), vec![(10, vec![3])]),
1681 ])),
1682 };
1683 let query = BitmapQuery::new(vec![
1684 BitmapTerm::new(vec![include(b"a"), exclude(b"e")]).unwrap(),
1685 BitmapTerm::new(vec![include(b"e")]).unwrap(),
1686 ])
1687 .unwrap();
1688 let (metrics_tx, metrics_rx) = std::sync::mpsc::channel();
1689 let stream = eval_bitmap_query_stream(
1690 source,
1691 query,
1692 0..(51 * BUCKET_SIZE),
1693 BUCKET_SIZE,
1694 ScanDirection::Ascending,
1695 100,
1696 gap_probe_policy(),
1697 move |observed| metrics_tx.send(observed).unwrap(),
1698 );
1699 let (items, _) = drain_marked(stream).await.unwrap();
1700
1701 assert_eq!(items, vec![1, 10 * BUCKET_SIZE + 3, 50 * BUCKET_SIZE + 1]);
1702 let metrics = metrics_rx.recv().expect("metrics callback ran");
1703 assert_eq!(metrics.buckets_discarded, 0);
1704 assert_eq!(metrics.leaf_seeks, 0);
1705 }
1706
1707 #[tokio::test]
1708 async fn exclude_leaf_drains_dead_rows_when_dragged() {
1709 let source = TestBucketSource {
1710 buckets: Arc::new(BTreeMap::from([
1711 (test_key(b"a"), vec![(0, vec![1]), (50, vec![1])]),
1712 (test_key(b"e"), vec![(10, vec![3])]),
1713 ])),
1714 };
1715 let query = BitmapQuery::new(vec![
1716 BitmapTerm::new(vec![include(b"a"), exclude(b"e")]).unwrap(),
1717 ])
1718 .unwrap();
1719 let (metrics_tx, metrics_rx) = std::sync::mpsc::channel();
1720 let stream = eval_bitmap_query_stream(
1721 source,
1722 query,
1723 0..(51 * BUCKET_SIZE),
1724 BUCKET_SIZE,
1725 ScanDirection::Ascending,
1726 100,
1727 gap_probe_policy(),
1728 move |observed| metrics_tx.send(observed).unwrap(),
1729 );
1730 let (items, _) = drain_marked(stream).await.unwrap();
1731
1732 assert_eq!(items, vec![1, 50 * BUCKET_SIZE + 1]);
1733 let metrics = metrics_rx.recv().expect("metrics callback ran");
1734 assert_eq!(metrics.buckets_discarded, 1);
1735 assert_eq!(metrics.leaf_seeks, 0);
1736 }
1737
1738 #[test]
1741 fn collapse_fault_outranks_budget_stop() {
1742 let collapsed = collapse(
1743 vec![
1744 LeafStop::BudgetExhausted,
1745 LeafStop::Fault(anyhow::anyhow!("storage boom")),
1746 ],
1747 7,
1748 );
1749 match collapsed {
1750 ScanStop::Fault(e) => assert!(e.to_string().contains("storage boom")),
1751 other => panic!("expected Fault to win, got {other:?}"),
1752 }
1753 }
1754
1755 #[test]
1757 fn collapse_fault_outranks_budget_stop_and_cancelled() {
1758 let collapsed = collapse(
1759 vec![
1760 LeafStop::BudgetExhausted,
1761 LeafStop::Cancelled,
1762 LeafStop::Fault(anyhow::anyhow!("storage boom")),
1763 ],
1764 11,
1765 );
1766 match collapsed {
1767 ScanStop::Fault(e) => assert!(e.to_string().contains("storage boom")),
1768 other => panic!("expected Fault to win, got {other:?}"),
1769 }
1770 }
1771
1772 #[test]
1775 fn collapse_budget_stop_outranks_cancelled() {
1776 assert!(matches!(
1777 collapse(vec![LeafStop::Cancelled, LeafStop::BudgetExhausted], 29),
1778 ScanStop::ScanLimit { scan_frontier: 29 }
1779 ));
1780 }
1781
1782 #[test]
1785 fn collapse_all_cancelled() {
1786 assert!(matches!(
1787 collapse(vec![LeafStop::Cancelled, LeafStop::Cancelled], 31),
1788 ScanStop::Cancelled
1789 ));
1790 }
1791
1792 #[test]
1795 fn collapse_combines_concurrent_faults() {
1796 let collapsed = collapse(
1797 vec![
1798 LeafStop::Fault(anyhow::anyhow!("boom one")),
1799 LeafStop::Fault(anyhow::anyhow!("boom two")),
1800 ],
1801 37,
1802 );
1803 match collapsed {
1804 ScanStop::Fault(e) => {
1805 let s = e.to_string();
1806 assert!(s.contains("boom one"), "missing first fault: {s}");
1807 assert!(s.contains("boom two"), "missing second fault: {s}");
1808 }
1809 other => panic!("expected combined Fault, got {other:?}"),
1810 }
1811 }
1812
1813 #[test]
1816 fn from_anyhow_funnels_to_leaf_fault() {
1817 match LeafStop::from(anyhow::anyhow!("leaf storage boom")) {
1818 LeafStop::Fault(e) => assert!(e.to_string().contains("leaf storage boom")),
1819 other => panic!("expected leaf Fault, got {other:?}"),
1820 }
1821 }
1822
1823 #[test]
1826 fn from_anyhow_funnels_to_scan_fault() {
1827 match ScanStop::from(anyhow::anyhow!("merged storage boom")) {
1828 ScanStop::Fault(e) => assert!(e.to_string().contains("merged storage boom")),
1829 other => panic!("expected scan Fault, got {other:?}"),
1830 }
1831 }
1832
1833 fn stream_items(all: &[Result<Watermarked<u64>, ScanStop>]) -> Vec<u64> {
1834 all.iter()
1835 .filter_map(|r| match r {
1836 Ok(Watermarked::Item(v)) => Some(*v),
1837 _ => None,
1838 })
1839 .collect()
1840 }
1841
1842 #[tokio::test]
1847 async fn absent_include_annihilates_term() {
1848 let source = TestBucketSource {
1849 buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
1850 };
1851 let query =
1852 BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"ghost")]).unwrap()]).unwrap();
1853
1854 let stream = eval_bitmap_query_stream(
1855 source,
1856 query,
1857 0..(2 * BUCKET_SIZE),
1858 BUCKET_SIZE,
1859 ScanDirection::Ascending,
1860 3,
1861 SkipPolicy::DRAIN_ONLY,
1862 |_| {},
1863 );
1864 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
1865 let items = stream_items(&all);
1866
1867 assert!(
1868 items.is_empty(),
1869 "absent include must annihilate: {items:?}"
1870 );
1871 }
1872
1873 #[tokio::test]
1876 async fn absent_include_annihilates_term_despite_present_include() {
1877 let source = TestBucketSource {
1878 buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
1879 };
1880 let query = BitmapQuery::new(vec![
1881 BitmapTerm::new(vec![include(b"a"), include(b"ghost")]).unwrap(),
1882 ])
1883 .unwrap();
1884
1885 let stream = eval_bitmap_query_stream(
1886 source,
1887 query,
1888 0..(2 * BUCKET_SIZE),
1889 BUCKET_SIZE,
1890 ScanDirection::Ascending,
1891 3,
1892 SkipPolicy::DRAIN_ONLY,
1893 |_| {},
1894 );
1895 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
1896 let items = stream_items(&all);
1897
1898 assert!(
1899 items.is_empty(),
1900 "absent include must annihilate: {items:?}"
1901 );
1902 }
1903
1904 #[tokio::test]
1907 async fn absent_exclude_is_noop() {
1908 let source = TestBucketSource {
1909 buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
1910 };
1911 let query = BitmapQuery::new(vec![
1912 BitmapTerm::new(vec![include(b"a"), exclude(b"ghost")]).unwrap(),
1913 ])
1914 .unwrap();
1915
1916 let stream = eval_bitmap_query_stream(
1917 source,
1918 query,
1919 0..(2 * BUCKET_SIZE),
1920 BUCKET_SIZE,
1921 ScanDirection::Ascending,
1922 3,
1923 SkipPolicy::DRAIN_ONLY,
1924 |_| {},
1925 );
1926 let all: Vec<Result<Watermarked<u64>, ScanStop>> = stream.collect().await;
1927
1928 assert_eq!(stream_items(&all), vec![1, 2]);
1929 }
1930}