use super::event::{Event, EventStatus};
use super::plugin::EventResponsePlugin;
use super::rule_manager::{RuleAdjustmentStrategy, RuleManager};
use mofa_kernel::plugin::{AgentPlugin, PluginContext, PluginResult};
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use tokio::sync::{RwLock, Semaphore};
pub struct EventHandlingEngine {
rule_manager: Arc<RuleManager>,
event_queue: Arc<RwLock<VecDeque<Event>>>,
plugins: Arc<RwLock<HashMap<String, Box<dyn EventResponsePlugin + Send + Sync>>>>,
plugin_context: PluginContext,
max_concurrent_handlers: usize,
semaphore: Arc<Semaphore>,
}
impl EventHandlingEngine {
pub fn new() -> Self {
Self {
rule_manager: Arc::new(RuleManager::new()),
event_queue: Arc::new(RwLock::new(VecDeque::new())),
plugins: Arc::new(RwLock::new(HashMap::new())),
plugin_context: PluginContext::new("event-handling-engine"),
max_concurrent_handlers: 10,
semaphore: Arc::new(Semaphore::new(10)),
}
}
pub fn with_max_concurrent_handlers(mut self, limit: usize) -> Self {
self.max_concurrent_handlers = limit;
self.semaphore = Arc::new(Semaphore::new(limit));
self
}
pub async fn set_rule_strategy(&self, strategy: RuleAdjustmentStrategy) {
self.rule_manager.set_strategy(strategy).await;
}
pub async fn register_plugin(&self, plugin: Box<dyn EventResponsePlugin + Send + Sync>) {
let plugin_id = plugin.metadata().id.to_string();
let mut plugins = self.plugins.write().await;
plugins.insert(plugin_id.clone(), plugin);
}
pub async fn register_plugins(&self, plugins: Vec<Box<dyn EventResponsePlugin + Send + Sync>>) {
for plugin in plugins {
self.register_plugin(plugin).await;
}
}
pub async fn submit_event(&self, event: Event) {
println!(
"Submitted new event: [{}] {} - {}",
event.priority, event.source, event.description
);
let mut queue = self.event_queue.write().await;
queue.push_back(event);
}
pub async fn process_next_event(&self) -> PluginResult<Option<Event>> {
let semaphore = self.semaphore.clone();
let _permit = semaphore.acquire().await.unwrap();
let mut queue = self.event_queue.write().await;
let mut event = match queue.pop_front() {
Some(event) => event,
None => return Ok(None),
};
event.update_status(EventStatus::Processing);
self.rule_manager.adjust_rules(&event).await?;
let mut plugins = self.plugins.write().await;
for (_plugin_id, plugin) in plugins.iter_mut() {
if plugin.can_handle(&event) {
println!(
"Processing event {} with plugin: {}",
event.id,
plugin.metadata().name
);
let processed_event = plugin.handle_event(event).await?;
println!(
"Event {} processed successfully by plugin {}",
processed_event.id,
plugin.metadata().name
);
return Ok(Some(processed_event));
}
}
println!("No plugin found to handle event: {}", event.id);
event.update_status(EventStatus::ManualInterventionNeeded);
Ok(Some(event))
}
pub async fn start(&self) -> PluginResult<()> {
println!(
"Starting event handling engine with {} concurrent handlers...",
self.max_concurrent_handlers
);
loop {
match self.process_next_event().await {
Ok(Some(event)) => {
if event.status == EventStatus::Resolved {
println!("Event resolved: {}", event.id);
} else {
println!("Event {} status: {:?}", event.id, event.status);
}
}
Ok(None) => {
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
}
Err(err) => {
println!("Error processing event: {}", err);
tokio::time::sleep(tokio::time::Duration::from_millis(1000)).await;
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::secretary::monitoring::event::*;
use crate::secretary::monitoring::plugins::{
NetworkAttackResponsePlugin, ServerFaultResponsePlugin,
};
#[tokio::test]
async fn test_event_handling() {
let engine = EventHandlingEngine::new();
let server_fault_plugin = ServerFaultResponsePlugin::new();
let network_attack_plugin = NetworkAttackResponsePlugin::new();
engine
.register_plugins(vec![
Box::new(server_fault_plugin) as Box<dyn EventResponsePlugin + Send + Sync>,
Box::new(network_attack_plugin) as Box<dyn EventResponsePlugin + Send + Sync>,
])
.await;
let server_fault_event = Event::new(
EventType::ServerFault,
EventPriority::High,
ImpactScope::Instance("web-server-01".to_string()),
"monitoring-agent".to_string(),
"Server unresponsive".to_string(),
serde_json::json!({ "server": "web-server-01" }),
);
let network_attack_event = Event::new(
EventType::NetworkAttack,
EventPriority::Emergency,
ImpactScope::Service("api-gateway".to_string()),
"ids".to_string(),
"DDoS attack".to_string(),
serde_json::json!({ "source_ip": "10.0.0.1" }),
);
engine.submit_event(server_fault_event).await;
engine.submit_event(network_attack_event).await;
let result = engine.process_next_event().await;
assert!(result.is_ok());
if let Ok(Some(event)) = result {
assert!(matches!(event.status, EventStatus::Resolved));
}
}
}