1use std::collections::HashMap;
10use std::sync::RwLock;
11use std::time::SystemTime;
12
13use serde::{Deserialize, Serialize};
14
15#[derive(Clone, Debug, Serialize, Deserialize)]
19pub struct EventDescriptor {
20 pub event_name: String,
22 pub topic: String,
24 #[serde(default)]
26 pub description: String,
27 pub schema: Option<serde_json::Value>,
29 pub source_module: Option<String>,
31 #[serde(default)]
33 pub tags: Vec<String>,
34 pub registered_at: SystemTime,
36 #[serde(default)]
38 pub publish_count: u64,
39 #[serde(default)]
41 pub subscriber_count: usize,
42}
43
44#[derive(Clone, Debug, Default, Serialize, Deserialize)]
46pub struct EventQuery {
47 pub name: Option<String>,
49 pub topic: Option<String>,
51 pub source_module: Option<String>,
53 #[serde(default)]
55 pub tags: Vec<String>,
56 pub search: Option<String>,
58 pub limit: Option<usize>,
60 pub offset: Option<usize>,
62}
63
64pub struct EventRegistry {
68 events: RwLock<HashMap<String, EventDescriptor>>,
69}
70
71impl EventRegistry {
72 pub fn new() -> Self {
74 Self {
75 events: RwLock::new(HashMap::new()),
76 }
77 }
78
79 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 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 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 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 pub fn contains(&self, event_name: &str) -> bool {
115 let events = self.events.read().unwrap();
116 events.contains_key(event_name)
117 }
118
119 pub fn list_all(&self) -> Vec<EventDescriptor> {
121 let events = self.events.read().unwrap();
122 events.values().cloned().collect()
123 }
124
125 pub fn list_names(&self) -> Vec<String> {
127 let events = self.events.read().unwrap();
128 events.keys().cloned().collect()
129 }
130
131 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 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 if let Some(ref topic) = query.topic {
153 if &e.topic != topic {
154 return false;
155 }
156 }
157 if let Some(ref module) = query.source_module {
159 if e.source_module.as_ref() != Some(module) {
160 return false;
161 }
162 }
163 if !query.tags.is_empty() {
165 if !query.tags.iter().any(|t| e.tags.contains(t)) {
166 return false;
167 }
168 }
169 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 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 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 pub fn count(&self) -> usize {
208 let events = self.events.read().unwrap();
209 events.len()
210 }
211
212 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}