Skip to main content

faucet_core/profiling/
spec.rs

1//! Config types for the top-level `profiling:` block (#708).
2
3use crate::anomaly::AnomalyMethod;
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6
7/// Default cap on profiled columns per dataset.
8pub const DEFAULT_MAX_COLUMNS: usize = 200;
9/// Default number of most-frequent values kept per categorical column.
10pub const DEFAULT_TOP_VALUES: usize = 10;
11/// Default distinct-count ceiling above which a column is treated as
12/// high-cardinality (no top values are published for it).
13pub const DEFAULT_CATEGORICAL_MAX_DISTINCT: u64 = 100;
14/// Default rolling-window size (successful runs kept as the baseline).
15pub const DEFAULT_WINDOW: u32 = 20;
16/// Default number of baseline runs required before drift detection fires.
17pub const DEFAULT_MIN_HISTORY: u32 = 5;
18/// Default minimum share a previously unseen categorical value must reach
19/// before it is reported as a new value.
20pub const DEFAULT_NEW_VALUE_MIN_SHARE: f64 = 0.05;
21/// Default population-stability-index threshold for categorical drift.
22pub const DEFAULT_PSI_THRESHOLD: f64 = 0.2;
23/// Hard ceiling on `top_values` (bounds the per-column sketch).
24pub const MAX_TOP_VALUES: usize = 1000;
25/// Hard ceiling on `max_columns` (bounds per-run memory).
26pub const MAX_MAX_COLUMNS: usize = 5000;
27
28/// Learned column profiles with drift detection. Every run profiles the
29/// records it wrote (null rate, distinct count, min / max / mean, string
30/// length, type mix, top values) into a bounded-memory sketch, compares the
31/// result against a rolling baseline of earlier successful runs, and reports
32/// statistically significant changes per column — no thresholds to write.
33#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
34#[serde(deny_unknown_fields)]
35pub struct ProfilingSpec {
36    /// Top-level columns to profile. Empty (the default) profiles every
37    /// column the run writes, up to `max_columns`.
38    #[serde(default, skip_serializing_if = "Vec::is_empty")]
39    pub columns: Vec<String>,
40
41    /// Columns to leave out, exact names or `prefix*` globs. Defaults to
42    /// `["_faucet_*"]` so the run-metadata columns are never profiled.
43    #[serde(default = "default_exclude")]
44    pub exclude: Vec<String>,
45
46    /// Maximum columns profiled per run (default 200). Columns beyond the cap
47    /// are skipped and the profile is marked truncated.
48    #[serde(default = "default_max_columns")]
49    pub max_columns: usize,
50
51    /// How many most-frequent values to keep per categorical column (default
52    /// 10). `0` disables top values entirely.
53    #[serde(default = "default_top_values")]
54    pub top_values: usize,
55
56    /// A column whose estimated distinct count exceeds this is treated as
57    /// high-cardinality: only null / type / length metrics are kept for it and
58    /// no values are published (default 100).
59    #[serde(default = "default_categorical_max_distinct")]
60    pub categorical_max_distinct: u64,
61
62    /// Rolling-window size: how many recent successful-run profiles form the
63    /// baseline (default 20; must be ≥ `min_history`).
64    #[serde(default = "default_window")]
65    pub window: u32,
66
67    /// Minimum baseline runs before drift detection starts (default 5; at
68    /// least 2).
69    #[serde(default = "default_min_history")]
70    pub min_history: u32,
71
72    /// How a numeric metric (null rate, distinct count, mean, …) is compared
73    /// against its baseline series. Default `zscore`.
74    #[serde(default)]
75    pub method: AnomalyMethod,
76
77    /// Detection threshold for numeric metrics. `zscore`: max |x − mean| / std
78    /// (default 3.0). `iqr`: Tukey fence multiplier (default 1.5).
79    #[serde(default, skip_serializing_if = "Option::is_none")]
80    pub sensitivity: Option<f64>,
81
82    /// A categorical value never seen in the baseline is reported once its
83    /// share of the run's values reaches this fraction (default 0.05).
84    #[serde(default = "default_new_value_min_share")]
85    pub new_value_min_share: f64,
86
87    /// Population stability index above which a categorical column's value
88    /// distribution counts as drifted (default 0.2; the conventional
89    /// "significant shift" threshold).
90    #[serde(default = "default_psi_threshold")]
91    pub psi_threshold: f64,
92
93    /// What a detected drift does to the run. Default `warn`.
94    #[serde(default)]
95    pub on_drift: OnProfileDrift,
96}
97
98/// The action taken when a column's profile drifts from its baseline.
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
100#[serde(rename_all = "snake_case")]
101pub enum OnProfileDrift {
102    /// Log a warning and count `faucet_profile_drift_total`; the run succeeds.
103    #[default]
104    Warn,
105    /// As `warn`, plus a `profile_drift` notification event per finding.
106    Notify,
107    /// As `notify`, and the run is reported as failed (the data is already
108    /// written; the failure marks the run so operators look).
109    Fail,
110}
111
112fn default_exclude() -> Vec<String> {
113    vec!["_faucet_*".to_string()]
114}
115fn default_max_columns() -> usize {
116    DEFAULT_MAX_COLUMNS
117}
118fn default_top_values() -> usize {
119    DEFAULT_TOP_VALUES
120}
121fn default_categorical_max_distinct() -> u64 {
122    DEFAULT_CATEGORICAL_MAX_DISTINCT
123}
124fn default_window() -> u32 {
125    DEFAULT_WINDOW
126}
127fn default_min_history() -> u32 {
128    DEFAULT_MIN_HISTORY
129}
130fn default_new_value_min_share() -> f64 {
131    DEFAULT_NEW_VALUE_MIN_SHARE
132}
133fn default_psi_threshold() -> f64 {
134    DEFAULT_PSI_THRESHOLD
135}
136
137impl Default for ProfilingSpec {
138    fn default() -> Self {
139        Self {
140            columns: Vec::new(),
141            exclude: default_exclude(),
142            max_columns: DEFAULT_MAX_COLUMNS,
143            top_values: DEFAULT_TOP_VALUES,
144            categorical_max_distinct: DEFAULT_CATEGORICAL_MAX_DISTINCT,
145            window: DEFAULT_WINDOW,
146            min_history: DEFAULT_MIN_HISTORY,
147            method: AnomalyMethod::default(),
148            sensitivity: None,
149            new_value_min_share: DEFAULT_NEW_VALUE_MIN_SHARE,
150            psi_threshold: DEFAULT_PSI_THRESHOLD,
151            on_drift: OnProfileDrift::default(),
152        }
153    }
154}
155
156impl ProfilingSpec {
157    /// The configured sensitivity, or the method's conventional default.
158    pub fn effective_sensitivity(&self) -> f64 {
159        self.sensitivity
160            .unwrap_or_else(|| self.method.default_sensitivity())
161    }
162
163    /// Fail-fast validation, surfaced at config-load time.
164    pub fn validate(&self) -> Result<(), String> {
165        if self.max_columns == 0 || self.max_columns > MAX_MAX_COLUMNS {
166            return Err(format!(
167                "max_columns must be between 1 and {MAX_MAX_COLUMNS}, got {}",
168                self.max_columns
169            ));
170        }
171        if self.top_values > MAX_TOP_VALUES {
172            return Err(format!(
173                "top_values must be at most {MAX_TOP_VALUES}, got {}",
174                self.top_values
175            ));
176        }
177        if self.min_history < 2 {
178            return Err(format!(
179                "min_history must be at least 2, got {}",
180                self.min_history
181            ));
182        }
183        if self.window < self.min_history {
184            return Err(format!(
185                "window ({}) must be >= min_history ({})",
186                self.window, self.min_history
187            ));
188        }
189        if let Some(s) = self.sensitivity
190            && (!s.is_finite() || s <= 0.0)
191        {
192            return Err(format!("sensitivity must be a finite number > 0, got {s}"));
193        }
194        if !self.new_value_min_share.is_finite()
195            || self.new_value_min_share <= 0.0
196            || self.new_value_min_share > 1.0
197        {
198            return Err(format!(
199                "new_value_min_share must be in (0, 1], got {}",
200                self.new_value_min_share
201            ));
202        }
203        if !self.psi_threshold.is_finite() || self.psi_threshold <= 0.0 {
204            return Err(format!(
205                "psi_threshold must be a finite number > 0, got {}",
206                self.psi_threshold
207            ));
208        }
209        if let Some(c) = self.columns.iter().find(|c| c.trim().is_empty()) {
210            return Err(format!("columns contains an empty name {c:?}"));
211        }
212        Ok(())
213    }
214
215    /// Whether `name` is selected for profiling by `columns` / `exclude`.
216    pub fn selects(&self, name: &str) -> bool {
217        if crate::diff::is_excluded(name, &self.exclude) {
218            return false;
219        }
220        self.columns.is_empty() || self.columns.iter().any(|c| c == name)
221    }
222}
223
224#[cfg(test)]
225mod tests {
226    use super::*;
227
228    #[test]
229    fn defaults_are_sane_and_valid() {
230        let s = ProfilingSpec::default();
231        assert!(s.validate().is_ok());
232        assert_eq!(s.exclude, vec!["_faucet_*"]);
233        assert_eq!(s.effective_sensitivity(), 3.0);
234        let parsed: ProfilingSpec = serde_json::from_str("{}").unwrap();
235        assert_eq!(parsed, s);
236    }
237
238    #[test]
239    fn validation_names_each_defect() {
240        let bad = |f: fn(&mut ProfilingSpec)| {
241            let mut s = ProfilingSpec::default();
242            f(&mut s);
243            s.validate().unwrap_err()
244        };
245        assert!(bad(|s| s.max_columns = 0).contains("max_columns"));
246        assert!(bad(|s| s.max_columns = MAX_MAX_COLUMNS + 1).contains("max_columns"));
247        assert!(bad(|s| s.top_values = MAX_TOP_VALUES + 1).contains("top_values"));
248        assert!(bad(|s| s.min_history = 1).contains("min_history"));
249        assert!(bad(|s| s.window = 3).contains("window"));
250        assert!(bad(|s| s.sensitivity = Some(0.0)).contains("sensitivity"));
251        assert!(bad(|s| s.sensitivity = Some(f64::NAN)).contains("sensitivity"));
252        assert!(bad(|s| s.new_value_min_share = 0.0).contains("new_value_min_share"));
253        assert!(bad(|s| s.new_value_min_share = 1.5).contains("new_value_min_share"));
254        assert!(bad(|s| s.psi_threshold = -1.0).contains("psi_threshold"));
255        assert!(bad(|s| s.columns = vec![" ".into()]).contains("empty name"));
256    }
257
258    #[test]
259    fn selection_honours_columns_and_exclude_globs() {
260        let mut s = ProfilingSpec::default();
261        assert!(s.selects("amount"));
262        assert!(!s.selects("_faucet_run_id"));
263        s.columns = vec!["amount".into()];
264        assert!(s.selects("amount"));
265        assert!(!s.selects("other"));
266        s.exclude.push("amount".into());
267        assert!(!s.selects("amount"));
268    }
269
270    #[test]
271    fn on_drift_and_sensitivity_parse() {
272        let s: ProfilingSpec = serde_json::from_value(
273            serde_json::json!({"on_drift": "fail", "method": "iqr", "sensitivity": 2}),
274        )
275        .unwrap();
276        assert_eq!(s.on_drift, OnProfileDrift::Fail);
277        assert_eq!(s.effective_sensitivity(), 2.0);
278        let s: ProfilingSpec =
279            serde_json::from_value(serde_json::json!({"method": "iqr"})).unwrap();
280        assert_eq!(s.effective_sensitivity(), 1.5);
281        assert!(serde_json::from_value::<ProfilingSpec>(serde_json::json!({"bogus": 1})).is_err());
282    }
283}