Skip to main content

datui_lib/analysis/
data_quality.rs

1use crate::analysis::statistics::collect_lazy;
2use color_eyre::Result;
3use color_eyre::eyre::Report;
4use polars::chunked_array::cast::CastOptions;
5use polars::prelude::*;
6use std::collections::BTreeMap;
7use std::sync::Arc;
8
9// A dataset-grain sample is spread across the whole scope (see `sampling::analysis_rows`).
10const DEFAULT_SAMPLE_ROWS: usize = 10_000;
11const DEFAULT_CHUNK_ROWS: usize = 1_000_000;
12const QUALITY_WINDOW_START: &str = "__datui_quality_window_start";
13pub const QUALITY_SOURCE_FILE_COLUMN: &str = "__datui_quality_source_file";
14/// How nearly unique a column must be before its repeats are worth naming: an id
15/// repeating twice in a million rows is a finding, a category repeating is not.
16pub const KEY_LIKE_UNIQUENESS: f64 = 0.95;
17/// Files named per drift observation, and values read from each of them. Both are the
18/// evidence, not the measurement: the counts above them cover every file.
19const MAX_EVIDENCE_FILES: usize = 20;
20const MAX_CONFLICT_EXAMPLES: usize = 5;
21/// Values kept per finding from the rows a run read, and groups of duplicate rows:
22/// enough to recognize the problem in the detail, which opens the rest.
23pub const MAX_FINDING_EXAMPLES: usize = 3;
24/// Window widths offered for time-window grain, in the order the plan cycles them.
25pub const QUALITY_WINDOW_WIDTHS: [&str; 4] = ["1h", "1d", "1w", "1mo"];
26
27#[derive(Debug, Clone, PartialEq, Eq, Default)]
28pub enum QualityScope {
29    #[default]
30    CurrentView,
31    WholeSource,
32    FirstRows(usize),
33    ViewRows {
34        start: usize,
35        end: usize,
36    },
37    SourceFiles(Vec<usize>),
38    SourcePartition {
39        column: String,
40        value: String,
41    },
42    SourceTimeRange {
43        column: String,
44        start: String,
45        end: String,
46    },
47}
48
49impl QualityScope {
50    pub fn label(&self) -> String {
51        match self {
52            Self::CurrentView => "current view".to_string(),
53            Self::WholeSource => "whole source".to_string(),
54            Self::FirstRows(rows) => format!(
55                "first {} rows of the view",
56                crate::numfmt::group_chrome(*rows)
57            ),
58            Self::ViewRows { start, end } => format!(
59                "view rows {}-{}",
60                crate::numfmt::group_chrome(*start),
61                crate::numfmt::group_chrome(*end)
62            ),
63            Self::SourceFiles(indices) => format!(
64                "source files {}",
65                indices
66                    .iter()
67                    .map(usize::to_string)
68                    .collect::<Vec<_>>()
69                    .join(",")
70            ),
71            Self::SourcePartition { column, value } => format!("source {column}={value}"),
72            Self::SourceTimeRange { column, start, end } => {
73                format!("source {column} {start}..{end}")
74            }
75        }
76    }
77
78    pub fn uses_source(&self) -> bool {
79        matches!(
80            self,
81            Self::WholeSource
82                | Self::SourceFiles(_)
83                | Self::SourcePartition { .. }
84                | Self::SourceTimeRange { .. }
85        )
86    }
87
88    pub fn command(&self) -> String {
89        match self {
90            Self::CurrentView => "view".to_string(),
91            Self::WholeSource => "source".to_string(),
92            Self::FirstRows(rows) => format!("rows 1..{rows}"),
93            Self::ViewRows { start, end } => format!("rows {start}..{end}"),
94            Self::SourceFiles(indices) => format!(
95                "files {}",
96                indices
97                    .iter()
98                    .map(usize::to_string)
99                    .collect::<Vec<_>>()
100                    .join(",")
101            ),
102            Self::SourcePartition { column, value } => format!("partition {column}={value}"),
103            Self::SourceTimeRange { column, start, end } => format!("time {column}={start}..{end}"),
104        }
105    }
106
107    pub fn parse_command(text: &str) -> Result<Self> {
108        let value = text.trim();
109        if value == "view" {
110            return Ok(Self::CurrentView);
111        }
112        if value == "source" {
113            return Ok(Self::WholeSource);
114        }
115        if let Some(range) = value.strip_prefix("rows ") {
116            let (start, end) = range
117                .split_once("..")
118                .ok_or_else(|| color_eyre::eyre::eyre!("use rows START..END"))?;
119            let start = start.parse::<usize>()?;
120            let end = end.parse::<usize>()?;
121            if start == 0 || end < start {
122                return Err(color_eyre::eyre::eyre!(
123                    "row range must be 1-based with END >= START"
124                ));
125            }
126            if start == 1 {
127                // The same rows as FirstRows(end); use the one spelling so the
128                // scope round-trips through the editor and keeps its cache entry.
129                return Ok(Self::FirstRows(end));
130            }
131            return Ok(Self::ViewRows { start, end });
132        }
133        if let Some(files) = value.strip_prefix("files ") {
134            let indices = files
135                .split(',')
136                .map(|part| part.trim().parse::<usize>())
137                .collect::<std::result::Result<Vec<_>, _>>()?;
138            if indices.is_empty() || indices.contains(&0) {
139                return Err(color_eyre::eyre::eyre!(
140                    "use 1-based file numbers, for example files 1,3"
141                ));
142            }
143            let mut indices = indices;
144            indices.sort_unstable();
145            indices.dedup();
146            return Ok(Self::SourceFiles(indices));
147        }
148        if let Some(partition) = value.strip_prefix("partition ") {
149            let (column, value) = partition
150                .split_once('=')
151                .ok_or_else(|| color_eyre::eyre::eyre!("use partition COLUMN=VALUE"))?;
152            if column.trim().is_empty() || value.trim().is_empty() {
153                return Err(color_eyre::eyre::eyre!(
154                    "partition column and value are required"
155                ));
156            }
157            return Ok(Self::SourcePartition {
158                column: column.trim().to_string(),
159                value: value.trim().to_string(),
160            });
161        }
162        if let Some(time) = value.strip_prefix("time ") {
163            let (column, range) = time
164                .split_once('=')
165                .ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
166            let (start, end) = range
167                .split_once("..")
168                .ok_or_else(|| color_eyre::eyre::eyre!("use time COLUMN=START..END"))?;
169            let (start, end) = (start.trim(), end.trim());
170            if column.trim().is_empty()
171                || parse_scope_time(start).is_none()
172                || parse_scope_time(end).is_none()
173                || parse_scope_time(end) <= parse_scope_time(start)
174            {
175                return Err(color_eyre::eyre::eyre!(
176                    "time range needs a column and increasing ISO dates or UTC timestamps"
177                ));
178            }
179            return Ok(Self::SourceTimeRange {
180                column: column.trim().to_string(),
181                start: start.to_string(),
182                end: end.to_string(),
183            });
184        }
185        Err(color_eyre::eyre::eyre!(
186            "use view, source, rows, files, partition, or time"
187        ))
188    }
189}
190
191pub(crate) fn parse_scope_time(text: &str) -> Option<i64> {
192    chrono::DateTime::parse_from_rfc3339(text)
193        .ok()
194        .map(|value| value.timestamp_micros())
195        .or_else(|| {
196            chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d")
197                .ok()
198                .and_then(|value| value.and_hms_opt(0, 0, 0))
199                .map(|value| value.and_utc().timestamp_micros())
200        })
201}
202
203#[derive(Debug, Clone, Default)]
204pub struct QualitySourceContext {
205    pub file_names: Vec<String>,
206    pub file_starts: Vec<usize>,
207    pub row_index_column: String,
208    /// Per file, in the order of `file_names`, its group in `drift_groups`. Empty when
209    /// the dataset's files all agree with its schema, which is nearly all of them.
210    pub file_group: Vec<u32>,
211    /// The distinct ways this dataset's files differ from its schema, as the footers
212    /// found them. Group 0 is always "nothing missing".
213    pub drift_groups: Arc<Vec<crate::formats::schema_union::DriftGroup>>,
214    /// Per file, the type it holds each unreadable column in (empty when all fit); kept
215    /// beside the groups as the only way back to values a conflict hides.
216    pub file_omitted: Vec<Vec<(PlSmallStr, DataType)>>,
217    /// Rows in the whole loaded source, which is what closes the last file's range.
218    pub dataset_rows: usize,
219    /// How many footers were read: below the file count on a huge dataset, where an
220    /// unread file looks like one missing nothing, so file counts are floors.
221    pub footers_read: usize,
222    /// How to read a column at a file's own type, for values a type conflict hides.
223    /// `None` when files agree, or the run's budget did not promise the reads.
224    pub conflict_scan: Option<QualityConflictScan>,
225}
226
227impl QualitySourceContext {
228    /// What the file at `file` is missing. Group 0 for a file that agrees with the
229    /// dataset's schema, and for a dataset whose files were never grouped.
230    fn group_of_file(&self, file: usize) -> Option<&crate::formats::schema_union::DriftGroup> {
231        let group = *self.file_group.get(file)? as usize;
232        self.drift_groups.get(group)
233    }
234
235    /// The files this dataset is missing something from, by 1-based inventory number,
236    /// paired with what each is missing. Only files that differ have an entry.
237    fn drifting_files(
238        &self,
239    ) -> impl Iterator<Item = (usize, &crate::formats::schema_union::DriftGroup)> {
240        (0..self.file_names.len()).filter_map(move |file| {
241            let group = self.group_of_file(file)?;
242            (!group.is_empty()).then_some((file, group))
243        })
244    }
245
246    /// Rows the file at `file` holds, from its footer.
247    fn file_rows(&self, file: usize) -> usize {
248        let Some(start) = self.file_starts.get(file) else {
249            return 0;
250        };
251        self.file_starts
252            .get(file + 1)
253            .copied()
254            .unwrap_or(self.dataset_rows)
255            .saturating_sub(*start)
256    }
257
258    /// The type the file at `file` holds `column` in, when that is not the type the
259    /// scan reads it as.
260    fn stored_type(&self, file: usize, column: &str) -> Option<&DataType> {
261        self.file_omitted
262            .get(file)?
263            .iter()
264            .find(|(name, _)| name.as_str() == column)
265            .map(|(_, dtype)| dtype)
266    }
267}
268
269/// Prepare the loaded source in the worker, keeping only a provenance index
270/// and replacing binary payloads before any value collection.
271pub fn prepare_source_quality_scan(
272    lf: LazyFrame,
273    source: Option<&QualitySourceContext>,
274) -> Result<LazyFrame> {
275    let schema = lf.clone().collect_schema()?;
276    let expressions = schema
277        .iter()
278        .filter_map(|(name, dtype)| {
279            let column = name.as_str();
280            if column == crate::formats::schema_union::DRIFT_COLUMN
281                && !source.is_some_and(|context| context.row_index_column == column)
282            {
283                return None;
284            }
285            Some(if matches!(dtype, DataType::Binary) {
286                lit(crate::table::binary_stub()).alias(column)
287            } else {
288                col(column)
289            })
290        })
291        .collect::<Vec<_>>();
292    let lf = lf.select(expressions);
293    Ok(
294        if source.is_some_and(|context| context.row_index_column == "__datui_quality_row") {
295            lf.with_row_index("__datui_quality_row", None)
296        } else {
297            lf
298        },
299    )
300}
301
302/// The rows of one partition value, a list (`2019,2021`), or an inclusive range
303/// (`2020..2022`), each read as the column's type (so a file's statistics can answer
304/// the predicate); `∅` is the null partition. A value that does not read as the type
305/// is an error.
306fn partition_predicate(column: &str, value: &str, schema: &Schema) -> Result<Expr> {
307    let dtype = schema
308        .get(column)
309        .ok_or_else(|| color_eyre::eyre::eyre!("partition column {column:?} is unavailable"))?;
310    let read = |text: &str| {
311        crate::typed_value::parse(text, dtype)
312            .map(lit)
313            .map_err(|why| color_eyre::eyre::eyre!("{column}: {why}"))
314    };
315    if let Some((start, end)) = value.split_once("..") {
316        let (start, end) = (start.trim(), end.trim());
317        if start.is_empty() || end.is_empty() {
318            return Err(color_eyre::eyre::eyre!(
319                "a partition range needs both ends, for example year=2020..2022"
320            ));
321        }
322        return Ok(col(column)
323            .gt_eq(read(start)?)
324            .and(col(column).lt_eq(read(end)?)));
325    }
326    value
327        .split(',')
328        .map(str::trim)
329        .filter(|value| !value.is_empty())
330        .map(|value| {
331            Ok(if value == "∅" {
332                col(column).is_null()
333            } else {
334                col(column).eq(read(value)?)
335            })
336        })
337        .reduce(|all, one| Ok(all?.or(one?)))
338        .ok_or_else(|| color_eyre::eyre::eyre!("name at least one partition value"))?
339}
340
341pub fn apply_quality_scope(
342    lf: LazyFrame,
343    scope: &QualityScope,
344    source: Option<&QualitySourceContext>,
345) -> Result<LazyFrame> {
346    match scope {
347        QualityScope::CurrentView | QualityScope::WholeSource => Ok(lf),
348        QualityScope::FirstRows(rows) => Ok(lf.slice(0, (*rows).min(u32::MAX as usize) as u32)),
349        QualityScope::ViewRows { start, end } => {
350            if *start == 0 || end < start {
351                return Err(color_eyre::eyre::eyre!("invalid 1-based view row range"));
352            }
353            let offset = i64::try_from(start - 1)?;
354            let length = end
355                .saturating_sub(*start)
356                .saturating_add(1)
357                .min(u32::MAX as usize) as u32;
358            Ok(lf.slice(offset, length))
359        }
360        QualityScope::SourceFiles(indices) => {
361            let source = source
362                .ok_or_else(|| color_eyre::eyre::eyre!("source-file positions are unavailable"))?;
363            let mut predicate: Option<Expr> = None;
364            for index in indices {
365                let file = index
366                    .checked_sub(1)
367                    .ok_or_else(|| color_eyre::eyre::eyre!("source file numbers start at 1"))?;
368                let start = *source.file_starts.get(file).ok_or_else(|| {
369                    color_eyre::eyre::eyre!("source file #{index} is unavailable")
370                })?;
371                let start = u32::try_from(start)?;
372                let mut range = col(&source.row_index_column).gt_eq(lit(start));
373                if let Some(end) = source.file_starts.get(*index) {
374                    range = range.and(col(&source.row_index_column).lt(lit(u32::try_from(*end)?)));
375                }
376                predicate = Some(match predicate {
377                    Some(previous) => previous.or(range),
378                    None => range,
379                });
380            }
381            Ok(lf.filter(
382                predicate.ok_or_else(|| color_eyre::eyre::eyre!("select at least one file"))?,
383            ))
384        }
385        QualityScope::SourcePartition { column, value } => {
386            let schema = lf.clone().collect_schema()?;
387            if !schema.contains(column.as_str()) {
388                return Err(color_eyre::eyre::eyre!(
389                    "partition column {column:?} is unavailable"
390                ));
391            }
392            Ok(lf.filter(partition_predicate(column, value, &schema)?))
393        }
394        QualityScope::SourceTimeRange { column, start, end } => {
395            let schema = lf.clone().collect_schema()?;
396            let dtype = schema
397                .get(column.as_str())
398                .ok_or_else(|| color_eyre::eyre::eyre!("time column {column:?} is unavailable"))?;
399            if !matches!(dtype, DataType::Date | DataType::Datetime(..)) {
400                return Err(color_eyre::eyre::eyre!(
401                    "{column:?} is not a date or datetime column"
402                ));
403            }
404            let start = parse_scope_time(start)
405                .ok_or_else(|| color_eyre::eyre::eyre!("invalid start time"))?;
406            let end =
407                parse_scope_time(end).ok_or_else(|| color_eyre::eyre::eyre!("invalid end time"))?;
408            if end <= start {
409                return Err(color_eyre::eyre::eyre!("time end must be after start"));
410            }
411            let value = col(column).cast(DataType::Datetime(TimeUnit::Microseconds, None));
412            Ok(lf.filter(
413                value
414                    .clone()
415                    .gt_eq(lit(start).cast(DataType::Datetime(TimeUnit::Microseconds, None)))
416                    .and(value.lt(lit(end).cast(DataType::Datetime(TimeUnit::Microseconds, None)))),
417            ))
418        }
419    }
420}
421
422#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
423pub enum QualityPage {
424    /// Everything a run needs, staged until Enter runs it. Not a tab: a report's
425    /// pages are tabs, and Setup is where a report comes from.
426    #[default]
427    Setup,
428    Overview,
429    Columns,
430    Segments,
431    Trends,
432    Detail,
433    /// One segment's columns beside the segment it is compared with.
434    SegmentDetail,
435    /// Each interval in each segment: the time between two roles.
436    Intervals,
437    /// One interval in one segment, every count it took and out of what.
438    IntervalDetail,
439    TimeRoles,
440    /// Which starts and ends the intervals are, chosen from the assigned roles.
441    IntervalPairs,
442    /// One bar of one Trends line: the segments it pools, what they hold and how
443    /// much of them was read.
444    TrendDetail,
445    /// The expected windows with no rows to show, by why.
446    Gaps,
447    /// Which time windows rows are expected in, edited from Setup.
448    ExpectedWindows,
449    /// What each column must hold, declared: the key and each column's rules.
450    Intent,
451}
452
453impl QualityPage {
454    /// The report's tabs, in the order ←→ walk them. A column's detail sits under
455    /// Columns, a segment's under Segments and an interval's under Intervals.
456    pub const TABS: [Self; 5] = [
457        Self::Overview,
458        Self::Columns,
459        Self::Segments,
460        Self::Trends,
461        Self::Intervals,
462    ];
463
464    pub fn tab(self) -> Self {
465        match self {
466            Self::Detail => Self::Columns,
467            Self::SegmentDetail => Self::Segments,
468            Self::IntervalDetail => Self::Intervals,
469            Self::TrendDetail | Self::Gaps => Self::Trends,
470            Self::TimeRoles | Self::IntervalPairs | Self::ExpectedWindows | Self::Intent => {
471                Self::Setup
472            }
473            page => page,
474        }
475    }
476
477    /// Setup and its editors, which stage a run rather than show one.
478    pub fn is_setup(self) -> bool {
479        self.tab() == Self::Setup
480    }
481
482    pub fn title(self) -> &'static str {
483        match self.tab() {
484            Self::Overview => "Overview",
485            Self::Columns => "Columns",
486            Self::Segments => "Segments",
487            Self::Trends => "Trends",
488            Self::Intervals => "Intervals",
489            _ => "Setup",
490        }
491    }
492}
493
494/// What an empty page is missing, which Enter opens in Setup.
495#[derive(Debug, Clone, Copy, PartialEq, Eq)]
496pub enum QualitySetup {
497    Grain,
498    TimeRoles,
499    Intervals,
500}
501
502impl QualitySetup {
503    pub fn label(self) -> &'static str {
504        match self {
505            Self::Grain => "Set grain",
506            Self::TimeRoles => "Time roles",
507            Self::Intervals => "Intervals",
508        }
509    }
510}
511
512/// Whether the Trends page can draw a column's measure across segments: that
513/// needs segments in an order, and more than one of them.
514pub fn shows_trend(plan: &DataQualityPlan, results: &DataQualityResults) -> bool {
515    matches!(
516        plan.grain,
517        QualityGrain::RowChunks(_) | QualityGrain::TimeWindows { .. } | QualityGrain::Partition(_)
518    ) && results.segments.len() + results.unsampled_segments.len() > 1
519}
520
521/// The plan setting a result page needs before it has anything to show, if any.
522/// Intervals need time roles, and ask only when there are dates to assign.
523pub fn page_setup(
524    page: QualityPage,
525    plan: &DataQualityPlan,
526    results: Option<&DataQualityResults>,
527    has_time_columns: bool,
528) -> Option<QualitySetup> {
529    let results = results?;
530    match page {
531        QualityPage::Segments if plan.grain == QualityGrain::Dataset => Some(QualitySetup::Grain),
532        QualityPage::Trends if !shows_trend(plan, results) => Some(QualitySetup::Grain),
533        // Roles that make no interval want a pair chosen; otherwise, roles. Pairs measuring
534        // nothing are fixed by neither, and the page says what is.
535        QualityPage::Intervals if results.temporal.is_empty() && has_time_columns => {
536            if plan.candidate_pairs().is_empty() {
537                Some(QualitySetup::TimeRoles)
538            } else if plan.interval_pairs().is_empty() {
539                Some(QualitySetup::Intervals)
540            } else {
541                None
542            }
543        }
544        _ => None,
545    }
546}
547
548#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
549pub enum QualityCompute {
550    Metadata,
551    #[default]
552    Sample,
553    Full,
554}
555
556impl QualityCompute {
557    pub fn label(self) -> &'static str {
558        match self {
559            Self::Metadata => "metadata",
560            Self::Sample => "sample",
561            Self::Full => "full",
562        }
563    }
564}
565
566#[derive(Debug, Clone, PartialEq, Eq, Default)]
567pub enum QualityGrain {
568    #[default]
569    Dataset,
570    File,
571    Partition(String),
572    RowChunks(usize),
573    TimeWindows {
574        column: String,
575        every: String,
576    },
577}
578
579impl QualityGrain {
580    /// How the rows are split, in words: "by day of date", "by year".
581    pub fn label(&self) -> String {
582        match self {
583            Self::Dataset => "whole dataset".to_string(),
584            Self::File => "by file".to_string(),
585            Self::Partition(column) => format!("by {column}"),
586            Self::RowChunks(rows) => {
587                format!("in chunks of {} rows", crate::numfmt::group_chrome(*rows))
588            }
589            Self::TimeWindows { column, every } => {
590                let unit = match every.as_str() {
591                    "1h" => "hour",
592                    "1d" => "day",
593                    "1w" => "week",
594                    "1mo" => "month",
595                    other => other,
596                };
597                // A column named for its unit would read "by day of day".
598                if column.eq_ignore_ascii_case(unit) {
599                    format!("by {unit} of the {column} column")
600                } else {
601                    format!("by {unit} of {column}")
602                }
603            }
604        }
605    }
606}
607
608#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
609pub enum QualityComparison {
610    #[default]
611    None,
612    Previous,
613    Baseline,
614}
615
616impl QualityComparison {
617    /// The comparison as a plan choice says it.
618    pub fn choice_label(self) -> &'static str {
619        match self {
620            Self::None => "none",
621            Self::Previous => "the segment before",
622            Self::Baseline => "a baseline segment (the first, or b on Segments)",
623        }
624    }
625
626    pub fn label(self) -> &'static str {
627        match self {
628            Self::None => "none",
629            Self::Previous => "previous",
630            Self::Baseline => "baseline (first)",
631        }
632    }
633}
634
635#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
636pub enum TemporalRole {
637    Event,
638    Effective,
639    PeriodEnd,
640    Created,
641    Published,
642    Received,
643    Processed,
644    ValidFrom,
645    ValidTo,
646}
647
648impl TemporalRole {
649    pub const ALL: [Self; 9] = [
650        Self::Event,
651        Self::Effective,
652        Self::PeriodEnd,
653        Self::Created,
654        Self::Published,
655        Self::Received,
656        Self::Processed,
657        Self::ValidFrom,
658        Self::ValidTo,
659    ];
660
661    pub fn label(self) -> &'static str {
662        match self {
663            Self::Event => "event",
664            Self::Effective => "effective/as-of",
665            Self::PeriodEnd => "period end",
666            Self::Created => "created",
667            Self::Published => "published",
668            Self::Received => "received",
669            Self::Processed => "processed",
670            Self::ValidFrom => "valid from",
671            Self::ValidTo => "valid to",
672        }
673    }
674}
675
676#[derive(Debug, Clone, PartialEq, Eq)]
677pub struct TemporalRoleAssignment {
678    pub role: TemporalRole,
679    pub column: String,
680    pub timezone: Option<String>,
681}
682
683/// The intervals measured when none are chosen: start role to end role, the pairs
684/// whose order the roles state. Others are chosen in Setup; a role in no interval
685/// measures nothing, as Setup says.
686pub const INTERVAL_PAIRS: [(TemporalRole, TemporalRole); 7] = [
687    (TemporalRole::Event, TemporalRole::Published),
688    (TemporalRole::Event, TemporalRole::Received),
689    (TemporalRole::PeriodEnd, TemporalRole::Published),
690    (TemporalRole::Published, TemporalRole::Received),
691    (TemporalRole::Received, TemporalRole::Processed),
692    (TemporalRole::Event, TemporalRole::Processed),
693    (TemporalRole::ValidFrom, TemporalRole::ValidTo),
694];
695
696/// `event to received`.
697pub fn interval_label((start, end): (TemporalRole, TemporalRole)) -> String {
698    format!("{} to {}", start.label(), end.label())
699}
700
701/// Which time puts an interval in a window under a time-window grain: the grain's
702/// column, or the interval's start or end (by end, counted on the day it finished).
703#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
704pub enum IntervalClock {
705    #[default]
706    Grain,
707    Start,
708    End,
709}
710
711impl IntervalClock {
712    pub const ALL: [Self; 3] = [Self::Grain, Self::Start, Self::End];
713
714    pub fn label(self) -> &'static str {
715        match self {
716            Self::Grain => "the grain's column",
717            Self::Start => "each interval's start",
718            Self::End => "each interval's end",
719        }
720    }
721}
722
723/// Whether text read as time is a date or a date with a time of day.
724#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
725pub enum TimeKind {
726    Date,
727    Datetime,
728}
729
730impl TimeKind {
731    pub fn label(self) -> &'static str {
732        match self {
733            Self::Date => "date",
734            Self::Datetime => "datetime",
735        }
736    }
737}
738
739/// The formats Setup offers for reading text as time, unambiguous first. Named, not
740/// inferred: every row reads the same way, and an unreadable value is counted.
741pub const TIME_FORMATS: [(TimeKind, &str); 16] = [
742    (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S"),
743    (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S"),
744    (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f"),
745    (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f"),
746    // ISO 8601 with an offset: `%#z` takes `Z`, `+05:00`, `-0500` and `+05`, and
747    // the values are read as instants in UTC.
748    (TimeKind::Datetime, "%Y-%m-%dT%H:%M:%S%.f%#z"),
749    (TimeKind::Datetime, "%Y-%m-%d %H:%M:%S%.f%#z"),
750    (TimeKind::Datetime, "%Y-%m-%d %H:%M"),
751    (TimeKind::Date, "%Y-%m-%d"),
752    (TimeKind::Date, "%Y%m%d"),
753    (TimeKind::Datetime, "%m/%d/%Y %H:%M:%S"),
754    (TimeKind::Datetime, "%m/%d/%Y %I:%M:%S %p"),
755    (TimeKind::Datetime, "%d/%m/%Y %H:%M:%S"),
756    (TimeKind::Datetime, "%d.%m.%Y %H:%M:%S"),
757    (TimeKind::Date, "%m/%d/%Y"),
758    (TimeKind::Date, "%d/%m/%Y"),
759    (TimeKind::Date, "%d.%m.%Y"),
760];
761
762/// A text column read as a date or datetime for one study: grain and time roles see
763/// the parsed value, other checks the stored text. Unreadable values count as
764/// unparsed, never as missing.
765#[derive(Debug, Clone, PartialEq, Eq, Hash)]
766pub struct TimeInterpretation {
767    pub column: String,
768    pub kind: TimeKind,
769    /// A strftime format, as Polars' `str.to_datetime` takes it. With an offset
770    /// (`%z`), the values are instants in UTC; without one, times with no zone.
771    pub format: String,
772}
773
774impl TimeInterpretation {
775    /// Whether the format reads an offset, so its values are instants in UTC
776    /// rather than times with no zone.
777    pub fn zoned(&self) -> bool {
778        self.format.contains('z')
779    }
780
781    /// `datetime %Y-%m-%d %H:%M:%S`.
782    pub fn label(&self) -> String {
783        format!("{} {}", self.kind.label(), self.format)
784    }
785
786    /// The column's values as time: null where the format does not read the text.
787    pub fn expr(&self) -> Expr {
788        let options = StrptimeOptions {
789            format: Some(PlSmallStr::from(self.format.as_str())),
790            strict: false,
791            exact: true,
792            cache: true,
793        };
794        // A categorical column holds codes; its values are read as the text they name.
795        let text = col(self.column.as_str()).cast(DataType::String).str();
796        match self.kind {
797            TimeKind::Date => text.to_date(options),
798            TimeKind::Datetime => text.to_datetime(
799                Some(TimeUnit::Microseconds),
800                None,
801                options,
802                lit(PlSmallStr::from_static("raise")),
803            ),
804        }
805    }
806
807    /// Rows holding text the format does not read.
808    pub fn unparsed(&self) -> Expr {
809        col(self.column.as_str())
810            .is_not_null()
811            .and(self.expr().is_null())
812    }
813
814    /// Whether the format reads `value`, the way a run will: for the examples Setup
815    /// shows beside each format, from rows already on screen.
816    pub fn reads(&self, value: &str) -> bool {
817        match self.kind {
818            TimeKind::Date => chrono::NaiveDate::parse_from_str(value, &self.format).is_ok(),
819            // An offset format needs the offset: without one there is no instant.
820            TimeKind::Datetime if self.zoned() => {
821                chrono::DateTime::parse_from_str(value, &self.format).is_ok()
822            }
823            TimeKind::Datetime => {
824                chrono::NaiveDateTime::parse_from_str(value, &self.format).is_ok()
825            }
826        }
827    }
828}
829
830/// What a Data Quality run is doing now. The worker names each stage as it enters
831/// it, and the progress view shows the latest.
832#[derive(Debug, Clone, Copy, PartialEq, Eq)]
833pub enum QualityStage {
834    Preparing,
835    CopyingSource,
836    ReusingSample,
837    ReadingSample,
838    CountingRows,
839    CountingSegments,
840    ProfilingColumns,
841    CheckingDuplicates,
842    CheckingKey,
843    CheckingSpellings,
844    ReadingConflicts,
845    ProfilingSegments,
846    ComputingIntervals,
847    CheckingSharedNulls,
848    /// An audio file's samples, read whole for clipping, runs of zeros and DC offset.
849    CheckingSignal,
850    Assembling,
851}
852
853impl QualityStage {
854    pub fn label(self) -> &'static str {
855        match self {
856            Self::Preparing => "Preparing the plan",
857            Self::CopyingSource => "Copying the source locally",
858            Self::ReusingSample => "Reusing the retained sample",
859            Self::ReadingSample => "Reading the sample",
860            Self::CountingRows => "Counting rows",
861            Self::CountingSegments => "Counting segment rows",
862            Self::ProfilingColumns => "Profiling columns",
863            Self::CheckingDuplicates => "Checking duplicate rows",
864            Self::CheckingKey => "Checking the declared key",
865            Self::CheckingSpellings => "Checking category spellings",
866            Self::ReadingConflicts => "Reading conflicting values",
867            Self::ProfilingSegments => "Profiling segments",
868            Self::ComputingIntervals => "Computing intervals",
869            Self::CheckingSharedNulls => "Checking columns missing together",
870            Self::CheckingSignal => "Checking the signal",
871            Self::Assembling => "Assembling the report",
872        }
873    }
874}
875
876/// A stage, whether it reads the source or works on rows already read, and whether
877/// a cancel stops it partway.
878#[derive(Debug, Clone, Copy, PartialEq, Eq)]
879pub struct QualityPhase {
880    pub stage: QualityStage,
881    pub reads_source: bool,
882    /// A cancel ends this stage within a batch. A stage that is one collect Polars
883    /// cannot watch runs to its end, and the screen says so while it does.
884    pub interruptible: bool,
885}
886
887/// What a run's reads of the source were seen to do, counted as the rows went by.
888/// Bytes and requests are not counted: a Polars scan does not report them.
889#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
890pub struct ObservedReads {
891    /// Stages that read the source.
892    pub reads: usize,
893    /// Of those, the ones whose rows were counted as they came.
894    pub counted: usize,
895    /// Rows the counted reads passed through from the scope, over every pass.
896    pub rows: usize,
897    /// The local copy a full scan's passes read instead of the source, when they did.
898    pub copy: Option<CopyRead>,
899}
900
901/// A local copy of a remote source that a full scan's passes read.
902#[derive(Debug, Clone, Copy, PartialEq, Eq)]
903pub struct CopyRead {
904    pub bytes: u64,
905    pub objects: usize,
906    /// This run fetched it; otherwise an earlier run did and this one reused it.
907    pub fetched: bool,
908}
909
910/// A run's line to the screen: stages as entered, rows its reads have seen, and a
911/// stop checked between stages and between read batches.
912#[derive(Clone, Default)]
913pub struct QualityWatch {
914    read: crate::analysis::sampling::ReadWatch,
915    report: Option<Arc<dyn Fn(QualityPhase) + Send + Sync>>,
916    /// The stage under way, and what the stages before it read.
917    last: Arc<std::sync::Mutex<(Option<QualityPhase>, ObservedReads)>>,
918    /// Set once the scope reads a local copy: its passes then read no source.
919    copy: Arc<std::sync::OnceLock<CopyRead>>,
920}
921
922impl std::fmt::Debug for QualityWatch {
923    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
924        f.debug_struct("QualityWatch")
925            .field("read", &self.read)
926            .finish_non_exhaustive()
927    }
928}
929
930impl QualityWatch {
931    /// A watch that hands each new stage to `report`.
932    pub fn new(report: impl Fn(QualityPhase) + Send + Sync + 'static) -> Self {
933        Self {
934            report: Some(Arc::new(report)),
935            ..Self::default()
936        }
937    }
938
939    pub fn cancel(&self) {
940        self.read.stop();
941    }
942
943    pub fn cancelled(&self) -> bool {
944        self.read.stopped()
945    }
946
947    /// The reads' side: the stop, and the rows the stage under way has seen.
948    pub fn read(&self) -> &crate::analysis::sampling::ReadWatch {
949        &self.read
950    }
951
952    /// What the run's reads were seen to do, the one under way included.
953    pub fn observed(&self) -> ObservedReads {
954        let Ok(last) = self.last.lock() else {
955            return ObservedReads::default();
956        };
957        let (phase, mut observed) = *last;
958        if phase.is_some_and(|phase| phase.reads_source) {
959            observed.reads += 1;
960            if let Some(rows) = self.read.rows_seen() {
961                observed.counted += 1;
962                observed.rows += rows;
963            }
964        }
965        observed.copy = self.copy.get().copied();
966        observed
967    }
968
969    /// The scope is read from a local copy from here on.
970    pub(crate) fn use_copy(&self, copy: CopyRead) {
971        let _ = self.copy.set(copy);
972    }
973
974    /// Whether a pass over the scope reads the source: not once it reads a copy.
975    fn scope_reads(&self, reads: bool) -> bool {
976        reads && self.copy.get().is_none()
977    }
978
979    /// `lf` watched: each batch reaching the top is counted, and after a cancel the next
980    /// fails the query (streaming stops within a batch; in memory the scope is one batch).
981    /// Projections and filters still push down.
982    fn watched(&self, lf: &LazyFrame) -> LazyFrame {
983        let read = self.read.clone();
984        lf.clone().map(
985            move |df: DataFrame| {
986                if read.stopped() {
987                    return Err(PolarsError::ComputeError(
988                        crate::analysis::sampling::CANCELLED.into(),
989                    ));
990                }
991                read.saw(df.height());
992                Ok(df)
993            },
994            OptFlags::PROJECTION_PUSHDOWN | OptFlags::PREDICATE_PUSHDOWN | OptFlags::STREAMING,
995            None,
996            Some("quality watch"),
997        )
998    }
999
1000    /// Enter `stage`: said once however often entered, refused after a cancel (runs stop
1001    /// between stages). Leaving a source-reading stage adds its rows to what was observed.
1002    pub(crate) fn stage(
1003        &self,
1004        stage: QualityStage,
1005        reads_source: bool,
1006        interruptible: bool,
1007    ) -> Result<()> {
1008        self.read.check()?;
1009        let phase = QualityPhase {
1010            stage,
1011            reads_source,
1012            interruptible,
1013        };
1014        let mut last = self
1015            .last
1016            .lock()
1017            .map_err(|_| Report::msg("quality progress lock failed"))?;
1018        let (previous, observed) = &mut *last;
1019        if *previous != Some(phase) {
1020            // Rows on screen are the stage's own, so each read counts from zero.
1021            let seen = self.read.restart();
1022            if previous.is_some_and(|phase| phase.reads_source) {
1023                observed.reads += 1;
1024                if let Some(rows) = seen {
1025                    observed.counted += 1;
1026                    observed.rows += rows;
1027                }
1028            }
1029            *previous = Some(phase);
1030            if let Some(report) = &self.report {
1031                report(phase);
1032            }
1033        }
1034        Ok(())
1035    }
1036
1037    /// A watched collect's error, or the cancel that caused it: stopped partway, it
1038    /// fails as a stopped sampler does.
1039    fn failed(&self, error: impl Into<Report>) -> Report {
1040        if self.cancelled() {
1041            Report::msg(crate::analysis::sampling::CANCELLED)
1042        } else {
1043            error.into()
1044        }
1045    }
1046}
1047
1048#[derive(Debug, Clone, PartialEq, Eq)]
1049pub struct DataQualityPlan {
1050    pub scope: QualityScope,
1051    pub compute: QualityCompute,
1052    /// How a dataset-grain sample picks its rows, from the shared analysis sample.
1053    pub method: crate::analysis::sampling::SampleMethod,
1054    /// Rows a dataset-grain sample keeps: the shared analysis sample's size.
1055    pub dataset_rows: usize,
1056    pub sample_seed: u64,
1057    pub grain: QualityGrain,
1058    pub comparison: QualityComparison,
1059    pub baseline_segment: Option<String>,
1060    pub temporal_roles: Vec<TemporalRoleAssignment>,
1061    /// The intervals chosen under Intervals, start role to end role. `None` until
1062    /// one is chosen: the suggested pairs the assigned roles make.
1063    pub intervals: Option<Vec<(TemporalRole, TemporalRole)>>,
1064    /// Which time puts an interval in a time window.
1065    pub interval_clock: IntervalClock,
1066    pub latency_threshold_seconds: Option<i64>,
1067    /// Text columns read as time for this study, by grain and roles only.
1068    pub time_formats: Vec<TimeInterpretation>,
1069    /// The time windows rows are expected in, when stated, making an empty window a gap.
1070    /// Applied to counted segments: changes the report, never the read.
1071    pub expected: Option<ExpectedWindows>,
1072    /// What the columns must hold, declared: the key and each column's rules. Read
1073    /// from the rows the run reads; it decides no rows, so it is the report's.
1074    pub intent: crate::analysis::quality_intent::DeclaredIntent,
1075}
1076
1077/// Which windows a study expects rows in, from Setup: every window of the grain or
1078/// weekdays only, within a range. Unset, no window is a gap.
1079#[derive(Debug, Clone, PartialEq, Eq, Default)]
1080pub struct ExpectedWindows {
1081    /// Only Monday to Friday's hours or days are expected.
1082    pub weekdays: bool,
1083    /// The first expected time, as typed: a date or a UTC timestamp. `None` starts
1084    /// at the first window the run found.
1085    pub from: Option<String>,
1086    /// The time expected windows end before. `None` ends after the last window the
1087    /// run found.
1088    pub before: Option<String>,
1089}
1090
1091impl ExpectedWindows {
1092    /// Whether `every` is a width whose windows can fall on a weekend: an hour or a
1093    /// day. A week or a month always holds weekdays.
1094    pub fn weekdays_apply(every: &str) -> bool {
1095        matches!(every, "1h" | "1d")
1096    }
1097
1098    /// The cadence, in Setup's words: "every day", "weekdays".
1099    pub fn cadence_label(&self, every: &str) -> String {
1100        if self.weekdays && Self::weekdays_apply(every) {
1101            "weekdays".to_string()
1102        } else {
1103            let unit = match every {
1104                "1h" => "hour",
1105                "1d" => "day",
1106                "1w" => "week",
1107                "1mo" => "month",
1108                other => other,
1109            };
1110            format!("every {unit}")
1111        }
1112    }
1113
1114    /// The range, in Setup's words: "2024-01-01 to before 2024-04-01", or the windows
1115    /// found where a side is not stated.
1116    pub fn range_label(&self) -> String {
1117        match (self.from.as_deref(), self.before.as_deref()) {
1118            (None, None) => "first to last window found".to_string(),
1119            (Some(from), None) => format!("{from} to the last window found"),
1120            (None, Some(before)) => format!("first window found to before {before}"),
1121            (Some(from), Some(before)) => format!("{from} to before {before}"),
1122        }
1123    }
1124
1125    /// Why the typed range cannot be read, if it cannot.
1126    pub fn problem(&self) -> Option<String> {
1127        let read = |text: &Option<String>| match text.as_deref() {
1128            None => Ok(None),
1129            Some(text) => parse_scope_time(text)
1130                .map(Some)
1131                .ok_or_else(|| format!("{text} is not a date or UTC timestamp")),
1132        };
1133        match (read(&self.from), read(&self.before)) {
1134            (Err(problem), _) | (_, Err(problem)) => Some(problem),
1135            (Ok(Some(from)), Ok(Some(before))) if before <= from => {
1136                Some("Before must be after From".to_string())
1137            }
1138            _ => None,
1139        }
1140    }
1141
1142    /// The typed range in microseconds since the epoch, each side when stated and
1143    /// readable.
1144    pub fn bounds(&self) -> (Option<i64>, Option<i64>) {
1145        (
1146            self.from.as_deref().and_then(parse_scope_time),
1147            self.before.as_deref().and_then(parse_scope_time),
1148        )
1149    }
1150}
1151
1152impl Default for DataQualityPlan {
1153    fn default() -> Self {
1154        Self {
1155            scope: QualityScope::CurrentView,
1156            compute: QualityCompute::Sample,
1157            method: crate::analysis::sampling::SampleMethod::Spread,
1158            dataset_rows: DEFAULT_SAMPLE_ROWS,
1159            sample_seed: 42_891,
1160            grain: QualityGrain::Dataset,
1161            comparison: QualityComparison::None,
1162            baseline_segment: None,
1163            temporal_roles: Vec::new(),
1164            intervals: None,
1165            interval_clock: IntervalClock::Grain,
1166            latency_threshold_seconds: None,
1167            time_formats: Vec::new(),
1168            expected: None,
1169            intent: crate::analysis::quality_intent::DeclaredIntent::default(),
1170        }
1171    }
1172}
1173
1174impl DataQualityPlan {
1175    /// Only a full scan asks first. Every grain reads the shared sample and cuts it
1176    /// into segments, so no grain reads more than the sample says.
1177    pub fn requires_confirmation(&self) -> bool {
1178        self.compute == QualityCompute::Full
1179    }
1180
1181    pub fn comparison_label(&self) -> String {
1182        if self.comparison == QualityComparison::Baseline {
1183            self.baseline_segment
1184                .as_ref()
1185                .map(|label| format!("baseline: {label}"))
1186                .unwrap_or_else(|| self.comparison.label().to_string())
1187        } else {
1188            self.comparison.label().to_string()
1189        }
1190    }
1191
1192    /// The shared analysis sample this plan carries, as the Sample form shows it.
1193    pub fn sample(&self) -> crate::analysis::sampling::Sample {
1194        crate::analysis::sampling::Sample {
1195            scope: self.scope.clone(),
1196            method: self.method.clone(),
1197            rows: self.dataset_rows,
1198            seed: self.sample_seed,
1199        }
1200    }
1201
1202    /// Take `sample` as this plan's rows: metadata-only stays so; else every row is a full
1203    /// read, fewer a sampled one. Equal rows per value of a column sets that column as
1204    /// the grain, only when the choice is new and no grain was set.
1205    pub fn adopt_sample(&mut self, sample: &crate::analysis::sampling::Sample) {
1206        if self.scope != sample.scope {
1207            self.baseline_segment = None;
1208        }
1209        self.scope = sample.scope.clone();
1210        self.sample_seed = sample.seed;
1211        self.dataset_rows = sample.rows;
1212        if self.compute != QualityCompute::Metadata {
1213            self.compute = if sample.method == crate::analysis::sampling::SampleMethod::EveryRow {
1214                QualityCompute::Full
1215            } else {
1216                QualityCompute::Sample
1217            };
1218        }
1219        if let crate::analysis::sampling::SampleMethod::PerPartition { column } = &sample.method
1220            && self.method != sample.method
1221            && self.grain == QualityGrain::Dataset
1222        {
1223            self.grain = QualityGrain::Partition(column.clone());
1224            self.baseline_segment = None;
1225        }
1226        self.method = sample.method.clone();
1227    }
1228
1229    /// How `column` is read as time, when it is text read through a format.
1230    pub fn time_format(&self, column: &str) -> Option<&TimeInterpretation> {
1231        self.time_formats
1232            .iter()
1233            .find(|interpretation| interpretation.column == column)
1234    }
1235
1236    /// A column's values as time: parsed through its format when it has one, and as
1237    /// stored otherwise.
1238    pub fn time_value(&self, column: &str) -> Expr {
1239        self.time_format(column)
1240            .map(TimeInterpretation::expr)
1241            .unwrap_or_else(|| col(column))
1242    }
1243
1244    /// Whether `column` of `schema` can be read as time: a date or time type, or text
1245    /// with a format.
1246    pub fn reads_as_time(&self, column: &str, schema: &Schema) -> bool {
1247        self.time_format(column).is_some() || schema.get(column).is_some_and(DataType::is_temporal)
1248    }
1249
1250    /// The column a role is assigned to.
1251    pub fn role_column(&self, role: TemporalRole) -> Option<&str> {
1252        self.temporal_roles
1253            .iter()
1254            .find(|assignment| assignment.role == role)
1255            .map(|assignment| assignment.column.as_str())
1256    }
1257
1258    /// The intervals this plan measures: the chosen ones, or until one is chosen the
1259    /// suggested pairs; either way only those whose two roles are assigned.
1260    pub fn interval_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
1261        let assigned = |(start, end): &(TemporalRole, TemporalRole)| {
1262            self.role_column(*start).is_some() && self.role_column(*end).is_some()
1263        };
1264        match &self.intervals {
1265            None => INTERVAL_PAIRS.into_iter().filter(assigned).collect(),
1266            Some(chosen) => chosen.iter().copied().filter(assigned).collect(),
1267        }
1268    }
1269
1270    /// Every start and end the assigned roles can make, the suggested pairs first,
1271    /// then the rest in role order: what Intervals in Setup lists.
1272    pub fn candidate_pairs(&self) -> Vec<(TemporalRole, TemporalRole)> {
1273        let roles = TemporalRole::ALL
1274            .into_iter()
1275            .filter(|role| self.role_column(*role).is_some())
1276            .collect::<Vec<_>>();
1277        let mut pairs = INTERVAL_PAIRS
1278            .into_iter()
1279            .filter(|(start, end)| roles.contains(start) && roles.contains(end))
1280            .collect::<Vec<_>>();
1281        for start in &roles {
1282            for end in &roles {
1283                if start != end && !pairs.contains(&(*start, *end)) {
1284                    pairs.push((*start, *end));
1285                }
1286            }
1287        }
1288        pairs
1289    }
1290
1291    /// Measure `pair`, or stop measuring it. The first choice makes the list
1292    /// explicit, starting from what was measured.
1293    pub fn toggle_interval(&mut self, pair: (TemporalRole, TemporalRole)) {
1294        let mut chosen = self.interval_pairs();
1295        match chosen.iter().position(|chosen| *chosen == pair) {
1296            Some(index) => {
1297                chosen.remove(index);
1298            }
1299            None => chosen.push(pair),
1300        }
1301        self.intervals = Some(chosen);
1302    }
1303
1304    /// Assigned roles that no measured interval uses: they measure nothing.
1305    pub fn unpaired_roles(&self) -> Vec<TemporalRole> {
1306        let pairs = self.interval_pairs();
1307        TemporalRole::ALL
1308            .into_iter()
1309            .filter(|role| {
1310                self.role_column(*role).is_some()
1311                    && !pairs
1312                        .iter()
1313                        .any(|(start, end)| start == role || end == role)
1314            })
1315            .collect()
1316    }
1317
1318    /// Whether the clock choice means anything: intervals cut into time windows.
1319    pub fn windows_intervals(&self) -> bool {
1320        matches!(self.grain, QualityGrain::TimeWindows { .. }) && !self.interval_pairs().is_empty()
1321    }
1322
1323    /// The grain an interval from `start` to `end` is cut by: the plan's, except
1324    /// that time windows go by the interval's own start or end when the clock says.
1325    pub fn interval_grain(&self, start: &str, end: &str) -> QualityGrain {
1326        match (&self.grain, self.interval_clock) {
1327            (QualityGrain::TimeWindows { every, .. }, IntervalClock::Start) => {
1328                QualityGrain::TimeWindows {
1329                    column: start.to_string(),
1330                    every: every.clone(),
1331                }
1332            }
1333            (QualityGrain::TimeWindows { every, .. }, IntervalClock::End) => {
1334                QualityGrain::TimeWindows {
1335                    column: end.to_string(),
1336                    every: every.clone(),
1337                }
1338            }
1339            (grain, _) => grain.clone(),
1340        }
1341    }
1342
1343    /// Whether `column` holds instants (a zoned type, or text read with an
1344    /// offset) rather than times with no zone; `None` when it is not read as time.
1345    pub fn zoned(&self, column: &str, schema: &Schema) -> Option<bool> {
1346        if let Some(format) = self.time_format(column) {
1347            return Some(format.zoned());
1348        }
1349        match schema.get(column)? {
1350            DataType::Datetime(_, zone) => Some(zone.is_some()),
1351            DataType::Date => Some(false),
1352            _ => None,
1353        }
1354    }
1355
1356    /// Whether `other` measures the same, differing at most in expected windows or
1357    /// comparison, both worked out from the report's counts
1358    /// ([`DataQualityResults::compare_segments`]).
1359    pub fn same_measurement(&self, other: &Self) -> bool {
1360        let measured = |plan: &Self| Self {
1361            expected: None,
1362            comparison: QualityComparison::None,
1363            baseline_segment: None,
1364            ..plan.clone()
1365        };
1366        measured(self) == measured(other)
1367    }
1368
1369    /// Whether `other` compares segments differently from this plan.
1370    pub fn compares_differently(&self, other: &Self) -> bool {
1371        self.comparison != other.comparison || self.baseline_segment != other.baseline_segment
1372    }
1373
1374    /// The windows this plan expects rows in: only on a time-window grain.
1375    pub fn expected_windows(&self) -> Option<&ExpectedWindows> {
1376        matches!(self.grain, QualityGrain::TimeWindows { .. })
1377            .then_some(self.expected.as_ref())
1378            .flatten()
1379    }
1380
1381    /// The next coarser grain for thin segments (hour→day→week→month, a larger row chunk);
1382    /// none for partitions and files.
1383    pub fn coarser_grain(&self) -> Option<QualityGrain> {
1384        match &self.grain {
1385            QualityGrain::TimeWindows { column, every } => {
1386                let coarser = match every.as_str() {
1387                    "1h" => "1d",
1388                    "1d" => "1w",
1389                    "1w" => "1mo",
1390                    _ => return None,
1391                };
1392                Some(QualityGrain::TimeWindows {
1393                    column: column.clone(),
1394                    every: coarser.to_string(),
1395                })
1396            }
1397            QualityGrain::RowChunks(rows) if *rows < DEFAULT_CHUNK_ROWS => {
1398                Some(QualityGrain::RowChunks(DEFAULT_CHUNK_ROWS))
1399            }
1400            _ => None,
1401        }
1402    }
1403}
1404
1405#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1406pub enum QualityPrecision {
1407    Metadata,
1408    Sampled,
1409    Exact,
1410}
1411
1412impl QualityPrecision {
1413    pub fn label(self) -> &'static str {
1414        match self {
1415            Self::Metadata => "metadata",
1416            Self::Sampled => "sampled",
1417            Self::Exact => "exact",
1418        }
1419    }
1420}
1421
1422#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1423pub enum QualityMetric {
1424    #[default]
1425    NullRate,
1426    EmptyRate,
1427    WhitespaceRate,
1428    NonFiniteRate,
1429    DistinctShare,
1430    IntegerParseShare,
1431    DecimalParseShare,
1432}
1433
1434impl QualityMetric {
1435    pub const ALL: [Self; 7] = [
1436        Self::NullRate,
1437        Self::EmptyRate,
1438        Self::WhitespaceRate,
1439        Self::NonFiniteRate,
1440        Self::DistinctShare,
1441        Self::IntegerParseShare,
1442        Self::DecimalParseShare,
1443    ];
1444
1445    pub fn label(self) -> &'static str {
1446        match self {
1447            Self::NullRate => "Null rate",
1448            Self::EmptyRate => "Empty rate",
1449            Self::WhitespaceRate => "Whitespace rate",
1450            Self::NonFiniteRate => "Non-finite rate",
1451            Self::DistinctShare => "Distinct share",
1452            Self::IntegerParseShare => "Integer parse share",
1453            Self::DecimalParseShare => "Decimal parse share",
1454        }
1455    }
1456
1457    /// The rows a rate is taken over: every row, or the rows with a value.
1458    pub fn denominator(self, column: &ColumnQualityProfile) -> usize {
1459        match self {
1460            Self::NullRate | Self::EmptyRate | Self::WhitespaceRate | Self::NonFiniteRate => {
1461                column.evaluated_rows
1462            }
1463            Self::DistinctShare | Self::IntegerParseShare | Self::DecimalParseShare => {
1464                column.non_null_rows()
1465            }
1466        }
1467    }
1468
1469    /// The measure's name in a table cell or a change: "nulls", "distinct".
1470    pub fn short_label(self) -> &'static str {
1471        match self {
1472            Self::NullRate => "nulls",
1473            Self::EmptyRate => "empty",
1474            Self::WhitespaceRate => "blank",
1475            Self::NonFiniteRate => "NaN/inf",
1476            Self::DistinctShare => "distinct",
1477            Self::IntegerParseShare => "integer parse",
1478            Self::DecimalParseShare => "decimal parse",
1479        }
1480    }
1481
1482    pub fn value(self, column: &ColumnQualityProfile) -> Option<f64> {
1483        let ratio = |numerator: usize, denominator: usize| {
1484            (denominator > 0).then(|| numerator as f64 / denominator as f64)
1485        };
1486        match self {
1487            Self::NullRate => ratio(column.null_count, column.evaluated_rows),
1488            Self::EmptyRate => ratio(column.empty_count?, column.evaluated_rows),
1489            Self::WhitespaceRate => ratio(column.whitespace_count?, column.evaluated_rows),
1490            Self::NonFiniteRate => ratio(
1491                column.nan_count?
1492                    + column.positive_infinity_count?
1493                    + column.negative_infinity_count?,
1494                column.evaluated_rows,
1495            ),
1496            Self::DistinctShare => ratio(column.distinct_count?, column.non_null_rows()),
1497            Self::IntegerParseShare => ratio(column.integer_parse_count?, column.non_null_rows()),
1498            Self::DecimalParseShare => ratio(column.decimal_parse_count?, column.non_null_rows()),
1499        }
1500    }
1501}
1502
1503#[derive(Debug, Clone)]
1504pub struct ColumnQualityProfile {
1505    pub name: String,
1506    pub dtype: DataType,
1507    pub evaluated_rows: usize,
1508    pub null_count: usize,
1509    pub empty_count: Option<usize>,
1510    pub whitespace_count: Option<usize>,
1511    pub nan_count: Option<usize>,
1512    pub positive_infinity_count: Option<usize>,
1513    pub negative_infinity_count: Option<usize>,
1514    pub distinct_count: Option<usize>,
1515    pub min: Option<String>,
1516    pub max: Option<String>,
1517    pub integer_parse_count: Option<usize>,
1518    pub decimal_parse_count: Option<usize>,
1519    pub date_parse_count: Option<usize>,
1520    pub datetime_parse_count: Option<usize>,
1521    /// Text values that parse as whole numbers and are written with a leading zero:
1522    /// the mark of a code (a ZIP, an account, an industry code) rather than a number.
1523    pub leading_zero_count: Option<usize>,
1524    pub dominant_value: Option<String>,
1525    pub dominant_count: Option<usize>,
1526    pub min_length: Option<usize>,
1527    pub max_length: Option<usize>,
1528}
1529
1530impl ColumnQualityProfile {
1531    /// A column with its null count at zero and nothing else measured.
1532    pub fn unmeasured(name: &str, dtype: DataType, evaluated_rows: usize) -> Self {
1533        Self {
1534            name: name.to_string(),
1535            dtype,
1536            evaluated_rows,
1537            null_count: 0,
1538            empty_count: None,
1539            whitespace_count: None,
1540            nan_count: None,
1541            positive_infinity_count: None,
1542            negative_infinity_count: None,
1543            distinct_count: None,
1544            min: None,
1545            max: None,
1546            integer_parse_count: None,
1547            decimal_parse_count: None,
1548            date_parse_count: None,
1549            datetime_parse_count: None,
1550            leading_zero_count: None,
1551            dominant_value: None,
1552            dominant_count: None,
1553            min_length: None,
1554            max_length: None,
1555        }
1556    }
1557
1558    pub fn non_null_rows(&self) -> usize {
1559        self.evaluated_rows.saturating_sub(self.null_count)
1560    }
1561
1562    pub fn uniqueness_rate(&self) -> Option<f64> {
1563        self.distinct_count
1564            .map(|count| rate(count, self.non_null_rows()))
1565    }
1566}
1567
1568#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1569pub enum ObservationKind {
1570    Nulls,
1571    Empty,
1572    Whitespace,
1573    NonFinite,
1574    Constant,
1575    ParseableText,
1576    DuplicateRows,
1577    CategoryVariants,
1578    /// Rows whose own file has no such column. Their cells are absent, not null, and
1579    /// no measurement over values can tell the two apart.
1580    Absent,
1581    /// Rows whose file holds the column in a type the dataset's schema cannot read, so
1582    /// the column is not read from that file at all.
1583    TypeConflict,
1584    /// A column whose values are nearly unique and still repeat: the shape of a key
1585    /// that is not quite one.
1586    KeyLike,
1587    /// Text read as time that the chosen format does not read.
1588    UnparsedTime,
1589    /// Rows sharing a value of the declared key.
1590    KeyRepeated,
1591    /// Rows with no value in some part of the declared key.
1592    KeyMissing,
1593    /// A column declared required, with no value.
1594    RequiredMissing,
1595    /// Values outside a column's declared allowed set.
1596    NotAllowed,
1597    /// Values outside a column's declared range.
1598    OutOfRange,
1599    /// Text declared to read as a number that does not.
1600    UnparsedNumber,
1601    /// Audio samples in runs at full scale: the waveform cut flat at the limit.
1602    Clipping,
1603    /// Audio samples in long runs of exact zeros: dropouts, or digital silence.
1604    ZeroRuns,
1605    /// An audio channel whose mean sits away from zero.
1606    DcOffset,
1607}
1608
1609/// One file behind a drift observation: what it holds, and what that costs the column.
1610#[derive(Debug, Clone, PartialEq, Eq)]
1611pub struct QualityFileEvidence {
1612    /// Position in the Scope page's file inventory, which numbers files from 1.
1613    pub number: usize,
1614    pub name: String,
1615    /// Rows this file holds, from its footer.
1616    pub rows: usize,
1617    /// The type this file holds the column in, when the scan cannot read it as the
1618    /// dataset's. `None` for a file that simply has no such column.
1619    pub stored_type: Option<String>,
1620    /// The first values this file holds, read at its own type and rendered as text.
1621    /// Empty until a full run reads them.
1622    pub examples: Vec<String>,
1623}
1624
1625#[derive(Debug, Clone)]
1626pub struct QualityObservation {
1627    pub kind: ObservationKind,
1628    pub column: String,
1629    pub affected_rows: usize,
1630    pub evaluated_rows: usize,
1631    /// What only the engine can say of a drift or audio measurement: the files, the
1632    /// runs. Empty for every kind the report phrases from the numbers itself.
1633    pub fact: String,
1634    pub normalized_category: Option<String>,
1635    /// The files behind an [`ObservationKind::Absent`] or
1636    /// [`ObservationKind::TypeConflict`] measurement, commonest first; empty for checks
1637    /// over values.
1638    pub files: Vec<QualityFileEvidence>,
1639    /// The format an [`ObservationKind::UnparsedTime`] measurement read the text with.
1640    pub time_format: Option<TimeInterpretation>,
1641    /// The values at or past which an audio sample is at full scale, for an
1642    /// [`ObservationKind::Clipping`] measurement's rows.
1643    pub full_scale: Option<(f64, f64)>,
1644}
1645
1646impl QualityObservation {
1647    /// The scope holding this observation's rows when they are a set of files (absent or
1648    /// conflicting cells have no value to filter on).
1649    pub fn evidence_scope(&self) -> Option<QualityScope> {
1650        if !matches!(
1651            self.kind,
1652            ObservationKind::Absent | ObservationKind::TypeConflict
1653        ) || self.files.is_empty()
1654        {
1655            return None;
1656        }
1657        Some(QualityScope::SourceFiles(
1658            self.files.iter().map(|file| file.number).collect(),
1659        ))
1660    }
1661
1662    /// The rows behind this observation as a predicate over the run's rows; the engine
1663    /// counts with the same expression where it can, so count and rows agree.
1664    pub fn evidence_predicate(&self, results: &DataQualityResults) -> Option<Expr> {
1665        let value = col(&self.column);
1666        match self.kind {
1667            ObservationKind::Nulls => Some(value.is_null()),
1668            ObservationKind::Empty => Some(value.eq(lit(""))),
1669            ObservationKind::Whitespace => Some(
1670                value
1671                    .clone()
1672                    .cast(DataType::String)
1673                    .str()
1674                    .strip_chars(lit(LiteralValue::untyped_null()))
1675                    .eq(lit(""))
1676                    .and(value.neq(lit(""))),
1677            ),
1678            ObservationKind::NonFinite => Some(
1679                value
1680                    .clone()
1681                    .is_nan()
1682                    .or(value.clone().eq(lit(f64::INFINITY)))
1683                    .or(value.eq(lit(f64::NEG_INFINITY))),
1684            ),
1685            ObservationKind::Constant => Some(value.is_not_null()),
1686            ObservationKind::CategoryVariants => Some(
1687                value
1688                    .cast(DataType::String)
1689                    .str()
1690                    .strip_chars(lit(LiteralValue::untyped_null()))
1691                    .str()
1692                    .to_lowercase()
1693                    .eq(lit(self.normalized_category.clone()?)),
1694            ),
1695            // Nearly unique yet repeating: the rows whose value repeats; nulls are outside the
1696            // measurement.
1697            ObservationKind::KeyLike => {
1698                Some(value.clone().is_duplicated().and(value.is_not_null()))
1699            }
1700            ObservationKind::UnparsedTime => {
1701                self.time_format.as_ref().map(TimeInterpretation::unparsed)
1702            }
1703            // Every sample at full scale, in a run or not: the runs are what is
1704            // counted, and the samples around them are what a look wants.
1705            ObservationKind::Clipping => self
1706                .full_scale
1707                .map(|(low, high)| value.clone().lt_eq(lit(low)).or(value.gt_eq(lit(high)))),
1708            ObservationKind::ZeroRuns => Some(value.eq(lit(0))),
1709            // The values that stop a cast: text the reading does not parse.
1710            ObservationKind::ParseableText => unparsed_text(
1711                results
1712                    .columns
1713                    .iter()
1714                    .find(|profile| profile.name == self.column)?,
1715            ),
1716            // Declared rules find their rows through what the run measured them with.
1717            ObservationKind::KeyRepeated => results.intent.as_ref()?.repeated_key(),
1718            ObservationKind::KeyMissing => results.intent.as_ref()?.missing_key(),
1719            ObservationKind::RequiredMissing => {
1720                results.intent.as_ref()?.required_missing(&self.column)
1721            }
1722            ObservationKind::NotAllowed => results.intent.as_ref()?.not_allowed(&self.column),
1723            ObservationKind::OutOfRange => results.intent.as_ref()?.out_of_range(&self.column),
1724            ObservationKind::UnparsedNumber => {
1725                results.intent.as_ref()?.unparsed_number(&self.column)
1726            }
1727            // Duplicates are rows equal to another, not a predicate; absent and conflicting rows
1728            // are named by files; an offset is in every sample.
1729            ObservationKind::DuplicateRows
1730            | ObservationKind::Absent
1731            | ObservationKind::TypeConflict
1732            | ObservationKind::DcOffset => None,
1733        }
1734    }
1735}
1736
1737#[derive(Debug, Clone)]
1738pub struct IdentityProfile {
1739    pub duplicate_groups: usize,
1740    pub extra_rows: usize,
1741    pub rows_involved: usize,
1742    pub evaluated_rows: usize,
1743    /// The most copied groups, from the rows the run kept. Empty after a full scan,
1744    /// which keeps no rows.
1745    pub examples: Vec<DuplicateExample>,
1746}
1747
1748/// One group of identical rows: how many there are, and the row, a value a column.
1749#[derive(Debug, Clone, PartialEq, Eq)]
1750pub struct DuplicateExample {
1751    pub copies: usize,
1752    /// Rendered for reading: text quoted, a null as `null`.
1753    pub values: Vec<String>,
1754}
1755
1756/// A few of the values behind one column's finding, from the rows the run kept.
1757#[derive(Debug, Clone, PartialEq, Eq)]
1758pub struct FindingExamples {
1759    pub kind: ObservationKind,
1760    pub column: String,
1761    /// Distinct values, first seen first, quoted.
1762    pub values: Vec<String>,
1763}
1764
1765#[derive(Debug, Clone)]
1766pub struct CategoryVariantGroup {
1767    pub column: String,
1768    pub normalized: String,
1769    pub variants: Vec<(String, usize)>,
1770    pub rows_involved: usize,
1771}
1772
1773#[derive(Debug, Clone)]
1774pub struct SegmentQualityProfile {
1775    pub label: String,
1776    pub total_rows: Option<usize>,
1777    pub evaluated_rows: usize,
1778    pub columns: Vec<ColumnQualityProfile>,
1779    pub null_cells: usize,
1780    pub null_rate: f64,
1781    pub compared_with: Option<String>,
1782    pub largest_change: Option<String>,
1783    /// How big `largest_change` is (points, or percent for a row count), to rank
1784    /// segments by; `None` when nothing clear moved.
1785    pub change_size: Option<f64>,
1786}
1787
1788/// One column's measure in a segment, and in the segment it is compared with.
1789#[derive(Debug, Clone, PartialEq)]
1790pub struct SegmentChange {
1791    pub column: String,
1792    pub metric: QualityMetric,
1793    pub before: Option<f64>,
1794    pub now: f64,
1795    /// The move is past sampling noise (always, on an exact profile) and a point
1796    /// or more.
1797    pub clear: bool,
1798}
1799
1800impl SegmentChange {
1801    /// Percentage points moved, when there is something to have moved from.
1802    pub fn change(&self) -> Option<f64> {
1803        self.before.map(|before| (self.now - before) * 100.0)
1804    }
1805}
1806
1807/// The order Segments lists its rows in: as they fall, or the clearest change
1808/// first (ties, and segments with no clear change, keep their order).
1809pub fn segment_order(results: &DataQualityResults, by_change: bool) -> Vec<usize> {
1810    let mut order = (0..results.segments.len()).collect::<Vec<_>>();
1811    if by_change {
1812        order.sort_by(|&left, &right| {
1813            let size = |index: usize| results.segments[index].change_size.unwrap_or(-1.0);
1814            size(right).total_cmp(&size(left))
1815        });
1816    }
1817    order
1818}
1819
1820/// Every column's measures in segment `index`, beside its comparison segment largest
1821/// move first, or alone worst first; zero on both sides is left out.
1822pub fn segment_changes(results: &DataQualityResults, index: usize) -> Vec<SegmentChange> {
1823    let Some(segment) = results.segments.get(index) else {
1824        return Vec::new();
1825    };
1826    let compared = segment
1827        .compared_with
1828        .as_ref()
1829        .and_then(|label| results.segments.iter().find(|other| &other.label == label));
1830    let mut changes = Vec::new();
1831    for column in &segment.columns {
1832        let prior = compared.and_then(|other| other.columns.iter().find(|c| c.name == column.name));
1833        for metric in CHANGE_MEASURES {
1834            let Some(now) = metric.value(column) else {
1835                continue;
1836            };
1837            let before = prior.and_then(|prior| metric.value(prior));
1838            if now == 0.0 && before.unwrap_or(0.0) == 0.0 {
1839                continue;
1840            }
1841            let clear = match (prior, before) {
1842                (Some(prior), Some(before)) => {
1843                    (now - before).abs() * 100.0 >= MATERIAL_CHANGE_PP
1844                        && (results.precision == QualityPrecision::Exact
1845                            || beyond_noise(
1846                                now,
1847                                metric.denominator(column),
1848                                before,
1849                                metric.denominator(prior),
1850                            ))
1851                }
1852                _ => false,
1853            };
1854            changes.push(SegmentChange {
1855                column: column.name.clone(),
1856                metric,
1857                before,
1858                now,
1859                clear,
1860            });
1861        }
1862    }
1863    if compared.is_some() {
1864        // What cleared the noise first, then the rest, each largest first.
1865        changes.sort_by(|left, right| {
1866            let size = |change: &SegmentChange| change.change().unwrap_or(0.0).abs();
1867            right
1868                .clear
1869                .cmp(&left.clear)
1870                .then_with(|| size(right).total_cmp(&size(left)))
1871        });
1872    } else {
1873        changes.sort_by(|left, right| right.now.total_cmp(&left.now));
1874    }
1875    changes
1876}
1877
1878#[derive(Debug, Clone)]
1879pub struct TemporalLatencyProfile {
1880    pub segment: String,
1881    pub start_role: TemporalRole,
1882    pub end_role: TemporalRole,
1883    pub start_column: String,
1884    pub end_column: String,
1885    /// Rows in the segment.
1886    pub evaluated_rows: usize,
1887    /// Rows with both endpoints present and read: what durations and their counts are out
1888    /// of (not rows less missing ones, since a row can miss both).
1889    pub paired_rows: usize,
1890    pub missing_start: usize,
1891    pub missing_end: usize,
1892    /// Text the start column's format did not read; not counted as missing.
1893    pub unparsed_start: usize,
1894    pub unparsed_end: usize,
1895    /// Durations below zero: the end before the start.
1896    pub negative_count: usize,
1897    /// Durations of exactly zero: the end at the start.
1898    pub zero_count: usize,
1899    pub p50_seconds: Option<i64>,
1900    pub p90_seconds: Option<i64>,
1901    pub p95_seconds: Option<i64>,
1902    pub p99_seconds: Option<i64>,
1903    pub max_seconds: Option<i64>,
1904    /// The threshold the breaches were counted against: `duration > threshold`,
1905    /// strictly, so a duration of exactly the threshold is not a breach.
1906    pub threshold_seconds: Option<i64>,
1907    pub above_threshold_count: Option<usize>,
1908}
1909
1910impl TemporalLatencyProfile {
1911    pub fn pair(&self) -> (TemporalRole, TemporalRole) {
1912        (self.start_role, self.end_role)
1913    }
1914
1915    /// `event to received`.
1916    pub fn label(&self) -> String {
1917        interval_label(self.pair())
1918    }
1919
1920    /// A validity period, valid from to valid to: an end before the start is a
1921    /// period that is not valid, and no end is a period still open.
1922    pub fn is_validity(&self) -> bool {
1923        self.pair() == (TemporalRole::ValidFrom, TemporalRole::ValidTo)
1924    }
1925
1926    /// How many rows `fact` counts, and out of how many. `None` for a fact this
1927    /// interval does not measure: unparsed text with no format, a threshold not set.
1928    pub fn count(&self, fact: IntervalFact, plan: &DataQualityPlan) -> Option<(usize, usize)> {
1929        let rows = self.evaluated_rows;
1930        let paired = self.paired_rows;
1931        match fact {
1932            IntervalFact::MissingStart => Some((self.missing_start, rows)),
1933            IntervalFact::MissingEnd => Some((self.missing_end, rows)),
1934            IntervalFact::UnparsedStart => plan
1935                .time_format(&self.start_column)
1936                .map(|_| (self.unparsed_start, rows)),
1937            IntervalFact::UnparsedEnd => plan
1938                .time_format(&self.end_column)
1939                .map(|_| (self.unparsed_end, rows)),
1940            IntervalFact::Negative => Some((self.negative_count, paired)),
1941            IntervalFact::Zero => Some((self.zero_count, paired)),
1942            IntervalFact::OverThreshold => self.above_threshold_count.map(|count| (count, paired)),
1943        }
1944    }
1945
1946    /// Whether this interval's segment is a value its rows can be found by, rather
1947    /// than a stretch of rows or a file.
1948    pub fn segment_opens(&self, plan: &DataQualityPlan) -> bool {
1949        let grain = plan.interval_grain(&self.start_column, &self.end_column);
1950        segment_predicate(plan, &grain, &self.segment, None).is_some()
1951    }
1952
1953    /// The rows behind `fact` in this interval's segment, as a predicate over `plan`'s
1954    /// scope; `None` for row-chunk or file segments, or an unmeasured fact. `schema`,
1955    /// when known, lets a partition segment compare in its column's type.
1956    pub fn evidence_predicate(
1957        &self,
1958        fact: IntervalFact,
1959        plan: &DataQualityPlan,
1960        schema: Option<&Schema>,
1961    ) -> Option<Expr> {
1962        self.count(fact, plan)?;
1963        let micros = || interval_micros(plan, &self.start_column, &self.end_column);
1964        let rows = match fact {
1965            IntervalFact::MissingStart => col(self.start_column.as_str()).is_null(),
1966            IntervalFact::MissingEnd => col(self.end_column.as_str()).is_null(),
1967            IntervalFact::UnparsedStart => plan.time_format(&self.start_column)?.unparsed(),
1968            IntervalFact::UnparsedEnd => plan.time_format(&self.end_column)?.unparsed(),
1969            IntervalFact::Negative => micros().lt(lit(0i64)),
1970            IntervalFact::Zero => micros().eq(lit(0i64)),
1971            IntervalFact::OverThreshold => {
1972                micros().gt(lit(self.threshold_seconds?.saturating_mul(1_000_000)))
1973            }
1974        };
1975        let grain = plan.interval_grain(&self.start_column, &self.end_column);
1976        Some(
1977            match segment_predicate(plan, &grain, &self.segment, schema)? {
1978                Some(segment) => segment.and(rows),
1979                None => rows,
1980            },
1981        )
1982    }
1983}
1984
1985/// What an interval's detail counts, each with the rows behind it.
1986#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1987pub enum IntervalFact {
1988    MissingStart,
1989    MissingEnd,
1990    UnparsedStart,
1991    UnparsedEnd,
1992    Negative,
1993    Zero,
1994    OverThreshold,
1995}
1996
1997impl IntervalFact {
1998    pub const ALL: [Self; 7] = [
1999        Self::MissingStart,
2000        Self::MissingEnd,
2001        Self::UnparsedStart,
2002        Self::UnparsedEnd,
2003        Self::Negative,
2004        Self::Zero,
2005        Self::OverThreshold,
2006    ];
2007
2008    /// The fact as its row in the detail names it. A validity period's missing end
2009    /// is an open period, and its negative duration one that ends before it starts.
2010    pub fn label(self, profile: &TemporalLatencyProfile) -> String {
2011        let validity = profile.is_validity();
2012        match self {
2013            Self::MissingStart => "Missing start".to_string(),
2014            Self::MissingEnd if validity => "Open, no end".to_string(),
2015            Self::MissingEnd => "Missing end".to_string(),
2016            Self::UnparsedStart => "Unparsed start".to_string(),
2017            Self::UnparsedEnd => "Unparsed end".to_string(),
2018            Self::Negative if validity => "Ends first".to_string(),
2019            Self::Negative => "Negative".to_string(),
2020            Self::Zero => "Zero".to_string(),
2021            Self::OverThreshold => format!(
2022                "Over {}",
2023                crate::analysis::analysis_modal::threshold_label(profile.threshold_seconds)
2024            ),
2025        }
2026    }
2027
2028    /// Short words for a list of rows: a view's label.
2029    pub fn short(self) -> &'static str {
2030        match self {
2031            Self::MissingStart => "missing start",
2032            Self::MissingEnd => "missing end",
2033            Self::UnparsedStart => "unparsed start",
2034            Self::UnparsedEnd => "unparsed end",
2035            Self::Negative => "negative",
2036            Self::Zero => "zero",
2037            Self::OverThreshold => "over threshold",
2038        }
2039    }
2040}
2041
2042/// Each value of `column` written as a segment label writes it (null stays null): a
2043/// cast to text formats floats and datetimes differently.
2044fn label_text(column: &str) -> Expr {
2045    col(column).map(
2046        |values| {
2047            let text = (0..values.len())
2048                .map(|row| {
2049                    let value = values.get(row)?;
2050                    Ok((!value.is_null()).then(|| crate::exact::str_value(&value).into_owned()))
2051                })
2052                .collect::<PolarsResult<StringChunked>>()?;
2053            Ok(text.with_name(values.name().clone()).into_column())
2054        },
2055        |_, field| Ok(Field::new(field.name().clone(), DataType::String)),
2056    )
2057}
2058
2059/// The rows of a partition segment labeled `value`: a plain comparison (answerable by
2060/// file statistics) where the type writes each value one way and `value` reads back
2061/// as that label, else each row's label (e.g. rounded floats).
2062fn partition_label_predicate(column: &str, value: &str, schema: Option<&Schema>) -> Expr {
2063    let writes_each_once = |dtype: &&DataType| {
2064        dtype.is_integer()
2065            || matches!(
2066                dtype,
2067                DataType::String
2068                    | DataType::Boolean
2069                    | DataType::Date
2070                    | DataType::Decimal(..)
2071                    | DataType::Categorical(..)
2072                    | DataType::Enum(..)
2073            )
2074    };
2075    let native = schema
2076        .and_then(|schema| schema.get(column))
2077        .filter(writes_each_once)
2078        .and_then(|dtype| crate::typed_value::parse(value, dtype).ok())
2079        .filter(|scalar| crate::exact::str_value(scalar.value()) == value);
2080    match native {
2081        Some(scalar) => col(column).eq(lit(scalar)),
2082        None => label_text(column).eq(lit(value.to_string())),
2083    }
2084}
2085
2086/// The rows of the segment labeled `label` under `grain` as a predicate: `Some(None)`
2087/// for the whole scope, `None` for row stretches or files. Read back from the label,
2088/// which names a partition value as keyed and a window's exact start.
2089fn segment_predicate(
2090    plan: &DataQualityPlan,
2091    grain: &QualityGrain,
2092    label: &str,
2093    schema: Option<&Schema>,
2094) -> Option<Option<Expr>> {
2095    match grain {
2096        QualityGrain::Dataset => Some(None),
2097        QualityGrain::Partition(column) => {
2098            let value = label.strip_prefix(&format!("{column}="))?;
2099            Some(Some(if value == "∅" {
2100                col(column.as_str()).is_null()
2101            } else {
2102                partition_label_predicate(column, value, schema)
2103            }))
2104        }
2105        QualityGrain::TimeWindows { column, every } => {
2106            let value = plan.time_value(column);
2107            // The rows in no window: nulls, and dates past the calendar.
2108            if label == time_window_label(column, every, None) {
2109                return Some(Some(time_window_start(value, every).is_null()));
2110            }
2111            let date = |text: &str| chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d").ok();
2112            let start = match every.as_str() {
2113                "1h" => chrono::NaiveDateTime::parse_from_str(label, "%Y-%m-%d %H:%M").ok()?,
2114                "1d" => date(label)?.and_hms_opt(0, 0, 0)?,
2115                "1w" => date(label.strip_prefix("week of ")?)?.and_hms_opt(0, 0, 0)?,
2116                "1mo" => date(&format!("{label}-01"))?.and_hms_opt(0, 0, 0)?,
2117                _ => return None,
2118            };
2119            Some(Some(
2120                time_window_start(value, every).eq(lit(start.and_utc().timestamp_micros())
2121                    .cast(DataType::Datetime(TimeUnit::Microseconds, None))),
2122            ))
2123        }
2124        QualityGrain::RowChunks(_) | QualityGrain::File => None,
2125    }
2126}
2127/// Columns null the same number of times, and the rows null in all at once: when equal
2128/// they go missing together, one fact rather than one per column.
2129#[derive(Debug, Clone, PartialEq, Eq)]
2130pub struct SharedNulls {
2131    pub columns: Vec<String>,
2132    pub null_rows: usize,
2133    pub rows_null_in_all: usize,
2134}
2135
2136impl SharedNulls {
2137    pub fn same_rows(&self) -> bool {
2138        self.rows_null_in_all == self.null_rows
2139    }
2140}
2141
2142#[derive(Debug, Clone)]
2143pub struct DataQualityResults {
2144    pub total_rows: Option<usize>,
2145    pub evaluated_rows: usize,
2146    pub precision: QualityPrecision,
2147    pub columns: Vec<ColumnQualityProfile>,
2148    pub observations: Vec<QualityObservation>,
2149    pub segments: Vec<SegmentQualityProfile>,
2150    pub temporal: Vec<TemporalLatencyProfile>,
2151    pub identity: Option<IdentityProfile>,
2152    pub category_variants: Vec<CategoryVariantGroup>,
2153    pub shared_nulls: Vec<SharedNulls>,
2154    /// How many source files' footers were compared, when the scope has files to
2155    /// compare. `None` means the checks that compare files could not run.
2156    pub source_files: Option<usize>,
2157    /// Rows an equal-per-value sample kept of each value. See [`crate::analysis::sampling::PerValue`].
2158    pub per_value: Option<usize>,
2159    /// How many of `source_files` had their footers read: fewer on a dataset too
2160    /// large to read every footer, where the file checks cover only those.
2161    pub footers_read: Option<usize>,
2162    /// What the run's reads of the source were seen to do. `None` for results no
2163    /// watched run produced.
2164    pub reads: Option<ObservedReads>,
2165    /// Values behind text findings, from the rows the run kept; empty after a full
2166    /// scan.
2167    pub examples: Vec<FindingExamples>,
2168    /// Segments a sampled run counted rows in but drew none from, in order; not in
2169    /// `segments`, which profile only what was read.
2170    pub unsampled_segments: Vec<UnsampledSegment>,
2171    /// What the declared column intent found; `None` when nothing was declared.
2172    pub intent: Option<Box<crate::analysis::quality_intent::IntentResults>>,
2173    /// What the rows were read from, as the run that measured them labeled it.
2174    pub source: Option<Box<crate::analysis::quality_export::SourceIdentity>>,
2175    /// The report and checks read from the rest, once built. Edits in place go
2176    /// through [`Self::edit`], which drops them.
2177    pub(crate) derived: crate::analysis::quality_report::ReportCache,
2178}
2179
2180/// A segment the scope has rows in and a sample drew none of.
2181#[derive(Debug, Clone, PartialEq, Eq)]
2182pub struct UnsampledSegment {
2183    pub label: String,
2184    /// The rows the scope holds in it, by exact count.
2185    pub total_rows: usize,
2186}
2187
2188impl DataQualityResults {
2189    /// Approximate memory the report holds, for budgeting: profiles per column and per
2190    /// segment column, plus kept text.
2191    pub fn estimated_bytes(&self) -> usize {
2192        let profile = |column: &ColumnQualityProfile| {
2193            std::mem::size_of::<ColumnQualityProfile>()
2194                + column.name.len()
2195                + column.min.as_ref().map_or(0, String::len)
2196                + column.max.as_ref().map_or(0, String::len)
2197                + column.dominant_value.as_ref().map_or(0, String::len)
2198        };
2199        let segments = self
2200            .segments
2201            .iter()
2202            .map(|segment| {
2203                std::mem::size_of::<SegmentQualityProfile>()
2204                    + segment.label.len()
2205                    + segment.columns.iter().map(profile).sum::<usize>()
2206            })
2207            .sum::<usize>();
2208        let observations = self
2209            .observations
2210            .iter()
2211            .map(|observation| {
2212                std::mem::size_of::<QualityObservation>()
2213                    + observation.fact.len()
2214                    + observation.column.len()
2215                    + observation.normalized_category.as_ref().map_or(0, String::len)
2216                    // A footer finding names every file it applies to.
2217                    + observation
2218                        .files
2219                        .iter()
2220                        .map(|file| {
2221                            std::mem::size_of::<QualityFileEvidence>()
2222                                + file.name.len()
2223                                + file.stored_type.as_ref().map_or(0, String::len)
2224                                + file.examples.iter().map(String::len).sum::<usize>()
2225                        })
2226                        .sum::<usize>()
2227            })
2228            .sum::<usize>();
2229        let unsampled = self
2230            .unsampled_segments
2231            .iter()
2232            .map(|segment| std::mem::size_of::<UnsampledSegment>() + segment.label.len())
2233            .sum::<usize>();
2234        let texts = |values: &[String]| {
2235            values
2236                .iter()
2237                .map(|value| std::mem::size_of::<String>() + value.len())
2238                .sum::<usize>()
2239        };
2240        // Spellings are whole values, as wide as the column's text is.
2241        let spellings = self
2242            .category_variants
2243            .iter()
2244            .map(|group| {
2245                std::mem::size_of::<CategoryVariantGroup>()
2246                    + group.column.len()
2247                    + group.normalized.len()
2248                    + group
2249                        .variants
2250                        .iter()
2251                        .map(|(variant, _)| std::mem::size_of::<(String, usize)>() + variant.len())
2252                        .sum::<usize>()
2253            })
2254            .sum::<usize>();
2255        let examples = self
2256            .examples
2257            .iter()
2258            .map(|found| std::mem::size_of::<FindingExamples>() + texts(&found.values))
2259            .sum::<usize>()
2260            + self.identity.as_ref().map_or(0, |identity| {
2261                identity
2262                    .examples
2263                    .iter()
2264                    .map(|example| std::mem::size_of::<DuplicateExample>() + texts(&example.values))
2265                    .sum()
2266            });
2267        let temporal = self
2268            .temporal
2269            .iter()
2270            .map(|latency| {
2271                std::mem::size_of::<TemporalLatencyProfile>()
2272                    + latency.segment.len()
2273                    + latency.start_column.len()
2274                    + latency.end_column.len()
2275            })
2276            .sum::<usize>();
2277        let shared = self
2278            .shared_nulls
2279            .iter()
2280            .map(|shared| std::mem::size_of::<SharedNulls>() + texts(&shared.columns))
2281            .sum::<usize>();
2282        // Declared intent keeps whole values: the extremes and the commonest misfits.
2283        let intent = self.intent.as_ref().map_or(0, |intent| {
2284            let counted = |values: &[(String, usize)]| {
2285                values
2286                    .iter()
2287                    .map(|(value, _)| std::mem::size_of::<(String, usize)>() + value.len())
2288                    .sum::<usize>()
2289            };
2290            std::mem::size_of::<crate::analysis::quality_intent::IntentResults>()
2291                + intent
2292                    .columns
2293                    .iter()
2294                    .map(|check| {
2295                        std::mem::size_of::<crate::analysis::quality_intent::ColumnCheck>()
2296                            + check.lowest.as_ref().map_or(0, String::len)
2297                            + check.highest.as_ref().map_or(0, String::len)
2298                            + counted(&check.outside_examples)
2299                            + counted(&check.unparsed_examples)
2300                    })
2301                    .sum::<usize>()
2302        });
2303        std::mem::size_of::<Self>()
2304            + self.columns.iter().map(profile).sum::<usize>()
2305            + segments
2306            + unsampled
2307            + observations
2308            + temporal
2309            + spellings
2310            + examples
2311            + shared
2312            + intent
2313    }
2314
2315    pub fn compare_segments(&mut self, plan: &DataQualityPlan) {
2316        apply_comparisons(
2317            &mut self.segments,
2318            plan.comparison,
2319            plan.baseline_segment.as_deref(),
2320            self.precision,
2321        );
2322    }
2323
2324    pub fn empty(total_rows: Option<usize>, schema: &Schema) -> Self {
2325        Self {
2326            total_rows,
2327            evaluated_rows: 0,
2328            precision: QualityPrecision::Metadata,
2329            columns: schema
2330                .iter()
2331                .map(|(name, dtype)| ColumnQualityProfile::unmeasured(name, dtype.clone(), 0))
2332                .collect(),
2333            observations: Vec::new(),
2334            segments: Vec::new(),
2335            temporal: Vec::new(),
2336            identity: None,
2337            category_variants: Vec::new(),
2338            shared_nulls: Vec::new(),
2339            source_files: None,
2340            per_value: None,
2341            footers_read: None,
2342            reads: None,
2343            examples: Vec::new(),
2344            unsampled_segments: Vec::new(),
2345            intent: None,
2346            source: None,
2347            derived: Default::default(),
2348        }
2349    }
2350
2351    /// The kept examples of `kind` in `column`.
2352    pub fn examples_of(&self, kind: ObservationKind, column: &str) -> &[String] {
2353        self.examples
2354            .iter()
2355            .find(|examples| examples.kind == kind && examples.column == column)
2356            .map(|examples| examples.values.as_slice())
2357            .unwrap_or_default()
2358    }
2359}
2360
2361/// The rows a sampled run read, kept beside its results. Keyed by acquisition (dataset,
2362/// view, scope, method, size, seed); other plan settings are the report's, so a run
2363/// changing only those re-cuts these rows instead of reading the source. Every column
2364/// and each row's position are kept.
2365#[derive(Debug, Clone)]
2366pub struct QualitySample {
2367    df: DataFrame,
2368    /// Where each row sat in the scope, in the order of `df`.
2369    positions: Vec<IdxSize>,
2370    precision: QualityPrecision,
2371    total_rows: Option<usize>,
2372    per_value: Option<crate::analysis::sampling::PerValue>,
2373    /// Rows of the whole scope by segment key, per grain and text-as-time format, keyed as
2374    /// `AnyValue::str_value` reads (`None` for null).
2375    counted: Vec<(SegmentKey, SegmentCounts)>,
2376    /// Grains whose count stopped at [`crate::analysis::sampling::MAX_COUNTED_KEYS`], so a run
2377    /// of one again says so rather than reading to find out.
2378    too_many: Vec<SegmentKey>,
2379}
2380
2381/// Rows by segment key, as a count read them.
2382type SegmentCounts = BTreeMap<Option<String>, usize>;
2383
2384/// What decides a segment count: the grain, and how its column was read as time.
2385type SegmentKey = (QualityGrain, Option<TimeInterpretation>);
2386
2387fn segment_key(plan: &DataQualityPlan) -> SegmentKey {
2388    let format = match &plan.grain {
2389        QualityGrain::TimeWindows { column, .. } => plan.time_format(column).cloned(),
2390        _ => None,
2391    };
2392    (plan.grain.clone(), format)
2393}
2394
2395/// How a remote full scan gets its rows, as Setup says: one fetch into a local copy
2396/// all passes read, an earlier copy, or a source pass per check.
2397#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2398pub enum CopyPlan {
2399    /// Not a full scan of a remote source read in place.
2400    #[default]
2401    NotApplicable,
2402    /// Every pass reads the source, for the reason given.
2403    Passes(NoCopy),
2404    /// The objects are fetched once into the cache directory first.
2405    Fetch { bytes: u64, objects: usize },
2406    /// A copy fetched earlier this session serves every pass.
2407    Kept { bytes: u64, objects: usize },
2408}
2409
2410/// Why a remote full scan reads the source in each pass instead of a local copy.
2411#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2412pub enum NoCopy {
2413    /// `analysis.quality_local_copy` is 0.
2414    Off,
2415    /// The open did not learn every object's size.
2416    SizeUnknown,
2417    /// A copy fetched this session did not read as the source.
2418    Unusable,
2419    /// The scope reads only some of the rows or columns: its passes may read less
2420    /// than the whole objects a copy would fetch.
2421    PartOfTheSource,
2422    /// Larger than `analysis.quality_local_copy`.
2423    TooLarge { bytes: u64, limit: u64 },
2424    /// More than the cache directory has free, or its free space is unknown.
2425    NoRoom { bytes: u64, free: Option<u64> },
2426}
2427
2428/// Where a run's exact segment totals come from, as Setup says before Run.
2429#[derive(Debug, Clone, PartialEq, Eq, Default)]
2430pub enum SegmentCount {
2431    /// Nothing to count: the grain's sizes are known (files, row chunks, the whole
2432    /// scope), the run reads every row, or it reads no values.
2433    #[default]
2434    NotNeeded,
2435    /// An equal-per-value sample by the grain's column counts every value as it reads.
2436    PerValue,
2437    /// The pass that reads the sample counts the grain's key as it streams.
2438    InSamplePass,
2439    /// Counted by an earlier run of the rows being reused.
2440    Retained,
2441    /// Summed from a finer window's count of the same column, which it nests in
2442    /// exactly. Holds the finer width.
2443    RolledUp(String),
2444    /// A read of the grain's column of its own, after the sample.
2445    CountPass,
2446    /// The grain had more keys than a count holds: a coarser grain is needed.
2447    TooMany,
2448}
2449
2450impl SegmentCount {
2451    /// Whether the count reads the source in a pass of its own.
2452    pub fn reads(&self) -> bool {
2453        *self == Self::CountPass
2454    }
2455}
2456
2457/// A window width as a cadence: `1d` is daily.
2458pub fn window_cadence(every: &str) -> &str {
2459    match every {
2460        "1h" => "hourly",
2461        "1d" => "daily",
2462        "1w" => "weekly",
2463        "1mo" => "monthly",
2464        other => other,
2465    }
2466}
2467
2468/// Whether windows of width `fine` nest exactly in `coarse`: hours in a day, days in a
2469/// week (Monday start) and in a month; weeks not in months. Windows are cut on the
2470/// zoneless stored clock (UTC for zoned; see `time_window_start`), so no day has
2471/// 23 or 25 hours and finer counts sum to the coarser.
2472pub fn window_nests(fine: &str, coarse: &str) -> bool {
2473    matches!(
2474        (fine, coarse),
2475        ("1h", "1d" | "1w" | "1mo") | ("1d", "1w" | "1mo")
2476    )
2477}
2478
2479impl QualitySample {
2480    /// Where a run of `plan` over these rows gets its segment totals.
2481    pub fn segment_count(&self, plan: &DataQualityPlan) -> SegmentCount {
2482        if self.precision != QualityPrecision::Sampled || !segments_need_count(plan) {
2483            return SegmentCount::NotNeeded;
2484        }
2485        if per_value_counts(plan, self.per_value.as_ref()) {
2486            return SegmentCount::PerValue;
2487        }
2488        let key = segment_key(plan);
2489        if self.too_many.contains(&key) {
2490            return SegmentCount::TooMany;
2491        }
2492        if self.counted.iter().any(|(counted, _)| *counted == key) {
2493            return SegmentCount::Retained;
2494        }
2495        match self.finer_count(&key) {
2496            Some(((QualityGrain::TimeWindows { every, .. }, _), _)) => {
2497                SegmentCount::RolledUp(every.clone())
2498            }
2499            _ => SegmentCount::CountPass,
2500        }
2501    }
2502
2503    /// A count of a finer window of the same column, read the same way, that `key`'s
2504    /// windows nest in exactly.
2505    fn finer_count(&self, key: &SegmentKey) -> Option<&(SegmentKey, SegmentCounts)> {
2506        let (QualityGrain::TimeWindows { column, every }, format) = key else {
2507            return None;
2508        };
2509        self.counted.iter().find(|((grain, counted_format), _)| {
2510            matches!(
2511                grain,
2512                QualityGrain::TimeWindows { column: counted, every: fine }
2513                    if counted == column && window_nests(fine, every)
2514            ) && counted_format == format
2515        })
2516    }
2517
2518    /// The rows themselves, as the sample every tool reads.
2519    pub fn df(&self) -> &DataFrame {
2520        &self.df
2521    }
2522
2523    /// Memory the rows, their positions and their counts hold, near enough to budget
2524    /// by.
2525    pub fn estimated_bytes(&self) -> usize {
2526        let counts = self
2527            .counted
2528            .iter()
2529            .flat_map(|(_, counts)| counts.keys())
2530            .map(|key| key.as_ref().map_or(0, String::len) + 64)
2531            .sum::<usize>();
2532        let per_value = self.per_value.as_ref().map_or(0, |per_value| {
2533            per_value
2534                .totals
2535                .keys()
2536                .map(|key| key.as_ref().map_or(0, String::len) + 64)
2537                .sum()
2538        });
2539        self.df.estimated_size()
2540            + self.positions.len() * std::mem::size_of::<IdxSize>()
2541            + counts
2542            + per_value
2543    }
2544
2545    /// `df`, cut from these rows, described as the sampler described them.
2546    pub fn analysis_rows(&self, df: DataFrame) -> crate::analysis::sampling::AnalysisRows {
2547        crate::analysis::sampling::AnalysisRows {
2548            sample_size: (self.precision == QualityPrecision::Sampled).then_some(df.height()),
2549            total_rows: self.total_rows.unwrap_or(df.height()),
2550            per_value: self.per_value.clone(),
2551            df,
2552        }
2553    }
2554}
2555
2556/// Whether a sampled run of `plan` counts its segments' rows: partitions and time
2557/// windows are counted for exact totals; files and row chunks are known without it.
2558pub fn segments_need_count(plan: &DataQualityPlan) -> bool {
2559    matches!(
2560        plan.grain,
2561        QualityGrain::Partition(_) | QualityGrain::TimeWindows { .. }
2562    )
2563}
2564
2565/// Whether sampling `plan`'s rows also counts its segments: an equal-per-value sample
2566/// by the grain's column counts every value as it streams.
2567pub fn sampler_counts_segments(plan: &DataQualityPlan) -> bool {
2568    matches!(
2569        (&plan.grain, &plan.method),
2570        (
2571            QualityGrain::Partition(column),
2572            crate::analysis::sampling::SampleMethod::PerPartition { column: sampled },
2573        ) if column == sampled
2574    )
2575}
2576
2577/// Where a run reading a new sample gets segment totals. Seeded runs of one file
2578/// (`may_read_blocks`) and the head see too few rows; any other sample is one streamed
2579/// pass that counts the grain's key as it goes.
2580pub fn fresh_segment_count(plan: &DataQualityPlan, may_read_blocks: bool) -> SegmentCount {
2581    if plan.compute != QualityCompute::Sample || !segments_need_count(plan) {
2582        return SegmentCount::NotNeeded;
2583    }
2584    if sampler_counts_segments(plan) {
2585        return SegmentCount::PerValue;
2586    }
2587    match plan.method {
2588        crate::analysis::sampling::SampleMethod::FirstRows => SegmentCount::CountPass,
2589        crate::analysis::sampling::SampleMethod::Spread if may_read_blocks => {
2590            SegmentCount::CountPass
2591        }
2592        _ => SegmentCount::InSamplePass,
2593    }
2594}
2595
2596/// The key a partition or time-window grain splits rows by, as both the count and
2597/// the segments read it.
2598fn segment_count_key(plan: &DataQualityPlan) -> Option<Expr> {
2599    match &plan.grain {
2600        QualityGrain::Partition(column) => Some(col(column.as_str())),
2601        QualityGrain::TimeWindows { column, every } => {
2602            Some(time_window_start(plan.time_value(column), every))
2603        }
2604        _ => None,
2605    }
2606}
2607
2608/// Whether a sample's own counts are `plan`'s segment totals: the sampler counted
2609/// them, and the sample kept what it counted.
2610fn per_value_counts(
2611    plan: &DataQualityPlan,
2612    per_value: Option<&crate::analysis::sampling::PerValue>,
2613) -> bool {
2614    sampler_counts_segments(plan) && per_value.is_some()
2615}
2616
2617#[cfg(test)]
2618pub(crate) fn compute_data_quality(
2619    lf: &LazyFrame,
2620    total_rows: Option<usize>,
2621    plan: &DataQualityPlan,
2622    source: Option<&QualitySourceContext>,
2623    polars_streaming: bool,
2624) -> Result<DataQualityResults> {
2625    compute_data_quality_kept(lf, total_rows, plan, source, polars_streaming, None)
2626        .map(|(results, _)| results)
2627}
2628
2629/// [`compute_data_quality`], cutting `kept` instead of reading when it serves the
2630/// plan, and returning the sample a sampled run read so the next run can do the same.
2631#[cfg(test)]
2632pub(crate) fn compute_data_quality_kept(
2633    lf: &LazyFrame,
2634    total_rows: Option<usize>,
2635    plan: &DataQualityPlan,
2636    source: Option<&QualitySourceContext>,
2637    polars_streaming: bool,
2638    kept: Option<&QualitySample>,
2639) -> Result<(DataQualityResults, Option<QualitySample>)> {
2640    let (results, kept) = compute_data_quality_watched(
2641        lf,
2642        total_rows,
2643        plan,
2644        source,
2645        polars_streaming,
2646        kept,
2647        &QualityWatch::default(),
2648    );
2649    results.map(|results| (results, kept))
2650}
2651
2652/// The data quality of `lf` by `plan`, cutting `kept` instead of reading when it
2653/// serves the plan. Names each stage to `watch` and stops between stages (or inside a
2654/// streamed read) once cancelled. A sampled run's rows come back even if it stopped:
2655/// the read is paid for, and the next run can cut it.
2656pub fn compute_data_quality_watched(
2657    lf: &LazyFrame,
2658    total_rows: Option<usize>,
2659    plan: &DataQualityPlan,
2660    source: Option<&QualitySourceContext>,
2661    polars_streaming: bool,
2662    kept: Option<&QualitySample>,
2663    watch: &QualityWatch,
2664) -> (Result<DataQualityResults>, Option<QualitySample>) {
2665    let mut acquired = None;
2666    let inputs = QualityInputs {
2667        lf,
2668        total_rows,
2669        plan,
2670        source,
2671        // Without the feature `collect_lazy` is one in-memory collect whatever the
2672        // setting, with no batch boundary for a cancel to stop at (#498).
2673        polars_streaming: polars_streaming && cfg!(feature = "streaming"),
2674        watch,
2675    };
2676    let results = profile_quality(inputs, kept, &mut acquired);
2677    (results, acquired)
2678}
2679
2680/// What a run is asked to measure, and how it reports.
2681#[derive(Clone, Copy)]
2682struct QualityInputs<'a> {
2683    lf: &'a LazyFrame,
2684    total_rows: Option<usize>,
2685    plan: &'a DataQualityPlan,
2686    source: Option<&'a QualitySourceContext>,
2687    polars_streaming: bool,
2688    watch: &'a QualityWatch,
2689}
2690
2691/// The run itself. A sampled run's rows go into `acquired` the moment they are
2692/// read, so they outlive a run that stops after.
2693fn profile_quality(
2694    inputs: QualityInputs<'_>,
2695    kept: Option<&QualitySample>,
2696    acquired: &mut Option<QualitySample>,
2697) -> Result<DataQualityResults> {
2698    let QualityInputs {
2699        lf,
2700        total_rows,
2701        plan,
2702        source,
2703        polars_streaming,
2704        watch,
2705    } = inputs;
2706    watch.stage(QualityStage::Preparing, false, false)?;
2707    let collected_schema = lf.clone().collect_schema()?;
2708    let schema = visible_schema(&collected_schema, source);
2709    // What the footers already said: which files have which columns. Free at every
2710    // compute budget, including the one that reads no values at all.
2711    if plan.compute == QualityCompute::Metadata {
2712        watch.stage(QualityStage::Assembling, false, false)?;
2713        let mut results = DataQualityResults::empty(total_rows, &schema);
2714        if let Some(source) = source {
2715            results.observations = drift_observations(source, None, polars_streaming, watch);
2716        }
2717        results.source_files = source.map(|source| source.file_names.len());
2718        results.footers_read = source.map(|source| source.footers_read);
2719        results.reads = Some(watch.observed());
2720        results.intent =
2721            crate::analysis::quality_intent::IntentResults::unmeasured(plan, &schema).map(Box::new);
2722        return Ok(results);
2723    }
2724    let grain_column = match &plan.grain {
2725        QualityGrain::Partition(column) | QualityGrain::TimeWindows { column, .. } => Some(column),
2726        _ => None,
2727    };
2728    if let Some(column) = grain_column
2729        && collected_schema.get(column).is_none()
2730    {
2731        return Err(Report::msg(format!(
2732            "Grain column {column} is not in scope {}; choose another grain or scope",
2733            plan.scope.label()
2734        )));
2735    }
2736    if let QualityGrain::TimeWindows { column, .. } = &plan.grain
2737        && !plan.reads_as_time(column, &collected_schema)
2738    {
2739        return Err(Report::msg(format!(
2740            "Grain column {column} is text; choose a format for it under Text as time"
2741        )));
2742    }
2743    if plan.compute == QualityCompute::Full {
2744        let total_rows = match total_rows {
2745            Some(rows) => rows,
2746            None => {
2747                // Unwatched: Parquet and IPC answer a count from their metadata, which
2748                // a watch between the count and the scan would turn into a read.
2749                watch.stage(QualityStage::CountingRows, watch.scope_reads(true), false)?;
2750                let count = collect_lazy(crate::table::row_count_lf(lf), polars_streaming)
2751                    .map_err(Report::from)?;
2752                let count_values = count
2753                    .get(0)
2754                    .ok_or_else(|| Report::msg("Data quality row count was not returned"))?;
2755                let Some(AnyValue::UInt64(rows)) = count_values.first() else {
2756                    return Err(Report::msg("Data quality row count was not UInt64"));
2757                };
2758                *rows as usize
2759            }
2760        };
2761        if total_rows == 0 && plan.scope != QualityScope::CurrentView {
2762            return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
2763        }
2764        return compute_full_quality(
2765            lf,
2766            total_rows,
2767            plan,
2768            source,
2769            &schema,
2770            polars_streaming,
2771            watch,
2772        );
2773    }
2774
2775    // The shared sampler, spread across the scope by default (a date-sorted file is not
2776    // judged by its start); every grain cuts segments from this one sample.
2777    let kept = acquired.insert(match kept {
2778        Some(kept) => {
2779            watch.stage(QualityStage::ReusingSample, false, false)?;
2780            kept.clone()
2781        }
2782        None => {
2783            // First rows are one collect; other methods stream in batches or read seeded runs,
2784            // stopping between them (without the streaming engine, batches follow the whole
2785            // read).
2786            let interruptible = plan.method != crate::analysis::sampling::SampleMethod::FirstRows
2787                && cfg!(feature = "streaming");
2788            watch.stage(QualityStage::ReadingSample, true, interruptible)?;
2789            read_quality_sample(lf, total_rows, plan, polars_streaming, watch)?
2790        }
2791    });
2792    let profile_df = kept.df.clone();
2793    let sample_positions = kept.positions.clone();
2794    let evaluated_rows = profile_df.height();
2795    let precision = kept.precision;
2796    let total_rows = kept.total_rows;
2797
2798    // Rows chosen by the sample that match nothing are a mistake to name, not an
2799    // empty report that reads as clean.
2800    if total_rows == Some(0) && plan.scope != QualityScope::CurrentView {
2801        return Err(crate::analysis::sampling::no_rows_error(&plan.scope));
2802    }
2803    let profile_df = attach_source_file(profile_df, source)?;
2804    watch.stage(QualityStage::ProfilingColumns, false, false)?;
2805    let mut columns = profile_columns(&profile_df, &schema, polars_streaming)?;
2806    // The full scan's Polars aggregations over the kept rows, so sample and scan are
2807    // measured alike and any sample size scales.
2808    let profile_lf = profile_df.clone().lazy();
2809    add_dominance_lazy(&profile_lf, &mut columns, polars_streaming)?;
2810    // Text read as time and the declared intent are counted over the rows in memory,
2811    // as the columns were.
2812    let mut formats = interpretation_exprs(plan, &collected_schema);
2813    formats.extend(crate::analysis::quality_intent::intent_exprs(plan, &schema));
2814    let unparsed = if formats.is_empty() {
2815        DataFrame::default()
2816    } else {
2817        collect_lazy(profile_lf.clone().select(formats), polars_streaming).map_err(Report::from)?
2818    };
2819    watch.stage(QualityStage::CheckingDuplicates, false, false)?;
2820    let identity = profile_identity_lazy(&profile_lf, &schema, evaluated_rows, polars_streaming)?;
2821    // The declared key's repeats among the rows in memory: a repeat among distinct
2822    // sampled rows is a repeat in the data, and no repeat says nothing past them.
2823    let repeats =
2824        crate::analysis::quality_intent::key_repeats(&profile_lf, plan, &schema, polars_streaming)?;
2825    let intent = crate::analysis::quality_intent::IntentResults::from_counts(
2826        plan,
2827        &schema,
2828        &unparsed,
2829        repeats,
2830        evaluated_rows,
2831        precision,
2832        Some(&profile_lf),
2833    )?
2834    .map(Box::new);
2835    watch.stage(QualityStage::CheckingSpellings, false, false)?;
2836    let category_variants = profile_category_variants_lazy(&profile_lf, &schema, polars_streaming)?;
2837    let mut observations = observations_from_profiles(&columns, precision);
2838    observations.extend(interpretation_observations(
2839        &unparsed,
2840        plan,
2841        &collected_schema,
2842    ));
2843    observations.extend(identity_observations(&identity, &category_variants));
2844    if let Some(intent) = &intent {
2845        observations.extend(intent.observations());
2846    }
2847    crate::analysis::quality_intent::supersede(&mut observations, plan);
2848    // The rows are in memory, so the detail can show a few of the values behind a
2849    // finding without reading anything again.
2850    let mut identity = identity;
2851    if identity.duplicate_groups > 0 {
2852        identity.examples = duplicate_examples(&profile_lf, &schema, polars_streaming)?;
2853    }
2854    let examples = finding_examples(&profile_lf, &columns, &observations, polars_streaming)?;
2855    // A sampled run does not promise the extra reads, so the counts come without the
2856    // values behind them.
2857    if let Some(source) = source {
2858        observations.extend(drift_observations(source, None, polars_streaming, watch));
2859    }
2860    let totals = {
2861        let mut totals = known_segment_totals(plan, total_rows, source);
2862        if precision == QualityPrecision::Sampled {
2863            totals.extend(sampled_segment_totals(
2864                lf,
2865                plan,
2866                kept,
2867                polars_streaming,
2868                watch,
2869            )?);
2870        }
2871        totals
2872    };
2873    watch.stage(QualityStage::ProfilingSegments, false, false)?;
2874    let (segments, unsampled_segments) = profile_segments(
2875        &profile_df,
2876        total_rows,
2877        plan,
2878        precision,
2879        &schema,
2880        SegmentSampleProvenance {
2881            positions: Some(sample_positions.as_slice()),
2882            totals: &totals,
2883        },
2884        polars_streaming,
2885    )?;
2886    watch.stage(QualityStage::ComputingIntervals, false, false)?;
2887    let temporal = profile_temporal(&profile_df, plan, Some(sample_positions.as_slice()))?;
2888    watch.stage(QualityStage::CheckingSharedNulls, false, false)?;
2889    let shared_nulls = profile_shared_nulls(&profile_df.lazy(), &columns, polars_streaming)?;
2890    let per_value = kept.per_value.as_ref().map(|per_value| per_value.kept);
2891    watch.stage(QualityStage::Assembling, false, false)?;
2892
2893    let results = DataQualityResults {
2894        total_rows,
2895        evaluated_rows,
2896        precision,
2897        columns,
2898        observations,
2899        segments,
2900        temporal,
2901        identity: Some(identity),
2902        category_variants,
2903        shared_nulls,
2904        source_files: source.map(|source| source.file_names.len()),
2905        per_value,
2906        footers_read: source.map(|source| source.footers_read),
2907        reads: Some(watch.observed()),
2908        examples,
2909        unsampled_segments,
2910        intent,
2911        source: None,
2912        derived: Default::default(),
2913    };
2914    Ok(results)
2915}
2916
2917/// Read the rows a sampled run measures, counting the grain's segments in the same
2918/// pass when the sampler streams every row.
2919fn read_quality_sample(
2920    lf: &LazyFrame,
2921    total_rows: Option<usize>,
2922    plan: &DataQualityPlan,
2923    polars_streaming: bool,
2924    watch: &QualityWatch,
2925) -> Result<QualitySample> {
2926    let sample = crate::analysis::sampling::Sample {
2927        scope: QualityScope::CurrentView,
2928        method: plan.method.clone(),
2929        rows: plan.dataset_rows,
2930        seed: plan.sample_seed,
2931    };
2932    let count = if sampler_counts_segments(plan) {
2933        None
2934    } else {
2935        segment_count_key(plan)
2936    };
2937    let sampled = crate::analysis::sampling::acquire(
2938        lf,
2939        &sample,
2940        total_rows,
2941        polars_streaming,
2942        Some(watch.read()),
2943        count.as_ref(),
2944    )?;
2945    let precision = if sampled.rows.sample_size.is_some() {
2946        QualityPrecision::Sampled
2947    } else {
2948        QualityPrecision::Exact
2949    };
2950    let mut kept = QualitySample {
2951        df: sampled.rows.df,
2952        positions: sampled.positions,
2953        precision,
2954        total_rows: Some(sampled.rows.total_rows),
2955        per_value: sampled.rows.per_value,
2956        counted: Vec::new(),
2957        too_many: Vec::new(),
2958    };
2959    match sampled.counted {
2960        Some(crate::analysis::sampling::Counted::Totals(totals)) => {
2961            kept.counted.push((segment_key(plan), totals));
2962        }
2963        Some(crate::analysis::sampling::Counted::TooMany) => kept.too_many.push(segment_key(plan)),
2964        None => {}
2965    }
2966    Ok(kept)
2967}
2968
2969/// Rows per segment of a sampled run, reading only what is uncounted: an
2970/// equal-per-value or streamed sample already counted its grain; a coarser window
2971/// sums a nesting finer count; anything else is counted once by a read of its key
2972/// and kept with the sample.
2973fn sampled_segment_totals(
2974    lf: &LazyFrame,
2975    plan: &DataQualityPlan,
2976    kept: &mut QualitySample,
2977    polars_streaming: bool,
2978    watch: &QualityWatch,
2979) -> Result<BTreeMap<String, usize>> {
2980    if !segments_need_count(plan) {
2981        return Ok(BTreeMap::new());
2982    }
2983    let labeled = |counts: &SegmentCounts| {
2984        counts
2985            .iter()
2986            .map(|(raw, rows)| (segment_label(&plan.grain, raw.as_deref()), *rows))
2987            .collect::<BTreeMap<_, _>>()
2988    };
2989    if per_value_counts(plan, kept.per_value.as_ref())
2990        && let Some(per_value) = &kept.per_value
2991    {
2992        return Ok(labeled(&per_value.totals));
2993    }
2994    let key = segment_key(plan);
2995    if kept.too_many.contains(&key) {
2996        return Err(too_many_segments(plan));
2997    }
2998    if let Some((_, counts)) = kept.counted.iter().find(|(counted, _)| *counted == key) {
2999        return Ok(labeled(counts));
3000    }
3001    if let Some((_, finer)) = kept.finer_count(&key)
3002        && let QualityGrain::TimeWindows { every, .. } = &plan.grain
3003    {
3004        let counts = roll_up_windows(finer, every)?;
3005        let totals = labeled(&counts);
3006        kept.counted.push((key, counts));
3007        return Ok(totals);
3008    }
3009    watch.stage(QualityStage::CountingSegments, true, polars_streaming)?;
3010    let counts = counted_segment_totals(&watch.watched(lf), plan, polars_streaming)
3011        .map_err(|error| watch.failed(error))?;
3012    if counts.len() > crate::analysis::sampling::MAX_COUNTED_KEYS {
3013        kept.too_many.push(key);
3014        return Err(too_many_segments(plan));
3015    }
3016    let totals = labeled(&counts);
3017    kept.counted.push((key, counts));
3018    Ok(totals)
3019}
3020
3021fn too_many_segments(plan: &DataQualityPlan) -> Report {
3022    Report::msg(format!(
3023        "More than {} segments {}; choose a coarser grain",
3024        crate::numfmt::group_chrome(crate::analysis::sampling::MAX_COUNTED_KEYS),
3025        plan.grain.label()
3026    ))
3027}
3028
3029/// A finer window's counts summed into `every`'s windows through the window
3030/// expression; exact only where [`window_nests`] says so (asked first).
3031fn roll_up_windows(finer: &SegmentCounts, every: &str) -> Result<SegmentCounts> {
3032    let mut rolled = SegmentCounts::new();
3033    let mut starts = Vec::with_capacity(finer.len());
3034    let mut rows = Vec::with_capacity(finer.len());
3035    for (raw, count) in finer {
3036        match raw {
3037            Some(raw) => {
3038                let start = chrono::NaiveDateTime::parse_from_str(raw, "%Y-%m-%d %H:%M:%S%.f")
3039                    .map_err(|_| Report::msg(format!("Window start {raw:?} is not a time")))?;
3040                starts.push(start.and_utc().timestamp_micros());
3041                rows.push(*count as u64);
3042            }
3043            // A row with no time is in no window at any width.
3044            None => *rolled.entry(None).or_default() += count,
3045        }
3046    }
3047    let finer = DataFrame::new(
3048        starts.len(),
3049        vec![
3050            Column::new("start".into(), starts)
3051                .cast(&DataType::Datetime(TimeUnit::Microseconds, None))?,
3052            Column::new("rows".into(), rows),
3053        ],
3054    )?;
3055    let coarse = finer
3056        .lazy()
3057        .select([time_window_start(col("start"), every), col("rows")])
3058        .collect()?;
3059    let (starts, rows) = (coarse.column("start")?, coarse.column("rows")?.u64()?);
3060    for (row, count) in rows.into_no_null_iter().enumerate() {
3061        let start = starts.get(row)?;
3062        let key = (!start.is_null()).then(|| crate::exact::str_value(&start).into_owned());
3063        *rolled.entry(key).or_default() += count as usize;
3064    }
3065    Ok(rolled)
3066}
3067
3068fn compute_full_quality(
3069    lf: &LazyFrame,
3070    total_rows: usize,
3071    plan: &DataQualityPlan,
3072    source: Option<&QualitySourceContext>,
3073    schema: &Schema,
3074    polars_streaming: bool,
3075    watch: &QualityWatch,
3076) -> Result<DataQualityResults> {
3077    let full_schema = lf.clone().collect_schema()?;
3078    // Every pass reads through the watch, so a cancel stops it within a batch on the
3079    // streaming engine, and the rows each pass traverses are counted.
3080    let lf = &watch.watched(lf);
3081    let failed = |error: Report| watch.failed(error);
3082    watch.stage(
3083        QualityStage::ProfilingColumns,
3084        watch.scope_reads(true),
3085        polars_streaming,
3086    )?;
3087    // Text read as time is counted in the same pass as every column's profile.
3088    let mut exprs = build_profile_exprs(schema);
3089    exprs.extend(interpretation_exprs(plan, &full_schema));
3090    // The declared intent's counts too: sums over the same rows, in the same pass.
3091    exprs.extend(crate::analysis::quality_intent::intent_exprs(plan, schema));
3092    let aggregate = collect_lazy(lf.clone().select(exprs), polars_streaming)
3093        .map_err(|error| watch.failed(error))?;
3094    let mut columns = parse_profiles(&aggregate, schema, total_rows);
3095    add_dominance_lazy(lf, &mut columns, polars_streaming).map_err(failed)?;
3096    watch.stage(
3097        QualityStage::CheckingDuplicates,
3098        watch.scope_reads(true),
3099        polars_streaming,
3100    )?;
3101    let identity =
3102        profile_identity_lazy(lf, schema, total_rows, polars_streaming).map_err(failed)?;
3103    let texts = schema
3104        .iter_values()
3105        .any(|dtype| matches!(dtype, DataType::String | DataType::Categorical(..)));
3106    watch.stage(
3107        QualityStage::CheckingSpellings,
3108        watch.scope_reads(texts),
3109        polars_streaming,
3110    )?;
3111    let category_variants =
3112        profile_category_variants_lazy(lf, schema, polars_streaming).map_err(failed)?;
3113    // The declared key is one grouping of its columns: a pass of its own, which
3114    // Setup counts among the passes before Run.
3115    let keyed = !plan.intent.key.is_empty();
3116    watch.stage(
3117        QualityStage::CheckingKey,
3118        watch.scope_reads(keyed),
3119        polars_streaming,
3120    )?;
3121    let repeats = crate::analysis::quality_intent::key_repeats(lf, plan, schema, polars_streaming)
3122        .map_err(failed)?;
3123    let intent = crate::analysis::quality_intent::IntentResults::from_counts(
3124        plan,
3125        schema,
3126        &aggregate,
3127        repeats,
3128        total_rows,
3129        QualityPrecision::Exact,
3130        None,
3131    )?
3132    .map(Box::new);
3133    let mut observations = observations_from_profiles(&columns, QualityPrecision::Exact);
3134    observations.extend(interpretation_observations(&aggregate, plan, &full_schema));
3135    observations.extend(identity_observations(&identity, &category_variants));
3136    if let Some(intent) = &intent {
3137        observations.extend(intent.observations());
3138    }
3139    crate::analysis::quality_intent::supersede(&mut observations, plan);
3140    // Only a run already reading every value pays for conflicting values, as its access
3141    // plan promised; files are read one by one so a cancel stops between them.
3142    if let Some(source) = source {
3143        if source.conflict_scan.is_some() {
3144            watch.stage(QualityStage::ReadingConflicts, true, true)?;
3145        }
3146        observations.extend(drift_observations(
3147            source,
3148            source.conflict_scan.as_ref(),
3149            polars_streaming,
3150            watch,
3151        ));
3152    }
3153    let whole = unsegmented(plan, source);
3154    watch.stage(
3155        QualityStage::ProfilingSegments,
3156        watch.scope_reads(!whole),
3157        polars_streaming,
3158    )?;
3159    let segments = if whole {
3160        // The whole scope is one segment, and its profile is the one just measured:
3161        // reading it again would be a second pass for the same numbers.
3162        vec![whole_segment(plan, total_rows, &columns, schema.len())]
3163    } else {
3164        profile_segments_lazy(lf, total_rows, plan, source, schema, polars_streaming)
3165            .map_err(failed)?
3166    };
3167    let intervals = !resolved_intervals(plan, &full_schema).is_empty();
3168    watch.stage(
3169        QualityStage::ComputingIntervals,
3170        watch.scope_reads(intervals),
3171        polars_streaming,
3172    )?;
3173    let temporal = profile_temporal_lazy(lf, plan, source, polars_streaming).map_err(failed)?;
3174    let shared = !shared_null_groups(&columns).is_empty();
3175    watch.stage(
3176        QualityStage::CheckingSharedNulls,
3177        watch.scope_reads(shared),
3178        polars_streaming,
3179    )?;
3180    let shared_nulls = profile_shared_nulls(lf, &columns, polars_streaming).map_err(failed)?;
3181    watch.stage(QualityStage::Assembling, false, false)?;
3182    Ok(DataQualityResults {
3183        total_rows: Some(total_rows),
3184        evaluated_rows: total_rows,
3185        precision: QualityPrecision::Exact,
3186        columns,
3187        observations,
3188        segments,
3189        temporal,
3190        identity: Some(identity),
3191        category_variants,
3192        shared_nulls,
3193        source_files: source.map(|source| source.file_names.len()),
3194        per_value: None,
3195        footers_read: source.map(|source| source.footers_read),
3196        reads: Some(watch.observed()),
3197        examples: Vec::new(),
3198        unsampled_segments: Vec::new(),
3199        intent,
3200        source: None,
3201        derived: Default::default(),
3202    })
3203}
3204
3205/// Columns sharing a nonzero null count, by the count: the sets worth checking for
3206/// rows null in all of them.
3207fn shared_null_groups(columns: &[ColumnQualityProfile]) -> Vec<(usize, Vec<String>)> {
3208    let mut by_count = BTreeMap::<usize, Vec<String>>::new();
3209    for profile in columns.iter().filter(|profile| profile.null_count > 0) {
3210        by_count
3211            .entry(profile.null_count)
3212            .or_default()
3213            .push(profile.name.clone());
3214    }
3215    by_count
3216        .into_iter()
3217        .filter(|(_, names)| names.len() > 1)
3218        .collect()
3219}
3220
3221/// For each set of two or more columns with the same nonzero null count, the rows null
3222/// in all of them: equal counts are only a hint. Reads just those columns once;
3223/// skipped when no counts are shared.
3224fn profile_shared_nulls(
3225    lf: &LazyFrame,
3226    columns: &[ColumnQualityProfile],
3227    polars_streaming: bool,
3228) -> Result<Vec<SharedNulls>> {
3229    let groups = shared_null_groups(columns);
3230    if groups.is_empty() {
3231        return Ok(Vec::new());
3232    }
3233    let exprs = groups
3234        .iter()
3235        .enumerate()
3236        .map(|(index, (_, names))| {
3237            names
3238                .iter()
3239                .map(|name| col(name.as_str()).is_null())
3240                .reduce(Expr::and)
3241                .expect("a group has two columns")
3242                .sum()
3243                .alias(format!("__quality_shared_null_{index}"))
3244        })
3245        .collect::<Vec<_>>();
3246    let counts = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
3247    Ok(groups
3248        .into_iter()
3249        .enumerate()
3250        .map(|(index, (null_rows, columns))| SharedNulls {
3251            columns,
3252            null_rows,
3253            rows_null_in_all: usize_value(&counts, &format!("__quality_shared_null_{index}")),
3254        })
3255        .collect())
3256}
3257
3258/// The most common value of every column in one pass (a pass per column would reread
3259/// a remote source per column; the access plan promises one).
3260fn add_dominance_lazy(
3261    lf: &LazyFrame,
3262    profiles: &mut [ColumnQualityProfile],
3263    polars_streaming: bool,
3264) -> Result<()> {
3265    if profiles.is_empty() {
3266        return Ok(());
3267    }
3268    const COUNT: &str = "__quality_value_count";
3269    let exprs = profiles
3270        .iter()
3271        .enumerate()
3272        .map(|(index, profile)| {
3273            col(&profile.name)
3274                .drop_nulls()
3275                .value_counts(true, true, COUNT, false)
3276                .first()
3277                .alias(format!("__quality_dominant_{index}"))
3278        })
3279        .collect::<Vec<_>>();
3280    let top = collect_lazy(lf.clone().select(exprs), polars_streaming).map_err(Report::from)?;
3281    for (index, profile) in profiles.iter_mut().enumerate() {
3282        let Ok(column) = top.column(&format!("__quality_dominant_{index}")) else {
3283            continue;
3284        };
3285        let Ok(fields) = column.struct_() else {
3286            continue;
3287        };
3288        let Ok(value) = fields.field_by_name(&profile.name) else {
3289            continue;
3290        };
3291        let Ok(counts) = fields.field_by_name(COUNT) else {
3292            continue;
3293        };
3294        profile.dominant_value = value
3295            .get(0)
3296            .ok()
3297            .filter(|value| !value.is_null())
3298            .map(|value| crate::exact::str_value(&value).to_string());
3299        profile.dominant_count = counts
3300            .get(0)
3301            .ok()
3302            .and_then(|value| value.try_extract::<u64>().ok())
3303            .map(|count| count as usize);
3304    }
3305    Ok(())
3306}
3307
3308fn profile_category_variants_lazy(
3309    lf: &LazyFrame,
3310    schema: &Schema,
3311    polars_streaming: bool,
3312) -> Result<Vec<CategoryVariantGroup>> {
3313    let normalized_name = "__quality_normalized";
3314    let original_name = "__quality_original";
3315    let count_name = "__quality_variant_rows";
3316    let variant_count_name = "__quality_variant_count";
3317    let mut result = Vec::new();
3318    for (name, dtype) in schema.iter() {
3319        if !matches!(dtype, DataType::String | DataType::Categorical(..)) {
3320            continue;
3321        }
3322        let original = text_expr(col(name.as_str()), dtype);
3323        let normalized = original
3324            .clone()
3325            .str()
3326            .strip_chars(lit(LiteralValue::untyped_null()))
3327            .str()
3328            .to_lowercase();
3329        let variant_count = col(original_name)
3330            .n_unique()
3331            .over([col(normalized_name)])?
3332            .alias(variant_count_name);
3333        let query = lf
3334            .clone()
3335            .filter(original.clone().is_not_null())
3336            .select([
3337                normalized.alias(normalized_name),
3338                original.alias(original_name),
3339            ])
3340            .group_by([col(normalized_name), col(original_name)])
3341            .agg([len().alias(count_name)])
3342            .with_columns([variant_count])
3343            .filter(col(variant_count_name).gt(lit(1u32)))
3344            .limit(1_000);
3345        let groups = collect_lazy(query, polars_streaming).map_err(Report::from)?;
3346        let mut by_normalized = BTreeMap::<String, Vec<(String, usize)>>::new();
3347        for row in 0..groups.height() {
3348            let Some(normalized) = string_value_at(&groups, normalized_name, row) else {
3349                continue;
3350            };
3351            let Some(original) = string_value_at(&groups, original_name, row) else {
3352                continue;
3353            };
3354            let count = usize_value_at(&groups, count_name, row);
3355            by_normalized
3356                .entry(normalized)
3357                .or_default()
3358                .push((original, count));
3359        }
3360        for (normalized, variants) in by_normalized {
3361            if variants.len() < 2 {
3362                continue;
3363            }
3364            let rows_involved = variants.iter().map(|(_, count)| count).sum();
3365            result.push(CategoryVariantGroup {
3366                column: name.to_string(),
3367                normalized,
3368                variants,
3369                rows_involved,
3370            });
3371            if result.len() >= 100 {
3372                return Ok(result);
3373            }
3374        }
3375    }
3376    Ok(result)
3377}
3378
3379fn profile_identity_lazy(
3380    lf: &LazyFrame,
3381    schema: &Schema,
3382    total_rows: usize,
3383    polars_streaming: bool,
3384) -> Result<IdentityProfile> {
3385    let keys = schema
3386        .iter_names()
3387        .map(|name| col(name.as_str()))
3388        .collect::<Vec<_>>();
3389    let duplicate_count = "__quality_duplicate_count";
3390    let grouped = lf
3391        .clone()
3392        .group_by(keys)
3393        .agg([len().alias(duplicate_count)])
3394        .filter(col(duplicate_count).gt(lit(1u32)))
3395        .select([
3396            len().alias("duplicate_groups"),
3397            (col(duplicate_count) - lit(1u32)).sum().alias("extra_rows"),
3398            col(duplicate_count).sum().alias("rows_involved"),
3399        ]);
3400    let summary = collect_lazy(grouped, polars_streaming).map_err(Report::from)?;
3401    Ok(IdentityProfile {
3402        duplicate_groups: usize_value(&summary, "duplicate_groups"),
3403        extra_rows: usize_value(&summary, "extra_rows"),
3404        rows_involved: usize_value(&summary, "rows_involved"),
3405        evaluated_rows: total_rows,
3406        examples: Vec::new(),
3407    })
3408}
3409
3410const DUPLICATE_COPIES: &str = "__datui_quality_copies";
3411
3412/// Groups of rows identical in every `keys` column with their copy counts, most copies
3413/// then first seen first: the duplicate check's grouping.
3414fn duplicate_groups(lf: LazyFrame, keys: &[PlSmallStr]) -> LazyFrame {
3415    lf.group_by_stable(keys.iter().map(|key| col(key.clone())).collect::<Vec<_>>())
3416        .agg([len().alias(DUPLICATE_COPIES)])
3417        .filter(col(DUPLICATE_COPIES).gt(lit(1u32)))
3418        .sort(
3419            [DUPLICATE_COPIES],
3420            SortMultipleOptions::default()
3421                .with_order_descending(true)
3422                .with_maintain_order(true),
3423        )
3424}
3425
3426/// The rows the duplicate check counted: rows equal to another in every `keys` column,
3427/// copies together, most first, in one pass. Each group's key is its rows, repeated
3428/// by count rather than looked up again.
3429pub fn duplicate_rows(
3430    lf: LazyFrame,
3431    keys: &[PlSmallStr],
3432    polars_streaming: bool,
3433) -> Result<DataFrame> {
3434    let groups =
3435        collect_lazy(duplicate_groups(lf, keys), polars_streaming).map_err(Report::from)?;
3436    let copies = groups
3437        .column(DUPLICATE_COPIES)?
3438        .cast(&DataType::UInt64)?
3439        .u64()?
3440        .into_no_null_iter()
3441        .collect::<Vec<_>>();
3442    let mut take = Vec::with_capacity(copies.iter().sum::<u64>() as usize);
3443    for (group, copies) in copies.into_iter().enumerate() {
3444        take.extend(std::iter::repeat_n(group as IdxSize, copies as usize));
3445    }
3446    let rows = groups.drop(DUPLICATE_COPIES)?;
3447    Ok(rows.take(&IdxCa::from_vec(PlSmallStr::EMPTY, take))?)
3448}
3449
3450/// The most copied groups of rows kept in memory, rendered for the detail.
3451fn duplicate_examples(
3452    lf: &LazyFrame,
3453    schema: &Schema,
3454    polars_streaming: bool,
3455) -> Result<Vec<DuplicateExample>> {
3456    let keys = schema.iter_names().cloned().collect::<Vec<_>>();
3457    let groups = collect_lazy(
3458        duplicate_groups(lf.clone(), &keys).limit(MAX_FINDING_EXAMPLES as IdxSize),
3459        polars_streaming,
3460    )
3461    .map_err(Report::from)?;
3462    Ok((0..groups.height())
3463        .map(|row| DuplicateExample {
3464            copies: usize_value_at(&groups, DUPLICATE_COPIES, row),
3465            values: keys
3466                .iter()
3467                .map(|key| {
3468                    groups
3469                        .column(key)
3470                        .and_then(|column| column.get(row))
3471                        .map(|value| example_text(&value))
3472                        .unwrap_or_else(|_| "null".to_string())
3473                })
3474                .collect(),
3475        })
3476        .collect())
3477}
3478
3479/// A value as the detail shows it: text quoted and cut, a null named.
3480fn example_text(value: &AnyValue<'_>) -> String {
3481    match value {
3482        AnyValue::Null => "null".to_string(),
3483        AnyValue::String(text) => crate::analysis::quality_report::quoted(text, 24),
3484        AnyValue::StringOwned(text) => crate::analysis::quality_report::quoted(text, 24),
3485        other => {
3486            let text = crate::exact::str_value(other).to_string();
3487            if crate::glyphs::display_width(&text) > 24 {
3488                format!(
3489                    "{}{}",
3490                    crate::glyphs::take_columns(&text, 23),
3491                    crate::glyphs::get().ellipsis
3492                )
3493            } else {
3494                text
3495            }
3496        }
3497    }
3498}
3499
3500/// A few distinct values behind each text finding a sample can show: text its
3501/// reading does not parse, and text its time format does not read.
3502fn finding_examples(
3503    lf: &LazyFrame,
3504    columns: &[ColumnQualityProfile],
3505    observations: &[QualityObservation],
3506    polars_streaming: bool,
3507) -> Result<Vec<FindingExamples>> {
3508    let mut examples = Vec::new();
3509    for observation in observations {
3510        let profile = columns
3511            .iter()
3512            .find(|profile| profile.name == observation.column);
3513        let failed = match observation.kind {
3514            ObservationKind::ParseableText => profile.and_then(unparsed_text),
3515            ObservationKind::UnparsedTime => observation
3516                .time_format
3517                .as_ref()
3518                .map(TimeInterpretation::unparsed),
3519            _ => None,
3520        };
3521        let (Some(failed), Some(profile)) = (failed, profile) else {
3522            continue;
3523        };
3524        // The first failures, told apart here: a few hundred bound the work, and a
3525        // value repeated that often is the example anyway.
3526        let values = text_expr(col(observation.column.as_str()), &profile.dtype)
3527            .filter(failed)
3528            .head(Some(256))
3529            .alias("values");
3530        let found =
3531            collect_lazy(lf.clone().select([values]), polars_streaming).map_err(Report::from)?;
3532        let mut values = Vec::new();
3533        for value in (0..found.height()).filter_map(|row| string_value_at(&found, "values", row)) {
3534            let value = crate::analysis::quality_report::quoted(&value, 24);
3535            if !values.contains(&value) {
3536                values.push(value);
3537            }
3538            if values.len() == MAX_FINDING_EXAMPLES {
3539                break;
3540            }
3541        }
3542        if !values.is_empty() {
3543            examples.push(FindingExamples {
3544                kind: observation.kind,
3545                column: observation.column.clone(),
3546                values,
3547            });
3548        }
3549    }
3550    Ok(examples)
3551}
3552
3553fn visible_schema(schema: &Schema, source: Option<&QualitySourceContext>) -> Schema {
3554    let mut visible = Schema::with_capacity(schema.len());
3555    for (name, dtype) in schema.iter() {
3556        if source.is_some_and(|context| name.as_str() == context.row_index_column) {
3557            continue;
3558        }
3559        visible.insert(name.clone(), dtype.clone());
3560    }
3561    visible
3562}
3563
3564fn attach_source_file(
3565    mut df: DataFrame,
3566    source: Option<&QualitySourceContext>,
3567) -> Result<DataFrame> {
3568    let Some(source) = source else {
3569        return Ok(df);
3570    };
3571    let rows = df.drop_in_place(&source.row_index_column)?;
3572    let rows = rows.u32()?;
3573    let names: Vec<Option<&str>> = rows
3574        .iter()
3575        .map(|row| {
3576            let row = row? as usize;
3577            let file = source
3578                .file_starts
3579                .partition_point(|start| *start <= row)
3580                .saturating_sub(1);
3581            source.file_names.get(file).map(String::as_str)
3582        })
3583        .collect();
3584    df.with_column(Column::new(QUALITY_SOURCE_FILE_COLUMN.into(), names))?;
3585    Ok(df)
3586}
3587
3588fn profile_columns(
3589    df: &DataFrame,
3590    schema: &Schema,
3591    polars_streaming: bool,
3592) -> Result<Vec<ColumnQualityProfile>> {
3593    let aggregate = collect_lazy(
3594        df.clone().lazy().select(build_profile_exprs(schema)),
3595        polars_streaming,
3596    )
3597    .map_err(Report::from)?;
3598    Ok(parse_profiles(&aggregate, schema, df.height()))
3599}
3600
3601fn identity_observations(
3602    identity: &IdentityProfile,
3603    variants: &[CategoryVariantGroup],
3604) -> Vec<QualityObservation> {
3605    let mut observations = Vec::new();
3606    if identity.duplicate_groups > 0 {
3607        observations.push(QualityObservation {
3608            kind: ObservationKind::DuplicateRows,
3609            column: "all columns".to_string(),
3610            affected_rows: identity.rows_involved,
3611            evaluated_rows: identity.evaluated_rows,
3612            fact: String::new(),
3613            normalized_category: None,
3614            files: Vec::new(),
3615            time_format: None,
3616            full_scale: None,
3617        });
3618    }
3619    observations.extend(variants.iter().map(|group| QualityObservation {
3620        kind: ObservationKind::CategoryVariants,
3621        column: group.column.clone(),
3622        affected_rows: group.rows_involved,
3623        evaluated_rows: identity.evaluated_rows,
3624        fact: String::new(),
3625        normalized_category: Some(group.normalized.clone()),
3626        files: Vec::new(),
3627        time_format: None,
3628        full_scale: None,
3629    }));
3630    observations
3631}
3632
3633#[derive(Debug)]
3634struct SegmentRows {
3635    label: String,
3636    indices: Vec<u32>,
3637}
3638
3639fn segment_rows(
3640    df: &DataFrame,
3641    plan: &DataQualityPlan,
3642    sample_positions: Option<&[IdxSize]>,
3643) -> Result<Vec<SegmentRows>> {
3644    let all_rows = || SegmentRows {
3645        label: "current view".to_string(),
3646        indices: (0..df.height() as u32).collect(),
3647    };
3648    let groups = match &plan.grain {
3649        QualityGrain::Dataset => vec![all_rows()],
3650        QualityGrain::RowChunks(size) => {
3651            let size = (*size).max(1);
3652            let mut chunks = BTreeMap::<usize, Vec<u32>>::new();
3653            for row in 0..df.height() {
3654                let position = sample_positions
3655                    .and_then(|positions| positions.get(row))
3656                    .copied()
3657                    .unwrap_or(row as IdxSize) as usize;
3658                chunks.entry(position / size).or_default().push(row as u32);
3659            }
3660            chunks
3661                .into_iter()
3662                .map(|(chunk, indices)| SegmentRows {
3663                    label: format!(
3664                        "rows {}-{}",
3665                        chunk.saturating_mul(size) + 1,
3666                        (chunk + 1).saturating_mul(size)
3667                    ),
3668                    indices,
3669                })
3670                .collect()
3671        }
3672        QualityGrain::Partition(column) => group_by_value(df, column, &format!("{column}="))?,
3673        QualityGrain::TimeWindows { column, every } => {
3674            group_by_time_window(df, plan, column, every)?
3675        }
3676        QualityGrain::File => {
3677            if df.column(QUALITY_SOURCE_FILE_COLUMN).is_ok() {
3678                group_by_value(df, QUALITY_SOURCE_FILE_COLUMN, "file ")?
3679            } else {
3680                vec![SegmentRows {
3681                    label: "file mapping unavailable for this view".to_string(),
3682                    indices: (0..df.height() as u32).collect(),
3683                }]
3684            }
3685        }
3686    };
3687    Ok(groups)
3688}
3689
3690fn group_by_value(df: &DataFrame, column: &str, prefix: &str) -> Result<Vec<SegmentRows>> {
3691    let values = df.column(column)?;
3692    let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
3693    let mut missing = Vec::new();
3694    for row in 0..df.height() {
3695        let value = values.get(row)?;
3696        if value.is_null() {
3697            missing.push(row as u32);
3698        } else {
3699            groups
3700                .entry(format!("{prefix}{}", crate::exact::str_value(&value)))
3701                .or_default()
3702                .push(row as u32);
3703        }
3704    }
3705    let mut result: Vec<SegmentRows> = groups
3706        .into_iter()
3707        .map(|(label, indices)| SegmentRows { label, indices })
3708        .collect();
3709    // Unplaced rows follow placed ones, as on the scanned path (sorting "∅" by codepoint
3710    // would misplace it).
3711    if !missing.is_empty() {
3712        result.push(SegmentRows {
3713            label: format!("{prefix}∅"),
3714            indices: missing,
3715        });
3716    }
3717    Ok(result)
3718}
3719
3720/// Where a row's window starts: the one expression both sampled and full paths
3721/// bucket by, so weeks start alike. A date past the calendar falls in no window, like
3722/// a null (truncating it overflows).
3723fn time_window_start(value: Expr, every: &str) -> Expr {
3724    value
3725        .map(
3726            |c| {
3727                Ok(
3728                    crate::exact::calendar_without_out_of_range(c.as_materialized_series())?
3729                        .map_or(c, Column::from),
3730                )
3731            },
3732            |_, field| Ok(field.clone()),
3733        )
3734        .cast(DataType::Datetime(TimeUnit::Microseconds, None))
3735        .dt()
3736        .truncate(lit(every.to_string()))
3737}
3738
3739fn group_by_time_window(
3740    df: &DataFrame,
3741    plan: &DataQualityPlan,
3742    column: &str,
3743    every: &str,
3744) -> Result<Vec<SegmentRows>> {
3745    let starts = df
3746        .clone()
3747        .lazy()
3748        .select([time_window_start(plan.time_value(column), every).alias(QUALITY_WINDOW_START)])
3749        .collect()?;
3750    let starts = starts.column(QUALITY_WINDOW_START)?;
3751    let mut groups: BTreeMap<String, Vec<u32>> = BTreeMap::new();
3752    let mut missing = Vec::new();
3753    for row in 0..df.height() {
3754        let value = starts.get(row)?;
3755        if value.is_null() {
3756            missing.push(row as u32);
3757        } else {
3758            groups
3759                .entry(crate::exact::str_value(&value).into_owned())
3760                .or_default()
3761                .push(row as u32);
3762        }
3763    }
3764    let mut result: Vec<SegmentRows> = groups
3765        .into_iter()
3766        .map(|(start, indices)| SegmentRows {
3767            label: time_window_label(column, every, Some(&start)),
3768            indices,
3769        })
3770        .collect();
3771    if !missing.is_empty() {
3772        result.push(SegmentRows {
3773            label: time_window_label(column, every, None),
3774            indices: missing,
3775        });
3776    }
3777    Ok(result)
3778}
3779
3780/// A window by where it starts, to the precision its width needs: an hour to the
3781/// minute, a day as its date, a week as the date it starts, a month as the month.
3782pub(crate) fn time_window_label(column: &str, every: &str, start: Option<&str>) -> String {
3783    let Some(start) = start else {
3784        return format!("{column} ∅");
3785    };
3786    let prefix = |length: usize| start.get(..length).unwrap_or(start).to_string();
3787    match every {
3788        "1h" => prefix(16),
3789        "1d" => prefix(10),
3790        "1w" => format!("week of {}", prefix(10)),
3791        "1mo" => prefix(7),
3792        _ => format!("{start} / {every}"),
3793    }
3794}
3795
3796fn value_epoch_micros(value: AnyValue<'_>) -> Option<i64> {
3797    // A one-row segment's column can be a scalar, whose values come back owned.
3798    match value.as_borrowed() {
3799        AnyValue::Date(days) => Some(i64::from(days) * 86_400_000_000),
3800        AnyValue::Datetime(value, TimeUnit::Nanoseconds, _) => Some(value / 1_000),
3801        AnyValue::Datetime(value, TimeUnit::Microseconds, _) => Some(value),
3802        AnyValue::Datetime(value, TimeUnit::Milliseconds, _) => Some(value * 1_000),
3803        _ => None,
3804    }
3805}
3806
3807fn take_rows(df: &DataFrame, indices: &[u32]) -> PolarsResult<DataFrame> {
3808    df.take(&UInt32Chunked::new("quality_rows".into(), indices.to_vec()))
3809}
3810
3811struct SegmentSampleProvenance<'a> {
3812    positions: Option<&'a [IdxSize]>,
3813    totals: &'a BTreeMap<String, usize>,
3814}
3815
3816/// Segment sizes known without reading: a file's rows from its footer (whole-file
3817/// scopes) and a row chunk's from the scope size. Partitions and windows are counted
3818/// with a grouped read of the grain's column; other sizes stay unknown on a sample.
3819fn counted_segment_totals(
3820    lf: &LazyFrame,
3821    plan: &DataQualityPlan,
3822    polars_streaming: bool,
3823) -> Result<SegmentCounts> {
3824    const KEY: &str = "__quality_count_key";
3825    const ROWS: &str = "__quality_count_rows";
3826    let Some(key) = segment_count_key(plan) else {
3827        return Ok(SegmentCounts::new());
3828    };
3829    let counts = collect_lazy(
3830        lf.clone()
3831            .select([key.alias(KEY)])
3832            .group_by([col(KEY)])
3833            .agg([len().alias(ROWS)]),
3834        polars_streaming,
3835    )
3836    .map_err(Report::from)?;
3837    let keys = counts.column(KEY)?;
3838    let mut totals = BTreeMap::new();
3839    for row in 0..counts.height() {
3840        let raw = keys.get(row)?;
3841        // Keyed as the key reads, as a streamed count keys it: named as a segment only
3842        // when a run asks, so a finer window's count can be summed into a coarser one.
3843        let raw = (!raw.is_null()).then(|| crate::exact::str_value(&raw).into_owned());
3844        totals.insert(raw, usize_value_at(&counts, ROWS, row));
3845    }
3846    Ok(totals)
3847}
3848
3849fn known_segment_totals(
3850    plan: &DataQualityPlan,
3851    total_rows: Option<usize>,
3852    source: Option<&QualitySourceContext>,
3853) -> BTreeMap<String, usize> {
3854    let mut totals = BTreeMap::new();
3855    match &plan.grain {
3856        QualityGrain::File => {
3857            let Some(source) = source else {
3858                return totals;
3859            };
3860            let whole_files = matches!(plan.scope, QualityScope::SourceFiles(_))
3861                || total_rows == Some(source.dataset_rows);
3862            if whole_files {
3863                for (index, name) in source.file_names.iter().enumerate() {
3864                    totals.insert(format!("file {name}"), source.file_rows(index));
3865                }
3866            }
3867        }
3868        QualityGrain::RowChunks(size) => {
3869            let (Some(total), size) = (total_rows, (*size).max(1)) else {
3870                return totals;
3871            };
3872            for chunk in 0..total.div_ceil(size) {
3873                let start = chunk * size;
3874                totals.insert(
3875                    format!("rows {}-{}", start + 1, (chunk + 1).saturating_mul(size)),
3876                    size.min(total - start),
3877                );
3878            }
3879        }
3880        _ => {}
3881    }
3882    totals
3883}
3884
3885fn profile_segments(
3886    df: &DataFrame,
3887    total_rows: Option<usize>,
3888    plan: &DataQualityPlan,
3889    precision: QualityPrecision,
3890    schema: &Schema,
3891    sample: SegmentSampleProvenance<'_>,
3892    polars_streaming: bool,
3893) -> Result<(Vec<SegmentQualityProfile>, Vec<UnsampledSegment>)> {
3894    let groups = segment_rows(df, plan, sample.positions)?;
3895    // Every segment in one grouped query: a query per segment would pay Polars' planning
3896    // thousands of times for a daily grain over years.
3897    let mut segment_of = vec![0u32; df.height()];
3898    for (index, group) in groups.iter().enumerate() {
3899        for row in &group.indices {
3900            segment_of[*row as usize] = index as u32;
3901        }
3902    }
3903    const SEGMENT: &str = "__quality_segment_index";
3904    let mut keyed = df.clone();
3905    keyed.with_column(Column::new(SEGMENT.into(), segment_of))?;
3906    let grouped = collect_lazy(
3907        keyed
3908            .lazy()
3909            .group_by([col(SEGMENT)])
3910            .agg(build_profile_exprs(schema)),
3911        polars_streaming,
3912    )
3913    .map_err(Report::from)?;
3914    let mut by_segment = vec![None; groups.len()];
3915    for row in 0..grouped.height() {
3916        let index = usize_value_at(&grouped, SEGMENT, row);
3917        if let Some(slot) = by_segment.get_mut(index) {
3918            *slot = Some(row);
3919        }
3920    }
3921    let mut profiles = Vec::with_capacity(groups.len());
3922    for (group, row) in groups.into_iter().zip(by_segment) {
3923        let evaluated_rows = group.indices.len();
3924        let Some(row) = row else {
3925            continue;
3926        };
3927        let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
3928        let null_cells = columns
3929            .iter()
3930            .map(|column| column.null_count)
3931            .sum::<usize>();
3932        let denominator = evaluated_rows.saturating_mul(columns.len());
3933        let known_segment_rows = sample.totals.get(&group.label);
3934        profiles.push(SegmentQualityProfile {
3935            label: group.label,
3936            total_rows: if let Some(total) = known_segment_rows {
3937                Some(*total)
3938            } else if matches!(plan.grain, QualityGrain::Dataset) {
3939                total_rows
3940            } else if precision == QualityPrecision::Exact {
3941                Some(evaluated_rows)
3942            } else {
3943                None
3944            },
3945            evaluated_rows,
3946            columns,
3947            null_cells,
3948            null_rate: rate(null_cells, denominator),
3949            compared_with: None,
3950            largest_change: None,
3951            change_size: None,
3952        });
3953    }
3954    order_segments(&mut profiles);
3955    apply_comparisons(
3956        &mut profiles,
3957        plan.comparison,
3958        plan.baseline_segment.as_deref(),
3959        precision,
3960    );
3961    // What the count found and the sample did not: kept apart, so a segment with
3962    // rows the sample missed is never read as one with none.
3963    let drawn = profiles
3964        .iter()
3965        .map(|profile| profile.label.as_str())
3966        .collect::<std::collections::HashSet<_>>();
3967    let mut unsampled = sample
3968        .totals
3969        .iter()
3970        .filter(|(label, rows)| **rows > 0 && !drawn.contains(label.as_str()))
3971        .map(|(label, rows)| UnsampledSegment {
3972            label: label.clone(),
3973            total_rows: *rows,
3974        })
3975        .collect::<Vec<_>>();
3976    unsampled.sort_by(|left, right| segment_cmp(&left.label, &right.label));
3977    Ok((profiles, unsampled))
3978}
3979
3980/// Segments in natural name order (year=9 before year=10), unplaced (`∅`) last, since
3981/// "previous" means the one a person would read before.
3982fn order_segments(segments: &mut [SegmentQualityProfile]) {
3983    segments.sort_by(|left, right| segment_cmp(&left.label, &right.label));
3984}
3985
3986/// The order of two segments by their labels, as [`order_segments`] puts them.
3987pub(crate) fn segment_cmp(left: &str, right: &str) -> std::cmp::Ordering {
3988    left.ends_with('∅')
3989        .cmp(&right.ends_with('∅'))
3990        .then_with(|| natural_cmp(left, right))
3991}
3992
3993/// Text compared with its runs of digits compared as numbers.
3994fn natural_cmp(left: &str, right: &str) -> std::cmp::Ordering {
3995    use std::cmp::Ordering;
3996    let (mut left, mut right) = (left, right);
3997    loop {
3998        let (Some(l), Some(r)) = (left.chars().next(), right.chars().next()) else {
3999            return left.len().cmp(&right.len());
4000        };
4001        if l.is_ascii_digit() && r.is_ascii_digit() {
4002            let digits = |text: &str| {
4003                text.find(|c: char| !c.is_ascii_digit())
4004                    .unwrap_or(text.len())
4005            };
4006            let (l_end, r_end) = (digits(left), digits(right));
4007            let (l_num, r_num) = (
4008                left[..l_end].trim_start_matches('0'),
4009                right[..r_end].trim_start_matches('0'),
4010            );
4011            let order = l_num.len().cmp(&r_num.len()).then_with(|| l_num.cmp(r_num));
4012            if order != Ordering::Equal {
4013                return order;
4014            }
4015            left = &left[l_end..];
4016            right = &right[r_end..];
4017        } else {
4018            if l != r {
4019                return l.cmp(&r);
4020            }
4021            left = &left[l.len_utf8()..];
4022            right = &right[r.len_utf8()..];
4023        }
4024    }
4025}
4026
4027fn profile_segments_lazy(
4028    lf: &LazyFrame,
4029    total_rows: usize,
4030    plan: &DataQualityPlan,
4031    source: Option<&QualitySourceContext>,
4032    schema: &Schema,
4033    polars_streaming: bool,
4034) -> Result<Vec<SegmentQualityProfile>> {
4035    if unsegmented(plan, source) {
4036        let aggregate = collect_lazy(
4037            lf.clone().select(build_profile_exprs(schema)),
4038            polars_streaming,
4039        )
4040        .map_err(Report::from)?;
4041        let columns = parse_profiles(&aggregate, schema, total_rows);
4042        return Ok(vec![whole_segment(
4043            plan,
4044            total_rows,
4045            &columns,
4046            schema.len(),
4047        )]);
4048    }
4049
4050    let (grouped_lf, group) = grouped_frame(lf, plan, source)?;
4051    let mut aggregates = vec![len().alias("__quality_segment_rows")];
4052    aggregates.extend(build_profile_exprs(schema));
4053    let grouped = collect_lazy(
4054        grouped_lf
4055            .group_by([group.alias("__quality_segment")])
4056            .agg(aggregates),
4057        polars_streaming,
4058    )
4059    .map_err(Report::from)?;
4060    let mut segments = Vec::with_capacity(grouped.height());
4061    let mut unassigned = Vec::with_capacity(grouped.height());
4062    for row in 0..grouped.height() {
4063        let evaluated_rows = usize_value_at(&grouped, "__quality_segment_rows", row);
4064        let columns = parse_profiles_at(&grouped, schema, evaluated_rows, row);
4065        let null_cells = columns
4066            .iter()
4067            .map(|column| column.null_count)
4068            .sum::<usize>();
4069        let denominator = evaluated_rows.saturating_mul(schema.len());
4070        let raw_label = string_value_at(&grouped, "__quality_segment", row);
4071        unassigned.push(raw_label.is_none());
4072        segments.push(SegmentQualityProfile {
4073            label: segment_label(&plan.grain, raw_label.as_deref()),
4074            total_rows: Some(evaluated_rows),
4075            evaluated_rows,
4076            columns,
4077            null_cells,
4078            null_rate: rate(null_cells, denominator),
4079            compared_with: None,
4080            largest_change: None,
4081            change_size: None,
4082        });
4083    }
4084    // Rows the grain could not place carry no order, so they follow the ones it could.
4085    let mut ordered = unassigned.into_iter().zip(segments).collect::<Vec<_>>();
4086    ordered.sort_by(|left, right| {
4087        left.0
4088            .cmp(&right.0)
4089            .then_with(|| natural_cmp(&left.1.label, &right.1.label))
4090    });
4091    let mut segments = ordered
4092        .into_iter()
4093        .map(|(_, segment)| segment)
4094        .collect::<Vec<_>>();
4095    if matches!(plan.grain, QualityGrain::RowChunks(_)) {
4096        for segment in &mut segments {
4097            segment.label = pretty_chunk_label(&segment.label);
4098        }
4099    }
4100    apply_comparisons(
4101        &mut segments,
4102        plan.comparison,
4103        plan.baseline_segment.as_deref(),
4104        QualityPrecision::Exact,
4105    );
4106    Ok(segments)
4107}
4108
4109/// Whether a full run's grain leaves the scope whole: the dataset grain, or files
4110/// where the view has lost which file a row came from.
4111fn unsegmented(plan: &DataQualityPlan, source: Option<&QualitySourceContext>) -> bool {
4112    matches!(plan.grain, QualityGrain::Dataset)
4113        || matches!(plan.grain, QualityGrain::File) && source.is_none()
4114}
4115
4116/// The scope as its one segment, from its columns' profile.
4117fn whole_segment(
4118    plan: &DataQualityPlan,
4119    total_rows: usize,
4120    columns: &[ColumnQualityProfile],
4121    column_count: usize,
4122) -> SegmentQualityProfile {
4123    let null_cells = columns
4124        .iter()
4125        .map(|column| column.null_count)
4126        .sum::<usize>();
4127    let denominator = total_rows.saturating_mul(column_count);
4128    SegmentQualityProfile {
4129        label: if matches!(plan.grain, QualityGrain::File) {
4130            "file mapping unavailable for this view".to_string()
4131        } else {
4132            "current view".to_string()
4133        },
4134        total_rows: Some(total_rows),
4135        evaluated_rows: total_rows,
4136        columns: columns.to_vec(),
4137        null_cells,
4138        null_rate: rate(null_cells, denominator),
4139        compared_with: None,
4140        largest_change: None,
4141        change_size: None,
4142    }
4143}
4144
4145fn grouped_frame(
4146    lf: &LazyFrame,
4147    plan: &DataQualityPlan,
4148    source: Option<&QualitySourceContext>,
4149) -> Result<(LazyFrame, Expr)> {
4150    match &plan.grain {
4151        QualityGrain::Dataset => Err(color_eyre::eyre::eyre!(
4152            "dataset grain does not need grouping"
4153        )),
4154        QualityGrain::Partition(column) => Ok((lf.clone(), col(column))),
4155        QualityGrain::RowChunks(size) => {
4156            let row = "__datui_quality_row";
4157            Ok((
4158                lf.clone().with_row_index(row, None),
4159                col(row).cast(DataType::UInt64) / lit((*size).max(1) as u64),
4160            ))
4161        }
4162        QualityGrain::TimeWindows { column, every } => Ok((
4163            lf.clone(),
4164            time_window_start(plan.time_value(column), every),
4165        )),
4166        QualityGrain::File => {
4167            let source = source
4168                .ok_or_else(|| color_eyre::eyre::eyre!("source-file mapping is unavailable"))?;
4169            let mut file = lit("unknown");
4170            for (start, name) in source.file_starts.iter().zip(source.file_names.iter()) {
4171                file = when(col(&source.row_index_column).gt_eq(lit(*start as u32)))
4172                    .then(lit(name.clone()))
4173                    .otherwise(file);
4174            }
4175            Ok((lf.clone(), file))
4176        }
4177    }
4178}
4179
4180fn pretty_chunk_label(label: &str) -> String {
4181    let Some(range) = label.strip_prefix("rows ") else {
4182        return label.to_string();
4183    };
4184    let Some((start, end)) = range.split_once('-') else {
4185        return label.to_string();
4186    };
4187    let start = start.trim_start_matches('0');
4188    let end = end.trim_start_matches('0');
4189    format!(
4190        "rows {}-{}",
4191        if start.is_empty() { "0" } else { start },
4192        if end.is_empty() { "0" } else { end }
4193    )
4194}
4195
4196fn segment_label(grain: &QualityGrain, raw: Option<&str>) -> String {
4197    match grain {
4198        QualityGrain::RowChunks(size) => {
4199            let raw = raw.unwrap_or("∅");
4200            raw.parse::<usize>()
4201                .map(|chunk| {
4202                    let start = chunk.saturating_mul(*size) + 1;
4203                    let end = start.saturating_add(*size).saturating_sub(1);
4204                    format!("rows {start:012}-{end:012}")
4205                })
4206                .unwrap_or_else(|_| format!("rows {raw}"))
4207        }
4208        QualityGrain::Partition(column) => format!("{column}={}", raw.unwrap_or("∅")),
4209        QualityGrain::TimeWindows { column, every } => time_window_label(column, every, raw),
4210        QualityGrain::File => format!("file {}", raw.unwrap_or("∅")),
4211        QualityGrain::Dataset => "current view".to_string(),
4212    }
4213}
4214
4215fn apply_comparisons(
4216    segments: &mut [SegmentQualityProfile],
4217    comparison: QualityComparison,
4218    baseline_segment: Option<&str>,
4219    precision: QualityPrecision,
4220) {
4221    for segment in segments.iter_mut() {
4222        segment.compared_with = None;
4223        segment.largest_change = None;
4224        segment.change_size = None;
4225    }
4226    let baseline_index = baseline_segment
4227        .and_then(|label| segments.iter().position(|segment| segment.label == label))
4228        .or_else(|| baseline_segment.is_none().then_some(0));
4229    if comparison == QualityComparison::Baseline && baseline_index.is_none() {
4230        for segment in segments {
4231            segment.largest_change = Some("selected baseline unavailable".to_string());
4232        }
4233        return;
4234    }
4235    for index in 0..segments.len() {
4236        let compared = match comparison {
4237            QualityComparison::None => None,
4238            QualityComparison::Previous if index > 0 => Some(index - 1),
4239            QualityComparison::Baseline if Some(index) != baseline_index => baseline_index,
4240            QualityComparison::Previous | QualityComparison::Baseline => None,
4241        };
4242        if let Some(other) = compared {
4243            let change = largest_material_change(&segments[index], &segments[other], precision);
4244            segments[index].compared_with = Some(segments[other].label.clone());
4245            if let Some((what, size)) = change {
4246                segments[index].largest_change = Some(what);
4247                segments[index].change_size = Some(size);
4248            }
4249        }
4250    }
4251}
4252
4253/// How far a measurement has to move between segments before it is worth naming, in
4254/// percentage points.
4255pub(crate) const MATERIAL_CHANGE_PP: f64 = 1.0;
4256
4257/// Standard errors apart two sampled rates must be before a difference is named: with
4258/// dozens of measures and thousands of segments, three would name noise daily.
4259const NOISE_Z: f64 = 4.0;
4260
4261/// Whether rates `a` of `n_a` rows and `b` of `n_b` rows differ by more than two
4262/// samples of those sizes would by chance (a two-proportion z-test).
4263pub fn beyond_noise(a: f64, n_a: usize, b: f64, n_b: usize) -> bool {
4264    if n_a == 0 || n_b == 0 {
4265        return false;
4266    }
4267    let (n_a, n_b) = (n_a as f64, n_b as f64);
4268    let pooled = (a * n_a + b * n_b) / (n_a + n_b);
4269    let error = (pooled * (1.0 - pooled) * (1.0 / n_a + 1.0 / n_b)).sqrt();
4270    error > 0.0 && (a - b).abs() / error >= NOISE_Z
4271}
4272
4273/// The rates segments are compared on. Not distinct share, which falls as a segment
4274/// grows.
4275const CHANGE_MEASURES: [QualityMetric; 4] = [
4276    QualityMetric::NullRate,
4277    QualityMetric::EmptyRate,
4278    QualityMetric::WhitespaceRate,
4279    QualityMetric::NonFiniteRate,
4280];
4281
4282/// The clearest move between two segments and its size: the sharpest single move,
4283/// not an average (one column going all-null is the finding). A halved or doubled row
4284/// count comes first. On a sample, only moves past sampling noise; range moves only
4285/// on exact profiles (a sample's min and max move with the draw).
4286fn largest_material_change(
4287    segment: &SegmentQualityProfile,
4288    baseline: &SegmentQualityProfile,
4289    precision: QualityPrecision,
4290) -> Option<(String, f64)> {
4291    if let (Some(now), Some(before)) = (segment.total_rows, baseline.total_rows)
4292        && before > 0
4293    {
4294        let ratio = now as f64 / before as f64;
4295        if !(0.5..2.0).contains(&ratio) {
4296            let percent = (ratio - 1.0) * 100.0;
4297            return Some((
4298                format!("rows {} ({percent:+.0}%)", crate::numfmt::group_chrome(now)),
4299                percent.abs(),
4300            ));
4301        }
4302    }
4303    let sampled = precision != QualityPrecision::Exact;
4304    let mut largest: Option<(f64, String)> = None;
4305    let mut range: Option<String> = None;
4306    for (index, column) in segment.columns.iter().enumerate() {
4307        // Both profiles walk the same schema, so columns line up; a search per column per
4308        // segment would be quadratic where thousands of segments meet hundreds of columns.
4309        let Some(prior) = baseline
4310            .columns
4311            .get(index)
4312            .filter(|other| other.name == column.name)
4313            .or_else(|| {
4314                baseline
4315                    .columns
4316                    .iter()
4317                    .find(|other| other.name == column.name)
4318            })
4319        else {
4320            continue;
4321        };
4322        for metric in CHANGE_MEASURES {
4323            let (Some(now), Some(before)) = (metric.value(column), metric.value(prior)) else {
4324                continue;
4325            };
4326            let change = (now - before) * 100.0;
4327            if change.abs() < MATERIAL_CHANGE_PP
4328                || sampled
4329                    && !beyond_noise(
4330                        now,
4331                        metric.denominator(column),
4332                        before,
4333                        metric.denominator(prior),
4334                    )
4335            {
4336                continue;
4337            }
4338            if largest
4339                .as_ref()
4340                .is_none_or(|(most, _)| change.abs() > most.abs())
4341            {
4342                largest = Some((change, format!("{} {}", column.name, metric.short_label())));
4343            }
4344        }
4345        if !sampled && range.is_none() && (column.min != prior.min || column.max != prior.max) {
4346            range = Some(format!(
4347                "{} range {} -> {}",
4348                column.name,
4349                range_label(prior),
4350                range_label(column)
4351            ));
4352        }
4353    }
4354    match (largest, range) {
4355        (Some((change, what)), _) => Some((format!("{what} {change:+.1} pp"), change.abs())),
4356        (None, Some(moved)) => Some((moved, 0.0)),
4357        (None, None) => None,
4358    }
4359}
4360
4361fn range_label(column: &ColumnQualityProfile) -> String {
4362    match (&column.min, &column.max) {
4363        (Some(min), Some(max)) => format!("{min}..{max}"),
4364        (Some(min), None) => format!("{min}.."),
4365        (None, Some(max)) => format!("..{max}"),
4366        (None, None) => "none".to_string(),
4367    }
4368}
4369
4370/// A role's column as a run reads it: the column's own name, where its time values
4371/// are, and where the rows are flagged whose text the column's format did not read.
4372struct TimedColumn {
4373    name: String,
4374    values: String,
4375    unparsed: Option<String>,
4376}
4377
4378/// One interval a run measures: its roles, the columns they sit on, and the grain
4379/// its rows are cut by.
4380struct ResolvedInterval {
4381    start_role: TemporalRole,
4382    end_role: TemporalRole,
4383    start: String,
4384    end: String,
4385    grain: QualityGrain,
4386}
4387
4388/// The measured intervals whose roles sit on columns readable as time; a role on text
4389/// without a format measures nothing and is left out.
4390fn resolved_intervals(plan: &DataQualityPlan, schema: &Schema) -> Vec<ResolvedInterval> {
4391    let usable = |role| {
4392        plan.role_column(role)
4393            .filter(|column| plan.reads_as_time(column, schema))
4394            .map(str::to_string)
4395    };
4396    plan.interval_pairs()
4397        .into_iter()
4398        .filter_map(|(start_role, end_role)| {
4399            let (start, end) = (usable(start_role)?, usable(end_role)?);
4400            let grain = plan.interval_grain(&start, &end);
4401            Some(ResolvedInterval {
4402                start_role,
4403                end_role,
4404                start,
4405                end,
4406                grain,
4407            })
4408        })
4409        .collect()
4410}
4411
4412/// The distinct grains `intervals` are cut by, in the order they first appear: one
4413/// grouping each, however many intervals share it.
4414fn interval_grains(intervals: &[ResolvedInterval]) -> Vec<QualityGrain> {
4415    let mut grains = Vec::new();
4416    for interval in intervals {
4417        if !grains.contains(&interval.grain) {
4418            grains.push(interval.grain.clone());
4419        }
4420    }
4421    grains
4422}
4423
4424/// How many groupings a run's intervals take: one per distinct grain. A full run
4425/// reads the scope once for each.
4426pub fn interval_passes(plan: &DataQualityPlan, schema: &Schema) -> usize {
4427    interval_grains(&resolved_intervals(plan, schema)).len()
4428}
4429
4430/// End minus start per row as a duration, null when either is missing or unread.
4431/// Dates are midnight; zoned times their UTC instant; zoneless times read as UTC.
4432fn interval_duration(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
4433    let as_time = |column: &str| {
4434        plan.time_value(column)
4435            .cast(DataType::Datetime(TimeUnit::Microseconds, None))
4436    };
4437    as_time(end) - as_time(start)
4438}
4439
4440/// [`interval_duration`] in microseconds, the unit its counts are taken in.
4441fn interval_micros(plan: &DataQualityPlan, start: &str, end: &str) -> Expr {
4442    interval_duration(plan, start, end)
4443        .dt()
4444        .total_microseconds(false)
4445}
4446
4447fn profile_temporal(
4448    df: &DataFrame,
4449    plan: &DataQualityPlan,
4450    sample_positions: Option<&[IdxSize]>,
4451) -> Result<Vec<TemporalLatencyProfile>> {
4452    // Resolved before grouping, as the lazy path does: the default plan assigns no roles,
4453    // and splitting into segments to learn that copies a frame per segment.
4454    let resolved = resolved_intervals(plan, df.schema());
4455    if resolved.is_empty() {
4456        return Ok(Vec::new());
4457    }
4458    // Text read as time is parsed once, beside the text, with a flag on the rows the
4459    // format did not read, so an unread value is told apart from a missing one.
4460    let mut parsed = Vec::new();
4461    let mut timed = |column: &str| {
4462        let Some(format) = plan.time_format(column) else {
4463            return TimedColumn {
4464                name: column.to_string(),
4465                values: column.to_string(),
4466                unparsed: None,
4467            };
4468        };
4469        let values = format!("__datui_quality_time::{column}");
4470        let unparsed = format!("__datui_quality_unparsed::{column}");
4471        if !parsed
4472            .iter()
4473            .any(|(name, _): &(String, Expr)| *name == values)
4474        {
4475            parsed.push((values.clone(), format.expr()));
4476            parsed.push((unparsed.clone(), format.unparsed()));
4477        }
4478        TimedColumn {
4479            name: column.to_string(),
4480            values,
4481            unparsed: Some(unparsed),
4482        }
4483    };
4484    let intervals = resolved
4485        .iter()
4486        .map(|interval| (interval, timed(&interval.start), timed(&interval.end)))
4487        .collect::<Vec<_>>();
4488    let df = if parsed.is_empty() {
4489        df.clone()
4490    } else {
4491        df.clone()
4492            .lazy()
4493            .with_columns(
4494                parsed
4495                    .into_iter()
4496                    .map(|(name, expr)| expr.alias(name))
4497                    .collect::<Vec<_>>(),
4498            )
4499            .collect()?
4500    };
4501    // Each grain's segments are cut once, whichever intervals share it.
4502    let mut cut: Vec<(QualityGrain, Vec<(String, DataFrame)>)> = Vec::new();
4503    for grain in interval_grains(&resolved) {
4504        let grain_plan = DataQualityPlan {
4505            grain: grain.clone(),
4506            ..plan.clone()
4507        };
4508        let segments = segment_rows(&df, &grain_plan, sample_positions)?
4509            .into_iter()
4510            .map(|group| Ok((group.label, take_rows(&df, &group.indices)?)))
4511            .collect::<Result<Vec<_>>>()?;
4512        cut.push((grain, segments));
4513    }
4514    // One interval's segments together, in their order, then the next interval's.
4515    let mut profiles = Vec::new();
4516    for (interval, start, end) in &intervals {
4517        let Some((_, segments)) = cut.iter().find(|(grain, _)| *grain == interval.grain) else {
4518            continue;
4519        };
4520        for (label, segment) in segments {
4521            profiles.push(latency_profile(
4522                segment,
4523                label,
4524                (interval.start_role, start),
4525                (interval.end_role, end),
4526                plan.latency_threshold_seconds,
4527            )?);
4528        }
4529    }
4530    Ok(profiles)
4531}
4532
4533fn profile_temporal_lazy(
4534    lf: &LazyFrame,
4535    plan: &DataQualityPlan,
4536    source: Option<&QualitySourceContext>,
4537    polars_streaming: bool,
4538) -> Result<Vec<TemporalLatencyProfile>> {
4539    let schema = lf.clone().collect_schema()?;
4540    let resolved = resolved_intervals(plan, &schema);
4541    if resolved.is_empty() {
4542        return Ok(Vec::new());
4543    }
4544
4545    let unparsed = |column: &str| {
4546        plan.time_format(column)
4547            .map(|format| format.unparsed().sum())
4548            .unwrap_or_else(|| lit(0u32))
4549    };
4550    let mut profiles = (0..resolved.len()).map(|_| Vec::new()).collect::<Vec<_>>();
4551    // One collect per grain the intervals are cut by: usually one, and one more for
4552    // each clock that differs.
4553    for grain in interval_grains(&resolved) {
4554        let mut expressions = vec![len().alias("__quality_temporal_rows")];
4555        let members = resolved
4556            .iter()
4557            .enumerate()
4558            .filter(|(_, interval)| interval.grain == grain)
4559            .map(|(index, _)| index)
4560            .collect::<Vec<_>>();
4561        for index in &members {
4562            let interval = &resolved[*index];
4563            let prefix = format!("latency::{index}::");
4564            let micros = interval_micros(plan, &interval.start, &interval.end);
4565            // Seconds as the sampled path takes them: whole seconds, toward zero.
4566            let seconds = interval_duration(plan, &interval.start, &interval.end)
4567                .dt()
4568                .total_seconds(false);
4569            expressions.extend([
4570                // Missing is the stored value; text the format did not read is
4571                // counted on its own.
4572                col(interval.start.as_str())
4573                    .is_null()
4574                    .sum()
4575                    .alias(format!("{prefix}missing_start")),
4576                col(interval.end.as_str())
4577                    .is_null()
4578                    .sum()
4579                    .alias(format!("{prefix}missing_end")),
4580                unparsed(&interval.start).alias(format!("{prefix}unparsed_start")),
4581                unparsed(&interval.end).alias(format!("{prefix}unparsed_end")),
4582                micros
4583                    .clone()
4584                    .is_not_null()
4585                    .sum()
4586                    .alias(format!("{prefix}paired")),
4587                micros
4588                    .clone()
4589                    .lt(lit(0i64))
4590                    .sum()
4591                    .alias(format!("{prefix}negative")),
4592                micros
4593                    .clone()
4594                    .eq(lit(0i64))
4595                    .sum()
4596                    .alias(format!("{prefix}zero")),
4597                seconds
4598                    .clone()
4599                    .quantile(lit(0.50), QuantileMethod::Nearest)
4600                    .alias(format!("{prefix}p50")),
4601                seconds
4602                    .clone()
4603                    .quantile(lit(0.90), QuantileMethod::Nearest)
4604                    .alias(format!("{prefix}p90")),
4605                seconds
4606                    .clone()
4607                    .quantile(lit(0.95), QuantileMethod::Nearest)
4608                    .alias(format!("{prefix}p95")),
4609                seconds
4610                    .clone()
4611                    .quantile(lit(0.99), QuantileMethod::Nearest)
4612                    .alias(format!("{prefix}p99")),
4613                seconds.max().alias(format!("{prefix}max")),
4614            ]);
4615            if let Some(threshold) = plan.latency_threshold_seconds {
4616                expressions.push(
4617                    micros
4618                        .gt(lit(threshold.saturating_mul(1_000_000)))
4619                        .sum()
4620                        .alias(format!("{prefix}above")),
4621                );
4622            }
4623        }
4624
4625        let grain_plan = DataQualityPlan {
4626            grain: grain.clone(),
4627            ..plan.clone()
4628        };
4629        let ungrouped = matches!(grain, QualityGrain::Dataset)
4630            || matches!(grain, QualityGrain::File) && source.is_none();
4631        let aggregate = if ungrouped {
4632            collect_lazy(lf.clone().select(expressions), polars_streaming).map_err(Report::from)?
4633        } else {
4634            let (grouped_lf, group) = grouped_frame(lf, &grain_plan, source)?;
4635            collect_lazy(
4636                grouped_lf
4637                    .group_by([group.alias("__quality_segment")])
4638                    .agg(expressions),
4639                polars_streaming,
4640            )
4641            .map_err(Report::from)?
4642        };
4643
4644        for row in 0..aggregate.height() {
4645            let segment = if ungrouped {
4646                if matches!(grain, QualityGrain::File) {
4647                    "file mapping unavailable for this view".to_string()
4648                } else {
4649                    "current view".to_string()
4650                }
4651            } else {
4652                let raw = string_value_at(&aggregate, "__quality_segment", row);
4653                segment_label(&grain, raw.as_deref())
4654            };
4655            let evaluated_rows = usize_value_at(&aggregate, "__quality_temporal_rows", row);
4656            for index in &members {
4657                let interval = &resolved[*index];
4658                let prefix = format!("latency::{index}::");
4659                let count =
4660                    |name: &str| usize_value_at(&aggregate, &format!("{prefix}{name}"), row);
4661                let seconds =
4662                    |name: &str| optional_i64_at(&aggregate, &format!("{prefix}{name}"), row);
4663                profiles[*index].push(TemporalLatencyProfile {
4664                    segment: segment.clone(),
4665                    start_role: interval.start_role,
4666                    end_role: interval.end_role,
4667                    start_column: interval.start.clone(),
4668                    end_column: interval.end.clone(),
4669                    evaluated_rows,
4670                    paired_rows: count("paired"),
4671                    missing_start: count("missing_start"),
4672                    missing_end: count("missing_end"),
4673                    unparsed_start: count("unparsed_start"),
4674                    unparsed_end: count("unparsed_end"),
4675                    negative_count: count("negative"),
4676                    zero_count: count("zero"),
4677                    p50_seconds: seconds("p50"),
4678                    p90_seconds: seconds("p90"),
4679                    p95_seconds: seconds("p95"),
4680                    p99_seconds: seconds("p99"),
4681                    max_seconds: seconds("max"),
4682                    threshold_seconds: plan.latency_threshold_seconds,
4683                    above_threshold_count: plan.latency_threshold_seconds.map(|_| count("above")),
4684                });
4685            }
4686        }
4687    }
4688    let mut ordered = Vec::new();
4689    for (interval, mut segments) in resolved.iter().zip(profiles) {
4690        segments.sort_by(|left, right| left.segment.cmp(&right.segment));
4691        // Strip the zero padding used for sorting chunks, as Segments does, so both name a
4692        // chunk alike.
4693        if matches!(interval.grain, QualityGrain::RowChunks(_)) {
4694            for profile in &mut segments {
4695                profile.segment = pretty_chunk_label(&profile.segment);
4696            }
4697        }
4698        ordered.extend(segments);
4699    }
4700    Ok(ordered)
4701}
4702
4703fn latency_profile(
4704    df: &DataFrame,
4705    segment: &str,
4706    (start_role, start): (TemporalRole, &TimedColumn),
4707    (end_role, end): (TemporalRole, &TimedColumn),
4708    threshold_seconds: Option<i64>,
4709) -> Result<TemporalLatencyProfile> {
4710    let starts = df.column(&start.values)?;
4711    let ends = df.column(&end.values)?;
4712    let flags = |column: &TimedColumn| {
4713        column
4714            .unparsed
4715            .as_ref()
4716            .map(|name| df.column(name))
4717            .transpose()
4718    };
4719    let (start_flags, end_flags) = (flags(start)?, flags(end)?);
4720    let unread = |flags: Option<&Column>, row: usize| -> Result<bool> {
4721        Ok(match flags {
4722            Some(flags) => flags.get(row)? == AnyValue::Boolean(true),
4723            None => false,
4724        })
4725    };
4726    let mut missing_start = 0;
4727    let mut missing_end = 0;
4728    let mut unparsed_start = 0;
4729    let mut unparsed_end = 0;
4730    let mut micros = Vec::new();
4731    for row in 0..df.height() {
4732        let start_at = value_epoch_micros(starts.get(row)?);
4733        let end_at = value_epoch_micros(ends.get(row)?);
4734        if start_at.is_none() {
4735            if unread(start_flags, row)? {
4736                unparsed_start += 1;
4737            } else {
4738                missing_start += 1;
4739            }
4740        }
4741        if end_at.is_none() {
4742            if unread(end_flags, row)? {
4743                unparsed_end += 1;
4744            } else {
4745                missing_end += 1;
4746            }
4747        }
4748        if let (Some(start_at), Some(end_at)) = (start_at, end_at) {
4749            micros.push(end_at - start_at);
4750        }
4751    }
4752    // Counted on the exact difference, so half a second early is early; the
4753    // percentiles are whole seconds.
4754    let negative_count = micros.iter().filter(|value| **value < 0).count();
4755    let zero_count = micros.iter().filter(|value| **value == 0).count();
4756    let above_threshold_count = threshold_seconds.map(|threshold| {
4757        let threshold = threshold.saturating_mul(1_000_000);
4758        micros.iter().filter(|value| **value > threshold).count()
4759    });
4760    let mut seconds = micros
4761        .iter()
4762        .map(|value| value / 1_000_000)
4763        .collect::<Vec<_>>();
4764    seconds.sort_unstable();
4765    let percentile = |percent: usize| {
4766        if seconds.is_empty() {
4767            None
4768        } else {
4769            let index = ((seconds.len() - 1) * percent + 50) / 100;
4770            seconds.get(index).copied()
4771        }
4772    };
4773    Ok(TemporalLatencyProfile {
4774        segment: segment.to_string(),
4775        start_role,
4776        end_role,
4777        start_column: start.name.clone(),
4778        end_column: end.name.clone(),
4779        evaluated_rows: df.height(),
4780        paired_rows: micros.len(),
4781        missing_start,
4782        missing_end,
4783        unparsed_start,
4784        unparsed_end,
4785        negative_count,
4786        zero_count,
4787        p50_seconds: percentile(50),
4788        p90_seconds: percentile(90),
4789        p95_seconds: percentile(95),
4790        p99_seconds: percentile(99),
4791        max_seconds: seconds.last().copied(),
4792        threshold_seconds,
4793        above_threshold_count,
4794    })
4795}
4796
4797/// Two counts per text column read as time, in whatever pass profiles the columns:
4798/// its non-null values, and those the format does not read.
4799fn interpretation_exprs(plan: &DataQualityPlan, schema: &Schema) -> Vec<Expr> {
4800    plan.time_formats
4801        .iter()
4802        .enumerate()
4803        .filter(|(_, format)| schema.get(&format.column).is_some())
4804        .flat_map(|(index, format)| {
4805            [
4806                col(format.column.as_str())
4807                    .is_not_null()
4808                    .sum()
4809                    .alias(format!("__datui_time::{index}::values")),
4810                format
4811                    .unparsed()
4812                    .sum()
4813                    .alias(format!("__datui_time::{index}::unparsed")),
4814            ]
4815        })
4816        .collect()
4817}
4818
4819/// Text the chosen format does not read, one observation per column that has any,
4820/// from the counts [`interpretation_exprs`] took.
4821fn interpretation_observations(
4822    counts: &DataFrame,
4823    plan: &DataQualityPlan,
4824    schema: &Schema,
4825) -> Vec<QualityObservation> {
4826    plan.time_formats
4827        .iter()
4828        .enumerate()
4829        .filter(|(_, format)| schema.get(&format.column).is_some())
4830        .filter_map(|(index, format)| {
4831            let values = optional_usize(counts, &format!("__datui_time::{index}::values"))?;
4832            let unparsed = optional_usize(counts, &format!("__datui_time::{index}::unparsed"))?;
4833            (unparsed > 0).then(|| QualityObservation {
4834                kind: ObservationKind::UnparsedTime,
4835                column: format.column.clone(),
4836                affected_rows: unparsed,
4837                evaluated_rows: values,
4838                fact: String::new(),
4839                normalized_category: None,
4840                files: Vec::new(),
4841                time_format: Some(format.clone()),
4842                full_scale: None,
4843            })
4844        })
4845        .collect()
4846}
4847
4848/// A categorical column stores integer codes, not text: `.str()` rejects it and
4849/// a numeric cast would measure the codes. Read its values as strings instead.
4850fn text_expr(column: Expr, dtype: &DataType) -> Expr {
4851    if matches!(dtype, DataType::Categorical(..)) {
4852        column.cast(DataType::String)
4853    } else {
4854        column
4855    }
4856}
4857
4858/// One count the profile pass takes per column: the alias suffix it is read back by,
4859/// its expression, and the field it fills. Pass and reader both use [`MEASURES`].
4860struct Measure {
4861    name: &'static str,
4862    applies: fn(&DataType) -> bool,
4863    expr: fn(Expr, &DataType) -> Expr,
4864    field: fn(&mut ColumnQualityProfile) -> &mut Option<usize>,
4865}
4866
4867fn is_text(dtype: &DataType) -> bool {
4868    matches!(dtype, DataType::String | DataType::Categorical(..))
4869}
4870
4871fn has_length(dtype: &DataType) -> bool {
4872    is_text(dtype) || matches!(dtype, DataType::List(_))
4873}
4874
4875/// Text that parses as `reading`, nulls not counted.
4876fn parse_count(column: Expr, dtype: &DataType, reading: TextReading) -> Expr {
4877    let text = text_expr(column, dtype);
4878    parses_as(text.clone(), reading)
4879        .and(text.is_not_null())
4880        .sum()
4881}
4882
4883/// A text value's length in characters, a list's in items.
4884fn length(column: Expr, dtype: &DataType) -> Expr {
4885    if matches!(dtype, DataType::List(_)) {
4886        column.list().len()
4887    } else {
4888        text_expr(column, dtype).str().len_chars()
4889    }
4890}
4891
4892const MEASURES: [Measure; 13] = [
4893    Measure {
4894        name: "distinct",
4895        applies: |_| true,
4896        expr: |column, _| column.clone().filter(column.is_not_null()).n_unique(),
4897        field: |profile| &mut profile.distinct_count,
4898    },
4899    Measure {
4900        name: "empty",
4901        applies: is_text,
4902        expr: |column, dtype| text_expr(column, dtype).eq(lit("")).sum(),
4903        field: |profile| &mut profile.empty_count,
4904    },
4905    Measure {
4906        name: "whitespace",
4907        applies: is_text,
4908        expr: |column, dtype| {
4909            let text = text_expr(column, dtype);
4910            text.clone()
4911                .str()
4912                .strip_chars(lit(LiteralValue::untyped_null()))
4913                .eq(lit(""))
4914                .and(text.neq(lit("")))
4915                .sum()
4916        },
4917        field: |profile| &mut profile.whitespace_count,
4918    },
4919    Measure {
4920        name: "parse_int",
4921        applies: is_text,
4922        // Polars' cast is not strict: text that is not a whole number becomes null.
4923        expr: |column, dtype| {
4924            let text = text_expr(column, dtype);
4925            text.clone()
4926                .cast(DataType::Int64)
4927                .is_not_null()
4928                .and(text.is_not_null())
4929                .sum()
4930        },
4931        field: |profile| &mut profile.integer_parse_count,
4932    },
4933    Measure {
4934        name: "leading_zero",
4935        applies: is_text,
4936        expr: |column, dtype| {
4937            let text = text_expr(column, dtype);
4938            text.clone()
4939                .str()
4940                .starts_with(lit("0"))
4941                .and(text.clone().str().len_chars().gt(lit(1u32)))
4942                .and(text.cast(DataType::Int64).is_not_null())
4943                .sum()
4944        },
4945        field: |profile| &mut profile.leading_zero_count,
4946    },
4947    Measure {
4948        name: "parse_decimal",
4949        applies: is_text,
4950        expr: |column, dtype| parse_count(column, dtype, TextReading::Decimal),
4951        field: |profile| &mut profile.decimal_parse_count,
4952    },
4953    Measure {
4954        name: "parse_date",
4955        applies: is_text,
4956        expr: |column, dtype| parse_count(column, dtype, TextReading::Date),
4957        field: |profile| &mut profile.date_parse_count,
4958    },
4959    Measure {
4960        name: "parse_datetime",
4961        applies: is_text,
4962        expr: |column, dtype| parse_count(column, dtype, TextReading::Datetime),
4963        field: |profile| &mut profile.datetime_parse_count,
4964    },
4965    Measure {
4966        name: "min_length",
4967        applies: has_length,
4968        expr: |column, dtype| length(column, dtype).min(),
4969        field: |profile| &mut profile.min_length,
4970    },
4971    Measure {
4972        name: "max_length",
4973        applies: has_length,
4974        expr: |column, dtype| length(column, dtype).max(),
4975        field: |profile| &mut profile.max_length,
4976    },
4977    Measure {
4978        name: "nan",
4979        applies: DataType::is_float,
4980        expr: |column, _| column.cast(DataType::Float64).is_nan().sum(),
4981        field: |profile| &mut profile.nan_count,
4982    },
4983    Measure {
4984        name: "pos_inf",
4985        applies: DataType::is_float,
4986        expr: |column, _| column.cast(DataType::Float64).eq(lit(f64::INFINITY)).sum(),
4987        field: |profile| &mut profile.positive_infinity_count,
4988    },
4989    Measure {
4990        name: "neg_inf",
4991        applies: DataType::is_float,
4992        expr: |column, _| {
4993            column
4994                .cast(DataType::Float64)
4995                .eq(lit(f64::NEG_INFINITY))
4996                .sum()
4997        },
4998        field: |profile| &mut profile.negative_infinity_count,
4999    },
5000];
5001
5002/// Every column's null count, range and [`MEASURES`], each aliased
5003/// `{column}::{measure}` for [`parse_profiles_at`].
5004fn build_profile_exprs(schema: &Schema) -> Vec<Expr> {
5005    let mut exprs = Vec::new();
5006    for (name, dtype) in schema.iter() {
5007        let column = col(name.as_str());
5008        let alias = |measure: &str| format!("{name}::{measure}");
5009        exprs.push(column.clone().null_count().alias(alias("null")));
5010        if supports_range(dtype) {
5011            exprs.push(column.clone().min().alias(alias("min")));
5012            exprs.push(column.clone().max().alias(alias("max")));
5013        }
5014        for measure in MEASURES.iter().filter(|measure| (measure.applies)(dtype)) {
5015            exprs.push((measure.expr)(column.clone(), dtype).alias(alias(measure.name)));
5016        }
5017    }
5018    exprs
5019}
5020
5021/// Whether text parses as `reading`: the test the profile counts with, so counts and
5022/// opened rows agree. Whole numbers count among decimals.
5023fn parses_as(text: Expr, reading: TextReading) -> Expr {
5024    // Named formats, not inference: "parses as an ISO date" has to mean the same
5025    // thing on every column, including one where nothing does.
5026    let strptime = |format: &str| StrptimeOptions {
5027        format: Some(PlSmallStr::from(format)),
5028        strict: false,
5029        exact: true,
5030        cache: true,
5031    };
5032    match reading {
5033        TextReading::WholeNumber | TextReading::Decimal => {
5034            text.cast(DataType::Float64).is_not_null()
5035        }
5036        TextReading::Date => text.str().to_date(strptime("%Y-%m-%d")).is_not_null(),
5037        TextReading::Datetime => [
5038            "%Y-%m-%d %H:%M:%S%.f",
5039            "%Y-%m-%dT%H:%M:%S%.f%#z",
5040            "%Y-%m-%dT%H:%M:%S%.f",
5041            "%Y-%m-%d %H:%M:%S",
5042            "%Y-%m-%dT%H:%M:%S%#z",
5043            "%Y-%m-%dT%H:%M:%S",
5044        ]
5045        .into_iter()
5046        .map(|format| {
5047            text.clone()
5048                .str()
5049                .to_datetime(
5050                    Some(TimeUnit::Microseconds),
5051                    None,
5052                    strptime(format),
5053                    lit(PlSmallStr::from_static("raise")),
5054                )
5055                .is_not_null()
5056        })
5057        .reduce(Expr::or)
5058        .expect("at least one datetime format"),
5059    }
5060}
5061
5062/// The rows of a parseable-text column its reading does not parse: non-null text
5063/// that stops a cast. `None` when the column has no reading.
5064pub fn unparsed_text(profile: &ColumnQualityProfile) -> Option<Expr> {
5065    let (_, reading) = text_reading(profile)?;
5066    let text = text_expr(col(profile.name.as_str()), &profile.dtype);
5067    Some(
5068        text.clone()
5069            .is_not_null()
5070            .and(parses_as(text, reading).not()),
5071    )
5072}
5073
5074fn supports_range(dtype: &DataType) -> bool {
5075    dtype.is_numeric()
5076        || dtype.is_temporal()
5077        || matches!(
5078            dtype,
5079            DataType::String | DataType::Categorical(..) | DataType::Boolean
5080        )
5081}
5082
5083fn parse_profiles(
5084    aggregate: &DataFrame,
5085    schema: &Schema,
5086    evaluated_rows: usize,
5087) -> Vec<ColumnQualityProfile> {
5088    parse_profiles_at(aggregate, schema, evaluated_rows, 0)
5089}
5090
5091fn parse_profiles_at(
5092    aggregate: &DataFrame,
5093    schema: &Schema,
5094    evaluated_rows: usize,
5095    row: usize,
5096) -> Vec<ColumnQualityProfile> {
5097    schema
5098        .iter()
5099        .map(|(name, dtype)| {
5100            let alias = |measure: &str| format!("{name}::{measure}");
5101            let mut profile = ColumnQualityProfile::unmeasured(name, dtype.clone(), evaluated_rows);
5102            profile.null_count = usize_value_at(aggregate, &alias("null"), row);
5103            profile.min = string_value_at(aggregate, &alias("min"), row);
5104            profile.max = string_value_at(aggregate, &alias("max"), row);
5105            for measure in &MEASURES {
5106                *(measure.field)(&mut profile) =
5107                    optional_usize_at(aggregate, &alias(measure.name), row);
5108            }
5109            profile
5110        })
5111        .collect()
5112}
5113
5114/// The share of non-null text that must parse before a column is said to hold numbers
5115/// or dates.
5116pub const TEXT_READING_SHARE: f64 = 0.95;
5117
5118/// What the values of a text column parse as, most specific first.
5119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5120pub enum TextReading {
5121    WholeNumber,
5122    Decimal,
5123    Datetime,
5124    Date,
5125}
5126
5127impl TextReading {
5128    pub fn label(self) -> &'static str {
5129        match self {
5130            Self::WholeNumber => "whole numbers",
5131            Self::Decimal => "decimal numbers",
5132            Self::Datetime => "ISO datetimes",
5133            Self::Date => "ISO dates",
5134        }
5135    }
5136
5137    pub fn is_number(self) -> bool {
5138        matches!(self, Self::WholeNumber | Self::Decimal)
5139    }
5140}
5141
5142/// The one typed reading a text column supports, with how many parse: the most
5143/// specific of overlapping readings. Whole only when every number is.
5144pub fn text_reading(profile: &ColumnQualityProfile) -> Option<(usize, TextReading)> {
5145    let non_null = profile.non_null_rows();
5146    if non_null == 0 {
5147        return None;
5148    }
5149    let enough = |count: Option<usize>| {
5150        count.filter(|parsed| *parsed as f64 >= non_null as f64 * TEXT_READING_SHARE)
5151    };
5152    if let Some(parsed) = enough(profile.decimal_parse_count) {
5153        let reading = if profile.integer_parse_count == Some(parsed) {
5154            TextReading::WholeNumber
5155        } else {
5156            TextReading::Decimal
5157        };
5158        return Some((parsed, reading));
5159    }
5160    [
5161        (profile.datetime_parse_count, TextReading::Datetime),
5162        (profile.date_parse_count, TextReading::Date),
5163    ]
5164    .into_iter()
5165    .find_map(|(count, reading)| enough(count).map(|parsed| (parsed, reading)))
5166}
5167
5168fn observations_from_profiles(
5169    columns: &[ColumnQualityProfile],
5170    precision: QualityPrecision,
5171) -> Vec<QualityObservation> {
5172    let mut observations = Vec::new();
5173    for profile in columns {
5174        if profile.null_count > 0 {
5175            observations.push(observation(
5176                ObservationKind::Nulls,
5177                profile,
5178                profile.null_count,
5179            ));
5180        }
5181        if let Some(count) = profile.empty_count.filter(|count| *count > 0) {
5182            observations.push(observation(ObservationKind::Empty, profile, count));
5183        }
5184        if let Some(count) = profile.whitespace_count.filter(|count| *count > 0) {
5185            observations.push(observation(ObservationKind::Whitespace, profile, count));
5186        }
5187        let non_finite = profile.nan_count.unwrap_or(0)
5188            + profile.positive_infinity_count.unwrap_or(0)
5189            + profile.negative_infinity_count.unwrap_or(0);
5190        if non_finite > 0 {
5191            observations.push(observation(ObservationKind::NonFinite, profile, non_finite));
5192        }
5193        if profile.distinct_count == Some(1) && profile.non_null_rows() > 0 {
5194            observations.push(observation(
5195                ObservationKind::Constant,
5196                profile,
5197                profile.non_null_rows(),
5198            ));
5199        }
5200        if let Some((parsed, _)) = text_reading(profile) {
5201            observations.push(observation(ObservationKind::ParseableText, profile, parsed));
5202        }
5203        // Near-unique yet repeating, from numbers already measured. Exact profiles only
5204        // (distinct counts do not extrapolate from samples), and only integers and text (a
5205        // float or timestamp is nearly unique by nature).
5206        if precision == QualityPrecision::Exact
5207            && (profile.dtype.is_integer()
5208                || matches!(profile.dtype, DataType::String | DataType::Categorical(..)))
5209            && let (Some(distinct), Some(uniqueness)) =
5210                (profile.distinct_count, profile.uniqueness_rate())
5211            && (KEY_LIKE_UNIQUENESS..1.0).contains(&uniqueness)
5212        {
5213            // Rows beyond one per value, as `DuplicateRows` counts extras (not all rows sharing a
5214            // value, which the drill-in opens).
5215            let extras = profile.non_null_rows().saturating_sub(distinct);
5216            if extras > 0 {
5217                observations.push(observation(ObservationKind::KeyLike, profile, extras));
5218            }
5219        }
5220    }
5221    observations
5222}
5223
5224/// How many one-column file reads a full run makes for conflict-hidden values, so the
5225/// access plan can promise them. From the footers, since it is asked every frame.
5226pub(crate) fn conflict_reads(
5227    file_group: &[u32],
5228    groups: &[crate::formats::schema_union::DriftGroup],
5229) -> usize {
5230    let mut per_column = BTreeMap::<&str, usize>::new();
5231    for group in file_group {
5232        let Some(group) = groups.get(*group as usize) else {
5233            continue;
5234        };
5235        for column in &group.unread {
5236            *per_column.entry(column.as_str()).or_default() += 1;
5237        }
5238    }
5239    per_column
5240        .values()
5241        .map(|files| (*files).min(MAX_EVIDENCE_FILES))
5242        .sum()
5243}
5244
5245/// Reads named columns of named files at each file's own type, the only way back to
5246/// conflict-hidden values: one column of the few disagreeing files.
5247#[derive(Clone)]
5248pub struct QualityConflictScan(pub crate::table::FileScan);
5249
5250impl std::fmt::Debug for QualityConflictScan {
5251    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
5252        f.write_str("QualityConflictScan")
5253    }
5254}
5255
5256/// What one column loses to disagreeing files: all counted, the largest few kept by
5257/// name, pruned as they arrive (one entry per file per column would cost hundreds of
5258/// megabytes on thousands of files).
5259#[derive(Default)]
5260struct DriftTally {
5261    files: usize,
5262    rows: usize,
5263    named: Vec<QualityFileEvidence>,
5264}
5265
5266impl DriftTally {
5267    fn add(&mut self, evidence: QualityFileEvidence) {
5268        self.files += 1;
5269        self.rows += evidence.rows;
5270        self.named.push(evidence);
5271        if self.named.len() > MAX_EVIDENCE_FILES * 2 {
5272            self.prune();
5273        }
5274    }
5275
5276    /// Largest first: the files that cost the column the most rows are the ones worth
5277    /// naming and worth reading values from.
5278    fn prune(&mut self) {
5279        self.named.sort_by(|left, right| {
5280            right
5281                .rows
5282                .cmp(&left.rows)
5283                .then_with(|| left.number.cmp(&right.number))
5284        });
5285        self.named.truncate(MAX_EVIDENCE_FILES);
5286    }
5287}
5288
5289/// Absent columns and type conflicts, from the footers already read. Measured over
5290/// the whole loaded source however the run was scoped (a scope's values say nothing
5291/// of columns its files lack); the detail pane notes the different denominator.
5292fn drift_observations(
5293    source: &QualitySourceContext,
5294    conflicts: Option<&QualityConflictScan>,
5295    polars_streaming: bool,
5296    watch: &QualityWatch,
5297) -> Vec<QualityObservation> {
5298    let mut absent = BTreeMap::<String, DriftTally>::new();
5299    let mut unread = BTreeMap::<String, DriftTally>::new();
5300    for (file, group) in source.drifting_files() {
5301        let evidence = |stored_type: Option<String>| QualityFileEvidence {
5302            number: file + 1,
5303            name: source
5304                .file_names
5305                .get(file)
5306                .cloned()
5307                .unwrap_or_else(|| format!("file {}", file + 1)),
5308            rows: source.file_rows(file),
5309            stored_type,
5310            examples: Vec::new(),
5311        };
5312        for column in &group.absent {
5313            absent
5314                .entry(column.to_string())
5315                .or_default()
5316                .add(evidence(None));
5317        }
5318        for column in &group.unread {
5319            let stored = source
5320                .stored_type(file, column)
5321                .map(|dtype| dtype.to_string());
5322            unread
5323                .entry(column.to_string())
5324                .or_default()
5325                .add(evidence(stored));
5326        }
5327    }
5328
5329    let mut observations = Vec::new();
5330    for (kind, columns) in [
5331        (ObservationKind::Absent, absent),
5332        (ObservationKind::TypeConflict, unread),
5333    ] {
5334        for (column, mut tally) in columns {
5335            tally.prune();
5336            let mut files = tally.named;
5337            if kind == ObservationKind::TypeConflict
5338                && let Some(scan) = conflicts
5339            {
5340                read_conflict_examples(scan, &column, &mut files, polars_streaming, watch);
5341            }
5342            let named = if tally.files > files.len() {
5343                format!(", largest {} named", files.len())
5344            } else {
5345                String::new()
5346            };
5347            let verb = match (kind, tally.files) {
5348                (ObservationKind::Absent, 1) => "has no such column",
5349                (ObservationKind::Absent, _) => "have no such column",
5350                (_, 1) => "holds a type the scan cannot read",
5351                (_, _) => "hold a type the scan cannot read",
5352            };
5353            // A file whose footer was not read looks exactly like one missing nothing,
5354            // so on a sampled dataset the count is a floor and has to say so.
5355            let sampled = if source.footers_read < source.file_names.len() {
5356                format!(", from {} footers read", source.footers_read)
5357            } else {
5358                String::new()
5359            };
5360            observations.push(QualityObservation {
5361                kind,
5362                column,
5363                affected_rows: tally.rows,
5364                evaluated_rows: source.dataset_rows,
5365                fact: format!(
5366                    "{} of {} files {verb}{sampled}{named}",
5367                    tally.files,
5368                    source.file_names.len()
5369                ),
5370                normalized_category: None,
5371                files,
5372                time_format: None,
5373                full_scale: None,
5374            });
5375        }
5376    }
5377    observations
5378}
5379
5380/// The first values each conflicting file holds, at its own type: one limited scan of
5381/// one column per file (a conflict is the file's, so its first values suffice). A
5382/// file that cannot be read keeps its count, losing only examples.
5383fn read_conflict_examples(
5384    scan: &QualityConflictScan,
5385    column: &str,
5386    files: &mut [QualityFileEvidence],
5387    polars_streaming: bool,
5388    watch: &QualityWatch,
5389) {
5390    let name = PlSmallStr::from(column);
5391    for file in files.iter_mut() {
5392        if watch.cancelled() {
5393            return;
5394        }
5395        let Ok(lf) = (scan.0)(
5396            std::slice::from_ref(&file.name),
5397            std::slice::from_ref(&name),
5398        ) else {
5399            continue;
5400        };
5401        let query = lf
5402            .select([crate::past_calendar::text_expr(
5403                col(column),
5404                CastOptions::NonStrict,
5405            )])
5406            .drop_nulls(None)
5407            .limit(MAX_CONFLICT_EXAMPLES as u32);
5408        let Ok(values) = collect_lazy(query, polars_streaming) else {
5409            continue;
5410        };
5411        file.examples = (0..values.height())
5412            .filter_map(|row| string_value_at(&values, column, row))
5413            .collect();
5414    }
5415}
5416
5417/// Read every sample of `audio` once and add clipping, zero runs and DC offset per
5418/// channel to `results`; only for a full run over every frame (the caller decides).
5419pub fn add_signal_observations(
5420    results: &mut DataQualityResults,
5421    audio: &crate::formats::audio::AudioSource,
5422    watch: &QualityWatch,
5423) -> Result<()> {
5424    watch.stage(QualityStage::CheckingSignal, true, true)?;
5425    let reports = audio
5426        .signal_report(&|| watch.cancelled())?
5427        .ok_or_else(|| Report::msg(crate::analysis::sampling::CANCELLED))?;
5428    results
5429        .observations
5430        .extend(signal_observations(&reports, audio.header().sample_rate));
5431    Ok(())
5432}
5433
5434/// Observations from [`crate::formats::audio::SignalReport`]s: a channel with runs at full
5435/// scale, runs of exact zeros, or a mean 1% of full scale or more from zero.
5436pub fn signal_observations(
5437    reports: &[crate::formats::audio::SignalReport],
5438    sample_rate: f64,
5439) -> Vec<QualityObservation> {
5440    let mut observations = Vec::new();
5441    let samples = |n: u64| {
5442        format!(
5443            "{} {}",
5444            crate::numfmt::group_chrome(n as usize),
5445            if n == 1 { "sample" } else { "samples" }
5446        )
5447    };
5448    let runs = |n: u64| if n == 1 { "run" } else { "runs" };
5449    for report in reports {
5450        let evaluated = report.frames as usize;
5451        let push = |observations: &mut Vec<QualityObservation>,
5452                    kind: ObservationKind,
5453                    affected: u64,
5454                    fact: String,
5455                    full_scale: Option<(f64, f64)>| {
5456            observations.push(QualityObservation {
5457                kind,
5458                column: report.channel.clone(),
5459                affected_rows: affected as usize,
5460                evaluated_rows: evaluated,
5461                fact,
5462                normalized_category: None,
5463                files: Vec::new(),
5464                time_format: None,
5465                full_scale,
5466            });
5467        };
5468        if report.clip_runs > 0 {
5469            push(
5470                &mut observations,
5471                ObservationKind::Clipping,
5472                report.in_clip_runs,
5473                format!(
5474                    "{} {} of {}+ samples at full scale; longest {}",
5475                    crate::numfmt::group_chrome(report.clip_runs as usize),
5476                    runs(report.clip_runs),
5477                    report.clip_run_min,
5478                    samples(report.longest_clip)
5479                ),
5480                Some(report.full_scale),
5481            );
5482        }
5483        if report.zero_runs > 0 {
5484            push(
5485                &mut observations,
5486                ObservationKind::ZeroRuns,
5487                report.in_zero_runs,
5488                format!(
5489                    "{} {} of exact zeros, {}+ samples; longest {} ({})",
5490                    crate::numfmt::group_chrome(report.zero_runs as usize),
5491                    runs(report.zero_runs),
5492                    report.zero_run_min,
5493                    samples(report.longest_zeros),
5494                    crate::widgets::info::clock(report.longest_zeros as f64 / sample_rate)
5495                ),
5496                None,
5497            );
5498        }
5499        let (low, high) = report.full_scale;
5500        let half_range = (high - low) / 2.0;
5501        let share = if half_range > 0.0 {
5502            report.mean.abs() / half_range
5503        } else {
5504            0.0
5505        };
5506        if share >= DC_OFFSET_SHARE {
5507            // Integers in their own units; float and normalized to four places.
5508            let mean = if half_range > 2.0 {
5509                format!("{:+.1}", report.mean)
5510            } else {
5511                format!("{:+.4}", report.mean)
5512            };
5513            push(
5514                &mut observations,
5515                ObservationKind::DcOffset,
5516                report.frames,
5517                format!("mean {mean} ({:.1}% of full scale)", share * 100.0),
5518                None,
5519            );
5520        }
5521    }
5522    observations
5523}
5524
5525/// A channel's mean, as a share of full scale, from which it is called DC offset: 1%,
5526/// -40 dBFS, well above any dither or noise floor.
5527const DC_OFFSET_SHARE: f64 = 0.01;
5528
5529fn observation(
5530    kind: ObservationKind,
5531    profile: &ColumnQualityProfile,
5532    affected_rows: usize,
5533) -> QualityObservation {
5534    QualityObservation {
5535        kind,
5536        column: profile.name.clone(),
5537        affected_rows,
5538        evaluated_rows: profile.evaluated_rows,
5539        fact: String::new(),
5540        normalized_category: None,
5541        files: Vec::new(),
5542        time_format: None,
5543        full_scale: None,
5544    }
5545}
5546
5547fn rate(numerator: usize, denominator: usize) -> f64 {
5548    if denominator == 0 {
5549        0.0
5550    } else {
5551        numerator as f64 / denominator as f64
5552    }
5553}
5554
5555fn optional_usize(df: &DataFrame, name: &str) -> Option<usize> {
5556    optional_usize_at(df, name, 0)
5557}
5558
5559fn optional_usize_at(df: &DataFrame, name: &str, row: usize) -> Option<usize> {
5560    let value = df.column(name).ok()?.get(row).ok()?;
5561    match value {
5562        AnyValue::UInt32(value) => Some(value as usize),
5563        AnyValue::UInt64(value) => Some(value as usize),
5564        AnyValue::Int32(value) => usize::try_from(value).ok(),
5565        AnyValue::Int64(value) => usize::try_from(value).ok(),
5566        _ => None,
5567    }
5568}
5569
5570fn usize_value(df: &DataFrame, name: &str) -> usize {
5571    optional_usize(df, name).unwrap_or(0)
5572}
5573
5574fn usize_value_at(df: &DataFrame, name: &str, row: usize) -> usize {
5575    optional_usize_at(df, name, row).unwrap_or(0)
5576}
5577
5578fn optional_i64_at(df: &DataFrame, name: &str, row: usize) -> Option<i64> {
5579    let value = df.column(name).ok()?.get(row).ok()?;
5580    match value {
5581        AnyValue::Int64(value) => Some(value),
5582        AnyValue::Int32(value) => Some(i64::from(value)),
5583        AnyValue::UInt64(value) => i64::try_from(value).ok(),
5584        AnyValue::UInt32(value) => Some(i64::from(value)),
5585        AnyValue::Float64(value) if value.is_finite() => Some(value.round() as i64),
5586        AnyValue::Float32(value) if value.is_finite() => Some(value.round() as i64),
5587        _ => None,
5588    }
5589}
5590
5591fn string_value_at(df: &DataFrame, name: &str, row: usize) -> Option<String> {
5592    let value = df.column(name).ok()?.get(row).ok()?;
5593    if value.is_null() {
5594        None
5595    } else {
5596        Some(crate::exact::str_value(&value).to_string())
5597    }
5598}
5599
5600#[cfg(test)]
5601pub(crate) mod fixtures;
5602#[cfg(test)]
5603mod tests;
5604
5605/// Intervals between time roles: what they count, and out of what.
5606#[cfg(test)]
5607mod temporal_tests;