Skip to main content

stygian_charon/
field_anomaly.rs

1//! T107 poisoned-data field-level anomaly detector.
2//!
3//! Catches the "tarpit / poisoned data / silent 200s" pattern from
4//! <https://web-scraping-guide.com/#post-extract>: a target returns
5//! clean `200` responses with subtly wrong field values
6//! (price drift, listing reorder, fabricated rows, stale
7//! snapshots). The detector observes every field value the pipeline
8//! publishes and emits [`AnomalyReport`](crate::field_anomaly::AnomalyReport)s when a value is
9//! statistically inconsistent with the rolling baseline.
10//!
11//! Two pieces:
12//!
13//! - [`FieldAnomalyDetector`](crate::field_anomaly::FieldAnomalyDetector) — the consumer-owned port trait.
14//! - [`StatisticalFieldAnomalyDetector`](crate::field_anomaly::StatisticalFieldAnomalyDetector) — the default adapter
15//!   implementing price-drift, outlier, listing-reorder,
16//!   staleness, and cardinality-shift detection.
17//!
18//! Hidden behind a `field-anomaly` cargo feature in the parent
19//! crate so existing charon users aren't forced to opt in.
20
21use std::collections::{BTreeMap, VecDeque};
22use std::fmt;
23
24use async_trait::async_trait;
25use serde::{Deserialize, Serialize};
26use thiserror::Error;
27
28// ── Value types ─────────────────────────────────────────────────────────────
29
30/// Stable identifier for a schema or data-contract version. Two
31/// observations with different `SchemaId` belong to different
32/// baselines — the detector resets state on schema change.
33#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
34#[serde(transparent)]
35pub struct SchemaId(pub String);
36
37impl SchemaId {
38    /// Borrow the inner string.
39    #[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/// Dot-path to a field within a record (`"price"`, `"items.0.sku"`).
64#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
65#[serde(transparent)]
66pub struct FieldPath(pub String);
67
68impl FieldPath {
69    /// Borrow the inner string.
70    #[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/// One observed field value.
95///
96/// Covers the four kinds of fields the detector reasons about:
97/// numeric (price/outlier), string (cardinality), ordered list
98/// (listing-reorder), and timestamp (staleness).
99#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
100#[serde(tag = "kind", rename_all = "snake_case")]
101pub enum FieldValue {
102    /// Numeric value (price, count, score, etc.).
103    Number(f64),
104    /// String value (sku, title, status, etc.).
105    Text(String),
106    /// Ordered list of strings (a page of listings, a search
107    /// result set, etc.). Used for listing-reorder detection.
108    OrderedList(Vec<String>),
109    /// Timestamp value (`published_at`, `expires_at`, etc.). Stored as
110    /// seconds since the Unix epoch for portability.
111    TimestampSeconds(i64),
112    /// Boolean value.
113    Bool(bool),
114    /// Null / absent — explicitly recorded so the detector can
115    /// distinguish "field missing" from "field present with empty
116    /// value".
117    Null,
118}
119
120impl FieldValue {
121    /// `true` if the value carries numeric content (`Number`).
122    #[must_use]
123    pub const fn is_number(&self) -> bool {
124        matches!(self, Self::Number(_))
125    }
126
127    /// `true` if the value carries a text payload.
128    #[must_use]
129    pub const fn is_text(&self) -> bool {
130        matches!(self, Self::Text(_))
131    }
132
133    /// `true` if the value is an ordered list of strings.
134    #[must_use]
135    pub const fn is_ordered_list(&self) -> bool {
136        matches!(self, Self::OrderedList(_))
137    }
138
139    /// `true` if the value is a timestamp.
140    #[must_use]
141    pub const fn is_timestamp(&self) -> bool {
142        matches!(self, Self::TimestampSeconds(_))
143    }
144}
145
146/// Per-field rolling baseline — the statistical summary the
147/// detector compares new observations against.
148#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
149pub struct SchemaBaseline {
150    /// Last N values observed for each numeric field, capped at
151    /// `StatisticalFieldAnomalyDetector::WINDOW_SIZE`.
152    pub numeric_history: BTreeMap<FieldPath, VecDeque<f64>>,
153    /// Last N ordered lists observed for each list-valued field.
154    pub list_history: BTreeMap<FieldPath, VecDeque<Vec<String>>>,
155    /// Last N cardinality counts observed for each text-valued
156    /// field. Used for [`AnomalySignal::CardinalityShift`].
157    pub cardinality_history: BTreeMap<FieldPath, VecDeque<usize>>,
158    /// The schema this baseline belongs to.
159    pub schema_id: Option<SchemaId>,
160}
161
162/// One anomaly emitted by [`FieldAnomalyDetector::observe`].
163#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
164pub struct AnomalyReport {
165    /// Field path the anomaly was observed on.
166    pub field: FieldPath,
167    /// Schema the field belongs to (for downstream filtering).
168    pub schema_id: SchemaId,
169    /// What kind of anomaly.
170    pub signal: AnomalySignal,
171    /// Severity rating (Info / Warning / Error).
172    pub severity: AnomalySeverity,
173    /// Human-readable explanation.
174    pub reason: String,
175}
176
177/// The five anomaly kinds the detector emits.
178#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
179#[serde(tag = "kind", rename_all = "snake_case")]
180pub enum AnomalySignal {
181    /// Numeric value lies more than `k * IQR` from the rolling
182    /// median. `ratio` is the observed / median ratio.
183    PriceDrift {
184        /// Observed / median.
185        ratio: f64,
186    },
187    /// Ordered-list Jaccard distance vs the previous ordering is
188    /// above the threshold. `jaccard` is the distance in `[0.0, 1.0]`.
189    ListingReorder {
190        /// Distance from the previous observation.
191        jaccard: f64,
192    },
193    /// Timestamp is older than the staleness threshold. `age_secs`
194    /// is `now - observed`.
195    Staleness {
196        /// Seconds since the published timestamp.
197        age_secs: i64,
198    },
199    /// Numeric z-score across the rolling window exceeds the
200    /// outlier threshold (default 3.0).
201    Outlier {
202        /// z-score of the observation.
203        z_score: f64,
204    },
205    /// Text-field unique-value count changed by more than the
206    /// cardinality-shift threshold (default 50%) between windows.
207    CardinalityShift {
208        /// Previous window's unique-count.
209        from: usize,
210        /// Current window's unique-count.
211        to: usize,
212    },
213    /// No anomaly — included for completeness so callers can use a
214    /// single return type.
215    None,
216}
217
218/// Severity of an [`AnomalyReport`].
219#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
220#[serde(rename_all = "snake_case")]
221pub enum AnomalySeverity {
222    /// Informational — the value is unusual but not necessarily
223    /// wrong.
224    Info,
225    /// Warning — likely anomaly; should be reviewed.
226    Warning,
227    /// Error — definitely anomaly; the record should be flagged.
228    Error,
229}
230
231/// Errors raised by [`FieldAnomalyDetector`] operations.
232#[derive(Debug, Error)]
233pub enum AnomalyError {
234    /// The detector could not be configured or initialised.
235    #[error("anomaly detector init failed: {0}")]
236    Init(String),
237    /// Observation failed (e.g. lock contention).
238    #[error("observation failed: {0}")]
239    Observe(String),
240}
241
242/// Port trait: observe a field value, baseline reset on schema change.
243#[async_trait]
244pub trait FieldAnomalyDetector: Send + Sync {
245    /// Stable name for diagnostics.
246    fn name(&self) -> &'static str;
247
248    /// Observe one field value. Returns [`AnomalyReport::signal`]
249    /// describing any anomaly detected, or [`AnomalySignal::None`]
250    /// if the value is consistent with the rolling baseline.
251    ///
252    /// # Errors
253    ///
254    /// Returns [`AnomalyError::Observe`] if the detector cannot
255    /// record the observation.
256    async fn observe(
257        &self,
258        schema_id: &SchemaId,
259        field: &FieldPath,
260        value: &FieldValue,
261    ) -> Result<AnomalyReport, AnomalyError>;
262
263    /// Return the current baseline for a schema.
264    async fn baseline(&self, schema_id: &SchemaId) -> Result<SchemaBaseline, AnomalyError>;
265
266    /// Drop the rolling state for a schema. Used by callers that
267    /// want to force a fresh baseline (e.g. after a backfill).
268    async fn reset(&self, schema_id: &SchemaId) -> Result<(), AnomalyError>;
269}
270
271// ── Default adapter ────────────────────────────────────────────────────────
272
273/// Default statistical detector. Rolling-window-based: keeps the
274/// last [`Self::WINDOW_SIZE`] observations per field per schema and
275/// applies the five heuristics in [`AnomalySignal`].
276#[derive(Debug)]
277pub struct StatisticalFieldAnomalyDetector {
278    /// Sliding-window size for each (schema, field) pair.
279    pub window_size: usize,
280    /// IQR multiplier for [`AnomalySignal::PriceDrift`]. Default 3.0.
281    pub iqr_multiplier: f64,
282    /// Jaccard-distance threshold for [`AnomalySignal::ListingReorder`]. Default 0.3.
283    pub jaccard_threshold: f64,
284    /// Staleness threshold in seconds for [`AnomalySignal::Staleness`]. Default 7 days.
285    pub staleness_threshold_secs: i64,
286    /// z-score threshold for [`AnomalySignal::Outlier`]. Default 3.0.
287    pub z_score_threshold: f64,
288    /// Cardinality-shift threshold as a fraction (0.5 = 50%).
289    pub cardinality_shift_fraction: f64,
290    /// Internal baseline state.
291    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    /// Default sliding-window size.
310    pub const WINDOW_SIZE: usize = 64;
311
312    /// Construct a new detector with the default tuning constants.
313    #[must_use]
314    pub fn new() -> Self {
315        Self::default()
316    }
317
318    /// Convert a `usize` count to `f64` for ratio math.
319    ///
320    /// Counts flowing through this helper are bounded by
321    /// [`Self::WINDOW_SIZE`] (max 64) — well within `u32`, so the
322    /// `usize -> u32 -> f64` chain is precision-safe on every
323    /// target.
324    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    /// Trim a window to `max` entries from the front.
330    ///
331    /// Not `const` because `VecDeque::len` / `pop_front` are not
332    /// const-stable yet (Rust 1.96).
333    fn cap_window<T>(window: &mut VecDeque<T>, max: usize) {
334        while window.len() > max {
335            window.pop_front();
336        }
337    }
338
339    /// Median of a window of f64 values. Returns `None` if empty.
340    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    /// First and third quartile of a window.
357    ///
358    /// Returns `None` for windows with fewer than 4 samples.
359    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        // split_at is the safe form of `&v[..n]` + `&v[n..]`.
367        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    /// Pick the median (odd-length) or midpoint (even-length) of an
374    /// already-sorted slice. Returns `None` for empty slices.
375    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    /// Mean of a window.
390    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    /// Standard deviation of a window. Returns `None` when fewer
398    /// than 2 samples.
399    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    /// Jaccard distance between two ordered lists of strings.
410    /// Distance = 1 - |intersection| / |union| for unordered, or
411    /// Kendall-tau-style for ordered. We use set-based Jaccard
412    /// (unordered) per the brief's "Jaccard > 0.3" criterion.
413    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    /// Evaluate one observation against the existing baseline for
429    /// `(schema_id, field)`. Pure logic — no I/O.
430    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            // z-score outlier detection (independent of IQR).
492            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/// Internal tuning snapshot — kept separate from the public
661/// detector struct so the async interface doesn't have to take
662/// `&self` in three places.
663#[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// ── Tests ──────────────────────────────────────────────────────────────────
736
737#[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        // Establish baseline (median ~10)
765        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        // Outlier
771        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        // Completely different set
835        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); // 1 year old
870        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        // Tight cluster
885        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        // Establish baseline of small strings
908        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        // Shift to huge strings
918        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        // Establish s1 baseline.
934        d.observe(&"s1".into(), &"price".into(), &FieldValue::Number(100.0))
935            .await
936            .unwrap();
937        // Switch to s2 — the baseline for s1 is preserved but
938        // s2 starts fresh.
939        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        // s1 baseline should still be retrievable.
945        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}