Skip to main content

sz_rust_core/
event.rs

1//! 事件系统(对齐 PHP `think\Event`)
2//!
3//! 对齐 PHP `think\Event`(272 行)+ `think\facade\Event` + `event()` 助手函数。
4//!
5//! ## PHP API 对齐映射
6//!
7//! | PHP `think\Event` 方法 | Rust `EventDispatcher` 方法 | 说明 |
8//! |------------------------|----------------------------|------|
9//! | `listenEvents(array $events)` | `listen_events(events)` | 批量注册事件监听 |
10//! | `listen(string $event, $listener, bool $first = false)` | `listen(event, listener, first)` | 注册事件监听 |
11//! | `hasListener(string $event): bool` | `has_listener(event)` | 是否存在事件监听 |
12//! | `remove(string $event): void` | `remove(event)` | 移除事件监听 |
13//! | `bind(array $events)` | `bind(events)` | 指定事件别名标识 |
14//! | `subscribe($subscriber)` | `subscribe(subscriber)` | 注册事件订阅者 |
15//! | `observe($observer, $prefix = '')` | `observe(observer, prefix)` | 自动注册观察者 |
16//! | `trigger($event, $params = null, bool $once = false)` | `trigger(event, params, once)` | 触发事件 |
17//! | `until($event, $params = null)` | `until(event, params)` | 只获取一个有效返回值 |
18//!
19//! ## PHP Listener 类型映射
20//!
21//! PHP `$listener` 可以是:
22//! - 闭包 `function($params) { ... }` → Rust `Arc<dyn Fn(¶ms) -> Result<Value>>`
23//! - 对象方法数组 `[$obj, 'method']` → Rust `Arc<dyn Listener>` trait
24//! - 字符串类名(调用 `handle` 方法) → Rust `Arc<dyn Listener>` trait
25//! - 静态方法字符串 `Class::method` → Rust `Arc<dyn Listener>` trait
26//!
27//! Rust 端统一用 `Listener` trait(`handle(¶ms) -> Result<Value>`)+ 闭包包装。
28
29use std::collections::HashMap;
30use std::sync::{Arc, RwLock};
31
32use serde_json::Value;
33use tokio::task::JoinHandle;
34
35/// 事件监听器 trait(对齐 PHP `Listener` 接口 / `handle` 方法)
36///
37/// PHP 中监听器通常是一个类,含 `handle($params)` 方法:
38/// ```php
39/// class UserLoginListener
40/// {
41///     public function handle($params)
42///     {
43///         // 处理事件
44///     }
45/// }
46/// ```
47///
48/// Rust 端实现 `Listener` trait 即可:
49/// ```ignore
50/// struct UserLoginListener;
51///
52/// impl Listener for UserLoginListener {
53///     fn handle(&self, params: &Value) -> Result<Value, EventError> {
54///         // 处理事件
55///         Ok(Value::Null)
56///     }
57/// }
58/// ```
59pub trait Listener: Send + Sync {
60    /// 处理事件(对齐 PHP `Listener::handle($params)`)
61    ///
62    /// 返回 `Ok(Value::Null)` 等价 PHP 无返回值(继续执行后续监听器)。
63    /// 返回 `Ok(其他值)` 等价 PHP 返回非 null 值(在 `once=true` 模式下会停止)。
64    /// 返回 `Err(_)` 等价 PHP 抛出异常(会停止后续监听器执行)。
65    fn handle(&self, params: &Value) -> Result<Value, EventError>;
66}
67
68/// 闭包监听器(对齐 PHP 闭包 `$listener = function($params) { ... }`)
69///
70/// PHP 允许直接传入闭包作为监听器:
71/// ```php
72/// Event::listen('UserLogin', function($params) {
73///     // 处理事件
74/// });
75/// ```
76///
77/// Rust 端用 `ClosureListener::new(closure)` 包装:
78/// ```ignore
79/// dispatcher.listen("UserLogin", ClosureListener::new(|params| {
80///     // 处理事件
81///     Ok(Value::Null)
82/// }), false);
83/// ```
84pub 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    /// 创建一个闭包监听器
96    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
110/// 事件订阅者 trait(对齐 PHP 订阅者 `subscribe(Event $event)` 方法)
111///
112/// PHP 订阅者类含 `subscribe(Event $event)` 方法,手动注册多个监听器:
113/// ```php
114/// class UserEventSubscriber
115/// {
116///     public function onUserLogin($params) { ... }
117///     public function onUserLogout($params) { ... }
118///
119///     public function subscribe(Event $event)
120///     {
121///         $event->listen('UserLogin', [$this, 'onUserLogin']);
122///         $event->listen('UserLogout', [$this, 'onUserLogout']);
123///     }
124/// }
125/// ```
126///
127/// Rust 端实现 `Subscriber` trait:
128/// ```ignore
129/// struct UserEventSubscriber;
130///
131/// impl Subscriber for UserEventSubscriber {
132///     fn subscribe(&self, dispatcher: &EventDispatcher) {
133///         dispatcher.listen("UserLogin", Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false);
134///         dispatcher.listen("UserLogout", Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false);
135///     }
136/// }
137/// ```
138pub trait Subscriber: Send + Sync {
139    /// 订阅事件(对齐 PHP `Subscriber::subscribe(Event $event)`)
140    fn subscribe(&self, dispatcher: &EventDispatcher);
141}
142
143/// 观察者 trait(对齐 PHP `observe` 智能订阅)
144///
145/// PHP `observe($observer)` 通过反射获取所有 `onXxx` 公开方法,
146/// 自动注册为 `Xxx` 事件的监听器:
147/// ```php
148/// class UserObserver
149/// {
150///     public function onLogin($params) { ... }
151///     public function onLogout($params) { ... }
152/// }
153/// // 自动注册 Login 和 Logout 事件
154/// Event::observe(new UserObserver());
155/// ```
156///
157/// Rust 端用 `Observer` trait 模拟(Rust 无反射,需手动声明事件映射):
158/// ```text
159/// struct UserObserver;
160///
161/// impl Observer for UserObserver {
162///     fn events(&self) -> &[(&str, Arc<dyn Listener>)] {
163///         &[
164///             ("Login", Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
165///             ("Logout", Arc::new(ClosureListener::new(|_| Ok(Value::Null))),
166///         ]
167///     }
168/// }
169/// ```
170pub trait Observer: Send + Sync {
171    /// 返回观察者监听的所有事件(对齐 PHP 反射 `onXxx` 方法 → `Xxx` 事件)
172    ///
173    /// 每个元素是 `(事件名, 监听器)`,等价 PHP `listen($prefix . $event, [$observer, 'on' . $event])`
174    fn events(&self) -> Vec<(&'static str, Arc<dyn Listener>)>;
175}
176
177/// 事件错误类型
178#[derive(Debug)]
179pub enum EventError {
180    /// 监听器执行失败
181    ListenerError(String),
182    /// 事件不存在
183    EventNotFound(String),
184    /// 参数错误
185    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
200/// 事件分发器(对齐 PHP `think\Event`)
201///
202/// 对齐 PHP `think\Event` 类(272 行),提供事件监听/触发/订阅/观察者全 API。
203///
204/// ## 线程安全
205///
206/// 内部用 `Arc<RwLock<>>` 保护 `listener` 和 `bind` 映射,允许并发读写。
207/// `Listener` 要求 `Send + Sync`,可在多线程环境中分发事件。
208///
209/// ## PHP 行为对齐
210///
211/// 1. **事件别名(bind)**:`listen('AppInit', ...)` 实际注册到 `event\AppInit::class`
212/// 2. **优先执行(first)**:`listen(event, listener, true)` 用 `array_unshift` 插入队首
213/// 3. **触发返回值**:`trigger` 返回所有监听器返回值数组;`once=true` 返回首个非 null
214/// 4. **false 停止**:监听器返回 `false` 时停止后续监听器执行
215/// 5. **once 停止**:`once=true` 时监听器返回非 null 值停止后续
216/// 6. **点号通配**:`trigger('User.login')` 同时触发 `User.login` 和 `User.*` 监听器
217/// 7. **array_unique**:`trigger` 对监听器列表去重(`SORT_REGULAR`)
218pub struct EventDispatcher {
219    /// 监听者映射:event => [listener1, listener2, ...](对齐 PHP `$listener`)
220    listener: RwLock<HashMap<String, Vec<Arc<dyn Listener>>>>,
221
222    /// 事件别名映射:alias => real_event(对齐 PHP `$bind`)
223    bind: RwLock<HashMap<String, String>>,
224}
225
226impl Default for EventDispatcher {
227    fn default() -> Self {
228        Self::new()
229    }
230}
231
232impl EventDispatcher {
233    /// 创建新的事件分发器(对齐 PHP `__construct(App $app)`,Rust 无需 App 容器)
234    pub fn new() -> Self {
235        Self {
236            listener: RwLock::new(HashMap::new()),
237            bind: RwLock::new(HashMap::new()),
238        }
239    }
240
241    /// 批量注册事件监听(对齐 PHP `listenEvents(array $events)`)
242    ///
243    /// PHP:
244    /// ```php
245    /// Event::listenEvents([
246    ///     'UserLogin' => [LoginListener1::class, LoginListener2::class],
247    ///     'UserLogout' => [LogoutListener::class],
248    /// ]);
249    /// ```
250    ///
251    /// Rust:
252    /// ```ignore
253    /// dispatcher.listen_events(vec![
254    ///     ("UserLogin".to_string(), vec![Arc::new(LoginListener1), Arc::new(LoginListener2)]),
255    ///     ("UserLogout".to_string(), vec![Arc::new(LogoutListener)]),
256    /// ]);
257    /// ```
258    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            // 应用事件别名(对齐 PHP `if (isset($this->bind[$event]))`)
264            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    /// 注册事件监听(对齐 PHP `listen(string $event, $listener, bool $first = false)`)
274    ///
275    /// PHP:
276    /// ```php
277    /// Event::listen('UserLogin', function($params) { ... });
278    /// Event::listen('UserLogin', UserLoginListener::class);
279    /// Event::listen('UserLogin', [UserLoginListener::class, 'handle'], true); // 优先执行
280    /// ```
281    ///
282    /// Rust:
283    /// ```ignore
284    /// dispatcher.listen("UserLogin", Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false);
285    /// dispatcher.listen("UserLogin", Arc::new(UserLoginListener), true); // 优先执行
286    /// ```
287    ///
288    /// **PHP 行为对齐**:
289    /// - `first=true` 时插入队首(`array_unshift`)
290    /// - `first=false` 时追加队尾(`$this->listener[$event][]`)
291    /// - 应用事件别名(`bind` 映射)
292    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        // 应用事件别名(对齐 PHP `if (isset($this->bind[$event]))`)
297        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            // 对齐 PHP `array_unshift($this->listener[$event], $listener)`
305            entry.insert(0, listener);
306        } else {
307            // 对齐 PHP `$this->listener[$event][] = $listener`
308            entry.push(listener);
309        }
310
311        self
312    }
313
314    /// 是否存在事件监听(对齐 PHP `hasListener(string $event): bool`)
315    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        // 应用事件别名(对齐 PHP `if (isset($this->bind[$event]))`)
320        let event = bind_map.get(event).map(|s| s.as_str()).unwrap_or(event);
321
322        listener_map.contains_key(event)
323    }
324
325    /// 移除事件监听(对齐 PHP `remove(string $event): void`)
326    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        // 应用事件别名(对齐 PHP `if (isset($this->bind[$event]))`)
331        let event = bind_map
332            .get(event)
333            .cloned()
334            .unwrap_or_else(|| event.to_string());
335
336        // 对齐 PHP `unset($this->listener[$event])`
337        listener_map.remove(&event);
338    }
339
340    /// 指定事件别名标识(对齐 PHP `bind(array $events)`)
341    ///
342    /// PHP:
343    /// ```php
344    /// Event::bind([
345    ///     'UserLogin' => 'app\event\UserLogin',
346    /// ]);
347    /// ```
348    ///
349    /// Rust:
350    /// ```ignore
351    /// dispatcher.bind(vec![("UserLogin".to_string(), "app\\event\\UserLogin".to_string())]);
352    /// ```
353    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    /// 注册事件订阅者(对齐 PHP `subscribe($subscriber)`)
362    ///
363    /// PHP:
364    /// ```php
365    /// Event::subscribe(UserEventSubscriber::class);
366    /// // 或
367    /// Event::subscribe(new UserEventSubscriber());
368    /// ```
369    ///
370    /// Rust:
371    /// ```ignore
372    /// dispatcher.subscribe(Arc::new(UserEventSubscriber));
373    /// ```
374    ///
375    /// **PHP 行为对齐**:
376    /// - 若订阅者有 `subscribe` 方法 → 手动订阅(调用 `$subscriber->subscribe($this)`)
377    /// - 否则 → 智能订阅(调用 `observe($subscriber)`)
378    ///
379    /// Rust 端统一通过 `Subscriber` trait 的 `subscribe` 方法手动订阅。
380    /// 若需智能订阅,用 `observe` 方法。
381    pub fn subscribe(&self, subscriber: Arc<dyn Subscriber>) -> &Self {
382        // 对齐 PHP `if (method_exists($subscriber, 'subscribe'))` → 手动订阅
383        subscriber.subscribe(self);
384        self
385    }
386
387    /// 自动注册事件观察者(对齐 PHP `observe($observer, string $prefix = '')`)
388    ///
389    /// PHP:
390    /// ```php
391    /// Event::observe(new UserObserver());
392    /// // 自动注册 UserObserver 的 onLogin() → 'Login' 事件
393    /// ```
394    ///
395    /// Rust:
396    /// ```ignore
397    /// dispatcher.observe(Arc::new(UserObserver), "");
398    /// ```
399    ///
400    /// **PHP 行为对齐**:
401    /// - 反射获取所有 `onXxx` 公开方法
402    /// - 注册 `listen($prefix . $event_name, [$observer, 'on' . $event_name])`
403    /// - 若有 `eventPrefix` 属性,用作前缀
404    ///
405    /// Rust 端通过 `Observer` trait 的 `events()` 方法声明事件映射,
406    /// 避免运行时反射(Rust 无反射)。
407    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            // 对齐 PHP `$this->listen($prefix . substr($name, 2), [$observer, $name])`
415            self.listen(&full_event, listener, false);
416        }
417        self
418    }
419
420    /// 收集事件的监听器列表(私有助手,应用 bind 别名 + 点号通配 + 去重)
421    ///
422    /// 提取自 `trigger` 和 `trigger_spawn` 的公共逻辑,避免代码冗余。
423    fn collect_listeners(&self, event: &str) -> Vec<Arc<dyn Listener>> {
424        let bind_map = self.bind.read().expect("锁被毒化");
425
426        // 应用事件别名(对齐 PHP `if (isset($this->bind[$event]))`)
427        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        // 对齐 PHP `$listeners = $this->listener[$event] ?? []`
437        let mut listeners: Vec<Arc<dyn Listener>> =
438            listener_map.get(&event).cloned().unwrap_or_default();
439
440        // 对齐 PHP 点号通配:`if (strpos($event, '.'))` → 触发 `prefix.*`
441        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                // 对齐 PHP `array_merge($listeners, $this->listener[$prefix . '.*'])`
446                listeners.extend(wildcard_listeners.clone());
447            }
448        }
449
450        drop(listener_map);
451
452        // 对齐 PHP `array_unique($listeners, SORT_REGULAR)`
453        // Rust 端用 Arc::ptr_eq 指针比较去重(等价 PHP 对象引用去重)
454        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    /// 触发事件(对齐 PHP `trigger($event, $params = null, bool $once = false)`)
468    ///
469    /// PHP:
470    /// ```php
471    /// $results = Event::trigger('UserLogin', ['user_id' => 123]);
472    /// $first = Event::trigger('UserLogin', ['user_id' => 123], true); // 只获取一个有效返回值
473    /// ```
474    ///
475    /// Rust:
476    /// ```ignore
477    /// let results = dispatcher.trigger("UserLogin", &json!({"user_id": 123}), false).unwrap();
478    /// let first = dispatcher.trigger("UserLogin", &json!({"user_id": 123}), true).unwrap();
479    /// ```
480    ///
481    /// **PHP 行为对齐**:
482    /// 1. 若 `$event` 是对象,取类名作为事件名,对象作为参数
483    /// 2. 应用事件别名(`bind` 映射)
484    /// 3. 点号通配:`User.login` 同时触发 `User.login` 和 `User.*`
485    /// 4. `array_unique` 对监听器去重
486    /// 5. 逐个调用 `dispatch`,返回值收集到 `$result`
487    /// 6. 监听器返回 `false` → 停止后续
488    /// 7. `once=true` 且监听器返回非 null → 停止后续
489    /// 8. `once=false` 返回所有返回值数组;`once=true` 返回最后一个非 null 返回值
490    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            // 对齐 PHP `$result[$key] = $this->dispatch($listener, $params)`
502            let result = listener.handle(params)?;
503            results.push(result.clone());
504
505            // 对齐 PHP `if (false === $result[$key] || (!is_null($result[$key]) && $once)) break`
506            // PHP `false` 在 Rust 中用 `Value::Bool(false)` 表示
507            if result == Value::Bool(false) {
508                break;
509            }
510            // PHP `!is_null($result[$key]) && $once`:非 null 且 once 模式 → 停止
511            if once && !result.is_null() {
512                break;
513            }
514        }
515
516        Ok(results)
517    }
518
519    /// 异步触发事件 — fire-and-forget(Rust 特有扩展,对齐 think-swoole 异步事件分发)
520    ///
521    /// 每个监听器在独立的 `tokio::task` 中并发执行,立即返回 `JoinHandle` 列表不等待。
522    /// 适用于非关键事件(日志、指标、通知),不阻塞当前请求。
523    ///
524    /// **注意**:必须在 tokio 运行时中调用(axum 服务器已提供运行时)。
525    /// 调用方可选择 `.await` JoinHandle 获取结果,或丢弃 JoinHandle 实现 fire-and-forget。
526    ///
527    /// **与同步 `trigger` 的差异**:
528    /// - 同步 `trigger`:逐个串行执行,支持 `once`/`false` 停止
529    /// - 异步 `trigger_spawn`:并发执行,不支持 `once`/`false` 停止(各监听器独立运行)
530    /// - 两者的 bind 别名 / 点号通配 / 去重逻辑完全一致(共用 `collect_listeners`)
531    ///
532    /// Rust:
533    /// ```ignore
534    /// // fire-and-forget(丢弃 JoinHandle)
535    /// dispatcher.trigger_spawn("UserLogin", &json!({"user_id": 123}));
536    ///
537    /// // 等待所有监听器完成
538    /// let handles = dispatcher.trigger_spawn("UserLogin", &json!({"user_id": 123}));
539    /// for handle in handles {
540    ///     let _ = handle.await;
541    /// }
542    /// ```
543    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(&params) })
556            })
557            .collect()
558    }
559
560    /// 异步触发事件并等待所有监听器完成(Rust 特有扩展,对齐 think-swoole 异步事件分发)
561    ///
562    /// 等价于 `trigger_spawn` + 逐个 `.await`,返回所有监听器的结果列表。
563    /// 监听器错误被收集到 `Vec` 中(不传播),包括 task panic 产生的 `JoinError`。
564    ///
565    /// **注意**:必须在 tokio 运行时中调用。监听器并发执行,结果顺序与注册顺序一致
566    /// (因为 `trigger_spawn` 返回的 JoinHandle 列表保持注册顺序)。
567    ///
568    /// Rust:
569    /// ```ignore
570    /// # async fn example() {
571    /// let results = dispatcher.trigger_async("UserLogin", &json!({"user_id": 123})).await;
572    /// for result in &results {
573    ///     if let Err(e) = result {
574    ///         eprintln!("Listener error: {}", e);
575    ///     }
576    /// }
577    /// # }
578    /// ```
579    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    /// 触发事件(只获取一个有效返回值)(对齐 PHP `until($event, $params = null)`)
601    ///
602    /// 等价 `trigger(event, params, true)`,返回最后一个非 null 返回值。
603    pub fn until(&self, event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
604        self.trigger(event, params, true)
605    }
606
607    /// 获取事件的所有监听器数量(PHP 无对应 API,Rust 扩展用于测试)
608    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
618/// 事件 facade(对齐 PHP `think\facade\Event`)
619///
620/// PHP `think\facade\Event` 是静态外观,委托到容器中的 `think\Event` 实例:
621/// ```php
622/// Event::listen('UserLogin', function($params) { ... });
623/// Event::trigger('UserLogin', ['user_id' => 123]);
624/// ```
625///
626/// Rust 端用全局 `once_cell::sync::Lazy<EventDispatcher>` 模拟:
627/// ```ignore
628/// event::facade::listen("UserLogin", Arc::new(ClosureListener::new(|_| Ok(Value::Null))), false);
629/// event::facade::trigger("UserLogin", &json!({"user_id": 123}), false).unwrap();
630/// ```
631pub mod facade {
632    use super::*;
633    use std::sync::OnceLock;
634
635    /// 全局事件分发器单例(对齐 PHP 容器中的 `think\Event` 实例)
636    static GLOBAL_DISPATCHER: OnceLock<EventDispatcher> = OnceLock::new();
637
638    /// 获取全局事件分发器(对齐 PHP `Facade::getFacadeClass()` 返回 `'event'`)
639    pub fn dispatcher() -> &'static EventDispatcher {
640        GLOBAL_DISPATCHER.get_or_init(EventDispatcher::new)
641    }
642
643    /// 注册事件监听(对齐 PHP `Event::listen(...)`)
644    pub fn listen(event: &str, listener: Arc<dyn Listener>, first: bool) {
645        dispatcher().listen(event, listener, first);
646    }
647
648    /// 批量注册事件监听(对齐 PHP `Event::listenEvents(...)`)
649    pub fn listen_events(events: Vec<(String, Vec<Arc<dyn Listener>>)>) {
650        dispatcher().listen_events(events);
651    }
652
653    /// 是否存在事件监听(对齐 PHP `Event::hasListener(...)`)
654    pub fn has_listener(event: &str) -> bool {
655        dispatcher().has_listener(event)
656    }
657
658    /// 移除事件监听(对齐 PHP `Event::remove(...)`)
659    pub fn remove(event: &str) {
660        dispatcher().remove(event);
661    }
662
663    /// 指定事件别名(对齐 PHP `Event::bind(...)`)
664    pub fn bind(events: Vec<(String, String)>) {
665        dispatcher().bind(events);
666    }
667
668    /// 注册事件订阅者(对齐 PHP `Event::subscribe(...)`)
669    pub fn subscribe(subscriber: Arc<dyn Subscriber>) {
670        dispatcher().subscribe(subscriber);
671    }
672
673    /// 自动注册事件观察者(对齐 PHP `Event::observe(...)`)
674    pub fn observe(observer: Arc<dyn Observer>, prefix: &str) {
675        dispatcher().observe(observer, prefix);
676    }
677
678    /// 触发事件(对齐 PHP `Event::trigger(...)`)
679    pub fn trigger(event: &str, params: &Value, once: bool) -> Result<Vec<Value>, EventError> {
680        dispatcher().trigger(event, params, once)
681    }
682
683    /// 触发事件(只获取一个有效返回值)(对齐 PHP `Event::until(...)`)
684    pub fn until(event: &str, params: &Value) -> Result<Vec<Value>, EventError> {
685        dispatcher().until(event, params)
686    }
687
688    /// 异步触发事件 — fire-and-forget(Rust 特有扩展)
689    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    /// 异步触发事件并等待所有监听器完成(Rust 特有扩展)
697    pub async fn trigger_async(event: &str, params: &Value) -> Vec<Result<Value, EventError>> {
698        dispatcher().trigger_async(event, params).await
699    }
700
701    /// 重置全局分发器(仅用于测试,PHP 无对应 API)
702    ///
703    /// **注意**:`OnceLock` 一旦初始化无法重置。测试中用独立 `EventDispatcher::new()` 实例
704    /// 而非全局 facade,避免测试间状态污染。
705    #[cfg(test)]
706    pub fn _reset_for_test() {
707        // OnceLock 无法重置,测试用独立实例
708        // 此函数仅作为文档说明,实际不执行任何操作
709    }
710}
711
712/// 触发事件助手函数(对齐 PHP `event($event, $args = null)`)
713///
714/// PHP `helper.php` 定义全局函数:
715/// ```php
716/// function event($event, $args = null)
717/// {
718///     return Event::trigger($event, $args);
719/// }
720/// ```
721///
722/// Rust 端用 `event_trigger` 函数模拟:
723/// ```ignore
724/// use sz_rust_core::event::event_trigger;
725///
726/// event_trigger("UserLogin", &json!({"user_id": 123}), false).unwrap();
727/// ```
728pub fn event_trigger(event: &str, params: &Value) -> Vec<Value> {
729    facade::trigger(event, params, false).unwrap_or_default()
730}
731
732/// 异步触发事件助手函数(Rust 特有扩展,对齐 think-swoole 异步事件分发)
733///
734/// 等价于 `facade::trigger_async(event, params).await`,返回所有监听器的结果列表。
735/// 监听器错误被收集到 `Vec` 中(不传播)。
736///
737/// Rust:
738/// ```ignore
739/// use sz_rust_core::event::event_trigger_async;
740///
741/// # async fn example() {
742/// let results = event_trigger_async("UserLogin", &json!({"user_id": 123})).await;
743/// # }
744/// ```
745pub 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    // ========================================================================
756    // 测试组 40: EventDispatcher 基础 API(listen/has_listener/remove)
757    // ========================================================================
758
759    #[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        // 对齐 PHP `listen(event, listener, true)` → array_unshift 插入队首
789        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, // 插入队首,应先执行
810        );
811
812        dispatcher.trigger("Test", &Value::Null, false).unwrap();
813        // order2 先执行(队首),order1 后执行(队尾)
814        assert_eq!(call_order.load(Ordering::SeqCst), 1);
815    }
816
817    #[test]
818    fn test_event_listen_events_batch() {
819        // 对齐 PHP `listenEvents(array $events)`
820        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        // 对齐 PHP `bind(['AppInit' => 'event\AppInit::class'])`
841        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        // 实际注册到 "app\\event\\AppInit"
853        assert!(!dispatcher.has_listener("AppInit_alias_check"));
854        assert!(dispatcher.has_listener("app\\event\\AppInit"));
855    }
856
857    // ========================================================================
858    // 测试组 41: EventDispatcher trigger 行为
859    // ========================================================================
860
861    #[test]
862    fn test_event_trigger_returns_all_results() {
863        // 对齐 PHP `trigger` 返回所有监听器返回值数组
864        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        // 对齐 PHP `once=true`:监听器返回非 null 值停止后续
883        let dispatcher = EventDispatcher::new();
884        dispatcher.listen(
885            "Test",
886            Arc::new(ClosureListener::new(|_| Ok(Value::Null))), // null 不停止
887            false,
888        );
889        dispatcher.listen(
890            "Test",
891            Arc::new(ClosureListener::new(|_| Ok(json!("stop")))), // 非 null 停止
892            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        // once=true:第二个返回 "stop" 停止,第三个不执行
907        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); // 第三个未执行
911    }
912
913    #[test]
914    fn test_event_trigger_false_stops_execution() {
915        // 对齐 PHP `if (false === $result[$key]) break`
916        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 停止
922            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); // 第二个未执行
938    }
939
940    #[test]
941    fn test_event_trigger_empty_event() {
942        // 无监听器的事件触发返回空数组
943        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        // 对齐 PHP `trigger($event, $params)` 传递参数
953        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", &params, false).unwrap();
968
969        assert_eq!(*received.lock().unwrap(), params);
970    }
971
972    // ========================================================================
973    // 测试组 42: EventDispatcher 点号通配
974    // ========================================================================
975
976    #[test]
977    fn test_event_dot_wildcard() {
978        // 对齐 PHP `if (strpos($event, '.'))` → 触发 `prefix.*`
979        let dispatcher = EventDispatcher::new();
980        let executed = Arc::new(AtomicUsize::new(0));
981
982        // 注册 User.login 和 User.* 监听器
983        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        // 触发 User.login 应同时触发 User.* 监听器
999        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        // 无通配监听器时,点号事件只触发具体监听器
1011        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    // ========================================================================
1026    // 测试组 43: EventDispatcher 去重
1027    // ========================================================================
1028
1029    #[test]
1030    fn test_event_dedup_same_listener_instance() {
1031        // 对齐 PHP `array_unique($listeners, SORT_REGULAR)`
1032        let dispatcher = EventDispatcher::new();
1033        let listener: Arc<dyn Listener> = Arc::new(ClosureListener::new(|_| Ok(json!(1))));
1034
1035        // 同一 Arc 实例注册两次
1036        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        // 去重后只执行一次
1041        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        // 不同监听器实例不去重
1048        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    // ========================================================================
1065    // 测试组 44: Subscriber 订阅者
1066    // ========================================================================
1067
1068    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        // 对齐 PHP `subscribe($subscriber)` → 调用 `$subscriber->subscribe($this)`
1099        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); // 1 + 10
1119    }
1120
1121    // ========================================================================
1122    // 测试组 45: Observer 观察者
1123    // ========================================================================
1124
1125    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        // 对齐 PHP `observe($observer)` → 反射 onXxx 方法自动注册
1155        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        // 对齐 PHP `observe($observer, 'User')` → 注册 `UserLogin` 事件
1176        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    // ========================================================================
1194    // 测试组 46: until 方法
1195    // ========================================================================
1196
1197    #[test]
1198    fn test_event_until_returns_first_non_null() {
1199        // 对齐 PHP `until($event, $params)` = `trigger($event, $params, true)`
1200        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        // until=true:第一个 null 不停止,第二个 "first_valid" 停止
1214        assert_eq!(results.len(), 2);
1215        assert_eq!(results[1], json!("first_valid"));
1216    }
1217
1218    // ========================================================================
1219    // 测试组 47: ClosureListener
1220    // ========================================================================
1221
1222    #[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    // ========================================================================
1241    // 测试组 48: 自定义 Listener trait 实现
1242    // ========================================================================
1243
1244    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        // 对齐 PHP 监听器抛出异常 → 停止后续执行
1266        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); // 第二个未执行
1289    }
1290
1291    // ========================================================================
1292    // 测试组 49: R5 PHP 行为对齐
1293    // ========================================================================
1294
1295    #[test]
1296    fn test_r5_php_event_listen_then_trigger() {
1297        // R5: 对齐 PHP Event::listen + Event::trigger
1298        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", &params, false).unwrap();
1313
1314        assert_eq!(*received.lock().unwrap(), params);
1315    }
1316
1317    #[test]
1318    fn test_r5_php_event_bind_alias_resolution() {
1319        // R5: 对齐 PHP Event::bind 别名解析
1320        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        // 实际注册到 "think\\event\\AppInit"
1333        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        // R5: 对齐 PHP `array_unshift` — first=true 插入队首
1343        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, // 插入队首
1364        );
1365
1366        dispatcher.trigger("Test", &Value::Null, false).unwrap();
1367        // o2 先执行(队首),o1 后执行(队尾)
1368        assert_eq!(*order.lock().unwrap(), vec![2, 1]);
1369    }
1370
1371    #[test]
1372    fn test_r5_php_event_trigger_returns_array() {
1373        // R5: 对齐 PHP `trigger` 返回所有返回值数组
1374        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        // R5: 对齐 PHP `until` = `trigger(once=true)` 返回最后一个非 null
1398        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        // 第一个 null 不停止,第二个 "first" 非空停止
1412        assert_eq!(results.len(), 2);
1413    }
1414
1415    #[test]
1416    fn test_r5_php_event_false_stops() {
1417        // R5: 对齐 PHP `if (false === $result[$key]) break`
1418        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        // R5: 对齐 PHP `array_merge($listeners, $this->listener[$prefix . '.*'])`
1444        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        // R5: 对齐 PHP `array_unique($listeners, SORT_REGULAR)`
1465        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        // 同一实例去重,只执行一次
1474        assert_eq!(results.len(), 1);
1475    }
1476
1477    #[test]
1478    fn test_r5_php_event_remove_clears_listeners() {
1479        // R5: 对齐 PHP `unset($this->listener[$event])`
1480        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        // R5: 对齐 PHP `hasListener` 应用 bind 别名
1498        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        // 查询别名应返回 true(PHP `hasListener('AppInit')` 解析到 `app\\event\\AppInit`)
1510        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        // R5: 对齐 PHP `listenEvents` 的 `array_merge`
1517        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        // array_merge:原有 [1] + 新增 [2, 3]
1533        assert_eq!(results, vec![json!(1), json!(2), json!(3)]);
1534    }
1535
1536    // ========================================================================
1537    // 测试组 50: trigger_spawn 基础(异步 fire-and-forget 分发)
1538    // ========================================================================
1539
1540    #[tokio::test]
1541    async fn test_trigger_spawn_returns_join_handles() {
1542        // trigger_spawn 返回 Vec<JoinHandle>,长度等于监听器数量
1543        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        // 等待所有完成
1559        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        // 无监听器的事件返回空 Vec
1568        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        // 参数透传到异步监听器
1576        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", &params);
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        // 所有监听器在异步模式都执行
1601        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    // ========================================================================
1625    // 测试组 51: trigger_async(异步分发 + 等待所有完成)
1626    // ========================================================================
1627
1628    #[tokio::test]
1629    async fn test_trigger_async_awaits_all() {
1630        // trigger_async 等待所有监听器完成,返回正确数量的结果
1631        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        // 结果顺序与监听器注册顺序一致(JoinHandle 列表保持注册顺序)
1656        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        // 监听器错误被收集到 Vec 中,不传播
1682        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        // 监听器 panic 被 tokio 捕获为 JoinError,trigger_async 转为 EventError
1711        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()); // panic → JoinError → EventError
1727    }
1728
1729    // ========================================================================
1730    // 测试组 52: 异步分发行为对齐(bind 别名 / 点号通配 / 去重)
1731    // ========================================================================
1732
1733    #[tokio::test]
1734    async fn test_trigger_spawn_applies_bind_alias() {
1735        // 异步模式应用 bind 别名
1736        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        // 异步模式应用点号通配
1757        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); // specific + wildcard
1771
1772        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        // 异步模式应用去重
1780        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); // 去重后只剩 1 个
1789    }
1790
1791    #[tokio::test]
1792    async fn test_trigger_spawn_correct_handle_count() {
1793        // 异步模式监听器数量正确(含 first 优先级插入)
1794        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, // 插入队首
1804        );
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        // 验证顺序:first 插入队首 → [2, 1, 3]
1815        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    // ========================================================================
1822    // 测试组 53: facade + helper 异步
1823    // ========================================================================
1824
1825    #[tokio::test]
1826    async fn test_facade_trigger_spawn() {
1827        // facade::trigger_spawn 异步分发
1828        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        // event_trigger_async 助手函数异步分发
1844        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}