use crate::{match_event, EventBus}; use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::sync::Arc; #[derive(Clone, Debug, PartialEq)] struct User { name: String, } #[derive(Clone, Debug, PartialEq)] struct UdpMetric { value: f32, } // ── on (callback sync) ──────────────────────────────────────────────────── #[tokio::test] async fn test_on_callback_sync() { let bus = Arc::new(EventBus::new()); let received = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&received); bus.on::("user-connected", move |user| { if user.name == "Alice" { flag.store(true, Ordering::SeqCst); } }); bus.emit( "user-connected", User { name: "Alice".into(), }, ); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(received.load(Ordering::SeqCst)); } #[tokio::test] async fn test_on_targeted_wakeup() { // Émettre sur "user-connected" ne doit pas réveiller "udp-metrics-updated" let bus = Arc::new(EventBus::new()); let metric_called = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&metric_called); bus.on::("udp-metrics-updated", move |_| { flag.store(true, Ordering::SeqCst); }); bus.emit( "user-connected", User { name: "Carol".into(), }, ); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(!metric_called.load(Ordering::SeqCst)); } #[tokio::test] async fn test_on_type_mismatch_ignored() { // Émettre un UdpMetric sur un topic écouté en User → handler pas appelé let bus = Arc::new(EventBus::new()); let called = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&called); bus.on::("mixed-topic", move |_| { flag.store(true, Ordering::SeqCst); }); bus.emit("mixed-topic", UdpMetric { value: 1.0 }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(!called.load(Ordering::SeqCst)); } #[tokio::test] async fn test_on_multiple_subscribers_same_topic() { let bus = Arc::new(EventBus::new()); let count = Arc::new(AtomicU32::new(0)); for _ in 0..3 { let c = Arc::clone(&count); bus.on::("user-connected", move |_| { c.fetch_add(1, Ordering::SeqCst); }); } bus.emit( "user-connected", User { name: "Grace".into(), }, ); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert_eq!(count.load(Ordering::SeqCst), 3); } // ── on_async ────────────────────────────────────────────────────────────── #[tokio::test] async fn test_on_async_callback() { let bus = Arc::new(EventBus::new()); let received = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&received); bus.on_async::("user-connected", move |user| { let f = Arc::clone(&flag); async move { if user.name == "Async" { f.store(true, Ordering::SeqCst); } } }); bus.emit( "user-connected", User { name: "Async".into(), }, ); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(received.load(Ordering::SeqCst)); } // ── on_raw + match_event! ───────────────────────────────────────────────── #[tokio::test] async fn test_on_raw_and_match_event_macro() { let bus = Arc::new(EventBus::new()); let mut rx = bus.on_raw("user-connected"); bus.emit( "user-connected", User { name: "Frank".into(), }, ); let evt = rx.recv().await.unwrap(); let mut received_name = String::new(); match_event!(evt, User => |u: User| { received_name = u.name.clone(); }, UdpMetric => |_m: UdpMetric| { panic!("mauvais type"); } ); assert_eq!(received_name, "Frank"); } // ── Utilitaires ──────────────────────────────────────────────────────────── #[test] fn test_topics_list() { let bus = EventBus::new(); // on_raw enregistre le canal (get_or_create) let _rx1 = bus.on_raw("user-connected"); let _rx2 = bus.on_raw("udp-metrics-updated"); let mut topics = bus.topics(); topics.sort(); assert_eq!(topics, vec!["udp-metrics-updated", "user-connected"]); } #[tokio::test] async fn test_emit_multiple_types_same_bus() { let bus = Arc::new(EventBus::new()); let user_ok = Arc::new(AtomicBool::new(false)); let metric_ok = Arc::new(AtomicBool::new(false)); let u = Arc::clone(&user_ok); let m = Arc::clone(&metric_ok); bus.on::("user-connected", move |user| { if user.name == "Bob" { u.store(true, Ordering::SeqCst); } }); bus.on::("udp-metrics-updated", move |metric| { if (metric.value - 3.14).abs() < 0.001 { m.store(true, Ordering::SeqCst); } }); bus.emit("user-connected", User { name: "Bob".into() }); bus.emit("udp-metrics-updated", UdpMetric { value: 3.14 }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(user_ok.load(Ordering::SeqCst)); assert!(metric_ok.load(Ordering::SeqCst)); }