Skip to main content

mito2/sst/parquet/
prefilter.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Prefilter framework for parquet reader.
16//!
17//! Prefilter optimization reduces I/O by reading only a subset of columns first
18//! (the prefilter phase), applying filters to compute a refined row selection,
19//! then reading the remaining columns with the refined selection.
20
21use std::collections::HashSet;
22use std::ops::{BitAnd, Range};
23use std::sync::Arc;
24
25use api::v1::SemanticType;
26use common_recordbatch::filter::SimpleFilterEvaluator;
27use datafusion_common::ScalarValue;
28use datafusion_expr::Expr;
29use datatypes::arrow::array::{Array, BinaryArray, BooleanArray, BooleanBufferBuilder};
30use datatypes::arrow::buffer::BooleanBuffer;
31use datatypes::arrow::datatypes::SchemaRef;
32use datatypes::arrow::record_batch::RecordBatch;
33use datatypes::prelude::ConcreteDataType;
34use datatypes::value::Value;
35use futures::StreamExt;
36use mito_codec::row_converter::{PrimaryKeyCodec, PrimaryKeyFilter, build_primary_key_codec};
37use parquet::arrow::ProjectionMask;
38use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
39use parquet::file::metadata::ParquetMetaData;
40use parquet::schema::types::SchemaDescriptor;
41use smallvec::{SmallVec, smallvec};
42use snafu::{OptionExt, ResultExt};
43use store_api::metadata::{RegionMetadata, RegionMetadataRef};
44use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
45use table::predicate::Predicate;
46
47use crate::cache::PrefilterKey;
48use crate::error::{
49    ComputeArrowSnafu, DecodeSnafu, EvalPartitionFilterSnafu, NewRecordBatchSnafu,
50    RecordBatchSnafu, Result, UnexpectedSnafu,
51};
52use crate::sst::parquet::file_range::PreFilterMode;
53use crate::sst::parquet::flat_format::FlatReadFormat;
54use crate::sst::parquet::format::{PrimaryKeyArray, StatValues};
55use crate::sst::parquet::reader::{
56    MaybeFilter, PhysicalFilterContext, RowGroupBuildContext, RowGroupReaderBuilder,
57    SimpleFilterContext,
58};
59
60pub(crate) fn matching_row_ranges_by_primary_key(
61    input: &RecordBatch,
62    pk_column_index: usize,
63    pk_filter: &mut dyn PrimaryKeyFilter,
64) -> Result<Vec<Range<usize>>> {
65    let pk_column = input.column(pk_column_index);
66    if let Some(pk_dict_array) = pk_column.as_any().downcast_ref::<PrimaryKeyArray>() {
67        matching_row_ranges_from_dict(pk_dict_array, input.num_rows(), pk_filter)
68    } else if let Some(pk_binary_array) = pk_column.as_any().downcast_ref::<BinaryArray>() {
69        matching_row_ranges_from_binary(pk_binary_array, input.num_rows(), pk_filter)
70    } else {
71        UnexpectedSnafu {
72            reason: format!(
73                "Primary key column is neither a dictionary nor a binary array, got {:?}",
74                pk_column.data_type()
75            ),
76        }
77        .fail()
78    }
79}
80
81/// If `pk_filter` matches `pk`, records the row range `start..end`, coalescing
82/// it with the previous range when adjacent.
83fn push_matched_range(
84    matched_row_ranges: &mut Vec<Range<usize>>,
85    pk_filter: &mut dyn PrimaryKeyFilter,
86    pk: &[u8],
87    start: usize,
88    end: usize,
89) -> Result<()> {
90    if pk_filter.matches(pk).context(DecodeSnafu)? {
91        if let Some(last) = matched_row_ranges.last_mut()
92            && last.end == start
93        {
94            last.end = end;
95        } else {
96            matched_row_ranges.push(start..end);
97        }
98    }
99    Ok(())
100}
101
102/// Computes matched row ranges from a dictionary-encoded `__primary_key` column.
103fn matching_row_ranges_from_dict(
104    pk_dict_array: &PrimaryKeyArray,
105    num_rows: usize,
106    pk_filter: &mut dyn PrimaryKeyFilter,
107) -> Result<Vec<Range<usize>>> {
108    let pk_values = pk_dict_array
109        .values()
110        .as_any()
111        .downcast_ref::<BinaryArray>()
112        .context(UnexpectedSnafu {
113            reason: "Primary key values are not binary array",
114        })?;
115    let keys = pk_dict_array.keys();
116    let key_values = keys.values();
117
118    if key_values.is_empty() {
119        return Ok(std::iter::once(0..num_rows).collect());
120    }
121
122    let mut matched_row_ranges: Vec<Range<usize>> = Vec::new();
123    let mut start = 0;
124    while start < key_values.len() {
125        let key = key_values[start];
126        let mut end = start + 1;
127        while end < key_values.len() && key_values[end] == key {
128            end += 1;
129        }
130
131        push_matched_range(
132            &mut matched_row_ranges,
133            pk_filter,
134            pk_values.value(key as usize),
135            start,
136            end,
137        )?;
138
139        start = end;
140    }
141
142    Ok(matched_row_ranges)
143}
144
145/// Computes matched row ranges from a plain binary `__primary_key` column.
146///
147/// The writer falls back to plain binary encoding when the `__primary_key`
148/// chunk exceeds the dictionary page size limit (see `should_read_pk_as_binary`
149/// in the parquet reader), so the prefilter pass must handle this form too.
150fn matching_row_ranges_from_binary(
151    pk_array: &BinaryArray,
152    num_rows: usize,
153    pk_filter: &mut dyn PrimaryKeyFilter,
154) -> Result<Vec<Range<usize>>> {
155    if pk_array.is_empty() {
156        return Ok(std::iter::once(0..num_rows).collect());
157    }
158
159    let mut matched_row_ranges: Vec<Range<usize>> = Vec::new();
160    let mut start = 0;
161    while start < pk_array.len() {
162        let value = pk_array.value(start);
163        let mut end = start + 1;
164        while end < pk_array.len() && pk_array.value(end) == value {
165            end += 1;
166        }
167
168        push_matched_range(&mut matched_row_ranges, pk_filter, value, start, end)?;
169
170        start = end;
171    }
172
173    Ok(matched_row_ranges)
174}
175
176/// Filters a flat-format record batch by primary key, returning only rows whose
177/// primary key matches the filter. Returns `None` if all rows are filtered out.
178pub(crate) fn prefilter_flat_batch_by_primary_key(
179    input: RecordBatch,
180    pk_column_index: usize,
181    pk_filter: &mut dyn PrimaryKeyFilter,
182) -> Result<Option<RecordBatch>> {
183    if input.num_rows() == 0 {
184        return Ok(Some(input));
185    }
186
187    let matched_row_ranges =
188        matching_row_ranges_by_primary_key(&input, pk_column_index, pk_filter)?;
189    if matched_row_ranges.is_empty() {
190        return Ok(None);
191    }
192
193    if matched_row_ranges.len() == 1
194        && matched_row_ranges[0].start == 0
195        && matched_row_ranges[0].end == input.num_rows()
196    {
197        return Ok(Some(input));
198    }
199
200    if matched_row_ranges.len() == 1 {
201        let span = &matched_row_ranges[0];
202        return Ok(Some(input.slice(span.start, span.end - span.start)));
203    }
204
205    let mut builder = BooleanBufferBuilder::new(input.num_rows());
206    builder.append_n(input.num_rows(), false);
207    for span in matched_row_ranges {
208        for i in span {
209            builder.set_bit(i, true);
210        }
211    }
212
213    let filtered = datatypes::arrow::compute::filter_record_batch(
214        &input,
215        &BooleanArray::new(builder.finish(), None),
216    )
217    .context(ComputeArrowSnafu)?;
218    if filtered.num_rows() == 0 {
219        Ok(None)
220    } else {
221        Ok(Some(filtered))
222    }
223}
224
225pub(crate) struct CachedPrimaryKeyFilter {
226    inner: Box<dyn PrimaryKeyFilter>,
227    last_primary_key: Vec<u8>,
228    last_match: Option<bool>,
229}
230
231impl CachedPrimaryKeyFilter {
232    pub(crate) fn new(inner: Box<dyn PrimaryKeyFilter>) -> Self {
233        Self {
234            inner,
235            last_primary_key: Vec::new(),
236            last_match: None,
237        }
238    }
239}
240
241impl PrimaryKeyFilter for CachedPrimaryKeyFilter {
242    fn matches(&mut self, pk: &[u8]) -> mito_codec::error::Result<bool> {
243        if let Some(last_match) = self.last_match
244            && self.last_primary_key == pk
245        {
246            return Ok(last_match);
247        }
248
249        let matched = self.inner.matches(pk)?;
250        self.last_primary_key.clear();
251        self.last_primary_key.extend_from_slice(pk);
252        self.last_match = Some(matched);
253        Ok(matched)
254    }
255}
256
257/// How the bulk-memtable read should apply each predicate.
258///
259/// Unlike the parquet reader, the bulk path has no prefilter pass; predicates
260/// either run row-wise inside the iterator or are pushed down to encoded-PK
261/// matching when the batch still carries the primary-key column.
262pub(crate) struct BulkFilterPlan {
263    /// Simple filters the iterator still has to evaluate row-wise on each batch.
264    pub(crate) remaining_simple_filters: Vec<SimpleFilterContext>,
265    /// Tag predicates lowered to encoded-PK filters. `None` when the batch
266    /// already exposes raw tag columns or there are no tag predicates.
267    pub(crate) pk_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
268}
269
270/// Builds an encoded-primary-key filter from the supported tag predicates.
271///
272/// Predicates on fields, timestamps, or unsupported expression shapes are intentionally
273/// omitted. Callers use this as a pruning filter and must preserve the full predicate for
274/// authoritative filtering later in the scan.
275pub(crate) fn build_primary_key_filter(
276    sst_metadata: &RegionMetadataRef,
277    expected_metadata: Option<&RegionMetadata>,
278    predicate: Option<&Predicate>,
279) -> Option<CachedPrimaryKeyFilter> {
280    let filters = simple_tag_filters(sst_metadata, expected_metadata, predicate)
281        .into_iter()
282        .map(|(_, filter)| filter)
283        .collect::<Vec<_>>();
284    if filters.is_empty() {
285        return None;
286    }
287
288    let codec = build_primary_key_codec(sst_metadata.as_ref());
289    let filter = codec.primary_key_filter(sst_metadata, Arc::new(filters));
290    Some(CachedPrimaryKeyFilter::new(filter))
291}
292
293/// Extracts simple tag filters that can be applied to encoded primary keys or series indexes.
294pub(crate) fn simple_tag_filters(
295    sst_metadata: &RegionMetadataRef,
296    expected_metadata: Option<&RegionMetadata>,
297    predicate: Option<&Predicate>,
298) -> Vec<(Expr, SimpleFilterEvaluator)> {
299    predicate
300        .into_iter()
301        .flat_map(|predicate| predicate.exprs())
302        .filter_map(|expr| {
303            SimpleFilterContext::new_opt(sst_metadata, expected_metadata, expr)
304                .map(|filter_ctx| (expr, filter_ctx))
305        })
306        .filter_map(|(expr, filter_ctx)| {
307            (filter_ctx.semantic_type() == SemanticType::Tag)
308                .then(|| {
309                    filter_ctx
310                        .filter()
311                        .as_filter()
312                        .cloned()
313                        .map(|filter| (expr.clone(), filter))
314                })
315                .flatten()
316        })
317        .collect()
318}
319
320/// How the parquet reader should apply each predicate.
321///
322/// The reader runs in two phases. Predicates routed into `prefilter_builder`
323/// execute on a reduced column set first to compute a refined row selection;
324/// `remaining_simple_filters` execute alongside the full projection on the
325/// normal read path. The contract for what is precise vs best-effort is
326/// documented on [`build_reader_filter_plan`].
327pub(crate) struct ReaderFilterPlan {
328    /// Simple filters that must run on the normal read path: predicates with
329    /// `Matched` / `Pruned` outcomes (which carry expected-metadata
330    /// compatibility decisions later phases rely on), and predicates whose
331    /// column cannot be read directly during the prefilter pass.
332    pub(crate) remaining_simple_filters: Vec<SimpleFilterContext>,
333    /// Pre-built state for the prefilter pass, or `None` when prefiltering is
334    /// not worthwhile (no prefilter columns selected, or the prefilter
335    /// projection would cover nearly the full read).
336    pub(crate) prefilter_builder: Option<PrefilterContextBuilder>,
337}
338
339pub(crate) fn build_bulk_filter_plan(
340    read_format: &FlatReadFormat,
341    predicate: Option<&Predicate>,
342) -> BulkFilterPlan {
343    let metadata = read_format.metadata();
344    // Bulk memtable only needs simple binary filters here. Any filter that
345    // cannot be reduced to a SimpleFilterContext stays out of this fast path.
346    let simple_filters: Vec<SimpleFilterContext> = predicate
347        .into_iter()
348        .flat_map(|predicate| {
349            predicate
350                .exprs()
351                .iter()
352                .filter_map(|expr| SimpleFilterContext::new_opt(metadata, None, expr))
353        })
354        .collect();
355
356    // PK prefilter only works when flat batches still carry the encoded PK
357    // column. If tags have already been expanded to raw columns, the iterator
358    // can apply those filters directly and there is nothing to extract here.
359    if read_format.batch_has_raw_pk_columns() || metadata.primary_key.is_empty() {
360        return BulkFilterPlan {
361            remaining_simple_filters: simple_filters,
362            pk_filters: None,
363        };
364    }
365
366    let mut remaining_simple_filters = Vec::new();
367    let mut pk_filters = Vec::new();
368
369    for filter_ctx in simple_filters {
370        // Split tag predicates that can be evaluated against the encoded PK
371        // from filters that still need normal row-wise evaluation later.
372        let pk_filter = filter_ctx.filter().as_filter().and_then(|filter| {
373            (filter_ctx.semantic_type() == SemanticType::Tag).then(|| filter.clone())
374        });
375
376        if let Some(pk_filter) = pk_filter {
377            pk_filters.push(pk_filter);
378        } else {
379            remaining_simple_filters.push(filter_ctx);
380        }
381    }
382
383    BulkFilterPlan {
384        remaining_simple_filters,
385        pk_filters: (!pk_filters.is_empty()).then_some(Arc::new(pk_filters)),
386    }
387}
388
389/// Splits a query [`Predicate`] into a [`ReaderFilterPlan`]: predicates that can run
390/// during the prefilter pass (on a reduced projection, to compute a refined row
391/// selection) versus predicates that must run on the normal read path (alongside the
392/// full projection).
393///
394/// The prefilter pass is *best-effort pruning*: a physical-filter predicate is silently
395/// dropped when [`PhysicalFilterContext::new_opt`] returns `None` (column not in the
396/// projected arrow schema). This is safe because the DataFusion `FilterExec` above the
397/// reader always re-applies the original predicate, so the prefilter pass is purely a
398/// pruning hint.
399///
400/// With predicate prefiltering enabled, tag and timestamp predicates that lower to
401/// [`SimpleFilterEvaluator`] are an exception — the engine enforces them precisely in
402/// the prefilter pass. A caller can postpone simple timestamp filters to the normal
403/// precise-filter path when the scan time range covers the SST. When predicate
404/// prefiltering is disabled, all simple filters remain on the normal path instead.
405#[allow(clippy::too_many_arguments)]
406pub(crate) fn build_reader_filter_plan(
407    predicate: Option<&Predicate>,
408    expected_metadata: Option<&RegionMetadata>,
409    pre_filter_mode: PreFilterMode,
410    enable_predicate_prefilter: bool,
411    postpone_time_index_filter: bool,
412    read_format: &FlatReadFormat,
413    codec: &Arc<dyn PrimaryKeyCodec>,
414    parquet_metadata: &ParquetMetaData,
415) -> ReaderFilterPlan {
416    let Some(predicate) = predicate else {
417        return ReaderFilterPlan {
418            remaining_simple_filters: Vec::new(),
419            prefilter_builder: None,
420        };
421    };
422
423    let metadata = read_format.metadata();
424    let mut prefilter_simple_filters = Vec::new();
425    let mut remaining_simple_filters = Vec::new();
426    let mut prefilter_physical_filters = Vec::new();
427    let mut primary_key_filters = Vec::new();
428    let mut pk_filter_contexts = Vec::new();
429
430    // `SkipFields` keeps field predicates in the normal read path to avoid a
431    // second read of projected field columns, while tags/timestamp can still
432    // participate in prefiltering.
433    let field_prefilter_enabled = pre_filter_mode == PreFilterMode::All;
434    // When true, tag columns are encoded in the primary key column and are NOT
435    // stored as separate parquet columns. Tag predicates must go through PK
436    // decoding rather than direct column reads.
437    let need_pk_prefilter = !read_format.batch_has_raw_pk_columns();
438
439    // Whether a column can be read directly from parquet for prefiltering,
440    // based on its semantic type and the current mode/format.
441    let can_direct_prefilter = |semantic_type: SemanticType| -> bool {
442        match semantic_type {
443            SemanticType::Tag => !need_pk_prefilter,
444            SemanticType::Field => field_prefilter_enabled,
445            SemanticType::Timestamp => true,
446        }
447    };
448
449    for expr in predicate.exprs() {
450        // Prefer cheap simple filters first. They also preserve `Matched` /
451        // `Pruned` states for columns that only exist in expected metadata.
452        if let Some(filter_ctx) = SimpleFilterContext::new_opt(metadata, expected_metadata, expr) {
453            if !enable_predicate_prefilter {
454                remaining_simple_filters.push(filter_ctx);
455                continue;
456            }
457
458            // `Matched` and `Pruned` come from expected-metadata compatibility
459            // and must stay in the main filter list so later phases keep that
460            // outcome.
461            let Some(filter) = filter_ctx.filter().as_filter() else {
462                remaining_simple_filters.push(filter_ctx);
463                continue;
464            };
465
466            if postpone_time_index_filter && filter_ctx.semantic_type() == SemanticType::Timestamp {
467                remaining_simple_filters.push(filter_ctx);
468                continue;
469            }
470
471            // If the column is stored as a separate parquet column and is already projected in the main read,
472            // we can evaluate the simple filter directly during prefilter.
473            let direct_prefilter = can_direct_prefilter(filter_ctx.semantic_type());
474            if direct_prefilter {
475                assert!(
476                    read_format
477                        .arrow_schema()
478                        .column_with_name(filter.column_name())
479                        .is_some(),
480                    "Column '{}' is not present in the arrow schema {:?}",
481                    filter.column_name(),
482                    read_format.arrow_schema(),
483                );
484                prefilter_simple_filters.push(filter_ctx);
485                continue;
486            }
487
488            // Otherwise try to filter through encoded-PK matching.
489            if need_pk_prefilter && filter_ctx.semantic_type() == SemanticType::Tag {
490                primary_key_filters.push(filter.clone());
491                pk_filter_contexts.push(filter_ctx);
492            } else {
493                remaining_simple_filters.push(filter_ctx);
494            }
495            continue;
496        }
497
498        if !enable_predicate_prefilter {
499            continue;
500        }
501
502        // Best-effort physical-filter prefilter (see fn-level doc): `new_opt`
503        // returning `None` means the column is not in the projected arrow
504        // schema, and dropping the predicate is safe because the upper
505        // `FilterExec` re-applies it.
506        if let Some(filter) =
507            PhysicalFilterContext::new_opt(metadata, expected_metadata, read_format, expr)
508            && can_direct_prefilter(filter.semantic_type())
509        {
510            prefilter_physical_filters.push(filter);
511        }
512    }
513
514    if !enable_predicate_prefilter {
515        return ReaderFilterPlan {
516            remaining_simple_filters,
517            prefilter_builder: None,
518        };
519    }
520
521    let pk_filter_expr_strs = (!pk_filter_contexts.is_empty()).then(|| {
522        let mut expr_strs = pk_filter_contexts
523            .iter()
524            .map(|filter_ctx| filter_ctx.expr_str().to_string())
525            .collect::<Vec<_>>();
526        expr_strs.sort();
527        SmallVec::from_vec(expr_strs)
528    });
529    let pk_filter_exprs =
530        (!primary_key_filters.is_empty()).then_some(Arc::new(primary_key_filters));
531    let schema_version = expected_metadata
532        .map(|metadata| metadata.schema_version)
533        .unwrap_or_else(|| read_format.metadata().schema_version);
534    let prefilter_builder = PrefilterContextBuilder::new(
535        read_format,
536        codec,
537        pk_filter_exprs,
538        pk_filter_expr_strs,
539        prefilter_simple_filters.clone(),
540        prefilter_physical_filters,
541        schema_version,
542        parquet_metadata,
543    );
544
545    if prefilter_builder.is_some() {
546        ReaderFilterPlan {
547            remaining_simple_filters,
548            prefilter_builder,
549        }
550    } else {
551        // If prefilter setup is not worthwhile, keep the original simple
552        // filters on the normal path so behavior is unchanged.
553        remaining_simple_filters.extend(prefilter_simple_filters);
554        remaining_simple_filters.extend(pk_filter_contexts);
555        ReaderFilterPlan {
556            remaining_simple_filters,
557            prefilter_builder: None,
558        }
559    }
560}
561
562/// Context for prefiltering a row group.
563pub(crate) struct PrefilterContext {
564    /// Optional PK filter for legacy primary-key-format parquet.
565    pk_filter: Option<Box<dyn PrimaryKeyFilter>>,
566    /// Simple filters that can be evaluated directly from the prefilter batch.
567    filters: Vec<SimpleFilterContext>,
568    /// Physical filters that can be evaluated directly from the prefilter batch.
569    /// Physical expressions are only applied in the prefilter phase.
570    physical_filters: Vec<PhysicalFilterContext>,
571    /// Region schema version used in per-filter cache keys.
572    schema_version: u64,
573    /// Sorted expression strings for the encoded-PK filter group.
574    pk_filter_expr_strs: Option<SmallVec<[String; 1]>>,
575    /// Arrow schema used to build narrowed prefilter projections.
576    arrow_schema: SchemaRef,
577    /// Simple filters already proven SQL-true by this row group's statistics.
578    proven_simple_filters: Vec<bool>,
579}
580
581/// Pre-built state for constructing [PrefilterContext] per row group.
582///
583/// Fields invariant across row groups (projection mask, codec, metadata, filters)
584/// are computed once. A fresh [PrefilterContext] with its own mutable PK filter
585/// is created via [PrefilterContextBuilder::build()] for each row group.
586pub(crate) struct PrefilterContextBuilder {
587    pk_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
588    pk_filter_expr_strs: Option<SmallVec<[String; 1]>>,
589    filters: Vec<SimpleFilterContext>,
590    physical_filters: Vec<PhysicalFilterContext>,
591    codec: Arc<dyn PrimaryKeyCodec>,
592    metadata: RegionMetadataRef,
593    schema_version: u64,
594    arrow_schema: SchemaRef,
595    /// Per-row-group simple filters already proven SQL-true by column statistics.
596    proven_simple_filters: Vec<Vec<bool>>,
597}
598
599impl PrefilterContextBuilder {
600    /// Creates a builder if prefiltering is applicable.
601    ///
602    /// Returns `None` if:
603    /// - The read format doesn't use flat layout
604    /// - No prefilter columns are selected
605    /// - Prefilter would read the full projection without any PK filter
606    #[allow(clippy::too_many_arguments)]
607    pub(crate) fn new(
608        read_format: &FlatReadFormat,
609        codec: &Arc<dyn PrimaryKeyCodec>,
610        primary_key_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
611        primary_key_filter_expr_strs: Option<SmallVec<[String; 1]>>,
612        filters: Vec<SimpleFilterContext>,
613        physical_filters: Vec<PhysicalFilterContext>,
614        schema_version: u64,
615        parquet_metadata: &ParquetMetaData,
616    ) -> Option<Self> {
617        let metadata = read_format.metadata();
618        let use_raw_tag_columns = read_format.batch_has_raw_pk_columns();
619        let pk_filters = (!use_raw_tag_columns)
620            .then_some(primary_key_filters)
621            .flatten()
622            .filter(|filters| !filters.is_empty());
623        let pk_filter_expr_strs = pk_filters
624            .is_some()
625            .then_some(primary_key_filter_expr_strs)
626            .flatten();
627
628        let mut prefilter_column_names = HashSet::new();
629        for filter_ctx in &filters {
630            if let MaybeFilter::Filter(filter) = filter_ctx.filter() {
631                prefilter_column_names.insert(filter.column_name().to_string());
632            }
633        }
634
635        if pk_filters.is_some() {
636            prefilter_column_names.insert(PRIMARY_KEY_COLUMN_NAME.to_string());
637        }
638
639        for filter_ctx in &physical_filters {
640            prefilter_column_names.insert(filter_ctx.column_name().to_string());
641        }
642
643        let prefilter_count =
644            compute_projection_count(&prefilter_column_names, read_format.arrow_schema());
645
646        if prefilter_count == 0 {
647            return None;
648        }
649
650        let total_count = read_format.parquet_read_columns().root_indices().len();
651        let remaining_count = total_count.saturating_sub(prefilter_count);
652        if pk_filters.is_none() && prefilter_count >= total_count {
653            return None;
654        }
655
656        if pk_filters.is_none()
657            && !should_use_prefilter(prefilter_count, remaining_count, total_count)
658        {
659            return None;
660        }
661
662        let proven_simple_filters =
663            simple_filter_stats_proofs(read_format, parquet_metadata.row_groups(), &filters);
664
665        Some(Self {
666            pk_filters,
667            pk_filter_expr_strs,
668            filters,
669            physical_filters,
670            codec: Arc::clone(codec),
671            metadata: metadata.clone(),
672            schema_version,
673            arrow_schema: read_format.arrow_schema().clone(),
674            proven_simple_filters,
675        })
676    }
677
678    /// Builds a [PrefilterContext] for a specific row group.
679    pub(crate) fn build(&self, row_group_idx: usize) -> PrefilterContext {
680        let pk_filter = self
681            .build_primary_key_filter()
682            .map(|filter| Box::new(filter) as Box<dyn PrimaryKeyFilter>);
683        PrefilterContext {
684            pk_filter,
685            filters: self.filters.clone(),
686            physical_filters: self.physical_filters.clone(),
687            schema_version: self.schema_version,
688            pk_filter_expr_strs: self.pk_filter_expr_strs.clone(),
689            arrow_schema: self.arrow_schema.clone(),
690            proven_simple_filters: self
691                .proven_simple_filters
692                .get(row_group_idx)
693                .cloned()
694                .unwrap_or_else(|| vec![false; self.filters.len()]),
695        }
696    }
697
698    /// Builds a fresh encoded-primary-key filter selected by this plan.
699    pub(crate) fn build_primary_key_filter(&self) -> Option<CachedPrimaryKeyFilter> {
700        self.pk_filters.as_ref().map(|pk_filters| {
701            let filter = self
702                .codec
703                .primary_key_filter(&self.metadata, Arc::clone(pk_filters));
704            CachedPrimaryKeyFilter::new(filter)
705        })
706    }
707}
708
709const PREFILTER_COLUMN_RATIO_THRESHOLD: f64 = 0.5;
710const PREFILTER_MIN_REMAINING_COLUMNS: usize = 2;
711
712/// Returns row-group-major proof bits for simple filters. Statistics are
713/// extracted once per eligible filter across all row groups.
714fn simple_filter_stats_proofs(
715    read_format: &FlatReadFormat,
716    row_groups: &[parquet::file::metadata::RowGroupMetaData],
717    filters: &[SimpleFilterContext],
718) -> Vec<Vec<bool>> {
719    let mut proofs = vec![vec![false; filters.len()]; row_groups.len()];
720    for (filter_idx, filter_ctx) in filters.iter().enumerate() {
721        let Some((filter, literal)) = eligible_simple_filter(read_format, filter_ctx) else {
722            continue;
723        };
724        let (StatValues::Values(mins), StatValues::Values(maxs), StatValues::Values(null_counts)) = (
725            read_format.min_values(row_groups, filter_ctx.column_id()),
726            read_format.max_values(row_groups, filter_ctx.column_id()),
727            read_format.null_counts(row_groups, filter_ctx.column_id()),
728        ) else {
729            continue;
730        };
731        for (row_group_idx, proof) in proofs.iter_mut().enumerate() {
732            proof[filter_idx] = simple_filter_is_true_by_values(
733                filter,
734                &literal,
735                stat_value_at(&mins, row_group_idx),
736                stat_value_at(&maxs, row_group_idx),
737                stat_value_at(&null_counts, row_group_idx),
738            );
739        }
740    }
741    proofs
742}
743
744fn eligible_simple_filter<'a>(
745    read_format: &FlatReadFormat,
746    filter_ctx: &'a SimpleFilterContext,
747) -> Option<(&'a SimpleFilterEvaluator, Value)> {
748    if filter_ctx.semantic_type() != SemanticType::Field {
749        return None;
750    }
751    let filter = filter_ctx.filter().as_filter()?;
752    let literal = filter.literal_value()?;
753    let column = read_format
754        .metadata()
755        .column_by_id(filter_ctx.column_id())?;
756    column_type_matches_literal(&column.column_schema.data_type, &literal)
757        .then_some((filter, literal))
758}
759
760fn simple_filter_is_true_by_values(
761    filter: &SimpleFilterEvaluator,
762    literal: &Value,
763    min: Option<Value>,
764    max: Option<Value>,
765    null_count: Option<Value>,
766) -> bool {
767    let (Some(min), Some(max), Some(null_count)) = (min, max, null_count) else {
768        return false;
769    };
770    if null_count != Value::UInt64(0)
771        || !same_supported_value_type(&min, literal)
772        || !same_supported_value_type(&max, literal)
773        || min > max
774    {
775        return false;
776    }
777
778    if filter.is_gt() {
779        min > *literal
780    } else if filter.is_gt_eq() {
781        min >= *literal
782    } else if filter.is_lt() {
783        max < *literal
784    } else if filter.is_lt_eq() {
785        max <= *literal
786    } else if filter.is_eq() {
787        min == *literal && max == *literal
788    } else if filter.is_not_eq() {
789        max < *literal || min > *literal
790    } else {
791        false
792    }
793}
794
795fn column_type_matches_literal(data_type: &ConcreteDataType, literal: &Value) -> bool {
796    matches!(
797        (data_type, literal),
798        (ConcreteDataType::Int32(_), Value::Int32(_))
799            | (ConcreteDataType::UInt32(_), Value::UInt32(_))
800            | (ConcreteDataType::Int64(_), Value::Int64(_))
801            | (ConcreteDataType::UInt64(_), Value::UInt64(_))
802    )
803}
804
805fn same_supported_value_type(left: &Value, right: &Value) -> bool {
806    matches!(
807        (left, right),
808        (Value::Int32(_), Value::Int32(_))
809            | (Value::UInt32(_), Value::UInt32(_))
810            | (Value::Int64(_), Value::Int64(_))
811            | (Value::UInt64(_), Value::UInt64(_))
812    )
813}
814
815fn stat_value_at(values: &datatypes::arrow::array::ArrayRef, index: usize) -> Option<Value> {
816    let scalar = ScalarValue::try_from_array(values, index).ok()?;
817    Value::try_from(scalar).ok()
818}
819
820/// Result of prefiltering a row group.
821pub(crate) struct PrefilterResult {
822    /// Refined row selection after prefiltering.
823    pub(crate) refined_selection: RowSelection,
824    /// Number of rows filtered out by prefiltering.
825    pub(crate) filtered_rows: usize,
826}
827
828/// Executes prefiltering on a row group.
829///
830/// Reads only the prefilter columns (currently the PK dictionary column),
831/// applies filters, and returns a refined [RowSelection].
832fn compute_projection_mask(
833    column_names: &HashSet<String>,
834    arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
835    parquet_schema: &SchemaDescriptor,
836) -> ProjectionMask {
837    ProjectionMask::roots(
838        parquet_schema,
839        projection_indices(column_names, arrow_schema),
840    )
841}
842
843fn compute_projection_count(
844    column_names: &HashSet<String>,
845    arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
846) -> usize {
847    projection_indices(column_names, arrow_schema).len()
848}
849
850fn projection_indices(
851    column_names: &HashSet<String>,
852    arrow_schema: &datatypes::arrow::datatypes::SchemaRef,
853) -> Vec<usize> {
854    let mut projection_indices: Vec<usize> = column_names
855        .iter()
856        .filter_map(|name| arrow_schema.column_with_name(name).map(|(index, _)| index))
857        .collect();
858    projection_indices.sort_unstable();
859    projection_indices.dedup();
860    projection_indices
861}
862
863fn should_use_prefilter(
864    prefilter_count: usize,
865    remaining_count: usize,
866    total_count: usize,
867) -> bool {
868    if remaining_count == 0 {
869        return false;
870    }
871
872    if remaining_count < PREFILTER_MIN_REMAINING_COLUMNS {
873        return false;
874    }
875
876    let ratio = prefilter_count as f64 / total_count as f64;
877    ratio <= PREFILTER_COLUMN_RATIO_THRESHOLD
878}
879
880pub(crate) async fn execute_prefilter(
881    prefilter_ctx: &mut PrefilterContext,
882    reader_builder: &RowGroupReaderBuilder,
883    build_ctx: &RowGroupBuildContext<'_>,
884) -> Result<PrefilterResult> {
885    let entries = build_prefilter_cache_entries(prefilter_ctx, reader_builder, build_ctx);
886
887    if entries.is_empty() {
888        return execute_prefilter_by_reading_columns(prefilter_ctx, reader_builder, build_ctx)
889            .await;
890    }
891
892    execute_prefilter_with_result_cache(prefilter_ctx, reader_builder, build_ctx, entries).await
893}
894
895async fn execute_prefilter_with_result_cache(
896    prefilter_ctx: &mut PrefilterContext,
897    reader_builder: &RowGroupReaderBuilder,
898    build_ctx: &RowGroupBuildContext<'_>,
899    entries: Vec<PrefilterEntry>,
900) -> Result<PrefilterResult> {
901    let non_cacheable_physical = non_cacheable_physical_filters(prefilter_ctx);
902    let mut hit_mask: Option<BooleanBuffer> = None;
903    let mut misses = Vec::new();
904    for entry in entries {
905        let Some(key) = &entry.key else {
906            misses.push(entry);
907            continue;
908        };
909
910        if let Some(mask) = reader_builder.cache_strategy().get_prefilter_result(key) {
911            hit_mask = Some(match hit_mask {
912                Some(hit_mask) => hit_mask.bitand(mask.as_ref()),
913                None => mask.as_ref().clone(),
914            });
915        } else {
916            misses.push(entry);
917        }
918    }
919
920    if misses.is_empty() && non_cacheable_physical.is_empty() {
921        let combined_mask = hit_mask.unwrap_or_else(|| BooleanBuffer::new_set(0));
922        let refined_selection =
923            refined_selection_from_mask(&combined_mask, &build_ctx.row_selection);
924        let rows_before_filter = rows_before_filter(reader_builder, build_ctx);
925        let filtered_rows = rows_before_filter.saturating_sub(refined_selection.row_count());
926        return Ok(PrefilterResult {
927            refined_selection,
928            filtered_rows,
929        });
930    }
931
932    let mut uncached_entries = misses;
933    uncached_entries.extend(
934        non_cacheable_physical
935            .iter()
936            .copied()
937            .map(|idx| PrefilterEntry::without_cache(PrefilterEntryKind::Physical(idx))),
938    );
939    let (uncached_mask, read_rows) =
940        build_prefilter_masks(prefilter_ctx, reader_builder, build_ctx, &uncached_entries).await?;
941
942    let final_mask = match (hit_mask, uncached_mask) {
943        (Some(hit_mask), Some(uncached_mask)) => hit_mask.bitand(&uncached_mask),
944        (Some(hit_mask), None) => hit_mask,
945        (None, Some(uncached_mask)) => uncached_mask,
946        (None, None) => BooleanBuffer::new_set(read_rows),
947    };
948    debug_assert_eq!(final_mask.len(), read_rows);
949    let rows_selected = final_mask.count_set_bits();
950    let filtered_rows = read_rows.saturating_sub(rows_selected);
951    let refined_selection = refined_selection_from_mask(&final_mask, &build_ctx.row_selection);
952
953    Ok(PrefilterResult {
954        refined_selection,
955        filtered_rows,
956    })
957}
958
959fn non_cacheable_physical_filters(prefilter_ctx: &PrefilterContext) -> Vec<usize> {
960    prefilter_ctx
961        .physical_filters
962        .iter()
963        .enumerate()
964        .filter_map(|(idx, filter)| (!filter.is_immutable()).then_some(idx))
965        .collect()
966}
967
968async fn build_prefilter_masks(
969    prefilter_ctx: &mut PrefilterContext,
970    reader_builder: &RowGroupReaderBuilder,
971    build_ctx: &RowGroupBuildContext<'_>,
972    entries: &[PrefilterEntry],
973) -> Result<(Option<BooleanBuffer>, usize)> {
974    let prefilter_column_names = prefilter_column_names_for_entries(prefilter_ctx, entries);
975    let parquet_schema = reader_builder
976        .parquet_metadata()
977        .file_metadata()
978        .schema_descr();
979    let projection = compute_projection_mask(
980        &prefilter_column_names,
981        &prefilter_ctx.arrow_schema,
982        parquet_schema,
983    );
984
985    let mut stream = reader_builder
986        .build_with_projection(
987            build_ctx.row_group_idx,
988            build_ctx.row_selection.clone(),
989            projection,
990            build_ctx.fetch_metrics,
991        )
992        .await?;
993
994    let mut cache_builders = entries
995        .iter()
996        .map(|entry| entry.key.is_some().then(|| BooleanBufferBuilder::new(0)))
997        .collect::<Vec<_>>();
998    let mut combined_builder = (!entries.is_empty()).then(|| BooleanBufferBuilder::new(0));
999    let mut rows_before_filter = 0usize;
1000
1001    while let Some(batch_result) = stream.next().await {
1002        let batch = batch_result?;
1003        let num_rows = batch.num_rows();
1004        if num_rows == 0 {
1005            continue;
1006        }
1007        rows_before_filter += num_rows;
1008
1009        let mut batch_mask = BooleanBuffer::new_set(num_rows);
1010        for (idx, entry) in entries.iter().enumerate() {
1011            let mask = eval_entry_mask(
1012                &batch,
1013                prefilter_ctx,
1014                entry.kind,
1015                reader_builder.file_path(),
1016            )?;
1017            batch_mask = batch_mask.bitand(&mask);
1018            if let Some(Some(builder)) = cache_builders.get_mut(idx) {
1019                builder.append_buffer(&mask);
1020            }
1021        }
1022        if let Some(builder) = &mut combined_builder {
1023            builder.append_buffer(&batch_mask);
1024        }
1025    }
1026
1027    for (entry, builder) in entries.iter().zip(cache_builders) {
1028        if let (Some(key), Some(mut builder)) = (&entry.key, builder) {
1029            reader_builder
1030                .cache_strategy()
1031                .put_prefilter_result(key.clone(), Arc::new(builder.finish()));
1032        }
1033    }
1034
1035    Ok((
1036        combined_builder.map(|mut builder| builder.finish()),
1037        rows_before_filter,
1038    ))
1039}
1040
1041fn prefilter_column_names_for_entries(
1042    prefilter_ctx: &PrefilterContext,
1043    entries: &[PrefilterEntry],
1044) -> HashSet<String> {
1045    let mut prefilter_column_names = HashSet::new();
1046    for entry in entries {
1047        match entry.kind {
1048            PrefilterEntryKind::Simple(idx) => {
1049                if let MaybeFilter::Filter(filter) = prefilter_ctx.filters[idx].filter() {
1050                    prefilter_column_names.insert(filter.column_name().to_string());
1051                }
1052            }
1053            PrefilterEntryKind::Physical(idx) => {
1054                prefilter_column_names.insert(
1055                    prefilter_ctx.physical_filters[idx]
1056                        .column_name()
1057                        .to_string(),
1058                );
1059            }
1060            PrefilterEntryKind::PkGroup => {
1061                prefilter_column_names.insert(PRIMARY_KEY_COLUMN_NAME.to_string());
1062            }
1063        }
1064    }
1065    prefilter_column_names
1066}
1067
1068async fn execute_prefilter_by_reading_columns(
1069    prefilter_ctx: &mut PrefilterContext,
1070    reader_builder: &RowGroupReaderBuilder,
1071    build_ctx: &RowGroupBuildContext<'_>,
1072) -> Result<PrefilterResult> {
1073    let entries = all_prefilter_entries(prefilter_ctx);
1074    if entries.is_empty() {
1075        return Ok(identity_prefilter_result(reader_builder, build_ctx));
1076    }
1077    let (mask, rows_before_filter) =
1078        build_prefilter_masks(prefilter_ctx, reader_builder, build_ctx, &entries).await?;
1079
1080    let final_mask = mask.unwrap_or_else(|| BooleanBuffer::new_set(rows_before_filter));
1081    let rows_selected = final_mask.count_set_bits();
1082    let filtered_rows = rows_before_filter.saturating_sub(rows_selected);
1083    let refined_selection = refined_selection_from_mask(&final_mask, &build_ctx.row_selection);
1084
1085    Ok(PrefilterResult {
1086        refined_selection,
1087        filtered_rows,
1088    })
1089}
1090
1091fn all_prefilter_entries(prefilter_ctx: &PrefilterContext) -> Vec<PrefilterEntry> {
1092    let mut entries = Vec::new();
1093    if prefilter_ctx.pk_filter.is_some() {
1094        entries.push(PrefilterEntry::without_cache(PrefilterEntryKind::PkGroup));
1095    }
1096    entries.extend(
1097        prefilter_ctx
1098            .filters
1099            .iter()
1100            .enumerate()
1101            .filter(|(idx, _)| {
1102                !prefilter_ctx
1103                    .proven_simple_filters
1104                    .get(*idx)
1105                    .copied()
1106                    .unwrap_or(false)
1107            })
1108            .map(|(idx, _)| PrefilterEntry::without_cache(PrefilterEntryKind::Simple(idx))),
1109    );
1110    entries.extend(
1111        prefilter_ctx
1112            .physical_filters
1113            .iter()
1114            .enumerate()
1115            .map(|(idx, _)| PrefilterEntry::without_cache(PrefilterEntryKind::Physical(idx))),
1116    );
1117    entries
1118}
1119
1120#[derive(Clone, Copy)]
1121enum PrefilterEntryKind {
1122    Simple(usize),
1123    Physical(usize),
1124    PkGroup,
1125}
1126
1127struct PrefilterEntry {
1128    kind: PrefilterEntryKind,
1129    key: Option<PrefilterKey>,
1130}
1131
1132impl PrefilterEntry {
1133    fn without_cache(kind: PrefilterEntryKind) -> Self {
1134        Self { kind, key: None }
1135    }
1136}
1137
1138fn build_prefilter_cache_entries(
1139    prefilter_ctx: &PrefilterContext,
1140    reader_builder: &RowGroupReaderBuilder,
1141    build_ctx: &RowGroupBuildContext<'_>,
1142) -> Vec<PrefilterEntry> {
1143    let row_selection = PrefilterKey::row_selection_snapshot(build_ctx.row_selection.as_ref());
1144    let file_id = reader_builder.file_handle().file_id().file_id();
1145    let row_group_idx = build_ctx.row_group_idx as u32;
1146    let mut entries = Vec::new();
1147
1148    for (idx, filter_ctx) in prefilter_ctx.filters.iter().enumerate() {
1149        if prefilter_ctx
1150            .proven_simple_filters
1151            .get(idx)
1152            .copied()
1153            .unwrap_or(false)
1154        {
1155            continue;
1156        }
1157        entries.push(PrefilterEntry {
1158            kind: PrefilterEntryKind::Simple(idx),
1159            key: Some(PrefilterKey::new(
1160                file_id,
1161                row_group_idx,
1162                row_selection.clone(),
1163                prefilter_ctx.schema_version,
1164                smallvec![filter_ctx.expr_str().to_string()],
1165            )),
1166        });
1167    }
1168
1169    for (idx, filter_ctx) in prefilter_ctx.physical_filters.iter().enumerate() {
1170        if !filter_ctx.is_immutable() {
1171            continue;
1172        }
1173        entries.push(PrefilterEntry {
1174            kind: PrefilterEntryKind::Physical(idx),
1175            key: Some(PrefilterKey::new(
1176                file_id,
1177                row_group_idx,
1178                row_selection.clone(),
1179                prefilter_ctx.schema_version,
1180                smallvec![filter_ctx.expr_str().to_string()],
1181            )),
1182        });
1183    }
1184
1185    if prefilter_ctx.pk_filter.is_some()
1186        && let Some(exprs) = &prefilter_ctx.pk_filter_expr_strs
1187    {
1188        entries.push(PrefilterEntry {
1189            kind: PrefilterEntryKind::PkGroup,
1190            key: Some(PrefilterKey::new(
1191                file_id,
1192                row_group_idx,
1193                row_selection,
1194                prefilter_ctx.schema_version,
1195                exprs.clone(),
1196            )),
1197        });
1198    }
1199
1200    entries
1201}
1202
1203fn identity_prefilter_result(
1204    reader_builder: &RowGroupReaderBuilder,
1205    build_ctx: &RowGroupBuildContext<'_>,
1206) -> PrefilterResult {
1207    let row_count = reader_builder
1208        .parquet_metadata()
1209        .row_group(build_ctx.row_group_idx)
1210        .num_rows() as usize;
1211    PrefilterResult {
1212        refined_selection: identity_row_selection(&build_ctx.row_selection, row_count),
1213        filtered_rows: 0,
1214    }
1215}
1216
1217fn identity_row_selection(
1218    original_selection: &Option<RowSelection>,
1219    row_count: usize,
1220) -> RowSelection {
1221    original_selection
1222        .clone()
1223        .unwrap_or_else(|| RowSelection::from(vec![RowSelector::select(row_count)]))
1224}
1225
1226fn rows_before_filter(
1227    reader_builder: &RowGroupReaderBuilder,
1228    build_ctx: &RowGroupBuildContext<'_>,
1229) -> usize {
1230    build_ctx.row_selection.as_ref().map_or_else(
1231        || {
1232            reader_builder
1233                .parquet_metadata()
1234                .row_group(build_ctx.row_group_idx)
1235                .num_rows() as usize
1236        },
1237        RowSelection::row_count,
1238    )
1239}
1240
1241fn refined_selection_from_mask(
1242    mask: &BooleanBuffer,
1243    original_selection: &Option<RowSelection>,
1244) -> RowSelection {
1245    if mask.is_empty() || mask.count_set_bits() == 0 {
1246        return RowSelection::from(vec![]);
1247    }
1248
1249    let prefilter_selection = RowSelection::from_filters(&[BooleanArray::from(mask.clone())]);
1250    match original_selection {
1251        Some(original) => original.and_then(&prefilter_selection),
1252        None => prefilter_selection,
1253    }
1254}
1255
1256fn eval_entry_mask(
1257    batch: &RecordBatch,
1258    prefilter_ctx: &mut PrefilterContext,
1259    kind: PrefilterEntryKind,
1260    file_path: &str,
1261) -> Result<BooleanBuffer> {
1262    match kind {
1263        PrefilterEntryKind::Simple(idx) => {
1264            eval_simple_filter_mask(batch, &prefilter_ctx.filters[idx], file_path)
1265        }
1266        PrefilterEntryKind::Physical(idx) => {
1267            eval_physical_filter_mask(batch, &prefilter_ctx.physical_filters[idx], file_path)
1268        }
1269        PrefilterEntryKind::PkGroup => {
1270            let pk_filter = prefilter_ctx.pk_filter.as_mut().context(UnexpectedSnafu {
1271                reason: "Missing primary key filter for prefilter cache entry",
1272            })?;
1273            primary_key_filter_mask(batch, pk_filter.as_mut())
1274        }
1275    }
1276}
1277
1278/// Evaluates an encoded-primary-key filter against a primary-key batch.
1279pub(crate) fn primary_key_filter_mask(
1280    batch: &RecordBatch,
1281    pk_filter: &mut dyn PrimaryKeyFilter,
1282) -> Result<BooleanBuffer> {
1283    let (pk_column_index, _) = batch
1284        .schema()
1285        .column_with_name(PRIMARY_KEY_COLUMN_NAME)
1286        .context(UnexpectedSnafu {
1287            reason: "Primary key column not found in prefilter batch",
1288        })?;
1289    let matched_row_ranges = matching_row_ranges_by_primary_key(batch, pk_column_index, pk_filter)?;
1290    let mut builder = BooleanBufferBuilder::new(batch.num_rows());
1291    builder.append_n(batch.num_rows(), false);
1292    for range in matched_row_ranges {
1293        for row in range {
1294            builder.set_bit(row, true);
1295        }
1296    }
1297    Ok(builder.finish())
1298}
1299
1300fn eval_simple_filter_mask(
1301    batch: &RecordBatch,
1302    filter_ctx: &SimpleFilterContext,
1303    file_path: &str,
1304) -> Result<BooleanBuffer> {
1305    let filter = match filter_ctx.filter() {
1306        MaybeFilter::Filter(filter) => filter,
1307        MaybeFilter::Matched => return Ok(BooleanBuffer::new_set(batch.num_rows())),
1308        MaybeFilter::Pruned => return Ok(BooleanBuffer::new_unset(batch.num_rows())),
1309    };
1310
1311    let (idx, _) = batch
1312        .schema()
1313        .column_with_name(filter.column_name())
1314        .with_context(|| UnexpectedSnafu {
1315            reason: format!(
1316                "Prefilter column '{}' (id {}) not found in batch for file {}",
1317                filter.column_name(),
1318                filter_ctx.column_id(),
1319                file_path
1320            ),
1321        })?;
1322    let column = batch.column(idx).clone();
1323    filter.evaluate_array(&column).context(RecordBatchSnafu)
1324}
1325
1326fn eval_physical_filter_mask(
1327    batch: &RecordBatch,
1328    filter_ctx: &PhysicalFilterContext,
1329    file_path: &str,
1330) -> Result<BooleanBuffer> {
1331    let filter = filter_ctx.filter();
1332
1333    let (idx, _) = batch
1334        .schema()
1335        .column_with_name(filter_ctx.column_name())
1336        .with_context(|| UnexpectedSnafu {
1337            reason: format!(
1338                "Prefilter physical column '{}' (id {}) not found in batch for file {}",
1339                filter_ctx.column_name(),
1340                filter_ctx.column_id(),
1341                file_path
1342            ),
1343        })?;
1344    let column = batch.column(idx).clone();
1345
1346    let record_batch = RecordBatch::try_new(filter_ctx.schema().clone(), vec![column])
1347        .context(NewRecordBatchSnafu)?;
1348    let evaluated = filter
1349        .evaluate(&record_batch)
1350        .context(EvalPartitionFilterSnafu)?;
1351    let array = evaluated
1352        .into_array(record_batch.num_rows())
1353        .context(EvalPartitionFilterSnafu)?;
1354    let boolean_array = array
1355        .as_any()
1356        .downcast_ref::<BooleanArray>()
1357        .context(UnexpectedSnafu {
1358            reason: "Failed to downcast physical filter result to BooleanArray",
1359        })?;
1360    // Treat null results as false (filtered out); value bits are not guaranteed
1361    // to be false for invalid entries.
1362    let mut result = boolean_array.values().clone();
1363    if let Some(nulls) = boolean_array.nulls() {
1364        result = result.bitand(nulls.inner());
1365    }
1366    Ok(result)
1367}
1368
1369#[cfg(test)]
1370mod tests {
1371    use std::sync::Arc;
1372    use std::sync::atomic::{AtomicUsize, Ordering};
1373
1374    use bytes::Bytes;
1375    use common_recordbatch::filter::SimpleFilterEvaluator;
1376    use datafusion_common::ScalarValue;
1377    use datafusion_expr::{col, lit};
1378    use datatypes::arrow::array::{
1379        ArrayRef, DictionaryArray, Int32Array, TimestampMillisecondArray, UInt8Array, UInt32Array,
1380        UInt64Array,
1381    };
1382    use datatypes::arrow::datatypes::{DataType, Field, Schema, UInt32Type};
1383    use datatypes::arrow::record_batch::RecordBatch;
1384    use datatypes::prelude::ConcreteDataType;
1385    use datatypes::value::Value;
1386    use mito_codec::row_converter::{PrimaryKeyFilter, build_primary_key_codec};
1387    use parquet::arrow::ArrowWriter;
1388    use parquet::arrow::arrow_reader::{ParquetRecordBatchReaderBuilder, RowSelector};
1389    use store_api::codec::PrimaryKeyEncoding;
1390    use store_api::metadata::RegionMetadataBuilder;
1391    use store_api::region_request::{AlterKind, ModifyColumnType};
1392
1393    use super::*;
1394    use crate::read::read_columns::ReadColumns;
1395    use crate::sst::internal_fields;
1396    use crate::sst::parquet::flat_format::{FlatReadFormat, primary_key_column_index};
1397    use crate::test_util::sst_util::{
1398        new_primary_key, new_record_batch_with_custom_sequence, sst_region_metadata,
1399        sst_region_metadata_with_encoding,
1400    };
1401
1402    struct CountingPrimaryKeyFilter {
1403        hits: Arc<AtomicUsize>,
1404        expected: Vec<u8>,
1405    }
1406
1407    impl PrimaryKeyFilter for CountingPrimaryKeyFilter {
1408        fn matches(&mut self, pk: &[u8]) -> mito_codec::error::Result<bool> {
1409            self.hits.fetch_add(1, Ordering::Relaxed);
1410            Ok(pk == self.expected.as_slice())
1411        }
1412    }
1413
1414    #[test]
1415    fn test_cached_primary_key_filter_reuses_previous_result() {
1416        let expected = new_primary_key(&["a", "x"]);
1417        let hits = Arc::new(AtomicUsize::new(0));
1418        let mut filter = CachedPrimaryKeyFilter::new(Box::new(CountingPrimaryKeyFilter {
1419            hits: Arc::clone(&hits),
1420            expected: expected.clone(),
1421        }));
1422
1423        assert!(filter.matches(expected.as_slice()).unwrap());
1424        assert!(filter.matches(expected.as_slice()).unwrap());
1425        assert!(
1426            !filter
1427                .matches(new_primary_key(&["b", "x"]).as_slice())
1428                .unwrap()
1429        );
1430
1431        assert_eq!(hits.load(Ordering::Relaxed), 2);
1432    }
1433
1434    fn new_test_filters(exprs: &[datafusion_expr::Expr]) -> Vec<SimpleFilterEvaluator> {
1435        exprs
1436            .iter()
1437            .filter_map(SimpleFilterEvaluator::try_new)
1438            .collect()
1439    }
1440
1441    fn new_simple_filter_contexts(
1442        metadata: &RegionMetadataRef,
1443        exprs: &[datafusion_expr::Expr],
1444    ) -> Vec<SimpleFilterContext> {
1445        exprs
1446            .iter()
1447            .filter_map(|expr| SimpleFilterContext::new_opt(metadata, None, expr))
1448            .collect()
1449    }
1450
1451    fn stats_metadata_with_options(
1452        values: &[Vec<u64>],
1453        nulls: Option<&[bool]>,
1454        writer_options: Option<parquet::file::properties::WriterProperties>,
1455    ) -> Arc<ParquetMetaData> {
1456        let first = new_record_batch_with_custom_sequence(&["a", "x"], 0, values[0].len(), 1);
1457        let mut bytes = Vec::new();
1458        let mut writer = ArrowWriter::try_new(&mut bytes, first.schema(), writer_options).unwrap();
1459        for (idx, values) in values.iter().enumerate() {
1460            let has_null = nulls.and_then(|nulls| nulls.get(idx)).copied() == Some(true);
1461            let batch = new_record_batch_with_custom_sequence(
1462                &["a", "x"],
1463                0,
1464                if has_null {
1465                    values.len().max(2)
1466                } else {
1467                    values.len()
1468                },
1469                1,
1470            );
1471            let mut columns = batch.columns().to_vec();
1472            columns[2] = Arc::new(if has_null {
1473                UInt64Array::from(vec![Some(values[0]), None])
1474            } else {
1475                UInt64Array::from_iter_values(values.iter().copied())
1476            });
1477            writer
1478                .write(&RecordBatch::try_new(batch.schema(), columns).unwrap())
1479                .unwrap();
1480            writer.flush().unwrap();
1481        }
1482        writer.close().unwrap();
1483        ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
1484            .unwrap()
1485            .metadata()
1486            .clone()
1487    }
1488
1489    fn stats_metadata(values: &[Vec<u64>], nulls: Option<&[bool]>) -> Arc<ParquetMetaData> {
1490        stats_metadata_with_options(values, nulls, None)
1491    }
1492
1493    fn int32_stats_metadata(values: &[i32]) -> Arc<ParquetMetaData> {
1494        let batch = new_record_batch_with_custom_sequence(&["a", "x"], 0, values.len(), 1);
1495        let mut fields = batch
1496            .schema()
1497            .fields()
1498            .iter()
1499            .map(|field| field.as_ref().clone())
1500            .collect::<Vec<_>>();
1501        fields[2].set_data_type(DataType::Int32);
1502        let mut columns = batch.columns().to_vec();
1503        columns[2] = Arc::new(Int32Array::from(values.to_vec()));
1504        let batch = RecordBatch::try_new(Arc::new(Schema::new(fields)), columns).unwrap();
1505        let mut bytes = Vec::new();
1506        let mut writer = ArrowWriter::try_new(&mut bytes, batch.schema(), None).unwrap();
1507        writer.write(&batch).unwrap();
1508        writer.close().unwrap();
1509        ParquetRecordBatchReaderBuilder::try_new(Bytes::from(bytes))
1510            .unwrap()
1511            .metadata()
1512            .clone()
1513    }
1514
1515    fn metadata_with_field_type(data_type: ConcreteDataType) -> RegionMetadataRef {
1516        let mut builder = RegionMetadataBuilder::from_existing(sst_region_metadata());
1517        builder
1518            .alter(AlterKind::ModifyColumnTypes {
1519                columns: vec![ModifyColumnType {
1520                    column_name: "field_0".to_string(),
1521                    target_type: data_type,
1522                }],
1523            })
1524            .unwrap();
1525        Arc::new(builder.build().unwrap())
1526    }
1527
1528    fn new_physical_filter_contexts(
1529        metadata: &RegionMetadataRef,
1530        read_format: &FlatReadFormat,
1531        exprs: &[datafusion_expr::Expr],
1532    ) -> Vec<PhysicalFilterContext> {
1533        exprs
1534            .iter()
1535            .filter_map(|expr| PhysicalFilterContext::new_opt(metadata, None, read_format, expr))
1536            .collect()
1537    }
1538
1539    fn new_raw_batch(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1540        assert_eq!(primary_keys.len(), field_values.len());
1541
1542        let metadata = Arc::new(sst_region_metadata());
1543        let arrow_schema = metadata.schema.arrow_schema();
1544        let field_column = arrow_schema
1545            .field(arrow_schema.index_of("field_0").unwrap())
1546            .clone();
1547        let time_index_column = arrow_schema
1548            .field(arrow_schema.index_of("ts").unwrap())
1549            .clone();
1550        let mut fields = vec![field_column, time_index_column];
1551        fields.extend(
1552            internal_fields()
1553                .into_iter()
1554                .map(|field| field.as_ref().clone()),
1555        );
1556        let schema = Arc::new(Schema::new(fields));
1557
1558        let mut dict_values = Vec::new();
1559        let mut keys = Vec::with_capacity(primary_keys.len());
1560        for pk in primary_keys {
1561            let key = dict_values
1562                .iter()
1563                .position(|existing: &&[u8]| existing == pk)
1564                .unwrap_or_else(|| {
1565                    dict_values.push(*pk);
1566                    dict_values.len() - 1
1567                });
1568            keys.push(key as u32);
1569        }
1570        let pk_array: ArrayRef = Arc::new(DictionaryArray::<UInt32Type>::new(
1571            UInt32Array::from(keys),
1572            Arc::new(BinaryArray::from_iter_values(dict_values.iter().copied())),
1573        ));
1574
1575        RecordBatch::try_new(
1576            schema,
1577            vec![
1578                Arc::new(UInt64Array::from(field_values.to_vec())),
1579                Arc::new(TimestampMillisecondArray::from_iter_values(
1580                    0..primary_keys.len() as i64,
1581                )),
1582                pk_array,
1583                Arc::new(UInt64Array::from(vec![1; primary_keys.len()])),
1584                Arc::new(UInt8Array::from(vec![1; primary_keys.len()])),
1585            ],
1586        )
1587        .unwrap()
1588    }
1589
1590    fn new_prefilter_batch(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1591        assert_eq!(primary_keys.len(), field_values.len());
1592
1593        let metadata = Arc::new(sst_region_metadata());
1594        let arrow_schema = metadata.schema.arrow_schema();
1595        let field_column = arrow_schema
1596            .field(arrow_schema.index_of("field_0").unwrap())
1597            .clone();
1598        let time_index_column = arrow_schema
1599            .field(arrow_schema.index_of("ts").unwrap())
1600            .clone();
1601        let schema = Arc::new(Schema::new(vec![
1602            field_column,
1603            time_index_column,
1604            internal_fields()[0].as_ref().clone(),
1605        ]));
1606
1607        let mut dict_values = Vec::new();
1608        let mut keys = Vec::with_capacity(primary_keys.len());
1609        for pk in primary_keys {
1610            let key = dict_values
1611                .iter()
1612                .position(|existing: &&[u8]| existing == pk)
1613                .unwrap_or_else(|| {
1614                    dict_values.push(*pk);
1615                    dict_values.len() - 1
1616                });
1617            keys.push(key as u32);
1618        }
1619        let pk_array: ArrayRef = Arc::new(DictionaryArray::<UInt32Type>::new(
1620            UInt32Array::from(keys),
1621            Arc::new(BinaryArray::from_iter_values(dict_values.iter().copied())),
1622        ));
1623
1624        RecordBatch::try_new(
1625            schema,
1626            vec![
1627                Arc::new(UInt64Array::from(field_values.to_vec())),
1628                Arc::new(TimestampMillisecondArray::from_iter_values(
1629                    0..primary_keys.len() as i64,
1630                )),
1631                pk_array,
1632            ],
1633        )
1634        .unwrap()
1635    }
1636
1637    fn new_prefilter_batch_binary_pk(primary_keys: &[&[u8]], field_values: &[u64]) -> RecordBatch {
1638        assert_eq!(primary_keys.len(), field_values.len());
1639
1640        let metadata = Arc::new(sst_region_metadata());
1641        let arrow_schema = metadata.schema.arrow_schema();
1642        let field_column = arrow_schema
1643            .field(arrow_schema.index_of("field_0").unwrap())
1644            .clone();
1645        let time_index_column = arrow_schema
1646            .field(arrow_schema.index_of("ts").unwrap())
1647            .clone();
1648        let schema = Arc::new(Schema::new(vec![
1649            field_column,
1650            time_index_column,
1651            Field::new(PRIMARY_KEY_COLUMN_NAME, DataType::Binary, false),
1652        ]));
1653
1654        let pk_array: ArrayRef =
1655            Arc::new(BinaryArray::from_iter_values(primary_keys.iter().copied()));
1656
1657        RecordBatch::try_new(
1658            schema,
1659            vec![
1660                Arc::new(UInt64Array::from(field_values.to_vec())),
1661                Arc::new(TimestampMillisecondArray::from_iter_values(
1662                    0..primary_keys.len() as i64,
1663                )),
1664                pk_array,
1665            ],
1666        )
1667        .unwrap()
1668    }
1669
1670    fn field_values(batch: &RecordBatch) -> Vec<u64> {
1671        batch
1672            .column(0)
1673            .as_any()
1674            .downcast_ref::<UInt64Array>()
1675            .unwrap()
1676            .values()
1677            .to_vec()
1678    }
1679
1680    fn remaining_simple_filter_columns(filters: &[SimpleFilterContext]) -> Vec<&str> {
1681        filters
1682            .iter()
1683            .map(|filter_ctx| filter_ctx.filter().as_filter().unwrap().column_name())
1684            .collect()
1685    }
1686
1687    #[test]
1688    fn test_prefilter_primary_key_drops_single_dictionary_batch() {
1689        let metadata = Arc::new(sst_region_metadata());
1690        let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("b"))]));
1691        let mut primary_key_filter =
1692            build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters);
1693        let pk_a = new_primary_key(&["a", "x"]);
1694        let batch = new_raw_batch(&[pk_a.as_slice(), pk_a.as_slice()], &[10, 11]);
1695        let pk_col_idx = primary_key_column_index(batch.num_columns());
1696
1697        let filtered =
1698            prefilter_flat_batch_by_primary_key(batch, pk_col_idx, primary_key_filter.as_mut())
1699                .unwrap();
1700
1701        assert!(filtered.is_none());
1702    }
1703
1704    #[test]
1705    fn test_prefilter_primary_key_builds_mask_for_fragmented_matches() {
1706        let metadata = Arc::new(sst_region_metadata());
1707        let filters = Arc::new(new_test_filters(&[col("tag_0")
1708            .eq(lit("a"))
1709            .or(col("tag_0").eq(lit("c")))]));
1710        let mut primary_key_filter =
1711            build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters);
1712        let pk_a = new_primary_key(&["a", "x"]);
1713        let pk_b = new_primary_key(&["b", "x"]);
1714        let pk_c = new_primary_key(&["c", "x"]);
1715        let pk_d = new_primary_key(&["d", "x"]);
1716        let batch = new_raw_batch(
1717            &[
1718                pk_a.as_slice(),
1719                pk_a.as_slice(),
1720                pk_b.as_slice(),
1721                pk_b.as_slice(),
1722                pk_c.as_slice(),
1723                pk_c.as_slice(),
1724                pk_d.as_slice(),
1725                pk_d.as_slice(),
1726            ],
1727            &[10, 11, 12, 13, 14, 15, 16, 17],
1728        );
1729        let pk_col_idx = primary_key_column_index(batch.num_columns());
1730
1731        let filtered =
1732            prefilter_flat_batch_by_primary_key(batch, pk_col_idx, primary_key_filter.as_mut())
1733                .unwrap()
1734                .unwrap();
1735
1736        assert_eq!(filtered.num_rows(), 4);
1737        assert_eq!(field_values(&filtered), vec![10, 11, 14, 15]);
1738    }
1739
1740    #[test]
1741    fn test_prefilter_builder_returns_none_without_selected_filters() {
1742        let metadata: RegionMetadataRef =
1743            Arc::new(sst_region_metadata_with_encoding(PrimaryKeyEncoding::Dense));
1744        let read_format = FlatReadFormat::new(
1745            metadata.clone(),
1746            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
1747            None,
1748            "test",
1749            false,
1750        )
1751        .unwrap();
1752        let codec = build_primary_key_codec(metadata.as_ref());
1753
1754        let builder = PrefilterContextBuilder::new(
1755            &read_format,
1756            &codec,
1757            None,
1758            None,
1759            Vec::new(),
1760            Vec::new(),
1761            metadata.schema_version,
1762            &stats_metadata(&[vec![1]], None),
1763        );
1764        assert!(builder.is_none());
1765    }
1766
1767    #[test]
1768    fn test_simple_filter_stats_uses_real_metadata() {
1769        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1770        let read_format = FlatReadFormat::new(
1771            metadata.clone(),
1772            ReadColumns::new(
1773                metadata
1774                    .column_metadatas
1775                    .iter()
1776                    .map(|column| column.column_id),
1777            ),
1778            None,
1779            "test",
1780            true,
1781        )
1782        .unwrap();
1783        let parquet_metadata = stats_metadata(&[vec![i64::MAX as u64 + 2]], None);
1784        let filters = new_simple_filter_contexts(
1785            &metadata,
1786            &[
1787                col("field_0").gt(lit(i64::MAX as u64 + 1)),
1788                col("field_0").gt(lit(i64::MAX as u64 + 2)),
1789                lit(i64::MAX as u64 + 1).lt(col("field_0")),
1790            ],
1791        );
1792
1793        let proofs =
1794            simple_filter_stats_proofs(&read_format, parquet_metadata.row_groups(), &filters);
1795        assert_eq!(proofs, vec![vec![true, false, true]]);
1796    }
1797
1798    #[test]
1799    fn test_simple_filter_stats_retain_narrow_or_incomplete_metadata() {
1800        let parquet_metadata = stats_metadata(&[vec![2]], Some(&[true]));
1801        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1802        let read_format = FlatReadFormat::new(
1803            metadata.clone(),
1804            ReadColumns::new(
1805                metadata
1806                    .column_metadatas
1807                    .iter()
1808                    .map(|column| column.column_id),
1809            ),
1810            None,
1811            "test",
1812            true,
1813        )
1814        .unwrap();
1815        let filter = new_simple_filter_contexts(&metadata, &[col("field_0").gt(lit(1_u64))]);
1816        assert_eq!(
1817            simple_filter_stats_proofs(&read_format, parquet_metadata.row_groups(), &filter),
1818            vec![vec![false]],
1819        );
1820        let missing_stats = stats_metadata_with_options(
1821            &[vec![2]],
1822            None,
1823            Some(
1824                parquet::file::properties::WriterProperties::builder()
1825                    .set_statistics_enabled(parquet::file::properties::EnabledStatistics::None)
1826                    .build(),
1827            ),
1828        );
1829        assert_eq!(
1830            simple_filter_stats_proofs(&read_format, missing_stats.row_groups(), &filter),
1831            vec![vec![false]],
1832        );
1833
1834        // Parquet INT32 stats decode as Value::Int32; without the actual SST
1835        // type gate this Int8 field would be incorrectly proven true.
1836        let narrow_metadata = metadata_with_field_type(ConcreteDataType::int8_datatype());
1837        let narrow_read_format = FlatReadFormat::new(
1838            narrow_metadata.clone(),
1839            ReadColumns::new(
1840                narrow_metadata
1841                    .column_metadatas
1842                    .iter()
1843                    .map(|column| column.column_id),
1844            ),
1845            None,
1846            "test",
1847            true,
1848        )
1849        .unwrap();
1850        let narrow_filter =
1851            new_simple_filter_contexts(&narrow_metadata, &[col("field_0").gt(lit(1_i32))]);
1852        let narrow_stats = int32_stats_metadata(&[2]);
1853        assert_eq!(
1854            simple_filter_stats_proofs(
1855                &narrow_read_format,
1856                narrow_stats.row_groups(),
1857                &narrow_filter,
1858            ),
1859            vec![vec![false]],
1860        );
1861    }
1862
1863    #[test]
1864    fn test_prefilter_builder_uses_row_group_proofs_by_index() {
1865        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
1866        let read_format = FlatReadFormat::new(
1867            metadata.clone(),
1868            ReadColumns::new(
1869                metadata
1870                    .column_metadatas
1871                    .iter()
1872                    .map(|column| column.column_id),
1873            ),
1874            None,
1875            "test",
1876            true,
1877        )
1878        .unwrap();
1879        let codec = build_primary_key_codec(metadata.as_ref());
1880        let builder = PrefilterContextBuilder::new(
1881            &read_format,
1882            &codec,
1883            None,
1884            None,
1885            new_simple_filter_contexts(&metadata, &[col("field_0").gt(lit(1_u64))]),
1886            Vec::new(),
1887            metadata.schema_version,
1888            &stats_metadata(&[vec![2], vec![1]], None),
1889        )
1890        .unwrap();
1891
1892        assert!(builder.build(0).proven_simple_filters[0]);
1893        assert!(!builder.build(1).proven_simple_filters[0]);
1894    }
1895
1896    #[test]
1897    fn test_simple_filter_stats_prove_integer_predicates() {
1898        macro_rules! assert_all_operators {
1899            ($literal:expr, $base:expr, $below:expr, $above:expr) => {{
1900                for (expr, min, max) in [
1901                    (col("x").gt(lit($literal)), $above.clone(), $above.clone()),
1902                    (col("x").gt_eq(lit($literal)), $base.clone(), $above.clone()),
1903                    (col("x").lt(lit($literal)), $below.clone(), $below.clone()),
1904                    (col("x").lt_eq(lit($literal)), $below.clone(), $base.clone()),
1905                    (col("x").eq(lit($literal)), $base.clone(), $base.clone()),
1906                    (
1907                        col("x").not_eq(lit($literal)),
1908                        $below.clone(),
1909                        $below.clone(),
1910                    ),
1911                ] {
1912                    let filter = SimpleFilterEvaluator::try_new(&expr).unwrap();
1913                    assert!(simple_filter_is_true_by_values(
1914                        &filter,
1915                        &$base,
1916                        Some(min),
1917                        Some(max),
1918                        Some(Value::UInt64(0)),
1919                    ));
1920                }
1921            }};
1922        }
1923
1924        assert_all_operators!(5_i32, Value::Int32(5), Value::Int32(-1), Value::Int32(6));
1925        assert_all_operators!(5_u32, Value::UInt32(5), Value::UInt32(4), Value::UInt32(6));
1926        assert_all_operators!(5_i64, Value::Int64(5), Value::Int64(-1), Value::Int64(6));
1927        assert_all_operators!(
1928            i64::MAX as u64 + 2,
1929            Value::UInt64(i64::MAX as u64 + 2),
1930            Value::UInt64(i64::MAX as u64 + 1),
1931            Value::UInt64(i64::MAX as u64 + 3)
1932        );
1933
1934        // Boundary cases must remain unproven: the min/max interval still
1935        // admits a row that violates the predicate.
1936        for (expr, min, max) in [
1937            (col("x").lt(lit(5_i64)), 1, 5),
1938            (col("x").gt(lit(5_i64)), 5, 9),
1939            (col("x").lt_eq(lit(5_i64)), 1, 6),
1940            (col("x").gt_eq(lit(5_i64)), 4, 9),
1941            (col("x").eq(lit(5_i64)), 5, 6),
1942            (col("x").not_eq(lit(5_i64)), 5, 9),
1943            (col("x").not_eq(lit(5_i64)), 1, 5),
1944            (col("x").not_eq(lit(5_i64)), 1, 9),
1945        ] {
1946            let filter = SimpleFilterEvaluator::try_new(&expr).unwrap();
1947            assert!(
1948                !simple_filter_is_true_by_values(
1949                    &filter,
1950                    &Value::Int64(5),
1951                    Some(Value::Int64(min)),
1952                    Some(Value::Int64(max)),
1953                    Some(Value::UInt64(0)),
1954                ),
1955                "{expr:?} must not be proven by stats {min}..={max}",
1956            );
1957        }
1958    }
1959
1960    #[test]
1961    fn test_simple_filter_stats_retain_unknown_or_unsupported_values() {
1962        let filter = SimpleFilterEvaluator::try_new(&col("x").gt(lit(1_i32))).unwrap();
1963        let literal = Value::Int32(1);
1964        for (min, max, null_count) in [
1965            (
1966                Some(Value::Int32(2)),
1967                Some(Value::Int32(3)),
1968                Some(Value::UInt64(1)),
1969            ),
1970            (Some(Value::Int32(2)), Some(Value::Int32(3)), None),
1971            (
1972                Some(Value::Null),
1973                Some(Value::Int32(3)),
1974                Some(Value::UInt64(0)),
1975            ),
1976            (None, Some(Value::Int32(3)), Some(Value::UInt64(0))),
1977            (Some(Value::Int32(2)), None, Some(Value::UInt64(0))),
1978            (
1979                Some(Value::Int64(2)),
1980                Some(Value::Int64(3)),
1981                Some(Value::UInt64(0)),
1982            ),
1983            (
1984                Some(Value::Int32(3)),
1985                Some(Value::Int32(2)),
1986                Some(Value::UInt64(0)),
1987            ),
1988        ] {
1989            assert!(!simple_filter_is_true_by_values(
1990                &filter, &literal, min, max, null_count
1991            ));
1992        }
1993
1994        let float_filter = SimpleFilterEvaluator::try_new(&col("x").gt(lit(1.0_f64))).unwrap();
1995        assert!(!simple_filter_is_true_by_values(
1996            &float_filter,
1997            &Value::Float64(1.0.into()),
1998            Some(Value::Float64(2.0.into())),
1999            Some(Value::Float64(3.0.into())),
2000            Some(Value::UInt64(0)),
2001        ));
2002    }
2003
2004    #[test]
2005    fn test_prefilter_entries_keep_only_unproven_simple_filter_indices() {
2006        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2007        let filters = new_simple_filter_contexts(
2008            &metadata,
2009            &[col("field_0").gt(lit(1_u64)), col("field_0").lt(lit(9_u64))],
2010        );
2011        let context = PrefilterContext {
2012            pk_filter: None,
2013            filters,
2014            physical_filters: Vec::new(),
2015            schema_version: metadata.schema_version,
2016            pk_filter_expr_strs: None,
2017            arrow_schema: metadata.schema.arrow_schema().clone(),
2018            proven_simple_filters: vec![true, false],
2019        };
2020
2021        let entries = all_prefilter_entries(&context);
2022        assert!(matches!(
2023            entries.as_slice(),
2024            [PrefilterEntry {
2025                kind: PrefilterEntryKind::Simple(1),
2026                ..
2027            }]
2028        ));
2029    }
2030
2031    #[test]
2032    fn test_identity_row_selection_preserves_input() {
2033        let sparse = RowSelection::from(vec![
2034            RowSelector::skip(2),
2035            RowSelector::select(3),
2036            RowSelector::skip(1),
2037        ]);
2038        assert_eq!(identity_row_selection(&Some(sparse.clone()), 6), sparse);
2039
2040        let empty = RowSelection::from(vec![]);
2041        assert_eq!(identity_row_selection(&Some(empty.clone()), 6), empty);
2042        assert_eq!(
2043            identity_row_selection(&None, 6),
2044            RowSelection::from(vec![RowSelector::select(6)])
2045        );
2046    }
2047
2048    #[test]
2049    fn test_should_use_prefilter() {
2050        assert!(should_use_prefilter(1, 5, 6));
2051        assert!(!should_use_prefilter(1, 0, 1));
2052        assert!(!should_use_prefilter(1, 1, 2));
2053        assert!(!should_use_prefilter(4, 3, 7));
2054        assert!(should_use_prefilter(3, 3, 6));
2055    }
2056
2057    #[test]
2058    fn test_build_bulk_filter_plan_classifies_filters_across_read_paths() {
2059        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2060            PrimaryKeyEncoding::Sparse,
2061        ));
2062        let legacy_read_format = FlatReadFormat::new(
2063            metadata.clone(),
2064            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2065            None,
2066            "memtable",
2067            false,
2068        )
2069        .unwrap();
2070        assert!(!legacy_read_format.batch_has_raw_pk_columns());
2071
2072        let plan = build_bulk_filter_plan(
2073            &legacy_read_format,
2074            Some(&Predicate::new(vec![
2075                col("tag_0").eq(lit("a")),
2076                col("field_0").gt(lit(1_u64)),
2077            ])),
2078        );
2079        assert_eq!(
2080            plan.pk_filters.as_ref().map(|filters| filters.len()),
2081            Some(1)
2082        );
2083        assert_eq!(
2084            remaining_simple_filter_columns(&plan.remaining_simple_filters),
2085            vec!["field_0"]
2086        );
2087
2088        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2089        let raw_pk_read_format = FlatReadFormat::new(
2090            metadata.clone(),
2091            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2092            None,
2093            "memtable",
2094            true,
2095        )
2096        .unwrap();
2097        assert!(raw_pk_read_format.batch_has_raw_pk_columns());
2098
2099        let tag_only_plan = build_bulk_filter_plan(
2100            &raw_pk_read_format,
2101            Some(&Predicate::new(vec![col("tag_0").eq(lit("a"))])),
2102        );
2103        assert!(tag_only_plan.pk_filters.is_none());
2104        assert_eq!(
2105            remaining_simple_filter_columns(&tag_only_plan.remaining_simple_filters),
2106            vec!["tag_0"]
2107        );
2108
2109        let field_only_plan = build_bulk_filter_plan(
2110            &raw_pk_read_format,
2111            Some(&Predicate::new(vec![col("field_0").gt(lit(1_u64))])),
2112        );
2113        assert!(field_only_plan.pk_filters.is_none());
2114        assert_eq!(
2115            remaining_simple_filter_columns(&field_only_plan.remaining_simple_filters),
2116            vec!["field_0"]
2117        );
2118    }
2119
2120    #[test]
2121    fn test_build_reader_filter_plan_classifies_filters_for_prefilter_modes() {
2122        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2123        let full_read_format = FlatReadFormat::new(
2124            metadata.clone(),
2125            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2126            None,
2127            "test",
2128            true,
2129        )
2130        .unwrap();
2131        let codec = build_primary_key_codec(metadata.as_ref());
2132
2133        let skip_fields_plan = build_reader_filter_plan(
2134            Some(&Predicate::new(vec![
2135                col("tag_0").eq(lit("a")),
2136                col("field_0").gt(lit(1_u64)),
2137            ])),
2138            None,
2139            PreFilterMode::SkipFields,
2140            true,
2141            false,
2142            &full_read_format,
2143            &codec,
2144            &stats_metadata(&[vec![1]], None),
2145        );
2146        assert!(skip_fields_plan.prefilter_builder.is_some());
2147        assert_eq!(
2148            remaining_simple_filter_columns(&skip_fields_plan.remaining_simple_filters),
2149            vec!["field_0"]
2150        );
2151
2152        let postponed_time_plan = build_reader_filter_plan(
2153            Some(&Predicate::new(vec![
2154                col("tag_0").eq(lit("a")),
2155                col("field_0").gt(lit(1_u64)),
2156                col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(1), None))),
2157            ])),
2158            None,
2159            PreFilterMode::SkipFields,
2160            true,
2161            true,
2162            &full_read_format,
2163            &codec,
2164            &stats_metadata(&[vec![1]], None),
2165        );
2166        assert!(postponed_time_plan.prefilter_builder.is_some());
2167        assert_eq!(
2168            remaining_simple_filter_columns(&postponed_time_plan.remaining_simple_filters),
2169            vec!["field_0", "ts"]
2170        );
2171
2172        let postponed_time_only_plan = build_reader_filter_plan(
2173            Some(&Predicate::new(vec![col("ts").gt_eq(lit(
2174                ScalarValue::TimestampMillisecond(Some(1), None),
2175            ))])),
2176            None,
2177            PreFilterMode::All,
2178            true,
2179            true,
2180            &full_read_format,
2181            &codec,
2182            &stats_metadata(&[vec![1]], None),
2183        );
2184        assert!(postponed_time_only_plan.prefilter_builder.is_none());
2185        assert_eq!(
2186            remaining_simple_filter_columns(&postponed_time_only_plan.remaining_simple_filters),
2187            vec!["ts"]
2188        );
2189
2190        let metric_metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2191            PrimaryKeyEncoding::Sparse,
2192        ));
2193        let field_0 = metric_metadata.column_by_name("field_0").unwrap().column_id;
2194        let ts = metric_metadata.time_index_column().column_id;
2195        let projected_read_format = FlatReadFormat::new(
2196            metric_metadata.clone(),
2197            ReadColumns::new([field_0, ts]),
2198            None,
2199            "test",
2200            true,
2201        )
2202        .unwrap();
2203        let metric_codec = build_primary_key_codec(metric_metadata.as_ref());
2204        let pk_prefilter_plan = build_reader_filter_plan(
2205            Some(&Predicate::new(vec![col("tag_0").eq(lit("a"))])),
2206            None,
2207            PreFilterMode::All,
2208            true,
2209            false,
2210            &projected_read_format,
2211            &metric_codec,
2212            &stats_metadata(&[vec![1]], None),
2213        );
2214        assert!(pk_prefilter_plan.prefilter_builder.is_some());
2215        assert!(
2216            pk_prefilter_plan
2217                .prefilter_builder
2218                .as_ref()
2219                .unwrap()
2220                .build_primary_key_filter()
2221                .is_some()
2222        );
2223        assert!(pk_prefilter_plan.remaining_simple_filters.is_empty());
2224
2225        let disabled_plan = build_reader_filter_plan(
2226            Some(&Predicate::new(vec![
2227                col("tag_0").eq(lit("a")),
2228                col("field_0").gt(lit(1_u64)),
2229                col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(1), None))),
2230            ])),
2231            None,
2232            PreFilterMode::All,
2233            false,
2234            true,
2235            &projected_read_format,
2236            &metric_codec,
2237            &stats_metadata(&[vec![1]], None),
2238        );
2239        assert!(disabled_plan.prefilter_builder.is_none());
2240        assert_eq!(
2241            remaining_simple_filter_columns(&disabled_plan.remaining_simple_filters),
2242            vec!["tag_0", "field_0", "ts"]
2243        );
2244    }
2245
2246    #[test]
2247    fn test_pk_filter_expr_strings_are_stable_under_expr_order() {
2248        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata_with_encoding(
2249            PrimaryKeyEncoding::Sparse,
2250        ));
2251        let read_format = FlatReadFormat::new(
2252            metadata.clone(),
2253            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2254            None,
2255            "test",
2256            false,
2257        )
2258        .unwrap();
2259        let codec = build_primary_key_codec(metadata.as_ref());
2260
2261        let expr_a = col("tag_0").eq(lit("a"));
2262        let expr_b = col("tag_1").eq(lit("x"));
2263        let plan_ab = build_reader_filter_plan(
2264            Some(&Predicate::new(vec![expr_a.clone(), expr_b.clone()])),
2265            None,
2266            PreFilterMode::All,
2267            true,
2268            false,
2269            &read_format,
2270            &codec,
2271            &stats_metadata(&[vec![1]], None),
2272        );
2273        let plan_b_a = build_reader_filter_plan(
2274            Some(&Predicate::new(vec![expr_b, expr_a])),
2275            None,
2276            PreFilterMode::All,
2277            true,
2278            false,
2279            &read_format,
2280            &codec,
2281            &stats_metadata(&[vec![1]], None),
2282        );
2283
2284        let exprs_ab = plan_ab.prefilter_builder.unwrap().pk_filter_expr_strs;
2285        let exprs_b_a = plan_b_a.prefilter_builder.unwrap().pk_filter_expr_strs;
2286        assert!(exprs_ab.is_some());
2287        assert_eq!(exprs_ab, exprs_b_a);
2288    }
2289
2290    #[test]
2291    fn test_simple_and_physical_contexts_preserve_expr_strings() {
2292        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2293        let read_format = FlatReadFormat::new(
2294            metadata.clone(),
2295            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2296            None,
2297            "test",
2298            true,
2299        )
2300        .unwrap();
2301
2302        let simple_expr = col("tag_0").eq(lit("a"));
2303        let simple = SimpleFilterContext::new_opt(&metadata, None, &simple_expr).unwrap();
2304        assert_eq!(simple.expr_str(), format!("{simple_expr:?}"));
2305
2306        let physical_expr = col("field_0").in_list(vec![lit(1_u64), lit(2_u64)], false);
2307        let physical =
2308            PhysicalFilterContext::new_opt(&metadata, None, &read_format, &physical_expr).unwrap();
2309        assert_eq!(physical.expr_str(), format!("{physical_expr:?}"));
2310    }
2311
2312    #[test]
2313    fn test_eval_simple_filter_mask_uses_flat_tag_columns_directly() {
2314        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2315        let filters = new_simple_filter_contexts(&metadata, &[col("tag_0").eq(lit("a"))]);
2316        let batch = new_record_batch_with_custom_sequence(&["a", "x"], 0, 4, 1);
2317
2318        let mask = eval_simple_filter_mask(&batch, &filters[0], "test").unwrap();
2319        assert_eq!(mask.count_set_bits(), 4);
2320    }
2321
2322    #[test]
2323    fn test_eval_simple_filter_mask_errors_on_missing_selected_column() {
2324        let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
2325        let filters = new_simple_filter_contexts(&metadata, &[col("tag_0").eq(lit("a"))]);
2326        let pk = new_primary_key(&["a", "x"]);
2327        let batch = new_raw_batch(&[pk.as_slice()], &[10]);
2328
2329        let err = eval_simple_filter_mask(&batch, &filters[0], "test").unwrap_err();
2330        let err = err.to_string();
2331        assert!(err.contains("Prefilter column"));
2332        assert!(err.contains("tag_0"));
2333    }
2334
2335    #[test]
2336    fn test_eval_physical_filter_mask_evaluates_physical_filters() {
2337        let metadata: RegionMetadataRef =
2338            Arc::new(sst_region_metadata_with_encoding(PrimaryKeyEncoding::Dense));
2339        let read_format = FlatReadFormat::new(
2340            metadata.clone(),
2341            ReadColumns::new(metadata.column_metadatas.iter().map(|c| c.column_id)),
2342            None,
2343            "test",
2344            false,
2345        )
2346        .unwrap();
2347        let expr = col("field_0").in_list(vec![lit(11_u64)], false);
2348        let physical_filters = new_physical_filter_contexts(&metadata, &read_format, &[expr]);
2349        let pk = new_primary_key(&["a", "x"]);
2350        let batch = new_raw_batch(&[pk.as_slice(), pk.as_slice(), pk.as_slice()], &[9, 10, 11]);
2351
2352        let mask = eval_physical_filter_mask(&batch, &physical_filters[0], "test").unwrap();
2353        assert_eq!(mask.count_set_bits(), 1);
2354    }
2355
2356    #[test]
2357    fn test_eval_pk_group_mask_finds_pk_column_by_name() {
2358        let metadata = Arc::new(sst_region_metadata());
2359        let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("a"))]));
2360        let mut pk_filter = Some(Box::new(CachedPrimaryKeyFilter::new(
2361            build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters),
2362        )) as Box<dyn PrimaryKeyFilter>);
2363        let pk_a = new_primary_key(&["a", "x"]);
2364        let pk_b = new_primary_key(&["b", "x"]);
2365        let batch = new_prefilter_batch(
2366            &[
2367                pk_a.as_slice(),
2368                pk_a.as_slice(),
2369                pk_b.as_slice(),
2370                pk_b.as_slice(),
2371            ],
2372            &[10, 11, 12, 13],
2373        );
2374
2375        let mask = primary_key_filter_mask(&batch, pk_filter.as_mut().unwrap().as_mut()).unwrap();
2376
2377        assert_eq!(mask.count_set_bits(), 2);
2378    }
2379
2380    #[test]
2381    fn test_eval_pk_group_mask_handles_binary_pk_column() {
2382        let metadata = Arc::new(sst_region_metadata());
2383        let filters = Arc::new(new_test_filters(&[col("tag_0").eq(lit("a"))]));
2384        let mut pk_filter = Some(Box::new(CachedPrimaryKeyFilter::new(
2385            build_primary_key_codec(metadata.as_ref()).primary_key_filter(&metadata, filters),
2386        )) as Box<dyn PrimaryKeyFilter>);
2387        let pk_a = new_primary_key(&["a", "x"]);
2388        let pk_b = new_primary_key(&["b", "x"]);
2389        let batch = new_prefilter_batch_binary_pk(
2390            &[
2391                pk_a.as_slice(),
2392                pk_a.as_slice(),
2393                pk_b.as_slice(),
2394                pk_b.as_slice(),
2395            ],
2396            &[10, 11, 12, 13],
2397        );
2398
2399        let mask = primary_key_filter_mask(&batch, pk_filter.as_mut().unwrap().as_mut()).unwrap();
2400
2401        assert_eq!(mask.count_set_bits(), 2);
2402    }
2403}