1use crate::analysis::data_quality::{
10 DataQualityPlan, ObservationKind, QualityObservation, QualityPrecision, TimeInterpretation,
11 TimeKind,
12};
13use crate::analysis::statistics::collect_lazy;
14use color_eyre::Result;
15use polars::prelude::*;
16
17pub const MAX_ALLOWED_VALUES: usize = 100;
19
20pub const MAX_INTENT_EXAMPLES: usize = 3;
23
24const PREFIX: &str = "__datui_intent::";
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
28pub enum NumberReading {
29 Whole,
30 Decimal,
31}
32
33impl NumberReading {
34 pub fn label(self) -> &'static str {
35 match self {
36 Self::Whole => "whole number",
37 Self::Decimal => "decimal",
38 }
39 }
40
41 fn dtype(self) -> DataType {
42 match self {
43 Self::Whole => DataType::Int64,
44 Self::Decimal => DataType::Float64,
45 }
46 }
47}
48
49#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)]
51pub struct ColumnIntent {
52 pub column: String,
53 pub required: bool,
55 pub allowed: Vec<String>,
57 pub min: Option<String>,
59 pub max: Option<String>,
60 pub number: Option<NumberReading>,
63}
64
65impl ColumnIntent {
66 pub fn new(column: &str) -> Self {
67 Self {
68 column: column.to_string(),
69 ..Self::default()
70 }
71 }
72
73 pub fn is_empty(&self) -> bool {
74 !self.required
75 && self.allowed.is_empty()
76 && self.min.is_none()
77 && self.max.is_none()
78 && self.number.is_none()
79 }
80
81 pub fn range_label(&self) -> Option<String> {
83 match (&self.min, &self.max) {
84 (Some(min), Some(max)) => Some(format!("{min} to {max}")),
85 (Some(min), None) => Some(format!("at least {min}")),
86 (None, Some(max)) => Some(format!("at most {max}")),
87 (None, None) => None,
88 }
89 }
90
91 pub fn allowed_label(&self, shown: usize) -> String {
93 let mut label = format_allowed(&self.allowed[..self.allowed.len().min(shown)]);
94 if self.allowed.len() > shown {
95 label.push_str(&format!(" +{} more", self.allowed.len() - shown));
96 }
97 label
98 }
99
100 pub fn rules(&self) -> Vec<String> {
102 let mut rules = Vec::new();
103 if self.required {
104 rules.push("required".to_string());
105 }
106 if let Some(number) = self.number {
107 rules.push(format!("read as {}", number.label()));
108 }
109 if !self.allowed.is_empty() {
110 rules.push(format!(
111 "one of {}",
112 crate::numfmt::group_chrome(self.allowed.len())
113 ));
114 }
115 if let Some(range) = self.range_label() {
116 rules.push(range);
117 }
118 rules
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)]
124pub struct DeclaredIntent {
125 pub key: Vec<String>,
127 pub columns: Vec<ColumnIntent>,
128}
129
130impl DeclaredIntent {
131 pub fn is_empty(&self) -> bool {
132 self.key.is_empty() && self.columns.is_empty()
133 }
134
135 pub fn column(&self, name: &str) -> Option<&ColumnIntent> {
136 self.columns.iter().find(|intent| intent.column == name)
137 }
138
139 pub fn set(&mut self, intent: ColumnIntent) {
142 match self
143 .columns
144 .iter()
145 .position(|known| known.column == intent.column)
146 {
147 Some(index) if intent.is_empty() => {
148 self.columns.remove(index);
149 }
150 Some(index) => self.columns[index] = intent,
151 None if intent.is_empty() => {}
152 None => self.columns.push(intent),
153 }
154 }
155
156 pub fn set_key(&mut self, column: &str, in_key: bool) {
158 let known = self.key.iter().any(|name| name == column);
159 if in_key && !known {
160 self.key.push(column.to_string());
161 } else if !in_key {
162 self.key.retain(|name| name != column);
163 }
164 }
165
166 pub fn declared_columns(&self) -> Vec<&str> {
168 let mut names = self.key.iter().map(String::as_str).collect::<Vec<_>>();
169 for intent in &self.columns {
170 if !names.contains(&intent.column.as_str()) {
171 names.push(intent.column.as_str());
172 }
173 }
174 names
175 }
176
177 pub fn summary(&self) -> String {
179 let mut parts = Vec::new();
180 if !self.key.is_empty() {
181 parts.push(format!("key {}", self.key.join(", ")));
182 }
183 match self.columns.len() {
184 0 => {}
185 1 => parts.push(format!(
186 "{}: {}",
187 self.columns[0].column,
188 self.columns[0].rules().join(", ")
189 )),
190 count => parts.push(format!("rules on {count} columns")),
191 }
192 parts.join(&format!(" {} ", crate::glyphs::get().middot))
193 }
194}
195
196#[derive(Debug, Clone, Copy, PartialEq, Eq)]
199pub enum ValueKind {
200 Number,
201 Date,
202 Datetime,
203 Text,
204 Boolean,
205 Other,
206}
207
208impl ValueKind {
209 pub fn of(
212 dtype: &DataType,
213 number: Option<NumberReading>,
214 time: Option<&TimeInterpretation>,
215 ) -> Self {
216 if let Some(time) = time {
217 return match time.kind {
218 TimeKind::Date => Self::Date,
219 TimeKind::Datetime => Self::Datetime,
220 };
221 }
222 if number.is_some() && is_text(dtype) {
223 return Self::Number;
224 }
225 match dtype {
226 DataType::Date => Self::Date,
227 DataType::Datetime(..) => Self::Datetime,
228 DataType::String | DataType::Categorical(..) => Self::Text,
229 DataType::Boolean => Self::Boolean,
230 dtype if dtype.is_primitive_numeric() || dtype.is_decimal() => Self::Number,
231 _ => Self::Other,
232 }
233 }
234
235 pub fn ranges(self) -> bool {
237 matches!(self, Self::Number | Self::Date | Self::Datetime)
238 }
239
240 pub fn bound_hint(self) -> &'static str {
242 match self {
243 Self::Number => "a number",
244 Self::Date => "a date, 2024-01-31",
245 Self::Datetime => "a date or 2024-01-31 08:00:00",
246 _ => "",
247 }
248 }
249}
250
251fn is_text(dtype: &DataType) -> bool {
252 matches!(dtype, DataType::String | DataType::Categorical(..))
253}
254
255pub fn allows_set(dtype: &DataType) -> bool {
258 is_text(dtype) || dtype.is_integer() || matches!(dtype, DataType::Boolean)
259}
260
261pub fn reads_as_number(dtype: &DataType) -> bool {
263 is_text(dtype)
264}
265
266#[derive(Debug, Clone, Copy, PartialEq)]
269enum Bound {
270 Number(f64),
271 Micros(i64),
272}
273
274impl Bound {
275 fn lit(self) -> Expr {
276 match self {
277 Self::Number(value) => lit(value),
278 Self::Micros(value) => lit(value),
279 }
280 }
281
282 fn value(self) -> f64 {
283 match self {
284 Self::Number(value) => value,
285 Self::Micros(value) => value as f64,
286 }
287 }
288}
289
290fn parse_bound(kind: ValueKind, text: &str, upper: bool) -> std::result::Result<Bound, String> {
293 let text = text.trim();
294 let midnight = |date: chrono::NaiveDate| {
295 date.and_hms_opt(0, 0, 0)
296 .map(|time| Bound::Micros(time.and_utc().timestamp_micros()))
297 };
298 let day_end = |date: chrono::NaiveDate| {
299 date.succ_opt()
300 .and_then(|next| next.and_hms_opt(0, 0, 0))
301 .map(|time| Bound::Micros(time.and_utc().timestamp_micros() - 1))
302 };
303 let parsed = match kind {
304 ValueKind::Number => text
305 .parse::<f64>()
306 .ok()
307 .filter(|value| value.is_finite())
308 .map(Bound::Number),
309 ValueKind::Date => chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d")
310 .ok()
311 .and_then(midnight),
312 ValueKind::Datetime => ["%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M"]
313 .into_iter()
314 .find_map(|format| chrono::NaiveDateTime::parse_from_str(text, format).ok())
315 .map(|time| Bound::Micros(time.and_utc().timestamp_micros()))
316 .or_else(|| {
317 chrono::NaiveDate::parse_from_str(text, "%Y-%m-%d")
318 .ok()
319 .and_then(|date| if upper { day_end(date) } else { midnight(date) })
320 }),
321 _ => None,
322 };
323 parsed.ok_or_else(|| format!("{text:?} is not {}", kind.bound_hint()))
324}
325
326pub fn parse_allowed(dtype: &DataType, text: &str) -> std::result::Result<Vec<String>, String> {
330 let values = split_allowed(text)?;
331 check_allowed(dtype, &values)?;
332 Ok(values)
333}
334
335fn split_allowed(text: &str) -> std::result::Result<Vec<String>, String> {
336 let mut values: Vec<String> = Vec::new();
337 let mut push = |value: String| {
338 if !values.contains(&value) {
339 values.push(value);
340 }
341 };
342 let mut chars = text.chars().peekable();
343 loop {
344 while chars.next_if(|c| c.is_whitespace()).is_some() {}
345 if chars.next_if_eq(&'"').is_some() {
346 let mut value = String::new();
347 loop {
348 match chars.next() {
349 Some('"') if chars.next_if_eq(&'"').is_some() => value.push('"'),
350 Some('"') => break,
351 Some(c) => value.push(c),
352 None => return Err("A quoted value has no closing quote".to_string()),
353 }
354 }
355 while chars.next_if(|c| c.is_whitespace()).is_some() {}
356 match chars.next() {
357 None | Some(',') => push(value),
358 Some(_) => return Err(format!("Put a comma after {value:?}")),
359 }
360 } else {
361 let mut value = String::new();
362 for c in chars.by_ref() {
363 if c == ',' {
364 break;
365 }
366 value.push(c);
367 }
368 let value = value.trim();
369 if !value.is_empty() {
370 push(value.to_string());
371 }
372 }
373 if chars.peek().is_none() {
374 break;
375 }
376 }
377 Ok(values)
378}
379
380fn check_allowed(dtype: &DataType, values: &[String]) -> std::result::Result<(), String> {
382 if dtype.is_integer() || matches!(dtype, DataType::Boolean) {
384 for value in values {
385 crate::typed_value::parse(value, dtype)?;
386 }
387 }
388 if values.len() > MAX_ALLOWED_VALUES {
389 return Err(format!("At most {MAX_ALLOWED_VALUES} allowed values"));
390 }
391 Ok(())
392}
393
394pub fn format_allowed(values: &[String]) -> String {
397 values
398 .iter()
399 .map(|value| {
400 let plain = !value.is_empty()
401 && value.trim() == value
402 && !value.contains(',')
403 && !value.starts_with('"');
404 if plain {
405 value.clone()
406 } else {
407 format!("\"{}\"", value.replace('"', "\"\""))
408 }
409 })
410 .collect::<Vec<_>>()
411 .join(", ")
412}
413
414pub fn check_intent(
417 intent: &ColumnIntent,
418 dtype: &DataType,
419 time: Option<&TimeInterpretation>,
420) -> std::result::Result<(), String> {
421 let kind = ValueKind::of(dtype, intent.number, time);
422 if !intent.allowed.is_empty() {
423 if !allows_set(dtype) {
424 return Err("Allowed values are for text, whole numbers or true/false".to_string());
425 }
426 check_allowed(dtype, &intent.allowed)?;
427 }
428 let bound = |text: &Option<String>, upper: bool| {
429 text.as_deref()
430 .map(|text| parse_bound(kind, text, upper))
431 .transpose()
432 };
433 if intent.min.is_some() || intent.max.is_some() {
434 if !kind.ranges() {
435 return Err("A range is for numbers, dates and times".to_string());
436 }
437 if let (Some(min), Some(max)) = (bound(&intent.min, false)?, bound(&intent.max, true)?)
438 && min.value() > max.value()
439 {
440 return Err("Minimum is above maximum".to_string());
441 }
442 }
443 if intent.number.is_some() && !reads_as_number(dtype) {
444 return Err("Only text is read as a number".to_string());
445 }
446 Ok(())
447}
448
449fn within_micros(value: Expr) -> Expr {
454 value.map(
455 |c| {
456 let limit = match c.dtype() {
457 DataType::Date => i64::MAX / 86_400_000_000,
458 DataType::Datetime(TimeUnit::Milliseconds, _) => i64::MAX / 1_000,
459 _ => return Ok(c),
460 };
461 let series = c.as_materialized_series();
462 let stored = series.to_physical_repr().cast(&DataType::Int64)?;
463 let stored = stored.i64()?;
464 let fits = |v: i64| (-limit..=limit).contains(&v);
465 if [stored.min(), stored.max()].into_iter().flatten().all(fits) {
466 return Ok(c);
467 }
468 let held = stored.apply_values(|v| v.clamp(-limit, limit));
469 Ok(held
470 .into_series()
471 .cast(c.dtype())?
472 .with_name(series.name().clone())
473 .into_column())
474 },
475 |_, field| Ok(field.clone()),
476 )
477}
478
479struct Measured<'a> {
482 intent: &'a ColumnIntent,
483 dtype: DataType,
484 time: Option<TimeInterpretation>,
485}
486
487impl Measured<'_> {
488 fn stored(&self) -> Expr {
489 col(self.intent.column.as_str())
490 }
491
492 fn kind(&self) -> ValueKind {
493 ValueKind::of(&self.dtype, self.intent.number, self.time.as_ref())
494 }
495
496 fn value(&self) -> Expr {
500 if let Some(time) = &self.time {
501 return time.expr();
502 }
503 match self.intent.number {
504 Some(number) if is_text(&self.dtype) => {
505 self.stored().cast(DataType::String).cast(number.dtype())
506 }
507 _ => self.stored(),
508 }
509 }
510
511 fn compared(&self) -> Option<Expr> {
514 let value = self.value();
515 Some(match self.kind() {
516 ValueKind::Number => value.cast(DataType::Float64),
517 ValueKind::Date => within_micros(value)
518 .cast(DataType::Datetime(TimeUnit::Microseconds, None))
519 .dt()
520 .timestamp(TimeUnit::Microseconds),
521 ValueKind::Datetime => within_micros(value).dt().timestamp(TimeUnit::Microseconds),
522 _ => return None,
523 })
524 }
525
526 fn bounds(&self) -> (Option<Bound>, Option<Bound>) {
527 let kind = self.kind();
528 let bound = |text: &Option<String>, upper: bool| {
529 text.as_deref()
530 .and_then(|text| parse_bound(kind, text, upper).ok())
531 };
532 (
533 bound(&self.intent.min, false),
534 bound(&self.intent.max, true),
535 )
536 }
537
538 fn below(&self) -> Option<Expr> {
539 Some(self.compared()?.lt(self.bounds().0?.lit()))
540 }
541
542 fn above(&self) -> Option<Expr> {
543 Some(self.compared()?.gt(self.bounds().1?.lit()))
544 }
545
546 fn out_of_range(&self) -> Option<Expr> {
547 match (self.below(), self.above()) {
548 (Some(below), Some(above)) => Some(below.or(above)),
549 (below, above) => below.or(above),
550 }
551 }
552
553 fn in_set(&self) -> Option<Expr> {
557 if self.intent.allowed.is_empty() || !allows_set(&self.dtype) {
558 return None;
559 }
560 let stored = self.stored();
561 self.intent
562 .allowed
563 .iter()
564 .filter_map(|value| {
565 let value = crate::typed_value::parse(value, &self.dtype).ok()?;
566 Some(stored.clone().eq(lit(value)))
567 })
568 .reduce(Expr::or)
569 }
570
571 fn outside(&self) -> Option<Expr> {
572 Some(self.stored().is_not_null().and(self.in_set()?.not()))
573 }
574
575 fn unparsed(&self) -> Option<Expr> {
577 if self.intent.number.is_none() || !is_text(&self.dtype) || self.time.is_some() {
578 return None;
579 }
580 Some(self.stored().is_not_null().and(self.value().is_null()))
581 }
582}
583
584fn measured<'a>(plan: &'a DataQualityPlan, schema: &Schema) -> Vec<Measured<'a>> {
585 plan.intent
586 .columns
587 .iter()
588 .filter_map(|intent| {
589 Some(Measured {
590 intent,
591 dtype: schema.get(&intent.column)?.clone(),
592 time: plan.time_format(&intent.column).cloned(),
593 })
594 })
595 .collect()
596}
597
598fn name(index: usize, what: &str) -> String {
599 format!("{PREFIX}{index}::{what}")
600}
601
602fn key_columns(plan: &DataQualityPlan, schema: &Schema) -> Option<Vec<String>> {
604 let key = &plan.intent.key;
605 (!key.is_empty() && key.iter().all(|column| schema.get(column).is_some())).then(|| key.clone())
606}
607
608fn any_null(key: &[String]) -> Expr {
609 key.iter()
610 .map(|column| col(column.as_str()).is_null())
611 .reduce(Expr::or)
612 .unwrap_or_else(|| lit(false))
613}
614
615pub(crate) fn intent_exprs(plan: &DataQualityPlan, schema: &Schema) -> Vec<Expr> {
618 let mut exprs = Vec::new();
619 for (index, rules) in measured(plan, schema).iter().enumerate() {
620 exprs.push(
621 rules
622 .stored()
623 .is_not_null()
624 .sum()
625 .alias(name(index, "values")),
626 );
627 if rules.intent.required {
628 exprs.push(rules.stored().is_null().sum().alias(name(index, "missing")));
629 }
630 if let Some(unparsed) = rules.unparsed() {
631 exprs.push(unparsed.sum().alias(name(index, "unparsed")));
632 }
633 if let Some(outside) = rules.outside() {
634 exprs.push(outside.sum().alias(name(index, "outside")));
635 }
636 if let Some(compared) = rules.compared().filter(|_| rules.out_of_range().is_some()) {
637 exprs.push(compared.is_not_null().sum().alias(name(index, "compared")));
638 let value = rules.value();
639 if let Some(below) = rules.below() {
640 exprs.push(below.clone().sum().alias(name(index, "below")));
641 exprs.push(
642 value
643 .clone()
644 .filter(below)
645 .min()
646 .alias(name(index, "lowest")),
647 );
648 }
649 if let Some(above) = rules.above() {
650 exprs.push(above.clone().sum().alias(name(index, "above")));
651 exprs.push(value.filter(above).max().alias(name(index, "highest")));
652 }
653 }
654 }
655 if let Some(key) = key_columns(plan, schema) {
656 exprs.push(any_null(&key).sum().alias(format!("{PREFIX}key::missing")));
657 }
658 exprs
659}
660
661pub(crate) fn key_repeats(
667 lf: &LazyFrame,
668 plan: &DataQualityPlan,
669 schema: &Schema,
670 polars_streaming: bool,
671) -> Result<Option<(usize, usize, usize)>> {
672 let Some(key) = key_columns(plan, schema) else {
673 return Ok(None);
674 };
675 const COUNT: &str = "__datui_intent_key_rows";
676 let columns = key
677 .iter()
678 .map(|name| col(name.as_str()))
679 .collect::<Vec<_>>();
680 let query = lf
681 .clone()
682 .select(columns.clone())
683 .filter(any_null(&key).not())
684 .group_by(columns)
685 .agg([len().alias(COUNT)])
686 .filter(col(COUNT).gt(lit(1u32)))
687 .select([
688 len().alias("groups"),
689 (col(COUNT) - lit(1u32)).sum().alias("extra"),
690 col(COUNT).sum().alias("involved"),
691 ]);
692 let summary = collect_lazy(query, polars_streaming)?;
693 Ok(Some((
694 count_at(&summary, "groups").unwrap_or(0),
695 count_at(&summary, "extra").unwrap_or(0),
696 count_at(&summary, "involved").unwrap_or(0),
697 )))
698}
699
700fn count_at(df: &DataFrame, name: &str) -> Option<usize> {
701 match df.column(name).ok()?.get(0).ok()? {
702 AnyValue::UInt32(value) => Some(value as usize),
703 AnyValue::UInt64(value) => Some(value as usize),
704 AnyValue::Int32(value) => usize::try_from(value).ok(),
705 AnyValue::Int64(value) => usize::try_from(value).ok(),
706 _ => None,
707 }
708}
709
710fn text_at(df: &DataFrame, name: &str) -> Option<String> {
711 let value = df.column(name).ok()?.get(0).ok()?;
712 (!value.is_null()).then(|| crate::exact::str_value(&value).into_owned())
713}
714
715fn commonest(lf: &LazyFrame, rows: Expr, value: Expr) -> Result<Vec<(String, usize)>> {
718 const VALUE: &str = "__datui_intent_value";
719 const COUNT: &str = "__datui_intent_count";
720 let top = lf
721 .clone()
722 .filter(rows)
723 .select([value.cast(DataType::String).alias(VALUE)])
724 .group_by([col(VALUE)])
725 .agg([len().alias(COUNT)])
726 .sort_by_exprs(
727 [col(COUNT), col(VALUE)],
728 SortMultipleOptions::default().with_order_descending_multi([true, false]),
729 )
730 .limit(MAX_INTENT_EXAMPLES as IdxSize)
731 .collect()?;
732 let (values, counts) = (top.column(VALUE)?, top.column(COUNT)?);
733 Ok((0..top.height())
734 .filter_map(|row| {
735 let value = values.get(row).ok()?;
736 let count = match counts.get(row).ok()? {
737 AnyValue::UInt32(count) => count as usize,
738 AnyValue::UInt64(count) => count as usize,
739 _ => return None,
740 };
741 Some((crate::exact::str_value(&value).into_owned(), count))
742 })
743 .collect())
744}
745
746#[derive(Debug, Clone, PartialEq, Eq)]
748pub struct KeyCheck {
749 pub columns: Vec<String>,
750 pub missing: usize,
752 pub groups: usize,
754 pub extra_rows: usize,
756 pub rows_involved: usize,
758}
759
760#[derive(Debug, Clone)]
762pub struct ColumnCheck {
763 pub intent: ColumnIntent,
764 pub dtype: DataType,
766 pub time: Option<TimeInterpretation>,
768 pub values: usize,
770 pub missing: Option<usize>,
772 pub unparsed: Option<usize>,
774 pub outside: Option<usize>,
776 pub compared: Option<usize>,
778 pub below: Option<usize>,
779 pub above: Option<usize>,
780 pub lowest: Option<String>,
782 pub highest: Option<String>,
783 pub outside_examples: Vec<(String, usize)>,
785 pub unparsed_examples: Vec<(String, usize)>,
787}
788
789impl ColumnCheck {
790 fn rules(&self) -> Measured<'_> {
791 Measured {
792 intent: &self.intent,
793 dtype: self.dtype.clone(),
794 time: self.time.clone(),
795 }
796 }
797
798 pub fn out_of_range(&self) -> Option<usize> {
800 match (self.below, self.above) {
801 (None, None) => None,
802 (below, above) => Some(below.unwrap_or(0) + above.unwrap_or(0)),
803 }
804 }
805}
806
807#[derive(Debug, Clone)]
809pub struct IntentResults {
810 pub measured: bool,
812 pub precision: QualityPrecision,
813 pub evaluated_rows: usize,
815 pub key: Option<KeyCheck>,
816 pub columns: Vec<ColumnCheck>,
817 pub declared: Vec<String>,
819 pub absent: Vec<String>,
821}
822
823impl IntentResults {
824 pub(crate) fn unmeasured(plan: &DataQualityPlan, schema: &Schema) -> Option<Self> {
826 if plan.intent.is_empty() {
827 return None;
828 }
829 Some(Self {
830 measured: false,
831 precision: QualityPrecision::Metadata,
832 evaluated_rows: 0,
833 key: None,
834 columns: Vec::new(),
835 declared: declared(plan),
836 absent: absent(plan, schema),
837 })
838 }
839
840 pub(crate) fn from_counts(
844 plan: &DataQualityPlan,
845 schema: &Schema,
846 counts: &DataFrame,
847 repeats: Option<(usize, usize, usize)>,
848 evaluated_rows: usize,
849 precision: QualityPrecision,
850 rows: Option<&LazyFrame>,
851 ) -> Result<Option<Self>> {
852 if plan.intent.is_empty() {
853 return Ok(None);
854 }
855 let mut columns = Vec::new();
856 for (index, rules) in measured(plan, schema).iter().enumerate() {
857 let count = |what: &str| count_at(counts, &name(index, what));
858 let examples = |predicate: Option<Expr>, value: Expr| match (rows, predicate) {
859 (Some(rows), Some(predicate)) => commonest(rows, predicate, value),
860 _ => Ok(Vec::new()),
861 };
862 let unparsed = count("unparsed");
863 let outside = count("outside");
864 columns.push(ColumnCheck {
865 intent: rules.intent.clone(),
866 dtype: rules.dtype.clone(),
867 time: rules.time.clone(),
868 values: count("values").unwrap_or(0),
869 missing: count("missing"),
870 unparsed,
871 outside,
872 compared: count("compared"),
873 below: count("below"),
874 above: count("above"),
875 lowest: text_at(counts, &name(index, "lowest")),
876 highest: text_at(counts, &name(index, "highest")),
877 outside_examples: if outside.unwrap_or(0) > 0 {
878 examples(rules.outside(), rules.stored())?
879 } else {
880 Vec::new()
881 },
882 unparsed_examples: if unparsed.unwrap_or(0) > 0 {
883 examples(rules.unparsed(), rules.stored())?
884 } else {
885 Vec::new()
886 },
887 });
888 }
889 let key = key_columns(plan, schema).map(|key| {
890 let (groups, extra_rows, rows_involved) = repeats.unwrap_or((0, 0, 0));
891 KeyCheck {
892 columns: key,
893 missing: count_at(counts, &format!("{PREFIX}key::missing")).unwrap_or(0),
894 groups,
895 extra_rows,
896 rows_involved,
897 }
898 });
899 Ok(Some(Self {
900 measured: true,
901 precision,
902 evaluated_rows,
903 key,
904 columns,
905 declared: declared(plan),
906 absent: absent(plan, schema),
907 }))
908 }
909
910 pub(crate) fn observations(&self) -> Vec<QualityObservation> {
913 let mut observations = Vec::new();
914 let push = |observations: &mut Vec<QualityObservation>,
915 kind: ObservationKind,
916 column: &str,
917 affected_rows: usize,
918 evaluated_rows: usize| {
919 observations.push(QualityObservation {
920 kind,
921 column: column.to_string(),
922 affected_rows,
923 evaluated_rows,
924 fact: String::new(),
925 normalized_category: None,
926 files: Vec::new(),
927 time_format: None,
928 full_scale: None,
929 });
930 };
931 if !self.measured {
932 return observations;
933 }
934 if let Some(key) = &self.key {
935 for column in &key.columns {
936 if key.rows_involved > 0 {
937 push(
938 &mut observations,
939 ObservationKind::KeyRepeated,
940 column,
941 key.rows_involved,
942 self.evaluated_rows,
943 );
944 }
945 if key.missing > 0 {
946 push(
947 &mut observations,
948 ObservationKind::KeyMissing,
949 column,
950 key.missing,
951 self.evaluated_rows,
952 );
953 }
954 }
955 }
956 for check in &self.columns {
957 let column = check.intent.column.as_str();
958 if let Some(missing) = check.missing.filter(|missing| *missing > 0) {
959 push(
960 &mut observations,
961 ObservationKind::RequiredMissing,
962 column,
963 missing,
964 self.evaluated_rows,
965 );
966 }
967 if let Some(unparsed) = check.unparsed.filter(|unparsed| *unparsed > 0) {
968 push(
969 &mut observations,
970 ObservationKind::UnparsedNumber,
971 column,
972 unparsed,
973 check.values,
974 );
975 }
976 if let Some(outside) = check.outside.filter(|outside| *outside > 0) {
977 push(
978 &mut observations,
979 ObservationKind::NotAllowed,
980 column,
981 outside,
982 check.values,
983 );
984 }
985 if let Some(out) = check.out_of_range().filter(|out| *out > 0) {
986 push(
987 &mut observations,
988 ObservationKind::OutOfRange,
989 column,
990 out,
991 check.compared.unwrap_or(0),
992 );
993 }
994 }
995 observations
996 }
997
998 pub fn repeated_key(&self) -> Option<Expr> {
1000 let key = self.key.as_ref()?;
1001 let columns = key
1002 .columns
1003 .iter()
1004 .map(|name| col(name.as_str()))
1005 .collect::<Vec<_>>();
1006 let shared = len().over(columns).ok()?.gt(lit(1u32));
1007 Some(any_null(&key.columns).not().and(shared))
1008 }
1009
1010 pub fn missing_key(&self) -> Option<Expr> {
1012 Some(any_null(&self.key.as_ref()?.columns))
1013 }
1014
1015 pub fn required_missing(&self, column: &str) -> Option<Expr> {
1017 Some(self.column(column)?.rules().stored().is_null())
1018 }
1019
1020 pub fn unparsed_number(&self, column: &str) -> Option<Expr> {
1022 self.column(column)?.rules().unparsed()
1023 }
1024
1025 pub fn not_allowed(&self, column: &str) -> Option<Expr> {
1027 self.column(column)?.rules().outside()
1028 }
1029
1030 pub fn out_of_range(&self, column: &str) -> Option<Expr> {
1032 self.column(column)?.rules().out_of_range()
1033 }
1034
1035 pub fn column(&self, column: &str) -> Option<&ColumnCheck> {
1037 self.columns
1038 .iter()
1039 .find(|check| check.intent.column == column)
1040 }
1041}
1042
1043fn declared(plan: &DataQualityPlan) -> Vec<String> {
1044 plan.intent
1045 .declared_columns()
1046 .into_iter()
1047 .map(str::to_string)
1048 .collect()
1049}
1050
1051fn absent(plan: &DataQualityPlan, schema: &Schema) -> Vec<String> {
1052 plan.intent
1053 .declared_columns()
1054 .into_iter()
1055 .filter(|column| schema.get(column).is_none())
1056 .map(str::to_string)
1057 .collect()
1058}
1059
1060pub(crate) fn supersede(observations: &mut Vec<QualityObservation>, plan: &DataQualityPlan) {
1063 if let [key] = plan.intent.key.as_slice() {
1064 observations.retain(|observation| {
1065 !(observation.kind == ObservationKind::KeyLike && &observation.column == key)
1066 });
1067 }
1068}
1069
1070#[cfg(test)]
1071mod tests;