use std::collections::{HashSet, VecDeque};
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use tokio::sync::broadcast;
use trusty_common::control_bus::{EventId, HarnessEvent};
use super::log::DurableLog;
pub(crate) const DEFAULT_CAPACITY: usize = 8192;
#[derive(Debug, Clone, Copy)]
pub(crate) struct EventBusConfig {
pub capacity: usize,
}
impl Default for EventBusConfig {
fn default() -> Self {
Self {
capacity: DEFAULT_CAPACITY,
}
}
}
#[allow(dead_code)]
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub(crate) struct EventBusMetrics {
pub ingested: u64,
pub deduped: u64,
pub evicted: u64,
pub log_dropped: u64,
}
#[derive(Debug, Default)]
struct Counters {
ingested: AtomicU64,
deduped: AtomicU64,
evicted: AtomicU64,
log_dropped: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum IngestOutcome {
Ingested,
Deduped,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct BusFrame {
pub event: HarnessEvent,
pub persisted: bool,
}
struct Ring {
capacity: usize,
events: VecDeque<HarnessEvent>,
ids: HashSet<EventId>,
}
pub(crate) struct EventBus {
ring: Mutex<Ring>,
counters: Counters,
sender: broadcast::Sender<BusFrame>,
next_seq: AtomicU64,
log: Option<DurableLog>,
}
impl EventBus {
pub(crate) fn new(config: EventBusConfig) -> Self {
Self::with_log(config, None, 1)
}
pub(crate) fn with_log(
config: EventBusConfig,
log: Option<DurableLog>,
recovered_next_seq: u64,
) -> Self {
let capacity = config.capacity.max(1);
let (sender, _receiver) = broadcast::channel(capacity);
Self {
ring: Mutex::new(Ring {
capacity,
events: VecDeque::with_capacity(capacity),
ids: HashSet::with_capacity(capacity),
}),
counters: Counters::default(),
sender,
next_seq: AtomicU64::new(recovered_next_seq.max(1)),
log,
}
}
pub(crate) fn ingest(&self, mut event: HarnessEvent) -> IngestOutcome {
let id = event.id;
let mut ring = self.lock_ring();
if !ring.ids.insert(id) {
drop(ring);
self.counters.deduped.fetch_add(1, Ordering::Relaxed);
return IngestOutcome::Deduped;
}
event.seq = self.next_seq.fetch_add(1, Ordering::Relaxed);
if ring.events.len() >= ring.capacity
&& let Some(evicted) = ring.events.pop_front()
{
ring.ids.remove(&evicted.id);
self.counters.evicted.fetch_add(1, Ordering::Relaxed);
}
ring.events.push_back(event.clone());
self.counters.ingested.fetch_add(1, Ordering::Relaxed);
if let Some(log) = &self.log
&& !log.enqueue(event.clone())
{
self.counters.log_dropped.fetch_add(1, Ordering::Relaxed);
tracing::warn!(
seq = event.seq,
"event-bus: durable-log writer is backpressured; this event \
will not be persisted"
);
}
let _ = self.sender.send(BusFrame {
event,
persisted: false,
});
drop(ring);
IngestOutcome::Ingested
}
#[allow(dead_code)]
pub(crate) fn subscribe(&self) -> broadcast::Receiver<BusFrame> {
self.sender.subscribe()
}
#[allow(dead_code)]
pub(crate) async fn replay_since(
&self,
since_seq: u64,
) -> Option<Result<Vec<super::log::ReplayItem>, super::log::LogError>> {
let log = self.log.as_ref()?;
Some(log.replay_since(since_seq).await)
}
#[allow(dead_code)]
pub(crate) fn metrics(&self) -> EventBusMetrics {
EventBusMetrics {
ingested: self.counters.ingested.load(Ordering::Relaxed),
deduped: self.counters.deduped.load(Ordering::Relaxed),
evicted: self.counters.evicted.load(Ordering::Relaxed),
log_dropped: self.counters.log_dropped.load(Ordering::Relaxed),
}
}
#[cfg(test)]
pub(crate) fn len(&self) -> usize {
self.lock_ring().events.len()
}
#[cfg(test)]
pub(crate) fn contains(&self, id: EventId) -> bool {
self.lock_ring().ids.contains(&id)
}
fn lock_ring(&self) -> std::sync::MutexGuard<'_, Ring> {
self.ring
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}