rs-eda 0.1.2

rust EDA event-driven framework
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
|--src
    |--engine
        |--mod.rs
    |--event
        |--mod.rs
    |--plugin
        |--mod.rs
    |--strategy
        |--mod.rs
    |--mod.rs


各文件代码
-------src/engine/mod.rs-------
use super::event::{Event, EventDispatcher};
use super::plugin::PluginOptions;
use super::strategy::StrategyOptions;
use std::collections::HashMap;

// 定义 EngineTrait 并实现必要的特性
pub trait EngineTrait: Send + Sync + 'static {}

// 明确实现 EngineTrait 的具体类型
impl EngineTrait for Engine {}

// 修改 PluginOptions 和 StrategyOptions 的定义,确保泛型参数实现 EngineTrait
pub struct Engine {
    event_dispatcher: EventDispatcher,
    tables_plugin: HashMap<String, Box<dyn PluginOptions<Engine> + Send + Sync + 'static>>,
    tables_strategy: HashMap<String, Box<dyn StrategyOptions<Engine> + Send + Sync + 'static>>,
}

impl Engine {
    pub fn new() -> Self {
        Self {
            event_dispatcher: EventDispatcher::new(),
            tables_plugin: HashMap::new(),
            tables_strategy: HashMap::new(),
        }
    }

    pub fn install<T: PluginOptions<Engine> + Send + Sync + 'static>(&mut self, plugin: T) {
        if self.tables_plugin.contains_key(plugin.name()) {
            println!("This plugin already exists: {}", plugin.name());
            return;
        }

        let validate_dep = |name: &str| {
            if !self.tables_plugin.contains_key(name) {
                println!(
                    "{} must install this plugin before: {}",
                    plugin.name(),
                    name
                );
                false
            } else {
                true
            }
        };

        if let Some(deps) = plugin.deps() {
            if let Some(deps_list) = deps.as_ref().downcast_ref::<Vec<String>>() {
                for dep in deps_list {
                    if !validate_dep(dep) {
                        return;
                    }
                }
            } else if let Some(dep) = deps.as_ref().downcast_ref::<String>() {
                if !validate_dep(dep) {
                    return;
                }
            }
        }

        plugin.install(self);
        self.tables_plugin
            .insert(plugin.name().to_string(), Box::new(plugin));
    }

    pub fn uninstall(&mut self, name: &str) {
        if !self.tables_plugin.contains_key(name) {
            return;
        }

        // 提前收集需要回滚的策略名称,避免不可变借用与可变借用冲突
        let mut strategies_to_rollback = Vec::new();
        for (_, strategy) in self.tables_strategy.iter() {
            if strategy.condition().iter().any(|s| s == name) {
                println!(
                    "engine auto rollback strategy: {} before uninstall plugin: {}.",
                    strategy.name(),
                    name
                );
                strategies_to_rollback.push(strategy.name().to_string());
            }
        }

        // 执行回滚操作
        for strategy_name in strategies_to_rollback {
            self.rollback(&strategy_name);
        }

        // 提前收集需要卸载的插件名称,避免不可变借用与可变借用冲突
        let mut plugins_to_uninstall = Vec::new();
        for (_, plugin) in self.tables_plugin.iter() {
            if let Some(deps) = plugin.deps() {
                if let Some(deps_list) = deps.as_ref().downcast_ref::<Vec<String>>() {
                    if deps_list.contains(&name.to_string()) {
                        println!(
                            "engine auto uninstall plugin: {} before uninstall plugin: {}.",
                            plugin.name(),
                            name
                        );
                        plugins_to_uninstall.push(plugin.name().to_string());
                    }
                } else if let Some(dep) = deps.as_ref().downcast_ref::<String>() {
                    if dep == name {
                        println!(
                            "engine auto uninstall plugin: {} before uninstall plugin: {}.",
                            plugin.name(),
                            name
                        );
                        plugins_to_uninstall.push(plugin.name().to_string());
                    }
                }
            }
        }

        // 执行插件卸载操作
        for plugin_name in plugins_to_uninstall {
            self.uninstall(&plugin_name);
        }

        // 卸载当前插件
        if let Some(plugin) = self.tables_plugin.remove(name) {
            plugin.dispose(self);
        }
    }

    pub fn exec<T: StrategyOptions<Engine> + Send + Sync + 'static>(&mut self, strategy: T) {
        if self.tables_strategy.contains_key(strategy.name()) {
            println!("This strategy already exists: {}", strategy.name());
            return;
        }

        for plugin_name in strategy.condition().iter() {
            if !self.tables_plugin.contains_key(plugin_name) {
                println!(
                    "{} does not meet the conditions for execution: {}",
                    strategy.name(),
                    plugin_name
                );
                return;
            }
        }

        strategy.exec(self);
        self.tables_strategy
            .insert(strategy.name().to_string(), Box::new(strategy));
    }

    pub fn rollback(&mut self, name: &str) {
        if !self.tables_strategy.contains_key(name) {
            return;
        }

        let strategy = self.tables_strategy.remove(name).unwrap();
        strategy.rollback(self); // 此处不再有类型不匹配问题
    }

    pub fn dispatch_event(&self, event: &Event) {
        self.event_dispatcher.dispatch_event(event);
    }
}

