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::builder(env.access_layer.clone(), mapper)
845            .with_files(files)
846            .with_append_mode(true)
847            .build();
848        let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
849        let pruner = Arc::new(Pruner::new_with_options(
850            stream_ctx,
851            1,
852            PrunerOptions {
853                retain_builders,
854                ..Default::default()
855            },
856        ));
857        (env, pruner)
858    }
859
860    /// Builds a minimal `PartitionRange` that references `file_index`.
861    /// `add_partition_ranges` will look up `stream_ctx.ranges[identifier]`
862    /// and find `row_group_indices[0] == RowGroupIndex { index: file_index,
863    /// row_group_index: 0 }` because `unordered_scan_ranges` with
864    /// `num_row_groups=1` produces one range per file.
865    fn file_partition_range(file_index: usize) -> PartitionRange {
866        PartitionRange {
867            start: Timestamp::new_millisecond(0),
868            end: Timestamp::new_millisecond(1001),
869            num_rows: 1024,
870            identifier: file_index,
871        }
872    }
873
874    async fn make_test_pruner_with_predicate(
875        num_files: usize,
876        row_groups_per_file: u64,
877        predicate_exprs: &[Expr],
878    ) -> (SchedulerEnv, Arc<Pruner>) {
879        let env = SchedulerEnv::new().await;
880        let metadata = Arc::new(metadata_with_primary_key(vec![0, 1], false));
881        let mapper = FlatProjectionMapper::new(&metadata, [0, 2, 3]).unwrap();
882        let predicate = PredicateGroup::new(&metadata, predicate_exprs).unwrap();
883
884        let files: Vec<FileHandle> = (0..num_files)
885            .map(|_| {
886                let meta = FileMeta {
887                    region_id: RegionId::new(123, 456),
888                    file_id: FileId::random(),
889                    time_range: (
890                        Timestamp::new_millisecond(0),
891                        Timestamp::new_millisecond(1000),
892                    ),
893                    num_row_groups: row_groups_per_file,
894                    num_rows: row_groups_per_file * 1024,
895                    level: 0,
896                    ..Default::default()
897                };
898                FileHandle::new(meta, new_noop_file_purger())
899            })
900            .collect();
901
902        let input = ScanInput::builder(env.access_layer.clone(), mapper)
903            .with_files(files)
904            .with_predicate(predicate)
905            .with_append_mode(true)
906            .build();
907        let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
908        let pruner = Arc::new(Pruner::new(stream_ctx, 1));
909        (env, pruner)
910    }
911
912    fn make_partition_metrics() -> PartitionMetrics {
913        let metrics_set = ExecutionPlanMetricsSet::new();
914        PartitionMetrics::new(
915            RegionId::new(123, 456),
916            0,
917            "test",
918            Instant::now(),
919            false,
920            &metrics_set,
921        )
922    }
923
924    #[test]
925    fn should_cache_builder_when_ranges_remain() {
926        let entry = FileBuilderEntry {
927            builder: None,
928            remaining_ranges: 3,
929            waiters: Vec::new(),
930        };
931        assert!(should_cache_builder(&entry));
932    }
933
934    #[test]
935    fn should_not_cache_builder_when_no_ranges_remain() {
936        let entry = FileBuilderEntry {
937            builder: None,
938            remaining_ranges: 0,
939            waiters: Vec::new(),
940        };
941        assert!(!should_cache_builder(&entry));
942    }
943
944    #[test]
945    fn should_not_cache_builder_when_already_cached() {
946        let entry = FileBuilderEntry {
947            builder: Some(Arc::new(FileRangeBuilder::default())),
948            remaining_ranges: 1,
949            waiters: Vec::new(),
950        };
951        assert!(!should_cache_builder(&entry));
952    }
953
954    #[test]
955    fn cache_builder_records_metrics() {
956        let mut entry = FileBuilderEntry {
957            builder: None,
958            remaining_ranges: 1,
959            waiters: Vec::new(),
960        };
961        let builder = Arc::new(FileRangeBuilder::default());
962        let mut reader_metrics = ReaderMetrics::default();
963
964        assert!(cache_builder_if_needed(
965            &mut entry,
966            &builder,
967            &mut reader_metrics,
968        ));
969        assert!(entry.builder.is_some());
970        assert_eq!(
971            reader_metrics.metadata_mem_size,
972            builder.memory_size() as isize
973        );
974        assert_eq!(reader_metrics.num_range_builders, 1);
975
976        if entry.builder.take().is_some() {
977            PRUNER_ACTIVE_BUILDERS.dec();
978        }
979    }
980
981    #[tokio::test]
982    async fn skip_file_range_decrements_and_clears_builder() {
983        let (_env, pruner) = make_test_pruner(1).await;
984
985        // Simulate 3 partition ranges for file 0.
986        let ranges: Vec<PartitionRange> = (0..3).map(|_| file_partition_range(0)).collect();
987        pruner.add_partition_ranges(&ranges);
988        assert_eq!(pruner.test_remaining_ranges(0), 3);
989
990        // Manually set a cached builder (simulating a previous cache hit).
991        {
992            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
993            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
994            PRUNER_ACTIVE_BUILDERS.inc();
995        }
996        assert!(pruner.test_has_builder(0));
997
998        // Skip all 3 ranges; the third should clear the builder.
999        let mut reader_metrics = ReaderMetrics::default();
1000        for i in 0..3 {
1001            let index = RowGroupIndex {
1002                index: 0,
1003                row_group_index: i as i64,
1004            };
1005            pruner.skip_file_range(index, &mut reader_metrics);
1006        }
1007
1008        assert_eq!(pruner.test_remaining_ranges(0), 0);
1009        assert!(!pruner.test_has_builder(0));
1010    }
1011
1012    #[tokio::test]
1013    async fn retained_builder_survives_after_last_range() {
1014        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1015        pruner.add_partition_ranges(&[file_partition_range(0)]);
1016        {
1017            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1018            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1019            PRUNER_ACTIVE_BUILDERS.inc();
1020        }
1021
1022        let mut reader_metrics = ReaderMetrics::default();
1023        pruner.skip_file_range(
1024            RowGroupIndex {
1025                index: 0,
1026                row_group_index: 0,
1027            },
1028            &mut reader_metrics,
1029        );
1030
1031        assert_eq!(pruner.test_remaining_ranges(0), 0);
1032        assert!(pruner.test_has_builder(0));
1033
1034        let partition_metrics = make_partition_metrics();
1035        let mut reader_metrics = ReaderMetrics::default();
1036        let _builder = pruner
1037            .get_file_builder(
1038                0,
1039                PreFilterMode::SkipFields,
1040                &partition_metrics,
1041                &mut reader_metrics,
1042            )
1043            .await
1044            .unwrap();
1045        assert_eq!(reader_metrics.filter_metrics.pruner_cache_hit, 1);
1046    }
1047
1048    #[tokio::test]
1049    async fn add_partition_ranges_keeps_retained_builder() {
1050        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1051        pruner.add_partition_ranges(&[file_partition_range(0)]);
1052        {
1053            let mut entry = pruner.inner.file_entries[0].lock().unwrap();
1054            entry.builder = Some(Arc::new(FileRangeBuilder::default()));
1055            PRUNER_ACTIVE_BUILDERS.inc();
1056        }
1057
1058        let mut reader_metrics = ReaderMetrics::default();
1059        pruner.skip_file_range(
1060            RowGroupIndex {
1061                index: 0,
1062                row_group_index: 0,
1063            },
1064            &mut reader_metrics,
1065        );
1066        assert!(pruner.test_has_builder(0));
1067
1068        pruner.add_partition_ranges(&[file_partition_range(0)]);
1069        assert!(pruner.test_has_builder(0));
1070        assert_eq!(pruner.test_remaining_ranges(0), 1);
1071    }
1072
1073    #[tokio::test]
1074    async fn retaining_mode_does_not_cache_after_skip_file_range_consumed_all() {
1075        let (_env, pruner) = make_test_pruner_with_retained_builders(1, true).await;
1076
1077        // Simulate one range for file 0.
1078        let ranges = vec![file_partition_range(0)];
1079        pruner.add_partition_ranges(&ranges);
1080        assert_eq!(pruner.test_remaining_ranges(0), 1);
1081
1082        // Simulate skip_file_range consuming the last range BEFORE the
1083        // background worker finishes. This mirrors the race: a dynamic filter
1084        // tightens and manifest-prune fast-skip zeros out remaining_ranges.
1085        let mut reader_metrics = ReaderMetrics::default();
1086        let index = RowGroupIndex {
1087            index: 0,
1088            row_group_index: 0,
1089        };
1090        pruner.skip_file_range(index, &mut reader_metrics);
1091        assert_eq!(pruner.test_remaining_ranges(0), 0);
1092        assert!(!pruner.test_has_builder(0));
1093
1094        // Now simulate the worker completing: check the caching guard.
1095        let entry = pruner.inner.file_entries[0].lock().unwrap();
1096        let should_cache = should_cache_builder(&entry);
1097        drop(entry);
1098
1099        assert!(!should_cache);
1100
1101        // Ensure the gauge was not incremented for a stale builder.
1102        // (skip_file_range already decremented it if there was one, but here
1103        // there was none, so the gauge should be at baseline.)
1104    }
1105
1106    #[tokio::test]
1107    async fn worker_caches_when_ranges_remain() {
1108        let (_env, pruner) = make_test_pruner(1).await;
1109
1110        // Simulate 2 ranges for file 0.
1111        let ranges: Vec<PartitionRange> = (0..2).map(|_| file_partition_range(0)).collect();
1112        pruner.add_partition_ranges(&ranges);
1113        assert_eq!(pruner.test_remaining_ranges(0), 2);
1114
1115        // Consume only 1 range.
1116        let mut reader_metrics = ReaderMetrics::default();
1117        let index = RowGroupIndex {
1118            index: 0,
1119            row_group_index: 0,
1120        };
1121        pruner.skip_file_range(index, &mut reader_metrics);
1122        assert_eq!(pruner.test_remaining_ranges(0), 1);
1123
1124        // The worker should still cache because remaining_ranges > 0.
1125        let entry = pruner.inner.file_entries[0].lock().unwrap();
1126        assert!(should_cache_builder(&entry));
1127    }
1128
1129    // ── Corner case: fast-skip across multiple row groups ─────────────
1130
1131    /// 1 file × 3 row groups, predicate `ts > 10000ms` prunes the file at
1132    /// manifest level. Fast-skipping each of the 3 row groups must return true
1133    /// and decrement remaining_ranges to 0.
1134    #[tokio::test]
1135    async fn try_skip_manifest_pruned_file_range_multi_row_groups() {
1136        let predicate_exprs: Vec<Expr> =
1137            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1138        let (_env, pruner) = make_test_pruner_with_predicate(1, 3, &predicate_exprs).await;
1139
1140        let ranges = pruner.inner.stream_ctx.partition_ranges();
1141        assert_eq!(ranges.len(), 3);
1142        pruner.add_partition_ranges(&ranges);
1143        assert_eq!(pruner.test_remaining_ranges(0), 3);
1144
1145        let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1146        let partition_metrics = make_partition_metrics();
1147
1148        // Fast-skip each of the 3 row groups.
1149        for rg in 0..3 {
1150            let index = RowGroupIndex {
1151                index: 0, // file_index == 0, no memtables
1152                row_group_index: rg,
1153            };
1154            let skipped =
1155                partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1156            assert!(skipped, "row group {} should be skipped", rg);
1157        }
1158
1159        // All refs consumed.
1160        assert_eq!(pruner.test_remaining_ranges(0), 0);
1161        // manifest_pruned_files is CAS'd exactly once (first call).
1162        assert!(pruner.test_is_manifest_pruned(0));
1163    }
1164
1165    /// A file whose manifest time range may contain matching rows must not be
1166    /// fast-skipped. This protects query correctness over metrics precision.
1167    #[tokio::test]
1168    async fn try_skip_manifest_pruned_file_range_keeps_overlapping_file() {
1169        let predicate_exprs: Vec<Expr> =
1170            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(500), None)))];
1171        let (_env, pruner) = make_test_pruner_with_predicate(1, 2, &predicate_exprs).await;
1172
1173        let ranges = pruner.inner.stream_ctx.partition_ranges();
1174        assert_eq!(ranges.len(), 2);
1175        pruner.add_partition_ranges(&ranges);
1176        assert_eq!(pruner.test_remaining_ranges(0), 2);
1177
1178        let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &ranges));
1179        let partition_metrics = make_partition_metrics();
1180        let range_meta = &pruner.inner.stream_ctx.ranges[ranges[0].identifier];
1181        let index = range_meta.row_group_indices[0];
1182
1183        let skipped =
1184            partition_pruner.try_skip_manifest_pruned_file_range(index, &partition_metrics);
1185
1186        assert!(!skipped);
1187        assert_eq!(pruner.test_remaining_ranges(0), 2);
1188        assert!(!pruner.test_is_manifest_pruned(0));
1189    }
1190
1191    // ── Corner case: add_partition_ranges resets the manifest-pruned flag ──
1192
1193    #[tokio::test]
1194    async fn add_partition_ranges_resets_manifest_pruned_flag() {
1195        let predicate_exprs: Vec<Expr> =
1196            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1197        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1198
1199        // Mark file 0 as manifest-pruned.
1200        let mut reader_metrics = ReaderMetrics::default();
1201        let marked = pruner
1202            .inner
1203            .try_mark_manifest_pruned(0, &mut reader_metrics);
1204        assert!(marked);
1205        assert!(pruner.test_is_manifest_pruned(0));
1206        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1207
1208        // Calling add_partition_ranges must reset the flag.
1209        let ranges = vec![file_partition_range(0)];
1210        pruner.add_partition_ranges(&ranges);
1211        assert!(!pruner.test_is_manifest_pruned(0));
1212        // remaining_ranges was also incremented.
1213        assert_eq!(pruner.test_remaining_ranges(0), 1);
1214    }
1215
1216    // ── Corner case: prune_file_directly short-circuits via manifest prune ──
1217
1218    #[tokio::test]
1219    async fn prune_file_directly_manifest_pruned_returns_empty_builder() {
1220        let predicate_exprs: Vec<Expr> =
1221            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1222        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1223
1224        // Ensure there is no cached builder yet.
1225        assert!(!pruner.test_has_builder(0));
1226
1227        let mut reader_metrics = ReaderMetrics::default();
1228        let builder = pruner
1229            .prune_file_directly(0, PreFilterMode::SkipFields, &mut reader_metrics)
1230            .await
1231            .unwrap();
1232
1233        // Should be the default (empty) builder.
1234        assert_eq!(
1235            builder.memory_size(),
1236            FileRangeBuilder::default().memory_size()
1237        );
1238        // builder must NOT be cached — the manifest-pruned flag is enough.
1239        assert!(!pruner.test_has_builder(0));
1240        // files_time_range_pruned was recorded.
1241        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1242    }
1243
1244    // ── Corner case: try_mark_manifest_pruned does not double-count ───
1245
1246    #[tokio::test]
1247    async fn try_mark_manifest_pruned_only_counts_first_cas() {
1248        let predicate_exprs: Vec<Expr> =
1249            vec![col("ts").gt(lit(ScalarValue::TimestampMillisecond(Some(10_000), None)))];
1250        let (_env, pruner) = make_test_pruner_with_predicate(1, 1, &predicate_exprs).await;
1251
1252        // First call: CAS succeeds, metric incremented.
1253        let mut reader_metrics = ReaderMetrics::default();
1254        let marked = pruner
1255            .inner
1256            .try_mark_manifest_pruned(0, &mut reader_metrics);
1257        assert!(marked);
1258        assert_eq!(reader_metrics.filter_metrics.files_time_range_pruned, 1);
1259        assert!(pruner.test_is_manifest_pruned(0));
1260
1261        // Second call: already true, no metric delta.
1262        let mut reader_metrics2 = ReaderMetrics::default();
1263        let marked2 = pruner
1264            .inner
1265            .try_mark_manifest_pruned(0, &mut reader_metrics2);
1266        assert!(marked2);
1267        assert_eq!(reader_metrics2.filter_metrics.files_time_range_pruned, 0);
1268    }
1269}