use crate::analysis::analysis_modal::AnalysisProgress;
use crate::analysis::quality_memory::{
KeptQualitySample, QUALITY_RELEASED_REMEMBERED, QualityCacheEntry, QualityCopyJob, RetainedCopy,
};
use crate::app::feedback::Confirm;
use crate::app::jobs::{Answer, Job, Progress};
use crate::export::export_modal::ExportFormat;
use crate::table::DataTableState;
use crate::{
App, AppEvent, QUALITY_RUN_WAITS, analysis::analysis_modal, analysis::data_quality,
analysis::quality_report, analysis::sampling, app::jobs, glyphs, numfmt, widgets,
};
use color_eyre::Result;
use polars::prelude::LazyFrame;
use std::path::{Path, PathBuf};
use std::sync::Arc;
#[derive(Default)]
pub struct QualityRuns {
pub(crate) cache: Vec<QualityCacheEntry>,
pub(crate) samples: Vec<KeptQualitySample>,
pub(crate) released: Vec<(u64, u64, sampling::Sample)>,
pub(crate) memory_budget: usize,
pub(crate) copies: Vec<RetainedCopy>,
pub(crate) copy_released: Option<u64>,
pub(crate) copy_unusable: Option<u64>,
pub(crate) copy_free: std::sync::Mutex<Option<(std::time::Instant, Option<u64>)>>,
pub(crate) evidence_return: Option<Box<DataTableState>>,
pub(crate) evidence_label: Option<String>,
}
impl QualityRuns {
pub(crate) fn reset_for_dataset(&mut self) {
self.cache.clear();
self.samples.clear();
self.released.clear();
self.copies.clear();
self.copy_released = None;
self.copy_unusable = None;
self.evidence_return = None;
self.evidence_label = None;
}
}
impl App {
pub(crate) fn open_quality_evidence(&mut self) -> Option<AppEvent> {
let (_, finding) = self.analysis_modal.selected_finding()?;
let results = self.analysis_modal.quality.results.as_ref()?;
let rows = finding.evidence(results).ok()?;
let sampled = results.precision == data_quality::QualityPrecision::Sampled;
let count = finding.evidence_count(results);
let label = format!(
"Data Quality / {} / {}",
finding.title,
quality_report::columns_label(&finding.columns, 40)
);
let what = format!(
"{} {} {}",
finding.title,
glyphs::get().middot,
quality_report::columns_label(&finding.columns, 40)
);
self.show_quality_rows(rows, label, sampled, what, count)
}
pub(crate) fn open_interval_evidence(&mut self) -> Option<AppEvent> {
let schema = self.data_table_state.as_ref().map(|state| state.schema());
let (predicate, label, count) = self
.analysis_modal
.interval_evidence(schema.map(|schema| schema.as_ref()))?;
let sampled = self
.analysis_modal
.quality
.results
.as_ref()
.is_some_and(|results| results.precision == data_quality::QualityPrecision::Sampled);
let what = label
.trim_start_matches("Data Quality / ")
.replace(" / ", &format!(" {} ", glyphs::get().middot));
self.show_quality_rows(
quality_report::EvidenceRows::Matching(predicate),
label,
sampled,
what,
Some(count),
)
}
pub(crate) fn quality_rows_kept(&self) -> Option<std::sync::Arc<data_quality::QualitySample>> {
self.analysis_modal.quality.results.as_ref()?;
let plan = self.analysis_modal.quality_result_plan();
if plan.compute != data_quality::QualityCompute::Sample {
return None;
}
self.kept_quality_sample(&plan.sample())
}
fn show_quality_rows(
&mut self,
rows: quality_report::EvidenceRows,
label: String,
sampled: bool,
what: String,
count: Option<usize>,
) -> Option<AppEvent> {
let plan = self.analysis_modal.quality_result_plan().clone();
let by_files = matches!(rows, quality_report::EvidenceRows::Files(_));
let label = if sampled {
format!("{label} / sampled")
} else {
label
};
if !by_files && self.quality_rows_kept().is_some() {
return self.read_sample_rows(plan.sample(), Some((rows, label)));
}
let state = self.data_table_state.as_ref()?;
let g = glyphs::get();
let rows_label = |rows: usize| {
format!(
"{} {}",
numfmt::group_chrome(rows),
if rows == 1 { "row" } else { "rows" }
)
};
let scope = match &rows {
quality_report::EvidenceRows::Files(files) => files.clone(),
_ => plan.scope.clone(),
};
let not_kept = if plan.compute == data_quality::QualityCompute::Full {
"a full scan keeps no rows"
} else {
"the rows read are no longer kept"
};
let (why, reads) = match &rows {
quality_report::EvidenceRows::Files(data_quality::QualityScope::SourceFiles(files)) => {
(
"their rows are in the files, not the report".to_string(),
format!(
"the {} named {}",
files.len(),
if files.len() == 1 { "file" } else { "files" }
),
)
}
_ if sampled => (
"the sampled rows are no longer kept".to_string(),
format!(
"the sample again: {} {} {}",
widgets::data_quality::compute_label(&plan),
g.middot,
widgets::data_quality::planned_read_label(state, &plan)
),
),
quality_report::EvidenceRows::Duplicates => (
not_kept.to_string(),
format!(
"every row of {}, once {} {}",
plan.scope.label(),
g.middot,
widgets::data_quality::scope_read_label(state, &plan)
),
),
_ => (
not_kept.to_string(),
format!(
"every row of {} to count them, then the rows on screen {} {}",
plan.scope.label(),
g.middot,
widgets::data_quality::scope_read_label(state, &plan)
),
),
};
let shows = match (&rows, count) {
(quality_report::EvidenceRows::Duplicates, Some(count)) => {
format!("{}, copies together", rows_label(count))
}
(_, Some(count)) => rows_label(count),
(_, None) => "the rows that match".to_string(),
};
let source = if state.is_remote_source() {
"remote, read only"
} else {
"local, read only"
};
self.analysis_modal.quality.evidence_read = Some(analysis_modal::EvidenceRead {
summary: vec![
("Rows", what),
("Why", why),
("Reads", reads),
("Shows", shows),
("Source", source.to_string()),
],
sample: (sampled && !by_files).then(|| plan.sample()),
scope,
rows,
label,
});
None
}
pub(crate) fn confirm_evidence_read(&mut self) -> Option<AppEvent> {
let staged = self.analysis_modal.quality.evidence_read.as_ref()?;
let kept = staged
.sample
.as_ref()
.is_some_and(|sample| self.kept_quality_sample(sample).is_some());
if !kept && self.read_waits_for_cancelled() {
return None;
}
let read = self.analysis_modal.quality.evidence_read.take()?;
if let Some(sample) = read.sample {
return self.read_sample_rows(sample, Some((read.rows, read.label)));
}
let predicate = match read.rows {
quality_report::EvidenceRows::Matching(predicate) => predicate,
quality_report::EvidenceRows::Files(_) => polars::prelude::lit(true),
quality_report::EvidenceRows::Duplicates => {
return self.read_duplicate_rows(read.scope, read.label);
}
};
self.open_quality_scope_rows(&read.scope, predicate, read.label)
}
fn read_duplicate_rows(
&mut self,
scope: data_quality::QualityScope,
label: String,
) -> Option<AppEvent> {
let state = self.data_table_state.as_ref()?;
let (lf, schema) = match state.quality_scope_frame(&scope) {
Ok(frame) => frame,
Err(error) => {
self.error_modal
.show(format!("Cannot open matching rows: {error}"));
return None;
}
};
let keys = schema.iter_names().cloned().collect::<Vec<_>>();
let columns = schema
.iter()
.map(|(name, dtype)| {
if matches!(dtype, polars::prelude::DataType::Binary) {
polars::prelude::lit(crate::table::binary_stub()).alias(name.clone())
} else {
polars::prelude::col(name.clone())
}
})
.collect::<Vec<_>>();
let streaming = self.app_config.performance.streaming;
self.analysis_modal.computing = Some(AnalysisProgress::new("Reading the rows that repeat"));
self.spawn_job(
Job::SampleRows,
Some("Reading the rows that repeat..."),
move |_| {
let df = data_quality::duplicate_rows(lf.select(columns), &keys, streaming)
.map_err(|error| format!("{error}"))?;
Ok(Answer::Sample { df, label })
},
);
None
}
fn open_quality_scope_rows(
&mut self,
scope: &data_quality::QualityScope,
predicate: polars::prelude::Expr,
label: String,
) -> Option<AppEvent> {
let state = self.data_table_state.as_ref()?;
let view = match state.quality_evidence_view(scope, predicate) {
Ok(view) => view,
Err(error) => {
self.error_modal
.show(format!("Cannot open matching rows: {error}"));
return None;
}
};
if let Some(original) = self.data_table_state.replace(view) {
self.quality.evidence_return = Some(Box::new(original));
self.quality.evidence_label = Some(label);
self.step_back();
self.forget_the_rows_read();
self.spawn_async_collect("Loading matching rows...");
}
None
}
pub(crate) fn return_from_quality_evidence(&mut self, reopen_analysis: bool) -> bool {
let Some(original) = self.quality.evidence_return.take() else {
return false;
};
self.jobs.advance();
self.counting.len_count_inflight = None;
self.data_table_state = Some(*original);
self.quality.evidence_label = None;
if reopen_analysis {
self.open_overlay(crate::Overlay::Analysis);
}
self.busy = false;
self.status_message = None;
true
}
pub(crate) fn restore_recent_quality_plan(&mut self) {
let Some(view_generation) = self
.data_table_state
.as_ref()
.map(DataTableState::len_generation)
else {
return;
};
if self.analysis_modal.quality.plan != data_quality::DataQualityPlan::default() {
return;
}
if let Some(cached) = self.quality.cache.iter().find(|entry| {
entry.dataset_generation == self.dataset_generation
&& entry.view_generation == view_generation
}) {
self.analysis_modal.quality.plan = cached.plan.clone();
}
}
pub(crate) fn quality_plan_context(&self) -> analysis_modal::PlanContext {
let Some(state) = self.data_table_state.as_ref() else {
return analysis_modal::PlanContext::default();
};
let plan = &self.analysis_modal.quality.plan;
let scope = &plan.scope;
let schema = state.schema();
let mut partitions = state.partition_columns().unwrap_or_default().to_vec();
if partitions.is_empty()
&& let Some(dir) = self.path.as_ref().filter(|path| path.is_dir())
{
partitions = crate::formats::readers::hive::discover_hive_partition_columns(dir)
.into_iter()
.filter(|column| schema.get(column).is_some())
.collect();
}
let mut time_columns: Vec<(String, bool)> = state
.quality_temporal_columns(scope)
.into_iter()
.map(|column| {
let has_time =
!matches!(schema.get(&column), Some(polars::prelude::DataType::Date));
(column, has_time)
})
.collect();
for format in &plan.time_formats {
if !time_columns
.iter()
.any(|(column, _)| *column == format.column)
{
time_columns.push((
format.column.clone(),
format.kind == data_quality::TimeKind::Datetime,
));
}
}
analysis_modal::PlanContext {
partitions,
time_columns,
files: state.quality_source_file_count() > 1,
text_columns: state
.quality_text_columns(scope)
.into_iter()
.map(|column| {
let examples = state.buffered_values(&column, 3);
(column, examples)
})
.collect(),
}
}
pub(crate) fn quality_time_candidates(&self) -> Vec<String> {
let Some(state) = self.data_table_state.as_ref() else {
return Vec::new();
};
let scope = &self.analysis_modal.quality.plan.scope;
let mut columns = state.quality_temporal_columns(scope);
columns.extend(state.quality_text_columns(scope));
columns
}
pub(crate) fn quality_intent_columns(&self) -> Vec<(String, polars::prelude::DataType)> {
self.data_table_state
.as_ref()
.map(|state| {
crate::widgets::quality_intent::intent_columns(
state.quality_schema(&self.analysis_modal.quality.plan.scope),
)
})
.unwrap_or_default()
}
fn quality_source_identity(
&self,
state: &DataTableState,
scope: &data_quality::QualityScope,
) -> crate::analysis::quality_export::SourceIdentity {
let format = self
.source
.original_file_format
.map(|format| format.as_str().to_string())
.or_else(|| {
self.path
.as_ref()
.and_then(|path| path.extension())
.and_then(|extension| extension.to_str())
.map(str::to_string)
});
let mut view = Vec::new();
if !scope.uses_source() {
if !state.get_active_query().is_empty() {
view.push(format!("query: {}", state.get_active_query()));
}
if !state.get_active_sql_query().is_empty() {
view.push(format!("SQL: {}", state.get_active_sql_query()));
}
if !state.get_active_fuzzy_query().is_empty() {
view.push(format!("text: {}", state.get_active_fuzzy_query()));
}
for (index, filter) in state.view_filters().iter().enumerate() {
let join = if index == 0 {
String::new()
} else {
format!("{} ", filter.logical_op.as_str())
};
view.push(format!(
"filter: {join}{} {} {}",
filter.column,
filter.operator.as_str(),
filter.value
));
}
if state.reshape_source().is_some() {
view.push("reshaped: pivot or melt".to_string());
}
}
let remote = state.is_remote_source();
let piped = self.reads_stdin();
let location = self.path.as_ref().map(|path| {
match std::path::absolute(path).ok().filter(|_| !remote && !piped) {
Some(path) => path.display().to_string(),
None => path.display().to_string(),
}
});
crate::analysis::quality_export::SourceIdentity {
location,
remote,
format,
view,
..crate::analysis::quality_export::SourceIdentity::default()
}
.with_files(state.quality_source_file_names())
}
pub(crate) fn quality_page_setup(&self) -> Option<data_quality::QualitySetup> {
let modal = &self.analysis_modal;
data_quality::page_setup(
modal.quality.page,
modal.quality_result_plan(),
modal.quality.results.as_ref(),
self.has_quality_time_columns(),
)
}
pub(crate) fn has_quality_time_columns(&self) -> bool {
!self.analysis_modal.quality.plan.time_formats.is_empty()
|| self.data_table_state.as_ref().is_some_and(|state| {
!state
.quality_temporal_columns(&self.analysis_modal.quality.plan.scope)
.is_empty()
})
}
pub(crate) fn quality_kept_serves(&self, plan: &data_quality::DataQualityPlan) -> bool {
plan.compute == data_quality::QualityCompute::Sample
&& self.kept_quality_sample(&plan.sample()).is_some()
}
pub(crate) fn quality_segment_count(
&self,
plan: &data_quality::DataQualityPlan,
) -> data_quality::SegmentCount {
if plan.compute != data_quality::QualityCompute::Sample {
return data_quality::SegmentCount::NotNeeded;
}
match self.kept_quality_sample(&plan.sample()) {
Some(kept) => kept.segment_count(plan),
None => data_quality::fresh_segment_count(plan, self.quality_may_read_blocks(plan)),
}
}
pub(crate) fn quality_released(&self, plan: &data_quality::DataQualityPlan) -> bool {
let Some(view_generation) = self
.data_table_state
.as_ref()
.map(DataTableState::len_generation)
else {
return false;
};
let sample = plan.sample();
plan.compute == data_quality::QualityCompute::Sample
&& self
.quality
.released
.iter()
.any(|(dataset, view, released)| {
*dataset == self.dataset_generation
&& *view == view_generation
&& *released == sample
})
}
fn quality_one_columnar_file(&self) -> bool {
let Some(state) = self.data_table_state.as_ref() else {
return false;
};
let columnar = matches!(
self.source.original_file_format,
Some(ExportFormat::Parquet | ExportFormat::Ipc)
) || self.path.as_ref().is_some_and(|path| {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| {
matches!(
extension.to_ascii_lowercase().as_str(),
"parquet" | "pq" | "arrow" | "arrows" | "ipc" | "feather"
)
})
});
columnar && state.loaded_file_count() == 1
}
pub(crate) fn quality_reads_blocks(&self, plan: &data_quality::DataQualityPlan) -> bool {
let Some(state) = self.data_table_state.as_ref() else {
return false;
};
self.quality_one_columnar_file()
&& match plan.scope {
data_quality::QualityScope::WholeSource => true,
data_quality::QualityScope::CurrentView => !state.changes_rows(),
data_quality::QualityScope::FirstRows(_)
| data_quality::QualityScope::ViewRows { .. } => {
state.source_file_count() == Some(1)
}
_ => false,
}
}
pub(crate) fn quality_may_read_blocks(&self, plan: &data_quality::DataQualityPlan) -> bool {
let Some(state) = self.data_table_state.as_ref() else {
return false;
};
self.quality_reads_blocks(plan)
|| (self.quality_one_columnar_file()
&& state.may_keep_scan_rows()
&& matches!(
plan.scope,
data_quality::QualityScope::CurrentView
| data_quality::QualityScope::FirstRows(_)
| data_quality::QualityScope::ViewRows { .. }
))
}
pub(crate) fn quality_cached(&self, plan: &data_quality::DataQualityPlan) -> bool {
let Some(view_generation) = self
.data_table_state
.as_ref()
.map(DataTableState::len_generation)
else {
return false;
};
self.quality.cache.iter().any(|entry| {
entry.dataset_generation == self.dataset_generation
&& entry.view_generation == view_generation
&& entry.plan.same_measurement(plan)
})
}
pub(crate) fn quality_kept_rows(&self) -> Option<widgets::data_quality::KeptRows> {
let kept = self
.quality
.samples
.iter()
.filter(|kept| kept.dataset_generation == self.dataset_generation)
.collect::<Vec<_>>();
let copy_bytes = self
.quality
.copies
.iter()
.filter(|kept| kept.dataset_generation == self.dataset_generation)
.map(|kept| kept.copy.bytes())
.sum::<u64>();
(!kept.is_empty() || copy_bytes > 0).then(|| widgets::data_quality::KeptRows {
samples: kept.len(),
rows: kept.iter().map(|kept| kept.rows.df().height()).sum(),
bytes: kept.iter().map(|kept| kept.rows.estimated_bytes()).sum(),
copy_bytes,
})
}
pub(crate) fn release_quality_rows(&mut self) {
let Some(kept) = self.quality_kept_rows() else {
self.flash_note("Nothing kept to release".to_string());
return;
};
for released in std::mem::take(&mut self.quality.samples) {
self.quality.released.retain(|(dataset, view, sample)| {
!(*dataset == released.dataset_generation
&& *view == released.view_generation
&& *sample == released.sample)
});
self.quality.released.insert(
0,
(
released.dataset_generation,
released.view_generation,
released.sample,
),
);
}
self.quality.released.truncate(QUALITY_RELEASED_REMEMBERED);
let generation = self.dataset_generation;
self.quality
.copies
.retain(|kept| kept.dataset_generation != generation);
if kept.copy_bytes > 0 {
self.quality.copy_released = Some(generation);
}
let rows = format!(
"{} kept {} ({})",
numfmt::group_chrome(kept.rows),
if kept.rows == 1 { "row" } else { "rows" },
crate::numfmt::bytes(kept.bytes as u64)
);
let copy = format!("the local copy ({})", crate::numfmt::bytes(kept.copy_bytes));
self.flash_note(match (kept.samples > 0, kept.copy_bytes > 0) {
(true, true) => format!("Released {rows} and {copy}; the next run reads again"),
(false, true) => format!("Released {copy}; the next full scan fetches again"),
_ => format!("Released {rows}; the next run reads again"),
});
}
fn kept_quality_sample(
&self,
sample: &sampling::Sample,
) -> Option<std::sync::Arc<data_quality::QualitySample>> {
self.kept_quality_entry(sample)
.map(|kept| kept.rows.clone())
}
fn kept_quality_entry(&self, sample: &sampling::Sample) -> Option<&KeptQualitySample> {
let view_generation = self.data_table_state.as_ref()?.len_generation();
self.quality.samples.iter().find(|kept| {
kept.dataset_generation == self.dataset_generation
&& kept.view_generation == view_generation
&& &kept.sample == sample
})
}
pub(crate) fn retain_quality_sample(&mut self, kept: &KeptQualitySample) {
if kept.dataset_generation != self.dataset_generation {
return;
}
self.quality.samples.retain(|entry| !entry.same_rows(kept));
self.quality.released.retain(|(dataset, view, sample)| {
!(*dataset == kept.dataset_generation
&& *view == kept.view_generation
&& *sample == kept.sample)
});
self.quality.samples.insert(0, kept.clone());
self.trim_quality_memory();
}
fn trim_quality_memory(&mut self) {
loop {
let used = self
.quality
.cache
.iter()
.map(|entry| entry.bytes)
.sum::<usize>()
+ self
.quality
.samples
.iter()
.map(|kept| kept.rows.estimated_bytes())
.sum::<usize>();
if used <= self.quality.memory_budget {
return;
}
let remakeable = self
.quality
.cache
.iter()
.enumerate()
.skip(1)
.rev()
.find(|(_, entry)| {
entry.plan.compute == data_quality::QualityCompute::Sample
&& self.quality.samples.iter().any(|kept| {
kept.dataset_generation == entry.dataset_generation
&& kept.view_generation == entry.view_generation
&& kept.sample == entry.plan.sample()
})
})
.map(|(index, _)| index);
if let Some(index) = remakeable {
self.quality.cache.remove(index);
} else if self.quality.samples.len() > 1 {
if let Some(released) = self.quality.samples.pop() {
self.quality.released.insert(
0,
(
released.dataset_generation,
released.view_generation,
released.sample,
),
);
self.quality.released.truncate(QUALITY_RELEASED_REMEMBERED);
}
} else if self.quality.cache.len() > 1 {
self.quality.cache.pop();
} else {
return;
}
}
}
pub(crate) fn restore_cached_quality(&mut self) -> bool {
let Some(view_generation) = self
.data_table_state
.as_ref()
.map(DataTableState::len_generation)
else {
return false;
};
let plan = self.analysis_modal.quality.plan.clone();
let Some(cached) = self.quality.cache.iter().find(|entry| {
entry.dataset_generation == self.dataset_generation
&& entry.view_generation == view_generation
&& entry.plan.same_measurement(&plan)
}) else {
return false;
};
let mut results = cached.results.clone();
if cached.plan != plan {
if cached.plan.compares_differently(&plan) {
results.compare_segments(&plan);
}
self.cache_quality_result(&results, plan.clone());
}
self.analysis_modal.quality.results = Some(results);
self.analysis_modal.quality.last_plan = Some(plan);
self.analysis_modal.quality.from_cache = true;
self.analysis_modal
.set_quality_page(data_quality::QualityPage::Overview);
true
}
pub(crate) fn cache_quality_result(
&mut self,
results: &data_quality::DataQualityResults,
plan: data_quality::DataQualityPlan,
) {
let Some(view_generation) = self
.data_table_state
.as_ref()
.map(DataTableState::len_generation)
else {
return;
};
self.quality.cache.retain(|entry| {
!(entry.dataset_generation == self.dataset_generation
&& entry.view_generation == view_generation
&& entry.plan.same_measurement(&plan))
});
self.quality.cache.insert(
0,
QualityCacheEntry {
dataset_generation: self.dataset_generation,
view_generation,
plan,
bytes: results.estimated_bytes(),
results: results.clone(),
},
);
self.trim_quality_memory();
}
pub(crate) fn quality_scope_on_copy(
lf: LazyFrame,
job: QualityCopyJob,
watch: &data_quality::QualityWatch,
fetch: impl FnOnce(
&[crate::cloud::local_copy::RemoteObject],
&Path,
) -> Result<crate::cloud::local_copy::LocalCopy>,
kept: impl FnOnce(Option<Arc<crate::cloud::local_copy::LocalCopy>>),
) -> Result<(LazyFrame, Option<Arc<crate::cloud::local_copy::LocalCopy>>)> {
let (copy, fetched) = match job {
QualityCopyJob::Source => return Ok((lf, None)),
QualityCopyJob::Kept(copy) => (copy, false),
QualityCopyJob::Fetch { objects, root } => {
watch.stage(data_quality::QualityStage::CopyingSource, true, true)?;
let copy = fetch(&objects, &root).map_err(|error| {
if watch.cancelled() {
color_eyre::eyre::eyre!(crate::analysis::sampling::CANCELLED)
} else {
error
}
})?;
(Arc::new(copy), true)
}
};
let local = copy.redirect(&lf).filter(|local| {
let schemas = (local.clone().collect_schema(), lf.clone().collect_schema());
matches!(schemas, (Ok(local), Ok(source)) if local == source)
});
let Some(local) = local else {
log::warn!(target: "datui", "local copy does not read as the source; reading the source");
kept(None);
return Ok((lf, None));
};
if fetched {
kept(Some(copy.clone()));
}
watch.use_copy(data_quality::CopyRead {
bytes: copy.bytes(),
objects: copy.objects(),
fetched,
});
Ok((local, Some(copy)))
}
#[cfg(feature = "cloud")]
fn fetch_quality_copy(
objects: &[crate::cloud::local_copy::RemoteObject],
root: &Path,
cloud: &crate::config::CloudConfig,
runtime: &tokio::runtime::Handle,
stop: &crate::analysis::sampling::ReadWatch,
) -> Result<crate::cloud::local_copy::LocalCopy> {
use crate::cloud::download::StreamError;
use object_store::ObjectStoreExt;
crate::cloud::local_copy::LocalCopy::fetch(root, objects, stop, |object, write| {
let url = object.url.as_str();
let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)
.map_err(StreamError::Write)?;
let (_, key) = Self::cloud_bucket_and_key(url).map_err(StreamError::Write)?;
let path = crate::cloud::cloud_browse::object_path(&key);
let listed = object.etag.clone();
let gone = crate::error_display::gone_since_opened_message(url);
let open = async move {
let got = store.get(&path).await.map_err(|e| match e {
object_store::Error::NotFound { .. } => gone.clone(),
e => e.to_string(),
})?;
if let (Some(listed), Some(fetched)) = (&listed, &got.meta.e_tag)
&& !crate::cloud::local_copy::same_etag(listed, fetched)
{
return Err(gone);
}
Ok((got.into_stream(), None))
};
let watch = stop.clone();
crate::cloud::download::stream_into(runtime, open, move || watch.stopped(), write)
})
}
fn quality_copies_root(&self) -> PathBuf {
self.cache
.cache_dir()
.join(crate::cloud::local_copy::COPIES_DIR)
}
fn quality_copy_limit(&self) -> u64 {
self.app_config.analysis.quality_local_copy.bytes()
}
pub fn quality_copy_bytes(&self) -> u64 {
self.quality
.copies
.iter()
.map(|kept| kept.copy.bytes())
.sum()
}
fn quality_copy_kept(&self) -> Option<&Arc<crate::cloud::local_copy::LocalCopy>> {
let state = self.data_table_state.as_ref()?;
self.quality
.copies
.iter()
.find(|kept| {
kept.dataset_generation == self.dataset_generation
&& state.each_remote_object().is_some_and(|mut objects| {
objects.all(|object| {
object.is_some_and(|object| kept.copy.covers(&object.url))
})
})
})
.map(|kept| &kept.copy)
}
fn quality_copy_free_space(&self) -> Option<u64> {
let root = self.quality_copies_root();
let Ok(mut cached) = self.quality.copy_free.lock() else {
return crate::cloud::local_copy::free_space(&root);
};
match *cached {
Some((asked, free)) if asked.elapsed() < std::time::Duration::from_secs(5) => free,
_ => {
let free = crate::cloud::local_copy::free_space(&root);
*cached = Some((std::time::Instant::now(), free));
free
}
}
}
pub(crate) fn quality_copy_plan(
&self,
plan: &data_quality::DataQualityPlan,
) -> data_quality::CopyPlan {
use data_quality::{CopyPlan, NoCopy};
let Some(state) = self.data_table_state.as_ref() else {
return CopyPlan::NotApplicable;
};
if plan.compute != data_quality::QualityCompute::Full || !state.is_remote_source() {
return CopyPlan::NotApplicable;
}
if !state.quality_reads_whole_source(&plan.scope) {
return CopyPlan::Passes(NoCopy::PartOfTheSource);
}
if let Some(copy) = self.quality_copy_kept() {
return CopyPlan::Kept {
bytes: copy.bytes(),
objects: copy.objects(),
};
}
let limit = self.quality_copy_limit();
if limit == 0 {
return CopyPlan::Passes(NoCopy::Off);
}
if self.quality.copy_unusable == Some(self.dataset_generation) {
return CopyPlan::Passes(NoCopy::Unusable);
}
let Some((bytes, objects)) = state.remote_objects_size() else {
return CopyPlan::Passes(NoCopy::SizeUnknown);
};
if bytes > limit {
return CopyPlan::Passes(NoCopy::TooLarge { bytes, limit });
}
let free = self.quality_copy_free_space();
if free.is_none_or(|free| bytes > free) {
return CopyPlan::Passes(NoCopy::NoRoom { bytes, free });
}
CopyPlan::Fetch { bytes, objects }
}
pub(crate) fn quality_copy_released(&self) -> bool {
self.quality.copy_released == Some(self.dataset_generation)
}
pub(crate) fn retain_quality_copy(
&mut self,
dataset_generation: u64,
copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>,
) {
if dataset_generation != self.dataset_generation {
return;
}
let Some(copy) = copy else {
self.quality
.copies
.retain(|kept| kept.dataset_generation != dataset_generation);
self.quality.copy_unusable = Some(dataset_generation);
return;
};
self.quality.copies.insert(
0,
RetainedCopy {
dataset_generation,
copy,
},
);
self.quality.copy_released = None;
let limit = self.quality_copy_limit();
while self.quality.copies.len() > 1 && self.quality_copy_bytes() > limit {
self.quality.copies.pop();
}
}
pub(crate) fn read_sample_rows(
&mut self,
sample: sampling::Sample,
evidence: Option<(quality_report::EvidenceRows, String)>,
) -> Option<AppEvent> {
let state = self.data_table_state.as_ref()?;
let (source, known_total) = Self::sample_source_for(state, &sample.scope);
let streaming = self.app_config.performance.streaming;
let kept = self.kept_quality_sample(&sample).map(|kept| {
let columns: Vec<_> = state
.schema()
.iter_names()
.filter(|name| kept.df().column(name.as_str()).is_ok())
.map(|name| polars::prelude::col(name.clone()))
.collect();
(kept, columns)
});
if kept.is_none() && self.read_waits_for_cancelled() {
return None;
}
self.analysis_modal.computing = Some(AnalysisProgress::new(if evidence.is_some() {
"Reading the matching sampled rows"
} else {
"Reading the sample"
}));
self.spawn_job(Job::SampleRows, Some("Reading the sample..."), move |_| {
let (rows, columns) = match kept {
Some((kept, columns)) => (Ok(kept.analysis_rows(kept.df().clone())), Some(columns)),
None => (
source
.cut(&sample.scope)
.and_then(|lf| sampling::read(&lf, &sample, known_total, streaming)),
None,
),
};
let shown = |df: polars::prelude::DataFrame| match &columns {
Some(columns) => polars::prelude::IntoLazy::lazy(df)
.select(columns.clone())
.collect()
.map_err(color_eyre::eyre::Report::from),
None => Ok(df),
};
let read = rows.and_then(|rows| {
let label = format!(
"Sample {} {}",
crate::glyphs::get().middot,
sample.outcome(
rows.total_rows,
rows.sample_size,
rows.per_value.as_ref().map(|per_value| per_value.kept),
)
);
match evidence {
Some((quality_report::EvidenceRows::Duplicates, label)) => {
let keys = rows
.df
.get_column_names()
.into_iter()
.filter(|name| !name.starts_with("__datui"))
.cloned()
.collect::<Vec<_>>();
let df = data_quality::duplicate_rows(
polars::prelude::IntoLazy::lazy(rows.df),
&keys,
streaming,
)?;
Ok((shown(df)?, label))
}
Some((quality_report::EvidenceRows::Matching(predicate), label)) => {
let df = polars::prelude::IntoLazy::lazy(rows.df)
.filter(predicate)
.collect()?;
Ok((shown(df)?, label))
}
Some((quality_report::EvidenceRows::Files(_), label)) => {
Ok((shown(rows.df)?, label))
}
None => Ok((shown(rows.df)?, label)),
}
});
let (df, label) = read.map_err(|error| format!("{error}"))?;
Ok(Answer::Sample { df, label })
});
None
}
pub(crate) fn sync_quality_plan(&mut self) {
let sample = self.analysis_modal.sample.clone();
self.analysis_modal.quality.plan.adopt_sample(&sample);
}
pub(crate) fn open_quality_setup(&mut self) {
use data_quality::QualityPage;
let modal = &mut self.analysis_modal;
if !modal.quality.page.is_setup() {
modal.quality.setup_return = modal.quality.page.tab();
}
if modal.quality.setup_before.is_none() {
modal.quality.setup_before = Some(modal.quality.plan.clone());
}
if modal.quality.page != QualityPage::Setup {
modal.set_quality_page(QualityPage::Setup);
modal.quality.plan_field = 0;
}
modal.focus = analysis_modal::AnalysisFocus::Main;
}
pub(crate) fn leave_quality_setup(&mut self) {
use data_quality::QualityPage;
let modal = &mut self.analysis_modal;
if let Some(before) = modal.quality.setup_before.take() {
modal.quality.plan = before;
}
modal.quality.setup_note = None;
modal.quality.picker = None;
if modal.quality.results.is_some() {
let back = match modal.quality.setup_return {
page if page.is_setup() => QualityPage::Overview,
page => page,
};
modal.set_quality_page(back);
} else {
modal.set_quality_page(QualityPage::Setup);
modal.focus = analysis_modal::AnalysisFocus::Sidebar;
}
}
fn quality_full_scan_question(&self, plan: &data_quality::DataQualityPlan) -> String {
let mut lines = vec![
"Run a full scan?".to_string(),
String::new(),
"Reads: every eligible row, up to the whole source".to_string(),
];
if let data_quality::CopyPlan::Fetch { bytes, .. } = self.quality_copy_plan(plan) {
lines.push(format!(
"Fetch: {} once, to a local copy",
crate::numfmt::bytes(bytes)
));
}
lines.push("Source writes: none".to_string());
lines.join("\n")
}
fn quality_setup_problem(&self) -> Option<String> {
let plan = &self.analysis_modal.quality.plan;
let schema = self.data_table_state.as_ref()?.quality_schema(&plan.scope);
match &plan.grain {
data_quality::QualityGrain::TimeWindows { column, .. }
if plan.compute != data_quality::QualityCompute::Metadata
&& !plan.reads_as_time(column, schema) =>
{
Some(format!(
"{column}: text, no format {} set Text as time",
crate::glyphs::get().middot
))
}
_ => None,
}
}
pub(crate) fn run_quality_setup(&mut self, confirmed: bool) -> Option<AppEvent> {
use data_quality::QualityPage;
if self.cancelled_analysis_running().is_some() {
self.analysis_modal.quality.setup_note = Some(QUALITY_RUN_WAITS.to_string());
return None;
}
if let Some(problem) = self.quality_setup_problem() {
self.analysis_modal.quality.setup_note = Some(problem);
return None;
}
let plan = &self.analysis_modal.quality.plan;
let here = (self.analysis_modal.quality.results.is_some()
&& self
.analysis_modal
.quality
.last_plan
.as_ref()
.is_some_and(|last| last.same_measurement(plan)))
|| self.quality_cached(plan);
if plan.requires_confirmation() && !here && !confirmed {
let message = self.quality_full_scan_question(plan);
self.confirmation_modal
.show(message, Confirm::QualityFullScan);
self.confirmation_modal.yes_label = "Run";
return None;
}
self.commit_quality_plan();
let modal = &mut self.analysis_modal;
if modal.quality.results.is_some()
&& modal.quality.last_plan.as_ref() == Some(&modal.quality.plan)
{
let back = match modal.quality.setup_return {
page if page.is_setup() => QualityPage::Overview,
page => page,
};
modal.set_quality_page(back);
return None;
}
if let (Some(results), Some(last)) = (
modal.quality.results.as_ref(),
modal.quality.last_plan.as_ref(),
) && last.same_measurement(&modal.quality.plan)
{
let mut results = results.clone();
let plan = modal.quality.plan.clone();
let page = if last.compares_differently(&plan) {
results.compare_segments(&plan);
QualityPage::Segments
} else {
QualityPage::Trends
};
modal.quality.results = Some(results.clone());
modal.quality.last_plan = Some(plan.clone());
modal.set_quality_page(page);
self.cache_quality_result(&results, plan);
return None;
}
if self.restore_cached_quality() {
return None;
}
self.analysis_modal.quality.from_cache = false;
let mut progress = AnalysisProgress::new("Preparing the plan");
if self.quality_kept_serves(&self.analysis_modal.quality.plan) {
progress.reuse = Some("Starts from rows a run already read".to_string());
}
self.analysis_modal.computing = Some(progress);
self.busy = true;
Some(AppEvent::AnalysisCompute(
analysis_modal::AnalysisTool::DataQuality,
))
}
fn commit_quality_plan(&mut self) {
let modal = &mut self.analysis_modal;
let sample = modal.quality.plan.sample();
if sample != modal.sample {
modal.describe_results = None;
modal.distribution_results = None;
modal.correlation_results = None;
}
modal.sample = sample;
modal.sample_dataset = Some(self.dataset_generation);
modal.sample_run_for = Some(self.dataset_generation);
modal.quality.setup_before = None;
modal.quality.setup_note = None;
modal.quality.picker = None;
}
pub(crate) fn run_quality_compute(&mut self) -> Option<AppEvent> {
if let Some(state) = &self.data_table_state {
let plan = self.analysis_modal.quality.plan.clone();
let source_scope = plan.scope.uses_source();
let (lf, source, cached_rows) = if source_scope {
let (lf, source) = state.data_quality_source_scan();
(lf, source, None)
} else {
let ordered = matches!(
plan.scope,
data_quality::QualityScope::FirstRows(_)
| data_quality::QualityScope::ViewRows { .. }
);
let (lf, source) = state.data_quality_scan(ordered);
let rows = state.num_rows_if_valid().map(|rows| match &plan.scope {
data_quality::QualityScope::CurrentView => rows,
data_quality::QualityScope::FirstRows(limit) => rows.min(*limit),
data_quality::QualityScope::ViewRows { start, end } => {
rows.min(*end).saturating_sub(start.saturating_sub(1))
}
_ => unreachable!(),
});
(lf, source, rows)
};
let streaming = state.polars_streaming();
let audio = (plan.compute == data_quality::QualityCompute::Full)
.then(|| state.window_for_quality(&plan.scope))
.flatten()
.and_then(crate::formats::audio::recording);
let view_generation = state.len_generation();
let dataset_generation = self.dataset_generation;
let kept_entry = self.kept_quality_entry(&plan.sample());
let kept = kept_entry.map(|kept| kept.rows.clone());
let kept_source = kept_entry
.filter(|_| plan.compute == data_quality::QualityCompute::Sample)
.map(|kept| kept.source.clone());
let mut identity = self.quality_source_identity(state, &plan.scope);
let copy_job = match self.quality_copy_plan(&plan) {
data_quality::CopyPlan::Kept { .. } => self
.quality_copy_kept()
.cloned()
.map_or(QualityCopyJob::Source, QualityCopyJob::Kept),
data_quality::CopyPlan::Fetch { .. } => match state.remote_objects() {
Some(objects) => QualityCopyJob::Fetch {
objects,
root: self.quality_copies_root(),
},
None => QualityCopyJob::Source,
},
_ => QualityCopyJob::Source,
};
#[cfg(feature = "cloud")]
let (cloud, runtime) = (self.app_config.cloud.clone(), self.runtime.clone());
let mut source = source;
if plan.compute == data_quality::QualityCompute::Full
&& let Some(source) = source.as_mut()
{
source.conflict_scan = state.quality_conflict_scan();
}
let started = self.start_job(
Job::Analysis(jobs::AnalysisRun::default()),
Some("Profiling data quality..."),
);
let ticket = started.ticket();
let phases = self.events.clone();
let watch = data_quality::QualityWatch::new(move |phase| {
let _ = phases.send(AppEvent::JobProgress {
ticket,
progress: Progress::QualityPhase(phase),
});
});
if let Some(progress) = self.analysis_modal.computing.as_mut() {
progress.read = Some(watch.read().clone());
}
if let Some(Job::Analysis(run)) = self.jobs.job_mut(ticket) {
run.watch = Some(watch.clone());
}
started.run(&self.runtime, move |worker| {
match kept_source {
Some(source) => identity = source,
None => identity.stat(),
}
let lf = if source_scope {
data_quality::prepare_source_quality_scan(lf, source.as_ref())
.map_err(|error| format!("{error}"))?
} else {
lf
};
let lf = data_quality::apply_quality_scope(lf, &plan.scope, source.as_ref())
.map_err(|error| format!("{error}"))?;
let fetch = |objects: &[crate::cloud::local_copy::RemoteObject], root: &Path| {
#[cfg(feature = "cloud")]
{
Self::fetch_quality_copy(objects, root, &cloud, &runtime, watch.read())
}
#[cfg(not(feature = "cloud"))]
{
let _ = (objects, root);
Err(color_eyre::eyre::eyre!("Built without cloud support"))
}
};
let kept_copy = |copy: Option<Arc<crate::cloud::local_copy::LocalCopy>>| {
worker.send(AppEvent::BackgroundQualityCopyKept {
dataset_generation,
copy,
});
};
let (lf, held) =
Self::quality_scope_on_copy(lf, copy_job, &watch, fetch, kept_copy)
.map_err(|error| format!("{error}"))?;
let (results, rows) = crate::analysis::data_quality::compute_data_quality_watched(
&lf,
cached_rows,
&plan,
source.as_ref(),
streaming,
kept.as_deref(),
&watch,
);
let results = match (results, audio) {
(Ok(mut results), Some(audio)) => {
crate::analysis::data_quality::add_signal_observations(
&mut results,
&audio,
&watch,
)
.map(|()| results)
}
(results, _) => results,
};
drop(held);
let kept = rows.map(|rows| KeptQualitySample {
dataset_generation,
view_generation,
sample: plan.sample(),
rows: std::sync::Arc::new(rows),
source: identity.clone(),
});
match results {
Ok(mut results) => {
results.source = Some(Box::new(identity));
Ok(Answer::DataQuality {
results: Box::new(results),
kept,
plan: Box::new(plan),
})
}
Err(error) => {
if let Some(kept) = kept {
worker.send(AppEvent::BackgroundQualitySampleKept { kept });
}
Err(format!("{error}"))
}
}
});
} else {
self.analysis_modal.computing = None;
self.busy = false;
}
None
}
}