-------src/event/mod.rs-------
use std::collections::{HashMap, HashSet};
use std::hash::{Hash, Hasher};
use std::ptr;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};

#[derive(Debug, Clone)]
pub struct Event {
    pub event_type: String,
    pub target: Option<String>,     // 可以指向目标对象
    pub attachment: Option<String>, // 附加数据
}

// 自定义包装器,实现 Eq 和 Hash
#[derive(Clone)]
struct ListenerWrapper(Arc<Mutex<dyn EventListener>>);

impl PartialEq for ListenerWrapper {
    fn eq(&self, other: &Self) -> bool {
        // 使用 std::ptr::eq 来比较地址
        ptr::eq(Arc::as_ptr(&self.0), Arc::as_ptr(&other.0))
    }
}

impl Eq for ListenerWrapper {}

impl Hash for ListenerWrapper {
    fn hash<H: Hasher>(&self, state: &mut H) {
        // 基于指针地址哈希
        Arc::as_ptr(&self.0).hash(state);
    }
}

pub trait EventListener: Send + Sync {
    fn call(&mut self, event: &Event);
}

type EventListeners = HashMap<String, HashSet<ListenerWrapper>>;

pub struct EventDispatcher {
    listeners: EventListeners,
}
impl EventDispatcher {
    pub fn new() -> Self {
        Self {
            listeners: HashMap::new(),
        }
    }

    pub fn add_event_listener(
        &mut self,
        event_type: String,
        listener: Arc<Mutex<dyn EventListener>>,
    ) {
        self.listeners
            .entry(event_type.clone())
            .or_insert_with(HashSet::new)
            .insert(ListenerWrapper(listener)); // 包装监听器
    }

    pub fn remove_event_listener(&mut self, event_type: &str, listener: &Arc<dyn EventListener>) {
        if let Some(listeners) = self.listeners.get_mut(event_type) {
            listeners.retain(|wrapped| {
                // 使用 addr_eq 比较地址,避免元数据的干扰
                !ptr::addr_eq(Arc::as_ptr(&wrapped.0), Arc::as_ptr(listener))
            });
        }
    }

    pub fn dispatch_event(&self, event: &Event) {
        if let Some(listeners) = self.listeners.get(&event.event_type) {
            for ListenerWrapper(listener) in listeners {
                let mut listener = listener.lock().unwrap(); // 锁定监听器
                listener.call(event);
            }
        }
    }

    pub fn once(&mut self, event_type: String, listener: Arc<Mutex<dyn EventListener>>) {
        let once_listener = Arc::new(Mutex::new(OnceListener::new(listener)));
        self.add_event_listener(event_type, once_listener);
    }

    pub fn on_throttle(
        &mut self,
        event_type: String,
        listener: Arc<Mutex<dyn EventListener>>,
        time: u64,
    ) {
        let throttle_listener = Arc::new(Mutex::new(ThrottleListener::new(listener, time)));
        self.add_event_listener(event_type, throttle_listener);
    }

    pub fn on_anti_shake(
        &mut self,
        event_type: String,
        listener: Arc<Mutex<dyn EventListener>>,
        time: u64,
    ) {
        let anti_shake_listener = Arc::new(Mutex::new(AntiShakeListener::new(listener, time)));
        self.add_event_listener(event_type, anti_shake_listener);
    }
}

// 一次性监听器
pub struct OnceListener {
    listener: Arc<Mutex<dyn EventListener>>,
}
impl OnceListener {
    pub fn new(listener: Arc<Mutex<dyn EventListener>>) -> Self {
        OnceListener { listener }
    }
}
impl EventListener for OnceListener {
    fn call(&mut self, event: &Event) {
        let mut listener = self.listener.lock().unwrap(); // 锁定并调用
        listener.call(event);
    }
}

