1use 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#[doc(hidden)]
16#[derive(Debug, Clone, Copy, PartialEq, Eq)]
17pub struct ParquetRangePlanningDiagnosticSnapshot {
18 pub max_concurrent_physical_range_requests: u64,
20 pub physical_range_request_waves_planned: u64,
22 pub successful_plan_time_micros: u64,
24}
25
26#[non_exhaustive]
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct DeltaScanMetricsSnapshot {
33 pub snapshot_version: u64,
35 pub parquet_backend: ParquetReaderBackend,
37 pub scan_partitions_planned: u64,
39 pub files_planned: u64,
41 pub add_actions_excluded_during_planning: Option<u64>,
43 pub estimated_input_rows: Option<u64>,
45 pub estimated_input_bytes: Option<u64>,
47 pub scan_partitions_started: u64,
49 pub scan_partitions_completed: u64,
51 pub file_tasks_started: u64,
53 pub file_tasks_completed: u64,
55 pub scheduler_batches_emitted: u64,
57 pub scheduler_rows_emitted: u64,
59 pub deletion_vector_payloads_loaded: u64,
61 pub deletion_vectors_applied: u64,
63 pub deletion_vector_rows_deleted: u64,
65 pub deletion_vector_failures: u64,
67 pub deletion_vector_coordinate_rejections: u64,
69 pub parquet_data_file_exact_ranges_requested: Option<u64>,
71 pub parquet_data_file_exact_range_bytes_requested: Option<u64>,
73 pub parquet_data_file_physical_range_requests_planned: Option<u64>,
75 pub parquet_data_file_physical_range_bytes_planned: Option<u64>,
77 pub parquet_data_file_cold_start_range_plans: Option<u64>,
79 pub parquet_data_file_cost_based_exact_range_plans: Option<u64>,
81 pub parquet_data_file_cost_based_merged_range_plans: Option<u64>,
83 pub parquet_data_file_store_delegated_range_plans: Option<u64>,
85 pub parquet_data_file_range_get_operations: Option<u64>,
87 pub parquet_data_file_full_get_operations: Option<u64>,
89 pub parquet_data_file_bytes_received: Option<u64>,
91 pub estimated_parquet_task_bytes_admitted: Option<u64>,
93}
94
95#[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 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}