From c96101ec3d088e5bc74ab52cc1defe850d76ee5f Mon Sep 17 00:00:00 2001 From: Nell Date: Mon, 13 Jul 2026 18:46:19 +0200 Subject: [PATCH] init --- Cargo.lock | 2 +- event_bus/Cargo.toml | 4 +- event_bus/benches/event_bus_throughput.rs | 134 ------------- event_bus/src/bus.rs | 231 +++++----------------- event_bus/src/lib.rs | 117 +---------- event_bus/src/tests.rs | 103 +--------- 6 files changed, 55 insertions(+), 536 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index ff67ddf..e9f1009 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1384,10 +1384,10 @@ name = "event_bus" version = "0.1.0" dependencies = [ "criterion", - "glob", "parking_lot", "tokio", "tracing", + "uuid", ] [[package]] diff --git a/event_bus/Cargo.toml b/event_bus/Cargo.toml index 90118c7..8cb0bbd 100644 --- a/event_bus/Cargo.toml +++ b/event_bus/Cargo.toml @@ -14,10 +14,10 @@ harness = false [dependencies] tokio = { version = "1.52.3", default-features = false, features = ["rt", "sync"] } -glob = "0.3.3" parking_lot = "0.12.5" tracing = "0.1" +uuid = { version = "1.23.5", features = ["v4"] } [dev-dependencies] tokio = { version = "1.52.3", default-features = false, features = ["rt", "rt-multi-thread", "macros", "time", "sync"] } -criterion = { version = "0.8.2", features = ["async_tokio"] } +criterion = { version = "0.8.2", features = ["async_tokio"] } \ No newline at end of file diff --git a/event_bus/benches/event_bus_throughput.rs b/event_bus/benches/event_bus_throughput.rs index 92cc0e7..d308759 100644 --- a/event_bus/benches/event_bus_throughput.rs +++ b/event_bus/benches/event_bus_throughput.rs @@ -9,7 +9,6 @@ use event_bus::EventBus; use tokio::runtime::Runtime; const TOPIC: &str = "bench-topic"; -const PATTERN: &str = "bench-*"; #[derive(Clone)] struct SmallEvent { @@ -473,82 +472,6 @@ fn bench_typed_callback(c: &mut Criterion) { group.finish(); } -fn bench_pattern_callback(c: &mut Criterion) { - let rt = runtime(); - - let mut group = c.benchmark_group("event_bus/pattern_callback"); - group.throughput(Throughput::Elements(1)); - - group.bench_function("small_struct", |b| { - b.to_async(&rt).iter_custom(|iters| async move { - let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024)); - - let received = Arc::new(AtomicU64::new(0)); - let handler_count = Arc::clone(&received); - - let subscription = bus.on_pattern::(PATTERN, move |topic, event| { - let _ = topic; - let _ = event.value; - handler_count.fetch_add(1, Ordering::Relaxed); - }); - - let start = Instant::now(); - - for i in 0..iters { - bus.emit(TOPIC, SmallEvent { value: i }); - } - - wait_until_received(&received, iters).await; - - let elapsed = start.elapsed(); - - subscription.abort(); - - elapsed - }); - }); - - group.bench_function("arc_payload_1kb", |b| { - b.to_async(&rt).iter_custom(|iters| async move { - let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024)); - let payload: Arc<[u8]> = Arc::from(vec![7_u8; 1024].into_boxed_slice()); - - let received = Arc::new(AtomicU64::new(0)); - let handler_count = Arc::clone(&received); - - let subscription = - bus.on_pattern::(PATTERN, move |topic, event| { - let _ = topic; - let _ = event.id; - let _ = event.payload.len(); - handler_count.fetch_add(1, Ordering::Relaxed); - }); - - let start = Instant::now(); - - for i in 0..iters { - bus.emit( - TOPIC, - ArcPayloadEvent { - id: i, - payload: Arc::clone(&payload), - }, - ); - } - - wait_until_received(&received, iters).await; - - let elapsed = start.elapsed(); - - subscription.abort(); - - elapsed - }); - }); - - group.finish(); -} - fn bench_multiple_subscribers(c: &mut Criterion) { let rt = runtime(); @@ -597,69 +520,12 @@ fn bench_multiple_subscribers(c: &mut Criterion) { group.finish(); } -fn bench_multiple_patterns(c: &mut Criterion) { - let rt = runtime(); - - let mut group = c.benchmark_group("event_bus/multiple_patterns"); - group.throughput(Throughput::Elements(1)); - - for pattern_count in [1_u64, 4, 16, 64, 256] { - group.bench_function(format!("{pattern_count}_patterns_one_match"), |b| { - b.to_async(&rt).iter_custom(|iters| async move { - let bus = Arc::new(EventBus::with_capacity(iters as usize + 1024)); - - let received = Arc::new(AtomicU64::new(0)); - let mut subscriptions = Vec::with_capacity(pattern_count as usize); - - for index in 0..pattern_count { - let pattern = if index == 0 { - PATTERN.to_string() - } else { - format!("unused-{index}-*") - }; - - let handler_count = Arc::clone(&received); - - let subscription = - bus.on_pattern::(&pattern, move |topic, event| { - let _ = topic; - let _ = event.value; - handler_count.fetch_add(1, Ordering::Relaxed); - }); - - subscriptions.push(subscription); - } - - let start = Instant::now(); - - for i in 0..iters { - bus.emit(TOPIC, SmallEvent { value: i }); - } - - wait_until_received(&received, iters).await; - - let elapsed = start.elapsed(); - - for subscription in subscriptions { - subscription.abort(); - } - - elapsed - }); - }); - } - - group.finish(); -} - criterion_group!( benches, bench_emit_no_subscriber, bench_raw_subscriber, bench_typed_callback, - bench_pattern_callback, bench_multiple_subscribers, - bench_multiple_patterns, ); criterion_main!(benches); diff --git a/event_bus/src/bus.rs b/event_bus/src/bus.rs index e584bc3..64a788b 100644 --- a/event_bus/src/bus.rs +++ b/event_bus/src/bus.rs @@ -2,12 +2,13 @@ use std::any::Any; use std::future::Future; use std::sync::Arc; -use glob::Pattern; use parking_lot::RwLock; use std::collections::HashMap; use tokio::sync::broadcast; use tokio::task::JoinHandle; +use tracing::log::kv::{Key, Value}; use tracing::{debug, trace, warn}; +use uuid::Uuid; /// Raw event type: an atomic reference-counted pointer to any value. pub type AnyEvent = Arc; @@ -15,6 +16,25 @@ pub type AnyEvent = Arc; /// Default buffer capacity for each broadcast channel. const DEFAULT_CAPACITY: usize = 64; +#[derive(Debug, Clone, PartialEq, Eq)] +enum ScopeValue { + String(String), + Uuid(Uuid), +} +impl ScopeValue { + fn into_string(self) -> String { + match self { + Self::String(value) => value, + Self::Uuid(value) => value.to_string(), + } + } +} +#[derive(Debug, Clone, PartialEq, Eq)] +struct Scope { + key: String, + value: ScopeValue, +} + /// The central event bus. /// /// Share it via `Arc` across modules. @@ -60,37 +80,10 @@ const DEFAULT_CAPACITY: usize = 64; /// # tokio::time::sleep(std::time::Duration::from_millis(10)).await; /// # }); /// ``` -/// -/// # Example — glob pattern with topic -/// ```rust,no_run -/// use std::sync::Arc; -/// use oxspeak_server_lib::event_bus::EventBus; -/// -/// #[derive(Clone, Debug)] -/// struct User { name: String } -/// -/// # tokio_test::block_on(async { -/// let bus = Arc::new(EventBus::new()); -/// -/// bus.on_pattern::("user-*", |topic, user| { -/// match topic.as_str() { -/// "user-created" => println!("Created : {:?}", user), -/// "user-deleted" => println!("Deleted : {:?}", user), -/// other => println!("{}: {:?}", other, user), -/// } -/// }); -/// -/// bus.emit("user-created", User { name: "Alice".into() }); -/// bus.emit("user-deleted", User { name: "Bob".into() }); -/// # tokio::time::sleep(std::time::Duration::from_millis(10)).await; -/// # }); -/// ``` #[derive(Debug)] pub struct EventBus { /// Channels indexed by exact topic. channels: RwLock>>, - /// Channels for glob-pattern subscriptions. - patterns: RwLock)>>, capacity: usize, } @@ -103,7 +96,6 @@ impl EventBus { ); Self { channels: RwLock::new(HashMap::new()), - patterns: RwLock::new(Vec::new()), capacity: DEFAULT_CAPACITY, } } @@ -113,7 +105,6 @@ impl EventBus { debug!("EventBus created with capacity {}", capacity); Self { channels: RwLock::new(HashMap::new()), - patterns: RwLock::new(Vec::new()), capacity, } } @@ -151,7 +142,6 @@ impl EventBus { /// Emits an event on a topic. /// /// - Pushes the event into the exact-topic channel (if subscribers exist). - /// - Pushes the event into all glob-pattern channels that match the topic. /// - If nobody is listening, the event is silently dropped. /// /// # Example @@ -167,7 +157,6 @@ impl EventBus { trace!(topic, "Emitting event"); let event: AnyEvent = Arc::new(event); - // Exact-topic subscribers if let Some(tx) = self.channels.read().get(topic) { let receiver_count = tx.receiver_count(); let _ = tx.send(Arc::clone(&event)); @@ -176,20 +165,35 @@ impl EventBus { receiver_count, "Event delivered to exact-topic channel" ); } + } - // Glob-pattern subscribers - let patterns = self.patterns.read(); - for (pattern, tx) in patterns.iter() { - if pattern.matches(topic) { - let receiver_count = tx.receiver_count(); - let _ = tx.send((topic.to_string(), Arc::clone(&event))); - trace!( - topic, - pattern = pattern.as_str(), - receiver_count, - "Event delivered to pattern channel" - ); - } + // todo : undocumented... + pub fn emit_scoped(&self, topic: &str, scopes: impl IntoIterator, event: T) + where + T: Any + Send + Sync + 'static, + { + let event: AnyEvent = Arc::new(event); + + // Émission sur le topic général. + self.emit_arc(topic, Arc::clone(&event)); + + // Émission sur chaque topic scoped. + for scope in scopes { + let scoped_topic = format!("{}:{}:{}", scope.key, scope.value.into_string(), topic); + + self.emit_arc(&scoped_topic, Arc::clone(&event)); + } + } + + // todo : undocumented... + fn emit_arc(&self, topic: &str, event: AnyEvent) { + trace!(topic, "Emitting event"); + + if let Some(tx) = self.channels.read().get(topic) { + let receiver_count = tx.receiver_count(); + let _ = tx.send(event); + + trace!(topic, receiver_count, "Event delivered to channel"); } } @@ -307,143 +311,6 @@ impl EventBus { }) } - /// Subscribes to all topics matching a glob pattern. - /// - /// The handler receives `(topic, value)` — the topic name is included to - /// distinguish `user-created` from `user-deleted`, for example. - /// - /// **No pre-registration required**: future topics are automatically covered. - /// Supports glob syntax: `*` (any sequence), `?` (one character), - /// `[abc]` (character class). - /// - /// Returns a [`JoinHandle`] to cancel the subscription if needed. - /// - /// # Example - /// ```rust,no_run - /// # use std::sync::Arc; - /// # use oxspeak_server_lib::event_bus::EventBus; - /// # #[derive(Clone, Debug)] struct User { name: String } - /// # let bus = Arc::new(EventBus::new()); - /// bus.on_pattern::("user-*", |topic, user| { - /// match topic.as_str() { - /// "user-created" => println!("Created : {:?}", user), - /// "user-deleted" => println!("Deleted : {:?}", user), - /// other => println!("{}: {:?}", other, user), - /// } - /// }); - /// - /// bus.emit("user-created", User { name: "Alice".into() }); - /// bus.emit("user-deleted", User { name: "Bob".into() }); - /// ``` - pub fn on_pattern(&self, pattern: &str, handler: F) -> JoinHandle<()> - where - T: Any + Send + Sync + Clone + 'static, - F: Fn(String, T) + Send + Sync + 'static, - { - let glob = Pattern::new(pattern).expect("invalid glob pattern"); - let (tx, mut rx) = broadcast::channel(self.capacity); - self.patterns.write().push((glob, tx)); - - let pattern_owned = pattern.to_string(); - debug!(pattern, "Sync pattern subscriber registered"); - - tokio::spawn(async move { - loop { - match rx.recv().await { - Ok((topic, evt)) => { - if let Some(typed) = evt.downcast_ref::() { - trace!( - topic, - pattern = pattern_owned, - "Sync pattern handler invoked" - ); - handler(topic, typed.clone()); - } - } - Err(broadcast::error::RecvError::Lagged(n)) => { - warn!( - pattern = pattern_owned, - skipped = n, - "Pattern subscriber lagged, messages dropped" - ); - } - Err(broadcast::error::RecvError::Closed) => { - debug!( - pattern = pattern_owned, - "Channel closed, sync pattern subscriber exiting" - ); - break; - } - } - } - }) - } - - /// Subscribes to all topics matching a glob pattern, with an **async** handler. - /// - /// Returns a [`JoinHandle`] to cancel the subscription if needed. - /// - /// # Example - /// ```rust,no_run - /// # use std::sync::Arc; - /// # use oxspeak_server_lib::event_bus::EventBus; - /// # #[derive(Clone, Debug)] struct User { name: String } - /// # let bus = Arc::new(EventBus::new()); - /// bus.on_pattern_async::("user-*", |topic, user| async move { - /// match topic.as_str() { - /// "user-created" => println!("(async) Created : {:?}", user), - /// "user-deleted" => println!("(async) Deleted : {:?}", user), - /// other => println!("(async) {}: {:?}", other, user), - /// } - /// }); - /// - /// bus.emit("user-created", User { name: "Alice".into() }); - /// ``` - pub fn on_pattern_async(&self, pattern: &str, handler: F) -> JoinHandle<()> - where - T: Any + Send + Sync + Clone + 'static, - F: Fn(String, T) -> Fut + Send + Sync + 'static, - Fut: Future + Send + 'static, - { - let glob = Pattern::new(pattern).expect("invalid glob pattern"); - let (tx, mut rx) = broadcast::channel(self.capacity); - self.patterns.write().push((glob, tx)); - - let pattern_owned = pattern.to_string(); - debug!(pattern, "Async pattern subscriber registered"); - - tokio::spawn(async move { - loop { - match rx.recv().await { - Ok((topic, evt)) => { - if let Some(typed) = evt.downcast_ref::() { - trace!( - topic, - pattern = pattern_owned, - "Async pattern handler invoked" - ); - handler(topic, typed.clone()).await; - } - } - Err(broadcast::error::RecvError::Lagged(n)) => { - warn!( - pattern = pattern_owned, - skipped = n, - "Pattern subscriber lagged, messages dropped" - ); - } - Err(broadcast::error::RecvError::Closed) => { - debug!( - pattern = pattern_owned, - "Channel closed, async pattern subscriber exiting" - ); - break; - } - } - } - }) - } - // ───────────────────────────────────────────────────────────────────────── // Subscription — low-level access (advanced use cases) // ───────────────────────────────────────────────────────────────────────── diff --git a/event_bus/src/lib.rs b/event_bus/src/lib.rs index 3501c15..8397568 100644 --- a/event_bus/src/lib.rs +++ b/event_bus/src/lib.rs @@ -1,115 +1,3 @@ -//! # Event Bus -//! -//! An asynchronous event bus for routing typed messages between modules -//! without direct coupling. -//! -//! ## Features -//! -//! - **Key → event mapping**: each topic (`&str`) is independent -//! - **Type-unrestricted**: any `T: Any + Send + Sync + Clone` -//! - **Targeted wake-up**: only subscribers of the matching topic are woken up -//! - **Callback API**: JavaScript-style — `bus.on("topic", |payload| { ... })` -//! - **Glob pattern**: `on_pattern("user-*", |topic, payload| { ... })` -//! - **Async handlers**: `on_async` and `on_pattern_async` -//! -//! ## Example — sync callback -//! -//! ```rust,no_run -//! use std::sync::Arc; -//! use oxspeak_server_lib::event_bus::EventBus; -//! -//! #[derive(Clone, Debug)] -//! struct User { name: String } -//! -//! # tokio_test::block_on(async { -//! let bus = Arc::new(EventBus::new()); -//! -//! bus.on::("user-connected", |user| { -//! println!("Connected: {:?}", user); -//! }); -//! -//! bus.emit("user-connected", User { name: "Alice".into() }); -//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await; -//! # }); -//! ``` -//! -//! ## Example — async callback -//! -//! ```rust,no_run -//! use std::sync::Arc; -//! use oxspeak_server_lib::event_bus::EventBus; -//! -//! #[derive(Clone, Debug)] -//! struct User { name: String } -//! -//! # tokio_test::block_on(async { -//! let bus = Arc::new(EventBus::new()); -//! -//! bus.on_async::("user-connected", |user| async move { -//! println!("(async) Connected: {:?}", user); -//! }); -//! -//! bus.emit("user-connected", User { name: "Bob".into() }); -//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await; -//! # }); -//! ``` -//! -//! ## Example — glob pattern (topic included in the callback) -//! -//! ```rust,no_run -//! use std::sync::Arc; -//! use oxspeak_server_lib::event_bus::EventBus; -//! -//! #[derive(Clone, Debug)] -//! struct User { name: String } -//! -//! # tokio_test::block_on(async { -//! let bus = Arc::new(EventBus::new()); -//! -//! bus.on_pattern::("user-*", |topic, user| { -//! match topic.as_str() { -//! "user-created" => println!("Created : {:?}", user), -//! "user-deleted" => println!("Deleted : {:?}", user), -//! other => println!("{}: {:?}", other, user), -//! } -//! }); -//! -//! bus.emit("user-created", User { name: "Alice".into() }); -//! bus.emit("user-deleted", User { name: "Bob".into() }); -//! # tokio::time::sleep(std::time::Duration::from_millis(10)).await; -//! # }); -//! ``` -//! -//! ## Example — multi-type with `match_event!` (advanced) -//! -//! ```rust,no_run -//! use std::sync::Arc; -//! use oxspeak_server_lib::event_bus::EventBus; -//! use oxspeak_server_lib::match_event; -//! -//! #[derive(Clone, Debug)] struct User { name: String } -//! #[derive(Clone, Debug)] struct UdpMetric { value: f32 } -//! -//! # tokio_test::block_on(async { -//! let bus = Arc::new(EventBus::new()); -//! let mut rx = bus.on_raw("mixed-topic"); -//! -//! bus.emit("mixed-topic", User { name: "Alice".into() }); -//! -//! if let Ok(evt) = rx.recv().await { -//! match_event!(evt, -//! User => |u| println!("User: {:?}", u), -//! UdpMetric => |m| println!("Metric: {:?}", m), -//! ); -//! } -//! # }); -//! ``` - -mod bus; - -// Public re-exports -pub use bus::{AnyEvent, EventBus}; - /// Downcasts an [`AnyEvent`] to one or more concrete types and executes /// the matching closure if the type matches. /// @@ -154,9 +42,8 @@ macro_rules! match_event { }; } -// ───────────────────────────────────────────────────────────────────────────── -// Tests -// ───────────────────────────────────────────────────────────────────────────── +mod bus; +pub use bus::{AnyEvent, EventBus}; #[cfg(test)] mod tests; diff --git a/event_bus/src/tests.rs b/event_bus/src/tests.rs index b26e0d6..4efabe8 100644 --- a/event_bus/src/tests.rs +++ b/event_bus/src/tests.rs @@ -1,6 +1,7 @@ use crate::{match_event, EventBus}; use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::sync::Arc; + #[derive(Clone, Debug, PartialEq)] struct User { name: String, @@ -126,108 +127,6 @@ async fn test_on_async_callback() { 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-*", 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-*", 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-*", 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-*", 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]