This commit is contained in:
2026-05-03 21:25:01 +02:00
parent 47c33a3a6c
commit 0de2e334ae
10 changed files with 879 additions and 1131 deletions
+293
View File
@@ -0,0 +1,293 @@
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, _>("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::<UdpMetric, _>("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::<User, _>("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, _>("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, _, _>("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_pattern ────────────────────────────────────────────────────────────
#[tokio::test]
async fn test_on_pattern_callback_receives_topic_and_payload() {
let bus = Arc::new(EventBus::new());
let created = Arc::new(AtomicBool::new(false));
let deleted = Arc::new(AtomicBool::new(false));
let c = Arc::clone(&created);
let d = Arc::clone(&deleted);
bus.on_pattern::<User, _>("user-*", move |topic, user| match topic.as_str() {
"user-created" if user.name == "Dave" => c.store(true, Ordering::SeqCst),
"user-deleted" if user.name == "Eve" => d.store(true, Ordering::SeqCst),
_ => {}
});
bus.emit(
"user-created",
User {
name: "Dave".into(),
},
);
bus.emit("user-deleted", User { name: "Eve".into() });
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert!(created.load(Ordering::SeqCst));
assert!(deleted.load(Ordering::SeqCst));
}
#[tokio::test]
async fn test_on_pattern_does_not_match_other_topics() {
let bus = Arc::new(EventBus::new());
let called = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&called);
bus.on_pattern::<User, _>("user-*", move |_, _| {
flag.store(true, Ordering::SeqCst);
});
bus.emit(
"server-created",
User {
name: "Ghost".into(),
},
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert!(!called.load(Ordering::SeqCst));
}
#[tokio::test]
async fn test_on_pattern_no_pre_registration_needed() {
// Le pattern est enregistré avant que le topic n'existe
let bus = Arc::new(EventBus::new());
let received = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&received);
bus.on_pattern::<User, _>("user-*", move |topic, user| {
if topic == "user-new-topic" && user.name == "Frank" {
flag.store(true, Ordering::SeqCst);
}
});
bus.emit(
"user-new-topic",
User {
name: "Frank".into(),
},
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert!(received.load(Ordering::SeqCst));
}
// ── on_pattern_async ──────────────────────────────────────────────────────
#[tokio::test]
async fn test_on_pattern_async_callback() {
let bus = Arc::new(EventBus::new());
let received = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&received);
bus.on_pattern_async::<User, _, _>("user-*", move |topic, user| {
let f = Arc::clone(&flag);
async move {
if topic == "user-created" && user.name == "Hank" {
f.store(true, Ordering::SeqCst);
}
}
});
bus.emit(
"user-created",
User {
name: "Hank".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, _>("user-connected", move |user| {
if user.name == "Bob" {
u.store(true, Ordering::SeqCst);
}
});
bus.on::<UdpMetric, _>("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));
}