// 节流监听器
pub struct ThrottleListener {
    listener: Arc<Mutex<dyn EventListener>>, // 用 Mutex 包装 listener
    state: Arc<Mutex<ThrottleState>>,        // 使用 Mutex 来管理状态
}
struct ThrottleState {
    last_called: Option<Instant>,
    delay: u64,
}
impl ThrottleListener {
    pub fn new(listener: Arc<Mutex<dyn EventListener>>, delay: u64) -> Self {
        ThrottleListener {
            listener,
            state: Arc::new(Mutex::new(ThrottleState {
                last_called: None,
                delay,
            })),
        }
    }
}
impl EventListener for ThrottleListener {
    fn call(&mut self, event: &Event) {
        let mut state = self.state.lock().unwrap();
        let now = Instant::now();
        if let Some(last_called) = state.last_called {
            if now.duration_since(last_called) < Duration::from_millis(state.delay) {
                return; // 如果距离上次调用太近,则不处理
            }
        }

        // 锁定并修改 listener
        let mut listener = self.listener.lock().unwrap();
        listener.call(event);

        state.last_called = Some(now);
    }
}

// 防抖监听器
pub struct AntiShakeListener {
    listener: Arc<Mutex<dyn EventListener>>, // 用 Mutex 包装 listener
    state: Arc<Mutex<AntiShakeState>>,       // 使用 Mutex 来管理状态
}
struct AntiShakeState {
    last_called: Option<Instant>,
    delay: u64,
}
impl AntiShakeListener {
    pub fn new(listener: Arc<Mutex<dyn EventListener>>, delay: u64) -> Self {
        AntiShakeListener {
            listener,
            state: Arc::new(Mutex::new(AntiShakeState {
                last_called: None,
                delay,
            })),
        }
    }
}
impl EventListener for AntiShakeListener {
    fn call(&mut self, event: &Event) {
        let mut state = self.state.lock().unwrap();
        let now = Instant::now();
        if let Some(last_called) = state.last_called {
            if now.duration_since(last_called) < Duration::from_millis(state.delay) {
                return; // 如果距离上次调用太近,则不处理
            }
        }

        // 锁定并修改 listener
        let mut listener = self.listener.lock().unwrap();
        listener.call(event);

        state.last_called = Some(now);
    }
}



-------src/plugin/mod.rs-------
use super::engine::EngineTrait;
use std::any::Any;
use std::sync::Arc;

pub trait PluginOptions<E: EngineTrait> {
    fn name(&self) -> &str;
    fn deps(&self) -> Option<Box<dyn Any>>;
    fn install(&self, engine: &mut E);
    fn dispose(&self, engine: &mut E);
}

// 使用结构体替代类型别名
pub struct Plugin<E: EngineTrait, P: Any = ()>(Box<dyn Fn(Option<P>) -> Arc<dyn PluginOptions<E>>>);
impl<E: EngineTrait, P: Any> Plugin<E, P> {
    pub fn new<F>(func: F) -> Self
    where
        F: Fn(Option<P>) -> Arc<dyn PluginOptions<E>> + 'static,
    {
        Plugin(Box::new(func))
    }

    pub fn call(&self, param: Option<P>) -> Arc<dyn PluginOptions<E>> {
        (self.0)(param)
    }
}

pub fn define_plugin<E: EngineTrait>(options: Arc<dyn PluginOptions<E>>) -> Plugin<E> {
    Plugin::new(move |_| options.clone())
}


-------src/strategy/mod.rs-------
use super::engine::EngineTrait;
use std::any::Any;
use std::sync::Arc;

// 定义 StrategyOptions
pub trait StrategyOptions<E: EngineTrait> {
    fn name(&self) -> &str;
    fn condition(&self) -> Vec<String>;
    fn exec(&self, engine: &mut E);
    fn rollback(&self, engine: &mut E);
}

// 使用结构体替代类型别名
pub struct Strategy<E: EngineTrait, P: Any = ()>(
    Box<dyn Fn(Option<P>) -> Arc<dyn StrategyOptions<E>>>,
);

impl<E: EngineTrait, P: Any> Strategy<E, P> {
    pub fn new<F>(f: F) -> Self
    where
        F: Fn(Option<P>) -> Arc<dyn StrategyOptions<E>> + 'static,
    {
        Strategy(Box::new(f))
    }

    pub fn call(&self, param: Option<P>) -> Arc<dyn StrategyOptions<E>> {
        (self.0)(param)
    }
}

