1use crate::analysis::analysis_modal::AnalysisProgress;
5use crate::analysis::quality_memory::{
6 KeptQualitySample, QUALITY_RELEASED_REMEMBERED, QualityCacheEntry, QualityCopyJob, RetainedCopy,
7};
8use crate::app::feedback::Confirm;
9use crate::app::jobs::{Answer, Job, Progress};
10use crate::export::export_modal::ExportFormat;
11use crate::table::DataTableState;
12use crate::{
13 App, AppEvent, QUALITY_RUN_WAITS, analysis::analysis_modal, analysis::data_quality,
14 analysis::quality_report, analysis::sampling, app::jobs, glyphs, numfmt, widgets,
15};
16use color_eyre::Result;
17use polars::prelude::LazyFrame;
18use std::path::{Path, PathBuf};
19use std::sync::Arc;
20
21#[derive(Default)]
24pub struct QualityRuns {
25 pub(crate) cache: Vec<QualityCacheEntry>,
27 pub(crate) samples: Vec<KeptQualitySample>,
29 pub(crate) released: Vec<(u64, u64, sampling::Sample)>,
31 pub(crate) memory_budget: usize,
33 pub(crate) copies: Vec<RetainedCopy>,
37 pub(crate) copy_released: Option<u64>,
39 pub(crate) copy_unusable: Option<u64>,
42 pub(crate) copy_free: std::sync::Mutex<Option<(std::time::Instant, Option<u64>)>>,
45 pub(crate) evidence_return: Option<Box<DataTableState>>,
47 pub(crate) evidence_label: Option<String>,
48}
49
50impl QualityRuns {
51 pub(crate) fn reset_for_dataset(&mut self) {
54 self.cache.clear();
55 self.samples.clear();
56 self.released.clear();
57 self.copies.clear();
58 self.copy_released = None;
59 self.copy_unusable = None;
60 self.evidence_return = None;
61 self.evidence_label = None;
62 }
63}
64
65impl App {
66 pub(crate) fn open_quality_evidence(&mut self) -> Option<AppEvent> {
69 let (_, finding) = self.analysis_modal.selected_finding()?;
70 let results = self.analysis_modal.quality.results.as_ref()?;
71 let rows = finding.evidence(results).ok()?;
72 let sampled = results.precision == data_quality::QualityPrecision::Sampled;
73 let count = finding.evidence_count(results);
74 let label = format!(
75 "Data Quality / {} / {}",
76 finding.title,
77 quality_report::columns_label(&finding.columns, 40)
78 );
79 let what = format!(
80 "{} {} {}",
81 finding.title,
82 glyphs::get().middot,
83 quality_report::columns_label(&finding.columns, 40)
84 );
85 self.show_quality_rows(rows, label, sampled, what, count)
86 }
87
88 pub(crate) fn open_interval_evidence(&mut self) -> Option<AppEvent> {
91 let schema = self.data_table_state.as_ref().map(|state| state.schema());
92 let (predicate, label, count) = self
93 .analysis_modal
94 .interval_evidence(schema.map(|schema| schema.as_ref()))?;
95 let sampled = self
96 .analysis_modal
97 .quality
98 .results
99 .as_ref()
100 .is_some_and(|results| results.precision == data_quality::QualityPrecision::Sampled);
101 let what = label
102 .trim_start_matches("Data Quality / ")
103 .replace(" / ", &format!(" {} ", glyphs::get().middot));
104 self.show_quality_rows(
105 quality_report::EvidenceRows::Matching(predicate),
106 label,
107 sampled,
108 what,
109 Some(count),
110 )
111 }
112
113 pub(crate) fn quality_rows_kept(&self) -> Option<std::sync::Arc<data_quality::QualitySample>> {
116 self.analysis_modal.quality.results.as_ref()?;
117 let plan = self.analysis_modal.quality_result_plan();
118 if plan.compute != data_quality::QualityCompute::Sample {
119 return None;
120 }
121 self.kept_quality_sample(&plan.sample())
122 }
123
124 fn show_quality_rows(
127 &mut self,
128 rows: quality_report::EvidenceRows,
129 label: String,
130 sampled: bool,
131 what: String,
132 count: Option<usize>,
133 ) -> Option<AppEvent> {
134 let plan = self.analysis_modal.quality_result_plan().clone();
135 let by_files = matches!(rows, quality_report::EvidenceRows::Files(_));
136 let label = if sampled {
137 format!("{label} / sampled")
138 } else {
139 label
140 };
141 if !by_files && self.quality_rows_kept().is_some() {
142 return self.read_sample_rows(plan.sample(), Some((rows, label)));
143 }
144 let state = self.data_table_state.as_ref()?;
145 let g = glyphs::get();
146 let rows_label = |rows: usize| {
147 format!(
148 "{} {}",
149 numfmt::group_chrome(rows),
150 if rows == 1 { "row" } else { "rows" }
151 )
152 };
153 let scope = match &rows {
154 quality_report::EvidenceRows::Files(files) => files.clone(),
155 _ => plan.scope.clone(),
156 };
157 let not_kept = if plan.compute == data_quality::QualityCompute::Full {
160 "a full scan keeps no rows"
161 } else {
162 "the rows read are no longer kept"
163 };
164 let (why, reads) = match &rows {
165 quality_report::EvidenceRows::Files(data_quality::QualityScope::SourceFiles(files)) => {
166 (
167 "their rows are in the files, not the report".to_string(),
168 format!(
169 "the {} named {}",
170 files.len(),
171 if files.len() == 1 { "file" } else { "files" }
172 ),
173 )
174 }
175 _ if sampled => (
176 "the sampled rows are no longer kept".to_string(),
177 format!(
178 "the sample again: {} {} {}",
179 widgets::data_quality::compute_label(&plan),
180 g.middot,
181 widgets::data_quality::planned_read_label(state, &plan)
182 ),
183 ),
184 quality_report::EvidenceRows::Duplicates => (
185 not_kept.to_string(),
186 format!(
187 "every row of {}, once {} {}",
188 plan.scope.label(),
189 g.middot,
190 widgets::data_quality::scope_read_label(state, &plan)
191 ),
192 ),
193 _ => (
195 not_kept.to_string(),
196 format!(
197 "every row of {} to count them, then the rows on screen {} {}",
198 plan.scope.label(),
199 g.middot,
200 widgets::data_quality::scope_read_label(state, &plan)
201 ),
202 ),
203 };
204 let shows = match (&rows, count) {
205 (quality_report::EvidenceRows::Duplicates, Some(count)) => {
206 format!("{}, copies together", rows_label(count))
207 }
208 (_, Some(count)) => rows_label(count),
209 (_, None) => "the rows that match".to_string(),
210 };
211 let source = if state.is_remote_source() {
212 "remote, read only"
213 } else {
214 "local, read only"
215 };
216 self.analysis_modal.quality.evidence_read = Some(analysis_modal::EvidenceRead {
217 summary: vec![
218 ("Rows", what),
219 ("Why", why),
220 ("Reads", reads),
221 ("Shows", shows),
222 ("Source", source.to_string()),
223 ],
224 sample: (sampled && !by_files).then(|| plan.sample()),
225 scope,
226 rows,
227 label,
228 });
229 None
230 }
231
232 pub(crate) fn confirm_evidence_read(&mut self) -> Option<AppEvent> {
234 let staged = self.analysis_modal.quality.evidence_read.as_ref()?;
236 let kept = staged
237 .sample
238 .as_ref()
239 .is_some_and(|sample| self.kept_quality_sample(sample).is_some());
240 if !kept && self.read_waits_for_cancelled() {
241 return None;
242 }
243 let read = self.analysis_modal.quality.evidence_read.take()?;
244 if let Some(sample) = read.sample {
245 return self.read_sample_rows(sample, Some((read.rows, read.label)));
246 }
247 let predicate = match read.rows {
248 quality_report::EvidenceRows::Matching(predicate) => predicate,
249 quality_report::EvidenceRows::Files(_) => polars::prelude::lit(true),
252 quality_report::EvidenceRows::Duplicates => {
253 return self.read_duplicate_rows(read.scope, read.label);
254 }
255 };
256 self.open_quality_scope_rows(&read.scope, predicate, read.label)
257 }
258
259 fn read_duplicate_rows(
262 &mut self,
263 scope: data_quality::QualityScope,
264 label: String,
265 ) -> Option<AppEvent> {
266 let state = self.data_table_state.as_ref()?;
267 let (lf, schema) = match state.quality_scope_frame(&scope) {
268 Ok(frame) => frame,
269 Err(error) => {
270 self.error_modal
271 .show(format!("Cannot open matching rows: {error}"));
272 return None;
273 }
274 };
275 let keys = schema.iter_names().cloned().collect::<Vec<_>>();
276 let columns = schema
279 .iter()
280 .map(|(name, dtype)| {
281 if matches!(dtype, polars::prelude::DataType::Binary) {
282 polars::prelude::lit(crate::table::binary_stub()).alias(name.clone())
283 } else {
284 polars::prelude::col(name.clone())
285 }
286 })
287 .collect::<Vec<_>>();
288 let streaming = self.app_config.performance.streaming;
289 self.analysis_modal.computing = Some(AnalysisProgress::new("Reading the rows that repeat"));
290 self.spawn_job(
291 Job::SampleRows,
292 Some("Reading the rows that repeat..."),
293 move |_| {
294 let df = data_quality::duplicate_rows(lf.select(columns), &keys, streaming)
295 .map_err(|error| format!("{error}"))?;
296 Ok(Answer::Sample { df, label })
297 },
298 );
299 None
300 }
301
302 fn open_quality_scope_rows(
304 &mut self,
305 scope: &data_quality::QualityScope,
306 predicate: polars::prelude::Expr,
307 label: String,
308 ) -> Option<AppEvent> {
309 let state = self.data_table_state.as_ref()?;
310 let view = match state.quality_evidence_view(scope, predicate) {
311 Ok(view) => view,
312 Err(error) => {
313 self.error_modal
314 .show(format!("Cannot open matching rows: {error}"));
315 return None;
316 }
317 };
318 if let Some(original) = self.data_table_state.replace(view) {
319 self.quality.evidence_return = Some(Box::new(original));
320 self.quality.evidence_label = Some(label);
321 self.step_back();
322 self.forget_the_rows_read();
323 self.spawn_async_collect("Loading matching rows...");
324 }
325 None
326 }
327
328 pub(crate) fn return_from_quality_evidence(&mut self, reopen_analysis: bool) -> bool {
329 let Some(original) = self.quality.evidence_return.take() else {
330 return false;
331 };
332 self.jobs.advance();
333 self.counting.len_count_inflight = None;
334 self.data_table_state = Some(*original);
335 self.quality.evidence_label = None;
336 if reopen_analysis {
337 self.open_overlay(crate::Overlay::Analysis);
338 }
339 self.busy = false;
340 self.status_message = None;
341 true
342 }
343
344 pub(crate) fn restore_recent_quality_plan(&mut self) {
345 let Some(view_generation) = self
346 .data_table_state
347 .as_ref()
348 .map(DataTableState::len_generation)
349 else {
350 return;
351 };
352 if self.analysis_modal.quality.plan != data_quality::DataQualityPlan::default() {
353 return;
354 }
355 if let Some(cached) = self.quality.cache.iter().find(|entry| {
356 entry.dataset_generation == self.dataset_generation
357 && entry.view_generation == view_generation
358 }) {
359 self.analysis_modal.quality.plan = cached.plan.clone();
360 }
361 }
362
363 pub(crate) fn quality_plan_context(&self) -> analysis_modal::PlanContext {
366 let Some(state) = self.data_table_state.as_ref() else {
367 return analysis_modal::PlanContext::default();
368 };
369 let plan = &self.analysis_modal.quality.plan;
370 let scope = &plan.scope;
371 let schema = state.schema();
372 let mut partitions = state.partition_columns().unwrap_or_default().to_vec();
373 if partitions.is_empty()
376 && let Some(dir) = self.path.as_ref().filter(|path| path.is_dir())
377 {
378 partitions = crate::formats::readers::hive::discover_hive_partition_columns(dir)
379 .into_iter()
380 .filter(|column| schema.get(column).is_some())
381 .collect();
382 }
383 let mut time_columns: Vec<(String, bool)> = state
385 .quality_temporal_columns(scope)
386 .into_iter()
387 .map(|column| {
388 let has_time =
389 !matches!(schema.get(&column), Some(polars::prelude::DataType::Date));
390 (column, has_time)
391 })
392 .collect();
393 for format in &plan.time_formats {
394 if !time_columns
395 .iter()
396 .any(|(column, _)| *column == format.column)
397 {
398 time_columns.push((
399 format.column.clone(),
400 format.kind == data_quality::TimeKind::Datetime,
401 ));
402 }
403 }
404 analysis_modal::PlanContext {
405 partitions,
406 time_columns,
407 files: state.quality_source_file_count() > 1,
408 text_columns: state
409 .quality_text_columns(scope)
410 .into_iter()
411 .map(|column| {
412 let examples = state.buffered_values(&column, 3);
413 (column, examples)
414 })
415 .collect(),
416 }
417 }
418
419 pub(crate) fn quality_time_candidates(&self) -> Vec<String> {
422 let Some(state) = self.data_table_state.as_ref() else {
423 return Vec::new();
424 };
425 let scope = &self.analysis_modal.quality.plan.scope;
426 let mut columns = state.quality_temporal_columns(scope);
427 columns.extend(state.quality_text_columns(scope));
428 columns
429 }
430
431 pub(crate) fn quality_intent_columns(&self) -> Vec<(String, polars::prelude::DataType)> {
433 self.data_table_state
434 .as_ref()
435 .map(|state| {
436 crate::widgets::quality_intent::intent_columns(
437 state.quality_schema(&self.analysis_modal.quality.plan.scope),
438 )
439 })
440 .unwrap_or_default()
441 }
442
443 fn quality_source_identity(
446 &self,
447 state: &DataTableState,
448 scope: &data_quality::QualityScope,
449 ) -> crate::analysis::quality_export::SourceIdentity {
450 let format = self
451 .source
452 .original_file_format
453 .map(|format| format.as_str().to_string())
454 .or_else(|| {
455 self.path
456 .as_ref()
457 .and_then(|path| path.extension())
458 .and_then(|extension| extension.to_str())
459 .map(str::to_string)
460 });
461 let mut view = Vec::new();
462 if !scope.uses_source() {
463 if !state.get_active_query().is_empty() {
464 view.push(format!("query: {}", state.get_active_query()));
465 }
466 if !state.get_active_sql_query().is_empty() {
467 view.push(format!("SQL: {}", state.get_active_sql_query()));
468 }
469 if !state.get_active_fuzzy_query().is_empty() {
470 view.push(format!("text: {}", state.get_active_fuzzy_query()));
471 }
472 for (index, filter) in state.view_filters().iter().enumerate() {
473 let join = if index == 0 {
474 String::new()
475 } else {
476 format!("{} ", filter.logical_op.as_str())
477 };
478 view.push(format!(
479 "filter: {join}{} {} {}",
480 filter.column,
481 filter.operator.as_str(),
482 filter.value
483 ));
484 }
485 if state.reshape_source().is_some() {
486 view.push("reshaped: pivot or melt".to_string());
487 }
488 }
489 let remote = state.is_remote_source();
490 let piped = self.reads_stdin();
493 let location = self.path.as_ref().map(|path| {
494 match std::path::absolute(path).ok().filter(|_| !remote && !piped) {
495 Some(path) => path.display().to_string(),
496 None => path.display().to_string(),
497 }
498 });
499 crate::analysis::quality_export::SourceIdentity {
500 location,
501 remote,
502 format,
503 view,
504 ..crate::analysis::quality_export::SourceIdentity::default()
505 }
506 .with_files(state.quality_source_file_names())
507 }
508
509 pub(crate) fn quality_page_setup(&self) -> Option<data_quality::QualitySetup> {
512 let modal = &self.analysis_modal;
513 data_quality::page_setup(
514 modal.quality.page,
515 modal.quality_result_plan(),
516 modal.quality.results.as_ref(),
517 self.has_quality_time_columns(),
518 )
519 }
520
521 pub(crate) fn has_quality_time_columns(&self) -> bool {
524 !self.analysis_modal.quality.plan.time_formats.is_empty()
525 || self.data_table_state.as_ref().is_some_and(|state| {
526 !state
527 .quality_temporal_columns(&self.analysis_modal.quality.plan.scope)
528 .is_empty()
529 })
530 }
531
532 pub(crate) fn quality_kept_serves(&self, plan: &data_quality::DataQualityPlan) -> bool {
535 plan.compute == data_quality::QualityCompute::Sample
536 && self.kept_quality_sample(&plan.sample()).is_some()
537 }
538
539 pub(crate) fn quality_segment_count(
542 &self,
543 plan: &data_quality::DataQualityPlan,
544 ) -> data_quality::SegmentCount {
545 if plan.compute != data_quality::QualityCompute::Sample {
546 return data_quality::SegmentCount::NotNeeded;
547 }
548 match self.kept_quality_sample(&plan.sample()) {
549 Some(kept) => kept.segment_count(plan),
550 None => data_quality::fresh_segment_count(plan, self.quality_may_read_blocks(plan)),
551 }
552 }
553
554 pub(crate) fn quality_released(&self, plan: &data_quality::DataQualityPlan) -> bool {
557 let Some(view_generation) = self
558 .data_table_state
559 .as_ref()
560 .map(DataTableState::len_generation)
561 else {
562 return false;
563 };
564 let sample = plan.sample();
565 plan.compute == data_quality::QualityCompute::Sample
566 && self
567 .quality
568 .released
569 .iter()
570 .any(|(dataset, view, released)| {
571 *dataset == self.dataset_generation
572 && *view == view_generation
573 && *released == sample
574 })
575 }
576
577 fn quality_one_columnar_file(&self) -> bool {
580 let Some(state) = self.data_table_state.as_ref() else {
581 return false;
582 };
583 let columnar = matches!(
584 self.source.original_file_format,
585 Some(ExportFormat::Parquet | ExportFormat::Ipc)
586 ) || self.path.as_ref().is_some_and(|path| {
587 path.extension()
588 .and_then(|extension| extension.to_str())
589 .is_some_and(|extension| {
590 matches!(
591 extension.to_ascii_lowercase().as_str(),
592 "parquet" | "pq" | "arrow" | "arrows" | "ipc" | "feather"
593 )
594 })
595 });
596 columnar && state.loaded_file_count() == 1
597 }
598
599 pub(crate) fn quality_reads_blocks(&self, plan: &data_quality::DataQualityPlan) -> bool {
605 let Some(state) = self.data_table_state.as_ref() else {
606 return false;
607 };
608 self.quality_one_columnar_file()
609 && match plan.scope {
610 data_quality::QualityScope::WholeSource => true,
611 data_quality::QualityScope::CurrentView => !state.changes_rows(),
612 data_quality::QualityScope::FirstRows(_)
614 | data_quality::QualityScope::ViewRows { .. } => {
615 state.source_file_count() == Some(1)
616 }
617 _ => false,
618 }
619 }
620
621 pub(crate) fn quality_may_read_blocks(&self, plan: &data_quality::DataQualityPlan) -> bool {
626 let Some(state) = self.data_table_state.as_ref() else {
627 return false;
628 };
629 self.quality_reads_blocks(plan)
630 || (self.quality_one_columnar_file()
631 && state.may_keep_scan_rows()
632 && matches!(
633 plan.scope,
634 data_quality::QualityScope::CurrentView
635 | data_quality::QualityScope::FirstRows(_)
636 | data_quality::QualityScope::ViewRows { .. }
637 ))
638 }
639
640 pub(crate) fn quality_cached(&self, plan: &data_quality::DataQualityPlan) -> bool {
643 let Some(view_generation) = self
644 .data_table_state
645 .as_ref()
646 .map(DataTableState::len_generation)
647 else {
648 return false;
649 };
650 self.quality.cache.iter().any(|entry| {
651 entry.dataset_generation == self.dataset_generation
652 && entry.view_generation == view_generation
653 && entry.plan.same_measurement(plan)
654 })
655 }
656
657 pub(crate) fn quality_kept_rows(&self) -> Option<widgets::data_quality::KeptRows> {
660 let kept = self
661 .quality
662 .samples
663 .iter()
664 .filter(|kept| kept.dataset_generation == self.dataset_generation)
665 .collect::<Vec<_>>();
666 let copy_bytes = self
667 .quality
668 .copies
669 .iter()
670 .filter(|kept| kept.dataset_generation == self.dataset_generation)
671 .map(|kept| kept.copy.bytes())
672 .sum::<u64>();
673 (!kept.is_empty() || copy_bytes > 0).then(|| widgets::data_quality::KeptRows {
674 samples: kept.len(),
675 rows: kept.iter().map(|kept| kept.rows.df().height()).sum(),
676 bytes: kept.iter().map(|kept| kept.rows.estimated_bytes()).sum(),
677 copy_bytes,
678 })
679 }
680
681 pub(crate) fn release_quality_rows(&mut self) {
685 let Some(kept) = self.quality_kept_rows() else {
686 self.flash_note("Nothing kept to release".to_string());
687 return;
688 };
689 for released in std::mem::take(&mut self.quality.samples) {
690 self.quality.released.retain(|(dataset, view, sample)| {
691 !(*dataset == released.dataset_generation
692 && *view == released.view_generation
693 && *sample == released.sample)
694 });
695 self.quality.released.insert(
696 0,
697 (
698 released.dataset_generation,
699 released.view_generation,
700 released.sample,
701 ),
702 );
703 }
704 self.quality.released.truncate(QUALITY_RELEASED_REMEMBERED);
705 let generation = self.dataset_generation;
707 self.quality
708 .copies
709 .retain(|kept| kept.dataset_generation != generation);
710 if kept.copy_bytes > 0 {
711 self.quality.copy_released = Some(generation);
712 }
713 let rows = format!(
714 "{} kept {} ({})",
715 numfmt::group_chrome(kept.rows),
716 if kept.rows == 1 { "row" } else { "rows" },
717 crate::numfmt::bytes(kept.bytes as u64)
718 );
719 let copy = format!("the local copy ({})", crate::numfmt::bytes(kept.copy_bytes));
720 self.flash_note(match (kept.samples > 0, kept.copy_bytes > 0) {
721 (true, true) => format!("Released {rows} and {copy}; the next run reads again"),
722 (false, true) => format!("Released {copy}; the next full scan fetches again"),
723 _ => format!("Released {rows}; the next run reads again"),
724 });
725 }
726
727 fn kept_quality_sample(
730 &self,
731 sample: &sampling::Sample,
732 ) -> Option<std::sync::Arc<data_quality::QualitySample>> {
733 self.kept_quality_entry(sample)
734 .map(|kept| kept.rows.clone())
735 }
736
737 fn kept_quality_entry(&self, sample: &sampling::Sample) -> Option<&KeptQualitySample> {
738 let view_generation = self.data_table_state.as_ref()?.len_generation();
739 self.quality.samples.iter().find(|kept| {
740 kept.dataset_generation == self.dataset_generation
741 && kept.view_generation == view_generation
742 && &kept.sample == sample
743 })
744 }
745
746 pub(crate) fn retain_quality_sample(&mut self, kept: &KeptQualitySample) {
749 if kept.dataset_generation != self.dataset_generation {
750 return;
751 }
752 self.quality.samples.retain(|entry| !entry.same_rows(kept));
753 self.quality.released.retain(|(dataset, view, sample)| {
754 !(*dataset == kept.dataset_generation
755 && *view == kept.view_generation
756 && *sample == kept.sample)
757 });
758 self.quality.samples.insert(0, kept.clone());
759 self.trim_quality_memory();
760 }
761
762 fn trim_quality_memory(&mut self) {
768 loop {
769 let used = self
770 .quality
771 .cache
772 .iter()
773 .map(|entry| entry.bytes)
774 .sum::<usize>()
775 + self
776 .quality
777 .samples
778 .iter()
779 .map(|kept| kept.rows.estimated_bytes())
780 .sum::<usize>();
781 if used <= self.quality.memory_budget {
782 return;
783 }
784 let remakeable = self
785 .quality
786 .cache
787 .iter()
788 .enumerate()
789 .skip(1)
790 .rev()
791 .find(|(_, entry)| {
792 entry.plan.compute == data_quality::QualityCompute::Sample
793 && self.quality.samples.iter().any(|kept| {
794 kept.dataset_generation == entry.dataset_generation
795 && kept.view_generation == entry.view_generation
796 && kept.sample == entry.plan.sample()
797 })
798 })
799 .map(|(index, _)| index);
800 if let Some(index) = remakeable {
801 self.quality.cache.remove(index);
802 } else if self.quality.samples.len() > 1 {
803 if let Some(released) = self.quality.samples.pop() {
804 self.quality.released.insert(
805 0,
806 (
807 released.dataset_generation,
808 released.view_generation,
809 released.sample,
810 ),
811 );
812 self.quality.released.truncate(QUALITY_RELEASED_REMEMBERED);
813 }
814 } else if self.quality.cache.len() > 1 {
815 self.quality.cache.pop();
816 } else {
817 return;
818 }
819 }
820 }
821
822 pub(crate) fn restore_cached_quality(&mut self) -> bool {
823 let Some(view_generation) = self
824 .data_table_state
825 .as_ref()
826 .map(DataTableState::len_generation)
827 else {
828 return false;
829 };
830 let plan = self.analysis_modal.quality.plan.clone();
831 let Some(cached) = self.quality.cache.iter().find(|entry| {
832 entry.dataset_generation == self.dataset_generation
833 && entry.view_generation == view_generation
834 && entry.plan.same_measurement(&plan)
835 }) else {
836 return false;
837 };
838 let mut results = cached.results.clone();
839 if cached.plan != plan {
840 if cached.plan.compares_differently(&plan) {
842 results.compare_segments(&plan);
843 }
844 self.cache_quality_result(&results, plan.clone());
845 }
846 self.analysis_modal.quality.results = Some(results);
847 self.analysis_modal.quality.last_plan = Some(plan);
848 self.analysis_modal.quality.from_cache = true;
849 self.analysis_modal
850 .set_quality_page(data_quality::QualityPage::Overview);
851 true
852 }
853
854 pub(crate) fn cache_quality_result(
855 &mut self,
856 results: &data_quality::DataQualityResults,
857 plan: data_quality::DataQualityPlan,
858 ) {
859 let Some(view_generation) = self
860 .data_table_state
861 .as_ref()
862 .map(DataTableState::len_generation)
863 else {
864 return;
865 };
866 self.quality.cache.retain(|entry| {
868 !(entry.dataset_generation == self.dataset_generation
869 && entry.view_generation == view_generation
870 && entry.plan.same_measurement(&plan))
871 });
872 self.quality.cache.insert(
873 0,
874 QualityCacheEntry {
875 dataset_generation: self.dataset_generation,
876 view_generation,
877 plan,
878 bytes: results.estimated_bytes(),
879 results: results.clone(),
880 },
881 );
882 self.trim_quality_memory();
883 }
884
885 pub(crate) fn quality_scope_on_copy(
890 lf: LazyFrame,
891 job: QualityCopyJob,
892 watch: &data_quality::QualityWatch,
893 fetch: impl FnOnce(
894 &[crate::cloud::local_copy::RemoteObject],
895 &Path,
896 ) -> Result<crate::cloud::local_copy::LocalCopy>,
897 kept: impl FnOnce(Option<Arc<crate::cloud::local_copy::LocalCopy>>),
898 ) -> Result<(LazyFrame, Option<Arc<crate::cloud::local_copy::LocalCopy>>)> {
899 let (copy, fetched) = match job {
900 QualityCopyJob::Source => return Ok((lf, None)),
901 QualityCopyJob::Kept(copy) => (copy, false),
902 QualityCopyJob::Fetch { objects, root } => {
903 watch.stage(data_quality::QualityStage::CopyingSource, true, true)?;
904 let copy = fetch(&objects, &root).map_err(|error| {
905 if watch.cancelled() {
906 color_eyre::eyre::eyre!(crate::analysis::sampling::CANCELLED)
907 } else {
908 error
909 }
910 })?;
911 (Arc::new(copy), true)
912 }
913 };
914 let local = copy.redirect(&lf).filter(|local| {
916 let schemas = (local.clone().collect_schema(), lf.clone().collect_schema());
917 matches!(schemas, (Ok(local), Ok(source)) if local == source)
918 });
919 let Some(local) = local else {
920 log::warn!(target: "datui", "local copy does not read as the source; reading the source");
921 kept(None);
922 return Ok((lf, None));
923 };
924 if fetched {
925 kept(Some(copy.clone()));
926 }
927 watch.use_copy(data_quality::CopyRead {
928 bytes: copy.bytes(),
929 objects: copy.objects(),
930 fetched,
931 });
932 Ok((local, Some(copy)))
933 }
934
935 #[cfg(feature = "cloud")]
938 fn fetch_quality_copy(
939 objects: &[crate::cloud::local_copy::RemoteObject],
940 root: &Path,
941 cloud: &crate::config::CloudConfig,
942 runtime: &tokio::runtime::Handle,
943 stop: &crate::analysis::sampling::ReadWatch,
944 ) -> Result<crate::cloud::local_copy::LocalCopy> {
945 use crate::cloud::download::StreamError;
946 use object_store::ObjectStoreExt;
947
948 crate::cloud::local_copy::LocalCopy::fetch(root, objects, stop, |object, write| {
949 let url = object.url.as_str();
950 let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)
951 .map_err(StreamError::Write)?;
952 let (_, key) = Self::cloud_bucket_and_key(url).map_err(StreamError::Write)?;
953 let path = crate::cloud::cloud_browse::object_path(&key);
954 let listed = object.etag.clone();
955 let gone = crate::error_display::gone_since_opened_message(url);
956 let open = async move {
957 let got = store.get(&path).await.map_err(|e| match e {
958 object_store::Error::NotFound { .. } => gone.clone(),
959 e => e.to_string(),
960 })?;
961 if let (Some(listed), Some(fetched)) = (&listed, &got.meta.e_tag)
964 && !crate::cloud::local_copy::same_etag(listed, fetched)
965 {
966 return Err(gone);
967 }
968 Ok((got.into_stream(), None))
969 };
970 let watch = stop.clone();
971 crate::cloud::download::stream_into(runtime, open, move || watch.stopped(), write)
972 })
973 }
974
975 fn quality_copies_root(&self) -> PathBuf {
977 self.cache
978 .cache_dir()
979 .join(crate::cloud::local_copy::COPIES_DIR)
980 }
981
982 fn quality_copy_limit(&self) -> u64 {
984 self.app_config.analysis.quality_local_copy.bytes()
985 }
986
987 pub fn quality_copy_bytes(&self) -> u64 {
989 self.quality
990 .copies
991 .iter()
992 .map(|kept| kept.copy.bytes())
993 .sum()
994 }
995
996 fn quality_copy_kept(&self) -> Option<&Arc<crate::cloud::local_copy::LocalCopy>> {
998 let state = self.data_table_state.as_ref()?;
999 self.quality
1000 .copies
1001 .iter()
1002 .find(|kept| {
1003 kept.dataset_generation == self.dataset_generation
1004 && state.each_remote_object().is_some_and(|mut objects| {
1005 objects.all(|object| {
1006 object.is_some_and(|object| kept.copy.covers(&object.url))
1007 })
1008 })
1009 })
1010 .map(|kept| &kept.copy)
1011 }
1012
1013 fn quality_copy_free_space(&self) -> Option<u64> {
1015 let root = self.quality_copies_root();
1016 let Ok(mut cached) = self.quality.copy_free.lock() else {
1017 return crate::cloud::local_copy::free_space(&root);
1018 };
1019 match *cached {
1020 Some((asked, free)) if asked.elapsed() < std::time::Duration::from_secs(5) => free,
1021 _ => {
1022 let free = crate::cloud::local_copy::free_space(&root);
1023 *cached = Some((std::time::Instant::now(), free));
1024 free
1025 }
1026 }
1027 }
1028
1029 pub(crate) fn quality_copy_plan(
1032 &self,
1033 plan: &data_quality::DataQualityPlan,
1034 ) -> data_quality::CopyPlan {
1035 use data_quality::{CopyPlan, NoCopy};
1036 let Some(state) = self.data_table_state.as_ref() else {
1037 return CopyPlan::NotApplicable;
1038 };
1039 if plan.compute != data_quality::QualityCompute::Full || !state.is_remote_source() {
1040 return CopyPlan::NotApplicable;
1041 }
1042 if !state.quality_reads_whole_source(&plan.scope) {
1043 return CopyPlan::Passes(NoCopy::PartOfTheSource);
1044 }
1045 if let Some(copy) = self.quality_copy_kept() {
1046 return CopyPlan::Kept {
1047 bytes: copy.bytes(),
1048 objects: copy.objects(),
1049 };
1050 }
1051 let limit = self.quality_copy_limit();
1052 if limit == 0 {
1053 return CopyPlan::Passes(NoCopy::Off);
1054 }
1055 if self.quality.copy_unusable == Some(self.dataset_generation) {
1056 return CopyPlan::Passes(NoCopy::Unusable);
1057 }
1058 let Some((bytes, objects)) = state.remote_objects_size() else {
1059 return CopyPlan::Passes(NoCopy::SizeUnknown);
1060 };
1061 if bytes > limit {
1062 return CopyPlan::Passes(NoCopy::TooLarge { bytes, limit });
1063 }
1064 let free = self.quality_copy_free_space();
1065 if free.is_none_or(|free| bytes > free) {
1066 return CopyPlan::Passes(NoCopy::NoRoom { bytes, free });
1067 }
1068 CopyPlan::Fetch { bytes, objects }
1069 }
1070
1071 pub(crate) fn quality_copy_released(&self) -> bool {
1073 self.quality.copy_released == Some(self.dataset_generation)
1074 }
1075
1076 pub(crate) fn retain_quality_copy(
1080 &mut self,
1081 dataset_generation: u64,
1082 copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>,
1083 ) {
1084 if dataset_generation != self.dataset_generation {
1085 return;
1086 }
1087 let Some(copy) = copy else {
1088 self.quality
1089 .copies
1090 .retain(|kept| kept.dataset_generation != dataset_generation);
1091 self.quality.copy_unusable = Some(dataset_generation);
1092 return;
1093 };
1094 self.quality.copies.insert(
1095 0,
1096 RetainedCopy {
1097 dataset_generation,
1098 copy,
1099 },
1100 );
1101 self.quality.copy_released = None;
1102 let limit = self.quality_copy_limit();
1103 while self.quality.copies.len() > 1 && self.quality_copy_bytes() > limit {
1104 self.quality.copies.pop();
1105 }
1106 }
1107
1108 pub(crate) fn read_sample_rows(
1111 &mut self,
1112 sample: sampling::Sample,
1113 evidence: Option<(quality_report::EvidenceRows, String)>,
1114 ) -> Option<AppEvent> {
1115 let state = self.data_table_state.as_ref()?;
1116 let (source, known_total) = Self::sample_source_for(state, &sample.scope);
1117 let streaming = self.app_config.performance.streaming;
1118 let kept = self.kept_quality_sample(&sample).map(|kept| {
1121 let columns: Vec<_> = state
1122 .schema()
1123 .iter_names()
1124 .filter(|name| kept.df().column(name.as_str()).is_ok())
1125 .map(|name| polars::prelude::col(name.clone()))
1126 .collect();
1127 (kept, columns)
1128 });
1129 if kept.is_none() && self.read_waits_for_cancelled() {
1130 return None;
1131 }
1132 self.analysis_modal.computing = Some(AnalysisProgress::new(if evidence.is_some() {
1133 "Reading the matching sampled rows"
1134 } else {
1135 "Reading the sample"
1136 }));
1137 self.spawn_job(Job::SampleRows, Some("Reading the sample..."), move |_| {
1138 let (rows, columns) = match kept {
1141 Some((kept, columns)) => (Ok(kept.analysis_rows(kept.df().clone())), Some(columns)),
1142 None => (
1143 source
1144 .cut(&sample.scope)
1145 .and_then(|lf| sampling::read(&lf, &sample, known_total, streaming)),
1146 None,
1147 ),
1148 };
1149 let shown = |df: polars::prelude::DataFrame| match &columns {
1150 Some(columns) => polars::prelude::IntoLazy::lazy(df)
1151 .select(columns.clone())
1152 .collect()
1153 .map_err(color_eyre::eyre::Report::from),
1154 None => Ok(df),
1155 };
1156 let read = rows.and_then(|rows| {
1157 let label = format!(
1158 "Sample {} {}",
1159 crate::glyphs::get().middot,
1160 sample.outcome(
1161 rows.total_rows,
1162 rows.sample_size,
1163 rows.per_value.as_ref().map(|per_value| per_value.kept),
1164 )
1165 );
1166 match evidence {
1167 Some((quality_report::EvidenceRows::Duplicates, label)) => {
1168 let keys = rows
1170 .df
1171 .get_column_names()
1172 .into_iter()
1173 .filter(|name| !name.starts_with("__datui"))
1174 .cloned()
1175 .collect::<Vec<_>>();
1176 let df = data_quality::duplicate_rows(
1177 polars::prelude::IntoLazy::lazy(rows.df),
1178 &keys,
1179 streaming,
1180 )?;
1181 Ok((shown(df)?, label))
1182 }
1183 Some((quality_report::EvidenceRows::Matching(predicate), label)) => {
1184 let df = polars::prelude::IntoLazy::lazy(rows.df)
1185 .filter(predicate)
1186 .collect()?;
1187 Ok((shown(df)?, label))
1188 }
1189 Some((quality_report::EvidenceRows::Files(_), label)) => {
1191 Ok((shown(rows.df)?, label))
1192 }
1193 None => Ok((shown(rows.df)?, label)),
1194 }
1195 });
1196 let (df, label) = read.map_err(|error| format!("{error}"))?;
1197 Ok(Answer::Sample { df, label })
1198 });
1199 None
1200 }
1201
1202 pub(crate) fn sync_quality_plan(&mut self) {
1205 let sample = self.analysis_modal.sample.clone();
1206 self.analysis_modal.quality.plan.adopt_sample(&sample);
1207 }
1208
1209 pub(crate) fn open_quality_setup(&mut self) {
1212 use data_quality::QualityPage;
1213 let modal = &mut self.analysis_modal;
1214 if !modal.quality.page.is_setup() {
1215 modal.quality.setup_return = modal.quality.page.tab();
1216 }
1217 if modal.quality.setup_before.is_none() {
1218 modal.quality.setup_before = Some(modal.quality.plan.clone());
1219 }
1220 if modal.quality.page != QualityPage::Setup {
1221 modal.set_quality_page(QualityPage::Setup);
1222 modal.quality.plan_field = 0;
1223 }
1224 modal.focus = analysis_modal::AnalysisFocus::Main;
1225 }
1226
1227 pub(crate) fn leave_quality_setup(&mut self) {
1230 use data_quality::QualityPage;
1231 let modal = &mut self.analysis_modal;
1232 if let Some(before) = modal.quality.setup_before.take() {
1233 modal.quality.plan = before;
1234 }
1235 modal.quality.setup_note = None;
1236 modal.quality.picker = None;
1237 if modal.quality.results.is_some() {
1238 let back = match modal.quality.setup_return {
1239 page if page.is_setup() => QualityPage::Overview,
1240 page => page,
1241 };
1242 modal.set_quality_page(back);
1243 } else {
1244 modal.set_quality_page(QualityPage::Setup);
1245 modal.focus = analysis_modal::AnalysisFocus::Sidebar;
1246 }
1247 }
1248
1249 fn quality_full_scan_question(&self, plan: &data_quality::DataQualityPlan) -> String {
1252 let mut lines = vec![
1253 "Run a full scan?".to_string(),
1254 String::new(),
1255 "Reads: every eligible row, up to the whole source".to_string(),
1256 ];
1257 if let data_quality::CopyPlan::Fetch { bytes, .. } = self.quality_copy_plan(plan) {
1258 lines.push(format!(
1259 "Fetch: {} once, to a local copy",
1260 crate::numfmt::bytes(bytes)
1261 ));
1262 }
1263 lines.push("Source writes: none".to_string());
1264 lines.join("\n")
1265 }
1266
1267 fn quality_setup_problem(&self) -> Option<String> {
1270 let plan = &self.analysis_modal.quality.plan;
1271 let schema = self.data_table_state.as_ref()?.quality_schema(&plan.scope);
1272 match &plan.grain {
1273 data_quality::QualityGrain::TimeWindows { column, .. }
1274 if plan.compute != data_quality::QualityCompute::Metadata
1275 && !plan.reads_as_time(column, schema) =>
1276 {
1277 Some(format!(
1278 "{column}: text, no format {} set Text as time",
1279 crate::glyphs::get().middot
1280 ))
1281 }
1282 _ => None,
1283 }
1284 }
1285
1286 pub(crate) fn run_quality_setup(&mut self, confirmed: bool) -> Option<AppEvent> {
1291 use data_quality::QualityPage;
1292 if self.cancelled_analysis_running().is_some() {
1293 self.analysis_modal.quality.setup_note = Some(QUALITY_RUN_WAITS.to_string());
1294 return None;
1295 }
1296 if let Some(problem) = self.quality_setup_problem() {
1297 self.analysis_modal.quality.setup_note = Some(problem);
1298 return None;
1299 }
1300 let plan = &self.analysis_modal.quality.plan;
1302 let here = (self.analysis_modal.quality.results.is_some()
1303 && self
1304 .analysis_modal
1305 .quality
1306 .last_plan
1307 .as_ref()
1308 .is_some_and(|last| last.same_measurement(plan)))
1309 || self.quality_cached(plan);
1310 if plan.requires_confirmation() && !here && !confirmed {
1311 let message = self.quality_full_scan_question(plan);
1313 self.confirmation_modal
1314 .show(message, Confirm::QualityFullScan);
1315 self.confirmation_modal.yes_label = "Run";
1316 return None;
1317 }
1318 self.commit_quality_plan();
1319 let modal = &mut self.analysis_modal;
1320 if modal.quality.results.is_some()
1321 && modal.quality.last_plan.as_ref() == Some(&modal.quality.plan)
1322 {
1323 let back = match modal.quality.setup_return {
1324 page if page.is_setup() => QualityPage::Overview,
1325 page => page,
1326 };
1327 modal.set_quality_page(back);
1328 return None;
1329 }
1330 if let (Some(results), Some(last)) = (
1333 modal.quality.results.as_ref(),
1334 modal.quality.last_plan.as_ref(),
1335 ) && last.same_measurement(&modal.quality.plan)
1336 {
1337 let mut results = results.clone();
1338 let plan = modal.quality.plan.clone();
1339 let page = if last.compares_differently(&plan) {
1340 results.compare_segments(&plan);
1341 QualityPage::Segments
1342 } else {
1343 QualityPage::Trends
1344 };
1345 modal.quality.results = Some(results.clone());
1346 modal.quality.last_plan = Some(plan.clone());
1347 modal.set_quality_page(page);
1348 self.cache_quality_result(&results, plan);
1349 return None;
1350 }
1351 if self.restore_cached_quality() {
1352 return None;
1353 }
1354 self.analysis_modal.quality.from_cache = false;
1355 let mut progress = AnalysisProgress::new("Preparing the plan");
1356 if self.quality_kept_serves(&self.analysis_modal.quality.plan) {
1357 progress.reuse = Some("Starts from rows a run already read".to_string());
1358 }
1359 self.analysis_modal.computing = Some(progress);
1360 self.busy = true;
1361 Some(AppEvent::AnalysisCompute(
1362 analysis_modal::AnalysisTool::DataQuality,
1363 ))
1364 }
1365
1366 fn commit_quality_plan(&mut self) {
1370 let modal = &mut self.analysis_modal;
1371 let sample = modal.quality.plan.sample();
1372 if sample != modal.sample {
1373 modal.describe_results = None;
1374 modal.distribution_results = None;
1375 modal.correlation_results = None;
1376 }
1377 modal.sample = sample;
1378 modal.sample_dataset = Some(self.dataset_generation);
1379 modal.sample_run_for = Some(self.dataset_generation);
1380 modal.quality.setup_before = None;
1381 modal.quality.setup_note = None;
1382 modal.quality.picker = None;
1383 }
1384
1385 pub(crate) fn run_quality_compute(&mut self) -> Option<AppEvent> {
1387 if let Some(state) = &self.data_table_state {
1388 let plan = self.analysis_modal.quality.plan.clone();
1389 let source_scope = plan.scope.uses_source();
1390 let (lf, source, cached_rows) = if source_scope {
1391 let (lf, source) = state.data_quality_source_scan();
1392 (lf, source, None)
1393 } else {
1394 let ordered = matches!(
1395 plan.scope,
1396 data_quality::QualityScope::FirstRows(_)
1397 | data_quality::QualityScope::ViewRows { .. }
1398 );
1399 let (lf, source) = state.data_quality_scan(ordered);
1400 let rows = state.num_rows_if_valid().map(|rows| match &plan.scope {
1401 data_quality::QualityScope::CurrentView => rows,
1402 data_quality::QualityScope::FirstRows(limit) => rows.min(*limit),
1403 data_quality::QualityScope::ViewRows { start, end } => {
1404 rows.min(*end).saturating_sub(start.saturating_sub(1))
1405 }
1406 _ => unreachable!(),
1407 });
1408 (lf, source, rows)
1409 };
1410 let streaming = state.polars_streaming();
1411 let audio = (plan.compute == data_quality::QualityCompute::Full)
1413 .then(|| state.window_for_quality(&plan.scope))
1414 .flatten()
1415 .and_then(crate::formats::audio::recording);
1416 let view_generation = state.len_generation();
1417 let dataset_generation = self.dataset_generation;
1418 let kept_entry = self.kept_quality_entry(&plan.sample());
1419 let kept = kept_entry.map(|kept| kept.rows.clone());
1420 let kept_source = kept_entry
1423 .filter(|_| plan.compute == data_quality::QualityCompute::Sample)
1424 .map(|kept| kept.source.clone());
1425 let mut identity = self.quality_source_identity(state, &plan.scope);
1426 let copy_job = match self.quality_copy_plan(&plan) {
1427 data_quality::CopyPlan::Kept { .. } => self
1428 .quality_copy_kept()
1429 .cloned()
1430 .map_or(QualityCopyJob::Source, QualityCopyJob::Kept),
1431 data_quality::CopyPlan::Fetch { .. } => match state.remote_objects() {
1432 Some(objects) => QualityCopyJob::Fetch {
1433 objects,
1434 root: self.quality_copies_root(),
1435 },
1436 None => QualityCopyJob::Source,
1437 },
1438 _ => QualityCopyJob::Source,
1439 };
1440 #[cfg(feature = "cloud")]
1441 let (cloud, runtime) = (self.app_config.cloud.clone(), self.runtime.clone());
1442 let mut source = source;
1445 if plan.compute == data_quality::QualityCompute::Full
1446 && let Some(source) = source.as_mut()
1447 {
1448 source.conflict_scan = state.quality_conflict_scan();
1449 }
1450 let started = self.start_job(
1453 Job::Analysis(jobs::AnalysisRun::default()),
1454 Some("Profiling data quality..."),
1455 );
1456 let ticket = started.ticket();
1457 let phases = self.events.clone();
1458 let watch = data_quality::QualityWatch::new(move |phase| {
1459 let _ = phases.send(AppEvent::JobProgress {
1460 ticket,
1461 progress: Progress::QualityPhase(phase),
1462 });
1463 });
1464 if let Some(progress) = self.analysis_modal.computing.as_mut() {
1465 progress.read = Some(watch.read().clone());
1466 }
1467 if let Some(Job::Analysis(run)) = self.jobs.job_mut(ticket) {
1468 run.watch = Some(watch.clone());
1469 }
1470 started.run(&self.runtime, move |worker| {
1471 match kept_source {
1473 Some(source) => identity = source,
1474 None => identity.stat(),
1475 }
1476 let lf = if source_scope {
1477 data_quality::prepare_source_quality_scan(lf, source.as_ref())
1478 .map_err(|error| format!("{error}"))?
1479 } else {
1480 lf
1481 };
1482 let lf = data_quality::apply_quality_scope(lf, &plan.scope, source.as_ref())
1483 .map_err(|error| format!("{error}"))?;
1484 let fetch = |objects: &[crate::cloud::local_copy::RemoteObject], root: &Path| {
1486 #[cfg(feature = "cloud")]
1487 {
1488 Self::fetch_quality_copy(objects, root, &cloud, &runtime, watch.read())
1489 }
1490 #[cfg(not(feature = "cloud"))]
1491 {
1492 let _ = (objects, root);
1493 Err(color_eyre::eyre::eyre!("Built without cloud support"))
1494 }
1495 };
1496 let kept_copy = |copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>| {
1497 worker.send(AppEvent::BackgroundQualityCopyKept {
1498 dataset_generation,
1499 copy,
1500 });
1501 };
1502 let (lf, held) =
1503 Self::quality_scope_on_copy(lf, copy_job, &watch, fetch, kept_copy)
1504 .map_err(|error| format!("{error}"))?;
1505 let (results, rows) = crate::analysis::data_quality::compute_data_quality_watched(
1506 &lf,
1507 cached_rows,
1508 &plan,
1509 source.as_ref(),
1510 streaming,
1511 kept.as_deref(),
1512 &watch,
1513 );
1514 let results = match (results, audio) {
1515 (Ok(mut results), Some(audio)) => {
1516 crate::analysis::data_quality::add_signal_observations(
1517 &mut results,
1518 &audio,
1519 &watch,
1520 )
1521 .map(|()| results)
1522 }
1523 (results, _) => results,
1524 };
1525 drop(held);
1528 let kept = rows.map(|rows| KeptQualitySample {
1529 dataset_generation,
1530 view_generation,
1531 sample: plan.sample(),
1532 rows: std::sync::Arc::new(rows),
1533 source: identity.clone(),
1534 });
1535 match results {
1536 Ok(mut results) => {
1537 results.source = Some(Box::new(identity));
1538 Ok(Answer::DataQuality {
1539 results: Box::new(results),
1540 kept,
1541 plan: Box::new(plan),
1542 })
1543 }
1544 Err(error) => {
1545 if let Some(kept) = kept {
1547 worker.send(AppEvent::BackgroundQualitySampleKept { kept });
1548 }
1549 Err(format!("{error}"))
1550 }
1551 }
1552 });
1553 } else {
1554 self.analysis_modal.computing = None;
1555 self.busy = false;
1556 }
1557 None
1558 }
1559}