appcore_core/
event_bus.rs1use crate::envelope::EventEnvelope;
14use crate::operational_journal::FileOperationalJournal;
15use parking_lot::Mutex;
16use std::sync::Arc;
17
18#[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 pub fn new() -> Self {
40 Self::default()
41 }
42
43 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 pub fn durability_error(&self) -> Option<String> {
56 self.journal_error.lock().clone()
57 }
58
59 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 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 pub fn len(&self) -> usize {
89 self.events.lock().len()
90 }
91
92 pub fn is_empty(&self) -> bool {
94 self.events.lock().is_empty()
95 }
96
97 pub fn events(&self) -> Vec<EventEnvelope> {
99 self.events.lock().clone()
100 }
101
102 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;