1use std::collections::HashMap;
2use std::sync::Arc;
3
4use parking_lot::RwLock;
5use serde::{Deserialize, Serialize};
6
7use crate::storage::Severity;
8
9#[derive(Debug, Clone, Serialize, Deserialize)]
11pub struct Alert {
12 pub id: String,
14 pub name: String,
16 pub description: String,
18 pub severity: Severity,
20 pub metric: String,
22 pub threshold: f64,
24 pub operator: AlertOperator,
26 pub enabled: bool,
28}
29
30#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
32pub enum AlertOperator {
33 GreaterThan,
34 LessThan,
35 Equals,
36 NotEquals,
37}
38
39#[derive(Debug, Clone, Serialize, Deserialize)]
41pub struct AlertState {
42 pub alert_id: String,
44 pub triggered: bool,
46 pub last_triggered: Option<String>,
48 pub trigger_count: u64,
50}
51
52pub struct AlertManager {
54 alerts: Arc<RwLock<Vec<Alert>>>,
55 states: Arc<RwLock<HashMap<String, AlertState>>>,
56}
57
58impl AlertManager {
59 #[must_use]
61 pub fn new() -> Self {
62 Self {
63 alerts: Arc::new(RwLock::new(Vec::new())),
64 states: Arc::new(RwLock::new(HashMap::new())),
65 }
66 }
67
68 pub fn add_alert(&self, alert: Alert) {
70 let mut alerts = self.alerts.write();
71 alerts.push(alert);
72 }
73
74 pub fn check_alerts(&self, metrics: &HashMap<String, f64>) -> Vec<Alert> {
76 let alerts = self.alerts.read();
77 let mut states = self.states.write();
78 let mut triggered = Vec::new();
79
80 for alert in alerts.iter() {
81 if !alert.enabled {
82 continue;
83 }
84
85 if let Some(value) = metrics.get(&alert.metric) {
86 let is_triggered = match alert.operator {
87 AlertOperator::GreaterThan => *value > alert.threshold,
88 AlertOperator::LessThan => *value < alert.threshold,
89 AlertOperator::Equals => (*value - alert.threshold).abs() < f64::EPSILON,
90 AlertOperator::NotEquals => (*value - alert.threshold).abs() > f64::EPSILON,
91 };
92
93 let state = states
94 .entry(alert.id.clone())
95 .or_insert_with(|| AlertState {
96 alert_id: alert.id.clone(),
97 triggered: false,
98 last_triggered: None,
99 trigger_count: 0,
100 });
101
102 if is_triggered && !state.triggered {
103 state.triggered = true;
104 state.last_triggered = Some(chrono::Utc::now().to_rfc3339());
105 state.trigger_count += 1;
106 triggered.push(alert.clone());
107 } else if !is_triggered {
108 state.triggered = false;
109 }
110 }
111 }
112
113 triggered
114 }
115
116 #[must_use]
118 pub fn alerts(&self) -> Vec<Alert> {
119 self.alerts.read().clone()
120 }
121
122 #[must_use]
124 pub fn states(&self) -> HashMap<String, AlertState> {
125 self.states.read().clone()
126 }
127}
128
129impl Default for AlertManager {
130 fn default() -> Self {
131 Self::new()
132 }
133}
134
135#[derive(Debug, Clone, Serialize, Deserialize)]
137pub struct ScheduledCrawl {
138 pub id: String,
140 pub url: String,
142 pub schedule: String,
144 pub config: serde_json::Value,
146 pub enabled: bool,
148 pub last_run: Option<String>,
150 pub next_run: Option<String>,
152}
153
154pub struct CrawlScheduler {
156 schedules: Arc<RwLock<Vec<ScheduledCrawl>>>,
157}
158
159impl CrawlScheduler {
160 #[must_use]
162 pub fn new() -> Self {
163 Self {
164 schedules: Arc::new(RwLock::new(Vec::new())),
165 }
166 }
167
168 pub fn add_schedule(&self, schedule: ScheduledCrawl) {
170 let mut schedules = self.schedules.write();
171 schedules.push(schedule);
172 }
173
174 #[must_use]
176 pub fn schedules(&self) -> Vec<ScheduledCrawl> {
177 self.schedules.read().clone()
178 }
179
180 #[must_use]
182 pub fn get_due_crawls(&self) -> Vec<ScheduledCrawl> {
183 self.schedules
184 .read()
185 .iter()
186 .filter(|s| s.enabled)
187 .cloned()
188 .collect()
189 }
190}
191
192impl Default for CrawlScheduler {
193 fn default() -> Self {
194 Self::new()
195 }
196}
197
198#[derive(Debug, Clone, Serialize, Deserialize)]
200pub struct TrendDataPoint {
201 pub timestamp: String,
203 pub metric: String,
205 pub value: f64,
207 pub crawl_id: String,
209}
210
211pub struct TrendTracker {
213 data: Arc<RwLock<Vec<TrendDataPoint>>>,
214}
215
216impl TrendTracker {
217 #[must_use]
219 pub fn new() -> Self {
220 Self {
221 data: Arc::new(RwLock::new(Vec::new())),
222 }
223 }
224
225 pub fn record(&self, point: TrendDataPoint) {
227 let mut data = self.data.write();
228 data.push(point);
229 }
230
231 #[must_use]
233 pub fn get_trend(&self, metric: &str) -> Vec<TrendDataPoint> {
234 self.data
235 .read()
236 .iter()
237 .filter(|p| p.metric == metric)
238 .cloned()
239 .collect()
240 }
241
242 #[must_use]
244 pub fn all_data(&self) -> Vec<TrendDataPoint> {
245 self.data.read().clone()
246 }
247
248 #[must_use]
250 pub fn average(&self, metric: &str) -> Option<f64> {
251 let values: Vec<f64> = self
252 .data
253 .read()
254 .iter()
255 .filter(|p| p.metric == metric)
256 .map(|p| p.value)
257 .collect();
258
259 if values.is_empty() {
260 None
261 } else {
262 Some(values.iter().sum::<f64>() / values.len() as f64)
263 }
264 }
265}
266
267impl Default for TrendTracker {
268 fn default() -> Self {
269 Self::new()
270 }
271}
272
273#[cfg(test)]
278mod tests {
279 use super::*;
280
281 #[test]
282 fn test_alert_manager() {
283 let manager = AlertManager::new();
284
285 let alert = Alert {
286 id: "test".to_string(),
287 name: "Test Alert".to_string(),
288 description: "Test alert".to_string(),
289 severity: Severity::Warning,
290 metric: "pages_crawled".to_string(),
291 threshold: 100.0,
292 operator: AlertOperator::GreaterThan,
293 enabled: true,
294 };
295
296 manager.add_alert(alert);
297 assert_eq!(manager.alerts().len(), 1);
298 }
299
300 #[test]
301 fn test_alert_triggering() {
302 let manager = AlertManager::new();
303
304 let alert = Alert {
305 id: "test".to_string(),
306 name: "Test Alert".to_string(),
307 description: "Test alert".to_string(),
308 severity: Severity::Warning,
309 metric: "pages_crawled".to_string(),
310 threshold: 100.0,
311 operator: AlertOperator::GreaterThan,
312 enabled: true,
313 };
314
315 manager.add_alert(alert);
316
317 let mut metrics = HashMap::new();
318 metrics.insert("pages_crawled".to_string(), 150.0);
319
320 let triggered = manager.check_alerts(&metrics);
321 assert_eq!(triggered.len(), 1);
322 assert_eq!(triggered[0].id, "test");
323 }
324
325 #[test]
326 fn test_scheduler() {
327 let scheduler = CrawlScheduler::new();
328
329 let schedule = ScheduledCrawl {
330 id: "test".to_string(),
331 url: "https://example.com".to_string(),
332 schedule: "daily".to_string(),
333 config: serde_json::json!({}),
334 enabled: true,
335 last_run: None,
336 next_run: None,
337 };
338
339 scheduler.add_schedule(schedule);
340 assert_eq!(scheduler.schedules().len(), 1);
341 assert_eq!(scheduler.get_due_crawls().len(), 1);
342 }
343
344 #[test]
345 fn test_trend_tracker() {
346 let tracker = TrendTracker::new();
347
348 tracker.record(TrendDataPoint {
349 timestamp: "2026-01-01".to_string(),
350 metric: "pages_crawled".to_string(),
351 value: 100.0,
352 crawl_id: "crawl1".to_string(),
353 });
354
355 tracker.record(TrendDataPoint {
356 timestamp: "2026-01-02".to_string(),
357 metric: "pages_crawled".to_string(),
358 value: 150.0,
359 crawl_id: "crawl2".to_string(),
360 });
361
362 assert_eq!(tracker.get_trend("pages_crawled").len(), 2);
363 assert!((tracker.average("pages_crawled").unwrap() - 125.0).abs() < 0.01);
364 }
365}