1use crate::hooks::HookContext;
84use crate::Value;
85use parking_lot::{Mutex, RwLock};
86use std::collections::HashMap;
87use std::sync::Arc;
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
97pub enum Event {
98 BeforeInsert,
100 AfterInsert,
102 BeforeUpdate,
104 AfterUpdate,
106 BeforeDelete,
108 AfterDelete,
110 AfterFind,
112 BeforeRestore,
114 AfterRestore,
116}
117
118impl Event {
119 pub fn is_before(&self) -> bool {
121 matches!(
122 self,
123 Event::BeforeInsert | Event::BeforeUpdate | Event::BeforeDelete | Event::BeforeRestore
124 )
125 }
126
127 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 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 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#[derive(Debug)]
174pub enum SubscriberError {
175 Failed {
177 subscriber: String,
179 reason: String,
181 },
182 Vetoed {
186 subscriber: String,
188 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
208pub type SubscriberResult<T> = Result<T, SubscriberError>;
210
211pub trait Observer: Send + Sync {
225 fn name(&self) -> &str {
227 "anonymous_observer"
228 }
229
230 fn before_insert(
232 &self,
233 _ctx: &HookContext,
234 _attrs: &mut HashMap<String, Value>,
235 ) -> SubscriberResult<()> {
236 Ok(())
237 }
238
239 fn after_insert(
241 &self,
242 _ctx: &HookContext,
243 _attrs: &HashMap<String, Value>,
244 ) -> SubscriberResult<()> {
245 Ok(())
246 }
247
248 fn before_update(
250 &self,
251 _ctx: &HookContext,
252 _attrs: &mut HashMap<String, Value>,
253 ) -> SubscriberResult<()> {
254 Ok(())
255 }
256
257 fn after_update(
259 &self,
260 _ctx: &HookContext,
261 _attrs: &HashMap<String, Value>,
262 ) -> SubscriberResult<()> {
263 Ok(())
264 }
265
266 fn before_delete(
268 &self,
269 _ctx: &HookContext,
270 _attrs: &HashMap<String, Value>,
271 ) -> SubscriberResult<()> {
272 Ok(())
273 }
274
275 fn after_delete(
277 &self,
278 _ctx: &HookContext,
279 _attrs: &HashMap<String, Value>,
280 ) -> SubscriberResult<()> {
281 Ok(())
282 }
283
284 fn after_find(
286 &self,
287 _ctx: &HookContext,
288 _attrs: &mut HashMap<String, Value>,
289 ) -> SubscriberResult<()> {
290 Ok(())
291 }
292}
293
294pub trait EventSubscriber: Send + Sync {
303 fn name(&self) -> &str {
305 "anonymous_subscriber"
306 }
307
308 fn subscribed_events(&self) -> Vec<Event>;
312
313 fn on_event(
325 &self,
326 event: Event,
327 ctx: &HookContext,
328 attrs: &HashMap<String, Value>,
329 ) -> SubscriberResult<()>;
330}
331
332pub struct EventDispatcher {
360 observers: RwLock<Vec<Arc<dyn Observer>>>,
361 subscribers: RwLock<Vec<Arc<dyn EventSubscriber>>>,
362 errors: RwLock<Vec<SubscriberError>>,
364 max_errors: usize,
369}
370
371const DEFAULT_MAX_ERRORS: usize = 1024;
373
374impl EventDispatcher {
375 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 pub fn with_max_errors(mut self, max_errors: usize) -> Self {
390 self.max_errors = max_errors;
391 self
392 }
393
394 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 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 pub fn clear(&self) {
413 self.observers.write().clear();
414 self.subscribers.write().clear();
415 self.errors.write().clear();
416 }
417
418 pub fn observer_count(&self) -> usize {
420 self.observers.read().len()
421 }
422
423 pub fn subscriber_count(&self) -> usize {
425 self.subscribers.read().len()
426 }
427
428 pub fn drain_errors(&self) -> Vec<SubscriberError> {
430 std::mem::take(&mut *self.errors.write())
431 }
432
433 pub fn error_count(&self) -> usize {
435 self.errors.read().len()
436 }
437
438 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 errors.remove(0);
455 }
456 errors.push(e);
457 }
458 }
459
460 pub fn dispatch(&self, event: Event, ctx: &HookContext, attrs: &HashMap<String, Value>) {
470 let mut local_errors: Vec<SubscriberError> = Vec::new();
471
472 let observers_snapshot: Vec<Arc<dyn Observer>> = {
474 let observers = self.observers.read();
475 observers.clone()
476 };
477 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 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 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 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 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 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#[derive(Clone)]
657pub struct AuditLogSubscriber {
658 logs: std::sync::Arc<Mutex<Vec<String>>>,
659}
660
661impl AuditLogSubscriber {
662 pub fn new() -> Self {
664 Self {
665 logs: std::sync::Arc::new(Mutex::new(Vec::new())),
666 }
667 }
668
669 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#[cfg(test)]
712mod tests {
713 use super::*;
714 use std::sync::{Arc, Mutex};
715
716 #[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 #[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 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 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 d.dispatch(Event::AfterInsert, &ctx, &attrs);
936 d.dispatch(Event::AfterUpdate, &ctx, &attrs);
938 d.dispatch(Event::AfterDelete, &ctx, &attrs);
940 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 assert_eq!(*called.lock().unwrap(), 1);
1036 }
1037
1038 #[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 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 #[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 #[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 let errors = d.drain_errors();
1183 assert!(errors.is_empty());
1184 }
1185
1186 #[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 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 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 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 #[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 assert_eq!(*second_called.lock().unwrap(), 0);
1359 }
1360
1361 #[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}