Skip to main content

datui_lib/analysis/
quality_runs.rs

1//! Data quality runs: the Setup plan, the run itself, the evidence rows it opens,
2//! and the samples, copies and cached results it keeps within its memory budget.
3
4use 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/// What Data Quality runs keep: results, samples and local copies within the memory budget, and
22/// the evidence view's way back.
23#[derive(Default)]
24pub struct QualityRuns {
25    /// Reports, newest first, within [`QUALITY_MEMORY_BUDGET`](crate::analysis::quality_memory::QUALITY_MEMORY_BUDGET).
26    pub(crate) cache: Vec<QualityCacheEntry>,
27    /// See [`KeptQualitySample`]. Newest first, within [`QUALITY_MEMORY_BUDGET`](crate::analysis::quality_memory::QUALITY_MEMORY_BUDGET).
28    pub(crate) samples: Vec<KeptQualitySample>,
29    /// Acquisitions the budget released, newest first: (dataset, view, sample).
30    pub(crate) released: Vec<(u64, u64, sampling::Sample)>,
31    /// [`QUALITY_MEMORY_BUDGET`](crate::analysis::quality_memory::QUALITY_MEMORY_BUDGET), smaller in a test that fills it.
32    pub(crate) memory_budget: usize,
33    /// Local copies full scans read instead of a remote source, newest first, within
34    /// `analysis.quality_local_copy`. Removed from disk when released, when the dataset
35    /// is reopened or replaced, and at exit.
36    pub(crate) copies: Vec<RetainedCopy>,
37    /// The dataset whose copy was released, so Setup says why Run fetches again.
38    pub(crate) copy_released: Option<u64>,
39    /// The dataset whose copy did not read as its source: full scans read the source,
40    /// and Setup says why.
41    pub(crate) copy_unusable: Option<u64>,
42    /// Free bytes in the cache directory and when asked: Setup redraws often, and this
43    /// only feeds a line of text.
44    pub(crate) copy_free: std::sync::Mutex<Option<(std::time::Instant, Option<u64>)>>,
45    /// The table an analysis drill replaced, shown again on Esc.
46    pub(crate) evidence_return: Option<Box<DataTableState>>,
47    pub(crate) evidence_label: Option<String>,
48}
49
50impl QualityRuns {
51    /// Forget the last dataset's runs: its objects may have changed, so the next run
52    /// fetches again.
53    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    /// The selected finding's rows: from kept rows at once, or staged as a read that
67    /// says what it reads and waits for Enter. Nothing for a finding with no rows.
68    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    /// The rows an interval's count counted: from kept rows, or staged as a read.
89    /// Nothing for a count of zero.
90    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    /// The rows the on-screen report measured while kept: a sampled run's, same
114    /// dataset, view and sample. `None` after a full scan or once released.
115    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    /// Open rows a result counted. Kept rows are cut in memory; otherwise the read is
125    /// staged with what it reads, and only Enter reads.
126    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        // A sample as large as the scope reads every row but was not a full scan: its rows
158        // were kept, and since released.
159        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            // The table counts matches (reading every row), then reads the rows it shows.
194            _ => (
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    /// Enter on a staged read: read the rows it named, as it said.
233    pub(crate) fn confirm_evidence_read(&mut self) -> Option<AppEvent> {
234        // Beside a cancelled run still reading, the read stays staged for later.
235        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            // A column a file lacks, or holds in a type the scan cannot read, has no value to
250            // filter on: its rows are those files' rows.
251            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    /// Every repeating row of `scope`, read in one pass off the UI thread and shown as
260    /// a table: a full scan's duplicate finding, once asked.
261    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        // Binary values grouped as the run grouped them: one stub for all, so groups the
277        // check counted are not split.
278        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    /// The rows of `scope` matching `predicate`, shown in place of the table until Esc.
303    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    /// What the data offers Setup's choices, from the schema and rows on screen; reads
364    /// nothing.
365    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        // A directory whose files agree opens as one scan with no partition columns; its
374        // directory names still count.
375        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        // Date and time columns, then text read as time: a window can split by either.
384        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    /// The columns a time role can take: date and time columns, then text (read
420    /// through a Text as time format).
421    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    /// The columns Column intent lists: the draft's scope's, from the schema.
432    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    /// What a run reads from, as the app knows without reading: the origin, its files,
444    /// and the view's effect on rows when the scope is the view.
445    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        // A local path made absolute so the report names the file anywhere; no filesystem
491        // access.
492        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    /// The Setup setting the current Data Quality page needs before it can show
510    /// anything; Enter opens it, and the footer says so.
511    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    /// Whether the scope has a column read as time (date/time, or text with a format):
522    /// what an empty Trends page points to.
523    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    /// Whether the retained rows are the rows `plan` reads, so the run starts from
533    /// them. Any grain works: every column and row position is kept.
534    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    /// Where a run of `plan` gets exact segment totals: retained rows' counts, the pass
540    /// reading a new sample, or a read of their own.
541    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    /// Whether `plan`'s rows were read this session and released to the budget, so Run
555    /// reads them again.
556    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    /// Whether the dataset is one Parquet or IPC file, the only kind the sampler reads
578    /// seeded runs of.
579    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    /// Whether a random sample of `plan` reads seeded runs of one file rather than
600    /// streaming every row, as Setup's Read says, judged from path and view (the
601    /// sampler's own test needs the built plan). Yes only where the scan is read as
602    /// loaded: the whole source, or a view selecting no rows (samples ignore the sort).
603    /// When unsure, Setup names the longer read.
604    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                // Read in the order on screen, sort included.
613                data_quality::QualityScope::FirstRows(_)
614                | data_quality::QualityScope::ViewRows { .. } => {
615                    state.source_file_count() == Some(1)
616                }
617                _ => false,
618            }
619    }
620
621    /// Whether a random sample of `plan` may read seeded runs: where
622    /// [`Self::quality_reads_blocks`] is sure, and wherever the view (a query too) may
623    /// still read the scan as loaded. Leans yes: a yes makes Setup name a count pass,
624    /// and a run that streams instead reads less than Setup said, never more.
625    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    /// Whether the session cache holds a report measuring what `plan` does on this
641    /// view; expected windows are checked against it, not read.
642    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    /// What runs kept for reuse on this dataset, as `d` in Setup releases it: sampled
658    /// rows in memory and a full scan's local copy. `None` when neither.
659    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    /// `d` in Setup: release every kept row and the full scan's local copy (removed
682    /// from disk). Runs that would reuse them read again, as Setup's Read says.
683    /// Reports stay: showing one reads nothing.
684    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        // A run still reading the copy holds it until it ends; then the files go.
706        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    /// Rows a sampled run read, when they are the rows `sample` names now (same
728    /// dataset, view, sample).
729    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    /// Keep what a run read, newest first, replacing an older copy of the same rows: a
747    /// later run returns them with its added counts.
748    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    /// Hold reports and retained rows to the memory budget
763    /// ([`QUALITY_MEMORY_BUDGET`](crate::analysis::quality_memory::QUALITY_MEMORY_BUDGET)). First to go: reports
764    /// whose rows are retained (remade without reading); then the oldest rows (reread
765    /// next run, as Setup says); last, row-less full-scan reports (dearest to remake).
766    /// The newest report and rows always stay.
767    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            // Compared as the plan compares, and kept under the windows it expects now.
841            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        // One report per measurement: a plan expecting other windows replaces it.
867        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    /// The scope a full scan reads: `lf` over a local copy of its remote objects,
886    /// fetched first when `job` says. `kept` hears the copy to keep, or `None` if it
887    /// does not read as the source (later full scans read the source). The copy is
888    /// returned for the caller to hold while the passes read it.
889    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        // The copy must read as the source does, or the source is read as before.
915        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    /// Copy `objects` under `root`, streaming each as it arrives; a cancel stops at the
936    /// next chunk and removes the partial copy.
937    #[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                // Rewritten since opened (perhaps same size): the copy would not be the dataset on
962                // screen.
963                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    /// Where Data Quality's local copies are written.
976    fn quality_copies_root(&self) -> PathBuf {
977        self.cache
978            .cache_dir()
979            .join(crate::cloud::local_copy::COPIES_DIR)
980    }
981
982    /// `analysis.quality_local_copy`, in bytes.
983    fn quality_copy_limit(&self) -> u64 {
984        self.app_config.analysis.quality_local_copy.bytes()
985    }
986
987    /// Bytes on disk in the copies kept.
988    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    /// The copy this dataset's objects were fetched into this session, while kept.
997    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    /// Free bytes where copies are written, asked at most every few seconds.
1014    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    /// How a run of `plan` gets a remote source's rows, from what the open learned;
1030    /// at most a stat of the cache directory.
1031    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    /// Whether this dataset's copy was released this session, so Run fetches again.
1072    pub(crate) fn quality_copy_released(&self) -> bool {
1073        self.quality.copy_released == Some(self.dataset_generation)
1074    }
1075
1076    /// Keep a fetched copy, newest first; older ones go past the budget, the newest
1077    /// always stays. `None` means the copy did not read as its source: any kept one
1078    /// goes too.
1079    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    /// Read `sample` off the UI thread and show its rows (all, or a finding's under its
1109    /// label). Redrawn from its seed, so these are the rows the tool measured.
1110    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        // The rows Data Quality just measured, if they are the rows asked: cut from memory
1119        // rather than redrawn.
1120        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            // Shown columns are the table's; a finding is cut from every column the run read
1139            // first, so duplicates are judged as the run judged them.
1140            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                        // Every column the run grouped by: the scope's own, not the kept row numbers.
1169                        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                    // Files are read from the scope, never from a sample.
1190                    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    /// Mirror the shared sample into the Data Quality plan, which carries it to the
1203    /// engine and the cache key. Metadata-only stays metadata-only.
1204    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    /// Open Setup with the plan staged: edits wait for Run, and Esc restores the plan.
1210    /// Reopening while open changes nothing.
1211    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    /// Esc on Setup: staged edits go and the report returns. With no report, Setup stays
1228    /// and the cursor goes to the tools.
1229    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    /// A full scan's confirmation: what it reads, what it fetches from a remote source,
1250    /// and that it writes nothing there.
1251    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    /// What stops Setup from running, on its own line: a time window on text with no
1268    /// format.
1269    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    /// Run from Setup, the one place a run starts: the draft becomes the plan and its
1287    /// sample every tool's; the report is shown if cached, else read once. Waits while
1288    /// a cancelled run still stops (two reads can run memory out). A full scan asks
1289    /// first; `confirmed` is that question's Yes, and Esc leaves the draft staged.
1290    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        // A report already here, on screen or cached, reads nothing: nothing to confirm.
1301        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            // Asked with the one confirmation; its Yes comes back here.
1312            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        // Only expected windows or the comparison changed: the report on screen holds every
1331        // count and segment needed, so it is relabeled, not reread.
1332        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    /// Commit the draft: Setup closes and its sample becomes every tool's. Other tools'
1367    /// results were of the old sample and go; Data Quality's last report stays,
1368    /// labeled, until the run replaces it.
1369    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    /// Run the data quality check on the plan Setup committed.
1386    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            // An audio file's signal checks read its samples whole: a full run's.
1412            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            // A sampled run on rows already read is labeled as their read was: the file may
1421            // have changed since.
1422            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            // Only a confirmed full scan pays to read values a type conflict hides, as its
1443            // access plan promised.
1444            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            // Each stage comes back as progress, so a cancelled run's are dropped; Esc stops the
1451            // run through the job's watch.
1452            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                // A stat of a local file as the run begins, not a read.
1472                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                // Held to the end of the run: the copy stays on disk while read, released or not.
1485                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                // Released before the answer goes out, so a `d` handled on arrival finds the app's
1526                // handle the last.
1527                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                        // Stopped after the sample was read: the read is kept.
1546                        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}