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