use dashmap::DashMap;
use std::any::{Any, TypeId};
use std::collections::HashMap;
use std::sync::Arc;
#[derive(Debug, Clone)]
pub struct ContextInitializingEvent {
pub config_sources_count: usize,
pub timestamp: std::time::SystemTime,
}
impl Event for ContextInitializingEvent {
fn name(&self) -> &'static str {
"ContextInitializing"
}
fn as_any(&self) -> &dyn Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn Any> {
self
}
}
#[derive(Debug, Clone)]
pub struct ContextInitializedEvent {
pub config_sources_count: usize,
pub timestamp: std::time::SystemTime,
}
impl Event for ContextInitializedEvent {
fn name(&self) -> &'static str {
"ContextInitialized"
}
fn as_any(&self) -> &dyn Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn Any> {
self
}
}
#[derive(Debug, Clone)]
pub struct ConfigurationChangedEvent {
pub key: String,
pub old_value: Option<String>,
pub new_value: String,
pub timestamp: std::time::SystemTime,
}
impl Event for ConfigurationChangedEvent {
fn name(&self) -> &'static str {
"ConfigurationChanged"
}
fn as_any(&self) -> &dyn Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn Any> {
self
}
}
pub trait Event: Any + Send + Sync {
fn name(&self) -> &'static str;
fn as_any(&self) -> &dyn Any;
fn into_any(self: Box<Self>) -> Box<dyn Any>;
}
pub trait EventListener<T: Event>: Send + Sync {
fn on_event(&self, event: &T);
}
pub trait AnyEventListener: Send + Sync {
fn handle_event(&self, event: &dyn Event) -> bool;
fn event_type_id(&self) -> TypeId;
}
pub trait ContextAwareEventListener<T: Event>: Send + Sync {
fn on_context_event(&self, event: &T, context: &crate::context::ApplicationContext);
}
pub trait AnyContextAwareEventListener: Send + Sync {
fn handle_context_event(
&self,
event: &dyn Event,
context: &crate::context::ApplicationContext,
) -> bool;
fn event_type_id(&self) -> TypeId;
}
struct TypedContextAwareEventListener<T: Event, L: ContextAwareEventListener<T>> {
listener: L,
_phantom: std::marker::PhantomData<T>,
}
impl<T: Event, L: ContextAwareEventListener<T>> TypedContextAwareEventListener<T, L> {
fn new(listener: L) -> Self {
Self {
listener,
_phantom: std::marker::PhantomData,
}
}
}
impl<T: Event + 'static, L: ContextAwareEventListener<T>> AnyContextAwareEventListener
for TypedContextAwareEventListener<T, L>
{
fn handle_context_event(
&self,
event: &dyn Event,
context: &crate::context::ApplicationContext,
) -> bool {
if let Some(typed_event) = event.as_any().downcast_ref::<T>() {
self.listener.on_context_event(typed_event, context);
true
} else {
false
}
}
fn event_type_id(&self) -> TypeId {
TypeId::of::<T>()
}
}
struct TypedEventListener<T: Event, L: EventListener<T>> {
listener: L,
_phantom: std::marker::PhantomData<T>,
}
impl<T: Event + 'static, L: EventListener<T>> TypedEventListener<T, L> {
fn new(listener: L) -> Self {
Self {
listener,
_phantom: std::marker::PhantomData,
}
}
}
impl<T: Event + 'static, L: EventListener<T>> AnyEventListener for TypedEventListener<T, L> {
fn handle_event(&self, event: &dyn Event) -> bool {
if let Some(typed_event) = event.as_any().downcast_ref::<T>() {
self.listener.on_event(typed_event);
true
} else {
false
}
}
fn event_type_id(&self) -> TypeId {
TypeId::of::<T>()
}
}
pub struct EventPublisher {
listeners: DashMap<TypeId, Vec<Arc<dyn AnyEventListener>>>,
context_aware_listeners: DashMap<TypeId, Vec<Arc<dyn AnyContextAwareEventListener>>>,
}
impl EventPublisher {
pub fn new() -> Self {
Self {
listeners: DashMap::new(),
context_aware_listeners: DashMap::new(),
}
}
pub fn subscribe_context_aware<
T: Event + 'static,
L: ContextAwareEventListener<T> + 'static,
>(
&self,
listener: L,
) {
let type_id = TypeId::of::<T>();
let typed_listener = Arc::new(TypedContextAwareEventListener::new(listener));
self.context_aware_listeners
.entry(type_id)
.or_insert_with(Vec::new)
.push(typed_listener);
}
pub fn subscribe<T: Event + 'static, L: EventListener<T> + 'static>(&self, listener: L) {
let type_id = TypeId::of::<T>();
let typed_listener = Arc::new(TypedEventListener::new(listener));
self.listeners
.entry(type_id)
.or_insert_with(Vec::new)
.push(typed_listener);
}
pub fn publish_with_context<T: Event + 'static>(
&self,
event: &T,
context: &crate::context::ApplicationContext,
) {
let type_id = TypeId::of::<T>();
if let Some(listeners) = self.listeners.get(&type_id) {
for listener in listeners.iter() {
listener.handle_event(event);
}
}
if let Some(context_listeners) = self.context_aware_listeners.get(&type_id) {
for listener in context_listeners.iter() {
listener.handle_context_event(event, context);
}
}
}
pub fn publish<T: Event + 'static>(&self, event: &T) {
let type_id = TypeId::of::<T>();
if let Some(listeners) = self.listeners.get(&type_id) {
for listener in listeners.iter() {
listener.handle_event(event);
}
}
}
pub fn listener_count<T: Event + 'static>(&self) -> usize {
let type_id = TypeId::of::<T>();
self.listeners
.get(&type_id)
.map(|listeners| listeners.len())
.unwrap_or(0)
}
pub fn clear_all_listeners(&mut self) {
self.listeners.clear();
}
pub fn listener_statistics(&self) -> HashMap<String, usize> {
let mut stats = HashMap::new();
for entry in self.listeners.iter() {
let type_id = entry.key();
let listeners = entry.value();
let type_name = format!("{:?}", type_id);
stats.insert(type_name, listeners.len());
}
stats
}
}
impl Default for EventPublisher {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Debug, Clone)]
struct TestEvent {
message: String,
}
impl Event for TestEvent {
fn name(&self) -> &'static str {
"TestEvent"
}
fn as_any(&self) -> &dyn Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn Any> {
self
}
}
#[derive(Debug, Clone)]
struct AnotherEvent {
value: i32,
}
impl Event for AnotherEvent {
fn name(&self) -> &'static str {
"AnotherEvent"
}
fn as_any(&self) -> &dyn Any {
self
}
fn into_any(self: Box<Self>) -> Box<dyn Any> {
self
}
}
static TEST_COUNTER: AtomicUsize = AtomicUsize::new(0);
struct TestListener;
impl EventListener<TestEvent> for TestListener {
fn on_event(&self, _event: &TestEvent) {
TEST_COUNTER.fetch_add(1, Ordering::SeqCst);
}
}
struct AnotherListener;
impl EventListener<AnotherEvent> for AnotherListener {
fn on_event(&self, _event: &AnotherEvent) {
TEST_COUNTER.fetch_add(10, Ordering::SeqCst);
}
}
#[test]
fn test_event_publisher_creation() {
let publisher = EventPublisher::new();
assert_eq!(publisher.listener_count::<TestEvent>(), 0);
}
#[test]
fn test_event_subscription() {
let publisher = EventPublisher::new();
publisher.subscribe(TestListener);
assert_eq!(publisher.listener_count::<TestEvent>(), 1);
assert_eq!(publisher.listener_count::<AnotherEvent>(), 0);
}
#[test]
fn test_event_publishing() {
TEST_COUNTER.store(0, Ordering::SeqCst);
let publisher = EventPublisher::new();
publisher.subscribe(TestListener);
let event = TestEvent {
message: "test".to_string(),
};
publisher.publish(&event);
assert_eq!(TEST_COUNTER.load(Ordering::SeqCst), 1);
}
#[test]
fn test_multiple_listeners_same_event() {
TEST_COUNTER.store(0, Ordering::SeqCst);
let publisher = EventPublisher::new();
publisher.subscribe(TestListener);
publisher.subscribe(TestListener);
let event = TestEvent {
message: "test".to_string(),
};
publisher.publish(&event);
assert_eq!(TEST_COUNTER.load(Ordering::SeqCst), 2);
assert_eq!(publisher.listener_count::<TestEvent>(), 2);
}
#[test]
fn test_different_event_types() {
TEST_COUNTER.store(0, Ordering::SeqCst);
let publisher = EventPublisher::new();
publisher.subscribe(TestListener);
publisher.subscribe(AnotherListener);
let test_event = TestEvent {
message: "test".to_string(),
};
let another_event = AnotherEvent { value: 42 };
publisher.publish(&test_event);
publisher.publish(&another_event);
assert_eq!(TEST_COUNTER.load(Ordering::SeqCst), 11);
}
#[test]
fn test_event_without_listeners() {
let publisher = EventPublisher::new();
let event = TestEvent {
message: "no listeners".to_string(),
};
publisher.publish(&event);
}
#[test]
fn test_clear_listeners() {
let mut publisher = EventPublisher::new();
publisher.subscribe(TestListener);
publisher.subscribe(AnotherListener);
assert_eq!(publisher.listener_count::<TestEvent>(), 1);
assert_eq!(publisher.listener_count::<AnotherEvent>(), 1);
publisher.clear_all_listeners();
assert_eq!(publisher.listener_count::<TestEvent>(), 0);
assert_eq!(publisher.listener_count::<AnotherEvent>(), 0);
}
#[test]
fn test_listener_statistics() {
let publisher = EventPublisher::new();
publisher.subscribe(TestListener);
publisher.subscribe(AnotherListener);
let stats = publisher.listener_statistics();
assert_eq!(stats.len(), 2);
}
#[test]
fn test_event_trait_methods() {
let event = TestEvent {
message: "test".to_string(),
};
assert_eq!(event.name(), "TestEvent");
assert!(event.as_any().is::<TestEvent>());
}
}