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;
15use parking_lot::Mutex;
16use std::sync::Arc;
17
18/// Bounded process-local store of recently emitted Runtime events.
19#[derive(Debug, Default)]
20pub struct EventBus {
21    events: Mutex<Vec<EventEnvelope>>,
22    journal: Mutex<Option<Arc<FileOperationalJournal>>>,
23    journal_error: Mutex<Option<String>>,
24}
25
26impl Clone for EventBus {
27    fn clone(&self) -> Self {
28        let guard = self.events.lock();
29        Self {
30            events: Mutex::new(guard.clone()),
31            journal: Mutex::new(self.journal.lock().clone()),
32            journal_error: Mutex::new(self.journal_error.lock().clone()),
33        }
34    }
35}
36
37impl EventBus {
38    /// Creates an empty event bus.
39    pub fn new() -> Self {
40        Self::default()
41    }
42
43    /// Attaches a durable journal and loads its retained event envelopes.
44    pub fn attach_journal(&self, journal: Arc<FileOperationalJournal>) {
45        let mut events = journal.events();
46        if events.len() > 10_000 {
47            events.drain(..events.len() - 10_000);
48        }
49        *self.events.lock() = events;
50        *self.journal.lock() = Some(journal);
51        *self.journal_error.lock() = None;
52    }
53
54    /// Returns the last durable journal failure, when persistence degraded.
55    pub fn durability_error(&self) -> Option<String> {
56        self.journal_error.lock().clone()
57    }
58
59    /// Appends one event and evicts the oldest item at the configured bound.
60    pub fn emit(&self, event: EventEnvelope) {
61        self.persist(&event);
62        let mut guard = self.events.lock();
63        if guard.len() >= 10000 {
64            let to_remove = guard.len().saturating_sub(9999);
65            if to_remove > 0 {
66                guard.drain(0..to_remove);
67            }
68        }
69        guard.push(event);
70    }
71
72    /// Appends multiple events and retains at most 10,000 recent items.
73    pub fn emit_many(&self, events: Vec<EventEnvelope>) {
74        for event in &events {
75            self.persist(event);
76        }
77        let mut guard = self.events.lock();
78        guard.extend(events);
79        if guard.len() > 10000 {
80            let to_remove = guard.len().saturating_sub(10000);
81            if to_remove > 0 {
82                guard.drain(0..to_remove);
83            }
84        }
85    }
86
87    /// Returns the current number of retained events.
88    pub fn len(&self) -> usize {
89        self.events.lock().len()
90    }
91
92    /// Reports whether no events are retained.
93    pub fn is_empty(&self) -> bool {
94        self.events.lock().is_empty()
95    }
96
97    /// Returns a point-in-time copy of retained events.
98    pub fn events(&self) -> Vec<EventEnvelope> {
99        self.events.lock().clone()
100    }
101
102    /// Removes all retained events.
103    pub fn clear(&self) {
104        self.events.lock().clear();
105    }
106
107    fn persist(&self, event: &EventEnvelope) {
108        if let Some(journal) = self.journal.lock().clone() {
109            if let Err(error) = journal.append_event(event.clone()) {
110                *self.journal_error.lock() = Some(crate::redact_text(&format!("{error:?}")));
111            }
112        }
113    }
114}
115
116#[cfg(test)]
117#[path = "event_bus_tests.rs"]
118mod tests;