Skip to main content

anycms_event/
registry.rs

1//! 事件注册表模块,提供事件类型的注册、查询和发现能力。
2//!
3//! 通过事件注册表,系统管理功能可以:
4//! - 发现系统中所有可用的事件类型
5//! - 查询事件的元数据(schema、描述、来源模块等)
6//! - 按条件搜索和过滤事件
7//! - 了解事件的发布/订阅状态
8
9use std::collections::HashMap;
10use std::sync::RwLock;
11use std::time::SystemTime;
12
13use serde::{Deserialize, Serialize};
14
15/// 注册事件的描述信息。
16///
17/// 包含事件的元数据,用于事件发现和管理。
18#[derive(Clone, Debug, Serialize, Deserialize)]
19pub struct EventDescriptor {
20    /// 事件唯一名称(如 "user.created")。
21    pub event_name: String,
22    /// 事件所属主题。
23    pub topic: String,
24    /// 事件描述。
25    #[serde(default)]
26    pub description: String,
27    /// 事件 payload 的 JSON Schema。
28    pub schema: Option<serde_json::Value>,
29    /// 注册此事件的来源模块。
30    pub source_module: Option<String>,
31    /// 事件标签,用于分类和过滤。
32    #[serde(default)]
33    pub tags: Vec<String>,
34    /// 事件首次注册的时间。
35    pub registered_at: SystemTime,
36    /// 事件累计发布次数。
37    #[serde(default)]
38    pub publish_count: u64,
39    /// 当前活跃订阅者数量。
40    #[serde(default)]
41    pub subscriber_count: usize,
42}
43
44/// 事件查询过滤器。
45#[derive(Clone, Debug, Default, Serialize, Deserialize)]
46pub struct EventQuery {
47    /// 按事件名称过滤(支持精确匹配或 `*` 前缀匹配)。
48    pub name: Option<String>,
49    /// 按主题过滤。
50    pub topic: Option<String>,
51    /// 按来源模块过滤。
52    pub source_module: Option<String>,
53    /// 按标签过滤(匹配任意一个)。
54    #[serde(default)]
55    pub tags: Vec<String>,
56    /// 文本搜索(在事件名称和描述中搜索)。
57    pub search: Option<String>,
58    /// 最大返回数量。
59    pub limit: Option<usize>,
60    /// 分页偏移。
61    pub offset: Option<usize>,
62}
63
64/// 事件注册表,跟踪已注册的事件类型。
65///
66/// 提供事件发现和查询能力,支持系统管理功能。
67pub struct EventRegistry {
68    events: RwLock<HashMap<String, EventDescriptor>>,
69}
70
71impl EventRegistry {
72    /// 创建一个空的事件注册表。
73    pub fn new() -> Self {
74        Self {
75            events: RwLock::new(HashMap::new()),
76        }
77    }
78
79    /// 注册一个事件描述符。如果已存在则更新。
80    pub fn register(&self, descriptor: EventDescriptor) {
81        let mut events = self.events.write().unwrap();
82        events.insert(descriptor.event_name.clone(), descriptor);
83    }
84
85    /// 使用基本信息注册事件。
86    pub fn register_simple(&self, event_name: &str, topic: &str) {
87        let descriptor = EventDescriptor {
88            event_name: event_name.to_string(),
89            topic: topic.to_string(),
90            description: String::new(),
91            schema: None,
92            source_module: None,
93            tags: Vec::new(),
94            registered_at: SystemTime::now(),
95            publish_count: 0,
96            subscriber_count: 0,
97        };
98        self.register(descriptor);
99    }
100
101    /// 注销一个事件类型。返回是否成功移除。
102    pub fn unregister(&self, event_name: &str) -> bool {
103        let mut events = self.events.write().unwrap();
104        events.remove(event_name).is_some()
105    }
106
107    /// 获取指定事件的描述符。
108    pub fn get(&self, event_name: &str) -> Option<EventDescriptor> {
109        let events = self.events.read().unwrap();
110        events.get(event_name).cloned()
111    }
112
113    /// 检查事件是否已注册。
114    pub fn contains(&self, event_name: &str) -> bool {
115        let events = self.events.read().unwrap();
116        events.contains_key(event_name)
117    }
118
119    /// 列出所有已注册的事件描述符。
120    pub fn list_all(&self) -> Vec<EventDescriptor> {
121        let events = self.events.read().unwrap();
122        events.values().cloned().collect()
123    }
124
125    /// 仅列出事件名称(轻量级)。
126    pub fn list_names(&self) -> Vec<String> {
127        let events = self.events.read().unwrap();
128        events.keys().cloned().collect()
129    }
130
131    /// 按条件查询事件。
132    ///
133    /// 支持按名称、主题、来源模块、标签过滤,以及文本搜索。
134    /// 支持分页(offset/limit)。
135    pub fn query(&self, query: EventQuery) -> Vec<EventDescriptor> {
136        let events = self.events.read().unwrap();
137        let mut results: Vec<EventDescriptor> = events
138            .values()
139            .filter(|e| {
140                // 名称过滤
141                if let Some(ref name) = query.name {
142                    if name.ends_with('*') {
143                        let prefix = &name[..name.len() - 1];
144                        if !e.event_name.starts_with(prefix) {
145                            return false;
146                        }
147                    } else if &e.event_name != name {
148                        return false;
149                    }
150                }
151                // 主题过滤
152                if let Some(ref topic) = query.topic {
153                    if &e.topic != topic {
154                        return false;
155                    }
156                }
157                // 来源模块过滤
158                if let Some(ref module) = query.source_module {
159                    if e.source_module.as_ref() != Some(module) {
160                        return false;
161                    }
162                }
163                // 标签过滤
164                if !query.tags.is_empty() {
165                    if !query.tags.iter().any(|t| e.tags.contains(t)) {
166                        return false;
167                    }
168                }
169                // 文本搜索
170                if let Some(ref search) = query.search {
171                    let search_lower = search.to_lowercase();
172                    if !e.event_name.to_lowercase().contains(&search_lower)
173                        && !e.description.to_lowercase().contains(&search_lower)
174                    {
175                        return false;
176                    }
177                }
178                true
179            })
180            .cloned()
181            .collect();
182
183        results.sort_by(|a, b| a.event_name.cmp(&b.event_name));
184
185        let offset = query.offset.unwrap_or(0);
186        let limit = query.limit.unwrap_or(usize::MAX);
187        results.into_iter().skip(offset).take(limit).collect()
188    }
189
190    /// 增加事件的发布计数。
191    pub fn increment_publish_count(&self, event_name: &str) {
192        let mut events = self.events.write().unwrap();
193        if let Some(desc) = events.get_mut(event_name) {
194            desc.publish_count += 1;
195        }
196    }
197
198    /// 更新事件的订阅者数量。
199    pub fn set_subscriber_count(&self, event_name: &str, count: usize) {
200        let mut events = self.events.write().unwrap();
201        if let Some(desc) = events.get_mut(event_name) {
202            desc.subscriber_count = count;
203        }
204    }
205
206    /// 获取已注册事件总数。
207    pub fn count(&self) -> usize {
208        let events = self.events.read().unwrap();
209        events.len()
210    }
211
212    /// 清空所有已注册事件。
213    pub fn clear(&self) {
214        let mut events = self.events.write().unwrap();
215        events.clear();
216    }
217}
218
219impl Default for EventRegistry {
220    fn default() -> Self {
221        Self::new()
222    }
223}
224
225#[cfg(test)]
226mod tests {
227    use super::*;
228
229    fn make_descriptor(name: &str, topic: &str, desc: &str) -> EventDescriptor {
230        EventDescriptor {
231            event_name: name.to_string(),
232            topic: topic.to_string(),
233            description: desc.to_string(),
234            schema: None,
235            source_module: None,
236            tags: Vec::new(),
237            registered_at: SystemTime::now(),
238            publish_count: 0,
239            subscriber_count: 0,
240        }
241    }
242
243    #[test]
244    fn test_register_and_get() {
245        let registry = EventRegistry::new();
246        let desc = make_descriptor("user.created", "user", "User created event");
247        registry.register(desc);
248
249        let got = registry.get("user.created").unwrap();
250        assert_eq!(got.event_name, "user.created");
251        assert_eq!(got.topic, "user");
252        assert_eq!(got.description, "User created event");
253    }
254
255    #[test]
256    fn test_register_simple() {
257        let registry = EventRegistry::new();
258        registry.register_simple("order.placed", "order");
259
260        assert!(registry.contains("order.placed"));
261        let got = registry.get("order.placed").unwrap();
262        assert_eq!(got.topic, "order");
263    }
264
265    #[test]
266    fn test_unregister() {
267        let registry = EventRegistry::new();
268        registry.register_simple("user.deleted", "user");
269        assert!(registry.unregister("user.deleted"));
270        assert!(!registry.contains("user.deleted"));
271    }
272
273    #[test]
274    fn test_list_all() {
275        let registry = EventRegistry::new();
276        registry.register_simple("user.created", "user");
277        registry.register_simple("order.placed", "order");
278
279        let all = registry.list_all();
280        assert_eq!(all.len(), 2);
281    }
282
283    #[test]
284    fn test_list_names() {
285        let registry = EventRegistry::new();
286        registry.register_simple("user.created", "user");
287        registry.register_simple("order.placed", "order");
288
289        let names = registry.list_names();
290        assert_eq!(names.len(), 2);
291        assert!(names.contains(&"user.created".to_string()));
292    }
293
294    #[test]
295    fn test_query_by_topic() {
296        let registry = EventRegistry::new();
297        registry.register(make_descriptor("user.created", "user", ""));
298        registry.register(make_descriptor("user.deleted", "user", ""));
299        registry.register(make_descriptor("order.placed", "order", ""));
300
301        let results = registry.query(EventQuery {
302            topic: Some("user".to_string()),
303            ..Default::default()
304        });
305        assert_eq!(results.len(), 2);
306    }
307
308    #[test]
309    fn test_query_by_name_prefix() {
310        let registry = EventRegistry::new();
311        registry.register_simple("user.created", "user");
312        registry.register_simple("user.deleted", "user");
313        registry.register_simple("order.placed", "order");
314
315        let results = registry.query(EventQuery {
316            name: Some("user.*".to_string()),
317            ..Default::default()
318        });
319        assert_eq!(results.len(), 2);
320    }
321
322    #[test]
323    fn test_query_search() {
324        let registry = EventRegistry::new();
325        registry.register(make_descriptor("user.created", "user", "A new user account was created"));
326        registry.register(make_descriptor("user.deleted", "user", "User account was deleted"));
327        registry.register(make_descriptor("order.placed", "order", "A new order was placed"));
328
329        let results = registry.query(EventQuery {
330            search: Some("created".to_string()),
331            ..Default::default()
332        });
333        assert_eq!(results.len(), 1);
334        assert_eq!(results[0].event_name, "user.created");
335    }
336
337    #[test]
338    fn test_query_pagination() {
339        let registry = EventRegistry::new();
340        for i in 0..10 {
341            registry.register_simple(&format!("event.{}", i), "test");
342        }
343
344        let page1 = registry.query(EventQuery {
345            limit: Some(3),
346            offset: Some(0),
347            ..Default::default()
348        });
349        assert_eq!(page1.len(), 3);
350
351        let page2 = registry.query(EventQuery {
352            limit: Some(3),
353            offset: Some(3),
354            ..Default::default()
355        });
356        assert_eq!(page2.len(), 3);
357    }
358
359    #[test]
360    fn test_increment_publish_count() {
361        let registry = EventRegistry::new();
362        registry.register_simple("user.created", "user");
363
364        registry.increment_publish_count("user.created");
365        registry.increment_publish_count("user.created");
366
367        let desc = registry.get("user.created").unwrap();
368        assert_eq!(desc.publish_count, 2);
369    }
370
371    #[test]
372    fn test_tags_filter() {
373        let registry = EventRegistry::new();
374        let mut desc1 = make_descriptor("user.created", "user", "");
375        desc1.tags = vec!["auth".to_string(), "user".to_string()];
376        let mut desc2 = make_descriptor("order.placed", "order", "");
377        desc2.tags = vec!["commerce".to_string()];
378        registry.register(desc1);
379        registry.register(desc2);
380
381        let results = registry.query(EventQuery {
382            tags: vec!["auth".to_string()],
383            ..Default::default()
384        });
385        assert_eq!(results.len(), 1);
386        assert_eq!(results[0].event_name, "user.created");
387    }
388}