use std::sync::Arc;
use std::sync::RwLock;
use serde_json::json;
use anycms_event::prelude::*;
use anycms_event::trigger::{RuleStorage, TriggerRule, TriggerRuleEngine, InMemoryRuleStorage};
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
struct OrderCreated {
order_id: String,
amount: f64,
customer: String,
}
impl Event for OrderCreated {
fn event_name() -> &'static str {
"order.created"
}
fn topic() -> &'static str {
"order"
}
}
pub struct LoggingRuleStorage {
inner: InMemoryRuleStorage,
operations: RwLock<Vec<String>>,
}
impl LoggingRuleStorage {
pub fn new() -> Self {
Self {
inner: InMemoryRuleStorage::new(),
operations: RwLock::new(Vec::new()),
}
}
fn log_operation(&self, operation: &str) {
let timestamp = {
let ops = self.operations.read().unwrap();
ops.len() + 1
};
let log_entry = format!("[#{}] {}", timestamp, operation);
tracing::debug!("规则存储操作: {}", operation);
let mut ops = self.operations.write().unwrap();
ops.push(log_entry);
}
pub fn get_operations(&self) -> Vec<String> {
self.operations.read().unwrap().clone()
}
pub fn clear_logs(&self) {
self.operations.write().unwrap().clear();
}
}
impl Default for LoggingRuleStorage {
fn default() -> Self {
Self::new()
}
}
impl RuleStorage for LoggingRuleStorage {
fn add(&self, rule: TriggerRule) {
tracing::info!("添加规则: id={}, name={}, pattern={}",
rule.id, rule.name, rule.event_pattern);
self.log_operation(&format!("ADD: {} ({})", rule.id, rule.name));
self.inner.add(rule);
}
fn remove(&self, rule_id: &str) -> Option<TriggerRule> {
if let Some(rule) = self.inner.remove(rule_id) {
tracing::info!("移除规则: id={}, name={}", rule_id, rule.name);
self.log_operation(&format!("REMOVE: {} ({})", rule.id, rule.name));
Some(rule)
} else {
tracing::warn!("尝试移除不存在的规则: id={}", rule_id);
self.log_operation(&format!("REMOVE: {} (not found)", rule_id));
None
}
}
fn get(&self, rule_id: &str) -> Option<TriggerRule> {
self.inner.get(rule_id)
}
fn update(&self, rule: TriggerRule) -> bool {
let success = self.inner.update(rule.clone());
if success {
tracing::info!("更新规则: id={}, name={}", rule.id, rule.name);
self.log_operation(&format!("UPDATE: {} ({})", rule.id, rule.name));
} else {
tracing::warn!("尝试更新不存在的规则: id={}", rule.id);
self.log_operation(&format!("UPDATE: {} (not found)", rule.id));
}
success
}
fn list(&self) -> Vec<TriggerRule> {
self.inner.list()
}
fn count(&self) -> usize {
self.inner.count()
}
}
#[allow(dead_code)]
struct DatabaseRuleStorage {
}
#[allow(dead_code)]
impl DatabaseRuleStorage {
pub fn new() -> Self {
Self {
}
}
}
#[allow(dead_code)]
impl RuleStorage for DatabaseRuleStorage {
fn add(&self, rule: TriggerRule) {
tracing::info!("[DB] 添加规则: {}", rule.id);
}
fn remove(&self, rule_id: &str) -> Option<TriggerRule> {
tracing::info!("[DB] 移除规则: {}", rule_id);
None
}
fn get(&self, rule_id: &str) -> Option<TriggerRule> {
tracing::info!("[DB] 获取规则: {}", rule_id);
None
}
fn update(&self, rule: TriggerRule) -> bool {
tracing::info!("[DB] 更新规则: {}", rule.id);
false
}
fn list(&self) -> Vec<TriggerRule> {
tracing::info!("[DB] 列出所有规则");
Vec::new()
}
fn count(&self) -> usize {
tracing::info!("[DB] 统计规则数量");
0
}
}
#[tokio::main]
async fn main() {
tracing_subscriber::fmt()
.with_max_level(tracing::Level::INFO)
.init();
println!("=== 自定义规则存储后端示例 ===\n");
let bus = EventBus::builder()
.capacity(256)
.build();
let custom_storage = Arc::new(LoggingRuleStorage::new());
println!("✓ 创建自定义存储: LoggingRuleStorage\n");
let trigger_engine = Arc::new(TriggerRuleEngine::with_storage(
bus.clone(),
custom_storage.clone(),
));
println!("✓ 创建触发规则引擎(使用自定义存储)\n");
trigger_engine.register_action("log", |ctx| async move {
println!(" >> 触发规则: {}", ctx.rule_name);
println!(" 事件: {}, 规则ID: {}", ctx.event_name, ctx.rule_id);
Ok(())
});
trigger_engine.register_action("notify", |ctx| async move {
println!(" >> 发送通知: {}", ctx.action_config["message"]);
Ok(())
});
println!("✓ 注册动作处理器: log, notify\n");
println!("--- 添加触发规则 ---");
trigger_engine.add_rule(TriggerRule {
id: "order-small".into(),
name: "小订单处理".into(),
event_pattern: "order.created".into(),
condition: Some(json!({"amount": {"$lt": 100}})),
action_type: "log".into(),
action_config: json!({}),
enabled: true,
priority: 10,
});
trigger_engine.add_rule(TriggerRule {
id: "order-medium".into(),
name: "中等订单处理".into(),
event_pattern: "order.created".into(),
condition: Some(json!({"amount": {"$gte": 100, "$lt": 500}})),
action_type: "notify".into(),
action_config: json!({"message": "收到中等订单"}),
enabled: true,
priority: 5,
});
trigger_engine.add_rule(TriggerRule {
id: "order-large".into(),
name: "大订单处理".into(),
event_pattern: "order.created".into(),
condition: Some(json!({"amount": {"$gte": 500}})),
action_type: "notify".into(),
action_config: json!({"message": "收到大订单!需要人工审核"}),
enabled: true,
priority: 0, });
println!();
println!("--- 验证规则一致性 ---");
let engine_rules = trigger_engine.list_rules();
let storage_rules = custom_storage.list();
println!(" 引擎中的规则数: {}", engine_rules.len());
println!(" 存储中的规则数: {}", storage_rules.len());
println!(" 规则数一致: {}", engine_rules.len() == storage_rules.len());
println!();
println!("--- 存储操作日志 ---");
for (i, log) in custom_storage.get_operations().iter().enumerate() {
println!(" {}. {}", i + 1, log);
}
println!();
println!("--- 测试规则触发 ---");
println!("\n[小订单] $50:");
let _ = trigger_engine
.process_event("order.created", &json!({"order_id": "ORD-001", "amount": 50, "customer": "Alice"}))
.await;
println!("\n[中等订单] $250:");
let _ = trigger_engine
.process_event("order.created", &json!({"order_id": "ORD-002", "amount": 250, "customer": "Bob"}))
.await;
println!("\n[大订单] $800:");
let _ = trigger_engine
.process_event("order.created", &json!({"order_id": "ORD-003", "amount": 800, "customer": "Charlie"}))
.await;
println!();
println!("--- 更新规则 ---");
trigger_engine.disable_rule("order-medium");
println!(" 禁用规则: order-medium");
let updated_rule = TriggerRule {
id: "order-small".into(),
name: "小订单自动处理(更新)".into(),
event_pattern: "order.created".into(),
condition: Some(json!({"amount": {"$lt": 50}})),
action_type: "log".into(),
action_config: json!({}),
enabled: true,
priority: 10,
};
trigger_engine.update_rule(updated_rule);
println!(" 更新规则: order-small");
println!();
println!("--- 更新后的操作日志 ---");
for (i, log) in custom_storage.get_operations().iter().enumerate() {
println!(" {}. {}", i + 1, log);
}
println!();
println!("--- 验证更新后的规则 ---");
let rules = trigger_engine.list_rules();
println!(" 当前规则列表(按优先级排序):");
for rule in &rules {
let status = if rule.enabled { "✓" } else { "✗" };
println!(" [{}] {} - {} (pattern: {}, priority: {})",
status, rule.id, rule.name, rule.event_pattern, rule.priority);
}
println!();
println!("--- 测试更新后的规则 ---");
println!("\n[小订单] $60 (不满足更新后的条件 < 50):");
let results = trigger_engine
.process_event("order.created", &json!({"amount": 60}))
.await;
println!(" 触发规则数: {}", results.len());
println!("\n[小订单] $40 (满足更新后的条件 < 50):");
let results = trigger_engine
.process_event("order.created", &json!({"amount": 40}))
.await;
println!(" 触发规则数: {}", results.len());
println!();
println!("--- 删除规则 ---");
let removed = trigger_engine.remove_rule("order-large");
if let Some(rule) = removed {
println!(" 已删除规则: {}", rule.name);
}
println!();
println!("--- 最终统计 ---");
println!(" 剩余规则数: {}", trigger_engine.rule_count());
println!(" 总操作日志数: {}", custom_storage.get_operations().len());
println!();
println!("=== 示例完成 ===");
println!("\n💡 提示:");
println!(" - 自定义 RuleStorage 可以实现数据库持久化");
println!(" - 装饰器模式可以在不修改原有代码的情况下添加功能");
println!(" - 引擎和存储之间的数据始终保持同步");
println!(" - 操作日志可用于审计和调试");
}