Skip to main content

mito2/read/
pruner.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//! Pruner for parallel file pruning across scanner partitions.
16
17use std::collections::{HashMap, HashSet};
18use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
19use std::sync::{Arc, Mutex};
20use std::time::Instant;
21
22use common_telemetry::debug;
23use smallvec::SmallVec;
24use snafu::ResultExt;
25use store_api::region_engine::PartitionRange;
26use store_api::storage::FileId;
27use tokio::sync::{mpsc, oneshot};
28use uuid::Uuid;
29
30use crate::error::{PruneFileSnafu, Result};
31use crate::metrics::PRUNER_ACTIVE_BUILDERS;
32use crate::read::range::{FileRangeBuilder, RowGroupIndex};
33use crate::read::scan_region::StreamContext;
34use crate::read::scan_util::{FileScanMetrics, PartitionMetrics, new_filter_metrics};
35use crate::sst::parquet::file_range::{FileRange, PreFilterMode};
36use crate::sst::parquet::reader::ReaderMetrics;
37
38/// Number of files to pre-fetch ahead of the current position.
39const PREFETCH_COUNT: usize = 8;
40
41/// Local pruner in a partition that supports prefetching files to prune.
42pub struct PartitionPruner {
43    pruner: Arc<Pruner>,
44    /// Files to prune, in the order to scan.
45    file_indices: Vec<usize>,
46    /// Per-file pre-filter mode lookup indexed by file_index.
47    pre_filter_modes: Vec<PreFilterMode>,
48    /// Current position for tracking pre-fetch progress.
49    current_position: AtomicUsize,
50}
51
52impl PartitionPruner {
53    /// Creates a new `PartitionPruner` for the given partition ranges.
54    pub fn new(pruner: Arc<Pruner>, partition_ranges: &[PartitionRange]) -> Self {
55        let num_files = pruner.inner.stream_ctx.input.num_files();
56        let mut file_indices = Vec::with_capacity(num_files);
57        let mut pre_filter_modes = vec![PreFilterMode::SkipFields; num_files];
58        let mut dedup_set = HashSet::with_capacity(pruner.inner.stream_ctx.input.num_files());
59
60        let num_memtables = pruner.inner.stream_ctx.input.num_memtables();
61        for part_range in partition_ranges {
62            let range_meta = &pruner.inner.stream_ctx.ranges[part_range.identifier];
63            let pre_filter_mode = pruner.inner.stream_ctx.range_pre_filter_mode(part_range);
64            for row_group_index in &range_meta.row_group_indices {
65                if pruner
66                    .inner
67                    .stream_ctx
68                    .is_file_range_index(*row_group_index)
69                {
70                    let file_index = row_group_index.index - num_memtables;
71                    if dedup_set.contains(&file_index) {
72                        continue;
73                    } else {
74                        file_indices.push(file_index);
75                        pre_filter_modes[file_index] = pre_filter_mode;
76                        dedup_set.insert(file_index);
77                    }
78                }
79            }
80        }
81
82        Self {
83            pruner,
84            file_indices,
85            pre_filter_modes,
86            current_position: AtomicUsize::new(0),
87        }
88    }
89
90    /// Excludes files replaced by another candidate source from prefetching.
91    pub(crate) fn excluding_files(mut self, excluded: &HashSet<usize>) -> Self {
92        self.file_indices.retain(|index| !excluded.contains(index));
93        self
94    }
95
96    /// Gets or creates the FileRangeBuilder for a file.
97    ///
98    /// This method also triggers pre-fetching of upcoming files in the background
99    /// to improve performance by overlapping I/O with computation.
100    pub async fn build_file_ranges(
101        &self,
102        index: RowGroupIndex,
103        partition_metrics: &PartitionMetrics,
104        reader_metrics: &mut ReaderMetrics,
105    ) -> Result<SmallVec<[FileRange; 2]>> {
106        let file_index = index.index - self.pruner.inner.stream_ctx.input.num_memtables();
107        let pre_filter_mode = self.pre_filter_mode(file_index);
108
109        // Delegate to underlying Pruner
110        let ranges = self
111            .pruner
112            .build_file_ranges(index, pre_filter_mode, partition_metrics, reader_metrics)
113            .await?;
114
115        // Find position and trigger pre-fetch for upcoming files
116        if let Some(pos) = self.file_indices.iter().position(|&idx| idx == file_index) {
117            let prev_pos = self.current_position.fetch_max(pos, Ordering::Relaxed);
118            if pos > prev_pos || prev_pos == 0 {
119                self.prefetch_upcoming_files(pos, partition_metrics);
120            }
121        }
122
123        Ok(ranges)
124    }
125
126    /// Checks whether the file range at `index` can be skipped because the
127    /// current predicate definitively prunes it at manifest level (no I/O).
128    ///
129    /// Returns `true` if the range was skipped. When skipped, this method
130    /// balances the pruner's per-file reference count and merges the resulting
131    /// reader metrics into `part_metrics`.
132    ///
133    /// Uses shared per-file state so repeated row groups can skip cheaply after
134    /// the first manifest-prune decision.
135    pub fn try_skip_manifest_pruned_file_range(
136        &self,
137        index: RowGroupIndex,
138        part_metrics: &PartitionMetrics,
139    ) -> bool {
140        let Some(file_index) = self.file_index(index) else {
141            return false;
142        };
143        let mut reader_metrics = ReaderMetrics::default();
144        let pruned = self
145            .pruner
146            .inner
147            .try_mark_manifest_pruned(file_index, &mut reader_metrics);
148        if pruned {
149            self.pruner.skip_file_range(index, &mut reader_metrics);
150            part_metrics.merge_reader_metrics(&reader_metrics, None);
151        }
152        pruned
153    }
154
155    /// Pre-fetches upcoming files starting from the given position.
156    fn prefetch_upcoming_files(&self, current_pos: usize, partition_metrics: &PartitionMetrics) {
157        let start = current_pos + 1;
158        let end = (start + PREFETCH_COUNT).min(self.file_indices.len());
159
160        for i in start..end {
161            let file_index = self.file_indices[i];
162            let pre_filter_mode = self.pre_filter_mode(file_index);
163            self.pruner.get_file_builder_background(
164                file_index,
165                pre_filter_mode,
166                Some(partition_metrics.clone()),
167            );
168        }
169    }
170
171    fn pre_filter_mode(&self, file_index: usize) -> PreFilterMode {
172        self.pre_filter_modes
173            .get(file_index)
174            .copied()
175            .unwrap_or(PreFilterMode::SkipFields)
176    }
177
178    fn file_index(&self, index: RowGroupIndex) -> Option<usize> {
179        self.pruner
180            .inner
181            .stream_ctx
182            .is_file_range_index(index)
183            .then(|| index.index - self.pruner.inner.stream_ctx.input.num_memtables())
184    }
185}
186
187/// Options to create a [`Pruner`].
188#[derive(Debug, Clone, Copy)]
189pub struct PrunerOptions {
190    /// Keeps file range builders until the query-scoped pruner is dropped.
191    pub retain_builders: bool,
192    /// Whether [`FileRange`]s execute the reduced-column predicate prefilter.
193    ///
194    /// The flag is baked into the cached [`FileRangeBuilder`], so it belongs to
195    /// the whole scan. All partitions sharing this pruner must agree on it.
196    /// Two-stage series scans disable this prefilter: candidate discovery applies
197    /// tag predicates, then the data phase calls [`FileRange::precise_filter_flat`]
198    /// with tag filtering skipped so the predicates are applied exactly once.
199    pub enable_predicate_prefilter: bool,
200}
201
202impl Default for PrunerOptions {
203    fn default() -> Self {
204        Self {
205            retain_builders: false,
206            enable_predicate_prefilter: true,
207        }
208    }
209}
210
211/// A pruner that prunes files for all partitions of a scanner.
212pub struct Pruner {
213    /// Channels to send requests to workers.
214    worker_senders: Vec<mpsc::Sender<PruneRequest>>,
215    inner: Arc<PrunerInner>,
216}
217
218struct PrunerInner {
219    /// Number of worker tasks.
220    num_workers: usize,
221    /// Per-file state (indexed by file_index).
222    file_entries: Vec<Mutex<FileBuilderEntry>>,
223    /// StreamContext containing all context needed for pruning.
224    stream_ctx: Arc<StreamContext>,
225    /// Positive manifest-prune cache shared across all scan partitions.
226    ///
227    /// SAFETY: cached positives are valid because dynamic filters only tighten;
228    /// negative decisions are not cached. Reset by `add_partition_ranges()` for
229    /// each fresh batch of partition ranges.
230    manifest_pruned_files: Vec<AtomicBool>,
231    /// Keeps file range builders until the query-scoped pruner is dropped.
232    retain_builders: bool,
233    /// Whether FileRanges should execute the reduced-column predicate prefilter.
234    enable_predicate_prefilter: bool,
235}
236
237impl Drop for PrunerInner {
238    fn drop(&mut self) {
239        let active_builders = self
240            .file_entries
241            .iter_mut()
242            .map(|entry| usize::from(entry.get_mut().is_ok_and(|entry| entry.builder.is_some())))
243            .sum::<usize>();
244        PRUNER_ACTIVE_BUILDERS.sub(active_builders as i64);
245    }
246}
247
248impl PrunerInner {
249    /// Checks whether manifest-level pruning proves this file is empty given the
250    /// current predicate. If true, CAS the shared cache from false→true and
251    /// record `files_time_range_pruned` in `reader_metrics`.
252    ///
253    /// Returns `true` if already cached or newly proven pruned.
254    fn try_mark_manifest_pruned(
255        &self,
256        file_index: usize,
257        reader_metrics: &mut ReaderMetrics,
258    ) -> bool {
259        if self.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
260            return true;
261        }
262        let file = &self.stream_ctx.input.files[file_index];
263        if !self.stream_ctx.input.can_manifest_prune_file(file) {
264            return false;
265        }
266        if self.manifest_pruned_files[file_index]
267            .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
268            .is_ok()
269        {
270            reader_metrics.filter_metrics.files_time_range_pruned += 1;
271        }
272        true
273    }
274}
275
276/// Per-file state tracking.
277struct FileBuilderEntry {
278    /// Cached builder for this file.
279    ///
280    /// A newly completed builder is cached only while `remaining_ranges > 0`.
281    /// In retaining mode, the builder may remain cached after
282    /// `remaining_ranges` reaches zero.
283    builder: Option<Arc<FileRangeBuilder>>,
284    /// Number of ranges remaining in the current initialized batch.
285    ///
286    /// This may be zero while `builder` is still present in retaining mode.
287    remaining_ranges: usize,
288    /// Waiters when pruning is in-progress.
289    waiters: Vec<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
290}
291
292/// Request to prune a file.
293struct PruneRequest {
294    /// Index of the file in ScanInput.files.
295    file_index: usize,
296    /// Pre-filter mode to use for the file.
297    pre_filter_mode: PreFilterMode,
298    /// Oneshot channel to send back the result.
299    response_tx: Option<oneshot::Sender<Result<Arc<FileRangeBuilder>>>>,
300    /// Partition metrics for merging reader metrics.
301    partition_metrics: Option<PartitionMetrics>,
302}
303
304impl Pruner {
305    /// Creates a new Pruner with N worker tasks.
306    ///
307    /// Initially all file_entries have `remaining_ranges = 0`.
308    /// Call `add_partition_ranges()` to initialize ref counts.
309    pub fn new(stream_ctx: Arc<StreamContext>, num_workers: usize) -> Self {
310        Self::new_with_options(stream_ctx, num_workers, PrunerOptions::default())
311    }
312
313    /// Creates a new pruner with the given options.
314    pub fn new_with_options(
315        stream_ctx: Arc<StreamContext>,
316        num_workers: usize,
317        options: PrunerOptions,
318    ) -> Self {
319        let PrunerOptions {
320            retain_builders,
321            enable_predicate_prefilter,
322        } = options;
323        let num_files = stream_ctx.input.num_files();
324        let file_entries: Vec<_> = (0..num_files)
325            .map(|_| {
326                Mutex::new(FileBuilderEntry {
327                    builder: None,
328                    remaining_ranges: 0,
329                    waiters: Vec::new(),
330                })
331            })
332            .collect();
333        let manifest_pruned_files: Vec<AtomicBool> =
334            (0..num_files).map(|_| AtomicBool::new(false)).collect();
335        // Create channels and collect senders
336        let mut worker_senders = Vec::with_capacity(num_workers);
337        let mut receivers = Vec::with_capacity(num_workers);
338        for _ in 0..num_workers {
339            let (tx, rx) = mpsc::channel::<PruneRequest>(64);
340            worker_senders.push(tx);
341            receivers.push(rx);
342        }
343
344        let inner = Arc::new(PrunerInner {
345            num_workers,
346            file_entries,
347            stream_ctx,
348            manifest_pruned_files,
349            retain_builders,
350            enable_predicate_prefilter,
351        });
352
353        // Keep pruning and prefetching on the runtime of the originating workload.
354        for (worker_id, rx) in receivers.into_iter().enumerate() {
355            let worker = Self::worker_loop(worker_id, rx, inner.clone());
356            if inner.stream_ctx.input.compaction {
357                common_runtime::spawn_compact(worker);
358            } else {
359                common_runtime::spawn_query(worker);
360            }
361        }
362
363        Self {
364            worker_senders,
365            inner,
366        }
367    }
368
369    /// Adds reference counts for partition ranges that the caller will scan.
370    ///
371    /// The caller must keep these ranges consistent with the ranges passed to
372    /// `PartitionPruner`. If called multiple times, all calls must describe
373    /// ranges from the same logical scan because their reference counts
374    /// accumulate in this pruner.
375    pub fn add_partition_ranges(&self, partition_ranges: &[PartitionRange]) {
376        // Reset manifest-prune results so the latest dynamic filters are visible.
377        for pruned in &self.inner.manifest_pruned_files {
378            pruned.store(false, Ordering::Relaxed);
379        }
380
381        // Add reference counts for each partition range
382        let num_memtables = self.inner.stream_ctx.input.num_memtables();
383        for part_range in partition_ranges {
384            let range_meta = &self.inner.stream_ctx.ranges[part_range.identifier];
385            for row_group_index in &range_meta.row_group_indices {
386                if self.inner.stream_ctx.is_file_range_index(*row_group_index) {
387                    let file_index = row_group_index.index - num_memtables;
388                    if file_index < self.inner.file_entries.len() {
389                        let mut entry = self.inner.file_entries[file_index].lock().unwrap();
390                        entry.remaining_ranges += 1;
391                    }
392                }
393            }
394        }
395    }
396
397    /// Gets or creates the FileRangeBuilder for a file, builds ranges,
398    /// and decrements its ref count. Non-retained builders are cleaned up at zero.
399    ///
400    /// Callers should invoke [add_partition_ranges()](Pruner::add_partition_ranges()) to initialize the
401    /// file entries and ref counts.
402    pub async fn build_file_ranges(
403        &self,
404        index: RowGroupIndex,
405        pre_filter_mode: PreFilterMode,
406        partition_metrics: &PartitionMetrics,
407        reader_metrics: &mut ReaderMetrics,
408    ) -> Result<SmallVec<[FileRange; 2]>> {
409        let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
410
411        // Get builder (from cache or by pruning)
412        let builder = self
413            .get_file_builder(
414                file_index,
415                pre_filter_mode,
416                partition_metrics,
417                reader_metrics,
418            )
419            .await?;
420
421        // Build ranges
422        let mut ranges = SmallVec::new();
423        builder.build_ranges(index.row_group_index, &mut ranges);
424
425        // Decrement ref count and clean up non-retained builders if needed.
426        self.decrement_and_maybe_clear(file_index, reader_metrics);
427
428        Ok(ranges)
429    }
430
431    /// Skips a file range that has been pruned before entering the file pruner.
432    ///
433    /// This keeps the pruner's per-file reference counts balanced with
434    /// `add_partition_ranges()`. It may also clear a cached builder when this was the
435    /// last remaining range for the file.
436    pub fn skip_file_range(&self, index: RowGroupIndex, reader_metrics: &mut ReaderMetrics) {
437        if !self.inner.stream_ctx.is_file_range_index(index) {
438            return;
439        }
440        let file_index = index.index - self.inner.stream_ctx.input.num_memtables();
441        self.decrement_and_maybe_clear(file_index, reader_metrics);
442    }
443
444    /// Gets or creates the FileRangeBuilder for a file.
445    async fn get_file_builder(
446        &self,
447        file_index: usize,
448        pre_filter_mode: PreFilterMode,
449        partition_metrics: &PartitionMetrics,
450        reader_metrics: &mut ReaderMetrics,
451    ) -> Result<Arc<FileRangeBuilder>> {
452        // Fast path: checks cache
453        {
454            let entry = self.inner.file_entries[file_index].lock().unwrap();
455            if let Some(builder) = &entry.builder {
456                reader_metrics.filter_metrics.pruner_cache_hit += 1;
457                return Ok(builder.clone());
458            }
459        }
460
461        reader_metrics.filter_metrics.pruner_cache_miss += 1;
462        let prune_start = Instant::now();
463        let file = &self.inner.stream_ctx.input.files[file_index];
464        let file_id = file.file_id().file_id();
465        let worker_idx = self.get_worker_idx(file_id);
466
467        let (response_tx, response_rx) = oneshot::channel();
468        let request = PruneRequest {
469            file_index,
470            pre_filter_mode,
471            response_tx: Some(response_tx),
472            partition_metrics: Some(partition_metrics.clone()),
473        };
474
475        let result = if self.worker_senders[worker_idx].send(request).await.is_err() {
476            common_telemetry::warn!("Worker channel closed, falling back to direct pruning");
477            // Worker channel closed, falls back to direct pruning
478            self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
479                .await
480        } else {
481            // Waits for response
482            match response_rx.await {
483                Ok(result) => result,
484                Err(_) => {
485                    common_telemetry::warn!(
486                        "Response channel closed, falling back to direct pruning"
487                    );
488                    // Channel closed, falls back to direct pruning
489                    self.prune_file_directly(file_index, pre_filter_mode, reader_metrics)
490                        .await
491                }
492            }
493        };
494        reader_metrics.filter_metrics.pruner_prune_cost += prune_start.elapsed();
495        result
496    }
497
498    /// Gets or creates the FileRangeBuilder for a file.
499    pub fn get_file_builder_background(
500        &self,
501        file_index: usize,
502        pre_filter_mode: PreFilterMode,
503        partition_metrics: Option<PartitionMetrics>,
504    ) {
505        // Fast path: checks cache
506        {
507            let entry = self.inner.file_entries[file_index].lock().unwrap();
508            if entry.builder.is_some() {
509                return;
510            }
511        }
512
513        let file = &self.inner.stream_ctx.input.files[file_index];
514        let file_id = file.file_id().file_id();
515        let worker_idx = self.get_worker_idx(file_id);
516
517        let request = PruneRequest {
518            file_index,
519            pre_filter_mode,
520            response_tx: None,
521            partition_metrics,
522        };
523
524        // Sends request to worker
525        let _ = self.worker_senders[worker_idx].try_send(request);
526    }
527
528    /// Returns whether this pruner builds file ranges with the reduced-column
529    /// predicate prefilter enabled.
530    pub fn predicate_prefilter_enabled(&self) -> bool {
531        self.inner.enable_predicate_prefilter
532    }
533
534    fn get_worker_idx(&self, file_id: FileId) -> usize {
535        let file_id_hash = Uuid::from(file_id).as_u128() as usize;
536        file_id_hash % self.inner.num_workers
537    }
538
539    /// Prunes a file directly without going through a worker.
540    /// Used as fallback when worker channels are closed.
541    async fn prune_file_directly(
542        &self,
543        file_index: usize,
544        pre_filter_mode: PreFilterMode,
545        reader_metrics: &mut ReaderMetrics,
546    ) -> Result<Arc<FileRangeBuilder>> {
547        // Check manifest-level prune first (shared cache, no I/O).
548        if self
549            .inner
550            .try_mark_manifest_pruned(file_index, reader_metrics)
551        {
552            let arc_builder = Arc::new(FileRangeBuilder::default());
553            // Do NOT cache an empty manifest-pruned builder; the cache flag
554            // already records the decision.
555            return Ok(arc_builder);
556        }
557
558        let file = &self.inner.stream_ctx.input.files[file_index];
559        let predicate = self.inner.stream_ctx.input.predicate_for_file(file);
560        let builder = self
561            .inner
562            .stream_ctx
563            .input
564            .prune_file_after_manifest_check(
565                file,
566                pre_filter_mode,
567                self.inner.enable_predicate_prefilter,
568                predicate,
569                reader_metrics,
570            )
571            .await?;
572
573        let arc_builder = Arc::new(builder);
574
575        // Cache only while the file still has remaining ranges. Retaining mode
576        // affects cleanup after the builder has been cached, not cache eligibility.
577        {
578            let mut entry = self.inner.file_entries[file_index].lock().unwrap();
579            cache_builder_if_needed(&mut entry, &arc_builder, reader_metrics);
580        }
581
582        Ok(arc_builder)
583    }
584
585    /// Decrements ref count and clears builder if no longer needed.
586    fn decrement_and_maybe_clear(&self, file_index: usize, reader_metrics: &mut ReaderMetrics) {
587        let mut entry = self.inner.file_entries[file_index].lock().unwrap();
588        entry.remaining_ranges = entry.remaining_ranges.saturating_sub(1);
589
590        if !self.inner.retain_builders
591            && entry.remaining_ranges == 0
592            && let Some(builder) = entry.builder.take()
593        {
594            PRUNER_ACTIVE_BUILDERS.dec();
595            reader_metrics.metadata_mem_size -= builder.memory_size() as isize;
596            reader_metrics.num_range_builders -= 1;
597        }
598    }
599
600    /// Worker loop that processes prune requests.
601    async fn worker_loop(
602        worker_id: usize,
603        mut rx: mpsc::Receiver<PruneRequest>,
604        inner: Arc<PrunerInner>,
605    ) {
606        let mut worker_cache_hit = 0;
607        let mut worker_cache_miss = 0;
608        let mut pruned_files = Vec::new();
609
610        while let Some(request) = rx.recv().await {
611            let PruneRequest {
612                file_index,
613                pre_filter_mode,
614                response_tx,
615                partition_metrics,
616            } = request;
617
618            // Check if already cached or in-progress
619            {
620                let entry = inner.file_entries[file_index].lock().unwrap();
621                if let Some(builder) = &entry.builder {
622                    // Cache hit - send immediately
623                    if let Some(response_tx) = response_tx {
624                        let _ = response_tx.send(Ok(builder.clone()));
625                    }
626                    worker_cache_hit += 1;
627                    continue;
628                }
629            }
630            worker_cache_miss += 1;
631
632            let file = &inner.stream_ctx.input.files[file_index];
633            pruned_files.push(file.file_id().file_id());
634            let explain_verbose = partition_metrics
635                .as_ref()
636                .map(|m| m.explain_verbose())
637                .unwrap_or(false);
638            let mut metrics = ReaderMetrics {
639                filter_metrics: new_filter_metrics(explain_verbose),
640                ..Default::default()
641            };
642
643            // Check manifest-level prune first (shared cache, no I/O).
644            let result = if inner.try_mark_manifest_pruned(file_index, &mut metrics) {
645                // Manifest-level pruning proved the file empty — produce a
646                // default builder without reading any parquet metadata.
647                Ok(FileRangeBuilder::default())
648            } else {
649                let predicate = inner.stream_ctx.input.predicate_for_file(file);
650                inner
651                    .stream_ctx
652                    .input
653                    .prune_file_after_manifest_check(
654                        file,
655                        pre_filter_mode,
656                        inner.enable_predicate_prefilter,
657                        predicate,
658                        &mut metrics,
659                    )
660                    .await
661            };
662
663            // Update state and notify waiters
664            let mut entry = inner.file_entries[file_index].lock().unwrap();
665            match result {
666                Ok(builder) => {
667                    let arc_builder = Arc::new(builder);
668                    let is_background = response_tx.is_none();
669
670                    // Cache only if the file still has remaining ranges. If
671                    // remaining_ranges == 0, a concurrent `skip_file_range` already consumed
672                    // all ranges and may have cleared a previously cached builder. Retaining
673                    // mode only preserves builders that were cached before the count reached 0.
674                    // Skip caching manifest-pruned empty builders; the cache flag is enough.
675                    let did_cache =
676                        if inner.manifest_pruned_files[file_index].load(Ordering::Relaxed) {
677                            false
678                        } else {
679                            cache_builder_if_needed(&mut entry, &arc_builder, &mut metrics)
680                        };
681
682                    // Notify all waiters
683                    for waiter in entry.waiters.drain(..) {
684                        let _ = waiter.send(Ok(arc_builder.clone()));
685                    }
686                    // Always respond to foreground caller, even if we did not cache.
687                    if let Some(response_tx) = response_tx {
688                        let _ = response_tx.send(Ok(arc_builder));
689                    }
690
691                    debug!(
692                        "Pruner worker {} pruned file_index: {}, file: {:?}, metrics: {:?}",
693                        worker_id,
694                        file_index,
695                        file.file_id(),
696                        metrics
697                    );
698
699                    // Merge metrics if this is a foreground request, or if the builder
700                    // was cached. Skip stale per-file metrics
701                    // for background requests that completed after the file was already
702                    // fully skipped.
703                    if (!is_background || did_cache)
704                        && let Some(part_metrics) = &partition_metrics
705                    {
706                        let per_file_metrics = if part_metrics.explain_verbose() {
707                            let file_id = file.file_id();
708                            let mut map = HashMap::new();
709                            map.insert(
710                                file_id,
711                                FileScanMetrics {
712                                    build_part_cost: metrics.build_cost,
713                                    ..Default::default()
714                                },
715                            );
716                            Some(map)
717                        } else {
718                            None
719                        };
720                        part_metrics.merge_reader_metrics(&metrics, per_file_metrics.as_ref());
721                    }
722                }
723                Err(e) => {
724                    let arc_error = Arc::new(e);
725                    for waiter in entry.waiters.drain(..) {
726                        let _ = waiter.send(Err(arc_error.clone()).context(PruneFileSnafu));
727                    }
728                    if let Some(response_tx) = response_tx {
729                        let _ = response_tx.send(Err(arc_error).context(PruneFileSnafu));
730                    }
731                }
732            }
733        }
734
735        common_telemetry::debug!(
736            "Pruner worker {} finished, cache_hit: {}, cache_miss: {}, files: {:?}",
737            worker_id,
738            worker_cache_hit,
739            worker_cache_miss,
740            pruned_files,
741        );
742    }
743}
744
745#[cfg(test)]
746impl Pruner {
747    /// Returns the remaining range count for a file (test-only).
748    pub(crate) fn test_remaining_ranges(&self, file_index: usize) -> usize {
749        self.inner.file_entries[file_index]
750            .lock()
751            .unwrap()
752            .remaining_ranges
753    }
754
755    /// Returns whether a cached builder exists for a file (test-only).
756    fn test_has_builder(&self, file_index: usize) -> bool {
757        self.inner.file_entries[file_index]
758            .lock()
759            .unwrap()
760            .builder
761            .is_some()
762    }
763
764    /// Returns the manifest-pruned flag for a file (test-only).
765    fn test_is_manifest_pruned(&self, file_index: usize) -> bool {
766        self.inner.manifest_pruned_files[file_index].load(Ordering::Relaxed)
767    }
768
769    /// Clears a cached builder for a file, simulating stale cleanup (test-only).
770    #[allow(dead_code)]
771    fn test_clear_builder(&self, file_index: usize) {
772        let mut entry = self.inner.file_entries[file_index].lock().unwrap();
773        if entry.builder.take().is_some() {
774            PRUNER_ACTIVE_BUILDERS.dec();
775        }
776    }
777}
778
779/// Returns true if a freshly pruned builder should be cached for this file.
780fn should_cache_builder(entry: &FileBuilderEntry) -> bool {
781    entry.builder.is_none() && entry.remaining_ranges > 0
782}
783
784/// Caches a freshly pruned builder if it is still referenced, and records the
785/// corresponding builder memory/count deltas for verbose metrics.
786fn cache_builder_if_needed(
787    entry: &mut FileBuilderEntry,
788    builder: &Arc<FileRangeBuilder>,
789    reader_metrics: &mut ReaderMetrics,
790) -> bool {
791    if should_cache_builder(entry) {
792        reader_metrics.metadata_mem_size += builder.memory_size() as isize;
793        reader_metrics.num_range_builders += 1;
794        entry.builder = Some(builder.clone());
795        PRUNER_ACTIVE_BUILDERS.inc();
796        true
797    } else {
798        false
799    }
800}
801
802#[cfg(test)]
803mod tests {
804    use common_time::Timestamp;
805    use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
806    use datafusion_common::ScalarValue;
807    use datafusion_expr::{Expr, col, lit};
808    use store_api::region_engine::PartitionRange;
809    use store_api::storage::{FileId, RegionId};
810
811    use super::*;
812    use crate::read::flat_projection::FlatProjectionMapper;
813    use crate::read::range::RowGroupIndex;
814    use crate::read::scan_region::{PredicateGroup, ScanInput};
815    use crate::read::scan_util::PartitionMetrics;
816    use crate::sst::file::{FileHandle, FileMeta};
817    use crate::sst::parquet::reader::ReaderMetrics;
818    use crate::test_util::memtable_util::metadata_with_primary_key;
819    use crate::test_util::new_noop_file_purger;
820    use crate::test_util::scheduler_util::SchedulerEnv;
821
822    async fn make_test_pruner(num_files: usize) -> (SchedulerEnv, Arc<Pruner>) {
823        make_test_pruner_with_retained_builders(num_files, false).await
824    }
825
826    async fn make_test_pruner_with_retained_builders(
827        num_files: usize,
828        retain_builders: bool,
829    ) -> (SchedulerEnv, Arc<Pruner>) {
830        let env = SchedulerEnv::new().await;
831        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
832        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
833
834        let files: Vec<FileHandle> = (0..num_files)
835            .map(|_| {
836                let meta = FileMeta {
837                    region_id: RegionId::new(123, 456),
838                    file_id: FileId::random(),
839                    time_range: (
840                        Timestamp::new_millisecond(0),
841                        Timestamp::new_millisecond(1000),
842                    ),
843                    num_row_groups: 1,
844                    num_rows: 1024,
845                    level: 0,
846                    ..Default::default()
847                };
848                FileHandle::new(meta, new_noop_file_purger())
849            })
850            .collect();
851
852        let input = ScanInput::builder(env.access_layer.clone(), mapper)
853            .with_files(files)
854            .with_append_mode(true)
855            .build();
856        let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
857        let pruner = Arc::new(Pruner::new_with_options(
858            stream_ctx,
859            1,
860            PrunerOptions {
861                retain_builders,
862                ..Default::default()
863            },
864        ));
865        (env, pruner)
866    }
867
868    /// Builds a minimal `PartitionRange` that references `file_index`.
869    /// `add_partition_ranges` will look up `stream_ctx.ranges[identifier]`
870    /// and find `row_group_indices[0] == RowGroupIndex { index: file_index,
871    /// row_group_index: 0 }` because `unordered_scan_ranges` with
872    /// `num_row_groups=1` produces one range per file.
873    fn file_partition_range(file_index: usize) -> PartitionRange {
874        PartitionRange {
875            start: Timestamp::new_millisecond(0),
876            end: Timestamp::new_millisecond(1001),
877            num_rows: 1024,
878            identifier: file_index,
879        }
880    }
881
882    async fn make_test_pruner_with_predicate(
883        num_files: usize,
884        row_groups_per_file: u64,
885        predicate_exprs: &[Expr],
886    ) -> (SchedulerEnv, Arc<Pruner>) {
887        let env = SchedulerEnv::new().await;
888        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
889        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
890        let predicate = PredicateGroup::new(&metadata, predicate_exprs).unwrap();
891
892        let files: Vec<FileHandle> = (0..num_files)
893            .map(|_| {
894                let meta = FileMeta {
895                    region_id: RegionId::new(123, 456),
896                    file_id: FileId::random(),
897                    time_range: (
898                        Timestamp::new_millisecond(0),
899                        Timestamp::new_millisecond(1000),
900                    ),
901                    num_row_groups: row_groups_per_file,
902                    num_rows: row_groups_per_file * 1024,
903                    level: 0,
904                    ..Default::default()
905                };
906                FileHandle::new(meta, new_noop_file_purger())
907            })
908            .collect();
909
910        let input = ScanInput::builder(env.access_layer.clone(), mapper)
911            .with_files(files)
912            .with_predicate(predicate)
913            .with_append_mode(true)
914            .build();
915        let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
916        let pruner = Arc::new(Pruner::new(stream_ctx, 1));
917        (env, pruner)
918    }
919
920    fn make_partition_metrics() -> PartitionMetrics {
921        let metrics_set = ExecutionPlanMetricsSet::new();
922        PartitionMetrics::new(
923            RegionId::new(123, 456),
924            0,
925            "test",
926            Instant::now(),
927            false,
928            &metrics_set,
929        )
930    }
931
932    #[test]
933    fn should_cache_builder_when_ranges_remain() {
934        let entry = FileBuilderEntry {
935            builder: None,
936            remaining_ranges: 3,
937            waiters: Vec::new(),
938        };
939        assert!(should_cache_builder(&entry));
940    }
941
942    #[test]
943    fn should_not_cache_builder_when_no_ranges_remain() {
944        let entry = FileBuilderEntry {
945            builder: None,
946            remaining_ranges: 0,
947            waiters: Vec::new(),
948        };
949        assert!(!should_cache_builder(&entry));
950    }
951
952    #[test]
953    fn should_not_cache_builder_when_already_cached() {
954        let entry = FileBuilderEntry {
955            builder: Some(Arc::new(FileRangeBuilder::default())),
956            remaining_ranges: 1,
957            waiters: Vec::new(),
958        };
959        assert!(!should_cache_builder(&entry));
960    }
961
962    #[test]
963    fn cache_builder_records_metrics() {
964        let mut entry = FileBuilderEntry {
965            builder: None,
966            remaining_ranges: 1,
967            waiters: Vec::new(),
968        };
969        let builder = Arc::new(FileRangeBuilder::default());
970        let mut reader_metrics = ReaderMetrics::default();
971
972        assert!(cache_builder_if_needed(
973            &mut entry,
974            &builder,
975            &mut reader_metrics,
976        ));
977        assert!(entry.builder.is_some());
978        assert_eq!(
979            reader_metrics.metadata_mem_size,
980            builder.memory_size() as isize
981        );
982        assert_eq!(reader_metrics.num_range_builders, 1);
983
984        if entry.builder.take().is_some() {
985            PRUNER_ACTIVE_BUILDERS.dec();
986        }
987    }
988
989    #[tokio::test]
990    async fn skip_file_range_decrements_and_clears_builder() {
991        let (_env, pruner) = make_test_pruner(1).await;
992
993        // Simulate 3 partition ranges for file 0.
994        let ranges: Vec<PartitionRange> = (0..3).map(|_| file_partition_range(0)).collect();
995        pruner.add_partition_ranges(&ranges);
996        assert_eq!(pruner.test_remaining_ranges(0), 3);
997
998        // Manually set a cached builder (simulating a previous cache hit).
999        {
1000            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1001            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1002            PRUNER_ACTIVE_BUILDERS.inc();
1003        }
1004        assert!(pruner.test_has_builder(0));
1005
1006        // Skip all 3 ranges; the third should clear the builder.
1007        let mut reader_metrics = ReaderMetrics::default();
1008        for i in 0..3 {
1009            let index = RowGroupIndex {
1010                index: 0,
1011                row_group_index: i as i64,
1012            };
1013            pruner.skip_file_range(index, &mut reader_metrics);
1014        }
1015
1016        assert_eq!(pruner.test_remaining_ranges(0), 0);
1017        assert!(!pruner.test_has_builder(0));
1018    }
1019
1020    #[tokio::test]
1021    async fn retained_builder_survives_after_last_range() {
1022        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1023        pruner.add_partition_ranges(&[file_partition_range(0)]);
1024        {
1025            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1026            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1027            PRUNER_ACTIVE_BUILDERS.inc();
1028        }
1029
1030        let mut reader_metrics = ReaderMetrics::default();
1031        pruner.skip_file_range(
1032            RowGroupIndex {
1033                index: 0,
1034                row_group_index: 0,
1035            },
1036            &mut reader_metrics,
1037        );
1038
1039        assert_eq!(pruner.test_remaining_ranges(0), 0);
1040        assert!(pruner.test_has_builder(0));
1041
1042        let partition_metrics = make_partition_metrics();
1043        let mut reader_metrics = ReaderMetrics::default();
1044        let _builder = pruner
1045            .get_file_builder(
1046                0,
1047                PreFilterMode::SkipFields,
1048                &partition_metrics,
1049                &mut reader_metrics,
1050            )
1051            .await
1052            .unwrap();
1053        assert_eq!(reader_metrics.filter_metrics.pruner_cache_hit, 1);
1054    }
1055
1056    #[tokio::test]
1057    async fn add_partition_ranges_keeps_retained_builder() {
1058        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1059        pruner.add_partition_ranges(&[file_partition_range(0)]);
1060        {
1061            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1062            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1063            PRUNER_ACTIVE_BUILDERS.inc();
1064        }
1065
1066        let mut reader_metrics = ReaderMetrics::default();
1067        pruner.skip_file_range(
1068            RowGroupIndex {
1069                index: 0,
1070                row_group_index: 0,
1071            },
1072            &mut reader_metrics,
1073        );
1074        assert!(pruner.test_has_builder(0));
1075
1076        pruner.add_partition_ranges(&[file_partition_range(0)]);
1077        assert!(pruner.test_has_builder(0));
1078        assert_eq!(pruner.test_remaining_ranges(0), 1);
1079    }
1080
1081    #[tokio::test]
1082    async fn retaining_mode_does_not_cache_after_skip_file_range_consumed_all() {
1083        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1084
1085        // Simulate one range for file 0.
1086        let ranges = vec![file_partition_range(0)];
1087        pruner.add_partition_ranges(&ranges);
1088        assert_eq!(pruner.test_remaining_ranges(0), 1);
1089
1090        // Simulate skip_file_range consuming the last range BEFORE the
1091        // background worker finishes. This mirrors the race: a dynamic filter
1092        // tightens and manifest-prune fast-skip zeros out remaining_ranges.
1093        let mut reader_metrics = ReaderMetrics::default();
1094        let index = RowGroupIndex {
1095            index: 0,
1096            row_group_index: 0,
1097        };
1098        pruner.skip_file_range(index, &mut reader_metrics);
1099        assert_eq!(pruner.test_remaining_ranges(0), 0);
1100        assert!(!pruner.test_has_builder(0));
1101
1102        // Now simulate the worker completing: check the caching guard.
1103        let entry = pruner.inner.file_entries[0].lock().unwrap();
1104        let should_cache = should_cache_builder(&entry);
1105        drop(entry);
1106
1107        assert!(!should_cache);
1108
1109        // Ensure the gauge was not incremented for a stale builder.
1110        // (skip_file_range already decremented it if there was one, but here
1111        // there was none, so the gauge should be at baseline.)
1112    }
1113
1114    #[tokio::test]
1115    async fn worker_caches_when_ranges_remain() {
1116        let (_env, pruner) = make_test_pruner(1).await;
1117
1118        // Simulate 2 ranges for file 0.
1119        let ranges: Vec<PartitionRange> = (0..2).map(|_| file_partition_range(0)).collect();
1120        pruner.add_partition_ranges(&ranges);
1121        assert_eq!(pruner.test_remaining_ranges(0), 2);
1122
1123        // Consume only 1 range.
1124        let mut reader_metrics = ReaderMetrics::default();
1125        let index = RowGroupIndex {
1126            index: 0,
1127            row_group_index: 0,
1128        };
1129        pruner.skip_file_range(index, &mut reader_metrics);
1130        assert_eq!(pruner.test_remaining_ranges(0), 1);
1131
1132        // The worker should still cache because remaining_ranges > 0.
1133        let entry = pruner.inner.file_entries[0].lock().unwrap();
1134        assert!(should_cache_builder(&entry));
1135    }
1136
1137    // ── Corner case: fast-skip across multiple row groups ─────────────
1138
1139    /// 1 file × 3 row groups, predicate `ts > 10000ms` prunes the file at
1140    /// manifest level. Fast-skipping each of the 3 row groups must return true
1141    /// and decrement remaining_ranges to 0.
1142    #[tokio::test]
1143    async fn try_skip_manifest_pruned_file_range_multi_row_groups() {
1144        let predicate_exprs: Vec<Expr> =
1145            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1146        let (_env, pruner) = make_test_pruner_with_predicate(1, 3, &predicate_exprs).await;
1147
1148        let ranges = pruner.inner.stream_ctx.partition_ranges();
1149        assert_eq!(ranges.len(), 3);
1150        pruner.add_partition_ranges(&ranges);
1151        assert_eq!(pruner.test_remaining_ranges(0), 3);
1152
1153        let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1154        let partition_metrics = make_partition_metrics();
1155
1156        // Fast-skip each of the 3 row groups.
1157        for rg in 0..3 {
1158            let index = RowGroupIndex {
1159                index: 0, // file_index == 0, no memtables
1160                row_group_index: rg,
1161            };
1162            let skipped =
1163                partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1164            assert!(skipped, "row group {} should be skipped", rg);
1165        }
1166
1167        // All refs consumed.
1168        assert_eq!(pruner.test_remaining_ranges(0), 0);
1169        // manifest_pruned_files is CAS'd exactly once (first call).
1170        assert!(pruner.test_is_manifest_pruned(0));
1171    }
1172
1173    /// A file whose manifest time range may contain matching rows must not be
1174    /// fast-skipped. This protects query correctness over metrics precision.
1175    #[tokio::test]
1176    async fn try_skip_manifest_pruned_file_range_keeps_overlapping_file() {
1177        let predicate_exprs: Vec<Expr> =
1178            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(500), None)))];
1179        let (_env, pruner) = make_test_pruner_with_predicate(1, 2, &predicate_exprs).await;
1180
1181        let ranges = pruner.inner.stream_ctx.partition_ranges();
1182        assert_eq!(ranges.len(), 2);
1183        pruner.add_partition_ranges(&ranges);
1184        assert_eq!(pruner.test_remaining_ranges(0), 2);
1185
1186        let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1187        let partition_metrics = make_partition_metrics();
1188        let range_meta = &pruner.inner.stream_ctx.ranges[ranges[0].identifier];
1189        let index = range_meta.row_group_indices[0];
1190
1191        let skipped =
1192            partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1193
1194        assert!(!skipped);
1195        assert_eq!(pruner.test_remaining_ranges(0), 2);
1196        assert!(!pruner.test_is_manifest_pruned(0));
1197    }
1198
1199    // ── Corner case: add_partition_ranges resets the manifest-pruned flag ──
1200
1201    #[tokio::test]
1202    async fn add_partition_ranges_resets_manifest_pruned_flag() {
1203        let predicate_exprs: Vec<Expr> =
1204            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1205        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1206
1207        // Mark file 0 as manifest-pruned.
1208        let mut reader_metrics = ReaderMetrics::default();
1209        let marked = pruner
1210            .inner
1211            .try_mark_manifest_pruned(0, &mut reader_metrics);
1212        assert!(marked);
1213        assert!(pruner.test_is_manifest_pruned(0));
1214        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1215
1216        // Calling add_partition_ranges must reset the flag.
1217        let ranges = vec![file_partition_range(0)];
1218        pruner.add_partition_ranges(&ranges);
1219        assert!(!pruner.test_is_manifest_pruned(0));
1220        // remaining_ranges was also incremented.
1221        assert_eq!(pruner.test_remaining_ranges(0), 1);
1222    }
1223
1224    // ── Corner case: prune_file_directly short-circuits via manifest prune ──
1225
1226    #[tokio::test]
1227    async fn prune_file_directly_manifest_pruned_returns_empty_builder() {
1228        let predicate_exprs: Vec<Expr> =
1229            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1230        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1231
1232        // Ensure there is no cached builder yet.
1233        assert!(!pruner.test_has_builder(0));
1234
1235        let mut reader_metrics = ReaderMetrics::default();
1236        let builder = pruner
1237            .prune_file_directly(0, PreFilterMode::SkipFields, &mut reader_metrics)
1238            .await
1239            .unwrap();
1240
1241        // Should be the default (empty) builder.
1242        assert_eq!(
1243            builder.memory_size(),
1244            FileRangeBuilder::default().memory_size()
1245        );
1246        // builder must NOT be cached — the manifest-pruned flag is enough.
1247        assert!(!pruner.test_has_builder(0));
1248        // files_time_range_pruned was recorded.
1249        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1250    }
1251
1252    // ── Corner case: try_mark_manifest_pruned does not double-count ───
1253
1254    #[tokio::test]
1255    async fn try_mark_manifest_pruned_only_counts_first_cas() {
1256        let predicate_exprs: Vec<Expr> =
1257            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1258        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1259
1260        // First call: CAS succeeds, metric incremented.
1261        let mut reader_metrics = ReaderMetrics::default();
1262        let marked = pruner
1263            .inner
1264            .try_mark_manifest_pruned(0, &mut reader_metrics);
1265        assert!(marked);
1266        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1267        assert!(pruner.test_is_manifest_pruned(0));
1268
1269        // Second call: already true, no metric delta.
1270        let mut reader_metrics2 = ReaderMetrics::default();
1271        let marked2 = pruner
1272            .inner
1273            .try_mark_manifest_pruned(0, &mut reader_metrics2);
1274        assert!(marked2);
1275        assert_eq!(reader_metrics2.filter_metrics.files_time_range_pruned, 0);
1276    }
1277}