1use crate::data_quality::{
7 DataQualityPlan, DataQualityResults, QualityComparison, QualityGrain, QualityMetric,
8 QualityPrecision, QualityScope, SegmentQualityProfile, beyond_noise, parse_scope_time,
9 segment_cmp, time_window_label,
10};
11use crate::quality_report::THIN_SEGMENT_ROWS;
12use chrono::{Datelike, Duration, Months, NaiveDate, NaiveDateTime, Timelike, Weekday};
13use std::collections::HashMap;
14use std::ops::Range;
15
16#[derive(Debug, Clone, Copy)]
19pub struct TrendSlot<'a> {
20 pub label: &'a str,
21 pub total: Option<usize>,
23 pub evaluated: usize,
25 pub profile: Option<&'a SegmentQualityProfile>,
26}
27
28pub fn trend_slots(results: &DataQualityResults) -> Vec<TrendSlot<'_>> {
30 let mut slots = results
31 .segments
32 .iter()
33 .map(|segment| TrendSlot {
34 label: &segment.label,
35 total: segment.total_rows,
36 evaluated: segment.evaluated_rows,
37 profile: Some(segment),
38 })
39 .chain(results.unsampled_segments.iter().map(|segment| TrendSlot {
40 label: &segment.label,
41 total: Some(segment.total_rows),
42 evaluated: 0,
43 profile: None,
44 }))
45 .collect::<Vec<_>>();
46 slots.sort_by(|left, right| segment_cmp(left.label, right.label));
49 slots
50}
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub enum TrendMeasure {
55 Rows,
57 SampledRows,
59 Column(usize),
61}
62
63#[derive(Debug, Clone, PartialEq)]
66pub struct TrendRow {
67 pub names: Vec<String>,
69 pub measure: TrendMeasure,
70 pub bars: Vec<Option<f64>>,
73 pub parts: Vec<(f64, f64)>,
76 pub low: f64,
77 pub high: f64,
78}
79
80impl TrendRow {
81 pub fn rows(&self) -> bool {
83 !matches!(self.measure, TrendMeasure::Column(_))
84 }
85}
86
87type Cell = Option<(f64, f64)>;
89
90#[derive(Debug, Clone, PartialEq, Eq)]
92pub struct TrendBar {
93 pub slots: Range<usize>,
95 pub unsampled: usize,
97 pub thin: usize,
100 pub evaluated: usize,
101 pub eligible: Option<usize>,
103}
104
105impl TrendBar {
106 pub fn segments(&self) -> usize {
107 self.slots.len()
108 }
109}
110
111#[derive(Debug, Clone)]
113pub struct TrendView<'a> {
114 pub slots: Vec<TrendSlot<'a>>,
115 pub lines: Vec<TrendRow>,
116 pub bars: Vec<TrendBar>,
117 pub per_bar: usize,
119 pub sampled: bool,
121}
122
123impl TrendView<'_> {
124 pub fn coverage(&self) -> (usize, usize) {
126 self.bars.iter().fold((0, 0), |(unsampled, thin), bar| {
127 (unsampled + bar.unsampled, thin + bar.thin)
128 })
129 }
130}
131
132pub fn trend_view(
139 results: &DataQualityResults,
140 metric: QualityMetric,
141 bars: usize,
142) -> TrendView<'_> {
143 let slots = trend_slots(results);
144 let sampled = matches!(
145 results.precision,
146 QualityPrecision::Sampled | QualityPrecision::Estimated
147 );
148 if slots.is_empty() || bars == 0 {
149 return TrendView {
150 slots,
151 lines: Vec::new(),
152 bars: Vec::new(),
153 per_bar: 1,
154 sampled,
155 };
156 }
157 let per_bar = slots.len().div_ceil(bars);
158 let ranges = (0..slots.len())
159 .step_by(per_bar)
160 .map(|start| start..(start + per_bar).min(slots.len()))
161 .collect::<Vec<_>>();
162 let pooled = ranges
163 .iter()
164 .map(|range| {
165 let slots = &slots[range.clone()];
166 TrendBar {
167 slots: range.clone(),
168 unsampled: slots.iter().filter(|slot| slot.profile.is_none()).count(),
169 thin: if sampled {
170 slots
171 .iter()
172 .filter(|slot| slot.profile.is_some() && slot.evaluated < THIN_SEGMENT_ROWS)
173 .count()
174 } else {
175 0
176 },
177 evaluated: slots.iter().map(|slot| slot.evaluated).sum(),
178 eligible: slots.iter().map(|slot| slot.total).sum(),
179 }
180 })
181 .collect::<Vec<_>>();
182 let summarize = |names: Vec<String>, measure: TrendMeasure, parts: Vec<(f64, f64)>| {
183 let bars = parts
184 .iter()
185 .map(|(part, whole)| (*whole > 0.0).then(|| part / whole))
186 .collect::<Vec<_>>();
187 let known = bars.iter().flatten().copied();
188 let low = known.clone().fold(f64::INFINITY, f64::min);
189 let high = known.fold(0.0, f64::max);
190 TrendRow {
191 names,
192 measure,
193 bars,
194 parts,
195 low: if low.is_finite() { low } else { 0.0 },
196 high,
197 }
198 };
199 let rows_line = |name: &str, measure: TrendMeasure, rows: &dyn Fn(&TrendSlot<'_>) -> usize| {
200 summarize(
201 vec![name.to_string()],
202 measure,
203 ranges
204 .iter()
205 .map(|range| {
206 let total = slots[range.clone()].iter().map(rows).sum::<usize>();
207 (total as f64, range.len() as f64)
208 })
209 .collect(),
210 )
211 };
212 let counted = slots.iter().all(|slot| slot.total.is_some());
215 let mut lines = Vec::new();
216 if counted {
217 lines.push(rows_line("rows", TrendMeasure::Rows, &|slot| {
218 slot.total.unwrap_or(0)
219 }));
220 }
221 if sampled || !counted {
222 lines.push(rows_line(
223 "sampled rows",
224 TrendMeasure::SampledRows,
225 &|slot| slot.evaluated,
226 ));
227 }
228 let Some(first) = results.segments.first() else {
229 return TrendView {
230 slots,
231 lines,
232 bars: pooled,
233 per_bar,
234 sampled,
235 };
236 };
237 let cell = |slot: &TrendSlot<'_>, index: usize| -> Cell {
239 let column = slot.profile?.columns.get(index)?;
240 let value = metric.value(column)?;
241 let rows = metric.denominator(column) as f64;
242 Some((value * rows, rows))
243 };
244 let mut columns: Vec<(Vec<Cell>, TrendRow)> = Vec::new();
245 for (index, profile) in first.columns.iter().enumerate() {
246 let cells = slots
247 .iter()
248 .map(|slot| cell(slot, index))
249 .collect::<Vec<_>>();
250 let row = summarize(
251 vec![profile.name.clone()],
252 TrendMeasure::Column(index),
253 pool(&cells, &ranges),
254 );
255 if row.high == 0.0 {
256 continue;
257 }
258 match columns.iter_mut().find(|(other, _)| *other == cells) {
261 Some((_, other)) => other.names.push(profile.name.clone()),
262 None => columns.push((cells, row)),
263 }
264 }
265 let canonical = (0..slots.len())
269 .step_by(slots.len().div_ceil(ORDER_BARS))
270 .map(|start| start..(start + slots.len().div_ceil(ORDER_BARS)).min(slots.len()))
271 .collect::<Vec<_>>();
272 let spread = |cells: &[Cell]| {
273 let known = pool(cells, &canonical)
274 .into_iter()
275 .filter(|(_, whole)| *whole > 0.0)
276 .map(|(part, whole)| part / whole)
277 .collect::<Vec<_>>();
278 let high = known.iter().copied().fold(0.0, f64::max);
279 let low = known.iter().copied().fold(high, f64::min);
280 (high - low, high)
281 };
282 let mut columns = columns
283 .into_iter()
284 .map(|(cells, row)| (spread(&cells), row))
285 .collect::<Vec<_>>();
286 columns.sort_by(
287 |((left_spread, left_high), _), ((right_spread, right_high), _)| {
288 right_spread
289 .total_cmp(left_spread)
290 .then_with(|| right_high.total_cmp(left_high))
291 },
292 );
293 lines.extend(columns.into_iter().map(|(_, row)| row));
294 TrendView {
295 slots,
296 lines,
297 bars: pooled,
298 per_bar,
299 sampled,
300 }
301}
302
303const ORDER_BARS: usize = 32;
305
306fn pool(cells: &[Cell], ranges: &[Range<usize>]) -> Vec<(f64, f64)> {
308 ranges
309 .iter()
310 .map(|range| {
311 cells[range.clone()]
312 .iter()
313 .flatten()
314 .fold((0.0, 0.0), |(part, whole), (p, w)| (part + p, whole + w))
315 })
316 .collect()
317}
318
319pub fn wilson_interval(count: f64, n: f64) -> Option<(f64, f64)> {
323 if n <= 0.0 {
324 return None;
325 }
326 const Z: f64 = 1.96;
327 let p = (count / n).clamp(0.0, 1.0);
328 let z2 = Z * Z;
329 let centre = (p + z2 / (2.0 * n)) / (1.0 + z2 / n);
330 let half = Z * (p * (1.0 - p) / n + z2 / (4.0 * n * n)).sqrt() / (1.0 + z2 / n);
331 Some(((centre - half).max(0.0), (centre + half).min(1.0)))
332}
333
334pub fn compared_bar(view: &TrendView<'_>, bar: usize, plan: &DataQualityPlan) -> Option<usize> {
337 let other = if plan.comparison == QualityComparison::Baseline {
338 let slot = match plan.baseline_segment.as_deref() {
339 Some(label) => view.slots.iter().position(|slot| slot.label == label)?,
340 None => 0,
341 };
342 view.bars
343 .iter()
344 .position(|candidate| candidate.slots.contains(&slot))?
345 } else {
346 bar.checked_sub(1)?
347 };
348 (other != bar).then_some(other)
349}
350
351#[derive(Debug, Clone, Copy, PartialEq)]
353pub struct BarChange {
354 pub before: f64,
355 pub now: f64,
356 pub clear: bool,
358}
359
360impl BarChange {
361 pub fn points(&self) -> f64 {
362 (self.now - self.before) * 100.0
363 }
364}
365
366pub fn bar_change(line: &TrendRow, bar: usize, other: usize, exact: bool) -> Option<BarChange> {
370 let now = (*line.bars.get(bar)?)?;
371 let before = (*line.bars.get(other)?)?;
372 let clear = !line.rows()
373 && (now - before).abs() * 100.0 >= crate::data_quality::MATERIAL_CHANGE_PP
374 && (exact
375 || beyond_noise(
376 now,
377 line.parts[bar].1 as usize,
378 before,
379 line.parts[other].1 as usize,
380 ));
381 Some(BarChange { before, now, clear })
382}
383
384pub fn window_start(label: &str, every: &str) -> Option<NaiveDateTime> {
388 let date = |text: &str| NaiveDate::parse_from_str(text, "%Y-%m-%d").ok();
389 match every {
390 "1h" => NaiveDateTime::parse_from_str(label, "%Y-%m-%d %H:%M").ok(),
391 "1d" => date(label)?.and_hms_opt(0, 0, 0),
392 "1w" => date(label.strip_prefix("week of ")?)?.and_hms_opt(0, 0, 0),
393 "1mo" => date(&format!("{label}-01"))?.and_hms_opt(0, 0, 0),
394 _ => None,
395 }
396}
397
398pub fn next_window(start: NaiveDateTime, every: &str) -> Option<NaiveDateTime> {
400 match every {
401 "1h" => start.checked_add_signed(Duration::hours(1)),
402 "1d" => start.checked_add_signed(Duration::days(1)),
403 "1w" => start.checked_add_signed(Duration::weeks(1)),
404 "1mo" => start.checked_add_months(Months::new(1)),
405 _ => None,
406 }
407}
408
409pub fn floor_window(time: NaiveDateTime, every: &str) -> NaiveDateTime {
412 let midnight = |date: NaiveDate| date.and_hms_opt(0, 0, 0).unwrap_or(time);
413 match every {
414 "1h" => time.date().and_hms_opt(time.hour(), 0, 0).unwrap_or(time),
415 "1w" => {
416 midnight(time.date() - Duration::days(i64::from(time.weekday().num_days_from_monday())))
417 }
418 "1mo" => midnight(time.date().with_day(1).unwrap_or(time.date())),
419 _ => midnight(time.date()),
420 }
421}
422
423pub fn window_label(column: &str, every: &str, start: NaiveDateTime) -> String {
425 time_window_label(
426 column,
427 every,
428 Some(&start.format("%Y-%m-%d %H:%M:%S").to_string()),
429 )
430}
431
432pub fn calendar_span(first: NaiveDateTime, last: NaiveDateTime, every: &str) -> String {
435 let end = next_window(last, every)
436 .and_then(|end| end.checked_sub_signed(Duration::minutes(1)))
437 .unwrap_or(last);
438 let text = |time: NaiveDateTime| {
439 if every == "1h" {
440 time.format("%Y-%m-%d %H:%M").to_string()
441 } else {
442 time.format("%Y-%m-%d").to_string()
443 }
444 };
445 let (first, end) = (text(first), text(end));
446 if first == end {
447 first
448 } else {
449 format!("{first} to {end}")
450 }
451}
452
453pub fn bar_span(view: &TrendView<'_>, bar: &TrendBar, grain: &QualityGrain) -> String {
457 let slots = &view.slots[bar.slots.clone()];
458 let (Some(first), Some(last)) = (slots.first(), slots.last()) else {
459 return String::new();
460 };
461 if let QualityGrain::TimeWindows { every, .. } = grain {
462 let starts = slots
463 .iter()
464 .filter_map(|slot| window_start(slot.label, every))
465 .collect::<Vec<_>>();
466 let undated = slots.len() - starts.len();
467 let calendar = match (starts.first(), starts.last()) {
468 (Some(first), Some(last)) => calendar_span(*first, *last, every),
469 _ => String::new(),
470 };
471 return match (calendar.is_empty(), undated) {
472 (_, 0) => calendar,
473 (true, _) => "rows with no time".to_string(),
474 (false, _) => format!("{calendar}, and rows with no time"),
475 };
476 }
477 if first.label == last.label {
478 first.label.to_string()
479 } else {
480 format!("{} to {}", first.label, last.label)
481 }
482}
483
484pub const MAX_EXPECTED_WINDOWS: usize = 20_000;
487
488pub const MAX_GAP_RUNS: usize = 500;
490
491#[derive(Debug, Clone, Copy, PartialEq, Eq)]
493pub enum GapKind {
494 Empty,
496 Unsampled,
498 OutOfScope,
501}
502
503impl GapKind {
504 pub fn label(self) -> &'static str {
505 match self {
506 Self::Empty => "empty",
507 Self::Unsampled => "not sampled",
508 Self::OutOfScope => "out of scope",
509 }
510 }
511}
512
513#[derive(Debug, Clone, PartialEq, Eq)]
515pub struct GapRun {
516 pub kind: GapKind,
517 pub first: NaiveDateTime,
519 pub last: NaiveDateTime,
520 pub windows: usize,
521 pub rows: Option<usize>,
523}
524
525#[derive(Debug, Clone, PartialEq, Eq)]
527pub struct GapCheck {
528 pub every: String,
529 pub column: String,
530 pub from: NaiveDateTime,
532 pub before: NaiveDateTime,
533 pub expected: usize,
535 pub weekend: usize,
537 pub with_rows: usize,
538 pub empty: usize,
539 pub unsampled: usize,
540 pub out_of_scope: usize,
541 pub counted: bool,
545 pub runs: Vec<GapRun>,
546 pub more_runs: usize,
548}
549
550impl GapCheck {
551 pub fn gaps(&self) -> usize {
552 self.empty + self.unsampled + self.out_of_scope
553 }
554}
555
556#[derive(Debug, Clone, PartialEq, Eq)]
558pub enum Gaps {
559 NoValues,
561 NoWindows,
563 TooMany {
565 windows: usize,
566 },
567 Checked(GapCheck),
568}
569
570pub fn expected_gaps(plan: &DataQualityPlan, results: &DataQualityResults) -> Option<Gaps> {
575 let expected = plan.expected_windows()?;
576 let QualityGrain::TimeWindows { column, every } = &plan.grain else {
577 return None;
578 };
579 if results.precision == QualityPrecision::Metadata {
580 return Some(Gaps::NoValues);
581 }
582 let slots = trend_slots(results);
583 let mut found = HashMap::new();
584 for slot in &slots {
585 if let Some(start) = window_start(slot.label, every) {
586 found.insert(start, (slot.evaluated, slot.total));
587 }
588 }
589 let (from, before) = expected.bounds();
590 let to_time =
591 |micros: i64| chrono::DateTime::from_timestamp_micros(micros).map(|at| at.naive_utc());
592 let from = match from.and_then(to_time) {
593 Some(from) => floor_window(from, every),
594 None => match found.keys().min() {
595 Some(first) => *first,
596 None => return Some(Gaps::NoWindows),
597 },
598 };
599 let before = match before.and_then(to_time) {
600 Some(before) => before,
601 None => match found
602 .keys()
603 .max()
604 .and_then(|last| next_window(*last, every))
605 {
606 Some(end) => end,
607 None => return Some(Gaps::NoWindows),
608 },
609 };
610 let windows = window_count(from, before, every);
611 if windows > MAX_EXPECTED_WINDOWS {
612 return Some(Gaps::TooMany { windows });
613 }
614 let scope = match &plan.scope {
616 QualityScope::SourceTimeRange {
617 column: scoped,
618 start,
619 end,
620 } if scoped == column => parse_scope_time(start)
621 .and_then(to_time)
622 .zip(parse_scope_time(end).and_then(to_time)),
623 _ => None,
624 };
625 let counted = results.precision == QualityPrecision::Exact
626 || slots.iter().all(|slot| slot.total.is_some());
627 let mut check = GapCheck {
628 every: every.clone(),
629 column: column.clone(),
630 from,
631 before,
632 expected: 0,
633 weekend: 0,
634 with_rows: 0,
635 empty: 0,
636 unsampled: 0,
637 out_of_scope: 0,
638 counted,
639 runs: Vec::new(),
640 more_runs: 0,
641 };
642 let weekdays = expected.weekdays && crate::data_quality::ExpectedWindows::weekdays_apply(every);
643 let mut start = from;
644 let mut open: Option<GapRun> = None;
645 while start < before {
646 let Some(end) = next_window(start, every) else {
647 break;
648 };
649 if weekdays && matches!(start.weekday(), Weekday::Sat | Weekday::Sun) {
650 check.weekend += 1;
651 start = end;
652 continue;
653 }
654 check.expected += 1;
655 let gap = match found.get(&start) {
658 Some((evaluated, _)) if *evaluated > 0 => None,
659 Some((_, total)) => Some((GapKind::Unsampled, *total)),
660 None if scope.is_some_and(|(first, last)| start < first || end > last) => {
661 Some((GapKind::OutOfScope, None))
662 }
663 None if counted => Some((GapKind::Empty, None)),
664 None => Some((GapKind::Unsampled, None)),
665 };
666 match gap {
667 None => {
668 check.with_rows += 1;
669 close_run(&mut check, open.take());
670 }
671 Some((kind, rows)) => {
672 match kind {
673 GapKind::Empty => check.empty += 1,
674 GapKind::Unsampled => check.unsampled += 1,
675 GapKind::OutOfScope => check.out_of_scope += 1,
676 }
677 match open.as_mut() {
678 Some(run) if run.kind == kind => {
679 run.last = start;
680 run.windows += 1;
681 run.rows = run.rows.zip(rows).map(|(a, b)| a + b);
682 }
683 _ => {
684 close_run(&mut check, open.take());
685 open = Some(GapRun {
686 kind,
687 first: start,
688 last: start,
689 windows: 1,
690 rows,
691 });
692 }
693 }
694 }
695 }
696 start = end;
697 }
698 close_run(&mut check, open);
699 Some(Gaps::Checked(check))
700}
701
702fn close_run(check: &mut GapCheck, run: Option<GapRun>) {
703 let Some(run) = run else {
704 return;
705 };
706 if check.runs.len() < MAX_GAP_RUNS {
707 check.runs.push(run);
708 } else {
709 check.more_runs += 1;
710 }
711}
712
713fn window_count(from: NaiveDateTime, before: NaiveDateTime, every: &str) -> usize {
715 if before <= from {
716 return 0;
717 }
718 let span = before - from;
719 let per = |unit: Duration| {
720 let (span, unit) = (span.num_seconds(), unit.num_seconds().max(1));
721 usize::try_from((span + unit - 1) / unit).unwrap_or(usize::MAX)
722 };
723 match every {
724 "1h" => per(Duration::hours(1)),
725 "1w" => per(Duration::weeks(1)),
726 "1mo" => {
727 let months =
728 |time: NaiveDateTime| i64::from(time.year()) * 12 + i64::from(time.month0());
729 let whole = months(before) - months(from);
730 let past = before > floor_window(before, "1mo");
731 usize::try_from(whole + i64::from(past)).unwrap_or(usize::MAX)
732 }
733 _ => per(Duration::days(1)),
734 }
735}
736
737#[cfg(test)]
738mod tests {
739 use super::*;
740 use crate::data_quality::{
741 ExpectedWindows, QualityCompute, UnsampledSegment, compute_data_quality,
742 };
743 use polars::prelude::{DataType, IntoLazy, LazyFrame, col, df};
744
745 fn at(text: &str) -> NaiveDateTime {
746 NaiveDateTime::parse_from_str(text, "%Y-%m-%d %H:%M").unwrap()
747 }
748
749 #[test]
752 fn window_labels_read_back_to_their_start() {
753 for (every, start, label) in [
754 ("1h", "2024-03-05 13:00", "2024-03-05 13:00"),
755 ("1d", "2024-03-05 00:00", "2024-03-05"),
756 ("1w", "2024-03-04 00:00", "week of 2024-03-04"),
757 ("1mo", "2024-03-01 00:00", "2024-03"),
758 ] {
759 assert_eq!(window_label("day", every, at(start)), label);
760 assert_eq!(window_start(label, every), Some(at(start)), "{label}");
761 assert_eq!(floor_window(at("2024-03-05 13:27"), every), at(start));
762 }
763 assert_eq!(window_start("day ∅", "1d"), None);
764 assert_eq!(
765 next_window(at("2024-01-01 00:00"), "1mo"),
766 Some(at("2024-02-01 00:00"))
767 );
768 assert_eq!(
769 window_count(at("2024-01-01 00:00"), at("2024-04-01 00:00"), "1mo"),
770 3
771 );
772 assert_eq!(
773 window_count(at("2024-01-01 00:00"), at("2024-04-02 00:00"), "1mo"),
774 4
775 );
776 assert_eq!(
777 window_count(at("2024-01-01 00:00"), at("2024-01-08 00:00"), "1d"),
778 7
779 );
780 assert_eq!(
781 calendar_span(at("2024-01-01 00:00"), at("2024-01-22 00:00"), "1w"),
782 "2024-01-01 to 2024-01-28"
783 );
784 assert_eq!(
785 calendar_span(at("2024-01-01 05:00"), at("2024-01-01 05:00"), "1h"),
786 "2024-01-01 05:00 to 2024-01-01 05:59"
787 );
788 }
789
790 #[test]
793 fn a_wilson_interval_stays_honest_at_the_edges() {
794 let (low, high) = wilson_interval(5.0, 100.0).unwrap();
795 assert!(
796 low < 0.05 && 0.05 < high && high - low < 0.12,
797 "{low} {high}"
798 );
799 let (low, high) = wilson_interval(0.0, 10.0).unwrap();
800 assert_eq!(low, 0.0);
801 assert!(high > 0.25, "ten rows say little: {high}");
802 let (low, high) = wilson_interval(10.0, 10.0).unwrap();
804 assert!((1.0 - high).abs() < 1e-12 && low < 0.75, "{low} {high}");
805 assert!(wilson_interval(0.0, 0.0).is_none());
806 }
807
808 fn weekdays() -> LazyFrame {
811 let monday = NaiveDate::from_ymd_opt(2024, 1, 1).unwrap();
812 let days = (0..56)
813 .map(|day| monday + Duration::days(day))
814 .filter(|day| day.weekday().num_days_from_monday() < 5)
815 .filter(|day| !(7..14).contains(&(*day - monday).num_days()))
816 .collect::<Vec<_>>();
817 let rows = days
818 .iter()
819 .flat_map(|day| std::iter::repeat_n(*day, 40))
820 .collect::<Vec<_>>();
821 let epoch = NaiveDate::from_ymd_opt(1970, 1, 1).unwrap();
822 let day = rows
823 .iter()
824 .map(|day| (*day - epoch).num_days() as i32)
825 .collect::<Vec<_>>();
826 let value = rows
827 .iter()
828 .enumerate()
829 .map(|(row, day)| {
830 ((*day - monday).num_days() < 35 || row % 2 == 0).then_some(row as i64)
831 })
832 .collect::<Vec<_>>();
833 df!("day" => day, "value" => value)
834 .unwrap()
835 .lazy()
836 .with_column(col("day").cast(DataType::Date))
837 }
838
839 fn daily(compute: QualityCompute, rows: usize) -> DataQualityPlan {
840 DataQualityPlan {
841 compute,
842 dataset_rows: rows,
843 grain: QualityGrain::TimeWindows {
844 column: "day".to_string(),
845 every: "1d".to_string(),
846 },
847 ..DataQualityPlan::default()
848 }
849 }
850
851 #[test]
854 fn gaps_appear_only_with_stated_windows() {
855 let frame = weekdays();
856 let plan = daily(QualityCompute::Full, 0);
857 let results = compute_data_quality(&frame, None, &plan, None, false).unwrap();
858 assert_eq!(expected_gaps(&plan, &results), None);
859 let whole = DataQualityPlan {
861 grain: QualityGrain::Dataset,
862 expected: Some(ExpectedWindows::default()),
863 ..plan.clone()
864 };
865 assert_eq!(expected_gaps(&whole, &results), None);
866 }
867
868 #[test]
871 fn an_exact_count_calls_a_window_empty() {
872 let frame = weekdays();
873 let mut plan = daily(QualityCompute::Full, 0);
874 plan.expected = Some(ExpectedWindows::default());
875 let results = compute_data_quality(&frame, None, &plan, None, false).unwrap();
876 let Some(Gaps::Checked(every_day)) = expected_gaps(&plan, &results) else {
877 panic!("checked");
878 };
879 assert_eq!(every_day.expected, 54, "first Monday to the last Friday");
880 assert_eq!(every_day.with_rows, 35);
881 assert_eq!(
882 every_day.empty, 19,
883 "seven days of week two and six weekends"
884 );
885 assert_eq!((every_day.unsampled, every_day.out_of_scope), (0, 0));
886 assert!(every_day.counted);
887
888 plan.expected = Some(ExpectedWindows {
889 weekdays: true,
890 ..ExpectedWindows::default()
891 });
892 let Some(Gaps::Checked(weekdays)) = expected_gaps(&plan, &results) else {
893 panic!("checked");
894 };
895 assert_eq!(weekdays.weekend, 14);
896 assert_eq!(weekdays.empty, 5);
897 assert_eq!(
898 weekdays.runs,
899 [GapRun {
900 kind: GapKind::Empty,
901 first: at("2024-01-08 00:00"),
902 last: at("2024-01-12 00:00"),
903 windows: 5,
904 rows: None,
905 }]
906 );
907 plan.expected = Some(ExpectedWindows {
909 weekdays: true,
910 from: Some("2024-01-01".to_string()),
911 before: Some("2024-03-04".to_string()),
912 });
913 let Some(Gaps::Checked(longer)) = expected_gaps(&plan, &results) else {
914 panic!("checked");
915 };
916 assert_eq!(longer.empty, 5 + 5);
917 assert_eq!(longer.runs.len(), 2);
918 }
919
920 #[test]
923 fn a_window_the_sample_missed_is_not_empty() {
924 let frame = weekdays();
925 let mut plan = daily(QualityCompute::Sample, 30);
926 plan.expected = Some(ExpectedWindows {
927 weekdays: true,
928 ..ExpectedWindows::default()
929 });
930 let results = compute_data_quality(&frame, Some(35 * 40), &plan, None, false).unwrap();
931 assert_eq!(results.precision, QualityPrecision::Sampled);
932 assert!(
933 !results.unsampled_segments.is_empty(),
934 "30 rows cannot reach 35 days"
935 );
936 assert!(
937 results
938 .unsampled_segments
939 .iter()
940 .all(|segment| segment.total_rows == 40)
941 );
942 let Some(Gaps::Checked(check)) = expected_gaps(&plan, &results) else {
943 panic!("checked");
944 };
945 assert!(check.counted);
946 assert_eq!(check.unsampled, results.unsampled_segments.len());
947 assert_eq!(check.empty, 5, "the missing week, from the count");
948 assert_eq!(check.with_rows + check.unsampled, 35);
949 let missed = check
950 .runs
951 .iter()
952 .filter(|run| run.kind == GapKind::Unsampled)
953 .map(|run| run.rows)
954 .collect::<Vec<_>>();
955 assert!(
956 missed
957 .iter()
958 .all(|rows| rows.is_some_and(|rows| rows % 40 == 0))
959 );
960
961 let view = trend_view(&results, QualityMetric::NullRate, 35);
963 assert_eq!(view.slots.len(), 35);
964 assert_eq!(view.per_bar, 1);
965 assert_eq!(view.coverage().0, results.unsampled_segments.len());
966 assert_eq!(view.lines[0].names, ["rows"]);
967 assert_eq!(view.lines[1].names, ["sampled rows"]);
968 let missed = view.bars.iter().position(|bar| bar.unsampled == 1).unwrap();
969 assert_eq!(view.lines[0].bars[missed], Some(40.0), "counted");
970 assert_eq!(view.lines[1].bars[missed], Some(0.0), "drawn none");
971 if let Some(value) = view.lines.iter().find(|line| line.names == ["value"]) {
972 assert_eq!(value.bars[missed], None);
973 }
974 }
975
976 #[test]
979 fn windows_outside_the_scope_are_out_of_scope() {
980 let frame = weekdays();
981 let mut plan = daily(QualityCompute::Full, 0);
982 plan.scope = QualityScope::SourceTimeRange {
983 column: "day".to_string(),
984 start: "2024-01-15".to_string(),
985 end: "2024-02-01".to_string(),
986 };
987 plan.expected = Some(ExpectedWindows {
988 weekdays: true,
989 from: Some("2024-01-01".to_string()),
990 before: Some("2024-02-05".to_string()),
991 });
992 let scoped = crate::data_quality::apply_quality_scope(frame, &plan.scope, None).unwrap();
993 let results = compute_data_quality(&scoped, None, &plan, None, false).unwrap();
994 let Some(Gaps::Checked(check)) = expected_gaps(&plan, &results) else {
995 panic!("checked");
996 };
997 assert_eq!(
998 check.out_of_scope,
999 10 + 2,
1000 "two weeks before, two days after"
1001 );
1002 assert_eq!(check.empty, 0);
1003 assert_eq!(check.with_rows, 13);
1004 assert_eq!(check.runs[0].kind, GapKind::OutOfScope);
1005
1006 plan.grain = QualityGrain::TimeWindows {
1009 column: "day".to_string(),
1010 every: "1w".to_string(),
1011 };
1012 plan.scope = QualityScope::SourceTimeRange {
1013 column: "day".to_string(),
1014 start: "2024-01-17".to_string(),
1015 end: "2024-02-07".to_string(),
1016 };
1017 let scoped =
1018 crate::data_quality::apply_quality_scope(weekdays(), &plan.scope, None).unwrap();
1019 let results = compute_data_quality(&scoped, None, &plan, None, false).unwrap();
1020 let Some(Gaps::Checked(check)) = expected_gaps(&plan, &results) else {
1021 panic!("checked");
1022 };
1023 assert_eq!(check.expected, 5, "January 1 to before February 5");
1024 assert_eq!(
1025 check.with_rows, 3,
1026 "the weeks of January 15, 22 and 29: {check:?}"
1027 );
1028 assert_eq!(check.out_of_scope, 2, "the weeks of January 1 and 8");
1029 }
1030
1031 #[test]
1033 fn the_windows_checked_are_bounded() {
1034 let frame = weekdays();
1035 let mut plan = DataQualityPlan {
1036 compute: QualityCompute::Full,
1037 grain: QualityGrain::TimeWindows {
1038 column: "day".to_string(),
1039 every: "1h".to_string(),
1040 },
1041 ..DataQualityPlan::default()
1042 };
1043 plan.expected = Some(ExpectedWindows {
1044 from: Some("2020-01-01".to_string()),
1045 before: Some("2024-01-01".to_string()),
1046 ..ExpectedWindows::default()
1047 });
1048 let results = compute_data_quality(&frame, None, &plan, None, false).unwrap();
1049 assert_eq!(
1050 expected_gaps(&plan, &results),
1051 Some(Gaps::TooMany { windows: 35_064 })
1052 );
1053 }
1054
1055 #[test]
1058 fn trend_lines_keep_their_order_at_any_width() {
1059 let rows: i32 = 60 * 20;
1060 let frame = df!(
1061 "day" => (0..rows).map(|row| 19_723 + row / 20).collect::<Vec<_>>(),
1062 "steady" => (0..rows).map(|row| (row % 10 != 0).then_some(1i64)).collect::<Vec<_>>(),
1063 "late" => (0..rows).map(|row| (row < 900 || row % 2 == 0).then_some(1i64)).collect::<Vec<_>>(),
1064 "spike" => (0..rows).map(|row| !(400..440).contains(&row)).map(|kept| kept.then_some(1i64)).collect::<Vec<_>>(),
1065 )
1066 .unwrap()
1067 .lazy()
1068 .with_column(col("day").cast(DataType::Date));
1069 let results =
1070 compute_data_quality(&frame, None, &daily(QualityCompute::Full, 0), None, false)
1071 .unwrap();
1072 let order = |bars: usize| {
1073 trend_view(&results, QualityMetric::NullRate, bars)
1074 .lines
1075 .into_iter()
1076 .map(|line| line.names)
1077 .collect::<Vec<_>>()
1078 };
1079 let first = order(8);
1080 assert_eq!(first.len(), 4, "rows and three columns: {first:?}");
1081 for bars in [1, 15, 23, 60, 200] {
1082 assert_eq!(order(bars), first, "{bars} bars");
1083 }
1084 }
1085
1086 #[test]
1088 fn missed_segments_keep_their_place() {
1089 let frame = weekdays();
1090 let plan = daily(QualityCompute::Full, 0);
1091 let mut results = compute_data_quality(&frame, None, &plan, None, false).unwrap();
1092 let moved = results.segments.remove(3);
1093 results.unsampled_segments.push(UnsampledSegment {
1094 label: moved.label.clone(),
1095 total_rows: 40,
1096 });
1097 let slots = trend_slots(&results);
1098 assert_eq!(slots[3].label, moved.label);
1099 assert!(slots[3].profile.is_none());
1100 }
1101}