1use crate::analysis::data_quality::{
7 DataQualityPlan, DataQualityResults, QualityComparison, QualityGrain, QualityMetric,
8 QualityPrecision, QualityScope, SegmentQualityProfile, beyond_noise, parse_scope_time,
9 segment_cmp,
10};
11use crate::analysis::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 segment_coverage(results: &DataQualityResults) -> (bool, usize, usize) {
136 let sampled = matches!(results.precision, QualityPrecision::Sampled);
137 let thin = if sampled {
138 results
139 .segments
140 .iter()
141 .filter(|segment| segment.evaluated_rows < THIN_SEGMENT_ROWS)
142 .count()
143 } else {
144 0
145 };
146 (sampled, results.unsampled_segments.len(), thin)
147}
148
149pub fn trend_view(
156 results: &DataQualityResults,
157 metric: QualityMetric,
158 bars: usize,
159) -> TrendView<'_> {
160 let slots = trend_slots(results);
161 let sampled = matches!(results.precision, QualityPrecision::Sampled);
162 if slots.is_empty() || bars == 0 {
163 return TrendView {
164 slots,
165 lines: Vec::new(),
166 bars: Vec::new(),
167 per_bar: 1,
168 sampled,
169 };
170 }
171 let per_bar = slots.len().div_ceil(bars);
172 let ranges = (0..slots.len())
173 .step_by(per_bar)
174 .map(|start| start..(start + per_bar).min(slots.len()))
175 .collect::<Vec<_>>();
176 let pooled = ranges
177 .iter()
178 .map(|range| {
179 let slots = &slots[range.clone()];
180 TrendBar {
181 slots: range.clone(),
182 unsampled: slots.iter().filter(|slot| slot.profile.is_none()).count(),
183 thin: if sampled {
184 slots
185 .iter()
186 .filter(|slot| slot.profile.is_some() && slot.evaluated < THIN_SEGMENT_ROWS)
187 .count()
188 } else {
189 0
190 },
191 evaluated: slots.iter().map(|slot| slot.evaluated).sum(),
192 eligible: slots.iter().map(|slot| slot.total).sum(),
193 }
194 })
195 .collect::<Vec<_>>();
196 let summarize = |names: Vec<String>, measure: TrendMeasure, parts: Vec<(f64, f64)>| {
197 let bars = parts
198 .iter()
199 .map(|(part, whole)| (*whole > 0.0).then(|| part / whole))
200 .collect::<Vec<_>>();
201 let known = bars.iter().flatten().copied();
202 let low = known.clone().fold(f64::INFINITY, f64::min);
203 let high = known.fold(0.0, f64::max);
204 TrendRow {
205 names,
206 measure,
207 bars,
208 parts,
209 low: if low.is_finite() { low } else { 0.0 },
210 high,
211 }
212 };
213 let rows_line = |name: &str, measure: TrendMeasure, rows: &dyn Fn(&TrendSlot<'_>) -> usize| {
214 summarize(
215 vec![name.to_string()],
216 measure,
217 ranges
218 .iter()
219 .map(|range| {
220 let total = slots[range.clone()].iter().map(rows).sum::<usize>();
221 (total as f64, range.len() as f64)
222 })
223 .collect(),
224 )
225 };
226 let counted = slots.iter().all(|slot| slot.total.is_some());
229 let mut lines = Vec::new();
230 if counted {
231 lines.push(rows_line("rows", TrendMeasure::Rows, &|slot| {
232 slot.total.unwrap_or(0)
233 }));
234 }
235 if sampled || !counted {
236 lines.push(rows_line(
237 "sampled rows",
238 TrendMeasure::SampledRows,
239 &|slot| slot.evaluated,
240 ));
241 }
242 let Some(first) = results.segments.first() else {
243 return TrendView {
244 slots,
245 lines,
246 bars: pooled,
247 per_bar,
248 sampled,
249 };
250 };
251 let cell = |slot: &TrendSlot<'_>, index: usize| -> Cell {
253 let column = slot.profile?.columns.get(index)?;
254 let value = metric.value(column)?;
255 let rows = metric.denominator(column) as f64;
256 Some((value * rows, rows))
257 };
258 let mut columns: Vec<(Vec<Cell>, TrendRow)> = Vec::new();
259 for (index, profile) in first.columns.iter().enumerate() {
260 let cells = slots
261 .iter()
262 .map(|slot| cell(slot, index))
263 .collect::<Vec<_>>();
264 let row = summarize(
265 vec![profile.name.clone()],
266 TrendMeasure::Column(index),
267 pool(&cells, &ranges),
268 );
269 if row.high == 0.0 {
270 continue;
271 }
272 match columns.iter_mut().find(|(other, _)| *other == cells) {
275 Some((_, other)) => other.names.push(profile.name.clone()),
276 None => columns.push((cells, row)),
277 }
278 }
279 let canonical = (0..slots.len())
283 .step_by(slots.len().div_ceil(ORDER_BARS))
284 .map(|start| start..(start + slots.len().div_ceil(ORDER_BARS)).min(slots.len()))
285 .collect::<Vec<_>>();
286 let spread = |cells: &[Cell]| {
287 let known = pool(cells, &canonical)
288 .into_iter()
289 .filter(|(_, whole)| *whole > 0.0)
290 .map(|(part, whole)| part / whole)
291 .collect::<Vec<_>>();
292 let high = known.iter().copied().fold(0.0, f64::max);
293 let low = known.iter().copied().fold(high, f64::min);
294 (high - low, high)
295 };
296 let mut columns = columns
297 .into_iter()
298 .map(|(cells, row)| (spread(&cells), row))
299 .collect::<Vec<_>>();
300 columns.sort_by(
301 |((left_spread, left_high), _), ((right_spread, right_high), _)| {
302 right_spread
303 .total_cmp(left_spread)
304 .then_with(|| right_high.total_cmp(left_high))
305 },
306 );
307 lines.extend(columns.into_iter().map(|(_, row)| row));
308 TrendView {
309 slots,
310 lines,
311 bars: pooled,
312 per_bar,
313 sampled,
314 }
315}
316
317const ORDER_BARS: usize = 32;
319
320fn pool(cells: &[Cell], ranges: &[Range<usize>]) -> Vec<(f64, f64)> {
322 ranges
323 .iter()
324 .map(|range| {
325 cells[range.clone()]
326 .iter()
327 .flatten()
328 .fold((0.0, 0.0), |(part, whole), (p, w)| (part + p, whole + w))
329 })
330 .collect()
331}
332
333pub fn wilson_interval(count: f64, n: f64) -> Option<(f64, f64)> {
337 if n <= 0.0 {
338 return None;
339 }
340 const Z: f64 = 1.96;
341 let p = (count / n).clamp(0.0, 1.0);
342 let z2 = Z * Z;
343 let centre = (p + z2 / (2.0 * n)) / (1.0 + z2 / n);
344 let half = Z * (p * (1.0 - p) / n + z2 / (4.0 * n * n)).sqrt() / (1.0 + z2 / n);
345 Some(((centre - half).max(0.0), (centre + half).min(1.0)))
346}
347
348pub fn compared_bar(view: &TrendView<'_>, bar: usize, plan: &DataQualityPlan) -> Option<usize> {
351 let other = if plan.comparison == QualityComparison::Baseline {
352 let slot = match plan.baseline_segment.as_deref() {
353 Some(label) => view.slots.iter().position(|slot| slot.label == label)?,
354 None => 0,
355 };
356 view.bars
357 .iter()
358 .position(|candidate| candidate.slots.contains(&slot))?
359 } else {
360 bar.checked_sub(1)?
361 };
362 (other != bar).then_some(other)
363}
364
365#[derive(Debug, Clone, Copy, PartialEq)]
367pub struct BarChange {
368 pub before: f64,
369 pub now: f64,
370 pub clear: bool,
372}
373
374impl BarChange {
375 pub fn points(&self) -> f64 {
376 (self.now - self.before) * 100.0
377 }
378}
379
380pub fn bar_change(line: &TrendRow, bar: usize, other: usize, exact: bool) -> Option<BarChange> {
384 let now = (*line.bars.get(bar)?)?;
385 let before = (*line.bars.get(other)?)?;
386 let clear = !line.rows()
387 && (now - before).abs() * 100.0 >= crate::analysis::data_quality::MATERIAL_CHANGE_PP
388 && (exact
389 || beyond_noise(
390 now,
391 line.parts[bar].1 as usize,
392 before,
393 line.parts[other].1 as usize,
394 ));
395 Some(BarChange { before, now, clear })
396}
397
398pub fn window_start(label: &str, every: &str) -> Option<NaiveDateTime> {
402 let date = |text: &str| NaiveDate::parse_from_str(text, "%Y-%m-%d").ok();
403 match every {
404 "1h" => NaiveDateTime::parse_from_str(label, "%Y-%m-%d %H:%M").ok(),
405 "1d" => date(label)?.and_hms_opt(0, 0, 0),
406 "1w" => date(label.strip_prefix("week of ")?)?.and_hms_opt(0, 0, 0),
407 "1mo" => date(&format!("{label}-01"))?.and_hms_opt(0, 0, 0),
408 _ => None,
409 }
410}
411
412pub fn next_window(start: NaiveDateTime, every: &str) -> Option<NaiveDateTime> {
414 match every {
415 "1h" => start.checked_add_signed(Duration::hours(1)),
416 "1d" => start.checked_add_signed(Duration::days(1)),
417 "1w" => start.checked_add_signed(Duration::weeks(1)),
418 "1mo" => start.checked_add_months(Months::new(1)),
419 _ => None,
420 }
421}
422
423pub fn floor_window(time: NaiveDateTime, every: &str) -> NaiveDateTime {
426 let midnight = |date: NaiveDate| date.and_hms_opt(0, 0, 0).unwrap_or(time);
427 match every {
428 "1h" => time.date().and_hms_opt(time.hour(), 0, 0).unwrap_or(time),
429 "1w" => {
430 midnight(time.date() - Duration::days(i64::from(time.weekday().num_days_from_monday())))
431 }
432 "1mo" => midnight(time.date().with_day(1).unwrap_or(time.date())),
433 _ => midnight(time.date()),
434 }
435}
436
437pub fn calendar_span(first: NaiveDateTime, last: NaiveDateTime, every: &str) -> String {
440 let end = next_window(last, every)
441 .and_then(|end| end.checked_sub_signed(Duration::minutes(1)))
442 .unwrap_or(last);
443 let text = |time: NaiveDateTime| {
444 if every == "1h" {
445 time.format("%Y-%m-%d %H:%M").to_string()
446 } else {
447 time.format("%Y-%m-%d").to_string()
448 }
449 };
450 let (first, end) = (text(first), text(end));
451 if first == end {
452 first
453 } else {
454 format!("{first} to {end}")
455 }
456}
457
458pub fn bar_span(view: &TrendView<'_>, bar: &TrendBar, grain: &QualityGrain) -> String {
462 let slots = &view.slots[bar.slots.clone()];
463 let (Some(first), Some(last)) = (slots.first(), slots.last()) else {
464 return String::new();
465 };
466 if let QualityGrain::TimeWindows { every, .. } = grain {
467 let starts = slots
468 .iter()
469 .filter_map(|slot| window_start(slot.label, every))
470 .collect::<Vec<_>>();
471 let undated = slots.len() - starts.len();
472 let calendar = match (starts.first(), starts.last()) {
473 (Some(first), Some(last)) => calendar_span(*first, *last, every),
474 _ => String::new(),
475 };
476 return match (calendar.is_empty(), undated) {
477 (_, 0) => calendar,
478 (true, _) => "rows with no time".to_string(),
479 (false, _) => format!("{calendar}, and rows with no time"),
480 };
481 }
482 if first.label == last.label {
483 first.label.to_string()
484 } else {
485 format!("{} to {}", first.label, last.label)
486 }
487}
488
489pub const MAX_EXPECTED_WINDOWS: usize = 20_000;
492
493pub const MAX_GAP_RUNS: usize = 500;
495
496#[derive(Debug, Clone, Copy, PartialEq, Eq)]
498pub enum GapKind {
499 Empty,
501 Unsampled,
503 OutOfScope,
506}
507
508impl GapKind {
509 pub fn label(self) -> &'static str {
510 match self {
511 Self::Empty => "empty",
512 Self::Unsampled => "not sampled",
513 Self::OutOfScope => "out of scope",
514 }
515 }
516}
517
518#[derive(Debug, Clone, PartialEq, Eq)]
520pub struct GapRun {
521 pub kind: GapKind,
522 pub first: NaiveDateTime,
524 pub last: NaiveDateTime,
525 pub windows: usize,
526 pub rows: Option<usize>,
528}
529
530#[derive(Debug, Clone, PartialEq, Eq)]
532pub struct GapCheck {
533 pub every: String,
534 pub column: String,
535 pub from: NaiveDateTime,
537 pub before: NaiveDateTime,
538 pub expected: usize,
540 pub weekend: usize,
542 pub with_rows: usize,
543 pub empty: usize,
544 pub unsampled: usize,
545 pub out_of_scope: usize,
546 pub counted: bool,
550 pub runs: Vec<GapRun>,
551 pub more_runs: usize,
553}
554
555impl GapCheck {
556 pub fn gaps(&self) -> usize {
557 self.empty + self.unsampled + self.out_of_scope
558 }
559}
560
561#[derive(Debug, Clone, PartialEq, Eq)]
563pub enum Gaps {
564 NoValues,
566 NoWindows,
568 TooMany {
570 windows: usize,
571 },
572 Checked(GapCheck),
573}
574
575pub fn expected_gaps(plan: &DataQualityPlan, results: &DataQualityResults) -> Option<Gaps> {
580 let expected = plan.expected_windows()?;
581 let QualityGrain::TimeWindows { column, every } = &plan.grain else {
582 return None;
583 };
584 if results.precision == QualityPrecision::Metadata {
585 return Some(Gaps::NoValues);
586 }
587 let slots = trend_slots(results);
588 let mut found = HashMap::new();
589 for slot in &slots {
590 if let Some(start) = window_start(slot.label, every) {
591 found.insert(start, (slot.evaluated, slot.total));
592 }
593 }
594 let (from, before) = expected.bounds();
595 let to_time =
596 |micros: i64| chrono::DateTime::from_timestamp_micros(micros).map(|at| at.naive_utc());
597 let from = match from.and_then(to_time) {
598 Some(from) => floor_window(from, every),
599 None => match found.keys().min() {
600 Some(first) => *first,
601 None => return Some(Gaps::NoWindows),
602 },
603 };
604 let before = match before.and_then(to_time) {
605 Some(before) => before,
606 None => match found
607 .keys()
608 .max()
609 .and_then(|last| next_window(*last, every))
610 {
611 Some(end) => end,
612 None => return Some(Gaps::NoWindows),
613 },
614 };
615 let windows = window_count(from, before, every);
616 if windows > MAX_EXPECTED_WINDOWS {
617 return Some(Gaps::TooMany { windows });
618 }
619 let scope = match &plan.scope {
621 QualityScope::SourceTimeRange {
622 column: scoped,
623 start,
624 end,
625 } if scoped == column => parse_scope_time(start)
626 .and_then(to_time)
627 .zip(parse_scope_time(end).and_then(to_time)),
628 _ => None,
629 };
630 let counted = results.precision == QualityPrecision::Exact
631 || slots.iter().all(|slot| slot.total.is_some());
632 let mut check = GapCheck {
633 every: every.clone(),
634 column: column.clone(),
635 from,
636 before,
637 expected: 0,
638 weekend: 0,
639 with_rows: 0,
640 empty: 0,
641 unsampled: 0,
642 out_of_scope: 0,
643 counted,
644 runs: Vec::new(),
645 more_runs: 0,
646 };
647 let weekdays =
648 expected.weekdays && crate::analysis::data_quality::ExpectedWindows::weekdays_apply(every);
649 let mut start = from;
650 let mut open: Option<GapRun> = None;
651 while start < before {
652 let Some(end) = next_window(start, every) else {
653 break;
654 };
655 if weekdays && matches!(start.weekday(), Weekday::Sat | Weekday::Sun) {
656 check.weekend += 1;
657 start = end;
658 continue;
659 }
660 check.expected += 1;
661 let gap = match found.get(&start) {
664 Some((evaluated, _)) if *evaluated > 0 => None,
665 Some((_, total)) => Some((GapKind::Unsampled, *total)),
666 None if scope.is_some_and(|(first, last)| start < first || end > last) => {
667 Some((GapKind::OutOfScope, None))
668 }
669 None if counted => Some((GapKind::Empty, None)),
670 None => Some((GapKind::Unsampled, None)),
671 };
672 match gap {
673 None => {
674 check.with_rows += 1;
675 close_run(&mut check, open.take());
676 }
677 Some((kind, rows)) => {
678 match kind {
679 GapKind::Empty => check.empty += 1,
680 GapKind::Unsampled => check.unsampled += 1,
681 GapKind::OutOfScope => check.out_of_scope += 1,
682 }
683 match open.as_mut() {
684 Some(run) if run.kind == kind => {
685 run.last = start;
686 run.windows += 1;
687 run.rows = run.rows.zip(rows).map(|(a, b)| a + b);
688 }
689 _ => {
690 close_run(&mut check, open.take());
691 open = Some(GapRun {
692 kind,
693 first: start,
694 last: start,
695 windows: 1,
696 rows,
697 });
698 }
699 }
700 }
701 }
702 start = end;
703 }
704 close_run(&mut check, open);
705 Some(Gaps::Checked(check))
706}
707
708fn close_run(check: &mut GapCheck, run: Option<GapRun>) {
709 let Some(run) = run else {
710 return;
711 };
712 if check.runs.len() < MAX_GAP_RUNS {
713 check.runs.push(run);
714 } else {
715 check.more_runs += 1;
716 }
717}
718
719fn window_count(from: NaiveDateTime, before: NaiveDateTime, every: &str) -> usize {
721 if before <= from {
722 return 0;
723 }
724 let span = before - from;
725 let per = |unit: Duration| {
726 let (span, unit) = (span.num_seconds(), unit.num_seconds().max(1));
727 usize::try_from((span + unit - 1) / unit).unwrap_or(usize::MAX)
728 };
729 match every {
730 "1h" => per(Duration::hours(1)),
731 "1w" => per(Duration::weeks(1)),
732 "1mo" => {
733 let months =
734 |time: NaiveDateTime| i64::from(time.year()) * 12 + i64::from(time.month0());
735 let whole = months(before) - months(from);
736 let past = before > floor_window(before, "1mo");
737 usize::try_from(whole + i64::from(past)).unwrap_or(usize::MAX)
738 }
739 _ => per(Duration::days(1)),
740 }
741}
742
743#[cfg(test)]
744mod tests;