use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use uuid::Uuid; use crate::EventBus; #[derive(Clone, Debug, PartialEq)] struct MessageCreatedEvent { channel_id: Uuid, content: String, } #[derive(Clone, Debug, PartialEq)] struct MessageUpdatedEvent { id: u64, content: String, } #[derive(Clone, Debug, PartialEq)] struct MessageDeletedEvent { id: u64, } #[derive(Clone, Debug, PartialEq)] struct UdpMetricEvent { value: f32, } // ── Sync Callbacks ────────────────────────────────────────────────────────── #[tokio::test] async fn test_on_callback_sync() { let bus = EventBus::new(); let received = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&received); bus.on::(move |event| { if event.content == "Hello" { flag.store(true, Ordering::SeqCst); } }); bus.emit(MessageCreatedEvent { channel_id: Uuid::new_v4(), content: "Hello".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(received.load(Ordering::SeqCst)); } // ── Async Callbacks (Exact User Requirement) ──────────────────────────────── #[tokio::test] async fn test_on_async_callback_turbofish() { let bus = EventBus::new(); let received_content = Arc::new(tokio::sync::Mutex::new(String::new())); let rc = Arc::clone(&received_content); // Exact syntax specified by the user: // event_bus.on_async::(|event| async move { ... }); bus.on_async::(move |event| { let rc = Arc::clone(&rc); async move { let mut lock = rc.lock().await; *lock = event.content; } }); // Exact syntax specified by the user: // event_bus.emit(MessageUpdatedEvent { ... }); bus.emit(MessageUpdatedEvent { id: 42, content: "Updated message content".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; let result = received_content.lock().await.clone(); assert_eq!(result, "Updated message content"); } #[tokio::test] async fn test_on_async_callback_type_inferred() { let bus = EventBus::new(); let flag = Arc::new(AtomicBool::new(false)); let f = Arc::clone(&flag); // Also supports inferring the event type from closure parameter bus.on_async(move |event: MessageUpdatedEvent| { let f = Arc::clone(&f); async move { if event.id == 99 { f.store(true, Ordering::SeqCst); } } }); bus.emit(MessageUpdatedEvent { id: 99, content: "Inferred".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(flag.load(Ordering::SeqCst)); } // ── Targeted Wake-Up & Type Isolation ─────────────────────────────────────── #[tokio::test] async fn test_on_targeted_wakeup() { let bus = EventBus::new(); let metric_called = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&metric_called); bus.on::(move |_| { flag.store(true, Ordering::SeqCst); }); // Emitting MessageCreatedEvent must never wake up UdpMetricEvent subscribers bus.emit(MessageCreatedEvent { channel_id: Uuid::new_v4(), content: "Ignore me".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(!metric_called.load(Ordering::SeqCst)); } #[tokio::test] async fn test_multiple_subscribers_same_type() { let bus = EventBus::new(); let count = Arc::new(AtomicU32::new(0)); for _ in 0..3 { let c = Arc::clone(&count); bus.on::(move |_| { c.fetch_add(1, Ordering::SeqCst); }); } bus.emit(MessageCreatedEvent { channel_id: Uuid::new_v4(), content: "Broadcast".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert_eq!(count.load(Ordering::SeqCst), 3); } #[tokio::test] async fn test_multiple_different_types_same_bus() { let bus = EventBus::new(); let msg_ok = Arc::new(AtomicBool::new(false)); let metric_ok = Arc::new(AtomicBool::new(false)); let m_flag = Arc::clone(&msg_ok); let u_flag = Arc::clone(&metric_ok); bus.on::(move |event| { if event.content == "Test" { m_flag.store(true, Ordering::SeqCst); } }); bus.on::(move |metric| { if (metric.value - 42.5).abs() < 0.001 { u_flag.store(true, Ordering::SeqCst); } }); bus.emit(MessageCreatedEvent { channel_id: Uuid::new_v4(), content: "Test".into(), }); bus.emit(UdpMetricEvent { value: 42.5 }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(msg_ok.load(Ordering::SeqCst)); assert!(metric_ok.load(Ordering::SeqCst)); } // ── Direct Stream / Receiver (No match_event! needed) ─────────────────────── #[tokio::test] async fn test_subscribe_direct_typed_receiver() { let bus = EventBus::new(); let mut rx = bus.subscribe::(); bus.emit(MessageUpdatedEvent { id: 123, content: "Direct typed".into(), }); let event = rx.recv().await.expect("failed to receive event"); // event is directly of type MessageUpdatedEvent, no downcast needed! assert_eq!(event.id, 123); assert_eq!(event.content, "Direct typed"); } // ── In-Handler Filtering (Direct Field Access) ────────────────────────────── #[tokio::test] async fn test_filter_by_field_in_subscriber() { let bus = EventBus::new(); let channel_a = Uuid::new_v4(); let channel_b = Uuid::new_v4(); let count_a = Arc::new(AtomicU32::new(0)); let count_b = Arc::new(AtomicU32::new(0)); let count_global = Arc::new(AtomicU32::new(0)); let ca = Arc::clone(&count_a); bus.on_async::(move |event| { let ca = Arc::clone(&ca); async move { if event.channel_id == channel_a { ca.fetch_add(1, Ordering::SeqCst); } } }); let cb = Arc::clone(&count_b); bus.on_async::(move |event| { let cb = Arc::clone(&cb); async move { if event.channel_id == channel_b { cb.fetch_add(1, Ordering::SeqCst); } } }); let cg = Arc::clone(&count_global); bus.on::(move |_| { cg.fetch_add(1, Ordering::SeqCst); }); // Emit event with channel_a bus.emit(MessageCreatedEvent { channel_id: channel_a, content: "For A".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; // Both channel_a handler and global handler processed it, but not channel_b assert_eq!(count_a.load(Ordering::SeqCst), 1); assert_eq!(count_b.load(Ordering::SeqCst), 0); assert_eq!(count_global.load(Ordering::SeqCst), 1); } // ── Async with Context ────────────────────────────────────────────────────── #[tokio::test] async fn test_on_async_with_context() { let bus = EventBus::new(); let prefix = Arc::new("Prefix: ".to_string()); let result = Arc::new(tokio::sync::Mutex::new(String::new())); let r = Arc::clone(&result); bus.on_async_with::(prefix, move |ctx, event| { let r = Arc::clone(&r); async move { let mut lock = r.lock().await; *lock = format!("{}{}", ctx, event.content); } }); bus.emit(MessageCreatedEvent { channel_id: Uuid::new_v4(), content: "Hello Context".into(), }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; let final_str = result.lock().await.clone(); assert_eq!(final_str, "Prefix: Hello Context"); } // ── Metrics, Utilities & Cleanup ──────────────────────────────────────────── #[test] fn test_subscriber_count_and_clear() { let bus = EventBus::new(); assert_eq!(bus.subscriber_count::(), 0); assert!(!bus.has_subscribers::()); let _sub = bus.subscribe::(); assert_eq!(bus.subscriber_count::(), 1); assert!(bus.has_subscribers::()); assert_eq!(bus.channel_count(), 1); bus.clear_type::(); assert_eq!(bus.subscriber_count::(), 0); assert_eq!(bus.channel_count(), 0); } #[tokio::test] async fn test_subscription_abort() { let bus = EventBus::new(); let count = Arc::new(AtomicU32::new(0)); let c = Arc::clone(&count); let handle = bus.on::(move |_| { c.fetch_add(1, Ordering::SeqCst); }); bus.emit(MessageDeletedEvent { id: 1 }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert_eq!(count.load(Ordering::SeqCst), 1); // Cancel the subscription handle.abort(); tokio::time::sleep(std::time::Duration::from_millis(10)).await; bus.emit(MessageDeletedEvent { id: 2 }); tokio::time::sleep(std::time::Duration::from_millis(20)).await; // Count should not increase after abort assert_eq!(count.load(Ordering::SeqCst), 1); } // ── One-Shot Listeners (wait_next & wait_for) ──────────────────────────────── #[tokio::test] async fn test_wait_next() { let bus = EventBus::new(); let b = bus.clone(); tokio::spawn(async move { tokio::time::sleep(std::time::Duration::from_millis(10)).await; b.emit(MessageUpdatedEvent { id: 777, content: "Next event".into(), }); }); let event = bus.wait_next::().await.unwrap(); assert_eq!(event.id, 777); assert_eq!(event.content, "Next event"); } #[tokio::test] async fn test_wait_for() { let bus = EventBus::new(); let b = bus.clone(); tokio::spawn(async move { tokio::time::sleep(std::time::Duration::from_millis(10)).await; b.emit(MessageUpdatedEvent { id: 1, content: "Ignore".into(), }); tokio::time::sleep(std::time::Duration::from_millis(10)).await; b.emit(MessageUpdatedEvent { id: 2, content: "Target".into(), }); }); let event = bus .wait_for::(|e| e.id == 2) .await .unwrap(); assert_eq!(event.id, 2); assert_eq!(event.content, "Target"); }