Skip to main content

sui_inverted_index/bitmap_query/
iter.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Synchronous iterator evaluator for DNF bitmap queries.
5//!
6//! A single flat driver merge-joins every leaf scan against one shared *floor*
7//! (the slowest leaf's position). At the floor bucket it evaluates the query —
8//! intersect each term's included dimensions, subtract its excluded ones, then
9//! union across terms — and emits a watermark at the floor. Because leaves only
10//! ever advance at the floor (peeked one bucket ahead), no branch can run ahead
11//! of the others: the resume cursor is always within one sparse read of every
12//! leaf, and there is no windowing/parking machinery to get wrong. Mirrors the
13//! async [`super::stream`] evaluator, which shares the per-bucket evaluation
14//! ([`eval_term_at_bucket`]) and is cross-checked against this one in tests.
15//!
16//! Budget accounting lives in the request layer for the iterator path (each
17//! backend leaf iterator charges its own per-request budget and yields an error
18//! on exhaustion), so this evaluator only propagates those errors.
19
20use std::collections::VecDeque;
21use std::ops::Range;
22
23use roaring::RoaringBitmap;
24
25use super::BitmapBucketIteratorSource;
26use super::BitmapQuery;
27use super::BucketItem;
28use super::DedupedQuery;
29use super::LeafHead;
30use super::LeafStop;
31use super::ScanDirection;
32use super::ScanStop;
33use super::SkipPolicy;
34use super::Watermarked;
35use super::WatermarkedBucket;
36use super::advance_in_direction;
37use super::bound_in_direction;
38use super::bucket_edges;
39use super::build_term_specs;
40use super::collapse;
41use super::count_on_floor_refs;
42use super::eval_term_at_bucket;
43use super::frontier_advanced;
44use super::leaf_skip_targets;
45use super::recompute_unreferenced;
46use super::strictly_before;
47use super::take_snapshot_bitmap;
48
49struct Leaf<I> {
50    iter: I,
51    peeked: Option<Option<BucketItem>>,
52    drained: u64,
53}
54
55impl<I: Iterator<Item = BucketItem>> Leaf<I> {
56    fn new(iter: I) -> Self {
57        Self {
58            iter,
59            peeked: None,
60            drained: 0,
61        }
62    }
63
64    fn peek(&mut self) -> Option<&BucketItem> {
65        if self.peeked.is_none() {
66            self.peeked = Some(self.iter.next());
67        }
68        self.peeked.as_ref().and_then(Option::as_ref)
69    }
70
71    fn next(&mut self) -> Option<BucketItem> {
72        match self.peeked.take() {
73            Some(item) => item,
74            None => self.iter.next(),
75        }
76    }
77}
78
79impl<I: super::SeekableBucketIterator> Leaf<I> {
80    fn seek_bucket(&mut self, bucket: u64) {
81        if matches!(self.peeked.as_ref(), Some(Some(Ok(_)))) {
82            self.drained += 1;
83        }
84        self.peeked = None;
85        self.iter.seek_bucket(bucket);
86    }
87}
88
89/// Evaluate a DNF `BitmapQuery` as an ordered iterator of marked bucket bitmaps.
90/// Output emits `Watermarked::Item((bucket_id, bitmap))` interleaved with
91/// `Watermarked::Watermark(p)` derived from the slowest leaf's progress.
92pub fn eval_bitmap_query_bucket_iter<'a, S>(
93    source: S,
94    query: BitmapQuery,
95    range: Range<u64>,
96    bucket_size: u64,
97    direction: ScanDirection,
98    policy: SkipPolicy,
99) -> impl Iterator<Item = WatermarkedBucket> + 'a
100where
101    S: BitmapBucketIteratorSource<'a>,
102{
103    // Build one leaf per unique dimension key. Every term addresses these
104    // deduplicated leaves by index, so a shared dimension is scanned once.
105    let DedupedQuery {
106        keys: unique_keys,
107        mut terms,
108    } = build_term_specs(query.terms);
109    let mut leaves: Vec<Leaf<S::Iter>> = Vec::with_capacity(unique_keys.len());
110    for key in unique_keys {
111        leaves.push(Leaf::new(source.scan_bucket_iter(
112            key,
113            range.clone(),
114            direction,
115        )));
116    }
117
118    let leaf_count = leaves.len();
119    let terminus = if direction.is_ascending() {
120        range.end
121    } else {
122        range.start
123    };
124    let request_floor = if direction.is_ascending() {
125        range.start
126    } else {
127        range.end
128    };
129    // `unreferenced[i]` retires a leaf once no satisfiable term references it
130    // or its scan is permanently exhausted.
131    let mut unreferenced = vec![false; leaf_count];
132    // `front[i]` is the furthest position proven safe for this leaf, either by
133    // physical scanning or by a conjunction's leapfrog bound. The slowest live
134    // front limits the resume cursor.
135    let mut front = vec![request_floor; leaf_count];
136    // The request floor is a baseline, not progress earned by evaluation.
137    let mut progress_frontier = Some(request_floor);
138    let mut done = false;
139    let mut pending: VecDeque<WatermarkedBucket> = VecDeque::new();
140
141    std::iter::from_fn(move || {
142        loop {
143            if let Some(out) = pending.pop_front() {
144                return Some(out);
145            }
146            if done {
147                return None;
148            }
149
150            // Peek every active leaf without consuming it, then classify its
151            // current bucket, EOF, or error state.
152            let mut class: Vec<Option<LeafHead>> = (0..leaf_count).map(|_| None).collect();
153            for i in 0..leaf_count {
154                if unreferenced[i] {
155                    continue;
156                }
157                match leaves[i].peek() {
158                    Some(Ok((bucket, _))) => {
159                        let (pre, _) = bucket_edges(*bucket, bucket_size, &range, direction);
160                        front[i] = advance_in_direction(front[i], pre, direction);
161                        class[i] = Some(LeafHead::Bucket(*bucket));
162                    }
163                    None => {
164                        front[i] = advance_in_direction(front[i], terminus, direction);
165                        class[i] = Some(LeafHead::Eof);
166                    }
167                    Some(Err(_)) => class[i] = Some(LeafHead::Error),
168                }
169            }
170
171            // An include at EOF makes its conjunction permanently empty. A
172            // deduplicated include may make several conjunctions unsatisfiable.
173            for term in terms.iter_mut() {
174                if !term.unsatisfiable
175                    && term
176                        .includes
177                        .iter()
178                        .any(|&i| matches!(class[i], Some(LeafHead::Eof)))
179                {
180                    term.unsatisfiable = true;
181                }
182            }
183            // A shared leaf remains live until no satisfiable conjunction
184            // references it, whether as an include or an exclude.
185            recompute_unreferenced(&terms, &class, &mut unreferenced);
186
187            // Leapfrog bounds advance logical progress before the physical
188            // iterator drains or seeks across the corresponding dead rows.
189            let targets = leaf_skip_targets(&terms, &class, &unreferenced, direction);
190            for (i, target) in targets.iter().enumerate() {
191                if let Some(target) = target {
192                    let (pre, _) = bucket_edges(*target, bucket_size, &range, direction);
193                    front[i] = advance_in_direction(front[i], pre, direction);
194                }
195            }
196
197            // Consume stop frames so they surface. `collapse` attaches this
198            // round's proven-safe floor to scan-limit errors.
199            let mut errors: Vec<LeafStop> = Vec::new();
200            for i in 0..leaf_count {
201                if !unreferenced[i] && matches!(class[i], Some(LeafHead::Error)) {
202                    match leaves[i].next() {
203                        Some(Err(error)) => errors.push(error),
204                        _ => unreachable!("peek classified Error"),
205                    }
206                }
207            }
208
209            let active: Vec<usize> = (0..leaf_count).filter(|&i| !unreferenced[i]).collect();
210            // Natural completion is represented by the caller's terminal
211            // boundary; only progress earned here is emitted.
212            if active.is_empty() {
213                done = true;
214                return None;
215            }
216
217            // The floor is the slowest active leaf's proven-safe position and
218            // therefore the furthest safe merged watermark.
219            let floor_pos = active
220                .iter()
221                .map(|&i| front[i])
222                .reduce(|a, b| bound_in_direction(a, b, direction))
223                .expect("active non-empty");
224            let collapsed = (!errors.is_empty()).then(|| collapse(errors, floor_pos));
225            let scan_limited = matches!(collapsed, Some(ScanStop::ScanLimit { .. }));
226            if !scan_limited && frontier_advanced(progress_frontier, floor_pos, direction) {
227                pending.push_back(Ok(Watermarked::Watermark(floor_pos)));
228                progress_frontier = Some(floor_pos);
229            }
230            if let Some(stop) = collapsed {
231                done = true;
232                pending.push_back(Err(stop));
233                continue;
234            }
235
236            // A target marks a leaf whose logical frontier is ahead of its
237            // physical iterator. Rows before that target are proven dead.
238            let lagging: Vec<usize> = active
239                .iter()
240                .copied()
241                .filter(|&i| targets[i].is_some())
242                .collect();
243            // Preserve the drain count while a leaf remains lagging so a moving
244            // target cannot repeatedly restart its probe allowance.
245            for i in 0..leaf_count {
246                if targets[i].is_none() {
247                    leaves[i].drained = 0;
248                }
249            }
250
251            // Lagging heads cannot participate yet. Select the next physical
252            // bucket from leaves that are ready for evaluation.
253            let eval_bucket = active
254                .iter()
255                .filter(|&&i| targets[i].is_none())
256                .filter_map(|&i| match class[i] {
257                    Some(LeafHead::Bucket(bucket)) => Some(bucket),
258                    _ => None,
259                })
260                .reduce(|a, b| bound_in_direction(a, b, direction))
261                .expect("a least term candidate leaf is not lagging");
262            // Ready leaves may be evaluated strictly before the nearest lagging
263            // target. Equality waits so the lagging leaf joins that snapshot.
264            let lagging_target = lagging
265                .iter()
266                .filter_map(|&i| targets[i])
267                .reduce(|a, b| bound_in_direction(a, b, direction));
268            let evaluate =
269                lagging_target.is_none_or(|target| strictly_before(eval_bucket, target, direction));
270
271            if evaluate {
272                let (_, post) = bucket_edges(eval_bucket, bucket_size, &range, direction);
273                // Consume each ready leaf at this bucket once. Shared terms
274                // reuse its snapshot instead of advancing its iterator twice.
275                let mut snapshot: Vec<Option<RoaringBitmap>> =
276                    (0..leaf_count).map(|_| None).collect();
277                let mut on_floor = vec![false; leaf_count];
278                for i in 0..leaf_count {
279                    if !unreferenced[i]
280                        && targets[i].is_none()
281                        && matches!(class[i], Some(LeafHead::Bucket(b)) if b == eval_bucket)
282                    {
283                        on_floor[i] = true;
284                        front[i] = advance_in_direction(front[i], post, direction);
285                        snapshot[i] = match leaves[i].next() {
286                            Some(Ok((_, bitmap))) => Some(bitmap),
287                            _ => None,
288                        };
289                    }
290                }
291                // Evaluate each conjunction from the shared snapshot, then
292                // union non-empty bitmaps to implement the top-level OR.
293                let mut remaining_refs = count_on_floor_refs(&terms, &on_floor);
294                let mut result: Option<RoaringBitmap> = None;
295                for term in &terms {
296                    if term.unsatisfiable {
297                        continue;
298                    }
299                    let includes = term
300                        .includes
301                        .iter()
302                        .map(|&i| {
303                            take_snapshot_bitmap(&mut snapshot, &mut remaining_refs, &on_floor, i)
304                        })
305                        .collect();
306                    let excludes = term
307                        .excludes
308                        .iter()
309                        .map(|&i| {
310                            take_snapshot_bitmap(&mut snapshot, &mut remaining_refs, &on_floor, i)
311                        })
312                        .collect();
313                    if let Some(bitmap) = eval_term_at_bucket(includes, excludes) {
314                        result = Some(match result {
315                            None => bitmap,
316                            Some(acc) => acc | bitmap,
317                        });
318                    }
319                }
320                if let Some(bitmap) = result {
321                    pending.push_back(Ok(Watermarked::Item((eval_bucket, bitmap))));
322                }
323                if frontier_advanced(progress_frontier, post, direction) {
324                    pending.push_back(Ok(Watermarked::Watermark(post)));
325                    progress_frontier = Some(post);
326                }
327            }
328
329            // Physically catch up lagging leaves only after emitting any safe
330            // earlier bucket, because catch-up itself can exhaust the budget.
331            for i in lagging {
332                let target = targets[i].expect("lagging leaf has target");
333                loop {
334                    let dead_row = matches!(
335                        leaves[i].peek(),
336                        Some(Ok((bucket, _))) if strictly_before(*bucket, target, direction)
337                    );
338                    if !dead_row {
339                        break;
340                    }
341                    if policy
342                        .drain_probe_rows
343                        .is_some_and(|probe| leaves[i].drained >= u64::from(probe.get()))
344                    {
345                        leaves[i].seek_bucket(target);
346                        break;
347                    }
348                    let discarded = leaves[i].next();
349                    debug_assert!(matches!(discarded, Some(Ok(_))));
350                    leaves[i].drained += 1;
351                }
352            }
353        }
354    })
355}
356
357#[cfg(test)]
358mod tests {
359    use std::collections::BTreeMap;
360    use std::sync::Arc;
361
362    use futures::StreamExt;
363
364    use super::*;
365    use crate::bitmap_query::BitmapScanBudget;
366    use crate::bitmap_query::BitmapTerm;
367    use crate::bitmap_query::eval_bitmap_query_bucket_stream;
368    use crate::bitmap_query::test_utils::*;
369
370    /// Collect a marked sequence into a comparable `(bucket_id, bits)` /
371    /// watermark form.
372    fn collect_marked(items: Vec<WatermarkedBucket>) -> Vec<Watermarked<(u64, Vec<u32>)>> {
373        items
374            .into_iter()
375            .map(|r| r.unwrap().map_item(|(b, bm)| (b, bm.iter().collect())))
376            .collect()
377    }
378
379    fn items_only(marked: &[Watermarked<(u64, Vec<u32>)>]) -> Vec<(u64, Vec<u32>)> {
380        marked
381            .iter()
382            .filter_map(|w| match w {
383                Watermarked::Item(it) => Some(it.clone()),
384                Watermarked::Watermark(_) => None,
385            })
386            .collect()
387    }
388
389    #[test]
390    fn eval_bitmap_query_bucket_iter_uses_iterator_source() {
391        let source = TestBucketSource {
392            buckets: Arc::new(BTreeMap::from([
393                (test_key(b"a"), vec![(0, vec![1, 2, 3]), (1, vec![5])]),
394                (test_key(b"b"), vec![(0, vec![2, 3]), (1, vec![5])]),
395                (test_key(b"c"), vec![(0, vec![3])]),
396            ])),
397        };
398        let query = BitmapQuery::new(vec![
399            BitmapTerm::new(vec![include(b"a"), include(b"b"), exclude(b"c")]).unwrap(),
400        ])
401        .unwrap();
402
403        let out = eval_bitmap_query_bucket_iter(
404            source,
405            query,
406            0..200_000,
407            BUCKET_SIZE,
408            ScanDirection::Ascending,
409            SkipPolicy::DRAIN_ONLY,
410        )
411        .collect::<Vec<_>>();
412        let out = items_only(&collect_marked(out));
413
414        assert_eq!(out, vec![(0, vec![2]), (1, vec![5])]);
415    }
416
417    /// Two terms share the same include `a`. The iter evaluator must collapse
418    /// them to a single backend scan and distribute its per-bucket bitmap to
419    /// both terms. Mirrors the stream-side
420    /// `shared_include_across_terms_scans_dimension_once` test so the dedup
421    /// invariant is exercised on both evaluators.
422    #[test]
423    fn shared_include_across_terms_scans_dimension_once() {
424        use crate::bitmap_query::test_utils::CountingBucketSource;
425
426        let source = CountingBucketSource::new(BTreeMap::from([
427            (test_key(b"a"), vec![(0, vec![1, 2, 3])]),
428            (test_key(b"b"), vec![(0, vec![1])]),
429            (test_key(b"c"), vec![(0, vec![2])]),
430        ]));
431        let query = BitmapQuery::new(vec![
432            BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
433            BitmapTerm::new(vec![include(b"a"), include(b"c")]).unwrap(),
434        ])
435        .unwrap();
436
437        let out = items_only(&collect_marked(
438            eval_bitmap_query_bucket_iter(
439                source.clone(),
440                query,
441                0..200_000,
442                BUCKET_SIZE,
443                ScanDirection::Ascending,
444                SkipPolicy::DRAIN_ONLY,
445            )
446            .collect(),
447        ));
448
449        // Bucket 0: term1 = a∩b = {1}; term2 = a∩c = {2}; union = {1, 2}.
450        assert_eq!(out, vec![(0, vec![1, 2])]);
451        assert_eq!(source.scan_count(&test_key(b"a")), 1);
452        assert_eq!(source.scan_count(&test_key(b"b")), 1);
453        assert_eq!(source.scan_count(&test_key(b"c")), 1);
454    }
455
456    /// The iterator evaluator must produce the exact same `Watermarked` sequence
457    /// — items AND progress watermarks — as the stream evaluator for the same
458    /// query, since both share the per-bucket DNF evaluation and floor logic.
459    #[tokio::test]
460    async fn eval_bitmap_query_bucket_iter_matches_stream_for_or_terms() {
461        let source = TestBucketSource {
462            buckets: Arc::new(BTreeMap::from([
463                (
464                    test_key(b"a"),
465                    vec![(0, vec![1, 2, 3]), (1, vec![5, 6]), (2, vec![9])],
466                ),
467                (
468                    test_key(b"b"),
469                    vec![(0, vec![2, 3]), (1, vec![6]), (2, vec![9, 10])],
470                ),
471                (test_key(b"c"), vec![(0, vec![3]), (2, vec![9])]),
472                (test_key(b"d"), vec![(1, vec![1, 8]), (2, vec![7])]),
473                (test_key(b"e"), vec![(1, vec![8])]),
474            ])),
475        };
476        let query = BitmapQuery::new(vec![
477            BitmapTerm::new(vec![include(b"a"), include(b"b"), exclude(b"c")]).unwrap(),
478            BitmapTerm::new(vec![include(b"d"), exclude(b"e")]).unwrap(),
479        ])
480        .unwrap();
481
482        for direction in [ScanDirection::Ascending, ScanDirection::Descending] {
483            for policy in [
484                SkipPolicy::DRAIN_ONLY,
485                SkipPolicy {
486                    drain_probe_rows: std::num::NonZeroU32::new(2),
487                },
488            ] {
489                let stream_out: Vec<_> = eval_bitmap_query_bucket_stream(
490                    source.clone(),
491                    query.clone(),
492                    0..300_000,
493                    BUCKET_SIZE,
494                    direction,
495                    BitmapScanBudget::new(1_000_000),
496                    policy,
497                )
498                .collect()
499                .await;
500                let iter_out: Vec<_> = eval_bitmap_query_bucket_iter(
501                    source.clone(),
502                    query.clone(),
503                    0..300_000,
504                    BUCKET_SIZE,
505                    direction,
506                    policy,
507                )
508                .collect();
509
510                assert_eq!(
511                    collect_marked(stream_out),
512                    collect_marked(iter_out),
513                    "iter and stream marked sequences diverged for {direction:?}"
514                );
515            }
516        }
517    }
518
519    /// Parity holds even when buckets are spread far apart (sparse gaps, leaves
520    /// leapfrogging) — the regime where a naive merge could drift between the two
521    /// evaluators.
522    #[tokio::test]
523    async fn eval_bitmap_query_bucket_iter_matches_stream_over_sparse_gaps() {
524        let source = TestBucketSource {
525            buckets: Arc::new(BTreeMap::from([
526                (
527                    test_key(b"a"),
528                    vec![(0, vec![1, 2, 3]), (5, vec![5, 6]), (9, vec![9])],
529                ),
530                (
531                    test_key(b"b"),
532                    vec![(0, vec![2, 3]), (5, vec![6]), (9, vec![9, 10])],
533                ),
534                (test_key(b"c"), vec![(0, vec![3]), (9, vec![9])]),
535                (test_key(b"d"), vec![(3, vec![1, 8]), (7, vec![7])]),
536                (test_key(b"e"), vec![(3, vec![8])]),
537            ])),
538        };
539        let query = BitmapQuery::new(vec![
540            BitmapTerm::new(vec![include(b"a"), include(b"b"), exclude(b"c")]).unwrap(),
541            BitmapTerm::new(vec![include(b"d"), exclude(b"e")]).unwrap(),
542        ])
543        .unwrap();
544
545        for direction in [ScanDirection::Ascending, ScanDirection::Descending] {
546            for policy in [
547                SkipPolicy::DRAIN_ONLY,
548                SkipPolicy {
549                    drain_probe_rows: std::num::NonZeroU32::new(2),
550                },
551            ] {
552                let stream_out: Vec<_> = eval_bitmap_query_bucket_stream(
553                    source.clone(),
554                    query.clone(),
555                    0..(10 * BUCKET_SIZE),
556                    BUCKET_SIZE,
557                    direction,
558                    BitmapScanBudget::new(1_000_000),
559                    policy,
560                )
561                .collect()
562                .await;
563                let iter_out: Vec<_> = eval_bitmap_query_bucket_iter(
564                    source.clone(),
565                    query.clone(),
566                    0..(10 * BUCKET_SIZE),
567                    BUCKET_SIZE,
568                    direction,
569                    policy,
570                )
571                .collect();
572
573                assert_eq!(
574                    collect_marked(stream_out),
575                    collect_marked(iter_out),
576                    "iter and stream diverged over sparse gaps for {direction:?}"
577                );
578            }
579        }
580    }
581
582    /// A sparse intersection that matches nothing in a gap still emits earned
583    /// progress, including the last bucket's post edge at the range terminus.
584    #[test]
585    fn intersect_emits_coalesced_watermarks_over_sparse_gap() {
586        let source = TestBucketSource {
587            buckets: Arc::new(BTreeMap::from([
588                (test_key(b"a"), vec![(0, vec![1]), (2, vec![5])]),
589                (test_key(b"b"), vec![(0, vec![1]), (2, vec![9])]),
590            ])),
591        };
592        let query = BitmapQuery::new(vec![
593            BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
594        ])
595        .unwrap();
596
597        let marked = collect_marked(
598            eval_bitmap_query_bucket_iter(
599                source,
600                query,
601                0..300_000,
602                BUCKET_SIZE,
603                ScanDirection::Ascending,
604                SkipPolicy::DRAIN_ONLY,
605            )
606            .collect(),
607        );
608
609        // Only bucket 0 intersects (member 1); bucket 2 disjoint -> dropped.
610        assert_eq!(items_only(&marked), vec![(0, vec![1])]);
611
612        // Watermarks must be non-decreasing; the last bucket earns the terminus.
613        let watermarks: Vec<u64> = marked
614            .iter()
615            .filter_map(|w| match w {
616                Watermarked::Watermark(p) => Some(*p),
617                Watermarked::Item(_) => None,
618            })
619            .collect();
620        assert!(
621            watermarks.windows(2).all(|w| w[0] <= w[1]),
622            "ascending watermarks must be non-decreasing: {watermarks:?}"
623        );
624        assert_eq!(
625            watermarks.last().copied(),
626            Some(300_000),
627            "last bucket's post edge must reach the range terminus"
628        );
629    }
630
631    #[test]
632    fn descending_request_floor_is_not_emitted_as_progress() {
633        let source = TestBucketSource {
634            buckets: Arc::new(BTreeMap::from([(
635                test_key(b"a"),
636                vec![(2, vec![0, 50_000])],
637            )])),
638        };
639        let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
640
641        let marked = collect_marked(
642            eval_bitmap_query_bucket_iter(
643                source,
644                query,
645                50..(2 * BUCKET_SIZE + 50_001),
646                BUCKET_SIZE,
647                ScanDirection::Descending,
648                SkipPolicy::DRAIN_ONLY,
649            )
650            .collect(),
651        );
652
653        assert_eq!(
654            marked,
655            vec![
656                Watermarked::Item((2, vec![0, 50_000])),
657                Watermarked::Watermark(2 * BUCKET_SIZE),
658            ],
659        );
660    }
661
662    #[test]
663    fn natural_completion_omits_terminus_but_retains_earned_progress() {
664        let source = TestBucketSource {
665            buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(3, vec![5])])])),
666        };
667        let query = BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"a")]).unwrap()]).unwrap();
668
669        let marked = collect_marked(
670            eval_bitmap_query_bucket_iter(
671                source,
672                query,
673                0..(5 * BUCKET_SIZE),
674                BUCKET_SIZE,
675                ScanDirection::Ascending,
676                SkipPolicy::DRAIN_ONLY,
677            )
678            .collect(),
679        );
680
681        assert_eq!(
682            marked,
683            vec![
684                Watermarked::Watermark(3 * BUCKET_SIZE),
685                Watermarked::Item((3, vec![5])),
686                Watermarked::Watermark(4 * BUCKET_SIZE),
687            ],
688        );
689    }
690
691    /// An unanchored term (`NOT x`, anchored on the synthesized universe leaf)
692    /// emits the complement at exclude-occupied buckets, full bitmaps at gap
693    /// buckets, and keeps emitting full buckets after the exclude leaf EOFs.
694    #[test]
695    fn unanchored_term_emits_complement_over_gaps_and_past_exclude_eof() {
696        let source = TestBucketSource {
697            buckets: Arc::new(BTreeMap::from([(
698                test_key(b"x"),
699                vec![(0, vec![1, 2]), (2, vec![5])],
700            )])),
701        };
702        let query = BitmapQuery::new(vec![
703            BitmapTerm::new(vec![include_universe(), exclude(b"x")]).unwrap(),
704        ])
705        .unwrap();
706
707        let items: Vec<(u64, RoaringBitmap)> = eval_bitmap_query_bucket_iter(
708            source,
709            query,
710            0..(4 * BUCKET_SIZE),
711            BUCKET_SIZE,
712            ScanDirection::Ascending,
713            SkipPolicy::DRAIN_ONLY,
714        )
715        .filter_map(|r| match r.unwrap() {
716            Watermarked::Item(it) => Some(it),
717            Watermarked::Watermark(_) => None,
718        })
719        .collect();
720
721        let complement = |bits: &[u32]| {
722            let mut bm = full_bucket();
723            for &b in bits {
724                bm.remove(b);
725            }
726            bm
727        };
728        let expected = vec![
729            (0, complement(&[1, 2])),
730            (1, full_bucket()),
731            (2, complement(&[5])),
732            (3, full_bucket()),
733        ];
734        assert_eq!(items, expected);
735    }
736
737    /// Iter/stream parity for a mixed DNF with an unanchored term, in both
738    /// directions — the universe leaf forces dense bucket coverage while the
739    /// anchored term stays sparse.
740    #[tokio::test]
741    async fn eval_bitmap_query_bucket_iter_matches_stream_for_unanchored_terms() {
742        let source = TestBucketSource {
743            buckets: Arc::new(BTreeMap::from([
744                (test_key(b"a"), vec![(1, vec![7])]),
745                (test_key(b"x"), vec![(0, vec![1, 2]), (3, vec![5])]),
746            ])),
747        };
748        let query = BitmapQuery::new(vec![
749            BitmapTerm::new(vec![include(b"a")]).unwrap(),
750            BitmapTerm::new(vec![include_universe(), exclude(b"x")]).unwrap(),
751        ])
752        .unwrap();
753
754        for direction in [ScanDirection::Ascending, ScanDirection::Descending] {
755            for policy in [
756                SkipPolicy::DRAIN_ONLY,
757                SkipPolicy {
758                    drain_probe_rows: std::num::NonZeroU32::new(2),
759                },
760            ] {
761                let stream_out: Vec<_> = eval_bitmap_query_bucket_stream(
762                    source.clone(),
763                    query.clone(),
764                    0..(5 * BUCKET_SIZE),
765                    BUCKET_SIZE,
766                    direction,
767                    BitmapScanBudget::new(1_000_000),
768                    policy,
769                )
770                .collect()
771                .await;
772                let iter_out: Vec<_> = eval_bitmap_query_bucket_iter(
773                    source.clone(),
774                    query.clone(),
775                    0..(5 * BUCKET_SIZE),
776                    BUCKET_SIZE,
777                    direction,
778                    policy,
779                )
780                .collect();
781
782                assert_eq!(
783                    collect_marked(stream_out),
784                    collect_marked(iter_out),
785                    "iter and stream diverged on unanchored terms for {direction:?}"
786                );
787            }
788        }
789    }
790
791    /// Budget exhaustion mid-dense-scan bundles the merged floor in the
792    /// terminal, and resuming from that frontier covers every remaining bucket
793    /// exactly once without relying on a terminal-round progress beacon.
794    #[tokio::test]
795    async fn unanchored_budget_exhaustion_resumes_at_terminal_frontier() {
796        let source = TestBucketSource {
797            buckets: Arc::new(BTreeMap::new()),
798        };
799        let query = BitmapQuery::new(vec![
800            BitmapTerm::new(vec![include_universe(), exclude(b"x")]).unwrap(),
801        ])
802        .unwrap();
803
804        let first: Vec<_> = eval_bitmap_query_bucket_stream(
805            source.clone(),
806            query.clone(),
807            0..(10 * BUCKET_SIZE),
808            BUCKET_SIZE,
809            ScanDirection::Ascending,
810            BitmapScanBudget::new(3),
811            SkipPolicy::DRAIN_ONLY,
812        )
813        .collect()
814        .await;
815
816        let mut covered: Vec<u64> = Vec::new();
817        let mut resume_from = None;
818        let mut limit_hit = false;
819        for item in first {
820            match item {
821                Ok(Watermarked::Item((bucket, bitmap))) => {
822                    assert_eq!(bitmap, full_bucket());
823                    covered.push(bucket);
824                }
825                Ok(Watermarked::Watermark(_)) => {}
826                Err(ScanStop::ScanLimit { scan_frontier }) => {
827                    resume_from = Some(scan_frontier);
828                    limit_hit = true;
829                }
830                Err(other) => panic!("expected ScanLimit, got {other:?}"),
831            }
832        }
833        assert!(limit_hit, "3-bucket budget cannot cover 10 dense buckets");
834        assert_eq!(covered, vec![0, 1, 2]);
835        let resume_from = resume_from.expect("ScanLimit must carry a resume frontier");
836        assert_eq!(resume_from, 3 * BUCKET_SIZE);
837
838        let resumed: Vec<_> = eval_bitmap_query_bucket_stream(
839            source,
840            query,
841            resume_from..(10 * BUCKET_SIZE),
842            BUCKET_SIZE,
843            ScanDirection::Ascending,
844            BitmapScanBudget::new(1_000_000),
845            SkipPolicy::DRAIN_ONLY,
846        )
847        .collect()
848        .await;
849        for item in resumed {
850            if let Watermarked::Item((bucket, _)) = item.unwrap() {
851                covered.push(bucket);
852            }
853        }
854        assert_eq!(covered, (0..10).collect::<Vec<_>>());
855    }
856    #[tokio::test]
857    async fn iter_seeks_lagging_leaf_natively() {
858        let source = CountingBucketSource::new(BTreeMap::from([
859            (test_key(b"a"), vec![(0, vec![1]), (50, vec![1])]),
860            (
861                test_key(b"b"),
862                (0..=50).map(|bucket| (bucket, vec![1])).collect(),
863            ),
864        ]));
865        let query = BitmapQuery::new(vec![
866            BitmapTerm::new(vec![include(b"a"), include(b"b")]).unwrap(),
867        ])
868        .unwrap();
869        let policy = SkipPolicy {
870            drain_probe_rows: std::num::NonZeroU32::new(2),
871        };
872
873        let stream_out = eval_bitmap_query_bucket_stream(
874            source.clone(),
875            query.clone(),
876            0..(51 * BUCKET_SIZE),
877            BUCKET_SIZE,
878            ScanDirection::Ascending,
879            BitmapScanBudget::new(1_000),
880            policy,
881        )
882        .collect()
883        .await;
884        let iter_out = eval_bitmap_query_bucket_iter(
885            source.clone(),
886            query,
887            0..(51 * BUCKET_SIZE),
888            BUCKET_SIZE,
889            ScanDirection::Ascending,
890            policy,
891        )
892        .collect();
893
894        assert_eq!(collect_marked(stream_out), collect_marked(iter_out));
895        assert_eq!(source.seek_count(&test_key(b"b")), 1);
896    }
897
898    /// Absent-dimension semantics: an include whose key has no rows at all annihilates its
899    /// conjunction (`∩ ∅ = ∅`). Pinned explicitly because this shape only arises when a queried key
900    /// was never written (e.g. a sender with no transactions), which live-cluster tests never
901    /// exercise.
902    #[test]
903    fn absent_include_annihilates_term() {
904        let source = TestBucketSource {
905            buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
906        };
907        let query =
908            BitmapQuery::new(vec![BitmapTerm::new(vec![include(b"ghost")]).unwrap()]).unwrap();
909
910        let out = eval_bitmap_query_bucket_iter(
911            source,
912            query,
913            0..200_000,
914            BUCKET_SIZE,
915            ScanDirection::Ascending,
916            SkipPolicy::DRAIN_ONLY,
917        )
918        .collect::<Vec<_>>();
919        let items = items_only(&collect_marked(out));
920
921        assert!(
922            items.is_empty(),
923            "absent include must annihilate: {items:?}"
924        );
925    }
926
927    /// A present include cannot rescue a conjunction whose other include is absent — the
928    /// intersection is still empty.
929    #[test]
930    fn absent_include_annihilates_term_despite_present_include() {
931        let source = TestBucketSource {
932            buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
933        };
934        let query = BitmapQuery::new(vec![
935            BitmapTerm::new(vec![include(b"a"), include(b"ghost")]).unwrap(),
936        ])
937        .unwrap();
938
939        let out = eval_bitmap_query_bucket_iter(
940            source,
941            query,
942            0..200_000,
943            BUCKET_SIZE,
944            ScanDirection::Ascending,
945            SkipPolicy::DRAIN_ONLY,
946        )
947        .collect::<Vec<_>>();
948        let items = items_only(&collect_marked(out));
949
950        assert!(
951            items.is_empty(),
952            "absent include must annihilate: {items:?}"
953        );
954    }
955
956    /// An exclude whose key has no rows subtracts nothing (`∖ ∅`): the present include's matches
957    /// pass through untouched.
958    #[test]
959    fn absent_exclude_is_noop() {
960        let source = TestBucketSource {
961            buckets: Arc::new(BTreeMap::from([(test_key(b"a"), vec![(0, vec![1, 2])])])),
962        };
963        let query = BitmapQuery::new(vec![
964            BitmapTerm::new(vec![include(b"a"), exclude(b"ghost")]).unwrap(),
965        ])
966        .unwrap();
967
968        let out = eval_bitmap_query_bucket_iter(
969            source,
970            query,
971            0..200_000,
972            BUCKET_SIZE,
973            ScanDirection::Ascending,
974            SkipPolicy::DRAIN_ONLY,
975        )
976        .collect::<Vec<_>>();
977
978        assert_eq!(items_only(&collect_marked(out)), vec![(0, vec![1, 2])]);
979    }
980}