Skip to main content

mito2/read/
scan_region.rs

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