faucet_core/profiling/
spec.rs1use crate::anomaly::AnomalyMethod;
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6
7pub const DEFAULT_MAX_COLUMNS: usize = 200;
9pub const DEFAULT_TOP_VALUES: usize = 10;
11pub const DEFAULT_CATEGORICAL_MAX_DISTINCT: u64 = 100;
14pub const DEFAULT_WINDOW: u32 = 20;
16pub const DEFAULT_MIN_HISTORY: u32 = 5;
18pub const DEFAULT_NEW_VALUE_MIN_SHARE: f64 = 0.05;
21pub const DEFAULT_PSI_THRESHOLD: f64 = 0.2;
23pub const MAX_TOP_VALUES: usize = 1000;
25pub const MAX_MAX_COLUMNS: usize = 5000;
27
28#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
34#[serde(deny_unknown_fields)]
35pub struct ProfilingSpec {
36 #[serde(default, skip_serializing_if = "Vec::is_empty")]
39 pub columns: Vec<String>,
40
41 #[serde(default = "default_exclude")]
44 pub exclude: Vec<String>,
45
46 #[serde(default = "default_max_columns")]
49 pub max_columns: usize,
50
51 #[serde(default = "default_top_values")]
54 pub top_values: usize,
55
56 #[serde(default = "default_categorical_max_distinct")]
60 pub categorical_max_distinct: u64,
61
62 #[serde(default = "default_window")]
65 pub window: u32,
66
67 #[serde(default = "default_min_history")]
70 pub min_history: u32,
71
72 #[serde(default)]
75 pub method: AnomalyMethod,
76
77 #[serde(default, skip_serializing_if = "Option::is_none")]
80 pub sensitivity: Option<f64>,
81
82 #[serde(default = "default_new_value_min_share")]
85 pub new_value_min_share: f64,
86
87 #[serde(default = "default_psi_threshold")]
91 pub psi_threshold: f64,
92
93 #[serde(default)]
95 pub on_drift: OnProfileDrift,
96}
97
98#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
100#[serde(rename_all = "snake_case")]
101pub enum OnProfileDrift {
102 #[default]
104 Warn,
105 Notify,
107 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 pub fn effective_sensitivity(&self) -> f64 {
159 self.sensitivity
160 .unwrap_or_else(|| self.method.default_sensitivity())
161 }
162
163 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 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}