Skip to main content

delta_arrow_reader/reader/
metrics.rs

1//! Delta scan planning and execution metrics.
2
3use std::{
4    fmt,
5    sync::{
6        Arc,
7        atomic::{AtomicU64, Ordering},
8    },
9};
10
11use super::options::MAX_CONCURRENT_PARQUET_RANGE_READS;
12use super::options::ParquetReaderBackend;
13
14/// Internal measurements used to validate Parquet range-planning benchmarks.
15#[doc(hidden)]
16#[derive(Debug, Clone, Copy, PartialEq, Eq)]
17pub struct ParquetRangePlanningDiagnosticSnapshot {
18    /// Maximum physical range requests executed concurrently by one plan.
19    pub max_concurrent_physical_range_requests: u64,
20    /// Sum of physical request waves selected by visible plans.
21    pub physical_range_request_waves_planned: u64,
22    /// Sum of successful multi-range plan execution time in microseconds.
23    pub successful_plan_time_micros: u64,
24}
25
26/// Immutable point-in-time copy of one Delta scan's metrics.
27///
28/// A snapshot taken while the scan is running may contain partial progress. It does not change
29/// after creation; call [`DeltaScanMetrics::snapshot`] again to observe later progress.
30#[non_exhaustive]
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct DeltaScanMetricsSnapshot {
33    /// Delta snapshot version selected for the scan.
34    pub snapshot_version: u64,
35    /// Parquet reader backend selected for the scan.
36    pub parquet_backend: ParquetReaderBackend,
37    /// Final execution partitions planned for the scan, including source repartitioning.
38    pub scan_partitions_planned: u64,
39    /// Data files selected during planning.
40    pub files_planned: u64,
41    /// Add actions excluded during planning, when known.
42    pub add_actions_excluded_during_planning: Option<u64>,
43    /// Estimated rows in selected input files before row filtering, when known for every file.
44    pub estimated_input_rows: Option<u64>,
45    /// Estimated bytes in selected input files, when every file reported a size.
46    pub estimated_input_bytes: Option<u64>,
47    /// Scan partitions whose execution started.
48    pub scan_partitions_started: u64,
49    /// Scan partitions that completed normally.
50    pub scan_partitions_completed: u64,
51    /// Data-file tasks that started, including independently read file ranges.
52    pub file_tasks_started: u64,
53    /// Data-file tasks that completed normally, including independently read file ranges.
54    pub file_tasks_completed: u64,
55    /// Batches emitted by the scheduler before final stream operations.
56    pub scheduler_batches_emitted: u64,
57    /// Rows emitted by the scheduler before final stream operations.
58    pub scheduler_rows_emitted: u64,
59    /// Deletion-vector payloads loaded.
60    pub deletion_vector_payloads_loaded: u64,
61    /// Deletion-vector masks applied.
62    pub deletion_vectors_applied: u64,
63    /// Rows removed by deletion-vector masks.
64    pub deletion_vector_rows_deleted: u64,
65    /// Deletion-vector read or masking failures.
66    pub deletion_vector_failures: u64,
67    /// Deletion-vector coordinate operations rejected by safety checks.
68    pub deletion_vector_coordinate_rejections: u64,
69    /// Normalized exact Parquet ranges requested through direct multi-range calls.
70    pub parquet_data_file_exact_ranges_requested: Option<u64>,
71    /// Bytes in those normalized exact Parquet ranges.
72    pub parquet_data_file_exact_range_bytes_requested: Option<u64>,
73    /// Physical range requests selected by the automatic planner.
74    pub parquet_data_file_physical_range_requests_planned: Option<u64>,
75    /// Bytes covered by those automatically planned physical range requests.
76    pub parquet_data_file_physical_range_bytes_planned: Option<u64>,
77    /// Automatic range plans selected without a usable transport estimate.
78    pub parquet_data_file_cold_start_range_plans: Option<u64>,
79    /// Automatic range plans where an estimate favored the normalized exact ranges.
80    pub parquet_data_file_cost_based_exact_range_plans: Option<u64>,
81    /// Automatic range plans where an estimate favored merging gaps between exact ranges.
82    pub parquet_data_file_cost_based_merged_range_plans: Option<u64>,
83    /// Range-planning decisions delegated to the store's own implementation.
84    pub parquet_data_file_store_delegated_range_plans: Option<u64>,
85    /// Direct Parquet ranged GET operations, or `None` for another backend.
86    pub parquet_data_file_range_get_operations: Option<u64>,
87    /// Direct Parquet full GET operations, or `None` for another backend.
88    pub parquet_data_file_full_get_operations: Option<u64>,
89    /// Direct Parquet payload bytes received, or `None` for another backend.
90    pub parquet_data_file_bytes_received: Option<u64>,
91    /// Estimated bytes admitted across direct Parquet tasks, or `None` for another backend.
92    pub estimated_parquet_task_bytes_admitted: Option<u64>,
93}
94
95/// Shared live metrics for one Delta scan.
96#[derive(Clone)]
97pub struct DeltaScanMetrics {
98    inner: Arc<DeltaScanMetricsInner>,
99}
100
101impl fmt::Debug for DeltaScanMetrics {
102    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
103        formatter
104            .debug_struct("DeltaScanMetrics")
105            .finish_non_exhaustive()
106    }
107}
108
109struct DeltaScanMetricsInner {
110    snapshot_version: u64,
111    parquet_backend: ParquetReaderBackend,
112    scan_partitions_planned: AtomicU64,
113    files_planned: u64,
114    add_actions_excluded_during_planning: Option<u64>,
115    estimated_input_rows: Option<u64>,
116    estimated_input_bytes: Option<u64>,
117    scan_partitions_started: AtomicU64,
118    scan_partitions_completed: AtomicU64,
119    file_tasks_started: AtomicU64,
120    file_tasks_completed: AtomicU64,
121    scheduler_batches_emitted: AtomicU64,
122    scheduler_rows_emitted: AtomicU64,
123    deletion_vector_payloads_loaded: AtomicU64,
124    deletion_vectors_applied: AtomicU64,
125    deletion_vector_rows_deleted: AtomicU64,
126    deletion_vector_failures: AtomicU64,
127    deletion_vector_coordinate_rejections: AtomicU64,
128    parquet_data_file_exact_ranges_requested: AtomicU64,
129    parquet_data_file_exact_range_bytes_requested: AtomicU64,
130    parquet_data_file_physical_range_requests_planned: AtomicU64,
131    parquet_data_file_physical_range_bytes_planned: AtomicU64,
132    parquet_data_file_cold_start_range_plans: AtomicU64,
133    parquet_data_file_cost_based_exact_range_plans: AtomicU64,
134    parquet_data_file_cost_based_merged_range_plans: AtomicU64,
135    parquet_data_file_store_delegated_range_plans: AtomicU64,
136    parquet_range_request_waves_planned: AtomicU64,
137    parquet_range_successful_plan_time_micros: AtomicU64,
138    parquet_data_file_range_get_operations: AtomicU64,
139    parquet_data_file_full_get_operations: AtomicU64,
140    parquet_data_file_bytes_received: AtomicU64,
141    estimated_parquet_task_bytes_admitted: AtomicU64,
142}
143
144#[allow(dead_code)]
145pub(crate) struct DeltaScanMetricsConfig {
146    pub(crate) snapshot_version: u64,
147    pub(crate) parquet_backend: ParquetReaderBackend,
148    pub(crate) scan_partitions_planned: usize,
149    pub(crate) files_planned: usize,
150    pub(crate) add_actions_excluded_during_planning: Option<u64>,
151    pub(crate) estimated_input_rows: Option<u64>,
152    pub(crate) estimated_input_bytes: Option<u64>,
153}
154
155impl DeltaScanMetrics {
156    #[allow(dead_code)]
157    pub(crate) fn new(config: DeltaScanMetricsConfig) -> Self {
158        Self {
159            inner: Arc::new(DeltaScanMetricsInner {
160                snapshot_version: config.snapshot_version,
161                parquet_backend: config.parquet_backend,
162                scan_partitions_planned: AtomicU64::new(usize_to_u64_saturating(
163                    config.scan_partitions_planned,
164                )),
165                files_planned: usize_to_u64_saturating(config.files_planned),
166                add_actions_excluded_during_planning: config.add_actions_excluded_during_planning,
167                estimated_input_rows: config.estimated_input_rows,
168                estimated_input_bytes: config.estimated_input_bytes,
169                scan_partitions_started: AtomicU64::new(0),
170                scan_partitions_completed: AtomicU64::new(0),
171                file_tasks_started: AtomicU64::new(0),
172                file_tasks_completed: AtomicU64::new(0),
173                scheduler_batches_emitted: AtomicU64::new(0),
174                scheduler_rows_emitted: AtomicU64::new(0),
175                deletion_vector_payloads_loaded: AtomicU64::new(0),
176                deletion_vectors_applied: AtomicU64::new(0),
177                deletion_vector_rows_deleted: AtomicU64::new(0),
178                deletion_vector_failures: AtomicU64::new(0),
179                deletion_vector_coordinate_rejections: AtomicU64::new(0),
180                parquet_data_file_exact_ranges_requested: AtomicU64::new(0),
181                parquet_data_file_exact_range_bytes_requested: AtomicU64::new(0),
182                parquet_data_file_physical_range_requests_planned: AtomicU64::new(0),
183                parquet_data_file_physical_range_bytes_planned: AtomicU64::new(0),
184                parquet_data_file_cold_start_range_plans: AtomicU64::new(0),
185                parquet_data_file_cost_based_exact_range_plans: AtomicU64::new(0),
186                parquet_data_file_cost_based_merged_range_plans: AtomicU64::new(0),
187                parquet_data_file_store_delegated_range_plans: AtomicU64::new(0),
188                parquet_range_request_waves_planned: AtomicU64::new(0),
189                parquet_range_successful_plan_time_micros: AtomicU64::new(0),
190                parquet_data_file_range_get_operations: AtomicU64::new(0),
191                parquet_data_file_full_get_operations: AtomicU64::new(0),
192                parquet_data_file_bytes_received: AtomicU64::new(0),
193                estimated_parquet_task_bytes_admitted: AtomicU64::new(0),
194            }),
195        }
196    }
197
198    /// Returns an immutable point-in-time copy of all scan metrics.
199    pub fn snapshot(&self) -> DeltaScanMetricsSnapshot {
200        let inner = self.inner.as_ref();
201        DeltaScanMetricsSnapshot {
202            snapshot_version: inner.snapshot_version,
203            parquet_backend: inner.parquet_backend,
204            scan_partitions_planned: load(&inner.scan_partitions_planned),
205            files_planned: inner.files_planned,
206            add_actions_excluded_during_planning: inner.add_actions_excluded_during_planning,
207            estimated_input_rows: inner.estimated_input_rows,
208            estimated_input_bytes: inner.estimated_input_bytes,
209            scan_partitions_started: load(&inner.scan_partitions_started),
210            scan_partitions_completed: load(&inner.scan_partitions_completed),
211            file_tasks_started: load(&inner.file_tasks_started),
212            file_tasks_completed: load(&inner.file_tasks_completed),
213            scheduler_batches_emitted: load(&inner.scheduler_batches_emitted),
214            scheduler_rows_emitted: load(&inner.scheduler_rows_emitted),
215            deletion_vector_payloads_loaded: load(&inner.deletion_vector_payloads_loaded),
216            deletion_vectors_applied: load(&inner.deletion_vectors_applied),
217            deletion_vector_rows_deleted: load(&inner.deletion_vector_rows_deleted),
218            deletion_vector_failures: load(&inner.deletion_vector_failures),
219            deletion_vector_coordinate_rejections: load(
220                &inner.deletion_vector_coordinate_rejections,
221            ),
222            parquet_data_file_exact_ranges_requested: self
223                .parquet_metric(&inner.parquet_data_file_exact_ranges_requested),
224            parquet_data_file_exact_range_bytes_requested: self
225                .parquet_metric(&inner.parquet_data_file_exact_range_bytes_requested),
226            parquet_data_file_physical_range_requests_planned: self
227                .parquet_metric(&inner.parquet_data_file_physical_range_requests_planned),
228            parquet_data_file_physical_range_bytes_planned: self
229                .parquet_metric(&inner.parquet_data_file_physical_range_bytes_planned),
230            parquet_data_file_cold_start_range_plans: self
231                .parquet_metric(&inner.parquet_data_file_cold_start_range_plans),
232            parquet_data_file_cost_based_exact_range_plans: self
233                .parquet_metric(&inner.parquet_data_file_cost_based_exact_range_plans),
234            parquet_data_file_cost_based_merged_range_plans: self
235                .parquet_metric(&inner.parquet_data_file_cost_based_merged_range_plans),
236            parquet_data_file_store_delegated_range_plans: self
237                .parquet_metric(&inner.parquet_data_file_store_delegated_range_plans),
238            parquet_data_file_range_get_operations: self
239                .parquet_metric(&inner.parquet_data_file_range_get_operations),
240            parquet_data_file_full_get_operations: self
241                .parquet_metric(&inner.parquet_data_file_full_get_operations),
242            parquet_data_file_bytes_received: self
243                .parquet_metric(&inner.parquet_data_file_bytes_received),
244            estimated_parquet_task_bytes_admitted: self
245                .parquet_metric(&inner.estimated_parquet_task_bytes_admitted),
246        }
247    }
248
249    fn parquet_metric(&self, counter: &AtomicU64) -> Option<u64> {
250        match self.inner.parquet_backend {
251            ParquetReaderBackend::Direct => Some(load(counter)),
252            ParquetReaderBackend::DeltaKernel => None,
253        }
254    }
255
256    pub(crate) fn parquet_range_planning_diagnostic_snapshot(
257        &self,
258    ) -> ParquetRangePlanningDiagnosticSnapshot {
259        ParquetRangePlanningDiagnosticSnapshot {
260            max_concurrent_physical_range_requests: usize_to_u64_saturating(
261                MAX_CONCURRENT_PARQUET_RANGE_READS,
262            ),
263            physical_range_request_waves_planned: load(
264                &self.inner.parquet_range_request_waves_planned,
265            ),
266            successful_plan_time_micros: load(
267                &self.inner.parquet_range_successful_plan_time_micros,
268            ),
269        }
270    }
271
272    #[allow(dead_code)]
273    pub(crate) fn record_scan_partitions_planned(&self, value: usize) {
274        self.inner
275            .scan_partitions_planned
276            .store(usize_to_u64_saturating(value), Ordering::Relaxed);
277    }
278
279    #[allow(dead_code)]
280    pub(crate) fn record_scan_partition_started(&self) {
281        saturating_fetch_add(&self.inner.scan_partitions_started, 1);
282    }
283
284    #[allow(dead_code)]
285    pub(crate) fn record_scan_partition_completed(&self) {
286        saturating_fetch_add(&self.inner.scan_partitions_completed, 1);
287    }
288
289    #[allow(dead_code)]
290    pub(crate) fn record_file_task_started(&self) {
291        saturating_fetch_add(&self.inner.file_tasks_started, 1);
292    }
293
294    #[allow(dead_code)]
295    pub(crate) fn record_file_task_completed(&self) {
296        saturating_fetch_add(&self.inner.file_tasks_completed, 1);
297    }
298
299    #[allow(dead_code)]
300    pub(crate) fn record_scheduler_batch_emitted(&self, rows: usize) {
301        saturating_fetch_add(&self.inner.scheduler_batches_emitted, 1);
302        saturating_fetch_add(
303            &self.inner.scheduler_rows_emitted,
304            usize_to_u64_saturating(rows),
305        );
306    }
307
308    #[allow(dead_code)]
309    pub(crate) fn record_deletion_vector_payload_loaded(&self) {
310        saturating_fetch_add(&self.inner.deletion_vector_payloads_loaded, 1);
311    }
312
313    #[allow(dead_code)]
314    pub(crate) fn record_deletion_vector_applied(&self) {
315        saturating_fetch_add(&self.inner.deletion_vectors_applied, 1);
316    }
317
318    #[allow(dead_code)]
319    pub(crate) fn record_deletion_vector_rows_deleted(&self, rows: usize) {
320        saturating_fetch_add(
321            &self.inner.deletion_vector_rows_deleted,
322            usize_to_u64_saturating(rows),
323        );
324    }
325
326    #[allow(dead_code)]
327    pub(crate) fn record_deletion_vector_failure(&self) {
328        saturating_fetch_add(&self.inner.deletion_vector_failures, 1);
329    }
330
331    #[allow(dead_code)]
332    pub(crate) fn record_deletion_vector_coordinate_rejection(&self) {
333        saturating_fetch_add(&self.inner.deletion_vector_coordinate_rejections, 1);
334    }
335
336    pub(crate) fn record_parquet_data_file_exact_ranges_requested(
337        &self,
338        range_count: usize,
339        bytes: u128,
340    ) {
341        saturating_fetch_add(
342            &self.inner.parquet_data_file_exact_ranges_requested,
343            usize_to_u64_saturating(range_count),
344        );
345        saturating_fetch_add(
346            &self.inner.parquet_data_file_exact_range_bytes_requested,
347            u128_to_u64_saturating(bytes),
348        );
349    }
350
351    pub(crate) fn record_parquet_data_file_physical_range_plan(
352        &self,
353        request_count: usize,
354        bytes: u128,
355    ) {
356        saturating_fetch_add(
357            &self.inner.parquet_data_file_physical_range_requests_planned,
358            usize_to_u64_saturating(request_count),
359        );
360        saturating_fetch_add(
361            &self.inner.parquet_data_file_physical_range_bytes_planned,
362            u128_to_u64_saturating(bytes),
363        );
364        saturating_fetch_add(
365            &self.inner.parquet_range_request_waves_planned,
366            usize_to_u64_saturating(request_count.div_ceil(MAX_CONCURRENT_PARQUET_RANGE_READS)),
367        );
368    }
369
370    pub(crate) fn record_parquet_range_successful_plan_time(&self, elapsed: std::time::Duration) {
371        saturating_fetch_add(
372            &self.inner.parquet_range_successful_plan_time_micros,
373            u128_to_u64_saturating(elapsed.as_micros()),
374        );
375    }
376
377    pub(crate) fn record_parquet_data_file_cold_start_range_plan(&self) {
378        saturating_fetch_add(&self.inner.parquet_data_file_cold_start_range_plans, 1);
379    }
380
381    pub(crate) fn record_parquet_data_file_cost_based_exact_range_plan(&self) {
382        saturating_fetch_add(
383            &self.inner.parquet_data_file_cost_based_exact_range_plans,
384            1,
385        );
386    }
387
388    pub(crate) fn record_parquet_data_file_cost_based_merged_range_plan(&self) {
389        saturating_fetch_add(
390            &self.inner.parquet_data_file_cost_based_merged_range_plans,
391            1,
392        );
393    }
394
395    pub(crate) fn record_parquet_data_file_store_delegated_range_plan(&self) {
396        saturating_fetch_add(&self.inner.parquet_data_file_store_delegated_range_plans, 1);
397    }
398
399    pub(crate) fn record_parquet_data_file_range_get_operation(&self) {
400        saturating_fetch_add(&self.inner.parquet_data_file_range_get_operations, 1);
401    }
402
403    pub(crate) fn record_parquet_data_file_full_get_operation(&self) {
404        saturating_fetch_add(&self.inner.parquet_data_file_full_get_operations, 1);
405    }
406
407    pub(crate) fn record_parquet_data_file_bytes_received(&self, bytes: usize) {
408        saturating_fetch_add(
409            &self.inner.parquet_data_file_bytes_received,
410            usize_to_u64_saturating(bytes),
411        );
412    }
413
414    pub(crate) fn record_estimated_parquet_task_bytes_admitted(&self, bytes: u64) {
415        saturating_fetch_add(&self.inner.estimated_parquet_task_bytes_admitted, bytes);
416    }
417}
418
419fn load(counter: &AtomicU64) -> u64 {
420    counter.load(Ordering::Relaxed)
421}
422
423#[allow(dead_code)]
424pub(crate) fn saturating_fetch_add(counter: &AtomicU64, value: u64) {
425    let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
426        Some(current.saturating_add(value))
427    });
428}
429
430fn usize_to_u64_saturating(value: usize) -> u64 {
431    u64::try_from(value).unwrap_or(u64::MAX)
432}
433
434fn u128_to_u64_saturating(value: u128) -> u64 {
435    u64::try_from(value).unwrap_or(u64::MAX)
436}
437
438#[cfg(test)]
439mod tests {
440    use std::{sync::atomic::Ordering, thread};
441
442    use super::{DeltaScanMetrics, DeltaScanMetricsConfig, saturating_fetch_add};
443    use crate::ParquetReaderBackend;
444
445    fn metrics(parquet_backend: ParquetReaderBackend) -> DeltaScanMetrics {
446        DeltaScanMetrics::new(DeltaScanMetricsConfig {
447            snapshot_version: 7,
448            parquet_backend,
449            scan_partitions_planned: 3,
450            files_planned: 5,
451            add_actions_excluded_during_planning: Some(2),
452            estimated_input_rows: Some(99),
453            estimated_input_bytes: Some(42),
454        })
455    }
456
457    #[test]
458    fn snapshot_has_context_zeroes_and_backend_availability() {
459        let direct = metrics(ParquetReaderBackend::Direct).snapshot();
460        assert_eq!(direct.snapshot_version, 7);
461        assert_eq!(direct.parquet_backend, ParquetReaderBackend::Direct);
462        assert_eq!(direct.scan_partitions_planned, 3);
463        assert_eq!(direct.files_planned, 5);
464        assert_eq!(direct.add_actions_excluded_during_planning, Some(2));
465        assert_eq!(direct.estimated_input_rows, Some(99));
466        assert_eq!(direct.estimated_input_bytes, Some(42));
467        assert_eq!(direct.scan_partitions_started, 0);
468        assert_eq!(direct.scan_partitions_completed, 0);
469        assert_eq!(direct.file_tasks_started, 0);
470        assert_eq!(direct.file_tasks_completed, 0);
471        assert_eq!(direct.scheduler_batches_emitted, 0);
472        assert_eq!(direct.scheduler_rows_emitted, 0);
473        assert_eq!(direct.deletion_vector_payloads_loaded, 0);
474        assert_eq!(direct.deletion_vectors_applied, 0);
475        assert_eq!(direct.deletion_vector_rows_deleted, 0);
476        assert_eq!(direct.deletion_vector_failures, 0);
477        assert_eq!(direct.deletion_vector_coordinate_rejections, 0);
478        assert_eq!(direct.parquet_data_file_exact_ranges_requested, Some(0));
479        assert_eq!(
480            direct.parquet_data_file_exact_range_bytes_requested,
481            Some(0)
482        );
483        assert_eq!(
484            direct.parquet_data_file_physical_range_requests_planned,
485            Some(0)
486        );
487        assert_eq!(
488            direct.parquet_data_file_physical_range_bytes_planned,
489            Some(0)
490        );
491        assert_eq!(direct.parquet_data_file_cold_start_range_plans, Some(0));
492        assert_eq!(
493            direct.parquet_data_file_cost_based_exact_range_plans,
494            Some(0)
495        );
496        assert_eq!(
497            direct.parquet_data_file_cost_based_merged_range_plans,
498            Some(0)
499        );
500        assert_eq!(
501            direct.parquet_data_file_store_delegated_range_plans,
502            Some(0)
503        );
504        assert_eq!(direct.parquet_data_file_range_get_operations, Some(0));
505        assert_eq!(direct.parquet_data_file_full_get_operations, Some(0));
506        assert_eq!(direct.parquet_data_file_bytes_received, Some(0));
507        assert_eq!(direct.estimated_parquet_task_bytes_admitted, Some(0));
508
509        let kernel = metrics(ParquetReaderBackend::DeltaKernel).snapshot();
510        assert_eq!(kernel.parquet_data_file_exact_ranges_requested, None);
511        assert_eq!(kernel.parquet_data_file_exact_range_bytes_requested, None);
512        assert_eq!(
513            kernel.parquet_data_file_physical_range_requests_planned,
514            None
515        );
516        assert_eq!(kernel.parquet_data_file_physical_range_bytes_planned, None);
517        assert_eq!(kernel.parquet_data_file_cold_start_range_plans, None);
518        assert_eq!(kernel.parquet_data_file_cost_based_exact_range_plans, None);
519        assert_eq!(kernel.parquet_data_file_cost_based_merged_range_plans, None);
520        assert_eq!(kernel.parquet_data_file_store_delegated_range_plans, None);
521        assert_eq!(kernel.parquet_data_file_range_get_operations, None);
522        assert_eq!(kernel.parquet_data_file_full_get_operations, None);
523        assert_eq!(kernel.parquet_data_file_bytes_received, None);
524        assert_eq!(kernel.estimated_parquet_task_bytes_admitted, None);
525    }
526
527    #[test]
528    fn debug_output_is_safe_and_redacted() {
529        assert_eq!(
530            format!("{:?}", metrics(ParquetReaderBackend::Direct)),
531            "DeltaScanMetrics { .. }"
532        );
533    }
534
535    #[test]
536    fn snapshot_maps_live_counters() {
537        let metrics = metrics(ParquetReaderBackend::Direct);
538        metrics.record_scan_partitions_planned(16);
539        let counters = [
540            &metrics.inner.scan_partitions_started,
541            &metrics.inner.scan_partitions_completed,
542            &metrics.inner.file_tasks_started,
543            &metrics.inner.file_tasks_completed,
544            &metrics.inner.scheduler_batches_emitted,
545            &metrics.inner.scheduler_rows_emitted,
546            &metrics.inner.deletion_vector_payloads_loaded,
547            &metrics.inner.deletion_vectors_applied,
548            &metrics.inner.deletion_vector_rows_deleted,
549            &metrics.inner.deletion_vector_failures,
550            &metrics.inner.deletion_vector_coordinate_rejections,
551            &metrics.inner.parquet_data_file_exact_ranges_requested,
552            &metrics.inner.parquet_data_file_exact_range_bytes_requested,
553            &metrics
554                .inner
555                .parquet_data_file_physical_range_requests_planned,
556            &metrics.inner.parquet_data_file_physical_range_bytes_planned,
557            &metrics.inner.parquet_data_file_cold_start_range_plans,
558            &metrics.inner.parquet_data_file_cost_based_exact_range_plans,
559            &metrics
560                .inner
561                .parquet_data_file_cost_based_merged_range_plans,
562            &metrics.inner.parquet_data_file_store_delegated_range_plans,
563            &metrics.inner.parquet_data_file_range_get_operations,
564            &metrics.inner.parquet_data_file_full_get_operations,
565            &metrics.inner.parquet_data_file_bytes_received,
566            &metrics.inner.estimated_parquet_task_bytes_admitted,
567        ];
568        for (index, counter) in counters.into_iter().enumerate() {
569            saturating_fetch_add(counter, u64::try_from(index + 1).expect("small test value"));
570        }
571
572        let snapshot = metrics.snapshot();
573        assert_eq!(snapshot.scan_partitions_planned, 16);
574        assert_eq!(snapshot.scan_partitions_started, 1);
575        assert_eq!(snapshot.scan_partitions_completed, 2);
576        assert_eq!(snapshot.file_tasks_started, 3);
577        assert_eq!(snapshot.file_tasks_completed, 4);
578        assert_eq!(snapshot.scheduler_batches_emitted, 5);
579        assert_eq!(snapshot.scheduler_rows_emitted, 6);
580        assert_eq!(snapshot.deletion_vector_payloads_loaded, 7);
581        assert_eq!(snapshot.deletion_vectors_applied, 8);
582        assert_eq!(snapshot.deletion_vector_rows_deleted, 9);
583        assert_eq!(snapshot.deletion_vector_failures, 10);
584        assert_eq!(snapshot.deletion_vector_coordinate_rejections, 11);
585        assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(12));
586        assert_eq!(
587            snapshot.parquet_data_file_exact_range_bytes_requested,
588            Some(13)
589        );
590        assert_eq!(
591            snapshot.parquet_data_file_physical_range_requests_planned,
592            Some(14)
593        );
594        assert_eq!(
595            snapshot.parquet_data_file_physical_range_bytes_planned,
596            Some(15)
597        );
598        assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(16));
599        assert_eq!(
600            snapshot.parquet_data_file_cost_based_exact_range_plans,
601            Some(17)
602        );
603        assert_eq!(
604            snapshot.parquet_data_file_cost_based_merged_range_plans,
605            Some(18)
606        );
607        assert_eq!(
608            snapshot.parquet_data_file_store_delegated_range_plans,
609            Some(19)
610        );
611        assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(20));
612        assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(21));
613        assert_eq!(snapshot.parquet_data_file_bytes_received, Some(22));
614        assert_eq!(snapshot.estimated_parquet_task_bytes_admitted, Some(23));
615    }
616
617    #[test]
618    fn cloned_handles_saturate_under_concurrent_updates() -> Result<(), &'static str> {
619        let metrics = metrics(ParquetReaderBackend::Direct);
620        metrics
621            .inner
622            .file_tasks_started
623            .store(u64::MAX - 1, Ordering::Relaxed);
624        let workers = (0..4)
625            .map(|_| {
626                let metrics = metrics.clone();
627                thread::spawn(move || {
628                    saturating_fetch_add(&metrics.inner.file_tasks_started, 1);
629                })
630            })
631            .collect::<Vec<_>>();
632
633        for worker in workers {
634            worker.join().map_err(|_| "metrics worker panicked")?;
635        }
636
637        assert_eq!(metrics.snapshot().file_tasks_started, u64::MAX);
638        Ok(())
639    }
640}