use super::event::{Event, EventType};
use super::plugin::{BaseEventResponsePlugin, EventResponseConfig, EventResponsePlugin};
use async_trait::async_trait;
use mofa_kernel::plugin::{PluginPriority, PluginResult};
use std::collections::HashMap;
use std::sync::RwLock;
pub struct ServerFaultResponsePlugin {
base: BaseEventResponsePlugin,
config: RwLock<EventResponseConfig>,
}
impl ServerFaultResponsePlugin {
pub fn new() -> Self {
let handled_event_types = vec![EventType::ServerFault];
let workflow_steps = vec![
"attempt_auto_restart".to_string(),
"notify_administrator".to_string(),
];
let base = BaseEventResponsePlugin::new(
"server-fault-responder",
"Server Fault Responder",
handled_event_types.clone(), workflow_steps,
)
.with_priority(PluginPriority::High) .with_max_impact_scope("instance");
let config = RwLock::new(EventResponseConfig {
handled_event_types,
priority: PluginPriority::High,
..Default::default()
});
Self { base, config }
}
async fn attempt_auto_restart(&self, server: &str) -> Result<bool, String> {
println!("Attempting to restart server: {}", server);
Ok(true)
}
async fn notify_administrator(&self, event: &Event) -> Result<(), String> {
println!("Notifying administrator about server fault:");
println!(" Event ID: {}", event.id);
println!(" Source: {}", event.source);
println!(" Description: {}", event.description);
Ok(())
}
}
#[async_trait]
impl EventResponsePlugin for ServerFaultResponsePlugin {
fn config(&self) -> &EventResponseConfig {
panic!("config() should not be called directly on this plugin");
}
async fn update_config(&mut self, config: EventResponseConfig) -> PluginResult<()> {
{
let mut current_config = self.config.write().unwrap();
*current_config = config.clone();
}
self.base.update_config(config).await
}
fn can_handle(&self, event: &Event) -> bool {
self.base.can_handle(event)
}
async fn handle_event(&mut self, event: Event) -> PluginResult<Event> {
self.base.handle_event(event).await
}
async fn execute_workflow(&self, event: &Event) -> PluginResult<HashMap<String, String>> {
let mut result = HashMap::new();
let server_name = event
.data
.get("server")
.and_then(|s| s.as_str())
.unwrap_or("unknown");
let restart_result = self.attempt_auto_restart(server_name).await;
match restart_result {
Ok(success) => {
result.insert(
"auto_restart".to_string(),
if success {
"success".to_string()
} else {
"failed".to_string()
},
);
}
Err(err) => {
result.insert("auto_restart".to_string(), format!("error: {}", err));
}
}
match self.notify_administrator(event).await {
Ok(_) => {
result.insert("notify_admin".to_string(), "success".to_string());
}
Err(err) => {
result.insert("notify_admin".to_string(), format!("error: {}", err));
}
}
result.insert(
"workflow_status".to_string(),
"server_fault_workflow_completed".to_string(),
);
Ok(result)
}
}
#[async_trait]
impl mofa_kernel::plugin::AgentPlugin for ServerFaultResponsePlugin {
fn metadata(&self) -> &mofa_kernel::plugin::PluginMetadata {
self.base.metadata()
}
fn state(&self) -> mofa_kernel::plugin::PluginState {
self.base.state()
}
async fn load(
&mut self,
ctx: &mofa_kernel::plugin::PluginContext,
) -> mofa_kernel::plugin::PluginResult<()> {
self.base.load(ctx).await
}
async fn init_plugin(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.init_plugin().await
}
async fn start(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.start().await
}
async fn stop(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.stop().await
}
async fn unload(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.unload().await
}
async fn execute(&mut self, input: String) -> mofa_kernel::plugin::PluginResult<String> {
self.base.execute(input).await
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn as_any_mut(&mut self) -> &mut dyn std::any::Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn std::any::Any> {
self
}
}
impl From<ServerFaultResponsePlugin> for Box<dyn EventResponsePlugin> {
fn from(plugin: ServerFaultResponsePlugin) -> Self {
Box::new(plugin)
}
}
impl From<ServerFaultResponsePlugin> for Box<dyn mofa_kernel::plugin::AgentPlugin> {
fn from(plugin: ServerFaultResponsePlugin) -> Self {
Box::new(plugin)
}
}
pub struct NetworkAttackResponsePlugin {
base: BaseEventResponsePlugin,
config: RwLock<EventResponseConfig>,
}
impl NetworkAttackResponsePlugin {
pub fn new() -> Self {
let handled_event_types = vec![EventType::NetworkAttack];
let workflow_steps = vec![
"block_attacking_ip".to_string(),
"analyze_attack_pattern".to_string(),
"notify_security_team".to_string(),
];
let base = BaseEventResponsePlugin::new(
"network-attack-responder",
"Network Attack Responder",
handled_event_types.clone(), workflow_steps,
)
.with_priority(PluginPriority::Critical) .with_max_impact_scope("system");
let config = RwLock::new(EventResponseConfig {
handled_event_types,
priority: PluginPriority::Critical,
..Default::default()
});
Self { base, config }
}
async fn block_attacking_ip(&self, ip: &str) -> Result<bool, String> {
println!("Blocking attacking IP: {}", ip);
Ok(true)
}
async fn analyze_attack_pattern(&self, event: &Event) -> Result<String, String> {
println!("Analyzing attack pattern for event: {}", event.id);
Ok("ddos_attack".to_string()) }
async fn notify_security_team(&self, event: &Event, attack_type: &str) -> Result<(), String> {
println!("Notifying security team about {} attack:", attack_type);
println!(" Event ID: {}", event.id);
println!(" Source IP: {:?}", event.data.get("source_ip"));
Ok(())
}
}
#[async_trait]
impl EventResponsePlugin for NetworkAttackResponsePlugin {
fn config(&self) -> &EventResponseConfig {
panic!("config() should not be called directly on this plugin");
}
async fn update_config(&mut self, config: EventResponseConfig) -> PluginResult<()> {
{
let mut current_config = self.config.write().unwrap();
*current_config = config.clone();
}
self.base.update_config(config).await
}
fn can_handle(&self, event: &Event) -> bool {
self.base.can_handle(event)
}
async fn handle_event(&mut self, event: Event) -> PluginResult<Event> {
self.base.handle_event(event).await
}
async fn execute_workflow(&self, event: &Event) -> PluginResult<HashMap<String, String>> {
let mut result = HashMap::new();
let source_ip = event
.data
.get("source_ip")
.and_then(|ip| ip.as_str())
.unwrap_or("unknown");
let block_result = self.block_attacking_ip(source_ip).await;
match block_result {
Ok(success) => {
result.insert(
"block_ip".to_string(),
if success {
"success".to_string()
} else {
"failed".to_string()
},
);
}
Err(err) => {
result.insert("block_ip".to_string(), format!("error: {}", err));
}
}
let analysis_result = self.analyze_attack_pattern(event).await;
let attack_type = match analysis_result {
Ok(attack) => {
result.insert("attack_analysis".to_string(), attack.clone());
attack
}
Err(err) => {
result.insert("attack_analysis".to_string(), format!("error: {}", err));
"unknown".to_string()
}
};
if let Err(err) = self.notify_security_team(event, &attack_type).await {
result.insert("notify_security".to_string(), format!("error: {}", err));
} else {
result.insert("notify_security".to_string(), "success".to_string());
}
result.insert(
"workflow_status".to_string(),
"network_attack_workflow_completed".to_string(),
);
Ok(result)
}
}
#[async_trait]
impl mofa_kernel::plugin::AgentPlugin for NetworkAttackResponsePlugin {
fn metadata(&self) -> &mofa_kernel::plugin::PluginMetadata {
self.base.metadata()
}
fn state(&self) -> mofa_kernel::plugin::PluginState {
self.base.state()
}
async fn load(
&mut self,
ctx: &mofa_kernel::plugin::PluginContext,
) -> mofa_kernel::plugin::PluginResult<()> {
self.base.load(ctx).await
}
async fn init_plugin(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.init_plugin().await
}
async fn start(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.start().await
}
async fn stop(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.stop().await
}
async fn unload(&mut self) -> mofa_kernel::plugin::PluginResult<()> {
self.base.unload().await
}
async fn execute(&mut self, input: String) -> mofa_kernel::plugin::PluginResult<String> {
self.base.execute(input).await
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn as_any_mut(&mut self) -> &mut dyn std::any::Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn std::any::Any> {
self
}
}
impl From<NetworkAttackResponsePlugin> for Box<dyn EventResponsePlugin> {
fn from(plugin: NetworkAttackResponsePlugin) -> Self {
Box::new(plugin)
}
}
impl From<NetworkAttackResponsePlugin> for Box<dyn mofa_kernel::plugin::AgentPlugin> {
fn from(plugin: NetworkAttackResponsePlugin) -> Self {
Box::new(plugin)
}
}