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