pub struct IntegrationEventBus { /* private fields */ }Expand description
Central event bus for integration events across all bounded contexts
The IntegrationEventBus is the heart of cross-module communication.
It receives events from publishers and dispatches them to registered handlers
based on pattern matching.
§Features
- Type-erased: Events are serialized to JSON, allowing loose coupling
- Pattern matching: Handlers subscribe to event patterns with wildcards
- Persistence: Optional event history for replay/audit
- Dead letter queue: Failed events are captured for debugging
§Thread Safety
The bus is fully thread-safe and can be cloned and shared across tasks.
Implementations§
Source§impl IntegrationEventBus
impl IntegrationEventBus
Sourcepub fn with_config(config: IntegrationBusConfig) -> Self
pub fn with_config(config: IntegrationBusConfig) -> Self
Create a new integration event bus with custom configuration
Sourcepub async fn publish<E: IntegrationEvent>(
&self,
event: E,
) -> Result<(), EventError>
pub async fn publish<E: IntegrationEvent>( &self, event: E, ) -> Result<(), EventError>
Publish an integration event
The event is serialized to JSON and wrapped in an envelope. It’s then broadcast to all subscribers and dispatched to matching handlers.
§Errors
Returns EventError::SerializationError if the event cannot be serialized.
Sourcepub async fn publish_envelope(
&self,
envelope: IntegrationEventEnvelope,
) -> Result<(), EventError>
pub async fn publish_envelope( &self, envelope: IntegrationEventEnvelope, ) -> Result<(), EventError>
Publish a pre-built envelope
Use this when you already have an envelope (e.g., forwarding from another bus).
Sourcepub async fn register_handler(&self, handler: Arc<dyn IntegrationEventHandler>)
pub async fn register_handler(&self, handler: Arc<dyn IntegrationEventHandler>)
Register an integration event handler
The handler will be called for events matching any of its patterns.
Sourcepub fn subscribe(&self) -> Receiver<IntegrationEventEnvelope>
pub fn subscribe(&self) -> Receiver<IntegrationEventEnvelope>
Subscribe to all integration events (returns a broadcast receiver)
Use this for monitoring, logging, or custom event processing.
Sourcepub async fn history(&self) -> Vec<IntegrationEventEnvelope>
pub async fn history(&self) -> Vec<IntegrationEventEnvelope>
Get event history (if persistence enabled)
Sourcepub async fn events_by_pattern(
&self,
pattern: &str,
) -> Vec<IntegrationEventEnvelope>
pub async fn events_by_pattern( &self, pattern: &str, ) -> Vec<IntegrationEventEnvelope>
Get events by type pattern
Sourcepub async fn events_for_aggregate(
&self,
aggregate_id: &str,
) -> Vec<IntegrationEventEnvelope>
pub async fn events_for_aggregate( &self, aggregate_id: &str, ) -> Vec<IntegrationEventEnvelope>
Get events for a specific aggregate
Sourcepub async fn dead_letters(&self) -> Vec<DeadLetterEntry>
pub async fn dead_letters(&self) -> Vec<DeadLetterEntry>
Get dead letter queue entries
Sourcepub async fn clear_dead_letters(&self)
pub async fn clear_dead_letters(&self)
Clear dead letter queue
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
Sourcepub async fn registered_patterns(&self) -> Vec<String>
pub async fn registered_patterns(&self) -> Vec<String>
Get registered patterns