Skip to main content

appcore_core/
event_bus.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: event_bus.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/05/31 13:38:42 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/07/23 23:50:45 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Bounded in-memory event bus for emitted command events.
12
13use crate::envelope::EventEnvelope;
14use crate::operational_journal::{FileOperationalJournal, OperationalJournalRecord};
15use crate::TraceContext;
16use parking_lot::Mutex;
17use serde::ser::SerializeSeq;
18use serde::{Serialize, Serializer};
19use std::collections::VecDeque;
20use std::sync::Arc;
21
22const MAX_RETAINED_EVENTS: usize = 10_000;
23/// Default aggregate memory budget for process-local emitted events.
24pub const DEFAULT_EVENT_BUS_MAX_BYTES: usize = 16 * 1024 * 1024;
25
26/// Point-in-time pressure metrics for a process-local event bus.
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct EventBusStats {
29    /// Retained event envelopes.
30    pub event_count: usize,
31    /// Estimated bytes retained by the current snapshot.
32    pub used_bytes: usize,
33    /// Highest retained byte count observed by this bus.
34    pub peak_bytes: usize,
35    /// Events evicted to maintain count or byte limits.
36    pub evictions: u64,
37    /// Individual events too large for the configured byte budget.
38    pub rejections: u64,
39    /// Aggregate configured byte budget.
40    pub max_bytes: usize,
41}
42
43/// Shared immutable point-in-time view of emitted events.
44#[derive(Clone, Debug)]
45pub struct EventBusSnapshot {
46    events: Arc<VecDeque<Arc<OperationalJournalRecord>>>,
47}
48
49impl EventBusSnapshot {
50    /// Returns the number of events captured by this snapshot.
51    #[must_use]
52    pub fn len(&self) -> usize {
53        self.events.len()
54    }
55
56    /// Reports whether this snapshot contains no events.
57    #[must_use]
58    pub fn is_empty(&self) -> bool {
59        self.events.is_empty()
60    }
61
62    /// Iterates over at most the newest `limit` events without cloning them.
63    pub fn recent(&self, limit: usize) -> impl Iterator<Item = &EventEnvelope> {
64        let start = self.events.len().saturating_sub(limit);
65        self.events.iter().skip(start).filter_map(event_record)
66    }
67}
68
69impl Serialize for EventBusSnapshot {
70    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
71    where
72        S: Serializer,
73    {
74        let mut sequence = serializer.serialize_seq(Some(self.events.len()))?;
75        for event in self.events.iter().filter_map(event_record) {
76            sequence.serialize_element(event)?;
77        }
78        sequence.end()
79    }
80}
81
82#[derive(Debug, Clone)]
83struct EventBusState {
84    events: Arc<VecDeque<Arc<OperationalJournalRecord>>>,
85    used_bytes: usize,
86    peak_bytes: usize,
87    evictions: u64,
88    rejections: u64,
89    max_bytes: usize,
90}
91
92impl EventBusState {
93    fn new(max_bytes: usize) -> Self {
94        Self {
95            events: Arc::default(),
96            used_bytes: 0,
97            peak_bytes: 0,
98            evictions: 0,
99            rejections: 0,
100            max_bytes: max_bytes.max(1),
101        }
102    }
103
104    fn push_shared(&mut self, record: Arc<OperationalJournalRecord>) {
105        let Some(event) = event_record(&record) else {
106            self.rejections = self.rejections.saturating_add(1);
107            return;
108        };
109        let bytes = event_retained_bytes(event);
110        if bytes > self.max_bytes {
111            self.rejections = self.rejections.saturating_add(1);
112            return;
113        }
114        while self.events.len() >= MAX_RETAINED_EVENTS {
115            self.pop_front();
116        }
117        while self.used_bytes.saturating_add(bytes) > self.max_bytes {
118            if !self.pop_front() {
119                self.rejections = self.rejections.saturating_add(1);
120                return;
121            }
122        }
123        Arc::make_mut(&mut self.events).push_back(record);
124        self.used_bytes = self.used_bytes.saturating_add(bytes);
125        self.peak_bytes = self.peak_bytes.max(self.used_bytes);
126    }
127
128    fn pop_front(&mut self) -> bool {
129        let Some(event) = Arc::make_mut(&mut self.events).pop_front() else {
130            return false;
131        };
132        if let Some(event) = event_record(&event) {
133            self.used_bytes = self.used_bytes.saturating_sub(event_retained_bytes(event));
134        }
135        self.evictions = self.evictions.saturating_add(1);
136        true
137    }
138
139    fn replace(&mut self, events: Vec<Arc<OperationalJournalRecord>>) {
140        self.clear();
141        for event in events {
142            self.push_shared(event);
143        }
144    }
145
146    fn clear(&mut self) {
147        self.events = Arc::default();
148        self.used_bytes = 0;
149    }
150
151    fn stats(&self) -> EventBusStats {
152        EventBusStats {
153            event_count: self.events.len(),
154            used_bytes: self.used_bytes,
155            peak_bytes: self.peak_bytes,
156            evictions: self.evictions,
157            rejections: self.rejections,
158            max_bytes: self.max_bytes,
159        }
160    }
161}
162
163/// Bounded process-local store of recently emitted Runtime events.
164#[derive(Debug)]
165pub struct EventBus {
166    state: Mutex<EventBusState>,
167    journal: Mutex<Option<Arc<FileOperationalJournal>>>,
168    journal_error: Mutex<Option<String>>,
169}
170
171impl Default for EventBus {
172    fn default() -> Self {
173        Self::with_max_bytes(DEFAULT_EVENT_BUS_MAX_BYTES)
174    }
175}
176
177impl Clone for EventBus {
178    fn clone(&self) -> Self {
179        Self {
180            state: Mutex::new(self.state.lock().clone()),
181            journal: Mutex::new(self.journal.lock().clone()),
182            journal_error: Mutex::new(self.journal_error.lock().clone()),
183        }
184    }
185}
186
187impl EventBus {
188    /// Creates an empty event bus.
189    pub fn new() -> Self {
190        Self::default()
191    }
192
193    /// Creates an empty event bus with an aggregate retained-byte budget.
194    pub fn with_max_bytes(max_bytes: usize) -> Self {
195        Self {
196            state: Mutex::new(EventBusState::new(max_bytes)),
197            journal: Mutex::new(None),
198            journal_error: Mutex::new(None),
199        }
200    }
201
202    /// Attaches a durable journal and shares its retained immutable event records.
203    pub fn attach_journal(&self, journal: Arc<FileOperationalJournal>) {
204        self.state.lock().replace(journal.shared_event_records());
205        *self.journal.lock() = Some(journal);
206        *self.journal_error.lock() = None;
207    }
208
209    /// Returns the last durable journal failure, when persistence degraded.
210    pub fn durability_error(&self) -> Option<String> {
211        self.journal_error.lock().clone()
212    }
213
214    /// Appends one event, sharing its record with an attached durable journal.
215    pub fn emit(&self, event: EventEnvelope) {
216        let event = Arc::new(OperationalJournalRecord::Event(event));
217        self.persist(Arc::clone(&event));
218        self.state.lock().push_shared(event);
219    }
220
221    /// Appends events and retains them within the count and aggregate-byte bounds.
222    pub fn emit_many(&self, events: Vec<EventEnvelope>) {
223        let events = events
224            .into_iter()
225            .map(|event| Arc::new(OperationalJournalRecord::Event(event)))
226            .collect::<Vec<_>>();
227        for event in &events {
228            self.persist(Arc::clone(event));
229        }
230        let mut state = self.state.lock();
231        for event in events {
232            state.push_shared(event);
233        }
234    }
235
236    /// Returns the current number of retained events.
237    pub fn len(&self) -> usize {
238        self.state.lock().events.len()
239    }
240
241    /// Reports whether no events are retained.
242    pub fn is_empty(&self) -> bool {
243        self.state.lock().events.is_empty()
244    }
245
246    /// Returns a point-in-time copy of retained events.
247    pub fn events(&self) -> Vec<EventEnvelope> {
248        self.snapshot().recent(usize::MAX).cloned().collect()
249    }
250
251    /// Returns a shared immutable snapshot without cloning event fields.
252    pub fn snapshot(&self) -> EventBusSnapshot {
253        EventBusSnapshot {
254            events: Arc::clone(&self.state.lock().events),
255        }
256    }
257
258    /// Returns current count, byte-pressure, eviction, and rejection metrics.
259    pub fn stats(&self) -> EventBusStats {
260        self.state.lock().stats()
261    }
262
263    /// Removes all retained events.
264    pub fn clear(&self) {
265        self.state.lock().clear();
266    }
267
268    fn persist(&self, event: Arc<OperationalJournalRecord>) {
269        if let Some(journal) = self.journal.lock().clone() {
270            if let Err(error) = journal.append_shared_event(event) {
271                *self.journal_error.lock() = Some(crate::redact_text(&format!("{error:?}")));
272            }
273        }
274    }
275}
276
277fn event_record(record: &Arc<OperationalJournalRecord>) -> Option<&EventEnvelope> {
278    match record.as_ref() {
279        OperationalJournalRecord::Event(event) => Some(event),
280        OperationalJournalRecord::Audit(_) => None,
281    }
282}
283
284fn event_retained_bytes(event: &EventEnvelope) -> usize {
285    std::mem::size_of::<EventEnvelope>()
286        .saturating_add(event.event_name.as_str().len())
287        .saturating_add(event.event_id.len())
288        .saturating_add(event.app_id.as_str().len())
289        .saturating_add(event.node_id.as_str().len())
290        .saturating_add(event.payload.len())
291        .saturating_add(trace_retained_bytes(event.trace.as_ref()))
292}
293
294fn trace_retained_bytes(trace: Option<&TraceContext>) -> usize {
295    let Some(trace) = trace else {
296        return 0;
297    };
298    std::mem::size_of::<TraceContext>()
299        .saturating_add(trace.trace_id.len())
300        .saturating_add(trace.span_id.len())
301        .saturating_add(trace.parent_span_id.as_ref().map_or(0, String::len))
302        .saturating_add(trace.originating_core_id.as_str().len())
303        .saturating_add(trace.current_core_id.as_str().len())
304        .saturating_add(trace.tenant_id.as_str().len())
305        .saturating_add(trace.command_id.as_ref().map_or(0, String::len))
306}
307
308#[cfg(test)]
309#[path = "event_bus_tests.rs"]
310mod tests;