1use std::collections::{BTreeMap, VecDeque};
22use std::fmt;
23
24use async_trait::async_trait;
25use serde::{Deserialize, Serialize};
26use thiserror::Error;
27
28#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
34#[serde(transparent)]
35pub struct SchemaId(pub String);
36
37impl SchemaId {
38 #[must_use]
40 pub fn as_str(&self) -> &str {
41 &self.0
42 }
43}
44
45impl fmt::Display for SchemaId {
46 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
47 f.write_str(&self.0)
48 }
49}
50
51impl From<String> for SchemaId {
52 fn from(s: String) -> Self {
53 Self(s)
54 }
55}
56
57impl From<&str> for SchemaId {
58 fn from(s: &str) -> Self {
59 Self(s.to_string())
60 }
61}
62
63#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
65#[serde(transparent)]
66pub struct FieldPath(pub String);
67
68impl FieldPath {
69 #[must_use]
71 pub fn as_str(&self) -> &str {
72 &self.0
73 }
74}
75
76impl fmt::Display for FieldPath {
77 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
78 f.write_str(&self.0)
79 }
80}
81
82impl From<String> for FieldPath {
83 fn from(s: String) -> Self {
84 Self(s)
85 }
86}
87
88impl From<&str> for FieldPath {
89 fn from(s: &str) -> Self {
90 Self(s.to_string())
91 }
92}
93
94#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
100#[serde(tag = "kind", rename_all = "snake_case")]
101pub enum FieldValue {
102 Number(f64),
104 Text(String),
106 OrderedList(Vec<String>),
109 TimestampSeconds(i64),
112 Bool(bool),
114 Null,
118}
119
120impl FieldValue {
121 #[must_use]
123 pub const fn is_number(&self) -> bool {
124 matches!(self, Self::Number(_))
125 }
126
127 #[must_use]
129 pub const fn is_text(&self) -> bool {
130 matches!(self, Self::Text(_))
131 }
132
133 #[must_use]
135 pub const fn is_ordered_list(&self) -> bool {
136 matches!(self, Self::OrderedList(_))
137 }
138
139 #[must_use]
141 pub const fn is_timestamp(&self) -> bool {
142 matches!(self, Self::TimestampSeconds(_))
143 }
144}
145
146#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
149pub struct SchemaBaseline {
150 pub numeric_history: BTreeMap<FieldPath, VecDeque<f64>>,
153 pub list_history: BTreeMap<FieldPath, VecDeque<Vec<String>>>,
155 pub cardinality_history: BTreeMap<FieldPath, VecDeque<usize>>,
158 pub schema_id: Option<SchemaId>,
160}
161
162#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
164pub struct AnomalyReport {
165 pub field: FieldPath,
167 pub schema_id: SchemaId,
169 pub signal: AnomalySignal,
171 pub severity: AnomalySeverity,
173 pub reason: String,
175}
176
177#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
179#[serde(tag = "kind", rename_all = "snake_case")]
180pub enum AnomalySignal {
181 PriceDrift {
184 ratio: f64,
186 },
187 ListingReorder {
190 jaccard: f64,
192 },
193 Staleness {
196 age_secs: i64,
198 },
199 Outlier {
202 z_score: f64,
204 },
205 CardinalityShift {
208 from: usize,
210 to: usize,
212 },
213 None,
216}
217
218#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
220#[serde(rename_all = "snake_case")]
221pub enum AnomalySeverity {
222 Info,
225 Warning,
227 Error,
229}
230
231#[derive(Debug, Error)]
233pub enum AnomalyError {
234 #[error("anomaly detector init failed: {0}")]
236 Init(String),
237 #[error("observation failed: {0}")]
239 Observe(String),
240}
241
242#[async_trait]
244pub trait FieldAnomalyDetector: Send + Sync {
245 fn name(&self) -> &'static str;
247
248 async fn observe(
257 &self,
258 schema_id: &SchemaId,
259 field: &FieldPath,
260 value: &FieldValue,
261 ) -> Result<AnomalyReport, AnomalyError>;
262
263 async fn baseline(&self, schema_id: &SchemaId) -> Result<SchemaBaseline, AnomalyError>;
265
266 async fn reset(&self, schema_id: &SchemaId) -> Result<(), AnomalyError>;
269}
270
271#[derive(Debug)]
277pub struct StatisticalFieldAnomalyDetector {
278 pub window_size: usize,
280 pub iqr_multiplier: f64,
282 pub jaccard_threshold: f64,
284 pub staleness_threshold_secs: i64,
286 pub z_score_threshold: f64,
288 pub cardinality_shift_fraction: f64,
290 state: parking_lot::Mutex<BTreeMap<SchemaId, SchemaBaseline>>,
292}
293
294impl Default for StatisticalFieldAnomalyDetector {
295 fn default() -> Self {
296 Self {
297 window_size: Self::WINDOW_SIZE,
298 iqr_multiplier: 3.0,
299 jaccard_threshold: 0.3,
300 staleness_threshold_secs: 7 * 24 * 60 * 60,
301 z_score_threshold: 3.0,
302 cardinality_shift_fraction: 0.5,
303 state: parking_lot::Mutex::new(BTreeMap::new()),
304 }
305 }
306}
307
308impl StatisticalFieldAnomalyDetector {
309 pub const WINDOW_SIZE: usize = 64;
311
312 #[must_use]
314 pub fn new() -> Self {
315 Self::default()
316 }
317
318 fn usize_to_f64(n: usize) -> f64 {
325 let n32 = u32::try_from(n).unwrap_or(u32::MAX);
326 f64::from(n32)
327 }
328
329 fn cap_window<T>(window: &mut VecDeque<T>, max: usize) {
334 while window.len() > max {
335 window.pop_front();
336 }
337 }
338
339 fn median(window: &VecDeque<f64>) -> Option<f64> {
341 if window.is_empty() {
342 return None;
343 }
344 let mut sorted: Vec<f64> = window.iter().copied().collect();
345 sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
346 let mid = sorted.len() / 2;
347 if sorted.len().is_multiple_of(2) {
348 let a = sorted.get(mid - 1).copied()?;
349 let b = sorted.get(mid).copied()?;
350 Some(f64::midpoint(a, b))
351 } else {
352 Some(*sorted.get(mid)?)
353 }
354 }
355
356 fn quartiles(window: &VecDeque<f64>) -> Option<(f64, f64)> {
360 if window.len() < 4 {
361 return None;
362 }
363 let mut sorted: Vec<f64> = window.iter().copied().collect();
364 sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
365 let mid = sorted.len() / 2;
366 let (lower_slice, upper_with_mid) = sorted.split_at(mid);
368 let q1 = Self::pick_midpoint(lower_slice);
369 let q3 = Self::pick_midpoint(upper_with_mid);
370 Some((q1?, q3?))
371 }
372
373 fn pick_midpoint(sorted: &[f64]) -> Option<f64> {
376 if sorted.is_empty() {
377 return None;
378 }
379 let mid = sorted.len() / 2;
380 if sorted.len().is_multiple_of(2) {
381 let a = sorted.get(mid - 1).copied()?;
382 let b = sorted.get(mid).copied()?;
383 Some(f64::midpoint(a, b))
384 } else {
385 Some(*sorted.get(mid)?)
386 }
387 }
388
389 fn mean(window: &VecDeque<f64>) -> Option<f64> {
391 if window.is_empty() {
392 return None;
393 }
394 Some(window.iter().sum::<f64>() / Self::usize_to_f64(window.len()))
395 }
396
397 fn stddev(window: &VecDeque<f64>) -> Option<f64> {
400 if window.len() < 2 {
401 return None;
402 }
403 let m = Self::mean(window)?;
404 let n = Self::usize_to_f64(window.len());
405 let variance = window.iter().map(|v| (v - m).powi(2)).sum::<f64>() / (n - 1.0);
406 Some(variance.sqrt())
407 }
408
409 fn jaccard_distance(a: &[String], b: &[String]) -> f64 {
414 if a.is_empty() && b.is_empty() {
415 return 0.0;
416 }
417 let sa: std::collections::BTreeSet<&str> = a.iter().map(String::as_str).collect();
418 let sb: std::collections::BTreeSet<&str> = b.iter().map(String::as_str).collect();
419 let intersection = sa.intersection(&sb).count();
420 let union = sa.union(&sb).count();
421 if union == 0 {
422 0.0
423 } else {
424 1.0 - (Self::usize_to_f64(intersection) / Self::usize_to_f64(union))
425 }
426 }
427
428 fn evaluate(
431 baseline: &mut SchemaBaseline,
432 schema_id: &SchemaId,
433 field: &FieldPath,
434 value: &FieldValue,
435 cfg: &TuningConfig,
436 now_secs: i64,
437 ) -> AnomalyReport {
438 match value {
439 FieldValue::Number(n) => Self::evaluate_number(baseline, schema_id, field, *n, cfg),
440 FieldValue::OrderedList(items) => {
441 Self::evaluate_list(baseline, schema_id, field, items, cfg)
442 }
443 FieldValue::Text(t) => Self::evaluate_text(baseline, schema_id, field, t, cfg),
444 FieldValue::TimestampSeconds(ts) => {
445 Self::evaluate_timestamp(schema_id, field, *ts, now_secs, cfg)
446 }
447 FieldValue::Bool(_) | FieldValue::Null => AnomalyReport {
448 field: field.clone(),
449 schema_id: schema_id.clone(),
450 signal: AnomalySignal::None,
451 severity: AnomalySeverity::Info,
452 reason: "value type not subject to statistical checks".to_string(),
453 },
454 }
455 }
456
457 fn evaluate_number(
458 baseline: &mut SchemaBaseline,
459 schema_id: &SchemaId,
460 field: &FieldPath,
461 n: f64,
462 cfg: &TuningConfig,
463 ) -> AnomalyReport {
464 let window = baseline.numeric_history.entry(field.clone()).or_default();
465 if let Some((q1, q3)) = Self::quartiles(window) {
466 let iqr = q3 - q1;
467 if iqr > 0.0 {
468 let median = Self::median(window).unwrap_or(n);
469 let deviation = (n - median).abs();
470 if deviation > cfg.iqr_multiplier * iqr {
471 let ratio = if median.abs() > f64::EPSILON {
472 n / median
473 } else {
474 1.0
475 };
476 window.push_back(n);
477 Self::cap_window(window, cfg.window_size);
478 return AnomalyReport {
479 field: field.clone(),
480 schema_id: schema_id.clone(),
481 signal: AnomalySignal::PriceDrift { ratio },
482 severity: AnomalySeverity::Warning,
483 reason: format!(
484 "value {n} deviates {deviation:.2} from median {median:.2} \
485 (IQR={iqr:.2}, k={:.1})",
486 cfg.iqr_multiplier
487 ),
488 };
489 }
490 }
491 if let Some(sd) = Self::stddev(window) {
493 let mean = Self::mean(window).unwrap_or(n);
494 if sd > 0.0 {
495 let z = (n - mean) / sd;
496 if z.abs() > cfg.z_score_threshold {
497 window.push_back(n);
498 Self::cap_window(window, cfg.window_size);
499 return AnomalyReport {
500 field: field.clone(),
501 schema_id: schema_id.clone(),
502 signal: AnomalySignal::Outlier { z_score: z },
503 severity: AnomalySeverity::Warning,
504 reason: format!(
505 "value {n} has z-score {z:.2} (mean={mean:.2}, sd={sd:.2})"
506 ),
507 };
508 }
509 }
510 }
511 }
512 window.push_back(n);
513 Self::cap_window(window, cfg.window_size);
514 AnomalyReport {
515 field: field.clone(),
516 schema_id: schema_id.clone(),
517 signal: AnomalySignal::None,
518 severity: AnomalySeverity::Info,
519 reason: "within IQR and z-score thresholds".to_string(),
520 }
521 }
522
523 fn evaluate_list(
524 baseline: &mut SchemaBaseline,
525 schema_id: &SchemaId,
526 field: &FieldPath,
527 items: &[String],
528 cfg: &TuningConfig,
529 ) -> AnomalyReport {
530 let window = baseline.list_history.entry(field.clone()).or_default();
531 let report = window.back().map_or_else(
532 || AnomalyReport {
533 field: field.clone(),
534 schema_id: schema_id.clone(),
535 signal: AnomalySignal::None,
536 severity: AnomalySeverity::Info,
537 reason: "first observation establishes baseline".to_string(),
538 },
539 |prev| {
540 let distance = Self::jaccard_distance(prev, items);
541 if distance > cfg.jaccard_threshold {
542 AnomalyReport {
543 field: field.clone(),
544 schema_id: schema_id.clone(),
545 signal: AnomalySignal::ListingReorder { jaccard: distance },
546 severity: AnomalySeverity::Warning,
547 reason: format!(
548 "Jaccard distance {distance:.2} exceeds threshold {:.2}",
549 cfg.jaccard_threshold
550 ),
551 }
552 } else {
553 AnomalyReport {
554 field: field.clone(),
555 schema_id: schema_id.clone(),
556 signal: AnomalySignal::None,
557 severity: AnomalySeverity::Info,
558 reason: format!("Jaccard distance {distance:.2} within threshold"),
559 }
560 }
561 },
562 );
563 window.push_back(items.to_vec());
564 Self::cap_window(window, cfg.window_size);
565 report
566 }
567
568 fn evaluate_text(
569 baseline: &mut SchemaBaseline,
570 schema_id: &SchemaId,
571 field: &FieldPath,
572 text: &str,
573 cfg: &TuningConfig,
574 ) -> AnomalyReport {
575 let cardinality = text.chars().filter(|c| !c.is_whitespace()).count();
576 let window = baseline
577 .cardinality_history
578 .entry(field.clone())
579 .or_default();
580 let report = if window.len() >= 2 {
581 let previous = window.back().copied().unwrap_or(cardinality);
582 let previous_window_avg = if window.len() >= cfg.window_size {
583 Self::usize_to_f64(window.iter().sum::<usize>()) / Self::usize_to_f64(window.len())
584 } else {
585 Self::usize_to_f64(previous)
586 };
587 let current = Self::usize_to_f64(cardinality);
588 let drift = if previous_window_avg > 0.0 {
589 (current - previous_window_avg).abs() / previous_window_avg
590 } else {
591 0.0
592 };
593 if drift > cfg.cardinality_shift_fraction {
594 AnomalyReport {
595 field: field.clone(),
596 schema_id: schema_id.clone(),
597 signal: AnomalySignal::CardinalityShift {
598 from: previous,
599 to: cardinality,
600 },
601 severity: AnomalySeverity::Warning,
602 reason: format!(
603 "cardinality shifted {drift:.0}% ({previous} -> {cardinality})"
604 ),
605 }
606 } else {
607 AnomalyReport {
608 field: field.clone(),
609 schema_id: schema_id.clone(),
610 signal: AnomalySignal::None,
611 severity: AnomalySeverity::Info,
612 reason: "cardinality within threshold".to_string(),
613 }
614 }
615 } else {
616 AnomalyReport {
617 field: field.clone(),
618 schema_id: schema_id.clone(),
619 signal: AnomalySignal::None,
620 severity: AnomalySeverity::Info,
621 reason: "establishing baseline".to_string(),
622 }
623 };
624 window.push_back(cardinality);
625 Self::cap_window(window, cfg.window_size);
626 report
627 }
628
629 fn evaluate_timestamp(
630 schema_id: &SchemaId,
631 field: &FieldPath,
632 ts: i64,
633 now_secs: i64,
634 cfg: &TuningConfig,
635 ) -> AnomalyReport {
636 let age = (now_secs - ts).max(0);
637 if age > cfg.staleness_threshold_secs {
638 AnomalyReport {
639 field: field.clone(),
640 schema_id: schema_id.clone(),
641 signal: AnomalySignal::Staleness { age_secs: age },
642 severity: AnomalySeverity::Warning,
643 reason: format!(
644 "age {age}s exceeds staleness threshold {}s",
645 cfg.staleness_threshold_secs
646 ),
647 }
648 } else {
649 AnomalyReport {
650 field: field.clone(),
651 schema_id: schema_id.clone(),
652 signal: AnomalySignal::None,
653 severity: AnomalySeverity::Info,
654 reason: format!("age {age}s within threshold"),
655 }
656 }
657 }
658}
659
660#[derive(Debug, Clone, Copy)]
664struct TuningConfig {
665 window_size: usize,
666 iqr_multiplier: f64,
667 jaccard_threshold: f64,
668 staleness_threshold_secs: i64,
669 z_score_threshold: f64,
670 cardinality_shift_fraction: f64,
671}
672
673impl TuningConfig {
674 const fn from_detector(d: &StatisticalFieldAnomalyDetector) -> Self {
675 Self {
676 window_size: d.window_size,
677 iqr_multiplier: d.iqr_multiplier,
678 jaccard_threshold: d.jaccard_threshold,
679 staleness_threshold_secs: d.staleness_threshold_secs,
680 z_score_threshold: d.z_score_threshold,
681 cardinality_shift_fraction: d.cardinality_shift_fraction,
682 }
683 }
684}
685
686#[async_trait]
687impl FieldAnomalyDetector for StatisticalFieldAnomalyDetector {
688 fn name(&self) -> &'static str {
689 "statistical"
690 }
691
692 async fn observe(
693 &self,
694 schema_id: &SchemaId,
695 field: &FieldPath,
696 value: &FieldValue,
697 ) -> Result<AnomalyReport, AnomalyError> {
698 let cfg = TuningConfig::from_detector(self);
699 let now = chrono::Utc::now();
700 let now_secs = now.timestamp();
701 let mut guard = self.state.lock();
702 let baseline = guard.entry(schema_id.clone()).or_default();
703 if baseline.schema_id.is_none() {
704 baseline.schema_id = Some(schema_id.clone());
705 }
706 if baseline.schema_id.as_ref() != Some(schema_id) {
707 *baseline = SchemaBaseline {
708 schema_id: Some(schema_id.clone()),
709 ..Default::default()
710 };
711 }
712 let report = Self::evaluate(baseline, schema_id, field, value, &cfg, now_secs);
713 drop(guard);
714 Ok(report)
715 }
716
717 async fn baseline(&self, schema_id: &SchemaId) -> Result<SchemaBaseline, AnomalyError> {
718 Ok(self
719 .state
720 .lock()
721 .get(schema_id)
722 .cloned()
723 .unwrap_or_else(|| SchemaBaseline {
724 schema_id: Some(schema_id.clone()),
725 ..Default::default()
726 }))
727 }
728
729 async fn reset(&self, schema_id: &SchemaId) -> Result<(), AnomalyError> {
730 self.state.lock().remove(schema_id);
731 Ok(())
732 }
733}
734
735#[cfg(test)]
738#[allow(
739 clippy::unwrap_used,
740 clippy::expect_used,
741 clippy::panic,
742 clippy::indexing_slicing
743)]
744mod tests {
745 use super::*;
746
747 fn det() -> StatisticalFieldAnomalyDetector {
748 StatisticalFieldAnomalyDetector::new()
749 }
750
751 #[tokio::test]
752 async fn first_observation_is_baseline() {
753 let d = det();
754 let report = d
755 .observe(&"s1".into(), &"price".into(), &FieldValue::Number(9.99))
756 .await
757 .unwrap();
758 assert!(matches!(report.signal, AnomalySignal::None));
759 }
760
761 #[tokio::test]
762 async fn price_drift_far_from_median_triggers() {
763 let d = det();
764 for v in [9.5, 10.0, 10.0, 10.5, 9.8, 10.2, 9.9, 10.1] {
766 d.observe(&"s1".into(), &"price".into(), &FieldValue::Number(v))
767 .await
768 .unwrap();
769 }
770 let report = d
772 .observe(&"s1".into(), &"price".into(), &FieldValue::Number(99.0))
773 .await
774 .unwrap();
775 assert!(
776 matches!(
777 report.signal,
778 AnomalySignal::PriceDrift { .. } | AnomalySignal::Outlier { .. }
779 ),
780 "expected PriceDrift or Outlier, got {:?}",
781 report.signal
782 );
783 }
784
785 #[tokio::test]
786 async fn value_within_iqr_passes() {
787 let d = det();
788 for v in [9.5, 10.0, 10.0, 10.5, 9.8, 10.2, 9.9, 10.1] {
789 let r = d
790 .observe(&"s1".into(), &"price".into(), &FieldValue::Number(v))
791 .await
792 .unwrap();
793 assert!(
794 matches!(r.signal, AnomalySignal::None),
795 "in-range value should pass: {r:?}"
796 );
797 }
798 }
799
800 #[tokio::test]
801 async fn identical_list_ordering_passes() {
802 let d = det();
803 let items1 = vec!["a".to_string(), "b".to_string(), "c".to_string()];
804 let items2 = vec!["a".to_string(), "b".to_string(), "c".to_string()];
805 d.observe(
806 &"s1".into(),
807 &"listings".into(),
808 &FieldValue::OrderedList(items1),
809 )
810 .await
811 .unwrap();
812 let r = d
813 .observe(
814 &"s1".into(),
815 &"listings".into(),
816 &FieldValue::OrderedList(items2),
817 )
818 .await
819 .unwrap();
820 assert!(matches!(r.signal, AnomalySignal::None));
821 }
822
823 #[tokio::test]
824 async fn reordered_list_with_high_jaccard_triggers() {
825 let d = det();
826 let items1 = vec![
827 "a".to_string(),
828 "b".to_string(),
829 "c".to_string(),
830 "d".to_string(),
831 "e".to_string(),
832 "f".to_string(),
833 ];
834 let items2 = vec![
836 "x".to_string(),
837 "y".to_string(),
838 "z".to_string(),
839 "w".to_string(),
840 "v".to_string(),
841 "u".to_string(),
842 ];
843 d.observe(
844 &"s1".into(),
845 &"listings".into(),
846 &FieldValue::OrderedList(items1),
847 )
848 .await
849 .unwrap();
850 let r = d
851 .observe(
852 &"s1".into(),
853 &"listings".into(),
854 &FieldValue::OrderedList(items2),
855 )
856 .await
857 .unwrap();
858 assert!(
859 matches!(r.signal, AnomalySignal::ListingReorder { .. }),
860 "expected ListingReorder, got {:?}",
861 r.signal
862 );
863 }
864
865 #[tokio::test]
866 async fn stale_timestamp_triggers() {
867 let d = det();
868 let now = chrono::Utc::now().timestamp();
869 let stale_ts = now - (365 * 24 * 60 * 60); let r = d
871 .observe(
872 &"s1".into(),
873 &"published_at".into(),
874 &FieldValue::TimestampSeconds(stale_ts),
875 )
876 .await
877 .unwrap();
878 assert!(matches!(r.signal, AnomalySignal::Staleness { .. }));
879 }
880
881 #[tokio::test]
882 async fn outlier_z_score_triggers() {
883 let d = det();
884 for v in [10.0, 10.1, 9.9, 10.0, 10.05, 9.95] {
886 d.observe(&"s1".into(), &"x".into(), &FieldValue::Number(v))
887 .await
888 .unwrap();
889 }
890 let r = d
891 .observe(&"s1".into(), &"x".into(), &FieldValue::Number(50.0))
892 .await
893 .unwrap();
894 assert!(
895 matches!(
896 r.signal,
897 AnomalySignal::Outlier { .. } | AnomalySignal::PriceDrift { .. }
898 ),
899 "expected Outlier or PriceDrift, got {:?}",
900 r.signal
901 );
902 }
903
904 #[tokio::test]
905 async fn cardinality_shift_triggers() {
906 let d = det();
907 for _ in 0..5 {
909 d.observe(
910 &"s1".into(),
911 &"title".into(),
912 &FieldValue::Text("hello world".to_string()),
913 )
914 .await
915 .unwrap();
916 }
917 let huge = "x".repeat(10_000);
919 let r = d
920 .observe(&"s1".into(), &"title".into(), &FieldValue::Text(huge))
921 .await
922 .unwrap();
923 assert!(
924 matches!(r.signal, AnomalySignal::CardinalityShift { .. }),
925 "expected CardinalityShift, got {:?}",
926 r.signal
927 );
928 }
929
930 #[tokio::test]
931 async fn detector_resets_on_schema_change() {
932 let d = det();
933 d.observe(&"s1".into(), &"price".into(), &FieldValue::Number(100.0))
935 .await
936 .unwrap();
937 let r = d
940 .observe(&"s2".into(), &"price".into(), &FieldValue::Number(9.99))
941 .await
942 .unwrap();
943 assert!(matches!(r.signal, AnomalySignal::None));
944 let b = d.baseline(&"s1".into()).await.unwrap();
946 assert_eq!(b.schema_id, Some("s1".into()));
947 }
948
949 #[tokio::test]
950 async fn reset_clears_schema_state() {
951 let d = det();
952 d.observe(&"s1".into(), &"price".into(), &FieldValue::Number(10.0))
953 .await
954 .unwrap();
955 d.reset(&"s1".into()).await.unwrap();
956 let b = d.baseline(&"s1".into()).await.unwrap();
957 assert!(b.numeric_history.is_empty());
958 }
959
960 #[tokio::test]
961 async fn field_value_type_helpers() {
962 assert!(FieldValue::Number(1.0).is_number());
963 assert!(!FieldValue::Number(1.0).is_text());
964 assert!(FieldValue::Text("x".into()).is_text());
965 assert!(FieldValue::OrderedList(vec!["a".into()]).is_ordered_list());
966 assert!(FieldValue::TimestampSeconds(0).is_timestamp());
967 }
968
969 #[tokio::test]
970 async fn bool_and_null_values_do_not_trigger_anomaly() {
971 let d = det();
972 let r1 = d
973 .observe(&"s1".into(), &"flag".into(), &FieldValue::Bool(true))
974 .await
975 .unwrap();
976 let r2 = d
977 .observe(&"s1".into(), &"x".into(), &FieldValue::Null)
978 .await
979 .unwrap();
980 assert!(matches!(r1.signal, AnomalySignal::None));
981 assert!(matches!(r2.signal, AnomalySignal::None));
982 }
983
984 #[test]
985 fn jaccard_distance_identical_sets_is_zero() {
986 let a = vec!["x".into(), "y".into(), "z".into()];
987 let b = vec!["x".into(), "y".into(), "z".into()];
988 let distance = StatisticalFieldAnomalyDetector::jaccard_distance(&a, &b);
989 assert!(
990 distance.abs() < 1e-9,
991 "identical sets should have distance 0, got {distance}"
992 );
993 }
994
995 #[test]
996 fn jaccard_distance_disjoint_sets_is_one() {
997 let a = vec!["x".into(), "y".into()];
998 let b = vec!["p".into(), "q".into()];
999 assert!((StatisticalFieldAnomalyDetector::jaccard_distance(&a, &b) - 1.0).abs() < 1e-9);
1000 }
1001
1002 #[test]
1003 fn median_odd_count() {
1004 let w: VecDeque<f64> = [3.0, 1.0, 2.0].into_iter().collect();
1005 assert_eq!(StatisticalFieldAnomalyDetector::median(&w), Some(2.0));
1006 }
1007
1008 #[test]
1009 fn median_even_count() {
1010 let w: VecDeque<f64> = [4.0, 1.0, 3.0, 2.0].into_iter().collect();
1011 let m = StatisticalFieldAnomalyDetector::median(&w).unwrap_or(0.0);
1012 assert!(
1013 (m - 2.5).abs() < 1e-9,
1014 "median of even-count window should be 2.5, got {m}"
1015 );
1016 }
1017
1018 #[test]
1019 fn schema_id_round_trip() {
1020 let s: SchemaId = "product-v1".into();
1021 assert_eq!(s.as_str(), "product-v1");
1022 assert_eq!(s.to_string(), "product-v1");
1023 }
1024
1025 #[test]
1026 fn field_path_round_trip() {
1027 let p: FieldPath = "items.0.sku".into();
1028 assert_eq!(p.as_str(), "items.0.sku");
1029 }
1030}