Skip to main content

sui_inverted_index/bitmap_query/
stream.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Async stream 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. Leaves only ever
10//! advance at the floor (peeked one bucket ahead, polled concurrently), so no
11//! branch can run ahead of the others: the resume cursor stays within one sparse
12//! read of every leaf, and there is no windowing/parking to get wrong. Consumed
13//! by streaming backends such as BigTable; the synchronous
14//! [`super::eval_bitmap_query_bucket_iter`] mirrors it and shares the per-bucket
15//! evaluation ([`eval_term_at_bucket`]).
16
17use 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/// Per-request bucket-scan accounting, delivered via the `on_metrics`
56/// callback passed to `eval_bitmap_query_stream`. Fires once when the
57/// eval pipeline is dropped (natural end, error, or consumer cancel).
58/// The sole exception is the budget-misconfig early-out, which errors
59/// before any scan is set up and emits nothing.
60#[derive(Clone, Copy, Debug, Eq, PartialEq)]
61pub struct BitmapScanMetrics {
62    /// Bucket rows charged against the budget, including dead rows drained
63    /// during gap catch-up.
64    pub buckets_evaluated: u64,
65    /// Charged dead bucket rows discarded while catching up lagging leaves.
66    pub buckets_discarded: u64,
67    /// Physical leaf scans abandoned and reopened at a later bucket.
68    pub leaf_seeks: u64,
69}
70
71/// Per-request evaluated-bucket budget shared across all dimension
72/// streams of one eval. Charges are post-poll — see
73/// `budgeted_bucket_stream`.
74#[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    /// Charge one bucket. Returns false on underflow.
93    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    /// Charge a leaf's mandatory first bucket: decrements the shared pool
102    /// when it can, but ALWAYS succeeds. The runtime guards `budget >=
103    /// unique_leaf_count`, but a *shared* atomic with concurrent leaves
104    /// gives no ordering guarantee — a sparse term can drain the pool
105    /// before a slower sibling leaf charges its first bucket, leaving that
106    /// leaf unable to report its first position (a cursorless `SCAN_LIMIT`).
107    /// Reserving the first bucket per leaf makes the `unique_leaf_count`
108    /// floor's promise — "every leaf reaches its first bucket" — hold.
109    /// Charging-when-possible keeps `buckets_evaluated` accurate in the
110    /// common `budget >> leaves` case; it only undercounts a first bucket
111    /// taken after the pool was already exhausted by other leaves.
112    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
141/// RAII guard: fires `on_metrics` exactly once on drop with the final
142/// `BitmapScanMetrics`. Held inside the boxed eval stream so the callback
143/// fires on natural end, error, or consumer cancel.
144struct 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
161/// Evaluate a DNF `BitmapQuery` against a backend-provided bitmap source.
162///
163/// `budget` caps evaluated buckets across all dimension scans (see scan-limit
164/// handling and [`BitmapScanMetrics`]). `on_metrics`
165/// fires exactly once when the eval stream is dropped.
166///
167/// Output emits `Watermarked::Item(absolute_member_id)` interleaved with
168/// `Watermarked::Watermark(p)` derived from the slowest leaf — sparse scans
169/// that match nothing still report progress at the rate sources advance.
170pub 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        // Misconfig guard: short-circuit before any scan setup. No
187        // `on_metrics` here — there's no scan to account for, and the
188        // error surfaces on its own. `on_metrics` is dropped uncalled.
189        return async_stream::stream! {
190            // Misconfiguration is a genuine fault, not a `ScanLimit` stop.
191            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    // Wrapping the guard inside `async_stream::stream!` keeps it alive
210    // for the stream's full lifetime; the callback fires when the
211    // consumer drops the returned `BoxStream`.
212    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
226/// Non-consuming peek of one leaf, paired with its index and the position it has
227/// now scanned to (`None` for an error head, which leaves the prior position).
228async 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
249/// Evaluate a DNF `BitmapQuery` as an ordered `WatermarkedBucketStream`.
250///
251/// The flat driver: each round peeks every active leaf concurrently, takes the
252/// slowest leaf's position as the floor (the merged watermark), evaluates the
253/// whole DNF at the floor bucket, and advances only the leaves sitting there.
254/// No leaf runs more than one peeked bucket ahead of the floor.
255pub(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            // A target marks a leaf whose logical frontier is ahead of its
372            // physical stream. Its rows before the target are proven dead, but
373            // the stream still needs to drain or seek across them.
374            let lagging: Vec<usize> = active
375                .iter()
376                .copied()
377                .filter(|&i| targets[i].is_some())
378                .collect();
379            // A leaf gets a fresh drain probe after it catches up. Keeping the
380            // count while it remains lagging prevents a moving target from
381            // repeatedly restarting the probe.
382            for i in 0..leaf_count {
383                if targets[i].is_none() {
384                    drained[i] = 0;
385                }
386            }
387
388            // Lagging heads are intentionally excluded here. Find the next
389            // physical bucket among leaves that are ready to participate in
390            // evaluation.
391            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            // Ready leaves may be evaluated before the nearest lagging target:
401            // any term that needs a lagging leaf is known to produce nothing
402            // there. Equality must wait so a leaf needed at the target is not
403            // omitted from that bucket's snapshot.
404            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                // Consume only ready leaves positioned at this bucket. Lagging
414                // leaves stay out of the snapshot until catch-up is observed in
415                // a later outer round.
416                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                // Evaluate each conjunction from the shared leaf snapshot, then
433                // union the non-empty bitmaps to implement the top-level OR.
434                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            // Physically catch up all lagging leaves concurrently. The
481            // evaluator waits for every catch-up future before starting the
482            // next classification round.
483            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                        // Stop on the landing row, an error, or EOF. Leaving it
493                        // buffered lets the next outer round classify it using
494                        // fresh term and target state.
495                        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                            // Drain short gaps from the open scan. Once the
505                            // probe is exhausted, abandon the scan and reopen at
506                            // the target instead.
507                            if policy
508                                .drain_probe_rows
509                                .is_some_and(|probe| *drained >= u64::from(probe.get()))
510                            {
511                                // The dead row in the peek slot was already
512                                // pulled through the budgeted stream. Replacing
513                                // the stream discards it without calling next().
514                                budget.note_discarded();
515                                // Narrow the replacement range so it includes
516                                // the target bucket and excludes the dead gap.
517                                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                                // This is still the same logical leaf, so the
532                                // reopened stream gets no second mandatory
533                                // first-row reservation.
534                                *leaf = budgeted_bucket_stream(raw, budget.clone(), false)
535                                    .boxed()
536                                    .peekable();
537                                budget.note_seek();
538                                break;
539                            }
540                            // Consuming a dead row keeps using the existing scan
541                            // and counts both the budgeted read and the discard.
542                            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
555/// Wrap a raw per-dimension bucket stream with the shared scan budget: charge
556/// one bucket per pull (the first via `take_first`, the rest via `try_take`),
557/// yielding [`LeafStop::BudgetExhausted`] when the pool is empty — never a
558/// silent EOF.
559fn 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
583/// Flatten marked bucket bitmaps into absolute member ids with
584/// edge-bucket trimming against `range`. Watermarks pass through
585/// unchanged.
586pub 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    /// Drain into (items, watermarks) parallel vecs. Order between the
709    /// two is lost; for ordering checks collect `Watermarked` directly.
710    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    /// Tightest-budget starvation. `term1 = (a AND b)` cannot match before
738    /// bucket 50, while independent `term2 = c` matches at bucket 40. The
739    /// logical jump lets bucket 40 evaluate before catch-up exhausts the shared
740    /// budget; the stopping round carries bucket 50's leading edge only in the
741    /// terminal frontier.
742    #[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        // Budget == unique_leaf_count: the runtime floor, the tightest starvation.
767        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    /// Flush-on-error: a budget error truncates the scan at the floor, but
817    /// matches at or below the floor were already emitted in earlier rounds.
818    /// Here `c` matches in the FIRST bucket — below where `(a AND b)` exhausts
819    /// the budget — so its item must be DELIVERED, and the resume cursor must
820    /// advance to that death floor rather than be pinned at the request floor
821    /// (the livelock this guards against). The delivered item stays below the
822    /// final watermark, so resuming from that cursor will not re-emit it.
823    ///
824    /// The logical frontier for `a` advances to bucket 50 before catch-up. The
825    /// independent bucket-0 match remains below that valid resume edge.
826    #[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        // c's bucket-0 match (member id 7) is delivered despite term1 dying.
863        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    /// Two terms share the same include literal `a`. Dedup must collapse them
919    /// to a single backend scan of `a` and distribute its per-bucket bitmap to
920    /// both terms — otherwise term 2 would see `a` already consumed by term 1
921    /// at the floor bucket and silently drop matches.
922    #[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        // Term 1: a ∩ b = {1}; term 2: a ∩ c = {2}; OR = {1, 2}. If `a` were
950        // not distributed to term 2, term 2 would be empty and items = [1].
951        assert_eq!(items, vec![1, 2]);
952        // The dedup property: `a` was scanned exactly once.
953        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    /// Same key appearing as include in one term and exclude in another:
959    /// dedup still collapses to one leaf, and snapshot-distribute clones so
960    /// both polarities see the bitmap.
961    #[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        // term 1: b AND NOT a -> {3}
970        // term 2: b AND a     -> {1, 2}
971        // OR -> {1, 2, 3}
972        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    /// Flattening trims edge buckets while preserving already-clamped
1054    /// evaluator watermarks and their ordering among items.
1055    #[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        // Defensive runtime guard: a per-request budget smaller than the
1101        // query's leaf count would produce a cursorless SCAN_LIMIT
1102        // (merged watermarks stay None until every child reports). The
1103        // eval surfaces this as a plain anyhow error — distinct from a
1104        // scan-limit stop — so the handler propagates it as Internal rather
1105        // than SCAN_LIMIT.
1106        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        // The misconfig guard short-circuits before any scan setup, so
1137        // no scan metrics are emitted (the callback is dropped uncalled).
1138        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        // Three include dimensions with several buckets each. Budget = 4
1147        // should be consumed across all per-dimension fetches before a
1148        // scan-limit stop surfaces from the merged eval stream.
1149        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        // All four buckets were evaluated through budgeted_bucket_stream
1189        // before the fifth try_take() failed and surfaced a scan-limit stop.
1190        assert_eq!(metrics.lock().unwrap().unwrap().buckets_evaluated, 4);
1191    }
1192
1193    /// Budget exhausting on an exclude leaf must NOT be mistaken for the
1194    /// exclude reaching its range terminus. With silent EOF semantics, includes
1195    /// past the exclude cutoff would leak unfiltered. With scan-limit errors,
1196    /// the error propagates and the eval pipeline short-circuits cleanly.
1197    #[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        // Budget = leaf count gives every leaf one bucket fetch (the
1217        // minimum the runtime guard allows; see
1218        // `scan_budget_below_unique_leaf_count_yields_misconfig_error`). Once
1219        // the budget exhausts mid-scan, the scan-limit stop propagates
1220        // without the driver mistaking the exclude leaf's error for a
1221        // natural EOF.
1222        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        // Must error, not return Ok with leaked include rows.
1235        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    /// Disjoint intersect with a large gap. The logical frontier jumps to the
1243    /// leading edge of bucket 100 before physical catch-up drains the dense
1244    /// sibling and exhausts the budget. The stopping round must not duplicate
1245    /// that frontier as an in-band beacon.
1246    #[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        // Items at the bits within each of the three populated buckets.
1336        assert_eq!(items, vec![1, 3 * BUCKET_SIZE + 2, 7 * BUCKET_SIZE + 3]);
1337        // The request floor is suppressed; sparse leading edges and every
1338        // post-bucket edge are emitted eagerly. The last bucket itself earns
1339        // the range terminus.
1340        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        // The exclusive upper request position is not earned progress: the
1376        // first beacon follows the highest bucket's items at its low edge.
1377        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        // The include dimensions never align, so the evaluator emits no items.
1392        // Retiring the term is a natural terminal boundary, not additional scan
1393        // progress; the last eager beacon is the final position both live leaves
1394        // actually established before one reached EOF.
1395        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        // Critical invariant: for each bucket, all Items come BEFORE the
1440        // post-bucket watermark. This is what makes the watermark safe as
1441        // a resume cursor — its arrival downstream proves the dominated
1442        // items also arrived in the same stream order.
1443        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        // The request floor is not scan progress. Each bucket's Items precede
1464        // its post-bucket watermark, and the shared edge between adjacent
1465        // buckets is emitted only once.
1466        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    /// A lone leaf budget stop becomes a merged scan limit carrying the exact
1524    /// floor supplied by the evaluator.
1525    #[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    /// Several leaf budget stops collapse to one scan limit without losing the
1534    /// evaluator's merged floor.
1535    #[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    /// A storage fault co-occurring with budget exhaustion must win: masking
1739    /// the fault as a graceful scan limit would silently corrupt results.
1740    #[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    /// A storage fault outranks both budget exhaustion and cancellation.
1756    #[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    /// Budget exhaustion outranks cancellation because it preserves a usable
1773    /// merged resume frontier.
1774    #[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    /// Cancellation remains cancellation when no higher-precedence leaf stop
1783    /// occurred in the evaluator round.
1784    #[test]
1785    fn collapse_all_cancelled() {
1786        assert!(matches!(
1787            collapse(vec![LeafStop::Cancelled, LeafStop::Cancelled], 31),
1788            ScanStop::Cancelled
1789        ));
1790    }
1791
1792    /// Several concurrent faults combine into one terminal fault that retains
1793    /// every leaf's message rather than dropping all but one.
1794    #[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    /// `From<anyhow::Error>` on the leaf channel preserves backend failures as
1814    /// leaf faults rather than manufacturing a terminal disposition.
1815    #[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    /// `From<anyhow::Error>` on the merged channel preserves backend failures
1824    /// as terminal faults rather than manufacturing a scan limit or cancel.
1825    #[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    /// Absent-dimension semantics (mirrors the iter-side tests): an include whose key has no rows
1843    /// at all annihilates its conjunction (`∩ ∅ = ∅`). Pinned explicitly because this shape only
1844    /// arises when a queried key was never written (e.g. a sender with no transactions), which
1845    /// live-cluster tests never exercise.
1846    #[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    /// A present include cannot rescue a conjunction whose other include is absent — the
1874    /// intersection is still empty.
1875    #[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    /// An exclude whose key has no rows subtracts nothing (`∖ ∅`): the present include's matches
1905    /// pass through untouched.
1906    #[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}