pub struct EventBus<E: DomainEvent> { /* private fields */ }Expand description
Generic event bus for publishing and subscribing to domain events
The EventBus provides:
- Type-safe publish/subscribe
- Event envelope with metadata
- Handler registration by event type
- Optional event history for replay
§Type Parameters
E: The domain event type (must implementDomainEvent)
§Example
ⓘ
use backbone_messaging::{EventBus, EventBusConfig, DomainEvent};
#[derive(Clone, Debug)]
struct UserCreated { user_id: String }
impl DomainEvent for UserCreated {
fn event_type(&self) -> &'static str { "UserCreated" }
fn aggregate_id(&self) -> &str { &self.user_id }
}
let bus = EventBus::<UserCreated>::new();
bus.publish(UserCreated { user_id: "123".into() }).await?;Implementations§
Source§impl<E: DomainEvent> EventBus<E>
impl<E: DomainEvent> EventBus<E>
Sourcepub fn with_config(config: EventBusConfig) -> Self
pub fn with_config(config: EventBusConfig) -> Self
Create a new event bus with custom configuration
Sourcepub async fn publish(&self, event: E) -> Result<(), EventError>
pub async fn publish(&self, event: E) -> Result<(), EventError>
Publish a domain event
Sourcepub async fn publish_all(&self, events: Vec<E>) -> Result<(), EventError>
pub async fn publish_all(&self, events: Vec<E>) -> Result<(), EventError>
Publish multiple domain events
Sourcepub async fn publish_envelope(
&self,
envelope: EventEnvelope<E>,
) -> Result<(), EventError>
pub async fn publish_envelope( &self, envelope: EventEnvelope<E>, ) -> Result<(), EventError>
Publish an event envelope (with metadata)
Sourcepub async fn register_handler(&self, handler: Arc<dyn EventHandler<E>>)
pub async fn register_handler(&self, handler: Arc<dyn EventHandler<E>>)
Register an event handler
Sourcepub fn subscribe(&self) -> Receiver<EventEnvelope<E>>
pub fn subscribe(&self) -> Receiver<EventEnvelope<E>>
Subscribe to all events (returns a broadcast receiver)
Sourcepub async fn history(&self) -> Vec<EventEnvelope<E>>
pub async fn history(&self) -> Vec<EventEnvelope<E>>
Get event history (if persistence enabled)
Sourcepub async fn events_for_aggregate(
&self,
aggregate_id: &str,
) -> Vec<EventEnvelope<E>>
pub async fn events_for_aggregate( &self, aggregate_id: &str, ) -> Vec<EventEnvelope<E>>
Get events for a specific aggregate
Sourcepub async fn events_by_type(&self, event_type: &str) -> Vec<EventEnvelope<E>>
pub async fn events_by_type(&self, event_type: &str) -> Vec<EventEnvelope<E>>
Get events by type
Sourcepub async fn events_in_range(
&self,
start: DateTime<Utc>,
end: DateTime<Utc>,
) -> Vec<EventEnvelope<E>>
pub async fn events_in_range( &self, start: DateTime<Utc>, end: DateTime<Utc>, ) -> Vec<EventEnvelope<E>>
Get events in a time range
Sourcepub async fn clear_history(&self)
pub async fn clear_history(&self)
Clear event history
Sourcepub async fn handler_count(&self) -> usize
pub async fn handler_count(&self) -> usize
Get handler count
Trait Implementations§
Source§impl<E: DomainEvent> Clone for EventBus<E>
impl<E: DomainEvent> Clone for EventBus<E>
Auto Trait Implementations§
impl<E> !RefUnwindSafe for EventBus<E>
impl<E> !UnwindSafe for EventBus<E>
impl<E> Freeze for EventBus<E>
impl<E> Send for EventBus<E>
impl<E> Sync for EventBus<E>
impl<E> Unpin for EventBus<E>
impl<E> UnsafeUnpin for EventBus<E>where
Sender<EventEnvelope<E>>: UnsafeUnpin,
Arc<RwLock<HashMap<String, Vec<Arc<dyn EventHandler<E>>>>>>: UnsafeUnpin,
Arc<RwLock<Vec<EventEnvelope<E>>>>: UnsafeUnpin,
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more