use std::sync::Arc;
use super::{
BackpressurePolicy, DurableAuditSink, EventBus, EventEmissionMode, MetricsSink, Result,
new_with_config_and_mode_and_metrics_impl, new_with_config_and_mode_impl,
validate_regulated_audit_sink,
};
#[derive(Clone)]
pub struct EventBusBuilder {
capacity: usize,
audit_sink: Option<Arc<dyn DurableAuditSink>>,
backpressure: BackpressurePolicy,
emission_mode: EventEmissionMode,
metrics_sink: Option<Arc<dyn MetricsSink>>,
}
impl std::fmt::Debug for EventBusBuilder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("EventBusBuilder")
.field("capacity", &self.capacity)
.field("audit_sink", &self.audit_sink.is_some())
.field("backpressure", &self.backpressure)
.field("emission_mode", &self.emission_mode)
.field("metrics_sink", &self.metrics_sink.is_some())
.finish()
}
}
impl Default for EventBusBuilder {
fn default() -> Self {
Self {
capacity: 256,
audit_sink: None,
backpressure: BackpressurePolicy::default(),
emission_mode: EventEmissionMode::StrictTransactional,
metrics_sink: None,
}
}
}
impl EventBusBuilder {
#[must_use]
pub fn capacity(mut self, capacity: usize) -> Self {
self.capacity = capacity;
self
}
#[must_use]
pub fn audit_sink(mut self, sink: Arc<dyn DurableAuditSink>) -> Self {
self.audit_sink = Some(sink);
self
}
#[must_use]
pub fn backpressure(mut self, backpressure: BackpressurePolicy) -> Self {
self.backpressure = backpressure;
self
}
#[must_use]
pub fn emission_mode(mut self, mode: EventEmissionMode) -> Self {
self.emission_mode = mode;
self
}
#[must_use]
pub fn metrics_sink(mut self, sink: Arc<dyn MetricsSink>) -> Self {
self.metrics_sink = Some(sink);
self
}
pub fn build(self) -> Result<EventBus> {
match self.metrics_sink {
Some(metrics) => new_with_config_and_mode_and_metrics_impl(
self.capacity,
self.audit_sink,
self.backpressure,
self.emission_mode,
metrics,
),
None => new_with_config_and_mode_impl(
self.capacity,
self.audit_sink,
self.backpressure,
self.emission_mode,
),
}
}
}
impl EventBus {
#[must_use]
pub fn builder() -> EventBusBuilder {
EventBusBuilder::default()
}
pub fn new(capacity: usize) -> Result<Self> {
Self::builder().capacity(capacity).build()
}
pub fn new_regulated(capacity: usize, audit_sink: Arc<dyn DurableAuditSink>) -> Result<Self> {
validate_regulated_audit_sink(&audit_sink)?;
Self::builder()
.capacity(capacity)
.audit_sink(audit_sink)
.backpressure(BackpressurePolicy::regulated())
.build()
}
#[cfg(feature = "testing")]
#[must_use]
pub fn new_for_testing() -> Self {
Self::builder()
.emission_mode(EventEmissionMode::BestEffort)
.build()
.expect("EventBus::new_for_testing: infallible BestEffort construction")
}
#[cfg(test)]
pub(crate) fn new_with_config_and_mode(
capacity: usize,
audit_sink: Option<Arc<dyn DurableAuditSink>>,
backpressure: BackpressurePolicy,
emission_mode: EventEmissionMode,
) -> Result<Self> {
new_with_config_and_mode_impl(capacity, audit_sink, backpressure, emission_mode)
}
#[cfg(test)]
pub(crate) fn new_with_config_and_mode_and_metrics(
capacity: usize,
audit_sink: Option<Arc<dyn DurableAuditSink>>,
backpressure: BackpressurePolicy,
emission_mode: EventEmissionMode,
metrics_sink: Arc<dyn MetricsSink>,
) -> Result<Self> {
new_with_config_and_mode_and_metrics_impl(
capacity,
audit_sink,
backpressure,
emission_mode,
metrics_sink,
)
}
#[cfg(test)]
pub(crate) fn new_strict_with_audit_fallback(
capacity: usize,
audit_sink: Arc<dyn DurableAuditSink>,
) -> Result<Self> {
Self::builder()
.capacity(capacity)
.audit_sink(audit_sink)
.emission_mode(EventEmissionMode::StrictWithAuditFallback)
.build()
}
}