Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
213 changes: 188 additions & 25 deletions wacore/src/types/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -325,62 +325,75 @@ impl EventHandler for ChannelEventHandler {
}
}

/// Immutable snapshot of the registered handlers. `dispatch` clones only the
/// outer `Arc` (one refcount bump, no `Vec` allocation), then drops the lock and
/// iterates the snapshot. Handler interest is re-evaluated per dispatch so a
/// handler whose `interest()` widens at runtime still receives the new kinds.
#[derive(Default)]
struct HandlerSnapshot {
handlers: Vec<Arc<dyn EventHandler>>,
}

#[derive(Default, Clone)]
pub struct CoreEventBus {
handlers: Arc<RwLock<Vec<Arc<dyn EventHandler>>>>,
// Copy-on-write: the snapshot is only swapped (under the lock) when a
// handler is added, which happens at startup. `dispatch` takes a cheap
// outer-Arc clone and then drops the lock, so a concurrent `add_handler`
// can never invalidate a snapshot a dispatch is iterating.
handlers: Arc<RwLock<Arc<HandlerSnapshot>>>,
}

impl CoreEventBus {
pub fn new() -> Self {
Self::default()
}

pub fn add_handler(&self, handler: Arc<dyn EventHandler>) {
fn snapshot(&self) -> Arc<HandlerSnapshot> {
self.handlers
.write()
.read()
.expect("RwLock should not be poisoned")
.push(handler);
.clone()
}

pub fn add_handler(&self, handler: Arc<dyn EventHandler>) {
let mut guard = self
.handlers
.write()
.expect("RwLock should not be poisoned");
let current = &**guard;
let mut handlers = Vec::with_capacity(current.handlers.len() + 1);
handlers.extend(current.handlers.iter().cloned());
handlers.push(handler);
*guard = Arc::new(HandlerSnapshot { handlers });
}

/// Returns true if there are any event handlers registered.
/// Useful for skipping expensive work when no one is listening.
pub fn has_handlers(&self) -> bool {
!self
.handlers
.read()
.expect("RwLock should not be poisoned")
.is_empty()
!self.snapshot().handlers.is_empty()
}

/// Whether any registered handler is interested in `kind`. Lets callers
/// skip producing an event nobody would receive (e.g. retaining a large
/// `HistorySync` blob when only message-only handlers are registered).
pub fn has_handler_for(&self, kind: EventKind) -> bool {
self.handlers
.read()
.expect("RwLock should not be poisoned")
self.snapshot()
.handlers
.iter()
.any(|h| h.interest().wants(kind))
}

pub fn dispatch(&self, event: Event) {
let handlers = self
.handlers
.read()
.expect("RwLock should not be poisoned")
.clone();
if handlers.is_empty() {
return;
}
// Skip materializing the event (Arc) and invoking handlers whose
// declared interest excludes this kind. A handler that subscribed to a
// few kinds never pays for boxing the events it ignores.
let snapshot = self.snapshot();
// Skip materializing the event (Arc) when no handler wants this kind. The
// interest is re-evaluated here (not read from a cached aggregate) so a
// handler whose interest() widens at runtime is never short-circuited out.
let kind = event.kind();
if !handlers.iter().any(|h| h.interest().wants(kind)) {
if !snapshot.handlers.iter().any(|h| h.interest().wants(kind)) {
return;
}
let event = Arc::new(event);
for handler in &handlers {
for handler in &snapshot.handlers {
if handler.interest().wants(kind) {
handler.handle_event(Arc::clone(&event));
}
Expand Down Expand Up @@ -1468,4 +1481,154 @@ mod tests {
bus2.dispatch(Event::Connected(Connected));
assert_eq!(CALLS.load(Ordering::SeqCst), 0);
}

#[test]
fn dispatch_respects_dynamically_widened_interest() {
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};

// A handler whose interest() widens after registration. dispatch must
// re-read interest each time (never a stale cached aggregate), so the
// newly-wanted kind is delivered.
struct Dynamic {
interest: Mutex<EventInterest>,
hits: AtomicUsize,
}
impl EventHandler for Dynamic {
fn handle_event(&self, _: Arc<Event>) {
self.hits.fetch_add(1, Ordering::SeqCst);
}
fn interest(&self) -> EventInterest {
*self.interest.lock().unwrap()
}
}

let bus = CoreEventBus::new();
let h = Arc::new(Dynamic {
interest: Mutex::new(EventInterest::of(&[EventKind::Message])),
hits: AtomicUsize::new(0),
});
bus.add_handler(h.clone());

// Not yet interested in Connected: dropped before materialization.
bus.dispatch(Event::Connected(Connected));
assert_eq!(h.hits.load(Ordering::SeqCst), 0);
assert!(!bus.has_handler_for(EventKind::Connected));

// Widen interest at runtime.
*h.interest.lock().unwrap() = EventInterest::ALL;
assert!(bus.has_handler_for(EventKind::Connected));
bus.dispatch(Event::Connected(Connected));
assert_eq!(
h.hits.load(Ordering::SeqCst),
1,
"a handler whose interest widened at runtime must receive the newly-wanted kind"
);
}

#[test]
fn aggregate_interest_and_has_handler_for() {
struct Narrow(EventInterest);
impl EventHandler for Narrow {
fn handle_event(&self, _: Arc<Event>) {}
fn interest(&self) -> EventInterest {
self.0
}
}

let bus = CoreEventBus::new();
// Empty bus: nothing is wanted and there are no handlers.
assert!(!bus.has_handlers());
assert!(!bus.has_handler_for(EventKind::Message));
assert!(!bus.has_handler_for(EventKind::Receipt));

bus.add_handler(Arc::new(Narrow(EventInterest::of(&[EventKind::Message]))));
assert!(bus.has_handlers());
assert!(bus.has_handler_for(EventKind::Message));
assert!(!bus.has_handler_for(EventKind::Receipt));

// has_handler_for is true once any registered handler wants the kind.
bus.add_handler(Arc::new(Narrow(EventInterest::of(&[EventKind::Receipt]))));
assert!(bus.has_handler_for(EventKind::Message));
assert!(bus.has_handler_for(EventKind::Receipt));
assert!(!bus.has_handler_for(EventKind::Connected));
}

#[test]
fn dispatch_preserves_handler_ordering() {
use std::sync::Mutex;

struct Tagged {
id: u32,
log: Arc<Mutex<Vec<u32>>>,
}
impl EventHandler for Tagged {
fn handle_event(&self, _: Arc<Event>) {
self.log.lock().unwrap().push(self.id);
}
}

let bus = CoreEventBus::new();
let log = Arc::new(Mutex::new(Vec::new()));
for id in 0..5u32 {
bus.add_handler(Arc::new(Tagged {
id,
log: log.clone(),
}));
}
bus.dispatch(Event::Connected(Connected));
// Copy-on-write rebuilds must keep registration order intact.
assert_eq!(*log.lock().unwrap(), vec![0, 1, 2, 3, 4]);
}

#[test]
fn dispatch_is_reentrancy_safe_against_concurrent_add() {
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};

// A handler that registers another handler while it is being dispatched.
// The snapshot taken by `dispatch` must outlive the swap, so the newly
// added handler is NOT invoked for the in-flight event and the iteration
// does not observe a mutated list.
struct AddsDuringDispatch {
bus: CoreEventBus,
invocations: Arc<AtomicUsize>,
added: Mutex<bool>,
}
impl EventHandler for AddsDuringDispatch {
fn handle_event(&self, _: Arc<Event>) {
self.invocations.fetch_add(1, Ordering::SeqCst);
let mut added = self.added.lock().unwrap();
if !*added {
*added = true;
struct Late(Arc<AtomicUsize>);
impl EventHandler for Late {
fn handle_event(&self, _: Arc<Event>) {
self.0.fetch_add(1, Ordering::SeqCst);
}
}
self.bus
.add_handler(Arc::new(Late(self.invocations.clone())));
}
}
}

let bus = CoreEventBus::new();
let invocations = Arc::new(AtomicUsize::new(0));
bus.add_handler(Arc::new(AddsDuringDispatch {
bus: bus.clone(),
invocations: invocations.clone(),
added: Mutex::new(false),
}));

// First dispatch: only the original handler runs, even though it adds a
// second handler mid-flight.
bus.dispatch(Event::Connected(Connected));
assert_eq!(invocations.load(Ordering::SeqCst), 1);
assert_eq!(bus.snapshot().handlers.len(), 2);

// Second dispatch sees both handlers (original adds nothing new now).
bus.dispatch(Event::Connected(Connected));
assert_eq!(invocations.load(Ordering::SeqCst), 3);
}
}
Loading