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