1use std::collections::HashMap;
30use std::sync::{Arc, RwLock};
31
32use serde_json::Value;
33use tokio::task::JoinHandle;
34
35pub trait Listener: Send + Sync {
60 fn handle(&self, params: &Value) -> Result<Value, EventError>;
66}
67
68pub struct ClosureListener<F>
85where
86 F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
87{
88 closure: F,
89}
90
91impl<F> ClosureListener<F>
92where
93 F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
94{
95 pub fn new(closure: F) -> Self {
97 Self { closure }
98 }
99}
100
101impl<F> Listener for ClosureListener<F>
102where
103 F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
104{
105 fn handle(&self, params: &Value) -> Result<Value, EventError> {
106 (self.closure)(params)
107 }
108}
109
110pub trait Subscriber: Send + Sync {
139 fn subscribe(&self, dispatcher: &EventDispatcher);
141}
142
143pub trait Observer: Send + Sync {
171 fn events(&self) -> Vec<(&'static str, Arc<dyn Listener>)>;
175}
176
177#[derive(Debug)]
179pub enum EventError {
180 ListenerError(String),
182 EventNotFound(String),
184 InvalidParams(String),
186}
187
188impl std::fmt::Display for EventError {
189 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
190 match self {
191 EventError::ListenerError(s) => write!(f, "Event listener error: {}", s),
192 EventError::EventNotFound(s) => write!(f, "Event not found: {}", s),
193 EventError::InvalidParams(s) => write!(f, "Invalid event params: {}", s),
194 }
195 }
196}
197
198impl std::error::Error for EventError {}
199
200pub struct EventDispatcher {
219 listener: RwLock<HashMap<String, Vec<Arc<dyn Listener>>>>,
221
222 bind: RwLock<HashMap<String, String>>,
224}
225
226impl Default for EventDispatcher {
227 fn default() -> Self {
228 Self::new()
229 }
230}
231
232impl EventDispatcher {
233 pub fn new() -> Self {
235 Self {
236 listener: RwLock::new(HashMap::new()),
237 bind: RwLock::new(HashMap::new()),
238 }
239 }
240
241 pub fn listen_events(&self, events: Vec<(String, Vec<Arc<dyn Listener>>)>) -> &Self {
259 let mut listener_map = self.listener.write().expect("锁被毒化");
260 let bind_map = self.bind.read().expect("锁被毒化");
261
262 for (event, listeners) in events {
263 let event = bind_map.get(&event).cloned().unwrap_or(event);
265
266 let entry = listener_map.entry(event).or_default();
267 entry.extend(listeners);
268 }
269
270 self
271 }
272
273 pub fn listen(&self, event: &str, listener: Arc<dyn Listener>, first: bool) -> &Self {
293 let mut listener_map = self.listener.write().expect("锁被毒化");
294 let bind_map = self.bind.read().expect("锁被毒化");
295
296 let event = bind_map
298 .get(event)
299 .cloned()
300 .unwrap_or_else(|| event.to_string());
301
302 let entry = listener_map.entry(event).or_default();
303 if first {
304 entry.insert(0, listener);
306 } else {
307 entry.push(listener);
309 }
310
311 self
312 }
313
314 pub fn has_listener(&self, event: &str) -> bool {
316 let listener_map = self.listener.read().expect("锁被毒化");
317 let bind_map = self.bind.read().expect("锁被毒化");
318
319 let event = bind_map.get(event).map(|s| s.as_str()).unwrap_or(event);
321
322 listener_map.contains_key(event)
323 }
324
325 pub fn remove(&self, event: &str) {
327 let mut listener_map = self.listener.write().expect("锁被毒化");
328 let bind_map = self.bind.read().expect("锁被毒化");
329
330 let event = bind_map
332 .get(event)
333 .cloned()
334 .unwrap_or_else(|| event.to_string());
335
336 listener_map.remove(&event);
338 }
339
340 pub fn bind(&self, events: Vec<(String, String)>) -> &Self {
354 let mut bind_map = self.bind.write().expect("锁被毒化");
355 for (alias, real_event) in events {
356 bind_map.insert(alias, real_event);
357 }
358 self
359 }
360
361 pub fn subscribe(&self, subscriber: Arc<dyn Subscriber>) -> &Self {
382 subscriber.subscribe(self);
384 self
385 }
386
387 pub fn observe(&self, observer: Arc<dyn Observer>, prefix: &str) -> &Self {
408 for (event, listener) in observer.events() {
409 let full_event = if prefix.is_empty() {
410 event.to_string()
411 } else {
412 format!("{}{}", prefix, event)
413 };
414 self.listen(&full_event, listener, false);
416 }
417 self
418 }
419
420 fn collect_listeners(&self, event: &str) -> Vec<Arc<dyn Listener>> {
424 let bind_map = self.bind.read().expect("锁被毒化");
425
426 let event = bind_map
428 .get(event)
429 .cloned()
430 .unwrap_or_else(|| event.to_string());
431
432 drop(bind_map);
433
434 let listener_map = self.listener.read().expect("锁被毒化");
435
436 let mut listeners: Vec<Arc<dyn Listener>> =
438 listener_map.get(&event).cloned().unwrap_or_default();
439
440 if let Some(dot_pos) = event.find('.') {
442 let prefix = &event[..dot_pos];
443 let wildcard = format!("{}.*", prefix);
444 if let Some(wildcard_listeners) = listener_map.get(&wildcard) {
445 listeners.extend(wildcard_listeners.clone());
447 }
448 }
449
450 drop(listener_map);
451
452 let mut seen: Vec<Arc<dyn Listener>> = Vec::new();
455 listeners.retain(|l| {
456 if seen.iter().any(|s| Arc::ptr_eq(s, l)) {
457 false
458 } else {
459 seen.push(l.clone());
460 true
461 }
462 });
463
464 listeners
465 }
466
467 pub fn trigger(
491 &self,
492 event: &str,
493 params: &Value,
494 once: bool,
495 ) -> Result<Vec<Value>, EventError> {
496 let listeners = self.collect_listeners(event);
497
498 let mut results: Vec<Value> = Vec::new();
499
500 for listener in &listeners {
501 let result = listener.handle(params)?;
503 results.push(result.clone());
504
505 if result == Value::Bool(false) {
508 break;
509 }
510 if once && !result.is_null() {
512 break;
513 }
514 }
515
516 Ok(results)
517 }
518
519 pub fn trigger_spawn(
544 &self,
545 event: &str,
546 params: &Value,
547 ) -> Vec<JoinHandle<Result<Value, EventError>>> {
548 let listeners = self.collect_listeners(event);
549 let params_owned = params.clone();
550
551 listeners
552 .into_iter()
553 .map(|listener| {
554 let params = params_owned.clone();
555 tokio::spawn(async move { listener.handle(¶ms) })
556 })
557 .collect()
558 }
559
560 pub async fn trigger_async(
580 &self,
581 event: &str,
582 params: &Value,
583 ) -> Vec<Result<Value, EventError>> {
584 let handles = self.trigger_spawn(event, params);
585 let mut results = Vec::with_capacity(handles.len());
586
587 for handle in handles {
588 match handle.await {
589 Ok(result) => results.push(result),
590 Err(join_err) => results.push(Err(EventError::ListenerError(format!(
591 "Task panicked: {}",
592 join_err
593 )))),
594 }
595 }
596
597 results
598 }
599
600 pub fn until(&self, event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
604 self.trigger(event, params, true)
605 }
606
607 pub fn listener_count(&self, event: &str) -> usize {
609 let listener_map = self.listener.read().expect("锁被毒化");
610 let bind_map = self.bind.read().expect("锁被毒化");
611
612 let event = bind_map.get(event).map(|s| s.as_str()).unwrap_or(event);
613
614 listener_map.get(event).map(|v| v.len()).unwrap_or(0)
615 }
616}
617
618pub mod facade {
632 use super::*;
633 use std::sync::OnceLock;
634
635 static GLOBAL_DISPATCHER: OnceLock<EventDispatcher> = OnceLock::new();
637
638 pub fn dispatcher() -> &'static EventDispatcher {
640 GLOBAL_DISPATCHER.get_or_init(EventDispatcher::new)
641 }
642
643 pub fn listen(event: &str, listener: Arc<dyn Listener>, first: bool) {
645 dispatcher().listen(event, listener, first);
646 }
647
648 pub fn listen_events(events: Vec<(String, Vec<Arc<dyn Listener>>)>) {
650 dispatcher().listen_events(events);
651 }
652
653 pub fn has_listener(event: &str) -> bool {
655 dispatcher().has_listener(event)
656 }
657
658 pub fn remove(event: &str) {
660 dispatcher().remove(event);
661 }
662
663 pub fn bind(events: Vec<(String, String)>) {
665 dispatcher().bind(events);
666 }
667
668 pub fn subscribe(subscriber: Arc<dyn Subscriber>) {
670 dispatcher().subscribe(subscriber);
671 }
672
673 pub fn observe(observer: Arc<dyn Observer>, prefix: &str) {
675 dispatcher().observe(observer, prefix);
676 }
677
678 pub fn trigger(event: &str, params: &Value, once: bool) -> Result<Vec<Value>, EventError> {
680 dispatcher().trigger(event, params, once)
681 }
682
683 pub fn until(event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
685 dispatcher().until(event, params)
686 }
687
688 pub fn trigger_spawn(
690 event: &str,
691 params: &Value,
692 ) -> Vec<JoinHandle<Result<Value, EventError>>> {
693 dispatcher().trigger_spawn(event, params)
694 }
695
696 pub async fn trigger_async(event: &str, params: &Value) -> Vec<Result<Value, EventError>> {
698 dispatcher().trigger_async(event, params).await
699 }
700
701 #[cfg(test)]
706 pub fn _reset_for_test() {
707 }
710}
711
712pub fn event_trigger(event: &str, params: &Value) -> Vec<Value> {
729 facade::trigger(event, params, false).unwrap_or_default()
730}
731
732pub async fn event_trigger_async(event: &str, params: &Value) -> Vec<Result<Value, EventError>> {
746 facade::trigger_async(event, params).await
747}
748
749#[cfg(test)]
750mod tests {
751 use super::*;
752 use serde_json::json;
753 use std::sync::atomic::{AtomicUsize, Ordering};
754
755 #[test]
760 fn test_event_listen_and_has_listener() {
761 let dispatcher = EventDispatcher::new();
762 assert!(!dispatcher.has_listener("UserLogin"));
763
764 dispatcher.listen(
765 "UserLogin",
766 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
767 false,
768 );
769 assert!(dispatcher.has_listener("UserLogin"));
770 }
771
772 #[test]
773 fn test_event_remove_listener() {
774 let dispatcher = EventDispatcher::new();
775 dispatcher.listen(
776 "UserLogin",
777 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
778 false,
779 );
780 assert!(dispatcher.has_listener("UserLogin"));
781
782 dispatcher.remove("UserLogin");
783 assert!(!dispatcher.has_listener("UserLogin"));
784 }
785
786 #[test]
787 fn test_event_listen_first_priority() {
788 let dispatcher = EventDispatcher::new();
790 let call_order = Arc::new(AtomicUsize::new(0));
791
792 let order1 = call_order.clone();
793 dispatcher.listen(
794 "Test",
795 Arc::new(ClosureListener::new(move |_| {
796 order1.store(1, Ordering::SeqCst);
797 Ok(Value::Null)
798 })),
799 false,
800 );
801
802 let order2 = call_order.clone();
803 dispatcher.listen(
804 "Test",
805 Arc::new(ClosureListener::new(move |_| {
806 order2.store(2, Ordering::SeqCst);
807 Ok(Value::Null)
808 })),
809 true, );
811
812 dispatcher.trigger("Test", &Value::Null, false).unwrap();
813 assert_eq!(call_order.load(Ordering::SeqCst), 1);
815 }
816
817 #[test]
818 fn test_event_listen_events_batch() {
819 let dispatcher = EventDispatcher::new();
821 dispatcher.listen_events(vec![
822 (
823 "UserLogin".to_string(),
824 vec![Arc::new(ClosureListener::new(|_| Ok(Value::Null)))],
825 ),
826 (
827 "UserLogout".to_string(),
828 vec![
829 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
830 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
831 ],
832 ),
833 ]);
834 assert!(dispatcher.has_listener("UserLogin"));
835 assert_eq!(dispatcher.listener_count("UserLogout"), 2);
836 }
837
838 #[test]
839 fn test_event_bind_alias() {
840 let dispatcher = EventDispatcher::new();
842 dispatcher.bind(vec![(
843 "AppInit".to_string(),
844 "app\\event\\AppInit".to_string(),
845 )]);
846
847 dispatcher.listen(
848 "AppInit",
849 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
850 false,
851 );
852 assert!(!dispatcher.has_listener("AppInit_alias_check"));
854 assert!(dispatcher.has_listener("app\\event\\AppInit"));
855 }
856
857 #[test]
862 fn test_event_trigger_returns_all_results() {
863 let dispatcher = EventDispatcher::new();
865 dispatcher.listen(
866 "Test",
867 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
868 false,
869 );
870 dispatcher.listen(
871 "Test",
872 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
873 false,
874 );
875
876 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
877 assert_eq!(results, vec![json!(1), json!(2)]);
878 }
879
880 #[test]
881 fn test_event_trigger_once_stops_at_non_null() {
882 let dispatcher = EventDispatcher::new();
884 dispatcher.listen(
885 "Test",
886 Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false,
888 );
889 dispatcher.listen(
890 "Test",
891 Arc::new(ClosureListener::new(|_| Ok(json!("stop")))), false,
893 );
894 let executed = Arc::new(AtomicUsize::new(0));
895 let exec_clone = executed.clone();
896 dispatcher.listen(
897 "Test",
898 Arc::new(ClosureListener::new(move |_| {
899 exec_clone.fetch_add(1, Ordering::SeqCst);
900 Ok(Value::Null)
901 })),
902 false,
903 );
904
905 let results = dispatcher.trigger("Test", &Value::Null, true).unwrap();
906 assert_eq!(results.len(), 2);
908 assert_eq!(results[0], Value::Null);
909 assert_eq!(results[1], json!("stop"));
910 assert_eq!(executed.load(Ordering::SeqCst), 0); }
912
913 #[test]
914 fn test_event_trigger_false_stops_execution() {
915 let dispatcher = EventDispatcher::new();
917 let executed = Arc::new(AtomicUsize::new(0));
918
919 dispatcher.listen(
920 "Test",
921 Arc::new(ClosureListener::new(|_| Ok(Value::Bool(false)))), false,
923 );
924 let exec_clone = executed.clone();
925 dispatcher.listen(
926 "Test",
927 Arc::new(ClosureListener::new(move |_| {
928 exec_clone.fetch_add(1, Ordering::SeqCst);
929 Ok(Value::Null)
930 })),
931 false,
932 );
933
934 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
935 assert_eq!(results.len(), 1);
936 assert_eq!(results[0], Value::Bool(false));
937 assert_eq!(executed.load(Ordering::SeqCst), 0); }
939
940 #[test]
941 fn test_event_trigger_empty_event() {
942 let dispatcher = EventDispatcher::new();
944 let results = dispatcher
945 .trigger("Nonexistent", &Value::Null, false)
946 .unwrap();
947 assert!(results.is_empty());
948 }
949
950 #[test]
951 fn test_event_trigger_passes_params() {
952 let dispatcher = EventDispatcher::new();
954 let received = Arc::new(std::sync::Mutex::new(Value::Null));
955 let recv_clone = received.clone();
956
957 dispatcher.listen(
958 "Test",
959 Arc::new(ClosureListener::new(move |params| {
960 *recv_clone.lock().unwrap() = params.clone();
961 Ok(Value::Null)
962 })),
963 false,
964 );
965
966 let params = json!({"user_id": 123, "action": "login"});
967 dispatcher.trigger("Test", ¶ms, false).unwrap();
968
969 assert_eq!(*received.lock().unwrap(), params);
970 }
971
972 #[test]
977 fn test_event_dot_wildcard() {
978 let dispatcher = EventDispatcher::new();
980 let executed = Arc::new(AtomicUsize::new(0));
981
982 dispatcher.listen(
984 "User.login",
985 Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
986 false,
987 );
988 let exec_clone = executed.clone();
989 dispatcher.listen(
990 "User.*",
991 Arc::new(ClosureListener::new(move |_| {
992 exec_clone.fetch_add(1, Ordering::SeqCst);
993 Ok(json!("wildcard"))
994 })),
995 false,
996 );
997
998 let results = dispatcher
1000 .trigger("User.login", &Value::Null, false)
1001 .unwrap();
1002 assert_eq!(results.len(), 2);
1003 assert_eq!(results[0], json!("specific"));
1004 assert_eq!(results[1], json!("wildcard"));
1005 assert_eq!(executed.load(Ordering::SeqCst), 1);
1006 }
1007
1008 #[test]
1009 fn test_event_dot_wildcard_no_wildcard_listener() {
1010 let dispatcher = EventDispatcher::new();
1012 dispatcher.listen(
1013 "User.login",
1014 Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
1015 false,
1016 );
1017
1018 let results = dispatcher
1019 .trigger("User.login", &Value::Null, false)
1020 .unwrap();
1021 assert_eq!(results.len(), 1);
1022 assert_eq!(results[0], json!("specific"));
1023 }
1024
1025 #[test]
1030 fn test_event_dedup_same_listener_instance() {
1031 let dispatcher = EventDispatcher::new();
1033 let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
1034
1035 dispatcher.listen("Test", listener.clone(), false);
1037 dispatcher.listen("Test", listener.clone(), false);
1038
1039 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1040 assert_eq!(results.len(), 1);
1042 assert_eq!(results[0], json!(1));
1043 }
1044
1045 #[test]
1046 fn test_event_no_dedup_different_listeners() {
1047 let dispatcher = EventDispatcher::new();
1049 dispatcher.listen(
1050 "Test",
1051 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1052 false,
1053 );
1054 dispatcher.listen(
1055 "Test",
1056 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1057 false,
1058 );
1059
1060 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1061 assert_eq!(results.len(), 2);
1062 }
1063
1064 struct TestSubscriber {
1069 login_count: Arc<AtomicUsize>,
1070 }
1071
1072 impl Subscriber for TestSubscriber {
1073 fn subscribe(&self, dispatcher: &EventDispatcher) {
1074 let count = self.login_count.clone();
1075 dispatcher.listen(
1076 "UserLogin",
1077 Arc::new(ClosureListener::new(move |_| {
1078 count.fetch_add(1, Ordering::SeqCst);
1079 Ok(Value::Null)
1080 })),
1081 false,
1082 );
1083
1084 let count2 = self.login_count.clone();
1085 dispatcher.listen(
1086 "UserLogout",
1087 Arc::new(ClosureListener::new(move |_| {
1088 count2.fetch_add(10, Ordering::SeqCst);
1089 Ok(Value::Null)
1090 })),
1091 false,
1092 );
1093 }
1094 }
1095
1096 #[test]
1097 fn test_event_subscriber_registers_multiple_listeners() {
1098 let dispatcher = EventDispatcher::new();
1100 let login_count = Arc::new(AtomicUsize::new(0));
1101
1102 let subscriber = Arc::new(TestSubscriber {
1103 login_count: login_count.clone(),
1104 });
1105 dispatcher.subscribe(subscriber);
1106
1107 assert!(dispatcher.has_listener("UserLogin"));
1108 assert!(dispatcher.has_listener("UserLogout"));
1109
1110 dispatcher
1111 .trigger("UserLogin", &Value::Null, false)
1112 .unwrap();
1113 assert_eq!(login_count.load(Ordering::SeqCst), 1);
1114
1115 dispatcher
1116 .trigger("UserLogout", &Value::Null, false)
1117 .unwrap();
1118 assert_eq!(login_count.load(Ordering::SeqCst), 11); }
1120
1121 struct TestObserver {
1126 counter: Arc<AtomicUsize>,
1127 }
1128
1129 impl Observer for TestObserver {
1130 fn events(&self) -> Vec<(&'static str, Arc<dyn Listener>)> {
1131 let c1 = self.counter.clone();
1132 let c2 = self.counter.clone();
1133 vec![
1134 (
1135 "Login",
1136 Arc::new(ClosureListener::new(move |_| {
1137 c1.fetch_add(1, Ordering::SeqCst);
1138 Ok(Value::Null)
1139 })),
1140 ),
1141 (
1142 "Logout",
1143 Arc::new(ClosureListener::new(move |_| {
1144 c2.fetch_add(100, Ordering::SeqCst);
1145 Ok(Value::Null)
1146 })),
1147 ),
1148 ]
1149 }
1150 }
1151
1152 #[test]
1153 fn test_event_observer_auto_registers() {
1154 let dispatcher = EventDispatcher::new();
1156 let counter = Arc::new(AtomicUsize::new(0));
1157
1158 let observer = Arc::new(TestObserver {
1159 counter: counter.clone(),
1160 });
1161 dispatcher.observe(observer, "");
1162
1163 assert!(dispatcher.has_listener("Login"));
1164 assert!(dispatcher.has_listener("Logout"));
1165
1166 dispatcher.trigger("Login", &Value::Null, false).unwrap();
1167 assert_eq!(counter.load(Ordering::SeqCst), 1);
1168
1169 dispatcher.trigger("Logout", &Value::Null, false).unwrap();
1170 assert_eq!(counter.load(Ordering::SeqCst), 101);
1171 }
1172
1173 #[test]
1174 fn test_event_observer_with_prefix() {
1175 let dispatcher = EventDispatcher::new();
1177 let counter = Arc::new(AtomicUsize::new(0));
1178
1179 let observer = Arc::new(TestObserver {
1180 counter: counter.clone(),
1181 });
1182 dispatcher.observe(observer, "User");
1183
1184 assert!(dispatcher.has_listener("UserLogin"));
1185 assert!(dispatcher.has_listener("UserLogout"));
1186
1187 dispatcher
1188 .trigger("UserLogin", &Value::Null, false)
1189 .unwrap();
1190 assert_eq!(counter.load(Ordering::SeqCst), 1);
1191 }
1192
1193 #[test]
1198 fn test_event_until_returns_first_non_null() {
1199 let dispatcher = EventDispatcher::new();
1201 dispatcher.listen(
1202 "Test",
1203 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
1204 false,
1205 );
1206 dispatcher.listen(
1207 "Test",
1208 Arc::new(ClosureListener::new(|_| Ok(json!("first_valid")))),
1209 false,
1210 );
1211
1212 let results = dispatcher.until("Test", &Value::Null).unwrap();
1213 assert_eq!(results.len(), 2);
1215 assert_eq!(results[1], json!("first_valid"));
1216 }
1217
1218 #[test]
1223 fn test_closure_listener_executes_closure() {
1224 let listener = ClosureListener::new(|params| {
1225 assert_eq!(params, &json!({"key": "value"}));
1226 Ok(json!("result"))
1227 });
1228
1229 let result = listener.handle(&json!({"key": "value"})).unwrap();
1230 assert_eq!(result, json!("result"));
1231 }
1232
1233 #[test]
1234 fn test_closure_listener_returns_null() {
1235 let listener = ClosureListener::new(|_| Ok(Value::Null));
1236 let result = listener.handle(&Value::Null).unwrap();
1237 assert!(result.is_null());
1238 }
1239
1240 struct CustomListener {
1245 id: i32,
1246 }
1247
1248 impl Listener for CustomListener {
1249 fn handle(&self, _params: &Value) -> Result<Value, EventError> {
1250 Ok(json!({"listener_id": self.id}))
1251 }
1252 }
1253
1254 #[test]
1255 fn test_custom_listener_trait_impl() {
1256 let dispatcher = EventDispatcher::new();
1257 dispatcher.listen("Test", Arc::new(CustomListener { id: 42 }), false);
1258
1259 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1260 assert_eq!(results, vec![json!({"listener_id": 42})]);
1261 }
1262
1263 #[test]
1264 fn test_listener_error_propagates() {
1265 let dispatcher = EventDispatcher::new();
1267 let executed = Arc::new(AtomicUsize::new(0));
1268
1269 dispatcher.listen(
1270 "Test",
1271 Arc::new(ClosureListener::new(|_| {
1272 Err(EventError::ListenerError("test error".to_string()))
1273 })),
1274 false,
1275 );
1276 let exec_clone = executed.clone();
1277 dispatcher.listen(
1278 "Test",
1279 Arc::new(ClosureListener::new(move |_| {
1280 exec_clone.fetch_add(1, Ordering::SeqCst);
1281 Ok(Value::Null)
1282 })),
1283 false,
1284 );
1285
1286 let result = dispatcher.trigger("Test", &Value::Null, false);
1287 assert!(result.is_err());
1288 assert_eq!(executed.load(Ordering::SeqCst), 0); }
1290
1291 #[test]
1296 fn test_r5_php_event_listen_then_trigger() {
1297 let dispatcher = EventDispatcher::new();
1299 let received = Arc::new(std::sync::Mutex::new(Value::Null));
1300 let recv_clone = received.clone();
1301
1302 dispatcher.listen(
1303 "UserLogin",
1304 Arc::new(ClosureListener::new(move |params| {
1305 *recv_clone.lock().unwrap() = params.clone();
1306 Ok(Value::Null)
1307 })),
1308 false,
1309 );
1310
1311 let params = json!({"user_id": 123, "username": "alice"});
1312 dispatcher.trigger("UserLogin", ¶ms, false).unwrap();
1313
1314 assert_eq!(*received.lock().unwrap(), params);
1315 }
1316
1317 #[test]
1318 fn test_r5_php_event_bind_alias_resolution() {
1319 let dispatcher = EventDispatcher::new();
1321 dispatcher.bind(vec![(
1322 "AppInit".to_string(),
1323 "think\\event\\AppInit".to_string(),
1324 )]);
1325
1326 dispatcher.listen(
1327 "AppInit",
1328 Arc::new(ClosureListener::new(|_| Ok(json!("init_called")))),
1329 false,
1330 );
1331
1332 assert!(dispatcher.has_listener("think\\event\\AppInit"));
1334 assert_eq!(dispatcher.listener_count("think\\event\\AppInit"), 1);
1335
1336 let results = dispatcher.trigger("AppInit", &Value::Null, false).unwrap();
1337 assert_eq!(results, vec![json!("init_called")]);
1338 }
1339
1340 #[test]
1341 fn test_r5_php_event_first_array_unshift() {
1342 let dispatcher = EventDispatcher::new();
1344 let order = Arc::new(std::sync::Mutex::new(Vec::new()));
1345
1346 let o1 = order.clone();
1347 dispatcher.listen(
1348 "Test",
1349 Arc::new(ClosureListener::new(move |_| {
1350 o1.lock().unwrap().push(1);
1351 Ok(Value::Null)
1352 })),
1353 false,
1354 );
1355
1356 let o2 = order.clone();
1357 dispatcher.listen(
1358 "Test",
1359 Arc::new(ClosureListener::new(move |_| {
1360 o2.lock().unwrap().push(2);
1361 Ok(Value::Null)
1362 })),
1363 true, );
1365
1366 dispatcher.trigger("Test", &Value::Null, false).unwrap();
1367 assert_eq!(*order.lock().unwrap(), vec![2, 1]);
1369 }
1370
1371 #[test]
1372 fn test_r5_php_event_trigger_returns_array() {
1373 let dispatcher = EventDispatcher::new();
1375 dispatcher.listen(
1376 "Test",
1377 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1378 false,
1379 );
1380 dispatcher.listen(
1381 "Test",
1382 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1383 false,
1384 );
1385 dispatcher.listen(
1386 "Test",
1387 Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
1388 false,
1389 );
1390
1391 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1392 assert_eq!(results, vec![json!(1), json!(2), json!(3)]);
1393 }
1394
1395 #[test]
1396 fn test_r5_php_event_until_returns_last_non_null() {
1397 let dispatcher = EventDispatcher::new();
1399 dispatcher.listen(
1400 "Test",
1401 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
1402 false,
1403 );
1404 dispatcher.listen(
1405 "Test",
1406 Arc::new(ClosureListener::new(|_| Ok(json!("first")))),
1407 false,
1408 );
1409
1410 let results = dispatcher.until("Test", &Value::Null).unwrap();
1411 assert_eq!(results.len(), 2);
1413 }
1414
1415 #[test]
1416 fn test_r5_php_event_false_stops() {
1417 let dispatcher = EventDispatcher::new();
1419 let executed = Arc::new(AtomicUsize::new(0));
1420
1421 dispatcher.listen(
1422 "Test",
1423 Arc::new(ClosureListener::new(|_| Ok(Value::Bool(false)))),
1424 false,
1425 );
1426 let exec_clone = executed.clone();
1427 dispatcher.listen(
1428 "Test",
1429 Arc::new(ClosureListener::new(move |_| {
1430 exec_clone.fetch_add(1, Ordering::SeqCst);
1431 Ok(Value::Null)
1432 })),
1433 false,
1434 );
1435
1436 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1437 assert_eq!(results.len(), 1);
1438 assert_eq!(results[0], Value::Bool(false));
1439 }
1440
1441 #[test]
1442 fn test_r5_php_event_dot_wildcard_merge() {
1443 let dispatcher = EventDispatcher::new();
1445 dispatcher.listen(
1446 "User.login",
1447 Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
1448 false,
1449 );
1450 dispatcher.listen(
1451 "User.*",
1452 Arc::new(ClosureListener::new(|_| Ok(json!("wildcard")))),
1453 false,
1454 );
1455
1456 let results = dispatcher
1457 .trigger("User.login", &Value::Null, false)
1458 .unwrap();
1459 assert_eq!(results, vec![json!("specific"), json!("wildcard")]);
1460 }
1461
1462 #[test]
1463 fn test_r5_php_event_array_unique_dedup() {
1464 let dispatcher = EventDispatcher::new();
1466 let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
1467
1468 dispatcher.listen("Test", listener.clone(), false);
1469 dispatcher.listen("Test", listener.clone(), false);
1470 dispatcher.listen("Test", listener.clone(), false);
1471
1472 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1473 assert_eq!(results.len(), 1);
1475 }
1476
1477 #[test]
1478 fn test_r5_php_event_remove_clears_listeners() {
1479 let dispatcher = EventDispatcher::new();
1481 dispatcher.listen(
1482 "Test",
1483 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1484 false,
1485 );
1486 assert!(dispatcher.has_listener("Test"));
1487
1488 dispatcher.remove("Test");
1489 assert!(!dispatcher.has_listener("Test"));
1490
1491 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1492 assert!(results.is_empty());
1493 }
1494
1495 #[test]
1496 fn test_r5_php_event_has_listener_with_bind() {
1497 let dispatcher = EventDispatcher::new();
1499 dispatcher.bind(vec![(
1500 "AppInit".to_string(),
1501 "app\\event\\AppInit".to_string(),
1502 )]);
1503 dispatcher.listen(
1504 "AppInit",
1505 Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
1506 false,
1507 );
1508
1509 assert!(dispatcher.has_listener("AppInit"));
1511 assert!(dispatcher.has_listener("app\\event\\AppInit"));
1512 }
1513
1514 #[test]
1515 fn test_r5_php_event_listen_events_batch_merge() {
1516 let dispatcher = EventDispatcher::new();
1518 dispatcher.listen(
1519 "Test",
1520 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1521 false,
1522 );
1523 dispatcher.listen_events(vec![(
1524 "Test".to_string(),
1525 vec![
1526 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1527 Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
1528 ],
1529 )]);
1530
1531 let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
1532 assert_eq!(results, vec![json!(1), json!(2), json!(3)]);
1534 }
1535
1536 #[tokio::test]
1541 async fn test_trigger_spawn_returns_join_handles() {
1542 let dispatcher = EventDispatcher::new();
1544 dispatcher.listen(
1545 "Test",
1546 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1547 false,
1548 );
1549 dispatcher.listen(
1550 "Test",
1551 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1552 false,
1553 );
1554
1555 let handles = dispatcher.trigger_spawn("Test", &Value::Null);
1556 assert_eq!(handles.len(), 2);
1557
1558 for handle in handles {
1560 let result = handle.await.unwrap().unwrap();
1561 assert!(result == json!(1) || result == json!(2));
1562 }
1563 }
1564
1565 #[tokio::test]
1566 async fn test_trigger_spawn_empty_event() {
1567 let dispatcher = EventDispatcher::new();
1569 let handles = dispatcher.trigger_spawn("Nonexistent", &Value::Null);
1570 assert!(handles.is_empty());
1571 }
1572
1573 #[tokio::test]
1574 async fn test_trigger_spawn_passes_params() {
1575 let dispatcher = EventDispatcher::new();
1577 let received = Arc::new(std::sync::Mutex::new(Value::Null));
1578 let recv_clone = received.clone();
1579
1580 dispatcher.listen(
1581 "Test",
1582 Arc::new(ClosureListener::new(move |params| {
1583 *recv_clone.lock().unwrap() = params.clone();
1584 Ok(Value::Null)
1585 })),
1586 false,
1587 );
1588
1589 let params = json!({"user_id": 123, "action": "login"});
1590 let handles = dispatcher.trigger_spawn("Test", ¶ms);
1591 for handle in handles {
1592 let _ = handle.await;
1593 }
1594
1595 assert_eq!(*received.lock().unwrap(), params);
1596 }
1597
1598 #[tokio::test]
1599 async fn test_trigger_spawn_all_listeners_execute() {
1600 let dispatcher = EventDispatcher::new();
1602 let counter = Arc::new(AtomicUsize::new(0));
1603
1604 for _ in 0..5 {
1605 let c = counter.clone();
1606 dispatcher.listen(
1607 "Test",
1608 Arc::new(ClosureListener::new(move |_| {
1609 c.fetch_add(1, Ordering::SeqCst);
1610 Ok(Value::Null)
1611 })),
1612 false,
1613 );
1614 }
1615
1616 let handles = dispatcher.trigger_spawn("Test", &Value::Null);
1617 for handle in handles {
1618 let _ = handle.await;
1619 }
1620
1621 assert_eq!(counter.load(Ordering::SeqCst), 5);
1622 }
1623
1624 #[tokio::test]
1629 async fn test_trigger_async_awaits_all() {
1630 let dispatcher = EventDispatcher::new();
1632 dispatcher.listen(
1633 "Test",
1634 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1635 false,
1636 );
1637 dispatcher.listen(
1638 "Test",
1639 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1640 false,
1641 );
1642 dispatcher.listen(
1643 "Test",
1644 Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
1645 false,
1646 );
1647
1648 let results = dispatcher.trigger_async("Test", &Value::Null).await;
1649 assert_eq!(results.len(), 3);
1650 assert!(results.iter().all(|r| r.is_ok()));
1651 }
1652
1653 #[tokio::test]
1654 async fn test_trigger_async_collects_results_in_order() {
1655 let dispatcher = EventDispatcher::new();
1657 dispatcher.listen(
1658 "Test",
1659 Arc::new(ClosureListener::new(|_| Ok(json!("first")))),
1660 false,
1661 );
1662 dispatcher.listen(
1663 "Test",
1664 Arc::new(ClosureListener::new(|_| Ok(json!("second")))),
1665 false,
1666 );
1667 dispatcher.listen(
1668 "Test",
1669 Arc::new(ClosureListener::new(|_| Ok(json!("third")))),
1670 false,
1671 );
1672
1673 let results = dispatcher.trigger_async("Test", &Value::Null).await;
1674 assert_eq!(results[0].as_ref().unwrap(), &json!("first"));
1675 assert_eq!(results[1].as_ref().unwrap(), &json!("second"));
1676 assert_eq!(results[2].as_ref().unwrap(), &json!("third"));
1677 }
1678
1679 #[tokio::test]
1680 async fn test_trigger_async_handles_errors() {
1681 let dispatcher = EventDispatcher::new();
1683 dispatcher.listen(
1684 "Test",
1685 Arc::new(ClosureListener::new(|_| Ok(json!("ok")))),
1686 false,
1687 );
1688 dispatcher.listen(
1689 "Test",
1690 Arc::new(ClosureListener::new(|_| {
1691 Err(EventError::ListenerError("test error".to_string()))
1692 })),
1693 false,
1694 );
1695 dispatcher.listen(
1696 "Test",
1697 Arc::new(ClosureListener::new(|_| Ok(json!("ok2")))),
1698 false,
1699 );
1700
1701 let results = dispatcher.trigger_async("Test", &Value::Null).await;
1702 assert_eq!(results.len(), 3);
1703 assert!(results[0].is_ok());
1704 assert!(results[1].is_err());
1705 assert!(results[2].is_ok());
1706 }
1707
1708 #[tokio::test]
1709 async fn test_trigger_async_handles_panic() {
1710 let dispatcher = EventDispatcher::new();
1712 dispatcher.listen(
1713 "Test",
1714 Arc::new(ClosureListener::new(|_| Ok(json!("ok")))),
1715 false,
1716 );
1717 dispatcher.listen(
1718 "Test",
1719 Arc::new(ClosureListener::new(|_| panic!("test panic"))),
1720 false,
1721 );
1722
1723 let results = dispatcher.trigger_async("Test", &Value::Null).await;
1724 assert_eq!(results.len(), 2);
1725 assert!(results[0].is_ok());
1726 assert!(results[1].is_err()); }
1728
1729 #[tokio::test]
1734 async fn test_trigger_spawn_applies_bind_alias() {
1735 let dispatcher = EventDispatcher::new();
1737 dispatcher.bind(vec![(
1738 "AppInit".to_string(),
1739 "app\\event\\AppInit".to_string(),
1740 )]);
1741 dispatcher.listen(
1742 "AppInit",
1743 Arc::new(ClosureListener::new(|_| Ok(json!("init_called")))),
1744 false,
1745 );
1746
1747 let handles = dispatcher.trigger_spawn("AppInit", &Value::Null);
1748 assert_eq!(handles.len(), 1);
1749
1750 let result = handles.into_iter().next().unwrap().await.unwrap().unwrap();
1751 assert_eq!(result, json!("init_called"));
1752 }
1753
1754 #[tokio::test]
1755 async fn test_trigger_spawn_applies_dot_wildcard() {
1756 let dispatcher = EventDispatcher::new();
1758 dispatcher.listen(
1759 "User.login",
1760 Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
1761 false,
1762 );
1763 dispatcher.listen(
1764 "User.*",
1765 Arc::new(ClosureListener::new(|_| Ok(json!("wildcard")))),
1766 false,
1767 );
1768
1769 let handles = dispatcher.trigger_spawn("User.login", &Value::Null);
1770 assert_eq!(handles.len(), 2); let results = dispatcher.trigger_async("User.login", &Value::Null).await;
1773 assert_eq!(results[0].as_ref().unwrap(), &json!("specific"));
1774 assert_eq!(results[1].as_ref().unwrap(), &json!("wildcard"));
1775 }
1776
1777 #[tokio::test]
1778 async fn test_trigger_spawn_deduplicates() {
1779 let dispatcher = EventDispatcher::new();
1781 let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
1782
1783 dispatcher.listen("Test", listener.clone(), false);
1784 dispatcher.listen("Test", listener.clone(), false);
1785 dispatcher.listen("Test", listener.clone(), false);
1786
1787 let handles = dispatcher.trigger_spawn("Test", &Value::Null);
1788 assert_eq!(handles.len(), 1); }
1790
1791 #[tokio::test]
1792 async fn test_trigger_spawn_correct_handle_count() {
1793 let dispatcher = EventDispatcher::new();
1795 dispatcher.listen(
1796 "Test",
1797 Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
1798 false,
1799 );
1800 dispatcher.listen(
1801 "Test",
1802 Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
1803 true, );
1805 dispatcher.listen(
1806 "Test",
1807 Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
1808 false,
1809 );
1810
1811 let handles = dispatcher.trigger_spawn("Test", &Value::Null);
1812 assert_eq!(handles.len(), 3);
1813
1814 let results = dispatcher.trigger_async("Test", &Value::Null).await;
1816 assert_eq!(results[0].as_ref().unwrap(), &json!(2));
1817 assert_eq!(results[1].as_ref().unwrap(), &json!(1));
1818 assert_eq!(results[2].as_ref().unwrap(), &json!(3));
1819 }
1820
1821 #[tokio::test]
1826 async fn test_facade_trigger_spawn() {
1827 facade::listen(
1829 "FacadeTest",
1830 Arc::new(ClosureListener::new(|_| Ok(json!("facade_spawn")))),
1831 false,
1832 );
1833
1834 let handles = facade::trigger_spawn("FacadeTest", &Value::Null);
1835 assert_eq!(handles.len(), 1);
1836
1837 let result = handles.into_iter().next().unwrap().await.unwrap().unwrap();
1838 assert_eq!(result, json!("facade_spawn"));
1839 }
1840
1841 #[tokio::test]
1842 async fn test_event_trigger_async_helper() {
1843 facade::listen(
1845 "HelperTest",
1846 Arc::new(ClosureListener::new(|_| Ok(json!("helper_async")))),
1847 false,
1848 );
1849
1850 let results = event_trigger_async("HelperTest", &Value::Null).await;
1851 assert_eq!(results.len(), 1);
1852 assert_eq!(results[0].as_ref().unwrap(), &json!("helper_async"));
1853 }
1854}