use std::collections::HashMap;
use std::sync::{Arc, RwLock};
use serde_json::Value;
use tokio::task::JoinHandle;
pub trait Listener: Send + Sync {
fn handle(&self, params: &Value) -> Result<Value, EventError>;
}
pub struct ClosureListener<F>
where
F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
{
closure: F,
}
impl<F> ClosureListener<F>
where
F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
{
pub fn new(closure: F) -> Self {
Self { closure }
}
}
impl<F> Listener for ClosureListener<F>
where
F: Fn(&Value) -> Result<Value, EventError> + Send + Sync + 'static,
{
fn handle(&self, params: &Value) -> Result<Value, EventError> {
(self.closure)(params)
}
}
pub trait Subscriber: Send + Sync {
fn subscribe(&self, dispatcher: &EventDispatcher);
}
pub trait Observer: Send + Sync {
fn events(&self) -> Vec<(&'static str, Arc<dyn Listener>)>;
}
#[derive(Debug)]
pub enum EventError {
ListenerError(String),
EventNotFound(String),
InvalidParams(String),
}
impl std::fmt::Display for EventError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
EventError::ListenerError(s) => write!(f, "Event listener error: {}", s),
EventError::EventNotFound(s) => write!(f, "Event not found: {}", s),
EventError::InvalidParams(s) => write!(f, "Invalid event params: {}", s),
}
}
}
impl std::error::Error for EventError {}
pub struct EventDispatcher {
listener: RwLock<HashMap<String, Vec<Arc<dyn Listener>>>>,
bind: RwLock<HashMap<String, String>>,
}
impl Default for EventDispatcher {
fn default() -> Self {
Self::new()
}
}
impl EventDispatcher {
pub fn new() -> Self {
Self {
listener: RwLock::new(HashMap::new()),
bind: RwLock::new(HashMap::new()),
}
}
pub fn listen_events(&self, events: Vec<(String, Vec<Arc<dyn Listener>>)>) -> &Self {
let mut listener_map = self.listener.write().expect("锁被毒化");
let bind_map = self.bind.read().expect("锁被毒化");
for (event, listeners) in events {
let event = bind_map.get(&event).cloned().unwrap_or(event);
let entry = listener_map.entry(event).or_default();
entry.extend(listeners);
}
self
}
pub fn listen(&self, event: &str, listener: Arc<dyn Listener>, first: bool) -> &Self {
let mut listener_map = self.listener.write().expect("锁被毒化");
let bind_map = self.bind.read().expect("锁被毒化");
let event = bind_map
.get(event)
.cloned()
.unwrap_or_else(|| event.to_string());
let entry = listener_map.entry(event).or_default();
if first {
entry.insert(0, listener);
} else {
entry.push(listener);
}
self
}
pub fn has_listener(&self, event: &str) -> bool {
let listener_map = self.listener.read().expect("锁被毒化");
let bind_map = self.bind.read().expect("锁被毒化");
let event = bind_map.get(event).map(|s| s.as_str()).unwrap_or(event);
listener_map.contains_key(event)
}
pub fn remove(&self, event: &str) {
let mut listener_map = self.listener.write().expect("锁被毒化");
let bind_map = self.bind.read().expect("锁被毒化");
let event = bind_map
.get(event)
.cloned()
.unwrap_or_else(|| event.to_string());
listener_map.remove(&event);
}
pub fn bind(&self, events: Vec<(String, String)>) -> &Self {
let mut bind_map = self.bind.write().expect("锁被毒化");
for (alias, real_event) in events {
bind_map.insert(alias, real_event);
}
self
}
pub fn subscribe(&self, subscriber: Arc<dyn Subscriber>) -> &Self {
subscriber.subscribe(self);
self
}
pub fn observe(&self, observer: Arc<dyn Observer>, prefix: &str) -> &Self {
for (event, listener) in observer.events() {
let full_event = if prefix.is_empty() {
event.to_string()
} else {
format!("{}{}", prefix, event)
};
self.listen(&full_event, listener, false);
}
self
}
fn collect_listeners(&self, event: &str) -> Vec<Arc<dyn Listener>> {
let bind_map = self.bind.read().expect("锁被毒化");
let event = bind_map
.get(event)
.cloned()
.unwrap_or_else(|| event.to_string());
drop(bind_map);
let listener_map = self.listener.read().expect("锁被毒化");
let mut listeners: Vec<Arc<dyn Listener>> =
listener_map.get(&event).cloned().unwrap_or_default();
if let Some(dot_pos) = event.find('.') {
let prefix = &event[..dot_pos];
let wildcard = format!("{}.*", prefix);
if let Some(wildcard_listeners) = listener_map.get(&wildcard) {
listeners.extend(wildcard_listeners.clone());
}
}
drop(listener_map);
let mut seen: Vec<Arc<dyn Listener>> = Vec::new();
listeners.retain(|l| {
if seen.iter().any(|s| Arc::ptr_eq(s, l)) {
false
} else {
seen.push(l.clone());
true
}
});
listeners
}
pub fn trigger(
&self,
event: &str,
params: &Value,
once: bool,
) -> Result<Vec<Value>, EventError> {
let listeners = self.collect_listeners(event);
let mut results: Vec<Value> = Vec::new();
for listener in &listeners {
let result = listener.handle(params)?;
results.push(result.clone());
if result == Value::Bool(false) {
break;
}
if once && !result.is_null() {
break;
}
}
Ok(results)
}
pub fn trigger_spawn(
&self,
event: &str,
params: &Value,
) -> Vec<JoinHandle<Result<Value, EventError>>> {
let listeners = self.collect_listeners(event);
let params_owned = params.clone();
listeners
.into_iter()
.map(|listener| {
let params = params_owned.clone();
tokio::spawn(async move { listener.handle(¶ms) })
})
.collect()
}
pub async fn trigger_async(
&self,
event: &str,
params: &Value,
) -> Vec<Result<Value, EventError>> {
let handles = self.trigger_spawn(event, params);
let mut results = Vec::with_capacity(handles.len());
for handle in handles {
match handle.await {
Ok(result) => results.push(result),
Err(join_err) => results.push(Err(EventError::ListenerError(format!(
"Task panicked: {}",
join_err
)))),
}
}
results
}
pub fn until(&self, event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
self.trigger(event, params, true)
}
pub fn listener_count(&self, event: &str) -> usize {
let listener_map = self.listener.read().expect("锁被毒化");
let bind_map = self.bind.read().expect("锁被毒化");
let event = bind_map.get(event).map(|s| s.as_str()).unwrap_or(event);
listener_map.get(event).map(|v| v.len()).unwrap_or(0)
}
}
pub mod facade {
use super::*;
use std::sync::OnceLock;
static GLOBAL_DISPATCHER: OnceLock<EventDispatcher> = OnceLock::new();
pub fn dispatcher() -> &'static EventDispatcher {
GLOBAL_DISPATCHER.get_or_init(EventDispatcher::new)
}
pub fn listen(event: &str, listener: Arc<dyn Listener>, first: bool) {
dispatcher().listen(event, listener, first);
}
pub fn listen_events(events: Vec<(String, Vec<Arc<dyn Listener>>)>) {
dispatcher().listen_events(events);
}
pub fn has_listener(event: &str) -> bool {
dispatcher().has_listener(event)
}
pub fn remove(event: &str) {
dispatcher().remove(event);
}
pub fn bind(events: Vec<(String, String)>) {
dispatcher().bind(events);
}
pub fn subscribe(subscriber: Arc<dyn Subscriber>) {
dispatcher().subscribe(subscriber);
}
pub fn observe(observer: Arc<dyn Observer>, prefix: &str) {
dispatcher().observe(observer, prefix);
}
pub fn trigger(event: &str, params: &Value, once: bool) -> Result<Vec<Value>, EventError> {
dispatcher().trigger(event, params, once)
}
pub fn until(event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
dispatcher().until(event, params)
}
pub fn trigger_spawn(
event: &str,
params: &Value,
) -> Vec<JoinHandle<Result<Value, EventError>>> {
dispatcher().trigger_spawn(event, params)
}
pub async fn trigger_async(event: &str, params: &Value) -> Vec<Result<Value, EventError>> {
dispatcher().trigger_async(event, params).await
}
#[cfg(test)]
pub fn _reset_for_test() {
}
}
pub fn event_trigger(event: &str, params: &Value) -> Vec<Value> {
facade::trigger(event, params, false).unwrap_or_default()
}
pub async fn event_trigger_async(event: &str, params: &Value) -> Vec<Result<Value, EventError>> {
facade::trigger_async(event, params).await
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn test_event_listen_and_has_listener() {
let dispatcher = EventDispatcher::new();
assert!(!dispatcher.has_listener("UserLogin"));
dispatcher.listen(
"UserLogin",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
assert!(dispatcher.has_listener("UserLogin"));
}
#[test]
fn test_event_remove_listener() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"UserLogin",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
assert!(dispatcher.has_listener("UserLogin"));
dispatcher.remove("UserLogin");
assert!(!dispatcher.has_listener("UserLogin"));
}
#[test]
fn test_event_listen_first_priority() {
let dispatcher = EventDispatcher::new();
let call_order = Arc::new(AtomicUsize::new(0));
let order1 = call_order.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
order1.store(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let order2 = call_order.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
order2.store(2, Ordering::SeqCst);
Ok(Value::Null)
})),
true, );
dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(call_order.load(Ordering::SeqCst), 1);
}
#[test]
fn test_event_listen_events_batch() {
let dispatcher = EventDispatcher::new();
dispatcher.listen_events(vec![
(
"UserLogin".to_string(),
vec![Arc::new(ClosureListener::new(|_| Ok(Value::Null)))],
),
(
"UserLogout".to_string(),
vec![
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
],
),
]);
assert!(dispatcher.has_listener("UserLogin"));
assert_eq!(dispatcher.listener_count("UserLogout"), 2);
}
#[test]
fn test_event_bind_alias() {
let dispatcher = EventDispatcher::new();
dispatcher.bind(vec![(
"AppInit".to_string(),
"app\\event\\AppInit".to_string(),
)]);
dispatcher.listen(
"AppInit",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
assert!(!dispatcher.has_listener("AppInit_alias_check"));
assert!(dispatcher.has_listener("app\\event\\AppInit"));
}
#[test]
fn test_event_trigger_returns_all_results() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results, vec![json!(1), json!(2)]);
}
#[test]
fn test_event_trigger_once_stops_at_non_null() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("stop")))), false,
);
let executed = Arc::new(AtomicUsize::new(0));
let exec_clone = executed.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
exec_clone.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, true).unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[0], Value::Null);
assert_eq!(results[1], json!("stop"));
assert_eq!(executed.load(Ordering::SeqCst), 0); }
#[test]
fn test_event_trigger_false_stops_execution() {
let dispatcher = EventDispatcher::new();
let executed = Arc::new(AtomicUsize::new(0));
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(Value::Bool(false)))), false,
);
let exec_clone = executed.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
exec_clone.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0], Value::Bool(false));
assert_eq!(executed.load(Ordering::SeqCst), 0); }
#[test]
fn test_event_trigger_empty_event() {
let dispatcher = EventDispatcher::new();
let results = dispatcher
.trigger("Nonexistent", &Value::Null, false)
.unwrap();
assert!(results.is_empty());
}
#[test]
fn test_event_trigger_passes_params() {
let dispatcher = EventDispatcher::new();
let received = Arc::new(std::sync::Mutex::new(Value::Null));
let recv_clone = received.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |params| {
*recv_clone.lock().unwrap() = params.clone();
Ok(Value::Null)
})),
false,
);
let params = json!({"user_id": 123, "action": "login"});
dispatcher.trigger("Test", ¶ms, false).unwrap();
assert_eq!(*received.lock().unwrap(), params);
}
#[test]
fn test_event_dot_wildcard() {
let dispatcher = EventDispatcher::new();
let executed = Arc::new(AtomicUsize::new(0));
dispatcher.listen(
"User.login",
Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
false,
);
let exec_clone = executed.clone();
dispatcher.listen(
"User.*",
Arc::new(ClosureListener::new(move |_| {
exec_clone.fetch_add(1, Ordering::SeqCst);
Ok(json!("wildcard"))
})),
false,
);
let results = dispatcher
.trigger("User.login", &Value::Null, false)
.unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[0], json!("specific"));
assert_eq!(results[1], json!("wildcard"));
assert_eq!(executed.load(Ordering::SeqCst), 1);
}
#[test]
fn test_event_dot_wildcard_no_wildcard_listener() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"User.login",
Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
false,
);
let results = dispatcher
.trigger("User.login", &Value::Null, false)
.unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0], json!("specific"));
}
#[test]
fn test_event_dedup_same_listener_instance() {
let dispatcher = EventDispatcher::new();
let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
dispatcher.listen("Test", listener.clone(), false);
dispatcher.listen("Test", listener.clone(), false);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0], json!(1));
}
#[test]
fn test_event_no_dedup_different_listeners() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results.len(), 2);
}
struct TestSubscriber {
login_count: Arc<AtomicUsize>,
}
impl Subscriber for TestSubscriber {
fn subscribe(&self, dispatcher: &EventDispatcher) {
let count = self.login_count.clone();
dispatcher.listen(
"UserLogin",
Arc::new(ClosureListener::new(move |_| {
count.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let count2 = self.login_count.clone();
dispatcher.listen(
"UserLogout",
Arc::new(ClosureListener::new(move |_| {
count2.fetch_add(10, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
}
}
#[test]
fn test_event_subscriber_registers_multiple_listeners() {
let dispatcher = EventDispatcher::new();
let login_count = Arc::new(AtomicUsize::new(0));
let subscriber = Arc::new(TestSubscriber {
login_count: login_count.clone(),
});
dispatcher.subscribe(subscriber);
assert!(dispatcher.has_listener("UserLogin"));
assert!(dispatcher.has_listener("UserLogout"));
dispatcher
.trigger("UserLogin", &Value::Null, false)
.unwrap();
assert_eq!(login_count.load(Ordering::SeqCst), 1);
dispatcher
.trigger("UserLogout", &Value::Null, false)
.unwrap();
assert_eq!(login_count.load(Ordering::SeqCst), 11); }
struct TestObserver {
counter: Arc<AtomicUsize>,
}
impl Observer for TestObserver {
fn events(&self) -> Vec<(&'static str, Arc<dyn Listener>)> {
let c1 = self.counter.clone();
let c2 = self.counter.clone();
vec![
(
"Login",
Arc::new(ClosureListener::new(move |_| {
c1.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
),
(
"Logout",
Arc::new(ClosureListener::new(move |_| {
c2.fetch_add(100, Ordering::SeqCst);
Ok(Value::Null)
})),
),
]
}
}
#[test]
fn test_event_observer_auto_registers() {
let dispatcher = EventDispatcher::new();
let counter = Arc::new(AtomicUsize::new(0));
let observer = Arc::new(TestObserver {
counter: counter.clone(),
});
dispatcher.observe(observer, "");
assert!(dispatcher.has_listener("Login"));
assert!(dispatcher.has_listener("Logout"));
dispatcher.trigger("Login", &Value::Null, false).unwrap();
assert_eq!(counter.load(Ordering::SeqCst), 1);
dispatcher.trigger("Logout", &Value::Null, false).unwrap();
assert_eq!(counter.load(Ordering::SeqCst), 101);
}
#[test]
fn test_event_observer_with_prefix() {
let dispatcher = EventDispatcher::new();
let counter = Arc::new(AtomicUsize::new(0));
let observer = Arc::new(TestObserver {
counter: counter.clone(),
});
dispatcher.observe(observer, "User");
assert!(dispatcher.has_listener("UserLogin"));
assert!(dispatcher.has_listener("UserLogout"));
dispatcher
.trigger("UserLogin", &Value::Null, false)
.unwrap();
assert_eq!(counter.load(Ordering::SeqCst), 1);
}
#[test]
fn test_event_until_returns_first_non_null() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("first_valid")))),
false,
);
let results = dispatcher.until("Test", &Value::Null).unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[1], json!("first_valid"));
}
#[test]
fn test_closure_listener_executes_closure() {
let listener = ClosureListener::new(|params| {
assert_eq!(params, &json!({"key": "value"}));
Ok(json!("result"))
});
let result = listener.handle(&json!({"key": "value"})).unwrap();
assert_eq!(result, json!("result"));
}
#[test]
fn test_closure_listener_returns_null() {
let listener = ClosureListener::new(|_| Ok(Value::Null));
let result = listener.handle(&Value::Null).unwrap();
assert!(result.is_null());
}
struct CustomListener {
id: i32,
}
impl Listener for CustomListener {
fn handle(&self, _params: &Value) -> Result<Value, EventError> {
Ok(json!({"listener_id": self.id}))
}
}
#[test]
fn test_custom_listener_trait_impl() {
let dispatcher = EventDispatcher::new();
dispatcher.listen("Test", Arc::new(CustomListener { id: 42 }), false);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results, vec![json!({"listener_id": 42})]);
}
#[test]
fn test_listener_error_propagates() {
let dispatcher = EventDispatcher::new();
let executed = Arc::new(AtomicUsize::new(0));
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| {
Err(EventError::ListenerError("test error".to_string()))
})),
false,
);
let exec_clone = executed.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
exec_clone.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let result = dispatcher.trigger("Test", &Value::Null, false);
assert!(result.is_err());
assert_eq!(executed.load(Ordering::SeqCst), 0); }
#[test]
fn test_r5_php_event_listen_then_trigger() {
let dispatcher = EventDispatcher::new();
let received = Arc::new(std::sync::Mutex::new(Value::Null));
let recv_clone = received.clone();
dispatcher.listen(
"UserLogin",
Arc::new(ClosureListener::new(move |params| {
*recv_clone.lock().unwrap() = params.clone();
Ok(Value::Null)
})),
false,
);
let params = json!({"user_id": 123, "username": "alice"});
dispatcher.trigger("UserLogin", ¶ms, false).unwrap();
assert_eq!(*received.lock().unwrap(), params);
}
#[test]
fn test_r5_php_event_bind_alias_resolution() {
let dispatcher = EventDispatcher::new();
dispatcher.bind(vec![(
"AppInit".to_string(),
"think\\event\\AppInit".to_string(),
)]);
dispatcher.listen(
"AppInit",
Arc::new(ClosureListener::new(|_| Ok(json!("init_called")))),
false,
);
assert!(dispatcher.has_listener("think\\event\\AppInit"));
assert_eq!(dispatcher.listener_count("think\\event\\AppInit"), 1);
let results = dispatcher.trigger("AppInit", &Value::Null, false).unwrap();
assert_eq!(results, vec![json!("init_called")]);
}
#[test]
fn test_r5_php_event_first_array_unshift() {
let dispatcher = EventDispatcher::new();
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
let o1 = order.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
o1.lock().unwrap().push(1);
Ok(Value::Null)
})),
false,
);
let o2 = order.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
o2.lock().unwrap().push(2);
Ok(Value::Null)
})),
true, );
dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(*order.lock().unwrap(), vec![2, 1]);
}
#[test]
fn test_r5_php_event_trigger_returns_array() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results, vec![json!(1), json!(2), json!(3)]);
}
#[test]
fn test_r5_php_event_until_returns_last_non_null() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("first")))),
false,
);
let results = dispatcher.until("Test", &Value::Null).unwrap();
assert_eq!(results.len(), 2);
}
#[test]
fn test_r5_php_event_false_stops() {
let dispatcher = EventDispatcher::new();
let executed = Arc::new(AtomicUsize::new(0));
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(Value::Bool(false)))),
false,
);
let exec_clone = executed.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
exec_clone.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0], Value::Bool(false));
}
#[test]
fn test_r5_php_event_dot_wildcard_merge() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"User.login",
Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
false,
);
dispatcher.listen(
"User.*",
Arc::new(ClosureListener::new(|_| Ok(json!("wildcard")))),
false,
);
let results = dispatcher
.trigger("User.login", &Value::Null, false)
.unwrap();
assert_eq!(results, vec![json!("specific"), json!("wildcard")]);
}
#[test]
fn test_r5_php_event_array_unique_dedup() {
let dispatcher = EventDispatcher::new();
let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
dispatcher.listen("Test", listener.clone(), false);
dispatcher.listen("Test", listener.clone(), false);
dispatcher.listen("Test", listener.clone(), false);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results.len(), 1);
}
#[test]
fn test_r5_php_event_remove_clears_listeners() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
assert!(dispatcher.has_listener("Test"));
dispatcher.remove("Test");
assert!(!dispatcher.has_listener("Test"));
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert!(results.is_empty());
}
#[test]
fn test_r5_php_event_has_listener_with_bind() {
let dispatcher = EventDispatcher::new();
dispatcher.bind(vec![(
"AppInit".to_string(),
"app\\event\\AppInit".to_string(),
)]);
dispatcher.listen(
"AppInit",
Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
false,
);
assert!(dispatcher.has_listener("AppInit"));
assert!(dispatcher.has_listener("app\\event\\AppInit"));
}
#[test]
fn test_r5_php_event_listen_events_batch_merge() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen_events(vec![(
"Test".to_string(),
vec![
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
],
)]);
let results = dispatcher.trigger("Test", &Value::Null, false).unwrap();
assert_eq!(results, vec![json!(1), json!(2), json!(3)]);
}
#[tokio::test]
async fn test_trigger_spawn_returns_join_handles() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
false,
);
let handles = dispatcher.trigger_spawn("Test", &Value::Null);
assert_eq!(handles.len(), 2);
for handle in handles {
let result = handle.await.unwrap().unwrap();
assert!(result == json!(1) || result == json!(2));
}
}
#[tokio::test]
async fn test_trigger_spawn_empty_event() {
let dispatcher = EventDispatcher::new();
let handles = dispatcher.trigger_spawn("Nonexistent", &Value::Null);
assert!(handles.is_empty());
}
#[tokio::test]
async fn test_trigger_spawn_passes_params() {
let dispatcher = EventDispatcher::new();
let received = Arc::new(std::sync::Mutex::new(Value::Null));
let recv_clone = received.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |params| {
*recv_clone.lock().unwrap() = params.clone();
Ok(Value::Null)
})),
false,
);
let params = json!({"user_id": 123, "action": "login"});
let handles = dispatcher.trigger_spawn("Test", ¶ms);
for handle in handles {
let _ = handle.await;
}
assert_eq!(*received.lock().unwrap(), params);
}
#[tokio::test]
async fn test_trigger_spawn_all_listeners_execute() {
let dispatcher = EventDispatcher::new();
let counter = Arc::new(AtomicUsize::new(0));
for _ in 0..5 {
let c = counter.clone();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(move |_| {
c.fetch_add(1, Ordering::SeqCst);
Ok(Value::Null)
})),
false,
);
}
let handles = dispatcher.trigger_spawn("Test", &Value::Null);
for handle in handles {
let _ = handle.await;
}
assert_eq!(counter.load(Ordering::SeqCst), 5);
}
#[tokio::test]
async fn test_trigger_async_awaits_all() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
false,
);
let results = dispatcher.trigger_async("Test", &Value::Null).await;
assert_eq!(results.len(), 3);
assert!(results.iter().all(|r| r.is_ok()));
}
#[tokio::test]
async fn test_trigger_async_collects_results_in_order() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("first")))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("second")))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("third")))),
false,
);
let results = dispatcher.trigger_async("Test", &Value::Null).await;
assert_eq!(results[0].as_ref().unwrap(), &json!("first"));
assert_eq!(results[1].as_ref().unwrap(), &json!("second"));
assert_eq!(results[2].as_ref().unwrap(), &json!("third"));
}
#[tokio::test]
async fn test_trigger_async_handles_errors() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("ok")))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| {
Err(EventError::ListenerError("test error".to_string()))
})),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("ok2")))),
false,
);
let results = dispatcher.trigger_async("Test", &Value::Null).await;
assert_eq!(results.len(), 3);
assert!(results[0].is_ok());
assert!(results[1].is_err());
assert!(results[2].is_ok());
}
#[tokio::test]
async fn test_trigger_async_handles_panic() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!("ok")))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| panic!("test panic"))),
false,
);
let results = dispatcher.trigger_async("Test", &Value::Null).await;
assert_eq!(results.len(), 2);
assert!(results[0].is_ok());
assert!(results[1].is_err()); }
#[tokio::test]
async fn test_trigger_spawn_applies_bind_alias() {
let dispatcher = EventDispatcher::new();
dispatcher.bind(vec![(
"AppInit".to_string(),
"app\\event\\AppInit".to_string(),
)]);
dispatcher.listen(
"AppInit",
Arc::new(ClosureListener::new(|_| Ok(json!("init_called")))),
false,
);
let handles = dispatcher.trigger_spawn("AppInit", &Value::Null);
assert_eq!(handles.len(), 1);
let result = handles.into_iter().next().unwrap().await.unwrap().unwrap();
assert_eq!(result, json!("init_called"));
}
#[tokio::test]
async fn test_trigger_spawn_applies_dot_wildcard() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"User.login",
Arc::new(ClosureListener::new(|_| Ok(json!("specific")))),
false,
);
dispatcher.listen(
"User.*",
Arc::new(ClosureListener::new(|_| Ok(json!("wildcard")))),
false,
);
let handles = dispatcher.trigger_spawn("User.login", &Value::Null);
assert_eq!(handles.len(), 2);
let results = dispatcher.trigger_async("User.login", &Value::Null).await;
assert_eq!(results[0].as_ref().unwrap(), &json!("specific"));
assert_eq!(results[1].as_ref().unwrap(), &json!("wildcard"));
}
#[tokio::test]
async fn test_trigger_spawn_deduplicates() {
let dispatcher = EventDispatcher::new();
let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
dispatcher.listen("Test", listener.clone(), false);
dispatcher.listen("Test", listener.clone(), false);
dispatcher.listen("Test", listener.clone(), false);
let handles = dispatcher.trigger_spawn("Test", &Value::Null);
assert_eq!(handles.len(), 1); }
#[tokio::test]
async fn test_trigger_spawn_correct_handle_count() {
let dispatcher = EventDispatcher::new();
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(1)))),
false,
);
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(2)))),
true, );
dispatcher.listen(
"Test",
Arc::new(ClosureListener::new(|_| Ok(json!(3)))),
false,
);
let handles = dispatcher.trigger_spawn("Test", &Value::Null);
assert_eq!(handles.len(), 3);
let results = dispatcher.trigger_async("Test", &Value::Null).await;
assert_eq!(results[0].as_ref().unwrap(), &json!(2));
assert_eq!(results[1].as_ref().unwrap(), &json!(1));
assert_eq!(results[2].as_ref().unwrap(), &json!(3));
}
#[tokio::test]
async fn test_facade_trigger_spawn() {
facade::listen(
"FacadeTest",
Arc::new(ClosureListener::new(|_| Ok(json!("facade_spawn")))),
false,
);
let handles = facade::trigger_spawn("FacadeTest", &Value::Null);
assert_eq!(handles.len(), 1);
let result = handles.into_iter().next().unwrap().await.unwrap().unwrap();
assert_eq!(result, json!("facade_spawn"));
}
#[tokio::test]
async fn test_event_trigger_async_helper() {
facade::listen(
"HelperTest",
Arc::new(ClosureListener::new(|_| Ok(json!("helper_async")))),
false,
);
let results = event_trigger_async("HelperTest", &Value::Null).await;
assert_eq!(results.len(), 1);
assert_eq!(results[0].as_ref().unwrap(), &json!("helper_async"));
}
}