use std::collections::HashMap;
use std::sync::Arc;
use parking_lot::RwLock;
use serde::{Deserialize, Serialize};
use crate::storage::Severity;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Alert {
pub id: String,
pub name: String,
pub description: String,
pub severity: Severity,
pub metric: String,
pub threshold: f64,
pub operator: AlertOperator,
pub enabled: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub enum AlertOperator {
GreaterThan,
LessThan,
Equals,
NotEquals,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AlertState {
pub alert_id: String,
pub triggered: bool,
pub last_triggered: Option<String>,
pub trigger_count: u64,
}
pub struct AlertManager {
alerts: Arc<RwLock<Vec<Alert>>>,
states: Arc<RwLock<HashMap<String, AlertState>>>,
}
impl AlertManager {
#[must_use]
pub fn new() -> Self {
Self {
alerts: Arc::new(RwLock::new(Vec::new())),
states: Arc::new(RwLock::new(HashMap::new())),
}
}
pub fn add_alert(&self, alert: Alert) {
let mut alerts = self.alerts.write();
alerts.push(alert);
}
pub fn check_alerts(&self, metrics: &HashMap<String, f64>) -> Vec<Alert> {
let alerts = self.alerts.read();
let mut states = self.states.write();
let mut triggered = Vec::new();
for alert in alerts.iter() {
if !alert.enabled {
continue;
}
if let Some(value) = metrics.get(&alert.metric) {
let is_triggered = match alert.operator {
AlertOperator::GreaterThan => *value > alert.threshold,
AlertOperator::LessThan => *value < alert.threshold,
AlertOperator::Equals => (*value - alert.threshold).abs() < f64::EPSILON,
AlertOperator::NotEquals => (*value - alert.threshold).abs() > f64::EPSILON,
};
let state = states
.entry(alert.id.clone())
.or_insert_with(|| AlertState {
alert_id: alert.id.clone(),
triggered: false,
last_triggered: None,
trigger_count: 0,
});
if is_triggered && !state.triggered {
state.triggered = true;
state.last_triggered = Some(chrono::Utc::now().to_rfc3339());
state.trigger_count += 1;
triggered.push(alert.clone());
} else if !is_triggered {
state.triggered = false;
}
}
}
triggered
}
#[must_use]
pub fn alerts(&self) -> Vec<Alert> {
self.alerts.read().clone()
}
#[must_use]
pub fn states(&self) -> HashMap<String, AlertState> {
self.states.read().clone()
}
}
impl Default for AlertManager {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ScheduledCrawl {
pub id: String,
pub url: String,
pub schedule: String,
pub config: serde_json::Value,
pub enabled: bool,
pub last_run: Option<String>,
pub next_run: Option<String>,
}
pub struct CrawlScheduler {
schedules: Arc<RwLock<Vec<ScheduledCrawl>>>,
}
impl CrawlScheduler {
#[must_use]
pub fn new() -> Self {
Self {
schedules: Arc::new(RwLock::new(Vec::new())),
}
}
pub fn add_schedule(&self, schedule: ScheduledCrawl) {
let mut schedules = self.schedules.write();
schedules.push(schedule);
}
#[must_use]
pub fn schedules(&self) -> Vec<ScheduledCrawl> {
self.schedules.read().clone()
}
#[must_use]
pub fn get_due_crawls(&self) -> Vec<ScheduledCrawl> {
self.schedules
.read()
.iter()
.filter(|s| s.enabled)
.cloned()
.collect()
}
}
impl Default for CrawlScheduler {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TrendDataPoint {
pub timestamp: String,
pub metric: String,
pub value: f64,
pub crawl_id: String,
}
pub struct TrendTracker {
data: Arc<RwLock<Vec<TrendDataPoint>>>,
}
impl TrendTracker {
#[must_use]
pub fn new() -> Self {
Self {
data: Arc::new(RwLock::new(Vec::new())),
}
}
pub fn record(&self, point: TrendDataPoint) {
let mut data = self.data.write();
data.push(point);
}
#[must_use]
pub fn get_trend(&self, metric: &str) -> Vec<TrendDataPoint> {
self.data
.read()
.iter()
.filter(|p| p.metric == metric)
.cloned()
.collect()
}
#[must_use]
pub fn all_data(&self) -> Vec<TrendDataPoint> {
self.data.read().clone()
}
#[must_use]
pub fn average(&self, metric: &str) -> Option<f64> {
let values: Vec<f64> = self
.data
.read()
.iter()
.filter(|p| p.metric == metric)
.map(|p| p.value)
.collect();
if values.is_empty() {
None
} else {
Some(values.iter().sum::<f64>() / values.len() as f64)
}
}
}
impl Default for TrendTracker {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_alert_manager() {
let manager = AlertManager::new();
let alert = Alert {
id: "test".to_string(),
name: "Test Alert".to_string(),
description: "Test alert".to_string(),
severity: Severity::Warning,
metric: "pages_crawled".to_string(),
threshold: 100.0,
operator: AlertOperator::GreaterThan,
enabled: true,
};
manager.add_alert(alert);
assert_eq!(manager.alerts().len(), 1);
}
#[test]
fn test_alert_triggering() {
let manager = AlertManager::new();
let alert = Alert {
id: "test".to_string(),
name: "Test Alert".to_string(),
description: "Test alert".to_string(),
severity: Severity::Warning,
metric: "pages_crawled".to_string(),
threshold: 100.0,
operator: AlertOperator::GreaterThan,
enabled: true,
};
manager.add_alert(alert);
let mut metrics = HashMap::new();
metrics.insert("pages_crawled".to_string(), 150.0);
let triggered = manager.check_alerts(&metrics);
assert_eq!(triggered.len(), 1);
assert_eq!(triggered[0].id, "test");
}
#[test]
fn test_scheduler() {
let scheduler = CrawlScheduler::new();
let schedule = ScheduledCrawl {
id: "test".to_string(),
url: "https://example.com".to_string(),
schedule: "daily".to_string(),
config: serde_json::json!({}),
enabled: true,
last_run: None,
next_run: None,
};
scheduler.add_schedule(schedule);
assert_eq!(scheduler.schedules().len(), 1);
assert_eq!(scheduler.get_due_crawls().len(), 1);
}
#[test]
fn test_trend_tracker() {
let tracker = TrendTracker::new();
tracker.record(TrendDataPoint {
timestamp: "2026-01-01".to_string(),
metric: "pages_crawled".to_string(),
value: 100.0,
crawl_id: "crawl1".to_string(),
});
tracker.record(TrendDataPoint {
timestamp: "2026-01-02".to_string(),
metric: "pages_crawled".to_string(),
value: 150.0,
crawl_id: "crawl2".to_string(),
});
assert_eq!(tracker.get_trend("pages_crawled").len(), 2);
assert!((tracker.average("pages_crawled").unwrap() - 125.0).abs() < 0.01);
}
}