Skip to main content

mito2/read/
scan_region.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//! Scans a region according to the scan request.
16
17use std::collections::{BTreeMap, HashSet};
18use std::fmt;
19use std::num::NonZeroU64;
20use std::sync::{Arc, OnceLock};
21use std::time::Instant;
22
23use api::v1::SemanticType;
24use common_error::ext::BoxedError;
25use common_recordbatch::SendableRecordBatchStream;
26use common_recordbatch::adapter::RegionQueryStatCounters;
27use common_recordbatch::filter::SimpleFilterEvaluator;
28use common_telemetry::tracing::Instrument;
29use common_telemetry::{debug, error, tracing, warn};
30use common_time::range::TimestampRange;
31use datafusion::execution::memory_pool::{MemoryPool, UnboundedMemoryPool};
32use datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr;
33use datafusion_common::pruning::PruningStatistics;
34use datafusion_common::{Column, ScalarValue};
35use datafusion_expr::Expr;
36use datafusion_expr::utils::expr_to_columns;
37use datatypes::arrow::array::{ArrayRef, BooleanArray, UInt64Array};
38use datatypes::extension::json::is_json2_extension_type;
39use datatypes::prelude::ConcreteDataType;
40use datatypes::types::json_type::JsonNativeType;
41use datatypes::value::timestamp_to_scalar_value;
42use futures::StreamExt;
43use partition::expr::PartitionExpr;
44use smallvec::SmallVec;
45use snafu::{OptionExt, ResultExt, ensure};
46use store_api::metadata::{RegionMetadata, RegionMetadataRef};
47use store_api::region_engine::{PartitionRange, RegionScannerRef};
48use store_api::storage::{
49    ColumnId, RegionId, ScanRequest, SequenceNumber, SequenceRange, TimeSeriesDistribution,
50    TimeSeriesRowSelector,
51};
52use table::predicate::{Predicate, build_time_range_predicate, extract_time_range_from_expr};
53use tokio::sync::{Semaphore, mpsc};
54use tokio_stream::wrappers::ReceiverStream;
55
56use crate::access_layer::AccessLayerRef;
57use crate::cache::CacheStrategy;
58use crate::config::DEFAULT_MAX_CONCURRENT_SCAN_FILES;
59use crate::error::{
60    InvalidPartitionExprSnafu, InvalidRequestSnafu, RegionSequenceDomainBrokenSnafu, Result,
61    SequenceRangeUnsupportedSnafu,
62};
63#[cfg(feature = "enterprise")]
64use crate::extension::{BoxedExtensionRange, BoxedExtensionRangeProvider};
65use crate::memtable::{MemtableRange, RangesOptions};
66use crate::metrics::READ_SST_COUNT;
67use crate::read::compat::{self, FlatCompatBatch};
68use crate::read::flat_projection::FlatProjectionMapper;
69use crate::read::range::{FileRangeBuilder, MemRangeBuilder, RangeMeta, RowGroupIndex};
70use crate::read::range_cache::{ScanRequestFingerprint, implied_time_range_from_exprs};
71use crate::read::read_columns::ReadColumns;
72use crate::read::seq_scan::SeqScan;
73use crate::read::series_scan::SeriesScan;
74use crate::read::stream::ScanBatchStream;
75use crate::read::unordered_scan::UnorderedScan;
76use crate::read::{BoxedRecordBatchStream, RecordBatch};
77use crate::region::options::MergeMode;
78use crate::region::version::VersionRef;
79use crate::series_index::SeriesIndexReadContext;
80use crate::sst::file::FileHandle;
81use crate::sst::index::bloom_filter::applier::{
82    BloomFilterIndexApplierBuilder, BloomFilterIndexApplierRef,
83};
84use crate::sst::index::fulltext_index::applier::FulltextIndexApplierRef;
85use crate::sst::index::fulltext_index::applier::builder::FulltextIndexApplierBuilder;
86use crate::sst::index::inverted_index::applier::InvertedIndexApplierRef;
87use crate::sst::index::inverted_index::applier::builder::InvertedIndexApplierBuilder;
88#[cfg(feature = "vector_index")]
89use crate::sst::index::vector_index::applier::{VectorIndexApplier, VectorIndexApplierRef};
90use crate::sst::parquet::Json2RewriteTargets;
91use crate::sst::parquet::file_range::PreFilterMode;
92use crate::sst::parquet::reader::ReaderMetrics;
93use crate::sst::primary_key::PrimaryKeyRangeMapper;
94
95#[cfg(feature = "vector_index")]
96const VECTOR_INDEX_OVERFETCH_MULTIPLIER: usize = 2;
97
98/// A scanner scans a region and returns a [SendableRecordBatchStream].
99pub(crate) enum Scanner {
100    /// Sequential scan.
101    Seq(SeqScan),
102    /// Unordered scan.
103    Unordered(UnorderedScan),
104    /// Per-series scan.
105    Series(SeriesScan),
106}
107
108impl Scanner {
109    /// Returns a [SendableRecordBatchStream] to retrieve scan results from all partitions.
110    #[tracing::instrument(level = tracing::Level::DEBUG, skip_all)]
111    pub(crate) async fn scan(&self) -> Result<SendableRecordBatchStream, BoxedError> {
112        match self {
113            Scanner::Seq(seq_scan) => seq_scan.build_stream(),
114            Scanner::Unordered(unordered_scan) => unordered_scan.build_stream().await,
115            Scanner::Series(series_scan) => series_scan.build_stream().await,
116        }
117    }
118
119    /// Create a stream of [`Batch`] by this scanner.
120    pub(crate) fn scan_batch(&self) -> Result<ScanBatchStream> {
121        match self {
122            Scanner::Seq(x) => x.scan_all_partitions(),
123            Scanner::Unordered(x) => x.scan_all_partitions(),
124            Scanner::Series(x) => x.scan_all_partitions(),
125        }
126    }
127}
128
129#[cfg(test)]
130impl Scanner {
131    /// Returns number of files to scan.
132    pub(crate) fn num_files(&self) -> usize {
133        match self {
134            Scanner::Seq(seq_scan) => seq_scan.input().num_files(),
135            Scanner::Unordered(unordered_scan) => unordered_scan.input().num_files(),
136            Scanner::Series(series_scan) => series_scan.input().num_files(),
137        }
138    }
139
140    /// Returns number of memtables to scan.
141    pub(crate) fn num_memtables(&self) -> usize {
142        match self {
143            Scanner::Seq(seq_scan) => seq_scan.input().num_memtables(),
144            Scanner::Unordered(unordered_scan) => unordered_scan.input().num_memtables(),
145            Scanner::Series(series_scan) => series_scan.input().num_memtables(),
146        }
147    }
148
149    /// Returns SST file ids to scan.
150    pub(crate) fn file_ids(&self) -> Vec<crate::sst::file::RegionFileId> {
151        match self {
152            Scanner::Seq(seq_scan) => seq_scan.input().file_ids(),
153            Scanner::Unordered(unordered_scan) => unordered_scan.input().file_ids(),
154            Scanner::Series(series_scan) => series_scan.input().file_ids(),
155        }
156    }
157
158    pub(crate) fn index_ids(&self) -> Vec<crate::sst::file::RegionIndexId> {
159        match self {
160            Scanner::Seq(seq_scan) => seq_scan.input().index_ids(),
161            Scanner::Unordered(unordered_scan) => unordered_scan.input().index_ids(),
162            Scanner::Series(series_scan) => series_scan.input().index_ids(),
163        }
164    }
165
166    pub(crate) fn snapshot_sequence(&self) -> Option<SequenceNumber> {
167        match self {
168            Scanner::Seq(seq_scan) => seq_scan.input().snapshot_sequence,
169            Scanner::Unordered(unordered_scan) => unordered_scan.input().snapshot_sequence,
170            Scanner::Series(series_scan) => series_scan.input().snapshot_sequence,
171        }
172    }
173
174    /// Sets the target partitions for the scanner. It can controls the parallelism of the scanner.
175    pub(crate) fn set_target_partitions(&mut self, target_partitions: usize) {
176        use store_api::region_engine::{PrepareRequest, RegionScanner};
177
178        let request = PrepareRequest::default().with_target_partitions(target_partitions);
179        match self {
180            Scanner::Seq(seq_scan) => seq_scan.prepare(request).unwrap(),
181            Scanner::Unordered(unordered_scan) => unordered_scan.prepare(request).unwrap(),
182            Scanner::Series(series_scan) => series_scan.prepare(request).unwrap(),
183        }
184    }
185}
186
187#[cfg_attr(doc, aquamarine::aquamarine)]
188/// Helper to scans a region by [ScanRequest].
189///
190/// [ScanRegion] collects SSTs and memtables to scan without actually reading them. It
191/// creates a [Scanner] to actually scan these targets in [Scanner::scan()].
192///
193/// ```mermaid
194/// classDiagram
195/// class ScanRegion {
196///     -VersionRef version
197///     -ScanRequest request
198///     ~scanner() Scanner
199///     ~seq_scan() SeqScan
200/// }
201/// class Scanner {
202///     <<enumeration>>
203///     SeqScan
204///     UnorderedScan
205///     +scan() SendableRecordBatchStream
206/// }
207/// class SeqScan {
208///     -ScanInput input
209///     +build() SendableRecordBatchStream
210/// }
211/// class UnorderedScan {
212///     -ScanInput input
213///     +build() SendableRecordBatchStream
214/// }
215/// class ScanInput {
216///     -ProjectionMapper mapper
217///     -Option~TimeRange~ time_range
218///     -Option~Predicate~ predicate
219///     -Vec~MemtableRef~ memtables
220///     -Vec~FileHandle~ files
221/// }
222/// class ProjectionMapper {
223///     ~output_schema() SchemaRef
224///     ~convert(Batch) RecordBatch
225/// }
226/// ScanRegion -- Scanner
227/// ScanRegion o-- ScanRequest
228/// Scanner o-- SeqScan
229/// Scanner o-- UnorderedScan
230/// SeqScan o-- ScanInput
231/// UnorderedScan o-- ScanInput
232/// Scanner -- SendableRecordBatchStream
233/// ScanInput o-- ProjectionMapper
234/// SeqScan -- SendableRecordBatchStream
235/// UnorderedScan -- SendableRecordBatchStream
236/// ```
237pub(crate) struct ScanRegion {
238    /// Version of the region at scan.
239    version: VersionRef,
240    /// Pinned index snapshot and its storage for candidate discovery.
241    series_index: Option<SeriesIndexReadContext>,
242    /// Access layer of the region.
243    access_layer: AccessLayerRef,
244    /// Scan request.
245    request: ScanRequest,
246    /// Cache.
247    cache_strategy: CacheStrategy,
248    /// Maximum number of SST files to scan concurrently.
249    max_concurrent_scan_files: usize,
250    /// Memory pool shared by internal scan operators across all queries.
251    scan_memory_pool: Arc<dyn MemoryPool>,
252    /// Whether to enable the experimental two-phase metric series scan.
253    experimental_series_scan_v2: bool,
254    /// Whether to ignore range indexes during scans.
255    ignore_range_index: bool,
256    /// Whether to ignore inverted index.
257    ignore_inverted_index: bool,
258    /// Whether to ignore fulltext index.
259    ignore_fulltext_index: bool,
260    /// Whether to ignore bloom filter.
261    ignore_bloom_filter: bool,
262    /// Start time of the scan task.
263    start_time: Option<Instant>,
264    /// Whether to filter out the deleted rows.
265    /// Usually true for normal read, and false for scan for compaction.
266    filter_deleted: bool,
267    /// Files and exact sequence capability selected together by the engine.
268    exact_selection: Option<(Vec<FileHandle>, Option<SequenceRange>)>,
269    /// Counters that should receive query-load metrics.
270    query_stat_counters: Option<RegionQueryStatCounters>,
271    #[cfg(feature = "enterprise")]
272    extension_range_provider: Option<BoxedExtensionRangeProvider>,
273}
274
275impl ScanRegion {
276    /// Creates a [ScanRegion].
277    pub(crate) fn new(
278        version: VersionRef,
279        access_layer: AccessLayerRef,
280        request: ScanRequest,
281        cache_strategy: CacheStrategy,
282    ) -> ScanRegion {
283        ScanRegion {
284            version,
285            series_index: None,
286            access_layer,
287            request,
288            cache_strategy,
289            max_concurrent_scan_files: DEFAULT_MAX_CONCURRENT_SCAN_FILES,
290            scan_memory_pool: Arc::new(UnboundedMemoryPool::default()),
291            experimental_series_scan_v2: false,
292            ignore_range_index: false,
293            ignore_inverted_index: false,
294            ignore_fulltext_index: false,
295            ignore_bloom_filter: false,
296            start_time: None,
297            filter_deleted: true,
298            exact_selection: None,
299            query_stat_counters: None,
300            #[cfg(feature = "enterprise")]
301            extension_range_provider: None,
302        }
303    }
304
305    /// Pins the series-index snapshot used by candidate discovery.
306    #[must_use]
307    pub(crate) fn with_series_index(
308        mut self,
309        series_index: Option<SeriesIndexReadContext>,
310    ) -> Self {
311        self.series_index = series_index;
312        self
313    }
314
315    /// Sets counters that should receive query-load metrics.
316    #[must_use]
317    pub(crate) fn with_query_stat_counters(mut self, counters: RegionQueryStatCounters) -> Self {
318        self.query_stat_counters = Some(counters);
319        self
320    }
321
322    /// Sets maximum number of SST files to scan concurrently.
323    #[must_use]
324    pub(crate) fn with_max_concurrent_scan_files(
325        mut self,
326        max_concurrent_scan_files: usize,
327    ) -> Self {
328        self.max_concurrent_scan_files = max_concurrent_scan_files;
329        self
330    }
331
332    /// Sets the memory pool shared by internal scan operators.
333    #[must_use]
334    pub(crate) fn with_scan_memory_pool(mut self, scan_memory_pool: Arc<dyn MemoryPool>) -> Self {
335        self.scan_memory_pool = scan_memory_pool;
336        self
337    }
338
339    /// Sets whether eligible metric series scans use the two-phase implementation.
340    #[must_use]
341    pub(crate) fn with_experimental_series_scan_v2(mut self, enabled: bool) -> Self {
342        self.experimental_series_scan_v2 = enabled;
343        self
344    }
345
346    /// Sets whether to ignore range indexes during scans.
347    #[must_use]
348    pub(crate) fn with_ignore_range_index(mut self, ignore: bool) -> Self {
349        self.ignore_range_index = ignore;
350        self
351    }
352
353    /// Sets whether to ignore inverted index.
354    #[must_use]
355    pub(crate) fn with_ignore_inverted_index(mut self, ignore: bool) -> Self {
356        self.ignore_inverted_index = ignore;
357        self
358    }
359
360    /// Sets whether to ignore fulltext index.
361    #[must_use]
362    pub(crate) fn with_ignore_fulltext_index(mut self, ignore: bool) -> Self {
363        self.ignore_fulltext_index = ignore;
364        self
365    }
366
367    /// Sets whether to ignore bloom filter.
368    #[must_use]
369    pub(crate) fn with_ignore_bloom_filter(mut self, ignore: bool) -> Self {
370        self.ignore_bloom_filter = ignore;
371        self
372    }
373
374    #[must_use]
375    pub(crate) fn with_start_time(mut self, now: Instant) -> Self {
376        self.start_time = Some(now);
377        self
378    }
379
380    pub(crate) fn set_filter_deleted(&mut self, filter_deleted: bool) {
381        self.filter_deleted = filter_deleted;
382    }
383
384    pub(crate) fn with_exact_selection(
385        mut self,
386        exact_selection: (Vec<FileHandle>, Option<SequenceRange>),
387    ) -> Self {
388        self.exact_selection = Some(exact_selection);
389        self
390    }
391
392    #[cfg(feature = "enterprise")]
393    pub(crate) fn set_extension_range_provider(
394        &mut self,
395        extension_range_provider: BoxedExtensionRangeProvider,
396    ) {
397        self.extension_range_provider = Some(extension_range_provider);
398    }
399
400    /// Returns a [Scanner] to scan the region.
401    #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
402    pub(crate) async fn scanner(self) -> Result<Scanner> {
403        if self.use_series_scan() {
404            self.series_scan().await.map(Scanner::Series)
405        } else if self.use_unordered_scan() {
406            // If table is append only and there is no series row selector, we use unordered scan in query.
407            // We still use seq scan in compaction.
408            self.unordered_scan().await.map(Scanner::Unordered)
409        } else {
410            self.seq_scan().await.map(Scanner::Seq)
411        }
412    }
413
414    /// Returns a [RegionScanner] to scan the region.
415    #[tracing::instrument(
416        level = tracing::Level::DEBUG,
417        skip_all,
418        fields(region_id = %self.region_id())
419    )]
420    pub(crate) async fn region_scanner(self) -> Result<RegionScannerRef> {
421        if self.use_series_scan() {
422            self.series_scan()
423                .await
424                .map(|scanner| Box::new(scanner) as _)
425        } else if self.use_unordered_scan() {
426            self.unordered_scan()
427                .await
428                .map(|scanner| Box::new(scanner) as _)
429        } else {
430            self.seq_scan().await.map(|scanner| Box::new(scanner) as _)
431        }
432    }
433
434    /// Scan sequentially.
435    #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
436    pub(crate) async fn seq_scan(self) -> Result<SeqScan> {
437        let input = self.scan_input().await?;
438        Ok(SeqScan::new(input))
439    }
440
441    /// Unordered scan.
442    #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
443    pub(crate) async fn unordered_scan(self) -> Result<UnorderedScan> {
444        let input = self.scan_input().await?;
445        Ok(UnorderedScan::new(input))
446    }
447
448    /// Scans by series.
449    #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
450    pub(crate) async fn series_scan(self) -> Result<SeriesScan> {
451        let experimental_series_scan_v2 = self.experimental_series_scan_v2;
452        let input = self.scan_input().await?;
453        Ok(SeriesScan::new(input, experimental_series_scan_v2))
454    }
455
456    /// Returns true if the region can use unordered scan for current request.
457    fn use_unordered_scan(&self) -> bool {
458        // We use unordered scan when:
459        // 1. The region is in append mode.
460        // 2. There is no series row selector.
461        // 3. The required distribution is None or TimeSeriesDistribution::TimeWindowed.
462        //
463        // We still use seq scan in compaction.
464        self.version.options.append_mode
465            && self.request.series_row_selector.is_none()
466            && (self.request.distribution.is_none()
467                || self.request.distribution == Some(TimeSeriesDistribution::TimeWindowed))
468    }
469
470    /// Returns true if the region can use series scan for current request.
471    fn use_series_scan(&self) -> bool {
472        self.request.distribution == Some(TimeSeriesDistribution::PerSeries)
473    }
474
475    /// Creates a scan input.
476    #[tracing::instrument(skip_all, fields(region_id = %self.region_id()))]
477    async fn scan_input(mut self) -> Result<ScanInput> {
478        let metadata = &self.version.metadata;
479        let sst_min_sequence = self.request.sst_min_sequence.and_then(NonZeroU64::new);
480        let time_range = self.build_time_range_predicate();
481        let predicate = PredicateGroup::new(metadata, &self.request.filters)?;
482
483        let read_col_ids =
484            self.build_read_col_ids(self.request.projection.as_deref(), &predicate)?;
485        let read_cols = self.build_read_columns(&read_col_ids)?;
486
487        // The mapper always computes projected column ids as the schema of SSTs may change.
488        let projection = self
489            .request
490            .projection
491            .clone()
492            .unwrap_or_else(|| (0..metadata.column_metadatas.len()).collect());
493        let mapper = FlatProjectionMapper::new_with_read_columns(metadata, projection, read_cols)?;
494        let mapper = if self.request.preserve_pk_dictionary_encoding {
495            mapper.with_pk_dictionary_encoding()
496        } else {
497            mapper
498        };
499
500        let (files, sequence_range) =
501            if let Some((files, sequence_range)) = self.exact_selection.take() {
502                (files, sequence_range)
503            } else {
504                exact_sequence_range(&self.request, &self.version)?
505            };
506        if sst_min_sequence.is_some() && sequence_range.is_some() {
507            return SequenceRangeUnsupportedSnafu {
508                region_id: self.region_id(),
509                min_seq: self.request.memtable_min_sequence.unwrap_or_default(),
510                max_seq: self.request.memtable_max_sequence.unwrap_or_default(),
511                reason:
512                    "sst_min_sequence pruning hint is incompatible with exact sequence-range reads"
513                        .to_string(),
514            }
515            .fail();
516        }
517        let memtables = self.version.memtables.list_memtables();
518        // Skip empty memtables and memtables out of time range.
519        let mut mem_range_builders = Vec::new();
520        let filter_mode = pre_filter_mode(
521            self.version.options.append_mode,
522            self.version.options.merge_mode(),
523        );
524
525        for m in memtables {
526            // check if memtable is empty by reading stats.
527            let Some((start, end)) = m.stats().time_range() else {
528                continue;
529            };
530            // The time range of the memtable is inclusive.
531            let memtable_range = TimestampRange::new_inclusive(Some(start), Some(end));
532            if !memtable_range.intersects(&time_range) {
533                continue;
534            }
535            let ranges_in_memtable = m.ranges(
536                Some(&read_col_ids),
537                RangesOptions::default()
538                    .with_predicate(predicate.clone())
539                    .with_sequence(SequenceRange::new(
540                        self.request.memtable_min_sequence,
541                        self.request.memtable_max_sequence,
542                    ))
543                    .with_pre_filter_mode(filter_mode),
544            )?;
545            mem_range_builders.extend(ranges_in_memtable.ranges.into_values().map(|v| {
546                let stats = v.stats().clone();
547                MemRangeBuilder::new(v, stats)
548            }));
549        }
550
551        let region_id = self.region_id();
552        debug!(
553            "Scan region {}, request: {:?}, time range: {:?}, memtables: {}, ssts_to_read: {}, append_mode: {}",
554            region_id,
555            self.request,
556            time_range,
557            mem_range_builders.len(),
558            files.len(),
559            self.version.options.append_mode,
560        );
561
562        let (non_field_filters, field_filters) = self.partition_by_field_filters();
563        let inverted_index_appliers = [
564            self.build_invereted_index_applier(&non_field_filters),
565            self.build_invereted_index_applier(&field_filters),
566        ];
567        let bloom_filter_appliers = [
568            self.build_bloom_filter_applier(&non_field_filters),
569            self.build_bloom_filter_applier(&field_filters),
570        ];
571        let fulltext_index_appliers = [
572            self.build_fulltext_index_applier(&non_field_filters),
573            self.build_fulltext_index_applier(&field_filters),
574        ];
575        #[cfg(feature = "vector_index")]
576        let vector_index_applier = self.build_vector_index_applier();
577        #[cfg(feature = "vector_index")]
578        let vector_index_k = self.request.vector_search.as_ref().map(|search| {
579            if self.request.filters.is_empty() {
580                search.k
581            } else {
582                search.k.saturating_mul(VECTOR_INDEX_OVERFETCH_MULTIPLIER)
583            }
584        });
585
586        let input = ScanInput::builder(self.access_layer, mapper)
587            .with_series_index(self.series_index)
588            .with_ignore_range_index(self.ignore_range_index)
589            .with_time_range(Some(time_range))
590            .with_predicate(predicate)
591            .with_memtables(mem_range_builders)
592            .with_files(files)
593            .with_primary_key_mapper(self.version.ssts.primary_key_mapper())
594            .with_cache(self.cache_strategy)
595            .with_inverted_index_appliers(inverted_index_appliers)
596            .with_bloom_filter_index_appliers(bloom_filter_appliers)
597            .with_fulltext_index_appliers(fulltext_index_appliers)
598            .with_max_concurrent_scan_files(self.max_concurrent_scan_files)
599            .with_scan_memory_pool(self.scan_memory_pool)
600            .with_start_time(self.start_time)
601            .with_append_mode(self.version.options.append_mode)
602            .with_filter_deleted(self.filter_deleted)
603            .with_merge_mode(self.version.options.merge_mode())
604            .with_series_row_selector(self.request.series_row_selector)
605            .with_distribution(self.request.distribution)
606            .with_explain_flat_format(
607                self.version.options.sst_format == Some(crate::sst::FormatType::Flat),
608            )
609            .with_snapshot_sequence(
610                self.request
611                    .snapshot_on_scan
612                    .then_some(self.request.memtable_max_sequence)
613                    .flatten(),
614            )
615            .with_sequence_range(sequence_range)
616            .with_query_stat_counters(self.query_stat_counters);
617        #[cfg(feature = "vector_index")]
618        let input = input
619            .with_vector_index_applier(vector_index_applier)
620            .with_vector_index_k(vector_index_k);
621
622        #[cfg(feature = "enterprise")]
623        let input = if !self.request.skip_sst_files
624            && let Some(provider) = self.extension_range_provider
625        {
626            if sequence_range.is_some() {
627                // Defense in depth: the engine already rejects exact
628                // sequence-range reads on follower regions with an extension
629                // provider, but if one ever reaches the reader, fail closed
630                // here rather than letting unfiltered extension streams bypass
631                // the row-level sequence filter and emit out-of-range rows.
632                return SequenceRangeUnsupportedSnafu {
633                    region_id,
634                    min_seq: self.request.memtable_min_sequence.unwrap_or_default(),
635                    max_seq: self.request.memtable_max_sequence.unwrap_or_default(),
636                    reason:
637                        "exact sequence-range reads are unsupported when an extension range provider is present"
638                            .to_string(),
639                }
640                .fail();
641            }
642            let ranges = provider
643                .find_extension_ranges(self.version.flushed_sequence, time_range, &self.request)
644                .await?;
645            debug!("Find extension ranges: {ranges:?}");
646            input.with_extension_ranges(ranges)
647        } else {
648            input
649        };
650        Ok(input.build())
651    }
652
653    /// Builds the deduplicated root column ids required by the projection and
654    /// predicate.
655    fn build_read_col_ids(
656        &self,
657        projection: Option<&[usize]>,
658        predicate: &PredicateGroup,
659    ) -> Result<Vec<ColumnId>> {
660        let metadata = &self.version.metadata;
661        let Some(projection) = projection else {
662            return Ok(metadata
663                .column_metadatas
664                .iter()
665                .map(|col| col.column_id)
666                .collect());
667        };
668
669        let mut read_col_ids = Vec::new();
670        let mut seen = HashSet::new();
671        for idx in projection {
672            let col_id = metadata
673                .column_metadatas
674                .get(*idx)
675                .with_context(|| InvalidRequestSnafu {
676                    region_id: metadata.region_id,
677                    reason: format!("projection index {} is out of bounds", idx),
678                })?
679                .column_id;
680            let inserted = seen.insert(col_id);
681            // DataFusion's logical `OptimizeProjections` rule deduplicates pushed-down
682            // projection indices via `RequiredIndices::compact()`.
683            debug_assert!(
684                inserted,
685                "projection contains duplicate column id: {}",
686                col_id
687            );
688            // Keep the projection order.
689            read_col_ids.push(col_id);
690        }
691
692        if projection.is_empty() {
693            let time_index = metadata.time_index_column().column_id;
694            if seen.insert(time_index) {
695                read_col_ids.push(time_index);
696            }
697        }
698
699        let mut extra_col_names = HashSet::new();
700        let mut cols = HashSet::new();
701
702        if let Some(p) = predicate.predicate_without_region() {
703            for expr in p.exprs() {
704                cols.clear();
705                if expr_to_columns(expr, &mut cols).is_err() {
706                    continue;
707                }
708                extra_col_names.extend(cols.iter().map(|col| col.name.clone()));
709            }
710        }
711
712        if let Some(expr) = predicate.region_partition_expr() {
713            expr.collect_column_names(&mut extra_col_names);
714        }
715
716        if !extra_col_names.is_empty() {
717            for col in &metadata.column_metadatas {
718                if extra_col_names.remove(&col.column_schema.name) && !seen.contains(&col.column_id)
719                {
720                    read_col_ids.push(col.column_id);
721                }
722            }
723            if !extra_col_names.is_empty() {
724                warn!(
725                    "Some columns in filters are not found in region {}: {:?}",
726                    metadata.region_id, extra_col_names
727                );
728            }
729        }
730        Ok(read_col_ids)
731    }
732
733    /// Builds logical read columns and attaches JSON2 target types when needed.
734    ///
735    /// The behavior about JSON2 is as follows:
736    /// - JSON2 columns without hints use Variant to read the whole column.
737    /// - Hints targeting non-JSON2 read columns are rejected.
738    fn build_read_columns(&self, col_ids: &[ColumnId]) -> Result<ReadColumns> {
739        let metadata = &self.version.metadata;
740        let json_type_hint = &self.request.json_type_hint;
741
742        let has_json2 = metadata
743            .schema
744            .arrow_schema()
745            .fields()
746            .iter()
747            .any(is_json2_extension_type);
748
749        if !has_json2 && json_type_hint.is_empty() {
750            return Ok(ReadColumns::new(col_ids.iter().copied()));
751        }
752
753        let mut json_target_types = BTreeMap::new();
754        for &col_id in col_ids {
755            let Some(col) = metadata.column_by_id(col_id) else {
756                continue;
757            };
758            let col_name = &col.column_schema.name;
759            let hint = json_type_hint.get(col_name);
760            if !col.column_schema.data_type.is_json2() {
761                ensure!(
762                    hint.is_none(),
763                    InvalidRequestSnafu {
764                        region_id: metadata.region_id,
765                        reason: format!(
766                            "JSON type hint targets non-JSON2 column {} (id: {}, type: {})",
767                            col_name, col_id, col.column_schema.data_type
768                        ),
769                    }
770                );
771                continue;
772            }
773            let target_type = hint.cloned().unwrap_or(JsonNativeType::Variant);
774            json_target_types.insert(col_id, target_type);
775        }
776        Ok(ReadColumns::new(col_ids.iter().copied()).with_json_target_types(json_target_types))
777    }
778
779    fn region_id(&self) -> RegionId {
780        self.version.metadata.region_id
781    }
782
783    /// Build time range predicate from filters.
784    fn build_time_range_predicate(&self) -> TimestampRange {
785        let time_index = self.version.metadata.time_index_column();
786        let unit = time_index
787            .column_schema
788            .data_type
789            .as_timestamp()
790            .expect("Time index must have timestamp-compatible type")
791            .unit();
792        build_time_range_predicate(&time_index.column_schema.name, unit, &self.request.filters)
793    }
794
795    /// Partitions filters into two groups: non-field filters and field filters.
796    /// Returns `(non_field_filters, field_filters)`.
797    fn partition_by_field_filters(&self) -> (Vec<Expr>, Vec<Expr>) {
798        let field_columns = self
799            .version
800            .metadata
801            .field_columns()
802            .map(|col| &col.column_schema.name)
803            .collect::<HashSet<_>>();
804
805        let mut columns = HashSet::new();
806
807        self.request.filters.iter().cloned().partition(|expr| {
808            columns.clear();
809            // `expr_to_columns` won't return error.
810            if expr_to_columns(expr, &mut columns).is_err() {
811                // If we can't extract columns, treat it as non-field filter
812                return true;
813            }
814            // Return true for non-field filters (partition puts true cases in first vec)
815            !columns
816                .iter()
817                .any(|column| field_columns.contains(&column.name))
818        })
819    }
820
821    /// Use the latest schema to build the inverted index applier.
822    fn build_invereted_index_applier(&self, filters: &[Expr]) -> Option<InvertedIndexApplierRef> {
823        if self.ignore_inverted_index {
824            return None;
825        }
826
827        let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
828        let inverted_index_cache = self.cache_strategy.inverted_index_cache().cloned();
829
830        let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
831
832        InvertedIndexApplierBuilder::new(
833            self.access_layer.table_dir().to_string(),
834            self.access_layer.path_type(),
835            self.access_layer.object_store().clone(),
836            self.version.metadata.as_ref(),
837            self.version.metadata.inverted_indexed_column_ids(
838                self.version
839                    .options
840                    .index_options
841                    .inverted_index
842                    .ignore_column_ids
843                    .iter(),
844            ),
845            self.access_layer.puffin_manager_factory().clone(),
846        )
847        .with_file_cache(file_cache)
848        .with_inverted_index_cache(inverted_index_cache)
849        .with_puffin_metadata_cache(puffin_metadata_cache)
850        .build(filters)
851        .inspect_err(|err| warn!(err; "Failed to build invereted index applier"))
852        .ok()
853        .flatten()
854        .map(Arc::new)
855    }
856
857    /// Use the latest schema to build the bloom filter index applier.
858    fn build_bloom_filter_applier(&self, filters: &[Expr]) -> Option<BloomFilterIndexApplierRef> {
859        if self.ignore_bloom_filter {
860            return None;
861        }
862
863        let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
864        let bloom_filter_index_cache = self.cache_strategy.bloom_filter_index_cache().cloned();
865        let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
866
867        BloomFilterIndexApplierBuilder::new(
868            self.access_layer.table_dir().to_string(),
869            self.access_layer.path_type(),
870            self.access_layer.object_store().clone(),
871            self.version.metadata.as_ref(),
872            self.access_layer.puffin_manager_factory().clone(),
873        )
874        .with_file_cache(file_cache)
875        .with_bloom_filter_index_cache(bloom_filter_index_cache)
876        .with_puffin_metadata_cache(puffin_metadata_cache)
877        .build(filters)
878        .inspect_err(|err| warn!(err; "Failed to build bloom filter index applier"))
879        .ok()
880        .flatten()
881        .map(Arc::new)
882    }
883
884    /// Use the latest schema to build the fulltext index applier.
885    fn build_fulltext_index_applier(&self, filters: &[Expr]) -> Option<FulltextIndexApplierRef> {
886        if self.ignore_fulltext_index {
887            return None;
888        }
889
890        let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
891        let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
892        let bloom_filter_index_cache = self.cache_strategy.bloom_filter_index_cache().cloned();
893        FulltextIndexApplierBuilder::new(
894            self.access_layer.table_dir().to_string(),
895            self.access_layer.path_type(),
896            self.access_layer.object_store().clone(),
897            self.access_layer.puffin_manager_factory().clone(),
898            self.version.metadata.as_ref(),
899        )
900        .with_file_cache(file_cache)
901        .with_puffin_metadata_cache(puffin_metadata_cache)
902        .with_bloom_filter_cache(bloom_filter_index_cache)
903        .build(filters)
904        .inspect_err(|err| warn!(err; "Failed to build fulltext index applier"))
905        .ok()
906        .flatten()
907        .map(Arc::new)
908    }
909
910    /// Build the vector index applier from vector search request.
911    #[cfg(feature = "vector_index")]
912    fn build_vector_index_applier(&self) -> Option<VectorIndexApplierRef> {
913        let vector_search = self.request.vector_search.as_ref()?;
914
915        let file_cache = self.cache_strategy.write_cache().map(|w| w.file_cache());
916        let puffin_metadata_cache = self.cache_strategy.puffin_metadata_cache().cloned();
917        let vector_index_cache = self.cache_strategy.vector_index_cache().cloned();
918
919        let applier = VectorIndexApplier::new(
920            self.access_layer.table_dir().to_string(),
921            self.access_layer.path_type(),
922            self.access_layer.object_store().clone(),
923            self.access_layer.puffin_manager_factory().clone(),
924            vector_search.column_id,
925            vector_search.query_vector.clone(),
926            vector_search.metric,
927        )
928        .with_file_cache(file_cache)
929        .with_puffin_metadata_cache(puffin_metadata_cache)
930        .with_vector_index_cache(vector_index_cache);
931
932        Some(Arc::new(applier))
933    }
934}
935
936/// Returns true if the time range of a SST `file` matches the `predicate`.
937fn file_in_range(file: &FileHandle, predicate: &TimestampRange) -> bool {
938    if predicate == &TimestampRange::min_to_max() {
939        return true;
940    }
941    // end timestamp of a SST is inclusive.
942    let (start, end) = file.time_range();
943    let file_ts_range = TimestampRange::new_inclusive(Some(start), Some(end));
944    file_ts_range.intersects(predicate)
945}
946
947/// Returns true if `time_range` contains the inclusive time range of `file`.
948fn time_range_covers_file(time_range: Option<&TimestampRange>, file: &FileHandle) -> bool {
949    let Some(time_range) = time_range else {
950        return false;
951    };
952    let (start, end) = file.time_range();
953    time_range.contains(&start) && time_range.contains(&end)
954}
955
956/// Common input for different scanners.
957pub struct ScanInput {
958    /// Pinned series-index snapshot and its storage, when configured.
959    pub(crate) series_index: Option<SeriesIndexReadContext>,
960    /// Whether to ignore range indexes while retaining series-index candidate discovery.
961    ignore_range_index: bool,
962    /// Region SST access layer.
963    access_layer: AccessLayerRef,
964    /// Maps projected Batches to RecordBatches.
965    pub(crate) mapper: Arc<FlatProjectionMapper>,
966    /// The columns to read from memtables and SSTs.
967    /// Notice this is different from the columns in `mapper` which are projected columns.
968    /// But this read columns might also include non-projected columns needed for filtering.
969    pub(crate) read_cols: ReadColumns,
970    /// Time range filter for time index.
971    pub(crate) time_range: Option<TimestampRange>,
972    /// Cached analysis of the finalized scan request.
973    scan_analysis: Option<ScanAnalysis>,
974    /// Predicate to push down.
975    pub(crate) predicate: PredicateGroup,
976    /// Region partition expr applied at read time.
977    region_partition_expr: Option<PartitionExpr>,
978    /// Memtable range builders for memtables in the time range..
979    pub(crate) memtables: Vec<MemRangeBuilder>,
980    /// Handles to SST files to scan.
981    pub(crate) files: Vec<FileHandle>,
982    /// Shares the pinned schema's encoded defaults across parallel range readers.
983    primary_key_mapper: OnceLock<Arc<PrimaryKeyRangeMapper>>,
984    /// Scan-wide hint for rows in an execution batch.
985    batch_size: usize,
986    /// Cache.
987    pub(crate) cache_strategy: CacheStrategy,
988    /// Ignores file not found error.
989    ignore_file_not_found: bool,
990    /// Maximum number of SST files to scan concurrently.
991    pub(crate) max_concurrent_scan_files: usize,
992    /// Memory pool shared by internal scan operators across all queries.
993    pub(crate) scan_memory_pool: Arc<dyn MemoryPool>,
994    /// Index appliers.
995    inverted_index_appliers: [Option<InvertedIndexApplierRef>; 2],
996    bloom_filter_index_appliers: [Option<BloomFilterIndexApplierRef>; 2],
997    fulltext_index_appliers: [Option<FulltextIndexApplierRef>; 2],
998    /// Vector index applier for KNN search.
999    #[cfg(feature = "vector_index")]
1000    pub(crate) vector_index_applier: Option<VectorIndexApplierRef>,
1001    /// Over-fetched k for vector index scan.
1002    #[cfg(feature = "vector_index")]
1003    pub(crate) vector_index_k: Option<usize>,
1004    /// Start time of the query.
1005    pub(crate) query_start: Option<Instant>,
1006    /// The region is using append mode.
1007    pub(crate) append_mode: bool,
1008    /// Whether to remove deletion markers.
1009    pub(crate) filter_deleted: bool,
1010    /// Mode to merge duplicate rows.
1011    pub(crate) merge_mode: MergeMode,
1012    /// Hint to select rows from time series.
1013    pub(crate) series_row_selector: Option<TimeSeriesRowSelector>,
1014    /// Hint for the required distribution of the scanner.
1015    pub(crate) distribution: Option<TimeSeriesDistribution>,
1016    /// Whether the region's configured SST format is flat.
1017    explain_flat_format: bool,
1018    /// Snapshot upper bound bound at scan open and propagated back to the caller.
1019    pub(crate) snapshot_sequence: Option<SequenceNumber>,
1020    /// Set only when the region preserves per-row sequences
1021    /// (`preserve_row_sequence` on an append-only table) and the scan
1022    /// carries an explicit exact range. When set, SST readers apply the same
1023    /// `(checkpoint, upper_bound]` row-level filter as memtables.
1024    pub(crate) sequence_range: Option<SequenceRange>,
1025    /// Whether this scan is for compaction.
1026    pub(crate) compaction: bool,
1027    /// Compaction-only JSON2 physical rewrite targets.
1028    json2_rewrite_targets: Json2RewriteTargets,
1029    /// Counters that should receive query-load metrics.
1030    pub(crate) query_stat_counters: Option<RegionQueryStatCounters>,
1031    #[cfg(feature = "enterprise")]
1032    extension_ranges: Vec<BoxedExtensionRange>,
1033}
1034
1035/// Cached output of scan request analysis: the optional range-cache fingerprint
1036/// plus the derived implied time range used by range cache and prefilter.
1037struct ScanAnalysis {
1038    /// `None` when the scan is not eligible for range caching.
1039    fingerprint: Option<ScanRequestFingerprint>,
1040    /// `Some(r)` = all time-only predicates are guaranteed true on `r` (in the
1041    /// column's `TimeUnit`).
1042    /// `None`    = there is no analyzable time-only predicate or at least one
1043    /// predicate could not be proven (e.g. `OR`), so the time-filter
1044    /// optimization is disabled for this scan.
1045    implied_time_range: Option<TimestampRange>,
1046}
1047
1048/// Builder for a finalized [ScanInput].
1049pub(crate) struct ScanInputBuilder {
1050    input: ScanInput,
1051}
1052
1053impl ScanInput {
1054    /// Creates a new [ScanInputBuilder].
1055    #[must_use]
1056    pub(crate) fn builder(
1057        access_layer: AccessLayerRef,
1058        mapper: FlatProjectionMapper,
1059    ) -> ScanInputBuilder {
1060        ScanInputBuilder {
1061            input: ScanInput {
1062                series_index: None,
1063                ignore_range_index: false,
1064                access_layer,
1065                read_cols: mapper.read_columns().clone(),
1066                mapper: Arc::new(mapper),
1067                time_range: None,
1068                scan_analysis: None,
1069                predicate: PredicateGroup::default(),
1070                region_partition_expr: None,
1071                memtables: Vec::new(),
1072                files: Vec::new(),
1073                primary_key_mapper: OnceLock::new(),
1074                batch_size: crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
1075                cache_strategy: CacheStrategy::Disabled,
1076                ignore_file_not_found: false,
1077                max_concurrent_scan_files: DEFAULT_MAX_CONCURRENT_SCAN_FILES,
1078                scan_memory_pool: Arc::new(UnboundedMemoryPool::default()),
1079                inverted_index_appliers: [None, None],
1080                bloom_filter_index_appliers: [None, None],
1081                fulltext_index_appliers: [None, None],
1082                #[cfg(feature = "vector_index")]
1083                vector_index_applier: None,
1084                #[cfg(feature = "vector_index")]
1085                vector_index_k: None,
1086                query_start: None,
1087                append_mode: false,
1088                filter_deleted: true,
1089                merge_mode: MergeMode::default(),
1090                series_row_selector: None,
1091                distribution: None,
1092                explain_flat_format: false,
1093                snapshot_sequence: None,
1094                sequence_range: None,
1095                compaction: false,
1096                json2_rewrite_targets: Arc::default(),
1097                query_stat_counters: None,
1098                #[cfg(feature = "enterprise")]
1099                extension_ranges: Vec::new(),
1100            },
1101        }
1102    }
1103
1104    /// Returns the scan-wide hint for rows in an execution batch.
1105    pub(crate) fn batch_size(&self) -> usize {
1106        self.batch_size
1107    }
1108
1109    /// Interprets file statistics using this scan's pinned schema.
1110    pub(crate) fn primary_key_mapper(&self) -> &PrimaryKeyRangeMapper {
1111        self.primary_key_mapper
1112            .get_or_init(|| Arc::new(PrimaryKeyRangeMapper::new(self.region_metadata().clone())))
1113    }
1114
1115    /// Returns the range implied by the range-cache time filters.
1116    pub(crate) fn implied_time_range(&self) -> Option<&TimestampRange> {
1117        self.scan_analysis
1118            .as_ref()
1119            .expect("ScanInput must be built")
1120            .implied_time_range
1121            .as_ref()
1122    }
1123
1124    /// Returns the fingerprint for partition range caching.
1125    pub(crate) fn scan_fingerprint(&self) -> Option<&ScanRequestFingerprint> {
1126        self.scan_analysis
1127            .as_ref()
1128            .expect("ScanInput must be built")
1129            .fingerprint
1130            .as_ref()
1131    }
1132}
1133
1134impl ScanInputBuilder {
1135    fn with_primary_key_mapper(mut self, mapper: Arc<PrimaryKeyRangeMapper>) -> Self {
1136        self.input.primary_key_mapper = OnceLock::from(mapper);
1137        self
1138    }
1139
1140    /// Sets whether to ignore range indexes during scans.
1141    #[must_use]
1142    pub(crate) fn with_ignore_range_index(mut self, ignore: bool) -> Self {
1143        self.input.ignore_range_index = ignore;
1144        self
1145    }
1146
1147    /// Sets the pinned series-index context for candidate discovery.
1148    #[must_use]
1149    pub(crate) fn with_series_index(
1150        mut self,
1151        series_index: Option<SeriesIndexReadContext>,
1152    ) -> Self {
1153        self.input.series_index = series_index;
1154        self
1155    }
1156
1157    /// Sets time range filter for time index.
1158    #[must_use]
1159    pub(crate) fn with_time_range(mut self, time_range: Option<TimestampRange>) -> Self {
1160        self.input.time_range = time_range;
1161        self
1162    }
1163
1164    /// Sets predicate to push down.
1165    #[must_use]
1166    pub(crate) fn with_predicate(mut self, predicate: PredicateGroup) -> Self {
1167        self.input.region_partition_expr = predicate.region_partition_expr().cloned();
1168        self.input.predicate = predicate;
1169        self
1170    }
1171
1172    /// Builds a finalized [ScanInput] and computes its scan analysis.
1173    #[must_use]
1174    pub(crate) fn build(mut self) -> ScanInput {
1175        let input = &self.input;
1176        let eligible = !input.compaction
1177            && !input.files.is_empty()
1178            && matches!(input.cache_strategy, CacheStrategy::EnableAll(_));
1179
1180        let metadata = input.region_metadata();
1181        let tag_names: HashSet<&str> = metadata
1182            .column_metadatas
1183            .iter()
1184            .filter(|col| col.semantic_type == SemanticType::Tag)
1185            .map(|col| col.column_schema.name.as_str())
1186            .collect();
1187
1188        let time_index = metadata.time_index_column();
1189        let time_index_name = time_index.column_schema.name.clone();
1190        let ts_col_unit = time_index
1191            .column_schema
1192            .data_type
1193            .as_timestamp()
1194            .expect("Time index must have timestamp-compatible type")
1195            .unit();
1196
1197        let exprs = input
1198            .predicate_group()
1199            .predicate_without_region()
1200            .map(|predicate| predicate.exprs())
1201            .unwrap_or_default();
1202
1203        let mut filters = Vec::new();
1204        let mut time_only_exprs: Vec<&Expr> = Vec::new();
1205        let mut has_tag_filter = false;
1206        let mut columns = HashSet::new();
1207
1208        for expr in exprs {
1209            columns.clear();
1210            let is_time_only = match expr_to_columns(expr, &mut columns) {
1211                Ok(()) if !columns.is_empty() => {
1212                    has_tag_filter |= columns
1213                        .iter()
1214                        .any(|col| tag_names.contains(col.name.as_str()));
1215                    columns.iter().all(|col| col.name == time_index_name)
1216                }
1217                _ => false,
1218            };
1219
1220            // Route time-only exprs that the legacy extractor recognizes into
1221            // `time_only_exprs` so the implication walker
1222            // (`implied_time_range_from_exprs`, called below) can attempt to drop
1223            // them from the cache key when the partition's `FileTimeRange` is fully
1224            // covered, then stringify them into the fingerprint's `time_filters`
1225            // bucket. Time-only exprs that the extractor doesn't recognize stay in
1226            // `filters` and never get stripped — conservatively correct.
1227            if is_time_only
1228                && extract_time_range_from_expr(&time_index_name, ts_col_unit, expr).is_some()
1229            {
1230                time_only_exprs.push(expr);
1231            } else {
1232                filters.push(expr.to_string());
1233            }
1234        }
1235
1236        let implied_time_range =
1237            implied_time_range_from_exprs(&time_index_name, ts_col_unit, &time_only_exprs);
1238
1239        // We only cache requests that have tag filters to avoid caching all series.
1240        let fingerprint = if eligible && has_tag_filter {
1241            let mut time_filters: Vec<String> =
1242                time_only_exprs.iter().map(|e| e.to_string()).collect();
1243
1244            // Ensure the filters are sorted for consistent fingerprinting.
1245            filters.sort_unstable();
1246            time_filters.sort_unstable();
1247            let read_columns = input.read_cols.clone();
1248            let fingerprint = crate::read::range_cache::ScanRequestFingerprintBuilder {
1249                read_column_types: read_columns
1250                    .column_ids_iter()
1251                    .map(|id| {
1252                        read_columns
1253                            .json_target_type(id)
1254                            .cloned()
1255                            .map(ConcreteDataType::json2)
1256                            .or_else(|| {
1257                                metadata
1258                                    .column_by_id(id)
1259                                    .map(|col| col.column_schema.data_type.clone())
1260                            })
1261                    })
1262                    .collect(),
1263                read_columns,
1264                filters,
1265                time_filters,
1266                series_row_selector: input.series_row_selector,
1267                append_mode: input.append_mode,
1268                filter_deleted: input.filter_deleted,
1269                merge_mode: input.merge_mode,
1270                sequence_range: input.sequence_range,
1271                partition_expr_version: metadata.partition_expr_version,
1272            }
1273            .build();
1274            Some(fingerprint)
1275        } else {
1276            None
1277        };
1278
1279        self.input.scan_analysis = Some(ScanAnalysis {
1280            fingerprint,
1281            implied_time_range,
1282        });
1283        self.input
1284    }
1285
1286    /// Sets memtable range builders.
1287    #[must_use]
1288    pub(crate) fn with_memtables(mut self, memtables: Vec<MemRangeBuilder>) -> Self {
1289        self.input.memtables = memtables;
1290        self
1291    }
1292
1293    /// Sets files to read.
1294    #[must_use]
1295    pub(crate) fn with_files(mut self, files: Vec<FileHandle>) -> Self {
1296        self.input.files = files;
1297        self
1298    }
1299
1300    /// Sets the scan-wide hint for rows in an execution batch.
1301    #[must_use]
1302    pub(crate) fn with_batch_size(mut self, batch_size: usize) -> Self {
1303        self.input.batch_size = batch_size;
1304        self
1305    }
1306
1307    /// Sets cache for this query.
1308    #[must_use]
1309    pub(crate) fn with_cache(mut self, cache: CacheStrategy) -> Self {
1310        self.input.cache_strategy = cache;
1311        self
1312    }
1313
1314    /// Ignores file not found error.
1315    #[must_use]
1316    pub(crate) fn with_ignore_file_not_found(mut self, ignore: bool) -> Self {
1317        self.input.ignore_file_not_found = ignore;
1318        self
1319    }
1320
1321    /// Sets maximum number of SST files to scan concurrently.
1322    #[must_use]
1323    pub(crate) fn with_max_concurrent_scan_files(
1324        mut self,
1325        max_concurrent_scan_files: usize,
1326    ) -> Self {
1327        self.input.max_concurrent_scan_files = max_concurrent_scan_files;
1328        self
1329    }
1330
1331    /// Sets the memory pool shared by internal scan operators.
1332    #[must_use]
1333    pub(crate) fn with_scan_memory_pool(mut self, scan_memory_pool: Arc<dyn MemoryPool>) -> Self {
1334        self.input.scan_memory_pool = scan_memory_pool;
1335        self
1336    }
1337
1338    /// Sets inverted index appliers.
1339    #[must_use]
1340    pub(crate) fn with_inverted_index_appliers(
1341        mut self,
1342        appliers: [Option<InvertedIndexApplierRef>; 2],
1343    ) -> Self {
1344        self.input.inverted_index_appliers = appliers;
1345        self
1346    }
1347
1348    /// Sets bloom filter appliers.
1349    #[must_use]
1350    pub(crate) fn with_bloom_filter_index_appliers(
1351        mut self,
1352        appliers: [Option<BloomFilterIndexApplierRef>; 2],
1353    ) -> Self {
1354        self.input.bloom_filter_index_appliers = appliers;
1355        self
1356    }
1357
1358    /// Sets fulltext index appliers.
1359    #[must_use]
1360    pub(crate) fn with_fulltext_index_appliers(
1361        mut self,
1362        appliers: [Option<FulltextIndexApplierRef>; 2],
1363    ) -> Self {
1364        self.input.fulltext_index_appliers = appliers;
1365        self
1366    }
1367
1368    /// Sets vector index applier for KNN search.
1369    #[cfg(feature = "vector_index")]
1370    #[must_use]
1371    pub(crate) fn with_vector_index_applier(
1372        mut self,
1373        applier: Option<VectorIndexApplierRef>,
1374    ) -> Self {
1375        self.input.vector_index_applier = applier;
1376        self
1377    }
1378
1379    /// Sets over-fetched k for vector index scan.
1380    #[cfg(feature = "vector_index")]
1381    #[must_use]
1382    pub(crate) fn with_vector_index_k(mut self, k: Option<usize>) -> Self {
1383        self.input.vector_index_k = k;
1384        self
1385    }
1386
1387    /// Sets start time of the query.
1388    #[must_use]
1389    pub(crate) fn with_start_time(mut self, now: Option<Instant>) -> Self {
1390        self.input.query_start = now;
1391        self
1392    }
1393
1394    #[must_use]
1395    pub(crate) fn with_append_mode(mut self, is_append_mode: bool) -> Self {
1396        self.input.append_mode = is_append_mode;
1397        self
1398    }
1399
1400    pub(crate) fn with_query_stat_counters(
1401        mut self,
1402        counters: Option<RegionQueryStatCounters>,
1403    ) -> Self {
1404        self.input.query_stat_counters = counters;
1405        self
1406    }
1407
1408    /// Sets whether to remove deletion markers during scan.
1409    #[must_use]
1410    pub(crate) fn with_filter_deleted(mut self, filter_deleted: bool) -> Self {
1411        self.input.filter_deleted = filter_deleted;
1412        self
1413    }
1414
1415    /// Sets the merge mode.
1416    #[must_use]
1417    pub(crate) fn with_merge_mode(mut self, merge_mode: MergeMode) -> Self {
1418        self.input.merge_mode = merge_mode;
1419        self
1420    }
1421
1422    /// Sets the distribution hint.
1423    #[must_use]
1424    pub(crate) fn with_distribution(
1425        mut self,
1426        distribution: Option<TimeSeriesDistribution>,
1427    ) -> Self {
1428        self.input.distribution = distribution;
1429        self
1430    }
1431
1432    /// Sets whether the region's configured SST format is flat for explain output.
1433    #[must_use]
1434    pub(crate) fn with_explain_flat_format(mut self, explain_flat_format: bool) -> Self {
1435        self.input.explain_flat_format = explain_flat_format;
1436        self
1437    }
1438
1439    /// Sets the time series row selector.
1440    #[must_use]
1441    pub(crate) fn with_series_row_selector(
1442        mut self,
1443        series_row_selector: Option<TimeSeriesRowSelector>,
1444    ) -> Self {
1445        self.input.series_row_selector = series_row_selector;
1446        self
1447    }
1448
1449    #[must_use]
1450    pub(crate) fn with_snapshot_sequence(
1451        mut self,
1452        snapshot_sequence: Option<SequenceNumber>,
1453    ) -> Self {
1454        self.input.snapshot_sequence = snapshot_sequence;
1455        self
1456    }
1457
1458    #[must_use]
1459    pub(crate) fn with_sequence_range(mut self, sequence_range: Option<SequenceRange>) -> Self {
1460        self.input.sequence_range = sequence_range;
1461        self
1462    }
1463
1464    /// Sets whether this scan is for compaction.
1465    #[must_use]
1466    pub(crate) fn with_compaction(mut self, compaction: bool) -> Self {
1467        self.input.compaction = compaction;
1468        self
1469    }
1470
1471    /// Sets compaction-only JSON2 physical rewrite targets.
1472    #[must_use]
1473    pub(crate) fn with_json2_rewrite_targets(mut self, targets: Json2RewriteTargets) -> Self {
1474        self.input.json2_rewrite_targets = targets;
1475        self
1476    }
1477
1478    #[cfg(feature = "enterprise")]
1479    #[must_use]
1480    pub(crate) fn with_extension_ranges(
1481        mut self,
1482        extension_ranges: Vec<BoxedExtensionRange>,
1483    ) -> Self {
1484        self.input.extension_ranges = extension_ranges;
1485        self
1486    }
1487}
1488
1489impl ScanInput {
1490    /// Builds memtable ranges to scan by `index`.
1491    pub(crate) fn build_mem_ranges(&self, index: RowGroupIndex) -> SmallVec<[MemtableRange; 2]> {
1492        let memtable = &self.memtables[index.index];
1493        let mut ranges = SmallVec::new();
1494        memtable.build_ranges(index.row_group_index, &mut ranges);
1495        ranges
1496    }
1497
1498    pub(crate) fn predicate_for_file(&self, file: &FileHandle) -> Option<Predicate> {
1499        if self.should_skip_region_partition(file) {
1500            self.predicate.predicate_without_region().cloned()
1501        } else {
1502            self.predicate.predicate().cloned()
1503        }
1504    }
1505
1506    fn should_skip_region_partition(&self, file: &FileHandle) -> bool {
1507        match (
1508            self.region_partition_expr.as_ref(),
1509            file.meta_ref().partition_expr.as_ref(),
1510        ) {
1511            (Some(region_expr), Some(file_expr)) => region_expr == file_expr,
1512            _ => false,
1513        }
1514    }
1515
1516    /// Tries to build file-level pruning statistics using only the [FileHandle]'s manifest-level
1517    /// time range, without reading any parquet metadata.
1518    ///
1519    /// Returns `None` if timestamp unit conversion overflows (conservative: keep the file).
1520    fn try_file_level_pruning_stats(&self, file: &FileHandle) -> Option<FileLevelPruningStats> {
1521        let (ts_min, ts_max) = file.time_range();
1522        let time_index = self.mapper.metadata().time_index_column();
1523        let time_index_unit = time_index.column_schema.data_type.as_timestamp()?.unit();
1524
1525        // Convert file timestamps to the time index column's unit. Use `convert_to_ceil` for
1526        // the upper bound to avoid accidentally shrinking the manifest range.
1527        let min_ts = ts_min.convert_to(time_index_unit)?;
1528        let max_ts = ts_max.convert_to_ceil(time_index_unit)?;
1529
1530        Some(FileLevelPruningStats {
1531            min_scalar: timestamp_to_scalar_value(time_index_unit, Some(min_ts.value())),
1532            max_scalar: timestamp_to_scalar_value(time_index_unit, Some(max_ts.value())),
1533            time_index_col_name: time_index.column_schema.name.clone(),
1534        })
1535    }
1536
1537    /// Checks whether a file can be definitively pruned using only its manifest-level
1538    /// time range and the current predicate, without reading any parquet metadata.
1539    ///
1540    /// Returns `true` if [PruningStatistics] proves the file cannot contain matching rows.
1541    #[inline]
1542    pub(crate) fn can_manifest_prune_file(&self, file: &FileHandle) -> bool {
1543        let predicate = self.predicate_for_file(file);
1544        self.manifest_prunes_file(file, predicate.as_ref())
1545    }
1546
1547    fn manifest_prunes_file(&self, file: &FileHandle, predicate: Option<&Predicate>) -> bool {
1548        if let Some(pred) = predicate
1549            && !pred.is_empty()
1550            && let Some(file_level_stats) = self.try_file_level_pruning_stats(file)
1551        {
1552            let pruning_results = pred.prune_with_stats(
1553                &file_level_stats,
1554                self.mapper.metadata().schema.arrow_schema(),
1555            );
1556            pruning_results.first() == Some(&false)
1557        } else {
1558            false
1559        }
1560    }
1561
1562    /// Prunes a file to scan and returns the builder to build readers.
1563    ///
1564    /// This is the public entry point used by direct tests and non-pruner callers.
1565    /// It performs its own manifest-level pruning check internally.
1566    #[tracing::instrument(
1567        skip_all,
1568        fields(
1569            region_id = %self.region_metadata().region_id,
1570            file_id = %file.file_id()
1571        )
1572    )]
1573    pub async fn prune_file(
1574        &self,
1575        file: &FileHandle,
1576        pre_filter_mode: PreFilterMode,
1577        reader_metrics: &mut ReaderMetrics,
1578    ) -> Result<FileRangeBuilder> {
1579        let predicate = self.predicate_for_file(file);
1580
1581        // Early file-level pruning using manifest time range before any parquet metadata access.
1582        if self.manifest_prunes_file(file, predicate.as_ref()) {
1583            reader_metrics.filter_metrics.files_time_range_pruned += 1;
1584            return Ok(FileRangeBuilder::default());
1585        }
1586
1587        self.prune_file_after_manifest_check(file, pre_filter_mode, true, predicate, reader_metrics)
1588            .await
1589    }
1590
1591    /// Second half of `prune_file` — performs the actual parquet metadata /
1592    /// reader setup. Callers that already performed manifest-level pruning
1593    /// (e.g. the `Pruner` via its shared `manifest_pruned_files` cache) should
1594    /// call this directly to avoid a redundant manifest check.
1595    ///
1596    /// `predicate` is the result of `self.predicate_for_file(file)` computed
1597    /// externally so the caller can reuse it if needed.
1598    pub(crate) async fn prune_file_after_manifest_check(
1599        &self,
1600        file: &FileHandle,
1601        pre_filter_mode: PreFilterMode,
1602        enable_predicate_prefilter: bool,
1603        predicate: Option<Predicate>,
1604        reader_metrics: &mut ReaderMetrics,
1605    ) -> Result<FileRangeBuilder> {
1606        let may_build_selective_row_selection = predicate.is_some();
1607        let postpone_time_index_filter = time_range_covers_file(self.implied_time_range(), file);
1608        let decode_pk_values = !self.compaction
1609            && self
1610                .mapper
1611                .read_columns()
1612                .column_ids_iter()
1613                .any(|column_id| self.mapper.metadata().primary_key.contains(&column_id));
1614        let reader = self
1615            .access_layer
1616            .read_sst(file.clone())
1617            .series_index(
1618                self.series_index
1619                    .clone()
1620                    .filter(|_| !self.ignore_range_index),
1621            )
1622            .predicate(predicate)
1623            .projection(Some(self.read_cols.clone()))
1624            .json2_rewrite_targets(self.json2_rewrite_targets.clone())
1625            .cache(self.cache_strategy.clone())
1626            .inverted_index_appliers(self.inverted_index_appliers.clone())
1627            .bloom_filter_index_appliers(self.bloom_filter_index_appliers.clone())
1628            .fulltext_index_appliers(self.fulltext_index_appliers.clone());
1629        let reader = reader.batch_size(self.batch_size);
1630        let reader = if !self.compaction && may_build_selective_row_selection {
1631            reader.deferred_optional_page_index()
1632        } else {
1633            reader
1634        };
1635        #[cfg(feature = "vector_index")]
1636        let reader = {
1637            let mut reader = reader;
1638            reader =
1639                reader.vector_index_applier(self.vector_index_applier.clone(), self.vector_index_k);
1640            reader
1641        };
1642        let res = reader
1643            .expected_metadata(Some(self.mapper.metadata().clone()))
1644            .compaction(self.compaction)
1645            .pre_filter_mode(pre_filter_mode)
1646            .enable_predicate_prefilter(enable_predicate_prefilter)
1647            .postpone_time_index_filter(postpone_time_index_filter)
1648            .decode_primary_key_values(decode_pk_values)
1649            .build_reader_input(reader_metrics)
1650            .await;
1651        let read_input = match res {
1652            Ok(x) => x,
1653            Err(e) => {
1654                if e.is_object_not_found() && self.ignore_file_not_found {
1655                    error!(e; "File to scan does not exist, region_id: {}, file: {}", file.region_id(), file.file_id());
1656                    return Ok(FileRangeBuilder::default());
1657                } else {
1658                    return Err(e);
1659                }
1660            }
1661        };
1662
1663        let Some((mut file_range_ctx, selection)) = read_input else {
1664            return Ok(FileRangeBuilder::default());
1665        };
1666
1667        let need_compat = !compat::has_same_columns_and_pk_encoding(
1668            &self.mapper,
1669            file_range_ctx.read_format(),
1670            self.compaction,
1671        );
1672        if need_compat {
1673            // They have different schema. We need to adapt the batch first so the
1674            // mapper can convert it.
1675            let compat = FlatCompatBatch::try_new(
1676                &self.mapper,
1677                file_range_ctx.read_format(),
1678                self.compaction,
1679            )?;
1680            file_range_ctx.set_compat_batch(compat);
1681        }
1682        Ok(FileRangeBuilder::new(Arc::new(file_range_ctx), selection))
1683    }
1684
1685    /// Scans flat sources (RecordBatch streams) in parallel.
1686    ///
1687    /// # Panics if the input doesn't allow parallel scan.
1688    #[tracing::instrument(
1689        skip(self, sources, semaphore),
1690        fields(
1691            region_id = %self.region_metadata().region_id,
1692            source_count = sources.len()
1693        )
1694    )]
1695    pub(crate) fn create_parallel_flat_sources(
1696        &self,
1697        sources: Vec<BoxedRecordBatchStream>,
1698        semaphore: Arc<Semaphore>,
1699        channel_size: usize,
1700    ) -> Result<Vec<BoxedRecordBatchStream>> {
1701        if sources.len() <= 1 {
1702            return Ok(sources);
1703        }
1704
1705        // Spawn a task for each source.
1706        let sources = sources
1707            .into_iter()
1708            .map(|source| {
1709                let (sender, receiver) = mpsc::channel(channel_size);
1710                self.spawn_flat_scan_task(source, semaphore.clone(), sender);
1711                let stream = Box::pin(ReceiverStream::new(receiver));
1712                Box::pin(stream) as _
1713            })
1714            .collect();
1715        Ok(sources)
1716    }
1717
1718    /// Spawns a task to scan a flat source (RecordBatch stream) asynchronously.
1719    #[tracing::instrument(
1720        skip(self, input, semaphore, sender),
1721        fields(region_id = %self.region_metadata().region_id)
1722    )]
1723    pub(crate) fn spawn_flat_scan_task(
1724        &self,
1725        mut input: BoxedRecordBatchStream,
1726        semaphore: Arc<Semaphore>,
1727        sender: mpsc::Sender<Result<RecordBatch>>,
1728    ) {
1729        let region_id = self.region_metadata().region_id;
1730        let span = tracing::info_span!(
1731            "ScanInput::parallel_scan_task",
1732            region_id = %region_id,
1733            stream_kind = "flat"
1734        );
1735        common_runtime::spawn_query(
1736            async move {
1737                loop {
1738                    // We release the permit before sending result to avoid the task waiting on
1739                    // the channel with the permit held.
1740                    let maybe_batch = {
1741                        // Safety: We never close the semaphore.
1742                        let _permit = semaphore.acquire().await.unwrap();
1743                        input.next().await
1744                    };
1745                    match maybe_batch {
1746                        Some(Ok(batch)) => {
1747                            let _ = sender.send(Ok(batch)).await;
1748                        }
1749                        Some(Err(e)) => {
1750                            let _ = sender.send(Err(e)).await;
1751                            break;
1752                        }
1753                        None => break,
1754                    }
1755                }
1756            }
1757            .instrument(span),
1758        );
1759    }
1760
1761    /// Whether physical source row counts equal the rows visible to this region,
1762    /// before applying query predicates.
1763    pub(crate) fn total_rows_is_exact(&self) -> bool {
1764        if self.region_partition_expr.is_none() {
1765            return true;
1766        }
1767
1768        // Extension ranges do not expose their partition expressions.
1769        #[cfg(feature = "enterprise")]
1770        if !self.extension_ranges.is_empty() {
1771            return false;
1772        }
1773
1774        // Repartition flushes existing memtables before entering staging. New
1775        // writes are routed under the target partition rules, whereas historical
1776        // SSTs can be shared by regions and still need partition filtering.
1777        self.files
1778            .iter()
1779            .all(|file| self.should_skip_region_partition(file))
1780    }
1781
1782    pub(crate) fn total_rows(&self) -> usize {
1783        let rows_in_files: usize = self.files.iter().map(|f| f.num_rows()).sum();
1784        let rows_in_memtables: usize = self.memtables.iter().map(|m| m.stats().num_rows()).sum();
1785
1786        let rows = rows_in_files + rows_in_memtables;
1787        #[cfg(feature = "enterprise")]
1788        let rows = rows
1789            + self
1790                .extension_ranges
1791                .iter()
1792                .map(|x| x.num_rows())
1793                .sum::<u64>() as usize;
1794        rows
1795    }
1796
1797    pub(crate) fn predicate_group(&self) -> &PredicateGroup {
1798        &self.predicate
1799    }
1800
1801    /// Returns number of memtables to scan.
1802    pub(crate) fn num_memtables(&self) -> usize {
1803        self.memtables.len()
1804    }
1805
1806    /// Returns number of SST files to scan.
1807    pub(crate) fn num_files(&self) -> usize {
1808        self.files.len()
1809    }
1810
1811    /// Gets the file handle from a row group index.
1812    pub(crate) fn file_from_index(&self, index: RowGroupIndex) -> &FileHandle {
1813        let file_index = index.index - self.num_memtables();
1814        &self.files[file_index]
1815    }
1816
1817    pub fn region_metadata(&self) -> &RegionMetadataRef {
1818        self.mapper.metadata()
1819    }
1820
1821    fn range_pre_filter_mode(&self, source_count: usize) -> PreFilterMode {
1822        if source_count <= 1 {
1823            // Duplicated rows in the same source is not a normal case and we don't provide
1824            // strict dedup semantic (last_row/last_non_null) for it. We expect the duplicated rows
1825            // are exactly identical in the same source so we use PreFilterMode::All for
1826            // performance reason.
1827            return PreFilterMode::All;
1828        }
1829
1830        pre_filter_mode(self.append_mode, self.merge_mode)
1831    }
1832}
1833
1834#[cfg(feature = "enterprise")]
1835impl ScanInput {
1836    #[cfg(feature = "enterprise")]
1837    pub(crate) fn extension_ranges(&self) -> &[BoxedExtensionRange] {
1838        &self.extension_ranges
1839    }
1840
1841    /// Get a boxed [ExtensionRange] by the index in all ranges.
1842    #[cfg(feature = "enterprise")]
1843    pub(crate) fn extension_range(&self, i: usize) -> &BoxedExtensionRange {
1844        &self.extension_ranges[i - self.num_memtables() - self.num_files()]
1845    }
1846}
1847
1848/// Lightweight [PruningStatistics] that only uses the file-level time range from manifest
1849/// metadata, avoiding any parquet metadata reads. Used for early file-level pruning before
1850/// accessing row-group-level statistics.
1851pub(crate) struct FileLevelPruningStats {
1852    /// Scalar value for the file's minimum timestamp in the time index column's unit.
1853    pub(crate) min_scalar: ScalarValue,
1854    /// Scalar value for the file's maximum timestamp in the time index column's unit.
1855    pub(crate) max_scalar: ScalarValue,
1856    /// Name of the time index column.
1857    pub(crate) time_index_col_name: String,
1858}
1859
1860impl PruningStatistics for FileLevelPruningStats {
1861    fn min_values(&self, column: &Column) -> Option<ArrayRef> {
1862        if column.name == self.time_index_col_name {
1863            ScalarValue::iter_to_array(std::iter::once(self.min_scalar.clone())).ok()
1864        } else {
1865            None
1866        }
1867    }
1868
1869    fn max_values(&self, column: &Column) -> Option<ArrayRef> {
1870        if column.name == self.time_index_col_name {
1871            ScalarValue::iter_to_array(std::iter::once(self.max_scalar.clone())).ok()
1872        } else {
1873            None
1874        }
1875    }
1876
1877    fn num_containers(&self) -> usize {
1878        1
1879    }
1880
1881    fn null_counts(&self, column: &Column) -> Option<ArrayRef> {
1882        if column.name == self.time_index_col_name {
1883            // The time index column is NOT NULL.
1884            Some(Arc::new(UInt64Array::from(vec![0u64])))
1885        } else {
1886            None
1887        }
1888    }
1889
1890    fn row_counts(&self) -> Option<ArrayRef> {
1891        None
1892    }
1893
1894    fn contained(&self, _column: &Column, _values: &HashSet<ScalarValue>) -> Option<BooleanArray> {
1895        None
1896    }
1897}
1898
1899#[cfg(test)]
1900impl ScanInput {
1901    /// Returns SST file ids to scan.
1902    pub(crate) fn file_ids(&self) -> Vec<crate::sst::file::RegionFileId> {
1903        self.files.iter().map(|file| file.file_id()).collect()
1904    }
1905
1906    pub(crate) fn index_ids(&self) -> Vec<crate::sst::file::RegionIndexId> {
1907        self.files.iter().map(|file| file.index_id()).collect()
1908    }
1909}
1910
1911fn pre_filter_mode(append_mode: bool, merge_mode: MergeMode) -> PreFilterMode {
1912    if append_mode {
1913        return PreFilterMode::All;
1914    }
1915
1916    match merge_mode {
1917        MergeMode::LastRow => PreFilterMode::SkipFields,
1918        MergeMode::LastNonNull => PreFilterMode::SkipFields,
1919    }
1920}
1921
1922/// Selects the SST files for a scan and, when requested, determines whether
1923/// the selected files support an exact sequence-range scan.
1924///
1925/// Files excluded by the request time range are not selected: `(C, H]` rows are
1926/// a subset of the query's time-range rows, so a time-pruned file cannot
1927/// contribute a row to `(C, H]`.
1928///
1929/// Unmarked local files use `FileMeta.sequence` as an admission barrier rather
1930/// than a row maximum. Region edits allocate this barrier; compaction inherits
1931/// the maximum input bound without assigning a new barrier. `C >= barrier`
1932/// proves Flow has already consumed the entire file, so such a file is excluded
1933/// before the capability check. An unknown bound cannot prove this exclusion.
1934///
1935/// A foreign file is different: the parquet reader virtualizes every row to its
1936/// target-local `FileMeta.sequence`. Consequently, a present sequence is the
1937/// only trust requirement for a foreign file; its source marker is irrelevant.
1938/// A foreign file without that barrier is retained as a failed-closed error so
1939/// it cannot be mistaken for a local sequence domain. This check is performed
1940/// even when another exact-range capability condition would return `None`.
1941pub(crate) fn exact_sequence_range(
1942    request: &ScanRequest,
1943    version: &crate::region::version::Version,
1944) -> Result<(Vec<FileHandle>, Option<SequenceRange>)> {
1945    if request.skip_sst_files {
1946        return Ok((Vec::new(), None));
1947    }
1948
1949    let time_index = version.metadata.time_index_column();
1950    let unit = time_index
1951        .column_schema
1952        .data_type
1953        .as_timestamp()
1954        .expect("Time index must have timestamp-compatible type")
1955        .unit();
1956    let time_range =
1957        build_time_range_predicate(&time_index.column_schema.name, unit, &request.filters);
1958    let min = request.memtable_min_sequence;
1959    let mut check_capability = request.exact_sequence_range && min.is_some();
1960    let sst_min_sequence = request.sst_min_sequence.and_then(NonZeroU64::new);
1961    let mut files = Vec::new();
1962    let mut files_allow_exact_range = true;
1963
1964    for file in version
1965        .ssts
1966        .levels()
1967        .iter()
1968        .flat_map(|level| level.files.values())
1969        .filter(|file| file_in_range(file, &time_range))
1970    {
1971        let meta = file.meta_ref();
1972        let selected = (!request.exact_sequence_range
1973            || min.is_none_or(|min| meta.sequence.is_none_or(|sequence| sequence.get() > min)))
1974            && match (sst_min_sequence, meta.sequence) {
1975                (Some(min_sequence), Some(file_sequence)) => file_sequence > min_sequence,
1976                // A missing file sequence is treated as newer than the SST
1977                // pruning hint, matching scan input construction.
1978                (Some(_), None) | (None, _) => true,
1979            };
1980        if !selected {
1981            continue;
1982        }
1983
1984        if let Some(min) = min
1985            && check_capability
1986        {
1987            if meta.region_id != version.metadata.region_id && meta.sequence.is_none() {
1988                return RegionSequenceDomainBrokenSnafu {
1989                    region_id: version.metadata.region_id,
1990                    file_region_id: meta.region_id,
1991                    file_id: meta.file_id,
1992                }
1993                .fail();
1994            }
1995            if meta.region_id == version.metadata.region_id
1996                && !file.is_effective_target_sequence_trusted(version.metadata.region_id)
1997                && meta.sequence.is_none_or(|barrier| barrier.get() > min)
1998            {
1999                // Match the capability helper's short-circuit behavior: once
2000                // an unadmitted local file is found, later files cannot change
2001                // the result or expose a foreign-domain error.
2002                files_allow_exact_range = false;
2003                check_capability = false;
2004            }
2005        }
2006        files.push(file.clone());
2007    }
2008
2009    let sequence_range = match (
2010        request.exact_sequence_range,
2011        min,
2012        version.options.preserve_row_sequence,
2013        request.memtable_max_sequence,
2014        files_allow_exact_range,
2015    ) {
2016        (true, Some(min), true, Some(max), true) => Some(SequenceRange::GtLtEq { min, max }),
2017        _ => None,
2018    };
2019    Ok((files, sequence_range))
2020}
2021
2022/// Context shared by different streams from a scanner.
2023/// It contains the input and ranges to scan.
2024pub struct StreamContext {
2025    /// Input memtables and files.
2026    pub input: ScanInput,
2027    /// Metadata for partition ranges.
2028    pub(crate) ranges: Vec<RangeMeta>,
2029    // Metrics:
2030    /// The start time of the query.
2031    pub(crate) query_start: Instant,
2032}
2033
2034impl StreamContext {
2035    /// Creates a new [StreamContext] for [SeqScan].
2036    pub(crate) fn seq_scan_ctx(input: ScanInput) -> Self {
2037        let query_start = input.query_start.unwrap_or_else(Instant::now);
2038        let ranges = RangeMeta::seq_scan_ranges(&input);
2039        READ_SST_COUNT.observe(input.num_files() as f64);
2040        Self {
2041            input,
2042            ranges,
2043            query_start,
2044        }
2045    }
2046
2047    /// Creates a new [StreamContext] for [UnorderedScan].
2048    pub(crate) fn unordered_scan_ctx(input: ScanInput) -> Self {
2049        let query_start = input.query_start.unwrap_or_else(Instant::now);
2050        let ranges = RangeMeta::unordered_scan_ranges(&input);
2051        READ_SST_COUNT.observe(input.num_files() as f64);
2052        Self {
2053            input,
2054            ranges,
2055            query_start,
2056        }
2057    }
2058
2059    /// Returns true if the index refers to a memtable.
2060    pub(crate) fn is_mem_range_index(&self, index: RowGroupIndex) -> bool {
2061        self.input.num_memtables() > index.index
2062    }
2063
2064    pub(crate) fn is_file_range_index(&self, index: RowGroupIndex) -> bool {
2065        !self.is_mem_range_index(index)
2066            && index.index < self.input.num_files() + self.input.num_memtables()
2067    }
2068
2069    pub(crate) fn range_pre_filter_mode(&self, part_range: &PartitionRange) -> PreFilterMode {
2070        let range_meta = &self.ranges[part_range.identifier];
2071        let source_count = range_meta.indices.len();
2072
2073        self.input.range_pre_filter_mode(source_count)
2074    }
2075
2076    /// Retrieves the partition ranges.
2077    pub(crate) fn partition_ranges(&self) -> Vec<PartitionRange> {
2078        self.ranges
2079            .iter()
2080            .enumerate()
2081            .map(|(idx, range_meta)| range_meta.new_partition_range(idx))
2082            .collect()
2083    }
2084
2085    /// Format the context for explain.
2086    pub(crate) fn format_for_explain(&self, verbose: bool, f: &mut fmt::Formatter) -> fmt::Result {
2087        let (mut num_mem_ranges, mut num_file_ranges, mut num_other_ranges) = (0, 0, 0);
2088        for range_meta in &self.ranges {
2089            for idx in &range_meta.row_group_indices {
2090                if self.is_mem_range_index(*idx) {
2091                    num_mem_ranges += 1;
2092                } else if self.is_file_range_index(*idx) {
2093                    num_file_ranges += 1;
2094                } else {
2095                    num_other_ranges += 1;
2096                }
2097            }
2098        }
2099        if verbose {
2100            write!(f, "{{")?;
2101        }
2102        write!(
2103            f,
2104            r#""partition_count":{{"count":{}, "mem_ranges":{}, "files":{}, "file_ranges":{}"#,
2105            self.ranges.len(),
2106            num_mem_ranges,
2107            self.input.num_files(),
2108            num_file_ranges,
2109        )?;
2110        if num_other_ranges > 0 {
2111            write!(f, r#", "other_ranges":{}"#, num_other_ranges)?;
2112        }
2113        write!(f, "}}")?;
2114
2115        if let Some(selector) = &self.input.series_row_selector {
2116            write!(f, ", \"selector\":\"{}\"", selector)?;
2117        }
2118        if let Some(distribution) = &self.input.distribution {
2119            write!(f, ", \"distribution\":\"{}\"", distribution)?;
2120        }
2121
2122        if verbose {
2123            self.format_verbose_content(f)?;
2124        }
2125
2126        Ok(())
2127    }
2128
2129    fn format_verbose_content(&self, f: &mut fmt::Formatter) -> fmt::Result {
2130        struct FileWrapper<'a> {
2131            file: &'a FileHandle,
2132        }
2133
2134        impl fmt::Debug for FileWrapper<'_> {
2135            fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2136                let (start, end) = self.file.time_range();
2137                write!(
2138                    f,
2139                    r#"{{"file_id":"{}","time_range_start":"{}::{}","time_range_end":"{}::{}","rows":{},"size":{},"index_size":{}}}"#,
2140                    self.file.file_id(),
2141                    start.value(),
2142                    start.unit(),
2143                    end.value(),
2144                    end.unit(),
2145                    self.file.num_rows(),
2146                    self.file.size(),
2147                    self.file.index_size()
2148                )
2149            }
2150        }
2151
2152        struct InputWrapper<'a> {
2153            input: &'a ScanInput,
2154        }
2155
2156        #[cfg(feature = "enterprise")]
2157        impl InputWrapper<'_> {
2158            fn format_extension_ranges(&self, f: &mut fmt::Formatter) -> fmt::Result {
2159                if self.input.extension_ranges.is_empty() {
2160                    return Ok(());
2161                }
2162
2163                let mut delimiter = "";
2164                write!(f, ", extension_ranges: [")?;
2165                for range in self.input.extension_ranges() {
2166                    write!(f, "{}{:?}", delimiter, range)?;
2167                    delimiter = ", ";
2168                }
2169                write!(f, "]")?;
2170                Ok(())
2171            }
2172        }
2173
2174        impl fmt::Debug for InputWrapper<'_> {
2175            fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
2176                let output_schema = self.input.mapper.output_schema();
2177                if !output_schema.is_empty() {
2178                    let names: Vec<_> = output_schema
2179                        .column_schemas()
2180                        .iter()
2181                        .map(|col| &col.name)
2182                        .collect();
2183                    write!(f, ", \"projection\": {:?}", names)?;
2184                }
2185                if let Some(predicate) = &self.input.predicate.predicate() {
2186                    if !predicate.exprs().is_empty() {
2187                        let exprs: Vec<_> =
2188                            predicate.exprs().iter().map(|e| e.to_string()).collect();
2189                        write!(f, ", \"filters\": {:?}", exprs)?;
2190                    }
2191                    if !predicate.dyn_filters().is_empty() {
2192                        let dyn_filters: Vec<_> = predicate
2193                            .dyn_filters()
2194                            .iter()
2195                            .map(|f| format!("{}", f))
2196                            .collect();
2197                        write!(f, ", \"dyn_filters\": {:?}", dyn_filters)?;
2198                    }
2199                }
2200                #[cfg(feature = "vector_index")]
2201                if let Some(vector_index_k) = self.input.vector_index_k {
2202                    write!(f, ", \"vector_index_k\": {}", vector_index_k)?;
2203                }
2204                if !self.input.files.is_empty() {
2205                    write!(f, ", \"files\": ")?;
2206                    f.debug_list()
2207                        .entries(self.input.files.iter().map(|file| FileWrapper { file }))
2208                        .finish()?;
2209                }
2210                write!(f, ", \"flat_format\": {}", self.input.explain_flat_format)?;
2211                #[cfg(feature = "enterprise")]
2212                self.format_extension_ranges(f)?;
2213
2214                Ok(())
2215            }
2216        }
2217
2218        write!(f, "{:?}", InputWrapper { input: &self.input })
2219    }
2220
2221    /// Add new dynamic filters to the predicates.
2222    /// Safe after stream creation; in-flight reads may still observe an older snapshot.
2223    pub(crate) fn add_dyn_filter_to_predicate(
2224        self: &Arc<Self>,
2225        filter_exprs: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
2226    ) -> Vec<bool> {
2227        let mut supported = Vec::with_capacity(filter_exprs.len());
2228        let filter_expr = filter_exprs
2229            .into_iter()
2230            .filter_map(|expr| {
2231                if let Ok(dyn_filter) = (expr as Arc<dyn std::any::Any + Send + Sync + 'static>)
2232                .downcast::<datafusion::physical_plan::expressions::DynamicFilterPhysicalExpr>()
2233            {
2234                supported.push(true);
2235                Some(dyn_filter)
2236            } else {
2237                supported.push(false);
2238                None
2239            }
2240            })
2241            .collect();
2242        self.input.predicate.add_dyn_filters(filter_expr);
2243        supported
2244    }
2245}
2246
2247/// Predicates to evaluate.
2248/// It only keeps filters that [SimpleFilterEvaluator] supports.
2249#[derive(Clone, Default)]
2250pub struct PredicateGroup {
2251    time_filters: Option<Arc<Vec<SimpleFilterEvaluator>>>,
2252    /// Predicate that includes request filters and region partition expr (if any).
2253    predicate_all: Predicate,
2254    /// Predicate that only includes request filters.
2255    predicate_without_region: Predicate,
2256    /// Region partition expression restored from metadata.
2257    region_partition_expr: Option<PartitionExpr>,
2258}
2259
2260impl PredicateGroup {
2261    /// Creates a new `PredicateGroup` from exprs according to the metadata.
2262    pub fn new(metadata: &RegionMetadata, exprs: &[Expr]) -> Result<Self> {
2263        let mut combined_exprs = exprs.to_vec();
2264        let mut region_partition_expr = None;
2265
2266        if let Some(expr_json) = metadata.partition_expr.as_ref()
2267            && !expr_json.is_empty()
2268            && let Some(expr) = PartitionExpr::from_json_str(expr_json)
2269                .context(InvalidPartitionExprSnafu { expr: expr_json })?
2270        {
2271            let logical_expr = expr
2272                .try_as_logical_expr()
2273                .context(InvalidPartitionExprSnafu {
2274                    expr: expr_json.clone(),
2275                })?;
2276
2277            combined_exprs.push(logical_expr);
2278            region_partition_expr = Some(expr);
2279        }
2280
2281        let mut time_filters = Vec::with_capacity(combined_exprs.len());
2282        // Columns in the expr.
2283        let mut columns = HashSet::new();
2284        for expr in &combined_exprs {
2285            columns.clear();
2286            let Some(filter) = Self::expr_to_filter(expr, metadata, &mut columns) else {
2287                continue;
2288            };
2289            time_filters.push(filter);
2290        }
2291        let time_filters = if time_filters.is_empty() {
2292            None
2293        } else {
2294            Some(Arc::new(time_filters))
2295        };
2296
2297        let predicate_all = Predicate::new(combined_exprs);
2298        let predicate_without_region = Predicate::new(exprs.to_vec());
2299
2300        Ok(Self {
2301            time_filters,
2302            predicate_all,
2303            predicate_without_region,
2304            region_partition_expr,
2305        })
2306    }
2307
2308    /// Returns time filters.
2309    pub(crate) fn time_filters(&self) -> Option<Arc<Vec<SimpleFilterEvaluator>>> {
2310        self.time_filters.clone()
2311    }
2312
2313    /// Returns predicate of all exprs (including region partition expr if present).
2314    pub(crate) fn predicate(&self) -> Option<&Predicate> {
2315        if self.predicate_all.is_empty() {
2316            None
2317        } else {
2318            Some(&self.predicate_all)
2319        }
2320    }
2321
2322    /// Returns predicate that excludes region partition expr.
2323    pub(crate) fn predicate_without_region(&self) -> Option<&Predicate> {
2324        if self.predicate_without_region.is_empty() {
2325            None
2326        } else {
2327            Some(&self.predicate_without_region)
2328        }
2329    }
2330
2331    /// Add dynamic filters in the predicates.
2332    pub(crate) fn add_dyn_filters(&self, dyn_filters: Vec<Arc<DynamicFilterPhysicalExpr>>) {
2333        self.predicate_all.add_dyn_filters(dyn_filters.clone());
2334        self.predicate_without_region.add_dyn_filters(dyn_filters);
2335    }
2336
2337    /// Removes dynamic filters while preserving the static and region predicates.
2338    pub(crate) fn clear_dyn_filters(&self) {
2339        self.predicate_all.clear_dyn_filters();
2340        self.predicate_without_region.clear_dyn_filters();
2341    }
2342
2343    /// Returns the region partition expr from metadata, if any.
2344    pub(crate) fn region_partition_expr(&self) -> Option<&PartitionExpr> {
2345        self.region_partition_expr.as_ref()
2346    }
2347
2348    fn expr_to_filter(
2349        expr: &Expr,
2350        metadata: &RegionMetadata,
2351        columns: &mut HashSet<Column>,
2352    ) -> Option<SimpleFilterEvaluator> {
2353        columns.clear();
2354        // `expr_to_columns` won't return error.
2355        // We still ignore these expressions for safety.
2356        expr_to_columns(expr, columns).ok()?;
2357        if columns.len() > 1 {
2358            // Simple filter doesn't support multiple columns.
2359            return None;
2360        }
2361        let column = columns.iter().next()?;
2362        let column_meta = metadata.column_by_name(&column.name)?;
2363        if column_meta.semantic_type == SemanticType::Timestamp {
2364            SimpleFilterEvaluator::try_new(expr)
2365        } else {
2366            None
2367        }
2368    }
2369}
2370
2371#[cfg(test)]
2372mod tests {
2373    use std::collections::BTreeMap;
2374    use std::sync::Arc;
2375
2376    use common_time::timestamp::{TimeUnit, Timestamp};
2377    use datafusion::physical_plan::expressions::{
2378        binary as physical_binary, col as physical_col, lit as physical_lit,
2379    };
2380    use datafusion_common::ScalarValue;
2381    use datafusion_expr::{Operator, col, lit};
2382    use datatypes::arrow::datatypes::{
2383        DataType as ArrowDataType, Field, Schema as ArrowSchema, TimeUnit as ArrowTimeUnit,
2384    };
2385    use datatypes::prelude::ConcreteDataType;
2386    use datatypes::schema::ColumnSchema;
2387    use datatypes::types::json_type::JsonObjectType;
2388    use datatypes::value::Value;
2389    use partition::expr::col as partition_col;
2390    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
2391    use store_api::storage::{RegionId, TimeSeriesDistribution, TimeSeriesRowSelector};
2392
2393    use super::*;
2394    use crate::cache::CacheManager;
2395    use crate::read::range_cache::ScanRequestFingerprintBuilder;
2396    use crate::sst::file::FileMeta;
2397    use crate::test_util::memtable_util::metadata_with_primary_key;
2398    use crate::test_util::scheduler_util::SchedulerEnv;
2399
2400    async fn new_scan_input(metadata: RegionMetadataRef, filters: Vec<Expr>) -> ScanInputBuilder {
2401        let env = SchedulerEnv::new().await;
2402        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2403        let predicate = PredicateGroup::new(metadata.as_ref(), &filters).unwrap();
2404        let file = FileHandle::new(
2405            crate::sst::file::FileMeta::default(),
2406            Arc::new(crate::sst::file_purger::NoopFilePurger),
2407        );
2408
2409        ScanInput::builder(env.access_layer.clone(), mapper)
2410            .with_predicate(predicate)
2411            .with_cache(CacheStrategy::EnableAll(Arc::new(
2412                CacheManager::builder()
2413                    .range_result_cache_size(1024)
2414                    .build(),
2415            )))
2416            .with_files(vec![file])
2417    }
2418
2419    #[tokio::test]
2420    async fn test_total_rows_is_exact_after_partition_filter() {
2421        let expr = partition_col("k0").gt_eq(Value::String("foo".into()));
2422        let other = partition_col("k0").gt_eq(Value::String("bar".into()));
2423        for (region_expr, file_exprs, exact) in [
2424            (None, vec![None], true),
2425            (None, vec![Some(expr.clone())], true),
2426            (Some(expr.clone()), vec![], true),
2427            (Some(expr.clone()), vec![Some(expr.clone())], true),
2428            (Some(expr.clone()), vec![None], false),
2429            (Some(expr.clone()), vec![Some(other)], false),
2430            (Some(expr.clone()), vec![Some(expr), None], false),
2431        ] {
2432            let mut builder =
2433                RegionMetadataBuilder::from_existing(metadata_with_primary_key(vec![0, 1], false));
2434            builder.partition_expr_json(region_expr.map(|expr| expr.as_json_str().unwrap()));
2435            let metadata = Arc::new(builder.build_without_validation().unwrap());
2436            let files = file_exprs
2437                .into_iter()
2438                .map(|partition_expr| {
2439                    FileHandle::new(
2440                        FileMeta {
2441                            partition_expr,
2442                            ..Default::default()
2443                        },
2444                        Arc::new(crate::sst::file_purger::NoopFilePurger),
2445                    )
2446                })
2447                .collect();
2448            let input = new_scan_input(metadata, vec![])
2449                .await
2450                .with_files(files)
2451                .build();
2452            assert_eq!(input.total_rows_is_exact(), exact);
2453        }
2454    }
2455
2456    /// Helper to create a timestamp millisecond literal.
2457    fn ts_lit(val: i64) -> datafusion_expr::Expr {
2458        lit(ScalarValue::TimestampMillisecond(Some(val), None))
2459    }
2460
2461    fn metadata_with_time_index_unit(unit: TimeUnit) -> RegionMetadataRef {
2462        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2463        builder
2464            .push_column_metadata(ColumnMetadata {
2465                column_schema: ColumnSchema::new(
2466                    "k0".to_string(),
2467                    ConcreteDataType::string_datatype(),
2468                    false,
2469                ),
2470                semantic_type: SemanticType::Tag,
2471                column_id: 0,
2472            })
2473            .push_column_metadata(ColumnMetadata {
2474                column_schema: ColumnSchema::new(
2475                    "k1".to_string(),
2476                    ConcreteDataType::uint32_datatype(),
2477                    false,
2478                ),
2479                semantic_type: SemanticType::Tag,
2480                column_id: 1,
2481            })
2482            .push_column_metadata(ColumnMetadata {
2483                column_schema: ColumnSchema::new(
2484                    "ts".to_string(),
2485                    ConcreteDataType::timestamp_datatype(unit),
2486                    false,
2487                ),
2488                semantic_type: SemanticType::Timestamp,
2489                column_id: 2,
2490            })
2491            .push_column_metadata(ColumnMetadata {
2492                column_schema: ColumnSchema::new(
2493                    "v0".to_string(),
2494                    ConcreteDataType::int64_datatype(),
2495                    true,
2496                ),
2497                semantic_type: SemanticType::Field,
2498                column_id: 3,
2499            })
2500            .primary_key(vec![0, 1]);
2501
2502        Arc::new(builder.build().unwrap())
2503    }
2504
2505    fn file_handle_with_time_range(start: Timestamp, end: Timestamp) -> FileHandle {
2506        FileHandle::new(
2507            FileMeta {
2508                time_range: (start, end),
2509                ..Default::default()
2510            },
2511            Arc::new(crate::sst::file_purger::NoopFilePurger),
2512        )
2513    }
2514
2515    #[test]
2516    fn test_time_range_covers_file() {
2517        let file = file_handle_with_time_range(
2518            Timestamp::new_millisecond(1000),
2519            Timestamp::new_millisecond(2000),
2520        );
2521
2522        assert!(!time_range_covers_file(None, &file));
2523        assert!(time_range_covers_file(
2524            Some(&TimestampRange::min_to_max()),
2525            &file
2526        ));
2527        assert!(time_range_covers_file(
2528            Some(&TimestampRange::new_inclusive(
2529                Some(Timestamp::new_millisecond(1000)),
2530                Some(Timestamp::new_millisecond(2000)),
2531            )),
2532            &file
2533        ));
2534        assert!(time_range_covers_file(
2535            TimestampRange::with_unit(500, 3000, TimeUnit::Millisecond).as_ref(),
2536            &file
2537        ));
2538        assert!(!time_range_covers_file(
2539            TimestampRange::with_unit(1000, 2000, TimeUnit::Millisecond).as_ref(),
2540            &file
2541        ));
2542        assert!(!time_range_covers_file(
2543            TimestampRange::with_unit(1001, 3000, TimeUnit::Millisecond).as_ref(),
2544            &file
2545        ));
2546        assert!(!time_range_covers_file(
2547            Some(&TimestampRange::empty()),
2548            &file
2549        ));
2550
2551        let seconds_file = file_handle_with_time_range(
2552            Timestamp::new(1, TimeUnit::Second),
2553            Timestamp::new(2, TimeUnit::Second),
2554        );
2555        assert!(time_range_covers_file(
2556            TimestampRange::with_unit(1000, 2001, TimeUnit::Millisecond).as_ref(),
2557            &seconds_file
2558        ));
2559    }
2560
2561    #[tokio::test]
2562    async fn test_scan_input_builder_computes_scan_analysis() {
2563        let metadata = metadata_with_time_index_unit(TimeUnit::Millisecond);
2564        let env = SchedulerEnv::new().await;
2565        let file = file_handle_with_time_range(
2566            Timestamp::new_millisecond(1000),
2567            Timestamp::new_millisecond(2000),
2568        );
2569
2570        let make_input = |filters: &[Expr]| {
2571            let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2572            ScanInput::builder(env.access_layer.clone(), mapper)
2573                .with_predicate(PredicateGroup::new(&metadata, filters).unwrap())
2574                .build()
2575        };
2576
2577        let covered = make_input(&[col("ts").gt_eq(ts_lit(500)), col("ts").lt(ts_lit(3000))]);
2578        assert!(time_range_covers_file(covered.implied_time_range(), &file));
2579
2580        // The convex hull of this disjunction covers the file, but most rows do
2581        // not satisfy it. The strict implied-range analysis must reject it.
2582        let disjoint = make_input(&[col("ts").eq(ts_lit(1000)).or(col("ts").eq(ts_lit(2000)))]);
2583        assert!(disjoint.implied_time_range().is_none());
2584        assert!(!time_range_covers_file(
2585            disjoint.implied_time_range(),
2586            &file
2587        ));
2588
2589        let replaced = ScanInput::builder(
2590            env.access_layer.clone(),
2591            FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap(),
2592        )
2593        .with_predicate(PredicateGroup::new(&metadata, &[col("ts").gt(ts_lit(3000))]).unwrap())
2594        .build();
2595        assert!(!time_range_covers_file(
2596            replaced.implied_time_range(),
2597            &file
2598        ));
2599    }
2600
2601    #[tokio::test]
2602    async fn test_scan_input_uses_explicit_batch_size() {
2603        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2604        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2605        let env = SchedulerEnv::new().await;
2606        let input = ScanInput::builder(env.access_layer.clone(), mapper).build();
2607        assert_eq!(
2608            crate::sst::parquet::DEFAULT_READ_BATCH_SIZE,
2609            input.batch_size()
2610        );
2611
2612        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2613        let input = ScanInput::builder(env.access_layer.clone(), mapper)
2614            .with_compaction(true)
2615            .with_batch_size(256)
2616            .build();
2617        assert_eq!(256, input.batch_size());
2618
2619        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
2620        let input = ScanInput::builder(env.access_layer.clone(), mapper)
2621            .with_batch_size(256)
2622            .with_compaction(true)
2623            .with_compaction(false)
2624            .build();
2625        assert_eq!(256, input.batch_size());
2626    }
2627
2628    #[tokio::test]
2629    async fn test_build_scan_fingerprint_for_eligible_scan() {
2630        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2631        let input = new_scan_input(
2632            metadata.clone(),
2633            vec![
2634                col("ts").gt_eq(ts_lit(1000)),
2635                col("k0").eq(lit("foo")),
2636                col("v0").gt(lit(1)),
2637            ],
2638        )
2639        .await
2640        .with_distribution(Some(TimeSeriesDistribution::PerSeries))
2641        .with_series_row_selector(Some(TimeSeriesRowSelector::LastRow { after_merge: true }))
2642        .with_merge_mode(MergeMode::LastNonNull)
2643        .with_filter_deleted(false)
2644        .build();
2645
2646        let fingerprint = input.scan_fingerprint().unwrap();
2647
2648        let expected = ScanRequestFingerprintBuilder {
2649            read_columns: input.read_cols.clone(),
2650            read_column_types: vec![
2651                metadata
2652                    .column_by_id(0)
2653                    .map(|col| col.column_schema.data_type.clone()),
2654                metadata
2655                    .column_by_id(2)
2656                    .map(|col| col.column_schema.data_type.clone()),
2657                metadata
2658                    .column_by_id(3)
2659                    .map(|col| col.column_schema.data_type.clone()),
2660            ],
2661            filters: vec![
2662                col("k0").eq(lit("foo")).to_string(),
2663                col("v0").gt(lit(1)).to_string(),
2664            ],
2665            time_filters: vec![col("ts").gt_eq(ts_lit(1000)).to_string()],
2666            series_row_selector: Some(TimeSeriesRowSelector::LastRow { after_merge: true }),
2667            append_mode: false,
2668            filter_deleted: false,
2669            merge_mode: MergeMode::LastNonNull,
2670            sequence_range: None,
2671            partition_expr_version: 0,
2672        }
2673        .build();
2674        assert_eq!(&expected, fingerprint);
2675        assert_eq!(
2676            input.series_row_selector,
2677            Some(TimeSeriesRowSelector::LastRow { after_merge: true })
2678        );
2679    }
2680
2681    #[tokio::test]
2682    async fn test_build_scan_fingerprint_requires_tag_filter() {
2683        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2684        let input = new_scan_input(
2685            metadata,
2686            vec![col("ts").gt_eq(lit(1000)), col("v0").gt(lit(1))],
2687        )
2688        .await
2689        .build();
2690
2691        assert!(input.scan_fingerprint().is_none());
2692    }
2693
2694    #[tokio::test]
2695    async fn test_build_scan_fingerprint_respects_scan_eligibility() {
2696        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2697        let filters = vec![col("k0").eq(lit("foo"))];
2698
2699        let disabled = ScanInput::builder(
2700            SchedulerEnv::new().await.access_layer.clone(),
2701            FlatProjectionMapper::new(&metadata, [0, 2, 3].into_iter()).unwrap(),
2702        )
2703        .with_predicate(PredicateGroup::new(metadata.as_ref(), &filters).unwrap())
2704        .build();
2705        assert!(disabled.scan_fingerprint().is_none());
2706
2707        let compaction = new_scan_input(metadata.clone(), filters.clone())
2708            .await
2709            .with_compaction(true)
2710            .build();
2711        assert!(compaction.scan_fingerprint().is_none());
2712
2713        // No files to read.
2714        let no_files = new_scan_input(metadata, filters)
2715            .await
2716            .with_files(vec![])
2717            .build();
2718        assert!(no_files.scan_fingerprint().is_none());
2719    }
2720
2721    #[tokio::test]
2722    async fn test_build_scan_fingerprint_tracks_schema_and_partition_expr_changes() {
2723        let base = metadata_with_primary_key(vec![0, 1], false);
2724        let mut builder = RegionMetadataBuilder::from_existing(base);
2725        let partition_expr = partition_col("k0")
2726            .gt_eq(Value::String("foo".into()))
2727            .as_json_str()
2728            .unwrap();
2729        builder.partition_expr_json(Some(partition_expr));
2730        let metadata = Arc::new(builder.build_without_validation().unwrap());
2731
2732        let input = new_scan_input(metadata.clone(), vec![col("k0").eq(lit("foo"))])
2733            .await
2734            .build();
2735        let fingerprint = input.scan_fingerprint().unwrap();
2736
2737        let expected = ScanRequestFingerprintBuilder {
2738            read_columns: input.read_cols.clone(),
2739            read_column_types: vec![
2740                metadata
2741                    .column_by_id(0)
2742                    .map(|col| col.column_schema.data_type.clone()),
2743                metadata
2744                    .column_by_id(2)
2745                    .map(|col| col.column_schema.data_type.clone()),
2746                metadata
2747                    .column_by_id(3)
2748                    .map(|col| col.column_schema.data_type.clone()),
2749            ],
2750            filters: vec![col("k0").eq(lit("foo")).to_string()],
2751            time_filters: vec![],
2752            series_row_selector: None,
2753            append_mode: false,
2754            filter_deleted: true,
2755            merge_mode: MergeMode::LastRow,
2756            sequence_range: None,
2757            partition_expr_version: metadata.partition_expr_version,
2758        }
2759        .build();
2760        assert_eq!(&expected, fingerprint);
2761        assert_ne!(0, metadata.partition_expr_version);
2762    }
2763
2764    #[tokio::test]
2765    async fn test_build_scan_fingerprint_uses_json_target_types() {
2766        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2767        builder
2768            .push_column_metadata(ColumnMetadata {
2769                column_schema: ColumnSchema::new(
2770                    "k0".to_string(),
2771                    ConcreteDataType::string_datatype(),
2772                    false,
2773                ),
2774                semantic_type: SemanticType::Tag,
2775                column_id: 0,
2776            })
2777            .push_column_metadata(ColumnMetadata {
2778                column_schema: ColumnSchema::new(
2779                    "ts".to_string(),
2780                    ConcreteDataType::timestamp_millisecond_datatype(),
2781                    false,
2782                ),
2783                semantic_type: SemanticType::Timestamp,
2784                column_id: 1,
2785            })
2786            .push_column_metadata(ColumnMetadata {
2787                column_schema: ColumnSchema::new(
2788                    "j".to_string(),
2789                    ConcreteDataType::json2(JsonNativeType::Variant),
2790                    true,
2791                ),
2792                semantic_type: SemanticType::Field,
2793                column_id: 2,
2794            })
2795            .primary_key(vec![0]);
2796        let metadata = Arc::new(builder.build().unwrap());
2797
2798        let make_input = |target_type| async {
2799            let env = SchedulerEnv::new().await;
2800            let read_cols = ReadColumns::new([0, 1, 2])
2801                .with_json_target_types(BTreeMap::from([(2, target_type)]));
2802            let mapper =
2803                FlatProjectionMapper::new_with_read_columns(&metadata, vec![0, 1, 2], read_cols)
2804                    .unwrap();
2805            let predicate =
2806                PredicateGroup::new(metadata.as_ref(), &[col("k0").eq(lit("foo"))]).unwrap();
2807            let file = FileHandle::new(
2808                FileMeta::default(),
2809                Arc::new(crate::sst::file_purger::NoopFilePurger),
2810            );
2811            ScanInput::builder(env.access_layer.clone(), mapper)
2812                .with_predicate(predicate)
2813                .with_cache(CacheStrategy::EnableAll(Arc::new(
2814                    CacheManager::builder()
2815                        .range_result_cache_size(1024)
2816                        .build(),
2817                )))
2818                .with_files(vec![file])
2819                .build()
2820        };
2821
2822        let int_target = JsonNativeType::i64();
2823        let string_target = JsonNativeType::String;
2824        let int_input = make_input(int_target.clone()).await;
2825        let string_input = make_input(string_target).await;
2826        let int_fingerprint = int_input.scan_fingerprint().unwrap();
2827        let string_fingerprint = string_input.scan_fingerprint().unwrap();
2828
2829        assert_ne!(int_fingerprint, string_fingerprint);
2830        assert_eq!(
2831            Some(&Some(ConcreteDataType::json2(int_target))),
2832            int_fingerprint.read_column_types().get(2)
2833        );
2834    }
2835
2836    #[tokio::test]
2837    async fn test_scan_input_rejects_json_type_hint_for_non_json2_column() {
2838        let mut builder = RegionMetadataBuilder::new(RegionId::new(123, 456));
2839        builder
2840            .push_column_metadata(ColumnMetadata {
2841                column_schema: ColumnSchema::new(
2842                    "k0".to_string(),
2843                    ConcreteDataType::string_datatype(),
2844                    false,
2845                ),
2846                semantic_type: SemanticType::Tag,
2847                column_id: 0,
2848            })
2849            .push_column_metadata(ColumnMetadata {
2850                column_schema: ColumnSchema::new(
2851                    "ts".to_string(),
2852                    ConcreteDataType::timestamp_millisecond_datatype(),
2853                    false,
2854                ),
2855                semantic_type: SemanticType::Timestamp,
2856                column_id: 1,
2857            })
2858            .push_column_metadata(ColumnMetadata {
2859                column_schema: ColumnSchema::new(
2860                    "j".to_string(),
2861                    ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::from([(
2862                        "a".to_string(),
2863                        JsonNativeType::i64(),
2864                    )]))),
2865                    true,
2866                ),
2867                semantic_type: SemanticType::Field,
2868                column_id: 2,
2869            })
2870            .push_column_metadata(ColumnMetadata {
2871                column_schema: ColumnSchema::new(
2872                    "v0".to_string(),
2873                    ConcreteDataType::int64_datatype(),
2874                    true,
2875                ),
2876                semantic_type: SemanticType::Field,
2877                column_id: 3,
2878            })
2879            .primary_key(vec![0]);
2880        let metadata = Arc::new(builder.build().unwrap());
2881        let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
2882            metadata.clone(),
2883            Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
2884            0,
2885            None,
2886        ));
2887        let version = Arc::new(
2888            crate::region::version::VersionBuilder::new(metadata.clone(), mutable).build(),
2889        );
2890        let env = SchedulerEnv::new().await;
2891        let request = ScanRequest {
2892            projection: Some(vec![0, 1, 2, 3]),
2893            json_type_hint: std::collections::HashMap::from([(
2894                "v0".to_string(),
2895                JsonNativeType::i64(),
2896            )]),
2897            ..Default::default()
2898        };
2899
2900        let err = ScanRegion::new(
2901            version,
2902            env.access_layer.clone(),
2903            request,
2904            CacheStrategy::Disabled,
2905        )
2906        .scan_input()
2907        .await;
2908        let Err(err) = err else {
2909            panic!("scan input should reject JSON type hint for non-JSON2 column");
2910        };
2911
2912        assert!(err.to_string().contains("non-JSON2 column v0"));
2913    }
2914
2915    #[test]
2916    fn test_update_dyn_filters_with_empty_base_predicates() {
2917        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2918        let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
2919        assert!(predicate_group.predicate().is_none());
2920        assert!(predicate_group.predicate_without_region().is_none());
2921
2922        let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(vec![], physical_lit(false)));
2923        predicate_group.add_dyn_filters(vec![dyn_filter]);
2924
2925        let predicate_all = predicate_group.predicate().unwrap();
2926        assert!(predicate_all.exprs().is_empty());
2927        assert_eq!(1, predicate_all.dyn_filters().len());
2928
2929        let predicate_without_region = predicate_group.predicate_without_region().unwrap();
2930        assert!(predicate_without_region.exprs().is_empty());
2931        assert_eq!(1, predicate_without_region.dyn_filters().len());
2932    }
2933
2934    #[test]
2935    fn test_clear_dyn_filters_preserves_predicate_group_static_filters() {
2936        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
2937        let static_filters = vec![col("k0").eq(lit("foo"))];
2938        let predicate_group = PredicateGroup::new(metadata.as_ref(), &static_filters).unwrap();
2939        let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new(vec![], physical_lit(true)));
2940        predicate_group.add_dyn_filters(vec![dynamic_filter.clone()]);
2941
2942        predicate_group.clear_dyn_filters();
2943        // Updating the old producer cannot add its wrapper to a new execution.
2944        dynamic_filter.update(physical_lit(false)).unwrap();
2945
2946        for predicate in [
2947            predicate_group.predicate().unwrap(),
2948            predicate_group.predicate_without_region().unwrap(),
2949        ] {
2950            assert_eq!(predicate.exprs(), static_filters);
2951            assert!(predicate.dyn_filters().is_empty());
2952        }
2953    }
2954
2955    #[test]
2956    fn test_file_level_pruning_stats_prunes_old_file() {
2957        let ts_col_name = "ts";
2958        let predicate = Predicate::new(vec![col(ts_col_name).gt(ts_lit(1000))]);
2959        let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
2960            ts_col_name,
2961            ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
2962            false,
2963        )]));
2964
2965        // File with time range [0ms, 500ms] is completely before `ts > 1000ms`.
2966        let stats = FileLevelPruningStats {
2967            min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2968            max_scalar: ScalarValue::TimestampMillisecond(Some(500), None),
2969            time_index_col_name: ts_col_name.to_string(),
2970        };
2971        assert_eq!(
2972            vec![false],
2973            predicate.prune_with_stats(&stats, &arrow_schema)
2974        );
2975
2976        // File with time range [0ms, 2000ms] overlaps `ts > 1000ms`, so keep it.
2977        let stats = FileLevelPruningStats {
2978            min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2979            max_scalar: ScalarValue::TimestampMillisecond(Some(2000), None),
2980            time_index_col_name: ts_col_name.to_string(),
2981        };
2982        assert_eq!(
2983            vec![true],
2984            predicate.prune_with_stats(&stats, &arrow_schema)
2985        );
2986    }
2987
2988    #[test]
2989    fn test_file_level_pruning_stats_no_predicate_keeps_all() {
2990        let predicate = Predicate::new(vec![]);
2991        assert!(predicate.is_empty());
2992
2993        let stats = FileLevelPruningStats {
2994            min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
2995            max_scalar: ScalarValue::TimestampMillisecond(Some(500), None),
2996            time_index_col_name: "ts".to_string(),
2997        };
2998        let arrow_schema = Arc::new(ArrowSchema::new(Vec::<Field>::new()));
2999        assert_eq!(
3000            vec![true],
3001            predicate.prune_with_stats(&stats, &arrow_schema)
3002        );
3003    }
3004
3005    #[tokio::test]
3006    async fn test_file_level_pruning_stats_ceil_max_unit_conversion() {
3007        let metadata = metadata_with_time_index_unit(TimeUnit::Millisecond);
3008        let input = new_scan_input(metadata, vec![]).await.build();
3009        let file = file_handle_with_time_range(
3010            Timestamp::new(1_000_001, TimeUnit::Nanosecond),
3011            Timestamp::new(1_000_001, TimeUnit::Nanosecond),
3012        );
3013
3014        let stats = input.try_file_level_pruning_stats(&file).unwrap();
3015        assert_eq!(
3016            ScalarValue::TimestampMillisecond(Some(1), None),
3017            stats.min_scalar
3018        );
3019        assert_eq!(
3020            ScalarValue::TimestampMillisecond(Some(2), None),
3021            stats.max_scalar
3022        );
3023
3024        // The actual max timestamp is slightly greater than 1ms. It must be kept for `ts > 1ms`.
3025        let predicate = Predicate::new(vec![col("ts").gt(ts_lit(1))]);
3026        assert_eq!(
3027            vec![true],
3028            predicate.prune_with_stats(&stats, input.mapper.metadata().schema.arrow_schema())
3029        );
3030    }
3031
3032    #[tokio::test]
3033    async fn test_file_level_pruning_stats_overflow_keeps_file() {
3034        let metadata = metadata_with_time_index_unit(TimeUnit::Nanosecond);
3035        let input = new_scan_input(metadata, vec![]).await.build();
3036        let file = file_handle_with_time_range(
3037            Timestamp::new(0, TimeUnit::Second),
3038            Timestamp::new(i64::MAX, TimeUnit::Second),
3039        );
3040
3041        assert!(input.try_file_level_pruning_stats(&file).is_none());
3042    }
3043
3044    #[test]
3045    fn test_file_level_pruning_stats_keeps_inclusive_boundary() {
3046        let ts_col_name = "ts";
3047        let predicate = Predicate::new(vec![col(ts_col_name).gt_eq(ts_lit(1000))]);
3048        let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
3049            ts_col_name,
3050            ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
3051            false,
3052        )]));
3053        let stats = FileLevelPruningStats {
3054            min_scalar: ScalarValue::TimestampMillisecond(Some(0), None),
3055            max_scalar: ScalarValue::TimestampMillisecond(Some(1000), None),
3056            time_index_col_name: ts_col_name.to_string(),
3057        };
3058
3059        assert_eq!(
3060            vec![true],
3061            predicate.prune_with_stats(&stats, &arrow_schema)
3062        );
3063    }
3064
3065    #[tokio::test]
3066    async fn test_file_level_pruning_with_dyn_filter_only_predicate() {
3067        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3068        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
3069        let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
3070        predicate_group.add_dyn_filters(vec![Arc::new(DynamicFilterPhysicalExpr::new(
3071            vec![],
3072            physical_lit(false),
3073        ))]);
3074        let input = ScanInput::builder(SchedulerEnv::new().await.access_layer.clone(), mapper)
3075            .with_predicate(predicate_group)
3076            .build();
3077        let file = file_handle_with_time_range(
3078            Timestamp::new_millisecond(0),
3079            Timestamp::new_millisecond(1000),
3080        );
3081        let mut reader_metrics = ReaderMetrics::default();
3082
3083        let builder = input
3084            .prune_file(&file, PreFilterMode::SkipFields, &mut reader_metrics)
3085            .await
3086            .unwrap();
3087
3088        assert_eq!(1, reader_metrics.filter_metrics.files_time_range_pruned);
3089        let mut ranges = SmallVec::new();
3090        builder.build_ranges(-1, &mut ranges);
3091        assert!(ranges.is_empty());
3092    }
3093
3094    #[tokio::test]
3095    async fn test_manifest_pruning_observes_dynamic_filter_update() {
3096        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3097        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
3098        let predicate_group = PredicateGroup::new(metadata.as_ref(), &[]).unwrap();
3099        let arrow_schema = metadata.schema.arrow_schema();
3100        let ts_expr = physical_col("ts", arrow_schema.as_ref()).unwrap();
3101        let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(
3102            vec![ts_expr.clone()],
3103            physical_lit(true),
3104        ));
3105        predicate_group.add_dyn_filters(vec![dyn_filter.clone()]);
3106        let input = ScanInput::builder(SchedulerEnv::new().await.access_layer.clone(), mapper)
3107            .with_predicate(predicate_group)
3108            .build();
3109        let file = file_handle_with_time_range(
3110            Timestamp::new_millisecond(0),
3111            Timestamp::new_millisecond(1000),
3112        );
3113
3114        assert!(!input.can_manifest_prune_file(&file));
3115
3116        let updated = physical_binary(
3117            ts_expr,
3118            Operator::Gt,
3119            physical_lit(ScalarValue::TimestampMillisecond(Some(1000), None)),
3120            arrow_schema.as_ref(),
3121        )
3122        .unwrap();
3123        dyn_filter.update(updated).unwrap();
3124
3125        assert!(input.can_manifest_prune_file(&file));
3126        input.predicate.clear_dyn_filters();
3127        assert!(!input.can_manifest_prune_file(&file));
3128    }
3129
3130    #[tokio::test]
3131    async fn test_range_pre_filter_mode() {
3132        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3133        let cases = [
3134            (true, MergeMode::LastRow, 1, PreFilterMode::All),
3135            (false, MergeMode::LastNonNull, 1, PreFilterMode::All),
3136            (false, MergeMode::LastRow, 2, PreFilterMode::SkipFields),
3137            (true, MergeMode::LastRow, 2, PreFilterMode::All),
3138        ];
3139
3140        for (append_mode, merge_mode, source_count, expected_mode) in cases {
3141            let input = new_scan_input(metadata.clone(), vec![])
3142                .await
3143                .with_append_mode(append_mode)
3144                .with_merge_mode(merge_mode)
3145                .build();
3146
3147            assert_eq!(expected_mode, input.range_pre_filter_mode(source_count));
3148        }
3149    }
3150
3151    #[test]
3152    fn test_exact_sequence_selection_preserves_range_and_fallback() {
3153        let request = ScanRequest {
3154            exact_sequence_range: true,
3155            memtable_min_sequence: Some(1),
3156            memtable_max_sequence: Some(2),
3157            ..Default::default()
3158        };
3159        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3160        let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
3161            metadata.clone(),
3162            Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
3163            0,
3164            None,
3165        ));
3166        let version = Arc::new(
3167            crate::region::version::VersionBuilder::new(metadata, mutable)
3168                .options(crate::region::options::RegionOptions {
3169                    preserve_row_sequence: true,
3170                    ..Default::default()
3171                })
3172                .build(),
3173        );
3174        let (files, range) = exact_sequence_range(&request, &version).unwrap();
3175        assert!(files.is_empty());
3176        assert_eq!(Some(SequenceRange::GtLtEq { min: 1, max: 2 }), range);
3177
3178        let skip_sst_request = ScanRequest {
3179            skip_sst_files: true,
3180            ..request
3181        };
3182        let (files, range) = exact_sequence_range(&skip_sst_request, &version).unwrap();
3183        assert!(files.is_empty());
3184        assert_eq!(None, range);
3185    }
3186
3187    #[test]
3188    fn test_exact_sequence_selection_checks_selected_files() {
3189        use std::num::NonZeroU64;
3190
3191        use crate::test_util::new_noop_file_purger;
3192
3193        let request = ScanRequest {
3194            exact_sequence_range: true,
3195            memtable_min_sequence: Some(5),
3196            memtable_max_sequence: Some(10),
3197            ..Default::default()
3198        };
3199        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
3200        let mutable = Arc::new(crate::memtable::time_partition::TimePartitions::new(
3201            metadata.clone(),
3202            Arc::new(crate::test_util::memtable_util::EmptyMemtableBuilder::default()),
3203            0,
3204            None,
3205        ));
3206        let target_region_id = metadata.region_id;
3207
3208        // Files at or below the cursor are excluded. Among the remaining files,
3209        // marked locals are exact-capable while unmarked locals deny exactness.
3210        for (sequence, marked, expected_files, expected_range) in [
3211            (NonZeroU64::new(5), false, false, true),
3212            (NonZeroU64::new(5), true, false, true),
3213            (NonZeroU64::new(100), false, true, false),
3214            (None, false, true, false),
3215            (None, true, true, true),
3216        ] {
3217            let version =
3218                crate::region::version::VersionBuilder::new(metadata.clone(), mutable.clone())
3219                    .options(crate::region::options::RegionOptions {
3220                        preserve_row_sequence: true,
3221                        ..Default::default()
3222                    })
3223                    .add_files(
3224                        new_noop_file_purger(),
3225                        [FileMeta {
3226                            region_id: target_region_id,
3227                            sequence,
3228                            preserve_row_sequence: marked,
3229                            ..Default::default()
3230                        }]
3231                        .into_iter(),
3232                    )
3233                    .build();
3234            let (files, range) = exact_sequence_range(&request, &version).unwrap();
3235            assert_eq!(expected_files, !files.is_empty());
3236            assert_eq!(expected_range, range.is_some());
3237        }
3238
3239        // Foreign files use the target-local barrier, regardless of their
3240        // source marker. A missing barrier fails closed.
3241        for (sequence, marked, expected_count, expected_error) in [
3242            (NonZeroU64::new(5), false, 0, false),
3243            (NonZeroU64::new(6), false, 1, false),
3244            (NonZeroU64::new(11), true, 1, false),
3245            (None, true, 1, true),
3246        ] {
3247            let file_id = store_api::storage::FileId::random();
3248            let version =
3249                crate::region::version::VersionBuilder::new(metadata.clone(), mutable.clone())
3250                    .options(crate::region::options::RegionOptions {
3251                        preserve_row_sequence: true,
3252                        ..Default::default()
3253                    })
3254                    .add_files(
3255                        new_noop_file_purger(),
3256                        [FileMeta {
3257                            region_id: RegionId::new(1, 2),
3258                            file_id,
3259                            sequence,
3260                            preserve_row_sequence: marked,
3261                            ..Default::default()
3262                        }]
3263                        .into_iter(),
3264                    )
3265                    .build();
3266            let result = exact_sequence_range(&request, &version);
3267            if expected_error {
3268                assert!(result.is_err());
3269            } else {
3270                let (files, range) = result.unwrap();
3271                assert_eq!(expected_count, files.len());
3272                assert!(range.is_some());
3273            }
3274        }
3275
3276        // Foreign-domain validation still takes precedence over preserve mode
3277        // and a missing upper bound once the file is selected.
3278        let request = ScanRequest {
3279            memtable_max_sequence: None,
3280            ..request
3281        };
3282        let version = crate::region::version::VersionBuilder::new(metadata, mutable)
3283            .options(crate::region::options::RegionOptions {
3284                preserve_row_sequence: false,
3285                ..Default::default()
3286            })
3287            .add_files(
3288                new_noop_file_purger(),
3289                [FileMeta {
3290                    region_id: RegionId::new(1, 2),
3291                    ..Default::default()
3292                }]
3293                .into_iter(),
3294            )
3295            .build();
3296        assert!(exact_sequence_range(&request, &version).is_err());
3297    }
3298}