use crate::envelope::EventEnvelope;
use crate::operational_journal::{FileOperationalJournal, OperationalJournalRecord};
use crate::TraceContext;
use parking_lot::Mutex;
use serde::ser::SerializeSeq;
use serde::{Serialize, Serializer};
use std::collections::VecDeque;
use std::sync::Arc;
const MAX_RETAINED_EVENTS: usize = 10_000;
pub const DEFAULT_EVENT_BUS_MAX_BYTES: usize = 16 * 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EventBusStats {
pub event_count: usize,
pub used_bytes: usize,
pub peak_bytes: usize,
pub evictions: u64,
pub rejections: u64,
pub max_bytes: usize,
}
#[derive(Clone, Debug)]
pub struct EventBusSnapshot {
events: Arc<VecDeque<Arc<OperationalJournalRecord>>>,
}
impl EventBusSnapshot {
#[must_use]
pub fn len(&self) -> usize {
self.events.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
pub fn recent(&self, limit: usize) -> impl Iterator<Item = &EventEnvelope> {
let start = self.events.len().saturating_sub(limit);
self.events.iter().skip(start).filter_map(event_record)
}
}
impl Serialize for EventBusSnapshot {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
let mut sequence = serializer.serialize_seq(Some(self.events.len()))?;
for event in self.events.iter().filter_map(event_record) {
sequence.serialize_element(event)?;
}
sequence.end()
}
}
#[derive(Debug, Clone)]
struct EventBusState {
events: Arc<VecDeque<Arc<OperationalJournalRecord>>>,
used_bytes: usize,
peak_bytes: usize,
evictions: u64,
rejections: u64,
max_bytes: usize,
}
impl EventBusState {
fn new(max_bytes: usize) -> Self {
Self {
events: Arc::default(),
used_bytes: 0,
peak_bytes: 0,
evictions: 0,
rejections: 0,
max_bytes: max_bytes.max(1),
}
}
fn push_shared(&mut self, record: Arc<OperationalJournalRecord>) {
let Some(event) = event_record(&record) else {
self.rejections = self.rejections.saturating_add(1);
return;
};
let bytes = event_retained_bytes(event);
if bytes > self.max_bytes {
self.rejections = self.rejections.saturating_add(1);
return;
}
while self.events.len() >= MAX_RETAINED_EVENTS {
self.pop_front();
}
while self.used_bytes.saturating_add(bytes) > self.max_bytes {
if !self.pop_front() {
self.rejections = self.rejections.saturating_add(1);
return;
}
}
Arc::make_mut(&mut self.events).push_back(record);
self.used_bytes = self.used_bytes.saturating_add(bytes);
self.peak_bytes = self.peak_bytes.max(self.used_bytes);
}
fn pop_front(&mut self) -> bool {
let Some(event) = Arc::make_mut(&mut self.events).pop_front() else {
return false;
};
if let Some(event) = event_record(&event) {
self.used_bytes = self.used_bytes.saturating_sub(event_retained_bytes(event));
}
self.evictions = self.evictions.saturating_add(1);
true
}
fn replace(&mut self, events: Vec<Arc<OperationalJournalRecord>>) {
self.clear();
for event in events {
self.push_shared(event);
}
}
fn clear(&mut self) {
self.events = Arc::default();
self.used_bytes = 0;
}
fn stats(&self) -> EventBusStats {
EventBusStats {
event_count: self.events.len(),
used_bytes: self.used_bytes,
peak_bytes: self.peak_bytes,
evictions: self.evictions,
rejections: self.rejections,
max_bytes: self.max_bytes,
}
}
}
#[derive(Debug)]
pub struct EventBus {
state: Mutex<EventBusState>,
journal: Mutex<Option<Arc<FileOperationalJournal>>>,
journal_error: Mutex<Option<String>>,
}
impl Default for EventBus {
fn default() -> Self {
Self::with_max_bytes(DEFAULT_EVENT_BUS_MAX_BYTES)
}
}
impl Clone for EventBus {
fn clone(&self) -> Self {
Self {
state: Mutex::new(self.state.lock().clone()),
journal: Mutex::new(self.journal.lock().clone()),
journal_error: Mutex::new(self.journal_error.lock().clone()),
}
}
}
impl EventBus {
pub fn new() -> Self {
Self::default()
}
pub fn with_max_bytes(max_bytes: usize) -> Self {
Self {
state: Mutex::new(EventBusState::new(max_bytes)),
journal: Mutex::new(None),
journal_error: Mutex::new(None),
}
}
pub fn attach_journal(&self, journal: Arc<FileOperationalJournal>) {
self.state.lock().replace(journal.shared_event_records());
*self.journal.lock() = Some(journal);
*self.journal_error.lock() = None;
}
pub fn durability_error(&self) -> Option<String> {
self.journal_error.lock().clone()
}
pub fn emit(&self, event: EventEnvelope) {
let event = Arc::new(OperationalJournalRecord::Event(event));
self.persist(Arc::clone(&event));
self.state.lock().push_shared(event);
}
pub fn emit_many(&self, events: Vec<EventEnvelope>) {
let events = events
.into_iter()
.map(|event| Arc::new(OperationalJournalRecord::Event(event)))
.collect::<Vec<_>>();
for event in &events {
self.persist(Arc::clone(event));
}
let mut state = self.state.lock();
for event in events {
state.push_shared(event);
}
}
pub fn len(&self) -> usize {
self.state.lock().events.len()
}
pub fn is_empty(&self) -> bool {
self.state.lock().events.is_empty()
}
pub fn events(&self) -> Vec<EventEnvelope> {
self.snapshot().recent(usize::MAX).cloned().collect()
}
pub fn snapshot(&self) -> EventBusSnapshot {
EventBusSnapshot {
events: Arc::clone(&self.state.lock().events),
}
}
pub fn stats(&self) -> EventBusStats {
self.state.lock().stats()
}
pub fn clear(&self) {
self.state.lock().clear();
}
fn persist(&self, event: Arc<OperationalJournalRecord>) {
if let Some(journal) = self.journal.lock().clone() {
if let Err(error) = journal.append_shared_event(event) {
*self.journal_error.lock() = Some(crate::redact_text(&format!("{error:?}")));
}
}
}
}
fn event_record(record: &Arc<OperationalJournalRecord>) -> Option<&EventEnvelope> {
match record.as_ref() {
OperationalJournalRecord::Event(event) => Some(event),
OperationalJournalRecord::Audit(_) => None,
}
}
fn event_retained_bytes(event: &EventEnvelope) -> usize {
std::mem::size_of::<EventEnvelope>()
.saturating_add(event.event_name.as_str().len())
.saturating_add(event.event_id.len())
.saturating_add(event.app_id.as_str().len())
.saturating_add(event.node_id.as_str().len())
.saturating_add(event.payload.len())
.saturating_add(trace_retained_bytes(event.trace.as_ref()))
}
fn trace_retained_bytes(trace: Option<&TraceContext>) -> usize {
let Some(trace) = trace else {
return 0;
};
std::mem::size_of::<TraceContext>()
.saturating_add(trace.trace_id.len())
.saturating_add(trace.span_id.len())
.saturating_add(trace.parent_span_id.as_ref().map_or(0, String::len))
.saturating_add(trace.originating_core_id.as_str().len())
.saturating_add(trace.current_core_id.as_str().len())
.saturating_add(trace.tenant_id.as_str().len())
.saturating_add(trace.command_id.as_ref().map_or(0, String::len))
}
#[cfg(test)]
#[path = "event_bus_tests.rs"]
mod tests;