1use crate::app::background::{LenCount, OwedCount};
5use crate::app::jobs::{Answer, Job};
6use crate::table::DataTableState;
7use crate::{App, AppEvent, logging};
8use std::sync::Arc;
9
10#[derive(Default)]
12pub struct Counting {
13 pub(crate) footer_progress: Arc<crate::formats::schema_union::FooterProgress>,
16 pub(crate) footers_this_frame: Option<(usize, usize)>,
20 pub(crate) loaded_ahead_from: Option<(u64, usize, usize, usize)>,
22 pub(crate) len_count_inflight: Option<u64>,
25 pub(crate) count_after_paint: Option<u64>,
29 #[cfg(test)]
31 pub(crate) counts_spawned: std::cell::Cell<usize>,
32 #[cfg(test)]
35 pub(crate) first_rows_asked: usize,
36 pub(crate) len_count_failed: Option<u64>,
39 pub(crate) end_after_count: Option<u64>,
42 pub(crate) end_when_the_footers_land: Option<u64>,
46 pub(crate) end_when_indexed: Option<u64>,
49 pub(crate) indexing_stop: Arc<std::sync::atomic::AtomicBool>,
51 pub(crate) indexing_lines: Option<Arc<crate::formats::lines::Lines>>,
53 pub(crate) indexing_paused: bool,
55 pub(crate) goto_when_indexed: Option<(u64, usize)>,
57 pub(crate) count_progress: Arc<crate::formats::schema_union::FooterProgress>,
59 pub(crate) exact_count_asked: Option<u64>,
62 pub(crate) count_after_stop: Option<u64>,
65 pub(crate) footers_held: Option<(u64, crate::table::FootersFound)>,
69 pub(crate) followed_fields_held: Option<(u64, Vec<polars::prelude::Field>)>,
71 pub(crate) reread_owed: Option<u64>,
75 pub(crate) listed_this_frame: Option<usize>,
78}
79
80impl Counting {
81 pub(crate) fn reset_for_dataset(
86 &mut self,
87 footers: Arc<crate::formats::schema_union::FooterProgress>,
88 ) {
89 self.stop_footer_pass();
90 self.footer_progress = footers;
91 self.end_when_the_footers_land = None;
92 self.end_after_count = None;
93 self.end_when_indexed = None;
94 self.goto_when_indexed = None;
95 }
96
97 pub(crate) fn stop_footer_pass(&self) {
99 self.footer_progress.cancel();
100 }
101
102 pub(crate) fn stop_indexing(&mut self) {
104 self.indexing_stop
105 .store(true, std::sync::atomic::Ordering::Relaxed);
106 if let Some(lines) = self.indexing_lines.take() {
107 lines.stop_indexing();
108 }
109 }
110
111 pub(crate) fn pause_indexing(&mut self) {
114 if self.indexing_lines.is_some() {
115 self.indexing_stop
116 .store(true, std::sync::atomic::Ordering::Relaxed);
117 self.indexing_paused = true;
118 }
119 }
120
121 pub(crate) fn markers(&self) -> CountMarkers {
123 CountMarkers {
124 len_count_inflight: self.len_count_inflight,
125 count_after_paint: self.count_after_paint,
126 len_count_failed: self.len_count_failed,
127 }
128 }
129
130 pub(crate) fn restore(&mut self, markers: CountMarkers) {
132 self.len_count_inflight = markers.len_count_inflight;
133 self.count_after_paint = markers.count_after_paint;
134 self.len_count_failed = markers.len_count_failed;
135 }
136}
137
138#[derive(Clone, Copy)]
140pub(crate) struct CountMarkers {
141 pub(crate) len_count_inflight: Option<u64>,
142 pub(crate) count_after_paint: Option<u64>,
143 pub(crate) len_count_failed: Option<u64>,
144}
145
146const INDEX_STEP: usize = 16 << 20;
149
150impl App {
151 pub fn row_count_pending(&self) -> bool {
154 self.counting.len_count_inflight.is_some()
158 || self.loading.awaiting_dataset()
159 || self.counting.reread_owed.is_some()
162 || self
163 .data_table_state
164 .as_ref()
165 .is_some_and(|state| state.counts_itself_later())
166 || self.lines_answer_owed()
169 }
170
171 fn lines_answer_owed(&self) -> bool {
174 let dataset = self.dataset_generation;
175 self.data_table_state
176 .as_ref()
177 .is_some_and(|state| state.lines_to_index().is_some())
178 && self
179 .jobs
180 .current(
181 |job| matches!(job, Job::IndexLines { dataset: asked } if *asked == dataset),
182 )
183 .is_some()
184 }
185
186 pub(crate) fn dataset_is_still_reading_its_footers(&self) -> bool {
187 self.data_table_state
188 .as_ref()
189 .is_some_and(|state| state.footers_pending().is_some())
190 }
191
192 pub(crate) fn first_rows_settled(&mut self) {
194 self.loading.first_rows_settled();
195 }
196
197 pub(crate) fn reread_after_the_footers_joined(&mut self) {
201 self.counting.reread_owed = None;
203 if self.counting.end_when_the_footers_land.take() == Some(self.dataset_generation) {
206 self.status_message = None;
207 if let Some(next) = self.jump_key(crate::Scroll::End) {
208 let _ = self.events.send(next);
210 return;
211 }
212 }
215 self.spawn_async_collect(Self::LOADING_BUFFER);
216 }
217
218 pub(crate) fn collect_when_the_work_allows(&mut self) {
221 let Some(&Job::OwedRows { dataset, .. }) = self.jobs.owed(Self::owed_rows) else {
222 return;
223 };
224 if dataset != self.dataset_generation {
225 self.jobs.take_owed(Self::owed_rows);
228 return;
229 }
230 if self.work_a_bump_would_strand() {
231 return;
232 }
233 let Some(Job::OwedRows { status, .. }) = self.jobs.take_owed(Self::owed_rows) else {
234 return;
235 };
236 if !self.spawn_async_collect(&status) {
237 self.busy = false;
238 self.status_message = None;
239 self.first_rows_settled();
242 }
243 }
244
245 pub(crate) fn reread_when_the_work_allows(&mut self) {
248 let Some(generation) = self.counting.reread_owed else {
249 return;
250 };
251 if generation != self.dataset_generation {
252 self.counting.reread_owed = None;
254 return;
255 }
256 if self.work_the_join_would_cancel() {
257 return;
258 }
259 self.reread_after_the_footers_joined();
260 }
261
262 fn retire_the_end_that_was_waiting(&mut self) {
265 self.counting.end_after_count = None;
266 self.take_down_the_counting_status();
267 }
268
269 fn take_down_the_counting_status(&mut self) {
272 if matches!(
273 self.status_message.as_deref(),
274 Some(Self::COUNTING_FOR_END | Self::INDEXING_FOR_ROW)
275 ) {
276 self.status_message = None;
277 }
278 }
279
280 pub(crate) fn fetch_too_young_to_mention(&self) -> bool {
282 self.status_message.as_deref() == Some(Self::LOADING_BUFFER)
283 && self
284 .rows_in_flight()
285 .is_some_and(|inflight| inflight.began.elapsed() < Self::A_FETCH_WORTH_SAYING)
286 }
287
288 pub(crate) fn work_the_join_would_cancel(&self) -> bool {
292 self.work_a_bump_would_strand() || self.chart_preparing()
293 }
294
295 pub(crate) fn join_held_footers(&mut self) -> bool {
300 let Some((generation, _)) = self.counting.footers_held.as_ref() else {
301 return false;
302 };
303 if *generation != self.dataset_generation {
304 self.counting.footers_held = None;
306 return false;
307 }
308 if self.data_table_state.is_none() || self.work_the_join_would_cancel() {
309 return false;
310 }
311 let Some((generation, found)) = self.counting.footers_held.take() else {
312 return false;
313 };
314 let state = self
315 .data_table_state
316 .as_mut()
317 .expect("checked just above, and nothing since takes it");
318 match state.join_dataset_schema(found) {
320 Ok(()) => true,
321 Err(found) => {
322 self.counting.footers_held = Some((generation, *found));
323 false
324 }
325 }
326 }
327
328 pub(crate) fn start_pending_footers(&mut self) {
332 let Some(join) = self
333 .data_table_state
334 .as_ref()
335 .and_then(|state| state.footers_pending())
336 else {
337 return;
338 };
339 let dataset = self.dataset_generation;
340 let progress = self.counting.footer_progress.clone();
341 self.spawn_job(Job::FootersJoin { dataset }, None, move |_| {
344 Ok(Answer::FootersJoined(join(&progress).map(Box::new)))
345 });
346 }
347
348 pub(crate) fn index_lines(&mut self) {
351 use std::sync::atomic::Ordering;
352 self.counting.indexing_stop.store(true, Ordering::Relaxed);
353 self.counting.indexing_paused = false;
354 let lines = self
355 .data_table_state
356 .as_ref()
357 .and_then(|state| state.lines_to_index().cloned());
358 if let Some(old) = self.counting.indexing_lines.take()
359 && lines.as_ref().is_none_or(|lines| !Arc::ptr_eq(lines, &old))
360 {
361 old.stop_indexing();
362 }
363 let Some(lines) = lines.filter(|lines| lines.resume_indexing()) else {
364 return;
365 };
366 self.counting.indexing_lines = Some(lines.clone());
367 let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
368 self.counting.indexing_stop = stop.clone();
369 let dataset = self.dataset_generation;
370 self.spawn_job(Job::IndexLines { dataset }, None, move |_| {
372 loop {
373 if stop.load(Ordering::Relaxed) {
375 return Err("stopped".to_string());
376 }
377 let done = logging::catch_panic(|| lines.index_more(INDEX_STEP)).unwrap_or(true);
379 if done {
380 lines.stop_indexing();
381 return Ok(Answer::LinesIndexed(lines.rows()));
382 }
383 }
384 });
385 }
386
387 pub(crate) fn lines_indexed(&mut self, generation: u64, rows: usize) {
390 if generation != self.dataset_generation {
391 return;
392 }
393 let Some(state) = self.data_table_state.as_mut() else {
394 return;
395 };
396 self.counting.indexing_lines = None;
397 if !state.lines_indexed(rows) {
398 if let Some(held) = self.quality.evidence_return.as_mut() {
401 held.lines_indexed(rows);
402 }
403 return;
404 }
405 if let Some((goto, row)) = self.counting.goto_when_indexed.take()
406 && goto == generation
407 {
408 self.take_down_the_counting_status();
409 let _ = self
410 .events
411 .send(AppEvent::Applied(crate::Applied::GoToLine(row)));
412 }
413 if self.counting.end_when_indexed.take() == Some(generation) {
414 self.take_down_the_counting_status();
415 if let Some(next) = self.jump_key(crate::Scroll::End) {
416 let _ = self.events.send(next);
417 return;
418 }
419 }
420 if self.in_normal_table_view() && !self.loading.awaiting_dataset() {
422 self.spawn_collect(None);
423 }
424 }
425
426 pub(crate) fn count_held_at_estimate(&self, state: &DataTableState) -> bool {
429 let limit = self.app_config.read.exact_count_files;
430 limit > 0
431 && state.files_to_count().is_some_and(|files| files > limit)
432 && self.counting.exact_count_asked != Some(self.dataset_generation)
433 && state.row_estimate(None).is_some()
434 }
435
436 pub(crate) fn row_estimate(&self) -> Option<crate::formats::schema_union::RowEstimate> {
439 self.data_table_state
440 .as_ref()?
441 .row_estimate(self.counting.footer_progress.estimate())
442 }
443
444 pub(crate) fn footers_counted(&self) -> Option<(usize, usize)> {
447 self.counting.len_count_inflight?;
448 self.counting
449 .count_progress
450 .reading()
451 .filter(|_| !self.counting.count_progress.is_cancelled())
452 }
453
454 pub(crate) fn count_exactly(&mut self) {
456 let Some(state) = self.data_table_state.as_ref() else {
457 return;
458 };
459 if state.is_num_rows_valid() {
460 return;
461 }
462 let generation = state.len_generation();
463 self.counting.exact_count_asked = Some(self.dataset_generation);
464 if state.counts_itself_later() {
466 return;
467 }
468 if self.counting.len_count_failed == Some(generation) {
470 self.counting.len_count_failed = None;
471 }
472 if self.counting.len_count_inflight == Some(generation)
474 && self.counting.count_progress.is_cancelled()
475 {
476 self.counting.count_after_stop = Some(generation);
477 return;
478 }
479 if self.counting.len_count_inflight != Some(generation) {
480 self.counting.len_count_inflight = Some(generation);
481 let job = LenCount::for_state(state);
482 self.spawn_count(job);
483 }
484 }
485
486 pub(crate) fn stop_count(&mut self) {
488 self.counting.count_progress.cancel();
489 }
490
491 pub(crate) fn footers_joined(
494 &mut self,
495 dataset: u64,
496 found: Option<crate::table::FootersFound>,
497 ) -> Option<AppEvent> {
498 if dataset == self.dataset_generation {
499 let Some(found) = found else {
500 if let Some(state) = self.data_table_state.as_mut() {
503 state.give_up_on_pending_footers();
504 }
505 self.counting.reread_owed = Some(dataset);
508 self.reread_when_the_work_allows();
509 return None;
510 };
511 self.counting.footers_held = Some((dataset, found));
512 if self.join_held_footers() {
513 self.reread_after_the_footers_joined();
514 }
515 }
516 None
517 }
518
519 pub(crate) fn spawn_count(&mut self, job: LenCount) {
522 #[cfg(test)]
523 self.counting
524 .counts_spawned
525 .set(self.counting.counts_spawned.get() + 1);
526 self.counting.count_progress = job.progress.clone();
527 let count = OwedCount::new(job, self.events.clone());
528 self.runtime
529 .spawn_blocking(move || count.answer(LenCount::run));
530 }
531
532 fn waited_on_rows_pending(&self, generation: u64) -> bool {
536 self.loading.awaiting_dataset()
537 || self.jobs.owed(Self::owed_rows).is_some()
538 || (self.rows_waited_on()
539 && self
540 .rows_in_flight()
541 .is_some_and(|inflight| inflight.dataset == generation))
542 }
543
544 pub fn count_waits_for_a_frame(&self) -> bool {
547 self.counting
548 .count_after_paint
549 .is_some_and(|generation| !self.waited_on_rows_pending(generation))
550 }
551
552 pub fn frame_painted(&mut self) {
554 self.pointer.painted();
555 self.count_what_was_painted();
556 if self.refresh_stale_live_matches() {
559 let _ = self.events.send(AppEvent::Wake);
560 }
561 if let Some(state) = &mut self.data_table_state
562 && state.needs_recollect
563 {
564 state.needs_recollect = false;
565 self.spawn_async_collect(App::LOADING_BUFFER);
566 }
567 }
568
569 fn count_what_was_painted(&mut self) {
572 let Some(generation) = self.counting.count_after_paint else {
573 return;
574 };
575 if self.waited_on_rows_pending(generation) {
576 return;
577 }
578 self.counting.count_after_paint = None;
579 let wanted = self
580 .data_table_state
581 .as_ref()
582 .filter(|state| state.len_generation() == generation && !state.is_num_rows_valid());
583 match wanted {
584 Some(state) => {
585 self.counting.len_count_inflight = Some(generation);
586 self.spawn_count(LenCount::for_state(state));
587 }
588 None => {
589 if self.counting.len_count_inflight == Some(generation) {
590 self.counting.len_count_inflight = None;
591 }
592 }
593 }
594 }
595
596 pub(crate) fn retire_a_count_the_rows_answered(&mut self) {
599 let Some(generation) = self.counting.count_after_paint else {
600 return;
601 };
602 let answered = self
603 .data_table_state
604 .as_ref()
605 .is_none_or(|state| state.len_generation() != generation || state.is_num_rows_valid());
606 if answered {
607 self.counting.count_after_paint = None;
608 if self.counting.len_count_inflight == Some(generation) {
609 self.counting.len_count_inflight = None;
610 }
611 }
612 }
613
614 pub(crate) fn counting_event(&mut self, event: AppEvent) -> Option<AppEvent> {
616 match event {
617 AppEvent::BackgroundLenReady {
618 len_generation,
619 num_rows,
620 file_row_groups,
621 } => {
622 if self.counting.len_count_inflight == Some(len_generation) {
623 self.counting.len_count_inflight = None;
624 }
625 if self.counting.len_count_failed == Some(len_generation) {
626 self.counting.len_count_failed = None;
627 }
628 if let Some(run) = self.prompt.query_running.as_mut() {
630 run.rollback
631 .count_landed(len_generation, num_rows, file_row_groups.as_deref());
632 if run.counts.len_count_inflight == Some(len_generation) {
633 run.counts.len_count_inflight = None;
634 }
635 }
636 if let Some(state) = self.data_table_state.as_mut()
639 && state.count_landed(len_generation, num_rows, file_row_groups.as_deref())
640 {
641 if self.counting.end_after_count == Some(len_generation) {
643 self.counting.end_after_count = None;
644 self.status_message = None;
645 return self.jump_key(crate::Scroll::End);
646 }
647 } else if self.counting.end_after_count == Some(len_generation) {
648 self.retire_the_end_that_was_waiting();
653 }
654 self.remember_a_downloads_shape();
655 None
656 }
657 AppEvent::FramePainted => {
658 self.frame_painted();
659 None
660 }
661 AppEvent::BackgroundLenFailed { len_generation } => {
662 if self.counting.len_count_inflight == Some(len_generation) {
663 self.counting.len_count_inflight = None;
664 }
665 if self.counting.count_after_stop.take() == Some(len_generation) {
666 self.count_exactly();
667 return None;
668 }
669 if let Some(run) = self.prompt.query_running.as_mut()
670 && run.counts.len_count_inflight == Some(len_generation)
671 {
672 run.counts.len_count_inflight = None;
673 run.counts.len_count_failed = Some(len_generation);
674 }
675 if self
680 .data_table_state
681 .as_ref()
682 .is_some_and(|state| state.len_generation() == len_generation)
683 {
684 self.counting.len_count_failed = Some(len_generation);
685 }
686 if self.counting.end_after_count == Some(len_generation) {
688 self.counting.end_after_count = None;
689 if self
690 .data_table_state
691 .as_ref()
692 .is_some_and(|state| state.len_generation() == len_generation)
693 {
694 self.status_message =
695 Some("Could not count the rows to find the end".to_string());
696 } else {
697 self.take_down_the_counting_status();
700 }
701 }
702 None
703 }
704 _ => unreachable!("not an event for counting_event"),
705 }
706 }
707
708 pub(crate) fn count_unfit(&mut self) {
711 let dataset = self.dataset_generation;
712 let Some(state) = self.data_table_state.as_ref() else {
713 return;
714 };
715 let read = state
716 .unfit_to_count()
717 .map(|(source, typed)| (source, typed, None));
718 let view = state
719 .changes_unfit_to_count()
720 .map(|(source, typed, version)| (source, typed, Some(version)));
721 let streaming = self.app_config.performance.streaming;
722 for (source, typed, version) in [read, view].into_iter().flatten() {
723 let running = self
724 .jobs
725 .current(|job| {
726 matches!(job, Job::UnfitCount { dataset: d, version: v }
727 if *d == dataset && *v == version)
728 })
729 .is_some();
730 if running {
731 continue;
732 }
733 self.spawn_job(Job::UnfitCount { dataset, version }, None, move |_| {
734 let counted = crate::analysis::statistics::collect_lazy(
735 crate::formats::column_types::unfit_frame(source, &typed),
736 streaming,
737 )
738 .map_err(|e| crate::error_display::user_message_from_polars(&e))?;
739 Ok(Answer::UnfitCounted(
740 crate::formats::column_types::unfit_counts(&counted, &typed),
741 ))
742 });
743 }
744 }
745
746 pub fn unfit_count_pending(&self) -> bool {
748 self.jobs
749 .current(|job| matches!(job, Job::UnfitCount { .. }))
750 .is_some()
751 }
752}