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 open = async move {
956 let got = store.get(&path).await.map_err(|e| e.to_string())?;
957 if let (Some(listed), Some(fetched)) = (&listed, &got.meta.e_tag)
960 && !crate::cloud::local_copy::same_etag(listed, fetched)
961 {
962 return Err("it changed since it opened. Open the dataset again".to_string());
963 }
964 Ok((got.into_stream(), None))
965 };
966 let watch = stop.clone();
967 crate::cloud::download::stream_into(runtime, open, move || watch.stopped(), write)
968 })
969 }
970
971 fn quality_copies_root(&self) -> PathBuf {
973 self.cache
974 .cache_dir()
975 .join(crate::cloud::local_copy::COPIES_DIR)
976 }
977
978 fn quality_copy_limit(&self) -> u64 {
980 self.app_config.analysis.quality_local_copy.bytes()
981 }
982
983 pub fn quality_copy_bytes(&self) -> u64 {
985 self.quality
986 .copies
987 .iter()
988 .map(|kept| kept.copy.bytes())
989 .sum()
990 }
991
992 fn quality_copy_kept(&self) -> Option<&Arc<crate::cloud::local_copy::LocalCopy>> {
994 let state = self.data_table_state.as_ref()?;
995 self.quality
996 .copies
997 .iter()
998 .find(|kept| {
999 kept.dataset_generation == self.dataset_generation
1000 && state.each_remote_object().is_some_and(|mut objects| {
1001 objects.all(|object| {
1002 object.is_some_and(|object| kept.copy.covers(&object.url))
1003 })
1004 })
1005 })
1006 .map(|kept| &kept.copy)
1007 }
1008
1009 fn quality_copy_free_space(&self) -> Option<u64> {
1011 let root = self.quality_copies_root();
1012 let Ok(mut cached) = self.quality.copy_free.lock() else {
1013 return crate::cloud::local_copy::free_space(&root);
1014 };
1015 match *cached {
1016 Some((asked, free)) if asked.elapsed() < std::time::Duration::from_secs(5) => free,
1017 _ => {
1018 let free = crate::cloud::local_copy::free_space(&root);
1019 *cached = Some((std::time::Instant::now(), free));
1020 free
1021 }
1022 }
1023 }
1024
1025 pub(crate) fn quality_copy_plan(
1028 &self,
1029 plan: &data_quality::DataQualityPlan,
1030 ) -> data_quality::CopyPlan {
1031 use data_quality::{CopyPlan, NoCopy};
1032 let Some(state) = self.data_table_state.as_ref() else {
1033 return CopyPlan::NotApplicable;
1034 };
1035 if plan.compute != data_quality::QualityCompute::Full || !state.is_remote_source() {
1036 return CopyPlan::NotApplicable;
1037 }
1038 if !state.quality_reads_whole_source(&plan.scope) {
1039 return CopyPlan::Passes(NoCopy::PartOfTheSource);
1040 }
1041 if let Some(copy) = self.quality_copy_kept() {
1042 return CopyPlan::Kept {
1043 bytes: copy.bytes(),
1044 objects: copy.objects(),
1045 };
1046 }
1047 let limit = self.quality_copy_limit();
1048 if limit == 0 {
1049 return CopyPlan::Passes(NoCopy::Off);
1050 }
1051 if self.quality.copy_unusable == Some(self.dataset_generation) {
1052 return CopyPlan::Passes(NoCopy::Unusable);
1053 }
1054 let Some((bytes, objects)) = state.remote_objects_size() else {
1055 return CopyPlan::Passes(NoCopy::SizeUnknown);
1056 };
1057 if bytes > limit {
1058 return CopyPlan::Passes(NoCopy::TooLarge { bytes, limit });
1059 }
1060 let free = self.quality_copy_free_space();
1061 if free.is_none_or(|free| bytes > free) {
1062 return CopyPlan::Passes(NoCopy::NoRoom { bytes, free });
1063 }
1064 CopyPlan::Fetch { bytes, objects }
1065 }
1066
1067 pub(crate) fn quality_copy_released(&self) -> bool {
1069 self.quality.copy_released == Some(self.dataset_generation)
1070 }
1071
1072 pub(crate) fn retain_quality_copy(
1076 &mut self,
1077 dataset_generation: u64,
1078 copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>,
1079 ) {
1080 if dataset_generation != self.dataset_generation {
1081 return;
1082 }
1083 let Some(copy) = copy else {
1084 self.quality
1085 .copies
1086 .retain(|kept| kept.dataset_generation != dataset_generation);
1087 self.quality.copy_unusable = Some(dataset_generation);
1088 return;
1089 };
1090 self.quality.copies.insert(
1091 0,
1092 RetainedCopy {
1093 dataset_generation,
1094 copy,
1095 },
1096 );
1097 self.quality.copy_released = None;
1098 let limit = self.quality_copy_limit();
1099 while self.quality.copies.len() > 1 && self.quality_copy_bytes() > limit {
1100 self.quality.copies.pop();
1101 }
1102 }
1103
1104 pub(crate) fn read_sample_rows(
1107 &mut self,
1108 sample: sampling::Sample,
1109 evidence: Option<(quality_report::EvidenceRows, String)>,
1110 ) -> Option<AppEvent> {
1111 let state = self.data_table_state.as_ref()?;
1112 let (source, known_total) = Self::sample_source_for(state, &sample.scope);
1113 let streaming = self.app_config.performance.streaming;
1114 let kept = self.kept_quality_sample(&sample).map(|kept| {
1117 let columns: Vec<_> = state
1118 .schema()
1119 .iter_names()
1120 .filter(|name| kept.df().column(name.as_str()).is_ok())
1121 .map(|name| polars::prelude::col(name.clone()))
1122 .collect();
1123 (kept, columns)
1124 });
1125 if kept.is_none() && self.read_waits_for_cancelled() {
1126 return None;
1127 }
1128 self.analysis_modal.computing = Some(AnalysisProgress::new(if evidence.is_some() {
1129 "Reading the matching sampled rows"
1130 } else {
1131 "Reading the sample"
1132 }));
1133 self.spawn_job(Job::SampleRows, Some("Reading the sample..."), move |_| {
1134 let (rows, columns) = match kept {
1137 Some((kept, columns)) => (Ok(kept.analysis_rows(kept.df().clone())), Some(columns)),
1138 None => (
1139 source
1140 .cut(&sample.scope)
1141 .and_then(|lf| sampling::read(&lf, &sample, known_total, streaming)),
1142 None,
1143 ),
1144 };
1145 let shown = |df: polars::prelude::DataFrame| match &columns {
1146 Some(columns) => polars::prelude::IntoLazy::lazy(df)
1147 .select(columns.clone())
1148 .collect()
1149 .map_err(color_eyre::eyre::Report::from),
1150 None => Ok(df),
1151 };
1152 let read = rows.and_then(|rows| {
1153 let label = format!(
1154 "Sample {} {}",
1155 crate::glyphs::get().middot,
1156 sample.outcome(
1157 rows.total_rows,
1158 rows.sample_size,
1159 rows.per_value.as_ref().map(|per_value| per_value.kept),
1160 )
1161 );
1162 match evidence {
1163 Some((quality_report::EvidenceRows::Duplicates, label)) => {
1164 let keys = rows
1166 .df
1167 .get_column_names()
1168 .into_iter()
1169 .filter(|name| !name.starts_with("__datui"))
1170 .cloned()
1171 .collect::<Vec<_>>();
1172 let df = data_quality::duplicate_rows(
1173 polars::prelude::IntoLazy::lazy(rows.df),
1174 &keys,
1175 streaming,
1176 )?;
1177 Ok((shown(df)?, label))
1178 }
1179 Some((quality_report::EvidenceRows::Matching(predicate), label)) => {
1180 let df = polars::prelude::IntoLazy::lazy(rows.df)
1181 .filter(predicate)
1182 .collect()?;
1183 Ok((shown(df)?, label))
1184 }
1185 Some((quality_report::EvidenceRows::Files(_), label)) => {
1187 Ok((shown(rows.df)?, label))
1188 }
1189 None => Ok((shown(rows.df)?, label)),
1190 }
1191 });
1192 let (df, label) = read.map_err(|error| format!("{error}"))?;
1193 Ok(Answer::Sample { df, label })
1194 });
1195 None
1196 }
1197
1198 pub(crate) fn sync_quality_plan(&mut self) {
1201 let sample = self.analysis_modal.sample.clone();
1202 self.analysis_modal.quality.plan.adopt_sample(&sample);
1203 }
1204
1205 pub(crate) fn open_quality_setup(&mut self) {
1208 use data_quality::QualityPage;
1209 let modal = &mut self.analysis_modal;
1210 if !modal.quality.page.is_setup() {
1211 modal.quality.setup_return = modal.quality.page.tab();
1212 }
1213 if modal.quality.setup_before.is_none() {
1214 modal.quality.setup_before = Some(modal.quality.plan.clone());
1215 }
1216 if modal.quality.page != QualityPage::Setup {
1217 modal.set_quality_page(QualityPage::Setup);
1218 modal.quality.plan_field = 0;
1219 }
1220 modal.focus = analysis_modal::AnalysisFocus::Main;
1221 }
1222
1223 pub(crate) fn leave_quality_setup(&mut self) {
1226 use data_quality::QualityPage;
1227 let modal = &mut self.analysis_modal;
1228 if let Some(before) = modal.quality.setup_before.take() {
1229 modal.quality.plan = before;
1230 }
1231 modal.quality.setup_note = None;
1232 modal.quality.picker = None;
1233 if modal.quality.results.is_some() {
1234 let back = match modal.quality.setup_return {
1235 page if page.is_setup() => QualityPage::Overview,
1236 page => page,
1237 };
1238 modal.set_quality_page(back);
1239 } else {
1240 modal.set_quality_page(QualityPage::Setup);
1241 modal.focus = analysis_modal::AnalysisFocus::Sidebar;
1242 }
1243 }
1244
1245 fn quality_full_scan_question(&self, plan: &data_quality::DataQualityPlan) -> String {
1248 let mut lines = vec![
1249 "Run a full scan?".to_string(),
1250 String::new(),
1251 "Reads: every eligible row, up to the whole source".to_string(),
1252 ];
1253 if let data_quality::CopyPlan::Fetch { bytes, .. } = self.quality_copy_plan(plan) {
1254 lines.push(format!(
1255 "Fetch: {} once, to a local copy",
1256 crate::numfmt::bytes(bytes)
1257 ));
1258 }
1259 lines.push("Source writes: none".to_string());
1260 lines.join("\n")
1261 }
1262
1263 fn quality_setup_problem(&self) -> Option<String> {
1266 let plan = &self.analysis_modal.quality.plan;
1267 let schema = self.data_table_state.as_ref()?.quality_schema(&plan.scope);
1268 match &plan.grain {
1269 data_quality::QualityGrain::TimeWindows { column, .. }
1270 if plan.compute != data_quality::QualityCompute::Metadata
1271 && !plan.reads_as_time(column, schema) =>
1272 {
1273 Some(format!(
1274 "{column}: text, no format {} set Text as time",
1275 crate::glyphs::get().middot
1276 ))
1277 }
1278 _ => None,
1279 }
1280 }
1281
1282 pub(crate) fn run_quality_setup(&mut self, confirmed: bool) -> Option<AppEvent> {
1287 use data_quality::QualityPage;
1288 if self.cancelled_analysis_running().is_some() {
1289 self.analysis_modal.quality.setup_note = Some(QUALITY_RUN_WAITS.to_string());
1290 return None;
1291 }
1292 if let Some(problem) = self.quality_setup_problem() {
1293 self.analysis_modal.quality.setup_note = Some(problem);
1294 return None;
1295 }
1296 let plan = &self.analysis_modal.quality.plan;
1298 let here = (self.analysis_modal.quality.results.is_some()
1299 && self
1300 .analysis_modal
1301 .quality
1302 .last_plan
1303 .as_ref()
1304 .is_some_and(|last| last.same_measurement(plan)))
1305 || self.quality_cached(plan);
1306 if plan.requires_confirmation() && !here && !confirmed {
1307 let message = self.quality_full_scan_question(plan);
1309 self.confirmation_modal
1310 .show(message, Confirm::QualityFullScan);
1311 self.confirmation_modal.yes_label = "Run";
1312 return None;
1313 }
1314 self.commit_quality_plan();
1315 let modal = &mut self.analysis_modal;
1316 if modal.quality.results.is_some()
1317 && modal.quality.last_plan.as_ref() == Some(&modal.quality.plan)
1318 {
1319 let back = match modal.quality.setup_return {
1320 page if page.is_setup() => QualityPage::Overview,
1321 page => page,
1322 };
1323 modal.set_quality_page(back);
1324 return None;
1325 }
1326 if let (Some(results), Some(last)) = (
1329 modal.quality.results.as_ref(),
1330 modal.quality.last_plan.as_ref(),
1331 ) && last.same_measurement(&modal.quality.plan)
1332 {
1333 let mut results = results.clone();
1334 let plan = modal.quality.plan.clone();
1335 let page = if last.compares_differently(&plan) {
1336 results.compare_segments(&plan);
1337 QualityPage::Segments
1338 } else {
1339 QualityPage::Trends
1340 };
1341 modal.quality.results = Some(results.clone());
1342 modal.quality.last_plan = Some(plan.clone());
1343 modal.set_quality_page(page);
1344 self.cache_quality_result(&results, plan);
1345 return None;
1346 }
1347 if self.restore_cached_quality() {
1348 return None;
1349 }
1350 self.analysis_modal.quality.from_cache = false;
1351 let mut progress = AnalysisProgress::new("Preparing the plan");
1352 if self.quality_kept_serves(&self.analysis_modal.quality.plan) {
1353 progress.reuse = Some("Starts from rows a run already read".to_string());
1354 }
1355 self.analysis_modal.computing = Some(progress);
1356 self.busy = true;
1357 Some(AppEvent::AnalysisCompute(
1358 analysis_modal::AnalysisTool::DataQuality,
1359 ))
1360 }
1361
1362 fn commit_quality_plan(&mut self) {
1366 let modal = &mut self.analysis_modal;
1367 let sample = modal.quality.plan.sample();
1368 if sample != modal.sample {
1369 modal.describe_results = None;
1370 modal.distribution_results = None;
1371 modal.correlation_results = None;
1372 }
1373 modal.sample = sample;
1374 modal.sample_dataset = Some(self.dataset_generation);
1375 modal.sample_run_for = Some(self.dataset_generation);
1376 modal.quality.setup_before = None;
1377 modal.quality.setup_note = None;
1378 modal.quality.picker = None;
1379 }
1380
1381 pub(crate) fn run_quality_compute(&mut self) -> Option<AppEvent> {
1383 if let Some(state) = &self.data_table_state {
1384 let plan = self.analysis_modal.quality.plan.clone();
1385 let source_scope = plan.scope.uses_source();
1386 let (lf, source, cached_rows) = if source_scope {
1387 let (lf, source) = state.data_quality_source_scan();
1388 (lf, source, None)
1389 } else {
1390 let ordered = matches!(
1391 plan.scope,
1392 data_quality::QualityScope::FirstRows(_)
1393 | data_quality::QualityScope::ViewRows { .. }
1394 );
1395 let (lf, source) = state.data_quality_scan(ordered);
1396 let rows = state.num_rows_if_valid().map(|rows| match &plan.scope {
1397 data_quality::QualityScope::CurrentView => rows,
1398 data_quality::QualityScope::FirstRows(limit) => rows.min(*limit),
1399 data_quality::QualityScope::ViewRows { start, end } => {
1400 rows.min(*end).saturating_sub(start.saturating_sub(1))
1401 }
1402 _ => unreachable!(),
1403 });
1404 (lf, source, rows)
1405 };
1406 let streaming = state.polars_streaming();
1407 let audio = (plan.compute == data_quality::QualityCompute::Full)
1409 .then(|| state.window_for_quality(&plan.scope))
1410 .flatten()
1411 .and_then(crate::formats::audio::recording);
1412 let view_generation = state.len_generation();
1413 let dataset_generation = self.dataset_generation;
1414 let kept_entry = self.kept_quality_entry(&plan.sample());
1415 let kept = kept_entry.map(|kept| kept.rows.clone());
1416 let kept_source = kept_entry
1419 .filter(|_| plan.compute == data_quality::QualityCompute::Sample)
1420 .map(|kept| kept.source.clone());
1421 let mut identity = self.quality_source_identity(state, &plan.scope);
1422 let copy_job = match self.quality_copy_plan(&plan) {
1423 data_quality::CopyPlan::Kept { .. } => self
1424 .quality_copy_kept()
1425 .cloned()
1426 .map_or(QualityCopyJob::Source, QualityCopyJob::Kept),
1427 data_quality::CopyPlan::Fetch { .. } => match state.remote_objects() {
1428 Some(objects) => QualityCopyJob::Fetch {
1429 objects,
1430 root: self.quality_copies_root(),
1431 },
1432 None => QualityCopyJob::Source,
1433 },
1434 _ => QualityCopyJob::Source,
1435 };
1436 #[cfg(feature = "cloud")]
1437 let (cloud, runtime) = (self.app_config.cloud.clone(), self.runtime.clone());
1438 let mut source = source;
1441 if plan.compute == data_quality::QualityCompute::Full
1442 && let Some(source) = source.as_mut()
1443 {
1444 source.conflict_scan = state.quality_conflict_scan();
1445 }
1446 let started = self.start_job(
1449 Job::Analysis(jobs::AnalysisRun::default()),
1450 Some("Profiling data quality..."),
1451 );
1452 let ticket = started.ticket();
1453 let phases = self.events.clone();
1454 let watch = data_quality::QualityWatch::new(move |phase| {
1455 let _ = phases.send(AppEvent::JobProgress {
1456 ticket,
1457 progress: Progress::QualityPhase(phase),
1458 });
1459 });
1460 if let Some(progress) = self.analysis_modal.computing.as_mut() {
1461 progress.read = Some(watch.read().clone());
1462 }
1463 if let Some(Job::Analysis(run)) = self.jobs.job_mut(ticket) {
1464 run.watch = Some(watch.clone());
1465 }
1466 started.run(&self.runtime, move |worker| {
1467 match kept_source {
1469 Some(source) => identity = source,
1470 None => identity.stat(),
1471 }
1472 let lf = if source_scope {
1473 data_quality::prepare_source_quality_scan(lf, source.as_ref())
1474 .map_err(|error| format!("{error}"))?
1475 } else {
1476 lf
1477 };
1478 let lf = data_quality::apply_quality_scope(lf, &plan.scope, source.as_ref())
1479 .map_err(|error| format!("{error}"))?;
1480 let fetch = |objects: &[crate::cloud::local_copy::RemoteObject], root: &Path| {
1482 #[cfg(feature = "cloud")]
1483 {
1484 Self::fetch_quality_copy(objects, root, &cloud, &runtime, watch.read())
1485 }
1486 #[cfg(not(feature = "cloud"))]
1487 {
1488 let _ = (objects, root);
1489 Err(color_eyre::eyre::eyre!("Built without cloud support"))
1490 }
1491 };
1492 let kept_copy = |copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>| {
1493 worker.send(AppEvent::BackgroundQualityCopyKept {
1494 dataset_generation,
1495 copy,
1496 });
1497 };
1498 let (lf, held) =
1499 Self::quality_scope_on_copy(lf, copy_job, &watch, fetch, kept_copy)
1500 .map_err(|error| format!("{error}"))?;
1501 let (results, rows) = crate::analysis::data_quality::compute_data_quality_watched(
1502 &lf,
1503 cached_rows,
1504 &plan,
1505 source.as_ref(),
1506 streaming,
1507 kept.as_deref(),
1508 &watch,
1509 );
1510 let results = match (results, audio) {
1511 (Ok(mut results), Some(audio)) => {
1512 crate::analysis::data_quality::add_signal_observations(
1513 &mut results,
1514 &audio,
1515 &watch,
1516 )
1517 .map(|()| results)
1518 }
1519 (results, _) => results,
1520 };
1521 drop(held);
1524 let kept = rows.map(|rows| KeptQualitySample {
1525 dataset_generation,
1526 view_generation,
1527 sample: plan.sample(),
1528 rows: std::sync::Arc::new(rows),
1529 source: identity.clone(),
1530 });
1531 match results {
1532 Ok(mut results) => {
1533 results.source = Some(Box::new(identity));
1534 Ok(Answer::DataQuality {
1535 results: Box::new(results),
1536 kept,
1537 plan: Box::new(plan),
1538 })
1539 }
1540 Err(error) => {
1541 if let Some(kept) = kept {
1543 worker.send(AppEvent::BackgroundQualitySampleKept { kept });
1544 }
1545 Err(format!("{error}"))
1546 }
1547 }
1548 });
1549 } else {
1550 self.analysis_modal.computing = None;
1551 self.busy = false;
1552 }
1553 None
1554 }
1555}