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