jerry-024 commented on code in PR #99:
URL: 
https://github.com/apache/paimon-vector-index/pull/99#discussion_r3985408372


##########
core/src/ivfflat_io.rs:
##########
@@ -1057,17 +1346,20 @@ fn scan_flat_rows(
         }
         let vector = &vectors[local_idx * d..(local_idx + 1) * d];
         let distance = if metric == MetricType::L2 {
-            if let Some(threshold) = heap.worst_distance() {
-                if fvec_l2sqr_scaled_exceeds(query, vector, 1.0, threshold) {
-                    continue;
-                }
+            let threshold = collector.cutoff();
+            // An infinite cutoff can never abandon a row, so entering the
+            // kernel would only cost a wasted SIMD pass over it.
+            if threshold.is_finite() && fvec_l2sqr_scaled_exceeds(query, 
vector, 1.0, threshold) {
+                collector.note_abandoned();
+                continue;
             }
             fvec_l2sqr(query, vector)
         } else {

Review Comment:
   <!-- dlf-review -->
   **[MAJOR] Avoid evaluating admitted candidates twice**
   
   With a finite cutoff, every candidate that is not abandoned by 
`fvec_l2sqr_scaled_exceeds` has already traversed the full vector, then 
`fvec_l2sqr` traverses it again to obtain the distance. Broad or high-hit-rate 
range searches can therefore nearly double the distance work in this hot loop. 
Please return the completed distance from the cutoff kernel when it does not 
abandon, while preserving the existing accumulation and rounding contract.



##########
core/src/range.rs:
##########
@@ -0,0 +1,1152 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Shared types for distance range search (distance bands).
+//!
+//! The internal interval is always half-open, `[lower, upper)`, expressed in 
the
+//! index's own distance space: squared distance for L2, `1 - cos` for cosine 
and
+//! `-inner_product` for inner product.
+//!
+//! # Unboundedness is a structural state, not a sentinel value
+//!
+//! Representing "unbounded" with a finite sentinel drops rows on two of the
+//! three metrics. `fvec_cosine_distance_with_norms` does not clamp, so two
+//! identical normalized vectors can produce about `-1.19e-7`, and a
+//! `lower = 0.0` sentinel would exclude the most similar rows. An 
inner-product
+//! distance of `-ip` can be exactly `f32::MAX`, and a half-open interval with
+//! `upper = f32::MAX` excludes precisely that value. L2's `lower = 0.0` 
happens
+//! to be safe because squared distances are non-negative, but that is a
+//! coincidence of one metric and should not become the representation for 
three.
+
+use std::io;
+
+use crate::distance::MetricType;
+
+/// One side's bound on an interval.
+#[derive(Debug, Clone, Copy, PartialEq)]
+pub enum Bound {
+    /// This side is unbounded. Comparisons read it as ∓∞ by direction; no
+    /// non-finite cut value is ever constructed.
+    Unbounded,
+    Finite(f32),
+}
+
+/// A validated internal interval `[lower, upper)` together with the metric it
+/// belongs to.
+#[derive(Debug, Clone, Copy, PartialEq)]
+pub struct DistanceBand {
+    lower: Bound,
+    upper: Bound,
+    metric: MetricType,
+}
+
+fn invalid(message: impl Into<String>) -> io::Error {
+    io::Error::new(io::ErrorKind::InvalidInput, message.into())
+}
+
+impl DistanceBand {
+    /// Validates and constructs a band, rejecting non-finite cuts, negative 
cuts
+    /// under squared L2, and inverted intervals.
+    ///
+    /// This has to fail loud rather than "quietly return no rows". Any caller 
can
+    /// pass an illegal value, and an empty result would be read upstream as
+    /// "this bucket genuinely has no matches" -- indistinguishable from a
+    /// correct answer, and so undetectable.
+    pub fn new(lower: Bound, upper: Bound, metric: MetricType) -> 
io::Result<Self> {
+        for (side, bound) in [("lower", lower), ("upper", upper)] {
+            if let Bound::Finite(value) = bound {
+                if !value.is_finite() {
+                    return Err(invalid(format!("{side} cut is not finite: 
{value}")));
+                }
+                if metric == MetricType::L2 && value < 0.0 {
+                    return Err(invalid(format!(
+                        "{side} cut must be non-negative for squared-L2: 
{value}"
+                    )));
+                }
+            }
+        }
+        if let (Bound::Finite(low), Bound::Finite(high)) = (lower, upper) {
+            if low > high {
+                return Err(invalid(format!("inverted band: [{low}, {high})")));
+            }
+        }
+        Ok(Self {
+            lower,
+            upper,
+            metric,
+        })
+    }
+
+    pub fn metric(&self) -> MetricType {
+        self.metric
+    }
+
+    pub fn lower(&self) -> Bound {
+        self.lower
+    }
+
+    pub fn upper(&self) -> Bound {
+        self.upper
+    }
+
+    /// An empty interval (`lower == upper`) is legal and returns zero rows.
+    pub fn is_empty(&self) -> bool {
+        matches!((self.lower, self.upper), (Bound::Finite(low), 
Bound::Finite(high)) if low >= high)
+    }
+
+    /// Membership test. Left-closed, right-open.
+    ///
+    /// A non-finite value is never a member, whichever side is unbounded.
+    /// Without this an unbounded side would admit `NaN` and the matching
+    /// infinity, because each side's comparison short-circuits to `true` --
+    /// which would sit oddly beside the collector, where a non-finite computed
+    /// distance fails loud rather than being committed. The two are not the
+    /// same response -- this is a total predicate, so it answers "not a 
member"
+    /// rather than raising -- but neither should ever call such a value a 
match,
+    /// and this is the public membership authority.
+    #[inline]
+    pub fn admit(&self, value: f32) -> bool {
+        if !value.is_finite() {
+            return false;
+        }
+        let above_lower = match self.lower {
+            Bound::Unbounded => true,
+            Bound::Finite(low) => value >= low,
+        };
+        let below_upper = match self.upper {
+            Bound::Unbounded => true,
+            Bound::Finite(high) => value < high,
+        };
+        above_lower && below_upper
+    }
+}
+
+/// Only L2 is certified so far. Cosine and inner-product bands return
+/// `Unsupported`, which means "we cannot serve this request, please fall 
back",
+/// not "the call has a bug".
+pub(crate) fn ensure_certified_metric(metric: MetricType) -> io::Result<()> {
+    if metric == MetricType::L2 {
+        return Ok(());
+    }
+    Err(io::Error::new(
+        io::ErrorKind::Unsupported,
+        format!("range search is not certified for metric {metric:?} yet"),
+    ))
+}
+
+/// The comparison operator for one endpoint of a distance predicate.
+///
+/// This is a Rust API type only. It has no stable numeric, wire, or C ABI
+/// representation; a binding layer must define and validate its own external
+/// values separately. Nothing outside this crate currently consumes it, so
+/// fixing a numbering here would freeze an ABI before it has been designed.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum CutOperator {
+    Ge,
+    Gt,
+    Le,
+    Lt,
+}
+
+/// One endpoint of a predicate. `value` is the already-folded literal from the
+/// right-hand side, passed through as-is: no squaring, no square root, no
+/// binary search on the caller's part.
+#[derive(Debug, Clone, Copy, PartialEq)]
+pub struct DistanceEndpoint {
+    pub value: f64,
+    pub op: CutOperator,
+}
+
+/// What a stored squared-L2 f32 distance displays as in SQL: `sqrtf`, then
+/// widened to double.
+///
+/// **L2 only.** The conversions for the other two metrics are not implemented
+/// here; see [`DistanceBand::from_endpoints`].
+#[inline]
+fn l2_public_value(stored_bits: u32) -> f64 {
+    f64::from(f32::from_bits(stored_bits).sqrt())
+}
+
+/// The first bit pattern satisfying `l2_public_value(bits) >= endpoint`.
+///
+/// The search interval is `0..=f32::MAX.to_bits()`, i.e. the **non-negative**
+/// f32 values. That is complete for L2, whose squared distances are
+/// non-negative, but not for cosine, which can go slightly negative — which is
+/// why cosine needs a conversion of its own.
+fn first_ge(endpoint: f64) -> Option<u32> {
+    let mut low = 0u32;
+    let mut high = f32::MAX.to_bits();
+    if l2_public_value(high) < endpoint {
+        return None;
+    }
+    while low < high {
+        let mid = low + (high - low) / 2;
+        if l2_public_value(mid) >= endpoint {
+            high = mid;
+        } else {
+            low = mid + 1;
+        }
+    }
+    Some(low)
+}
+
+/// The first bit pattern satisfying `l2_public_value(bits) > endpoint`.
+fn first_gt(endpoint: f64) -> Option<u32> {
+    let mut low = 0u32;
+    let mut high = f32::MAX.to_bits();
+    if l2_public_value(high) <= endpoint {
+        return None;
+    }
+    while low < high {
+        let mid = low + (high - low) / 2;
+        if l2_public_value(mid) > endpoint {
+            high = mid;
+        } else {
+            low = mid + 1;
+        }
+    }
+    Some(low)
+}
+
+fn unsupported_endpoint(value: f64) -> io::Error {
+    io::Error::new(
+        io::ErrorKind::Unsupported,
+        format!("endpoint {value} has no representable cut; do not push this 
predicate down"),
+    )
+}
+
+impl DistanceBand {
+    /// Derives a band from a predicate's two endpoints. `None` on either side
+    /// means that side is unbounded.
+    ///
+    /// Side and operator must agree: `lower` accepts only `Ge`/`Gt` and 
`upper`
+    /// only `Le`/`Lt`. A mismatch is an error rather than a reinterpretation —
+    /// "the lower bound is `<`" is meaningless, and guessing the intent would
+    /// only mask a dispatch bug in the caller.
+    ///
+    /// # L2 only
+    ///
+    /// Cosine and inner product return `Unsupported`. Do **not** wave them
+    /// through the L2 path; both conversions differ:
+    ///
+    /// * **cosine**: the internal value `1 - cos` can be slightly below zero
+    ///   because `distance.rs` does not clamp it, and `first_ge`/`first_gt`
+    ///   search only the non-negative f32 bit range, where a negative value
+    ///   simply cannot be found. An ordered key spanning the **whole** finite
+    ///   f32 axis, negatives included, is required to binary search it.
+    /// * **inner product**: the internal distance is `-inner_product`, so the
+    ///   public value **decreases** as the internal one increases. The 
endpoint
+    ///   must first be negated, and **the side and its open/closed sense
+    ///   flipped together** (`>= e` becomes `<= -e` on the internal value), or
+    ///   the two bounds end up completely reversed.
+    ///
+    /// A metric-certification change must implement both conversions and add 
the
+    /// "three metrics × four operators × boundary ULP" tests, rather than 
merely
+    /// relaxing [`ensure_certified_metric`].
+    pub fn from_endpoints(
+        lower: Option<DistanceEndpoint>,
+        upper: Option<DistanceEndpoint>,
+        metric: MetricType,
+    ) -> io::Result<Self> {
+        // Order matters, and the rule is that **caller bugs outrank capability
+        // gaps**. A non-finite endpoint and an operator on the wrong side are
+        // both caller bugs (`InvalidInput`); an uncertified metric is a
+        // capability gap (`Unsupported`, so the caller should fall back).
+        // Reporting either bug as `Unsupported` would let an FFI caller
+        // silently fall back and never see it.
+        for ep in [lower, upper].into_iter().flatten() {
+            if !ep.value.is_finite() {
+                return Err(invalid(format!("endpoint is not finite: {}", 
ep.value)));
+            }
+        }
+        // Reduce each side's operator to the primitive it needs, which also
+        // rejects a mismatched side before the metric is consulted. A `lower`
+        // end is closed, so its cut is the first admitted value; an `upper` 
end
+        // is open, so its cut is the first excluded one.
+        type CutFinder = fn(f64) -> Option<u32>;
+        let lower_finder: Option<(DistanceEndpoint, CutFinder)> = match lower {
+            None => None,
+            Some(ep) => Some((
+                ep,
+                match ep.op {
+                    CutOperator::Ge => first_ge,
+                    CutOperator::Gt => first_gt,
+                    CutOperator::Le | CutOperator::Lt => {
+                        return Err(invalid(format!(
+                            "lower endpoint accepts Ge or Gt, got {:?}",
+                            ep.op
+                        )))
+                    }
+                },
+            )),
+        };
+        let upper_finder: Option<(DistanceEndpoint, CutFinder)> = match upper {
+            None => None,
+            Some(ep) => Some((
+                ep,
+                match ep.op {
+                    CutOperator::Lt => first_ge,
+                    CutOperator::Le => first_gt,
+                    CutOperator::Ge | CutOperator::Gt => {
+                        return Err(invalid(format!(
+                            "upper endpoint accepts Le or Lt, got {:?}",
+                            ep.op
+                        )))
+                    }
+                },
+            )),
+        };
+        if metric != MetricType::L2 {
+            return Err(io::Error::new(
+                io::ErrorKind::Unsupported,
+                format!("endpoint derivation is not implemented for metric 
{metric:?}"),
+            ));
+        }
+        let resolve = |side: Option<(DistanceEndpoint, CutFinder)>| -> 
io::Result<Bound> {
+            match side {
+                None => Ok(Bound::Unbounded),
+                Some((ep, find)) => Ok(Bound::Finite(f32::from_bits(
+                    find(ep.value).ok_or_else(|| 
unsupported_endpoint(ep.value))?,
+                ))),
+            }
+        };
+        let lower_bound = resolve(lower_finder)?;
+        let upper_bound = resolve(upper_finder)?;
+        Self::new(lower_bound, upper_bound, metric)
+    }
+}
+
+/// Why a scan stopped.
+///
+/// A Rust API type only, like [`CutOperator`]: no numeric or C ABI
+/// representation is fixed here, because no binding layer exists yet to
+/// consume one.
+#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
+pub enum StopReason {
+    /// Every probed list was walked to completion; an empty band lands here 
too.
+    #[default]
+    Exhausted,
+    /// Stopped early on consecutive empty buckets.
+    EmptyBuckets,
+    /// Hit the configured result cap.
+    LimitReached,
+}
+
+/// Per-query statistics. Fields stay private so they can be extended.
+#[derive(Debug, Clone, Copy, Default)]
+pub struct RangeSearchStats {
+    lists_probed: usize,
+    rows_scanned: usize,
+    rows_committed: usize,
+    early_abandoned: usize,
+    stop_reason: StopReason,
+}
+
+impl RangeSearchStats {
+    /// The number of **logical ranks** probed. A rank dropped under budget
+    /// pressure and re-read later is not counted twice.
+    pub fn lists_probed(&self) -> usize {
+        self.lists_probed
+    }
+    /// Rows read and at least partially evaluated, **including** rows 
abandoned
+    /// early.
+    pub fn rows_scanned(&self) -> usize {
+        self.rows_scanned
+    }
+    pub fn rows_committed(&self) -> usize {
+        self.rows_committed
+    }
+    pub fn early_abandoned(&self) -> usize {
+        self.early_abandoned
+    }
+    /// Always [`StopReason::Exhausted`] in this version. The other two 
variants
+    /// become reachable once early stopping and result caps land.
+    pub fn stop_reason(&self) -> StopReason {
+        self.stop_reason
+    }
+}
+
+/// Call-level statistics: the counters that are one-per-call rather than
+/// per-query. Per-query counters live in [`RangeSearchStats`].
+#[derive(Debug, Clone, Copy, Default)]
+pub struct RangeSearchCallStats {
+    list_reads: usize,
+}
+
+impl RangeSearchCallStats {
+    /// The number of **first read attempts of non-empty unique lists** (a
+    /// logical measure): empty lists do not count, because the existing reader
+    /// issues no payload I/O for them, and the several chunks of an oversized
+    /// list count once. There is no re-reading, so this is also the actual
+    /// number of list reads.
+    pub fn list_reads(&self) -> usize {
+        self.list_reads
+    }
+}
+
+/// One query's view of the result.
+#[derive(Debug)]
+pub struct QueryResult<'a> {
+    pub labels: &'a [i64],
+    pub distances: &'a [f32],
+    /// Always `false` in this version: there is no result cap yet, so nothing
+    /// can set it. Do not branch on it until caps land.
+    pub limit_reached: bool,
+    pub stats: &'a RangeSearchStats,
+}
+
+/// The CSR-shaped batch result: the repository's existing flat
+/// `(Vec<i64>, Vec<f32>)` convention plus the one thing variable-length 
results
+/// require, `lims`. The `lims` / `labels` / `distances` names are aligned with
+/// Faiss's `RangeSearchResult`; its `nq` is derivable here and so is not 
stored.
+#[derive(Debug)]
+pub struct RangeSearchResult {
+    lims: Vec<usize>,
+    labels: Vec<i64>,
+    distances: Vec<f32>,
+    limit_reached: Vec<u8>,
+    stats: Vec<RangeSearchStats>,
+    call_stats: RangeSearchCallStats,
+}
+
+impl RangeSearchResult {
+    pub fn query_count(&self) -> usize {
+        self.lims.len() - 1
+    }
+
+    /// Panics when `i >= query_count()`, rather than masking a caller's
+    /// out-of-range index with an empty result.
+    pub fn query(&self, i: usize) -> QueryResult<'_> {
+        let (start, end) = (self.lims[i], self.lims[i + 1]);
+        QueryResult {
+            labels: &self.labels[start..end],
+            distances: &self.distances[start..end],
+            limit_reached: self.limit_reached[i] != 0,
+            stats: &self.stats[i],
+        }
+    }
+
+    pub fn call_stats(&self) -> &RangeSearchCallStats {
+        &self.call_stats
+    }
+
+    pub fn lims(&self) -> &[usize] {
+        &self.lims
+    }
+
+    pub fn labels(&self) -> &[i64] {
+        &self.labels
+    }
+
+    pub fn distances(&self) -> &[f32] {
+        &self.distances
+    }
+}
+
+/// Accumulates a CSR result.
+///
+/// Rows are staged **per query** and flattened in query order only in
+/// [`build`]. A batch scan is list-major: one list fans out to several 
queries,
+/// so different queries' rows are produced interleaved. Appending into a 
single
+/// global array would let `query(0)` pick up another query's rows, so each 
query
+/// gets its own bucket.
+///
+/// [`build`]: RangeResultBuilder::build
+pub(crate) struct RangeResultBuilder {
+    rows: Vec<Vec<(i64, f32)>>,
+    limit_reached: Vec<u8>,
+    stats: Vec<RangeSearchStats>,
+    call_stats: RangeSearchCallStats,
+}
+
+impl RangeResultBuilder {
+    pub(crate) fn new(nq: usize) -> Self {
+        Self {
+            rows: vec![Vec::new(); nq],
+            limit_reached: vec![0; nq],
+            stats: vec![RangeSearchStats::default(); nq],
+            call_stats: RangeSearchCallStats::default(),
+        }
+    }
+
+    /// Hands over one query's whole batch of rows, avoiding a row-by-row copy
+    /// and a second resident allocation.
+    pub(crate) fn take_rows(&mut self, query: usize, rows: Vec<(i64, f32)>) {
+        self.stats[query].rows_committed += rows.len();
+        if self.rows[query].is_empty() {
+            self.rows[query] = rows; // the usual path: take ownership outright
+        } else {
+            self.rows[query].extend(rows);
+        }
+    }
+
+    /// Records one **list read attempt** (a logical measure). The call-level
+    /// fields are module-private, so other modules can only write them here.
+    ///
+    /// Accounting rules, each asserted by a test:
+    /// * a unique list counts once even when it fans out to several queries
+    /// * an oversized list read as several streamed chunks still counts once
+    /// * empty lists do not count, since the existing reader issues no payload
+    ///   I/O for them
+    pub(crate) fn record_list_read(&mut self) {
+        self.call_stats.list_reads += 1;
+    }
+
+    pub(crate) fn record_scanned(&mut self, query: usize, rows: usize) {
+        self.stats[query].rows_scanned += rows;
+    }
+
+    pub(crate) fn record_early_abandoned(&mut self, query: usize, rows: usize) 
{
+        self.stats[query].early_abandoned += rows;
+    }
+
+    pub(crate) fn record_lists_probed(&mut self, query: usize, lists: usize) {
+        self.stats[query].lists_probed += lists;
+    }
+
+    /// Flattens into CSR in query order. No `finish_query` is needed: the 
order
+    /// comes from the index into `rows` rather than from the order of calls, 
so
+    /// a list-major parallel merge cannot cross queries.
+    pub(crate) fn build(self) -> RangeSearchResult {
+        let total: usize = self.rows.iter().map(Vec::len).sum();
+        let mut lims = Vec::with_capacity(self.rows.len() + 1);
+        let mut labels = Vec::with_capacity(total);
+        let mut distances = Vec::with_capacity(total);
+        lims.push(0);
+        for query_rows in &self.rows {
+            for (id, distance) in query_rows {
+                labels.push(*id);
+                distances.push(*distance);
+            }
+            lims.push(labels.len());
+        }
+        RangeSearchResult {
+            lims,
+            labels,
+            distances,
+            limit_reached: self.limit_reached,
+            stats: self.stats,
+            call_stats: self.call_stats,
+        }
+    }
+}
+
+// Only `push_row` is gated behind cfg(test). The production path hands over a
+// whole batch with `take_rows`, while `push_row` exists purely so the 
builder's
+// unit tests can construct rows one at a time; a pub(crate) method referenced
+// only from tests is dead code in the lib target and would fail
+// `clippy --all-targets -- -D warnings`. The remaining methods (record_*,
+// build) are all called from the production path and must **not** be gated, or
+// the non-test lib target loses them.
+#[cfg(test)]
+impl RangeResultBuilder {
+    fn push_row(&mut self, query: usize, id: i64, distance: f32) {
+        self.rows[query].push((id, distance));
+        self.stats[query].rows_committed += 1;
+    }
+}
+
+/// The probe width for a range search. Both variants ship from the first
+/// version and the enum is `#[non_exhaustive]`: `Auto` fails loud with
+/// `Unsupported` until it is implemented, rather than being silently ignored, 
so
+/// the public enum is stable from day one and does not have to grow a variant
+/// later.
+///
+/// The existing `SearchWidth` is deliberately not reused: it carries a
+/// `DiskAnnLSearch` variant, and range search does not support DiskANN, so
+/// reusing it would drag a permanently illegal variant into the new API.
+#[non_exhaustive]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum RangeSearchWidth {
+    Fixed {
+        nprobe: usize,
+    },
+    Auto {
+        initial: usize,
+        growth_factor: usize,
+        max_width: usize,
+    },
+}
+
+/// Range search parameters. **Fields are private**: knobs introduced by later
+/// work would break struct-literal construction if they were `pub` fields, 
while
+/// declaring them now without implementing them would mean accepting a 
silently
+/// ineffective parameter. Private fields plus `new()` and `with_*` remove that
+/// dilemma.
+#[derive(Debug, Clone, Copy)]

Review Comment:
   <!-- dlf-review -->
   **[MINOR] Do not expose an `Auto` mode that cannot succeed**
   
   Every well-formed `RangeSearchWidth::Auto` value passes field validation 
only to return `Unsupported`, so the public variant and `with_width` path 
currently have no successful use. Since the enum is already `non_exhaustive`, 
please expose only the working fixed-`nprobe` constructor now and add `Auto` 
when automatic widening is implemented.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to