Skip to main content

sz_orm_core/
observer.rs

1//! Observer + Event Subscriber — 模型生命周期观察者模式
2//!
3//! 对应文档 6.8 节改进项 32(Observer)+ 33(Event Subscriber)。
4//!
5//! # 核心概念
6//!
7//! - **Observer**:观察者接口,订阅模型生命周期事件(INSERT/UPDATE/DELETE/FIND)
8//! - **EventSubscriber**:事件订阅者,按事件类型订阅(比 Observer 更细粒度)
9//! - **EventDispatcher**:事件分发器,管理 Observer 与 EventSubscriber 的注册和分发
10//!
11//! # 与 Behaviors 的区别
12//!
13//! | 特性 | Behaviors (behaviors.rs) | Observer (本模块) |
14//! |------|--------------------------|-------------------|
15//! | 注册方式 | Model 内部声明 | 外部注册到 Dispatcher |
16//! | 解耦程度 | Model 与 Behavior 强耦合 | 完全解耦,Model 无需感知 |
17//! | 适用场景 | 字段自动填充(时间戳/操作人) | 审计日志、缓存失效、外部通知 |
18//! | 事件粒度 | 4 个生命周期事件 | 可订阅特定事件类型 |
19//!
20//! # 设计灵感
21//!
22//! - Doctrine `EventSubscriber` / `LifecycleCallback`
23//! - Hibernate `EntityListener` / `@PostPersist`
24//! - Laravel Eloquent `Observer` 类
25//! - Rails ActiveRecord `Callbacks` + `Observers`
26//!
27//! # 使用示例
28//!
29//! ```
30//! use sz_orm_core::observer::{
31//!     Event, EventDispatcher, EventSubscriber, Observer, SubscriberResult,
32//! };
33//! use sz_orm_core::hooks::HookContext;
34//! use std::collections::HashMap;
35//! use std::sync::{Arc, Mutex};
36//! use sz_orm_core::Value;
37//!
38//! // 1. 审计日志订阅者(订阅所有事件)
39//! struct AuditLogSubscriber {
40//!     logs: Arc<Mutex<Vec<String>>>,
41//! }
42//!
43//! impl EventSubscriber for AuditLogSubscriber {
44//!     fn subscribed_events(&self) -> Vec<Event> {
45//!         vec![Event::AfterInsert, Event::AfterUpdate, Event::AfterDelete]
46//!     }
47//!
48//!     fn on_event(&self, event: Event, ctx: &HookContext, attrs: &HashMap<String, Value>) -> SubscriberResult<()> {
49//!         let mut logs = self.logs.lock().unwrap();
50//!         logs.push(format!("{:?} on attrs with {} fields", event, attrs.len()));
51//!         Ok(())
52//!     }
53//! }
54//!
55//! // 2. 缓存失效订阅者(仅订阅写入事件)
56//! struct CacheInvalidationSubscriber;
57//!
58//! impl EventSubscriber for CacheInvalidationSubscriber {
59//!     fn subscribed_events(&self) -> Vec<Event> {
60//!         vec![Event::AfterUpdate, Event::AfterDelete]
61//!     }
62//!
63//!     fn on_event(&self, event: Event, _ctx: &HookContext, attrs: &HashMap<String, Value>) -> SubscriberResult<()> {
64//!         // 失效缓存逻辑...
65//!         let _ = (event, attrs);
66//!         Ok(())
67//!     }
68//! }
69//!
70//! // 3. 注册并触发事件
71//! let logs = Arc::new(Mutex::new(Vec::new()));
72//! let mut dispatcher = EventDispatcher::new();
73//! dispatcher.subscribe(Box::new(AuditLogSubscriber { logs: logs.clone() }));
74//! dispatcher.subscribe(Box::new(CacheInvalidationSubscriber));
75//!
76//! let ctx = HookContext::default();
77//! let attrs = HashMap::new();
78//! dispatcher.dispatch(Event::AfterInsert, &ctx, &attrs);
79//!
80//! assert_eq!(logs.lock().unwrap().len(), 1);
81//! ```
82
83use crate::hooks::HookContext;
84use crate::Value;
85use parking_lot::{Mutex, RwLock};
86use std::collections::HashMap;
87use std::sync::Arc;
88
89// ============================================================================
90// Event — 事件类型
91// ============================================================================
92
93/// 模型生命周期事件类型
94///
95/// 与 `hooks::HookEvent` 类似但简化为运行时分发用的事件枚举。
96#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
97pub enum Event {
98    /// 插入前
99    BeforeInsert,
100    /// 插入后
101    AfterInsert,
102    /// 更新前
103    BeforeUpdate,
104    /// 更新后
105    AfterUpdate,
106    /// 删除前
107    BeforeDelete,
108    /// 删除后
109    AfterDelete,
110    /// 单行查询后
111    AfterFind,
112    /// 软删除恢复前
113    BeforeRestore,
114    /// 软删除恢复后
115    AfterRestore,
116}
117
118impl Event {
119    /// 是否为 before 事件
120    pub fn is_before(&self) -> bool {
121        matches!(
122            self,
123            Event::BeforeInsert | Event::BeforeUpdate | Event::BeforeDelete | Event::BeforeRestore
124        )
125    }
126
127    /// 是否为 after 事件
128    pub fn is_after(&self) -> bool {
129        matches!(
130            self,
131            Event::AfterInsert
132                | Event::AfterUpdate
133                | Event::AfterDelete
134                | Event::AfterFind
135                | Event::AfterRestore
136        )
137    }
138
139    /// 是否为写入事件(INSERT/UPDATE/DELETE)
140    pub fn is_write_event(&self) -> bool {
141        matches!(
142            self,
143            Event::BeforeInsert
144                | Event::AfterInsert
145                | Event::BeforeUpdate
146                | Event::AfterUpdate
147                | Event::BeforeDelete
148                | Event::AfterDelete
149        )
150    }
151
152    /// 事件名称(用于日志与错误信息)
153    pub fn name(&self) -> &'static str {
154        match self {
155            Event::BeforeInsert => "before_insert",
156            Event::AfterInsert => "after_insert",
157            Event::BeforeUpdate => "before_update",
158            Event::AfterUpdate => "after_update",
159            Event::BeforeDelete => "before_delete",
160            Event::AfterDelete => "after_delete",
161            Event::AfterFind => "after_find",
162            Event::BeforeRestore => "before_restore",
163            Event::AfterRestore => "after_restore",
164        }
165    }
166}
167
168// ============================================================================
169// SubscriberError — 订阅者错误
170// ============================================================================
171
172/// 订阅者错误类型
173#[derive(Debug)]
174pub enum SubscriberError {
175    /// 订阅者执行失败(携带错误描述)
176    Failed {
177        /// 订阅者名称
178        subscriber: String,
179        /// 错误描述
180        reason: String,
181    },
182    /// 中止后续订阅者执行(用于 veto 模式)
183    ///
184    /// 例如:before_insert 钩子拒绝该次插入
185    Vetoed {
186        /// 订阅者名称
187        subscriber: String,
188        /// 拒绝原因
189        reason: String,
190    },
191}
192
193impl std::fmt::Display for SubscriberError {
194    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
195        match self {
196            SubscriberError::Failed { subscriber, reason } => {
197                write!(f, "Subscriber `{}` failed: {}", subscriber, reason)
198            }
199            SubscriberError::Vetoed { subscriber, reason } => {
200                write!(f, "Subscriber `{}` vetoed: {}", subscriber, reason)
201            }
202        }
203    }
204}
205
206impl std::error::Error for SubscriberError {}
207
208/// 订阅者结果类型
209pub type SubscriberResult<T> = Result<T, SubscriberError>;
210
211// ============================================================================
212// Observer — 模型观察者 trait
213// ============================================================================
214
215/// 模型观察者 trait
216///
217/// 与 `EventSubscriber` 不同,`Observer` 默认订阅所有事件。
218/// 适合需要监控所有生命周期事件的场景(如审计日志)。
219///
220/// # 实现要点
221///
222/// - 所有方法默认实现为 no-op,按需 override
223/// - 任何方法返回 `Err(SubscriberError::Vetoed)` 会中止 before 事件的后续执行
224pub trait Observer: Send + Sync {
225    /// 观察者名称(用于日志与错误信息)
226    fn name(&self) -> &str {
227        "anonymous_observer"
228    }
229
230    /// 插入前
231    fn before_insert(
232        &self,
233        _ctx: &HookContext,
234        _attrs: &mut HashMap<String, Value>,
235    ) -> SubscriberResult<()> {
236        Ok(())
237    }
238
239    /// 插入后
240    fn after_insert(
241        &self,
242        _ctx: &HookContext,
243        _attrs: &HashMap<String, Value>,
244    ) -> SubscriberResult<()> {
245        Ok(())
246    }
247
248    /// 更新前
249    fn before_update(
250        &self,
251        _ctx: &HookContext,
252        _attrs: &mut HashMap<String, Value>,
253    ) -> SubscriberResult<()> {
254        Ok(())
255    }
256
257    /// 更新后
258    fn after_update(
259        &self,
260        _ctx: &HookContext,
261        _attrs: &HashMap<String, Value>,
262    ) -> SubscriberResult<()> {
263        Ok(())
264    }
265
266    /// 删除前
267    fn before_delete(
268        &self,
269        _ctx: &HookContext,
270        _attrs: &HashMap<String, Value>,
271    ) -> SubscriberResult<()> {
272        Ok(())
273    }
274
275    /// 删除后
276    fn after_delete(
277        &self,
278        _ctx: &HookContext,
279        _attrs: &HashMap<String, Value>,
280    ) -> SubscriberResult<()> {
281        Ok(())
282    }
283
284    /// 单行查询后
285    fn after_find(
286        &self,
287        _ctx: &HookContext,
288        _attrs: &mut HashMap<String, Value>,
289    ) -> SubscriberResult<()> {
290        Ok(())
291    }
292}
293
294// ============================================================================
295// EventSubscriber — 事件订阅者 trait
296// ============================================================================
297
298/// 事件订阅者 trait
299///
300/// 与 `Observer` 不同,`EventSubscriber` 只接收订阅的特定事件。
301/// 适合只关心特定事件的场景(如缓存失效仅关心 UPDATE/DELETE)。
302pub trait EventSubscriber: Send + Sync {
303    /// 订阅者名称
304    fn name(&self) -> &str {
305        "anonymous_subscriber"
306    }
307
308    /// 返回订阅的事件列表
309    ///
310    /// 仅当事件在此列表中时,`on_event` 才会被调用。
311    fn subscribed_events(&self) -> Vec<Event>;
312
313    /// 事件回调
314    ///
315    /// # 参数
316    /// - `event`:触发的事件
317    /// - `ctx`:钩子上下文
318    /// - `attrs`:当前属性(before 事件可修改)
319    ///
320    /// # 返回
321    /// - `Ok(())`:继续执行后续订阅者
322    /// - `Err(SubscriberError::Vetoed)`:中止 before 事件的后续执行
323    /// - `Err(SubscriberError::Failed)`:记录错误,继续执行后续订阅者
324    fn on_event(
325        &self,
326        event: Event,
327        ctx: &HookContext,
328        attrs: &HashMap<String, Value>,
329    ) -> SubscriberResult<()>;
330}
331
332// ============================================================================
333// EventDispatcher — 事件分发器
334// ============================================================================
335
336/// 事件分发器
337///
338/// 管理 `Observer` 与 `EventSubscriber` 的注册与分发。
339///
340/// # 分发顺序
341///
342/// 1. 先按注册顺序调用所有 `Observer` 的对应方法
343/// 2. 再按注册顺序调用所有订阅了该事件的 `EventSubscriber`
344///
345/// # 错误处理
346///
347/// - `before_*` 事件中任何订阅者返回 `Err(Vetoed)` 会立即中止后续执行
348/// - `after_*` 事件中的错误仅记录,不影响后续执行
349///
350/// # 线程安全
351///
352/// 内部使用 `RwLock<Vec<Arc<...>>>`,支持多线程并发。
353///
354/// # 死锁防护(v0.2.1 修复 Critical C-3)
355///
356/// `dispatch` / `dispatch_before_mut` 在调用用户回调前会先 clone 一份
357/// `Vec<Arc<dyn ...>>` 快照并释放读锁,避免持读锁调用用户代码——
358/// 否则用户回调中若尝试注册新订阅者(需要写锁)会自我死锁。
359pub struct EventDispatcher {
360    observers: RwLock<Vec<Arc<dyn Observer>>>,
361    subscribers: RwLock<Vec<Arc<dyn EventSubscriber>>>,
362    /// 错误收集(非致命错误,不影响流程)
363    errors: RwLock<Vec<SubscriberError>>,
364    /// 错误缓冲区最大容量(防止内存无限增长)
365    ///
366    /// 当 errors 长度达到此上限时,新增错误会以 FIFO 方式淘汰最早错误。
367    /// 默认 1024,可通过 `with_max_errors` 调整。
368    max_errors: usize,
369}
370
371/// 默认错误缓冲区容量
372const DEFAULT_MAX_ERRORS: usize = 1024;
373
374impl EventDispatcher {
375    /// 创建空的事件分发器
376    pub fn new() -> Self {
377        Self {
378            observers: RwLock::new(Vec::new()),
379            subscribers: RwLock::new(Vec::new()),
380            errors: RwLock::new(Vec::new()),
381            max_errors: DEFAULT_MAX_ERRORS,
382        }
383    }
384
385    /// 设置错误缓冲区最大容量
386    ///
387    /// 当 errors 达到此容量时,新增错误会淘汰最早错误(FIFO)。
388    /// 设置为 0 表示无限制(不推荐,可能导致内存泄漏)。
389    pub fn with_max_errors(mut self, max_errors: usize) -> Self {
390        self.max_errors = max_errors;
391        self
392    }
393
394    /// 注册 Observer
395    ///
396    /// 接收 `Box<dyn Observer>`(向后兼容),内部转 `Arc<dyn Observer>` 存储,
397    /// 以便 dispatch 时可以 cheap clone 快照后释放读锁。
398    pub fn add_observer(&self, observer: Box<dyn Observer>) {
399        let arc: Arc<dyn Observer> = Arc::from(observer);
400        self.observers.write().push(arc);
401    }
402
403    /// 注册 EventSubscriber
404    ///
405    /// 接收 `Box<dyn EventSubscriber>`(向后兼容),内部转 `Arc<dyn EventSubscriber>` 存储。
406    pub fn subscribe(&self, subscriber: Box<dyn EventSubscriber>) {
407        let arc: Arc<dyn EventSubscriber> = Arc::from(subscriber);
408        self.subscribers.write().push(arc);
409    }
410
411    /// 清空所有注册
412    pub fn clear(&self) {
413        self.observers.write().clear();
414        self.subscribers.write().clear();
415        self.errors.write().clear();
416    }
417
418    /// 返回已注册的 Observer 数量
419    pub fn observer_count(&self) -> usize {
420        self.observers.read().len()
421    }
422
423    /// 返回已注册的 EventSubscriber 数量
424    pub fn subscriber_count(&self) -> usize {
425        self.subscribers.read().len()
426    }
427
428    /// 取出收集到的非致命错误(清空内部缓冲)
429    pub fn drain_errors(&self) -> Vec<SubscriberError> {
430        std::mem::take(&mut *self.errors.write())
431    }
432
433    /// 返回当前错误缓冲区中的错误数量
434    pub fn error_count(&self) -> usize {
435        self.errors.read().len()
436    }
437
438    /// 将本地错误批量写入 errors 缓冲区,遵循 max_errors 限制(FIFO 淘汰)
439    ///
440    /// - `max_errors = 0` 表示无限制
441    /// - 否则当 errors 达到上限时,淘汰最早错误以腾出空间
442    fn push_errors(&self, new_errors: Vec<SubscriberError>) {
443        if new_errors.is_empty() {
444            return;
445        }
446        let mut errors = self.errors.write();
447        if self.max_errors == 0 {
448            errors.extend(new_errors);
449            return;
450        }
451        for e in new_errors {
452            if errors.len() >= self.max_errors {
453                // FIFO 淘汰最早错误
454                errors.remove(0);
455            }
456            errors.push(e);
457        }
458    }
459
460    /// 分发事件(after_* 事件,attrs 不可变)
461    ///
462    /// 错误仅记录,不影响后续订阅者执行。
463    ///
464    /// # 实现要点
465    ///
466    /// - **v0.2.1 修复 Critical C-3**:调用用户回调前先 clone `Vec<Arc<...>>` 快照
467    ///   并释放读锁,避免持读锁调用用户代码(防止死锁)
468    /// - 错误先收集到本地 `Vec`,循环结束后一次性批量写入 `self.errors`
469    pub fn dispatch(&self, event: Event, ctx: &HookContext, attrs: &HashMap<String, Value>) {
470        let mut local_errors: Vec<SubscriberError> = Vec::new();
471
472        // 1. 调用 Observers — 持读锁仅 clone 快照,立即释放
473        let observers_snapshot: Vec<Arc<dyn Observer>> = {
474            let observers = self.observers.read();
475            observers.clone()
476        };
477        // 释放读锁后调用用户代码
478        for observer in observers_snapshot.iter() {
479            let result = match event {
480                Event::AfterInsert => observer.after_insert(ctx, attrs),
481                Event::AfterUpdate => observer.after_update(ctx, attrs),
482                Event::AfterDelete => observer.after_delete(ctx, attrs),
483                _ => Ok(()),
484            };
485            if let Err(e) = result {
486                local_errors.push(e);
487            }
488        }
489
490        // 2. 调用 EventSubscribers — 持读锁仅 clone 快照,立即释放
491        let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
492            let subscribers = self.subscribers.read();
493            subscribers.clone()
494        };
495        for subscriber in subscribers_snapshot.iter() {
496            if !subscriber.subscribed_events().contains(&event) {
497                continue;
498            }
499            if let Err(e) = subscriber.on_event(event, ctx, attrs) {
500                local_errors.push(e);
501            }
502        }
503
504        if !local_errors.is_empty() {
505            self.push_errors(local_errors);
506        }
507    }
508
509    /// 分发 before 事件(attrs 可变)
510    ///
511    /// 任何订阅者返回 `Err(Vetoed)` 会立即中止并返回错误。
512    ///
513    /// # 实现要点
514    ///
515    /// - Vetoed 时直接返回该错误,**不会再次调用 `on_event`**(避免订阅者副作用翻倍)
516    /// - **v0.2.1 修复 Critical C-3**:调用用户回调前先 clone 快照并释放读锁
517    /// - 错误先收集到本地 `Vec`,避免持读锁时获取写锁造成死锁
518    pub fn dispatch_before_mut(
519        &self,
520        event: Event,
521        ctx: &HookContext,
522        attrs: &mut HashMap<String, Value>,
523    ) -> SubscriberResult<()> {
524        let mut local_errors: Vec<SubscriberError> = Vec::new();
525        let mut vetoed: Option<SubscriberError> = None;
526
527        // 1. 调用 Observers — 持读锁仅 clone 快照,立即释放
528        let observers_snapshot: Vec<Arc<dyn Observer>> = {
529            let observers = self.observers.read();
530            observers.clone()
531        };
532        for observer in observers_snapshot.iter() {
533            let result = match event {
534                Event::BeforeInsert => observer.before_insert(ctx, attrs),
535                Event::BeforeUpdate => observer.before_update(ctx, attrs),
536                _ => Ok(()),
537            };
538            match result {
539                Ok(()) => {}
540                Err(e @ SubscriberError::Vetoed { .. }) => {
541                    vetoed = Some(e);
542                    break;
543                }
544                Err(e) => local_errors.push(e),
545            }
546        }
547
548        // 2. 调用 EventSubscribers(仅当未被 vetoed)— 持读锁仅 clone 快照,立即释放
549        if vetoed.is_none() {
550            let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
551                let subscribers = self.subscribers.read();
552                subscribers.clone()
553            };
554            for subscriber in subscribers_snapshot.iter() {
555                if !subscriber.subscribed_events().contains(&event) {
556                    continue;
557                }
558                match subscriber.on_event(event, ctx, attrs) {
559                    Ok(()) => {}
560                    Err(e @ SubscriberError::Vetoed { .. }) => {
561                        vetoed = Some(e);
562                        break;
563                    }
564                    Err(e) => local_errors.push(e),
565                }
566            }
567        }
568
569        if !local_errors.is_empty() {
570            self.push_errors(local_errors);
571        }
572
573        if let Some(e) = vetoed {
574            return Err(e);
575        }
576        Ok(())
577    }
578
579    /// 分发 after_find 事件(attrs 可变,用于修改读出的数据)
580    ///
581    /// # 实现要点
582    ///
583    /// - **v0.2.1 修复 Critical C-3**:调用用户回调前先 clone 快照并释放读锁
584    /// - 错误先收集到本地 `Vec`,避免持读锁时获取写锁造成死锁
585    pub fn dispatch_after_find(
586        &self,
587        ctx: &HookContext,
588        attrs: &mut HashMap<String, Value>,
589    ) -> SubscriberResult<()> {
590        let mut local_errors: Vec<SubscriberError> = Vec::new();
591
592        let observers_snapshot: Vec<Arc<dyn Observer>> = {
593            let observers = self.observers.read();
594            observers.clone()
595        };
596        for observer in observers_snapshot.iter() {
597            if let Err(e) = observer.after_find(ctx, attrs) {
598                local_errors.push(e);
599            }
600        }
601
602        let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
603            let subscribers = self.subscribers.read();
604            subscribers.clone()
605        };
606        for subscriber in subscribers_snapshot.iter() {
607            if !subscriber.subscribed_events().contains(&Event::AfterFind) {
608                continue;
609            }
610            if let Err(e) = subscriber.on_event(Event::AfterFind, ctx, attrs) {
611                local_errors.push(e);
612            }
613        }
614
615        if !local_errors.is_empty() {
616            self.push_errors(local_errors);
617        }
618
619        Ok(())
620    }
621}
622
623impl Default for EventDispatcher {
624    fn default() -> Self {
625        Self::new()
626    }
627}
628
629// ============================================================================
630// 内置订阅者实现
631// ============================================================================
632
633// -------------------- AuditLogSubscriber --------------------
634
635/// 审计日志订阅者
636///
637/// 记录所有写入操作到内部日志缓冲,可用于调试或审计。
638///
639/// # 示例
640///
641/// ```
642/// use sz_orm_core::observer::{EventDispatcher, AuditLogSubscriber, Event};
643/// use sz_orm_core::hooks::HookContext;
644/// use std::collections::HashMap;
645///
646/// let audit = AuditLogSubscriber::new();
647/// let mut dispatcher = EventDispatcher::new();
648/// dispatcher.subscribe(Box::new(audit.clone()));
649///
650/// let ctx = HookContext::default();
651/// let attrs = HashMap::new();
652/// dispatcher.dispatch(Event::AfterInsert, &ctx, &attrs);
653///
654/// assert_eq!(audit.logs().lock().len(), 1);
655/// ```
656#[derive(Clone)]
657pub struct AuditLogSubscriber {
658    logs: std::sync::Arc<Mutex<Vec<String>>>,
659}
660
661impl AuditLogSubscriber {
662    /// 创建审计日志订阅者
663    pub fn new() -> Self {
664        Self {
665            logs: std::sync::Arc::new(Mutex::new(Vec::new())),
666        }
667    }
668
669    /// 获取日志列表(用于断言)
670    pub fn logs(&self) -> &std::sync::Arc<Mutex<Vec<String>>> {
671        &self.logs
672    }
673}
674
675impl Default for AuditLogSubscriber {
676    fn default() -> Self {
677        Self::new()
678    }
679}
680
681impl EventSubscriber for AuditLogSubscriber {
682    fn name(&self) -> &str {
683        "audit_log"
684    }
685
686    fn subscribed_events(&self) -> Vec<Event> {
687        vec![Event::AfterInsert, Event::AfterUpdate, Event::AfterDelete]
688    }
689
690    fn on_event(
691        &self,
692        event: Event,
693        ctx: &HookContext,
694        attrs: &HashMap<String, Value>,
695    ) -> SubscriberResult<()> {
696        let mut logs = self.logs.lock();
697        logs.push(format!(
698            "event={} operator={:?} field_count={}",
699            event.name(),
700            ctx.operator_id,
701            attrs.len()
702        ));
703        Ok(())
704    }
705}
706
707// ============================================================================
708// 单元测试
709// ============================================================================
710
711#[cfg(test)]
712mod tests {
713    use super::*;
714    use std::sync::{Arc, Mutex};
715
716    // ===== Event 测试 =====
717
718    #[test]
719    fn test_event_is_before_after() {
720        assert!(Event::BeforeInsert.is_before());
721        assert!(!Event::BeforeInsert.is_after());
722        assert!(Event::AfterInsert.is_after());
723        assert!(!Event::AfterInsert.is_before());
724    }
725
726    #[test]
727    fn test_event_is_write_event() {
728        assert!(Event::BeforeInsert.is_write_event());
729        assert!(Event::AfterUpdate.is_write_event());
730        assert!(Event::BeforeDelete.is_write_event());
731        assert!(!Event::AfterFind.is_write_event());
732    }
733
734    #[test]
735    fn test_event_name() {
736        assert_eq!(Event::BeforeInsert.name(), "before_insert");
737        assert_eq!(Event::AfterDelete.name(), "after_delete");
738        assert_eq!(Event::AfterFind.name(), "after_find");
739    }
740
741    // ===== EventDispatcher 基础测试 =====
742
743    #[test]
744    fn test_new_dispatcher_is_empty() {
745        let d = EventDispatcher::new();
746        assert_eq!(d.observer_count(), 0);
747        assert_eq!(d.subscriber_count(), 0);
748    }
749
750    #[test]
751    fn test_add_observer() {
752        struct DummyObserver;
753        impl Observer for DummyObserver {}
754
755        let d = EventDispatcher::new();
756        d.add_observer(Box::new(DummyObserver));
757        assert_eq!(d.observer_count(), 1);
758    }
759
760    #[test]
761    fn test_subscribe() {
762        struct DummySubscriber;
763        impl EventSubscriber for DummySubscriber {
764            fn subscribed_events(&self) -> Vec<Event> {
765                vec![Event::AfterInsert]
766            }
767            fn on_event(
768                &self,
769                _event: Event,
770                _ctx: &HookContext,
771                _attrs: &HashMap<String, Value>,
772            ) -> SubscriberResult<()> {
773                Ok(())
774            }
775        }
776
777        let d = EventDispatcher::new();
778        d.subscribe(Box::new(DummySubscriber));
779        assert_eq!(d.subscriber_count(), 1);
780    }
781
782    #[test]
783    fn test_clear() {
784        struct DummyObserver;
785        impl Observer for DummyObserver {}
786
787        let d = EventDispatcher::new();
788        d.add_observer(Box::new(DummyObserver));
789        d.clear();
790        assert_eq!(d.observer_count(), 0);
791    }
792
793    // ===== Observer 触发测试 =====
794
795    /// 计数 Observer,用于测试
796    struct CountingObserver {
797        insert_count: Arc<Mutex<u32>>,
798        update_count: Arc<Mutex<u32>>,
799        delete_count: Arc<Mutex<u32>>,
800    }
801
802    impl Observer for CountingObserver {
803        fn name(&self) -> &str {
804            "counting"
805        }
806
807        fn after_insert(
808            &self,
809            _ctx: &HookContext,
810            _attrs: &HashMap<String, Value>,
811        ) -> SubscriberResult<()> {
812            *self.insert_count.lock().unwrap() += 1;
813            Ok(())
814        }
815
816        fn after_update(
817            &self,
818            _ctx: &HookContext,
819            _attrs: &HashMap<String, Value>,
820        ) -> SubscriberResult<()> {
821            *self.update_count.lock().unwrap() += 1;
822            Ok(())
823        }
824
825        fn after_delete(
826            &self,
827            _ctx: &HookContext,
828            _attrs: &HashMap<String, Value>,
829        ) -> SubscriberResult<()> {
830            *self.delete_count.lock().unwrap() += 1;
831            Ok(())
832        }
833    }
834
835    #[test]
836    fn test_observer_triggered_on_dispatch() {
837        let insert = Arc::new(Mutex::new(0u32));
838        let update = Arc::new(Mutex::new(0u32));
839        let delete = Arc::new(Mutex::new(0u32));
840
841        let observer = CountingObserver {
842            insert_count: insert.clone(),
843            update_count: update.clone(),
844            delete_count: delete.clone(),
845        };
846
847        let d = EventDispatcher::new();
848        d.add_observer(Box::new(observer));
849
850        let ctx = HookContext::default();
851        let attrs = HashMap::new();
852
853        d.dispatch(Event::AfterInsert, &ctx, &attrs);
854        d.dispatch(Event::AfterInsert, &ctx, &attrs);
855        d.dispatch(Event::AfterUpdate, &ctx, &attrs);
856        d.dispatch(Event::AfterDelete, &ctx, &attrs);
857
858        assert_eq!(*insert.lock().unwrap(), 2);
859        assert_eq!(*update.lock().unwrap(), 1);
860        assert_eq!(*delete.lock().unwrap(), 1);
861    }
862
863    #[test]
864    fn test_observer_before_event_can_modify_attrs() {
865        struct TimestampInjector;
866        impl Observer for TimestampInjector {
867            fn before_insert(
868                &self,
869                _ctx: &HookContext,
870                attrs: &mut HashMap<String, Value>,
871            ) -> SubscriberResult<()> {
872                attrs.insert(
873                    "created_at".to_string(),
874                    Value::String("2026-07-19".to_string()),
875                );
876                Ok(())
877            }
878        }
879
880        let d = EventDispatcher::new();
881        d.add_observer(Box::new(TimestampInjector));
882
883        let ctx = HookContext::default();
884        let mut attrs = HashMap::new();
885        d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs)
886            .unwrap();
887
888        assert_eq!(
889            attrs.get("created_at"),
890            Some(&Value::String("2026-07-19".to_string()))
891        );
892    }
893
894    // ===== EventSubscriber 测试 =====
895
896    /// 只订阅 AfterInsert 的订阅者
897    struct InsertOnlySubscriber {
898        called: Arc<Mutex<u32>>,
899    }
900
901    impl EventSubscriber for InsertOnlySubscriber {
902        fn name(&self) -> &str {
903            "insert_only"
904        }
905
906        fn subscribed_events(&self) -> Vec<Event> {
907            vec![Event::AfterInsert]
908        }
909
910        fn on_event(
911            &self,
912            _event: Event,
913            _ctx: &HookContext,
914            _attrs: &HashMap<String, Value>,
915        ) -> SubscriberResult<()> {
916            *self.called.lock().unwrap() += 1;
917            Ok(())
918        }
919    }
920
921    #[test]
922    fn test_subscriber_only_called_for_subscribed_events() {
923        let called = Arc::new(Mutex::new(0u32));
924        let subscriber = InsertOnlySubscriber {
925            called: called.clone(),
926        };
927
928        let d = EventDispatcher::new();
929        d.subscribe(Box::new(subscriber));
930
931        let ctx = HookContext::default();
932        let attrs = HashMap::new();
933
934        // AfterInsert 应触发
935        d.dispatch(Event::AfterInsert, &ctx, &attrs);
936        // AfterUpdate 不应触发(未订阅)
937        d.dispatch(Event::AfterUpdate, &ctx, &attrs);
938        // AfterDelete 不应触发
939        d.dispatch(Event::AfterDelete, &ctx, &attrs);
940        // 再触发一次 AfterInsert
941        d.dispatch(Event::AfterInsert, &ctx, &attrs);
942
943        assert_eq!(*called.lock().unwrap(), 2);
944    }
945
946    #[test]
947    fn test_subscriber_veto_aborts_before_event() {
948        struct VetoSubscriber;
949        impl EventSubscriber for VetoSubscriber {
950            fn name(&self) -> &str {
951                "veto"
952            }
953            fn subscribed_events(&self) -> Vec<Event> {
954                vec![Event::BeforeInsert]
955            }
956            fn on_event(
957                &self,
958                _event: Event,
959                _ctx: &HookContext,
960                _attrs: &HashMap<String, Value>,
961            ) -> SubscriberResult<()> {
962                Err(SubscriberError::Vetoed {
963                    subscriber: "veto".to_string(),
964                    reason: "Business rule violation".to_string(),
965                })
966            }
967        }
968
969        let d = EventDispatcher::new();
970        d.subscribe(Box::new(VetoSubscriber));
971
972        let ctx = HookContext::default();
973        let mut attrs = HashMap::new();
974        let result = d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs);
975
976        assert!(matches!(result, Err(SubscriberError::Vetoed { .. })));
977    }
978
979    #[test]
980    fn test_subscriber_failed_does_not_abort_after_event() {
981        struct FailingSubscriber;
982        impl EventSubscriber for FailingSubscriber {
983            fn name(&self) -> &str {
984                "failing"
985            }
986            fn subscribed_events(&self) -> Vec<Event> {
987                vec![Event::AfterInsert]
988            }
989            fn on_event(
990                &self,
991                _event: Event,
992                _ctx: &HookContext,
993                _attrs: &HashMap<String, Value>,
994            ) -> SubscriberResult<()> {
995                Err(SubscriberError::Failed {
996                    subscriber: "failing".to_string(),
997                    reason: "Connection lost".to_string(),
998                })
999            }
1000        }
1001
1002        struct CountingSubscriber {
1003            called: Arc<Mutex<u32>>,
1004        }
1005        impl EventSubscriber for CountingSubscriber {
1006            fn name(&self) -> &str {
1007                "counting"
1008            }
1009            fn subscribed_events(&self) -> Vec<Event> {
1010                vec![Event::AfterInsert]
1011            }
1012            fn on_event(
1013                &self,
1014                _event: Event,
1015                _ctx: &HookContext,
1016                _attrs: &HashMap<String, Value>,
1017            ) -> SubscriberResult<()> {
1018                *self.called.lock().unwrap() += 1;
1019                Ok(())
1020            }
1021        }
1022
1023        let called = Arc::new(Mutex::new(0u32));
1024        let d = EventDispatcher::new();
1025        d.subscribe(Box::new(FailingSubscriber));
1026        d.subscribe(Box::new(CountingSubscriber {
1027            called: called.clone(),
1028        }));
1029
1030        let ctx = HookContext::default();
1031        let attrs = HashMap::new();
1032        d.dispatch(Event::AfterInsert, &ctx, &attrs);
1033
1034        // 即使 FailingSubscriber 失败,CountingSubscriber 仍应被调用
1035        assert_eq!(*called.lock().unwrap(), 1);
1036    }
1037
1038    // ===== AuditLogSubscriber 测试 =====
1039
1040    #[test]
1041    fn test_audit_log_subscriber() {
1042        let audit = AuditLogSubscriber::new();
1043        let audit_clone = audit.clone();
1044
1045        let d = EventDispatcher::new();
1046        d.subscribe(Box::new(audit_clone));
1047
1048        let ctx = HookContext {
1049            operator_id: Some(42),
1050            ..Default::default()
1051        };
1052        let mut attrs = HashMap::new();
1053        attrs.insert("name".to_string(), Value::String("alice".to_string()));
1054
1055        d.dispatch(Event::AfterInsert, &ctx, &attrs);
1056        d.dispatch(Event::AfterUpdate, &ctx, &attrs);
1057        d.dispatch(Event::AfterDelete, &ctx, &attrs);
1058        // AfterFind 不在订阅列表,不应记录
1059        d.dispatch_after_find(&ctx, &mut attrs).unwrap();
1060
1061        let logs = audit.logs().lock();
1062        assert_eq!(logs.len(), 3);
1063        assert!(logs[0].contains("event=after_insert"));
1064        assert!(logs[0].contains("operator=Some(42)"));
1065        assert!(logs[0].contains("field_count=1"));
1066    }
1067
1068    // ===== 多订阅者协同测试 =====
1069
1070    #[test]
1071    fn test_multiple_subscribers_and_observers() {
1072        let sub1_called = Arc::new(Mutex::new(0u32));
1073        let sub2_called = Arc::new(Mutex::new(0u32));
1074        let obs_called = Arc::new(Mutex::new(0u32));
1075
1076        struct Sub1(Arc<Mutex<u32>>);
1077        impl EventSubscriber for Sub1 {
1078            fn name(&self) -> &str {
1079                "sub1"
1080            }
1081            fn subscribed_events(&self) -> Vec<Event> {
1082                vec![Event::AfterInsert]
1083            }
1084            fn on_event(
1085                &self,
1086                _e: Event,
1087                _c: &HookContext,
1088                _a: &HashMap<String, Value>,
1089            ) -> SubscriberResult<()> {
1090                *self.0.lock().unwrap() += 1;
1091                Ok(())
1092            }
1093        }
1094
1095        struct Sub2(Arc<Mutex<u32>>);
1096        impl EventSubscriber for Sub2 {
1097            fn name(&self) -> &str {
1098                "sub2"
1099            }
1100            fn subscribed_events(&self) -> Vec<Event> {
1101                vec![Event::AfterInsert, Event::AfterUpdate]
1102            }
1103            fn on_event(
1104                &self,
1105                _e: Event,
1106                _c: &HookContext,
1107                _a: &HashMap<String, Value>,
1108            ) -> SubscriberResult<()> {
1109                *self.0.lock().unwrap() += 1;
1110                Ok(())
1111            }
1112        }
1113
1114        struct Obs(Arc<Mutex<u32>>);
1115        impl Observer for Obs {
1116            fn name(&self) -> &str {
1117                "obs"
1118            }
1119            fn after_insert(
1120                &self,
1121                _c: &HookContext,
1122                _a: &HashMap<String, Value>,
1123            ) -> SubscriberResult<()> {
1124                *self.0.lock().unwrap() += 1;
1125                Ok(())
1126            }
1127        }
1128
1129        let d = EventDispatcher::new();
1130        d.subscribe(Box::new(Sub1(sub1_called.clone())));
1131        d.subscribe(Box::new(Sub2(sub2_called.clone())));
1132        d.add_observer(Box::new(Obs(obs_called.clone())));
1133
1134        let ctx = HookContext::default();
1135        let attrs = HashMap::new();
1136
1137        d.dispatch(Event::AfterInsert, &ctx, &attrs);
1138
1139        assert_eq!(*sub1_called.lock().unwrap(), 1);
1140        assert_eq!(*sub2_called.lock().unwrap(), 1);
1141        assert_eq!(*obs_called.lock().unwrap(), 1);
1142    }
1143
1144    // ===== 错误收集测试 =====
1145
1146    #[test]
1147    fn test_drain_errors() {
1148        struct ErrSub;
1149        impl EventSubscriber for ErrSub {
1150            fn name(&self) -> &str {
1151                "err"
1152            }
1153            fn subscribed_events(&self) -> Vec<Event> {
1154                vec![Event::AfterInsert]
1155            }
1156            fn on_event(
1157                &self,
1158                _e: Event,
1159                _c: &HookContext,
1160                _a: &HashMap<String, Value>,
1161            ) -> SubscriberResult<()> {
1162                Err(SubscriberError::Failed {
1163                    subscriber: "err".to_string(),
1164                    reason: "test".to_string(),
1165                })
1166            }
1167        }
1168
1169        let d = EventDispatcher::new();
1170        d.subscribe(Box::new(ErrSub));
1171
1172        let ctx = HookContext::default();
1173        let attrs = HashMap::new();
1174        d.dispatch(Event::AfterInsert, &ctx, &attrs);
1175        d.dispatch(Event::AfterInsert, &ctx, &attrs);
1176
1177        let errors = d.drain_errors();
1178        assert_eq!(errors.len(), 2);
1179        assert!(matches!(errors[0], SubscriberError::Failed { .. }));
1180
1181        // drain 后内部应为空
1182        let errors = d.drain_errors();
1183        assert!(errors.is_empty());
1184    }
1185
1186    // ===== max_errors 限制测试(防内存无限增长) =====
1187
1188    #[test]
1189    fn test_max_errors_limits_buffer_size() {
1190        struct ErrSub;
1191        impl EventSubscriber for ErrSub {
1192            fn name(&self) -> &str {
1193                "err"
1194            }
1195            fn subscribed_events(&self) -> Vec<Event> {
1196                vec![Event::AfterInsert]
1197            }
1198            fn on_event(
1199                &self,
1200                _e: Event,
1201                _c: &HookContext,
1202                _a: &HashMap<String, Value>,
1203            ) -> SubscriberResult<()> {
1204                Err(SubscriberError::Failed {
1205                    subscriber: "err".to_string(),
1206                    reason: "test".to_string(),
1207                })
1208            }
1209        }
1210
1211        // 设置 max_errors = 3,触发 5 次错误,应只保留最新 3 个
1212        let d = EventDispatcher::new().with_max_errors(3);
1213        d.subscribe(Box::new(ErrSub));
1214
1215        let ctx = HookContext::default();
1216        let attrs = HashMap::new();
1217        for _ in 0..5 {
1218            d.dispatch(Event::AfterInsert, &ctx, &attrs);
1219        }
1220
1221        assert_eq!(d.error_count(), 3);
1222        let errors = d.drain_errors();
1223        assert_eq!(errors.len(), 3);
1224    }
1225
1226    #[test]
1227    fn test_max_errors_zero_means_unlimited() {
1228        struct ErrSub;
1229        impl EventSubscriber for ErrSub {
1230            fn name(&self) -> &str {
1231                "err"
1232            }
1233            fn subscribed_events(&self) -> Vec<Event> {
1234                vec![Event::AfterInsert]
1235            }
1236            fn on_event(
1237                &self,
1238                _e: Event,
1239                _c: &HookContext,
1240                _a: &HashMap<String, Value>,
1241            ) -> SubscriberResult<()> {
1242                Err(SubscriberError::Failed {
1243                    subscriber: "err".to_string(),
1244                    reason: "test".to_string(),
1245                })
1246            }
1247        }
1248
1249        let d = EventDispatcher::new().with_max_errors(0);
1250        d.subscribe(Box::new(ErrSub));
1251
1252        let ctx = HookContext::default();
1253        let attrs = HashMap::new();
1254        for _ in 0..10 {
1255            d.dispatch(Event::AfterInsert, &ctx, &attrs);
1256        }
1257
1258        assert_eq!(d.error_count(), 10);
1259    }
1260
1261    #[test]
1262    fn test_max_errors_fifo_eviction_order() {
1263        // 验证 FIFO 淘汰:保留的是最新错误
1264        struct CounterSub(Arc<Mutex<u32>>);
1265        impl EventSubscriber for CounterSub {
1266            fn name(&self) -> &str {
1267                "counter"
1268            }
1269            fn subscribed_events(&self) -> Vec<Event> {
1270                vec![Event::AfterInsert]
1271            }
1272            fn on_event(
1273                &self,
1274                _e: Event,
1275                _c: &HookContext,
1276                _a: &HashMap<String, Value>,
1277            ) -> SubscriberResult<()> {
1278                let mut n = self.0.lock().unwrap();
1279                *n += 1;
1280                Err(SubscriberError::Failed {
1281                    subscriber: "counter".to_string(),
1282                    reason: format!("call-{}", *n),
1283                })
1284            }
1285        }
1286
1287        let counter = Arc::new(Mutex::new(0u32));
1288        let d = EventDispatcher::new().with_max_errors(2);
1289        d.subscribe(Box::new(CounterSub(counter.clone())));
1290
1291        let ctx = HookContext::default();
1292        let attrs = HashMap::new();
1293        for _ in 0..4 {
1294            d.dispatch(Event::AfterInsert, &ctx, &attrs);
1295        }
1296
1297        let errors = d.drain_errors();
1298        assert_eq!(errors.len(), 2);
1299        // 应保留最新的两个(call-3, call-4)
1300        match &errors[0] {
1301            SubscriberError::Failed { reason, .. } => assert_eq!(reason, "call-3"),
1302            other => panic!("expected Failed, got {:?}", other),
1303        }
1304        match &errors[1] {
1305            SubscriberError::Failed { reason, .. } => assert_eq!(reason, "call-4"),
1306            other => panic!("expected Failed, got {:?}", other),
1307        }
1308    }
1309
1310    // ===== before 事件 Veto 测试 =====
1311
1312    #[test]
1313    fn test_veto_aborts_subsequent_observers() {
1314        let second_called = Arc::new(Mutex::new(0u32));
1315
1316        struct VetoObs;
1317        impl Observer for VetoObs {
1318            fn name(&self) -> &str {
1319                "veto"
1320            }
1321            fn before_insert(
1322                &self,
1323                _c: &HookContext,
1324                _a: &mut HashMap<String, Value>,
1325            ) -> SubscriberResult<()> {
1326                Err(SubscriberError::Vetoed {
1327                    subscriber: "veto".to_string(),
1328                    reason: "no".to_string(),
1329                })
1330            }
1331        }
1332
1333        struct CountingObs(Arc<Mutex<u32>>);
1334        impl Observer for CountingObs {
1335            fn name(&self) -> &str {
1336                "counting"
1337            }
1338            fn before_insert(
1339                &self,
1340                _c: &HookContext,
1341                _a: &mut HashMap<String, Value>,
1342            ) -> SubscriberResult<()> {
1343                *self.0.lock().unwrap() += 1;
1344                Ok(())
1345            }
1346        }
1347
1348        let d = EventDispatcher::new();
1349        d.add_observer(Box::new(VetoObs));
1350        d.add_observer(Box::new(CountingObs(second_called.clone())));
1351
1352        let ctx = HookContext::default();
1353        let mut attrs = HashMap::new();
1354        let result = d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs);
1355
1356        assert!(result.is_err());
1357        // 第二个 observer 不应被调用
1358        assert_eq!(*second_called.lock().unwrap(), 0);
1359    }
1360
1361    // ===== Display 测试 =====
1362
1363    #[test]
1364    fn test_error_display() {
1365        let e = SubscriberError::Failed {
1366            subscriber: "test".to_string(),
1367            reason: "boom".to_string(),
1368        };
1369        assert!(e.to_string().contains("test"));
1370        assert!(e.to_string().contains("boom"));
1371
1372        let e = SubscriberError::Vetoed {
1373            subscriber: "vetoer".to_string(),
1374            reason: "rejected".to_string(),
1375        };
1376        assert!(e.to_string().contains("vetoer"));
1377        assert!(e.to_string().contains("rejected"));
1378    }
1379}