// 定义 define_strategy 函数
pub fn define_strategy<E: EngineTrait>(options: Arc<dyn StrategyOptions<E>>) -> Strategy<E> {
    Strategy::new(move |_| options.clone())
}


-------src/mod.rs-------
pub mod engine;
pub mod event;
pub mod plugin;
pub mod strategy;



















1. 复杂的调试和追踪
优化措施:

日志和追踪系统:引入完善的日志和追踪机制,帮助开发者跟踪事件流和插件的执行过程。可以使用如log或tracing等Rust库来记录详细的日志信息。这样即使系统出现问题,开发者也能更快速地定位到问题所在。

事件上下文:每个事件可以携带更多上下文信息,比如调用栈、事件的发起者等,帮助开发者理解事件的生命周期和处理流程。

调试工具:开发一套调试工具,允许开发者查看事件调度的顺序、插件的执行情况以及事件处理的中间状态。

2. 性能开销
优化措施:

批量处理和延迟执行:对于频繁触发的事件,可以使用批量处理或者延迟执行策略,减少性能开销。例如,多个事件的处理可以合并在一个批次中,避免每次触发事件时都进行单独的计算。

事件去重和缓存:在事件派发之前检查是否已经有相同的事件正在处理中,避免重复计算。可以对常见的事件类型和处理结果进行缓存,减少重复计算的成本。

优化数据结构:对于插件和事件监听器的管理,可以选择更高效的数据结构,例如使用HashMap而不是Vec来快速查找监听器,或者使用BTreeMap来保持插件的有序性。

3. 设计上的复杂性
优化措施:

简化核心框架:在框架核心部分,避免过度抽象和复杂设计,只保留最常用的功能和最简洁的接口,让开发者能够快速上手。比如,可以提供一个简化版的事件系统和插件管理系统,默认开启简单的事件监听和插件安装卸载功能,允许开发者根据需求扩展。

模块化和清晰的文档:将框架分解为多个模块,确保每个模块都可以独立使用,并且为每个模块提供清晰的文档和示例代码。开发者可以从最基础的模块开始,逐步学习和使用更高级的功能。

提供高层API:为开发者提供更高级别的API封装,将一些常见操作进行简化封装,使得框架的使用变得更加直观。例如,为常见的插件和策略使用场景提供默认的实现,避免开发者每次都从头实现。

4. 插件和策略的耦合性
优化措施:

依赖注入(DI):采用依赖注入的方式来解耦插件和策略,使它们能够更加灵活地工作。例如,可以使用一些Rust库如shaku或者tikv来实现依赖注入,从而减少插件之间的直接依赖。

插件依赖关系管理:提供一个明确的插件依赖声明和版本控制机制,确保插件和策略之间的依赖关系清晰且不会相互干扰。可以为插件提供版本号,确保只有在插件兼容的情况下才会被安装。

插件生命周期管理:优化插件和策略的生命周期管理,确保它们的安装、卸载、回滚等操作不相互干扰。可以通过事件驱动的方式自动处理依赖关系,减少手动管理插件和策略的复杂性。

5. 事件顺序问题
优化措施:

事件优先级:为事件定义优先级,允许开发者控制事件的处理顺序。优先级可以在事件添加到监听队列时指定,或者在事件派发时动态调整。高优先级的事件可以先被处理,低优先级的事件则可以被延迟处理。

事件调度器:实现一个更强大的事件调度器,允许开发者自定义事件的处理顺序和调度策略。可以使用消息队列的方式管理事件的调度顺序,确保事件的处理不受到冲突。

事件批量处理:在某些情况下,可以将多个事件组合成一个批次进行处理,这样可以确保事件的处理顺序和一致性,避免因为并发处理而引发的问题。

6. 过度依赖插件系统
优化措施:

插件依赖图:为插件创建一个依赖图,明确展示插件之间的依赖关系。插件的安装和卸载可以基于这个依赖图来进行,避免错误的插件安装顺序和插件卸载时出现的依赖问题。

插件的隔离性:通过沙箱或容器化的方式隔离插件,使得每个插件的运行时环境相对独立,避免插件之间直接共享状态和数据。这样即使一个插件出现问题,也不会影响到其他插件或系统的稳定性。

插件生命周期管理:完善插件的生命周期管理,确保每个插件的生命周期和资源管理都是清晰的。例如,插件可以在运行时通过注册钩子或回调函数来通知框架进行清理,防止内存泄漏或资源占用。

请给框架进行这6点优化,并给出优化后框架的完整代码,注意,优化一定要在保证原框架优点的基础上