appcore_core/
event_bus.rs1use 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;
23pub const DEFAULT_EVENT_BUS_MAX_BYTES: usize = 16 * 1024 * 1024;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct EventBusStats {
29 pub event_count: usize,
31 pub used_bytes: usize,
33 pub peak_bytes: usize,
35 pub evictions: u64,
37 pub rejections: u64,
39 pub max_bytes: usize,
41}
42
43#[derive(Clone, Debug)]
45pub struct EventBusSnapshot {
46 events: Arc<VecDeque<Arc<OperationalJournalRecord>>>,
47}
48
49impl EventBusSnapshot {
50 #[must_use]
52 pub fn len(&self) -> usize {
53 self.events.len()
54 }
55
56 #[must_use]
58 pub fn is_empty(&self) -> bool {
59 self.events.is_empty()
60 }
61
62 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#[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 pub fn new() -> Self {
190 Self::default()
191 }
192
193 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 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 pub fn durability_error(&self) -> Option<String> {
211 self.journal_error.lock().clone()
212 }
213
214 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 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 pub fn len(&self) -> usize {
238 self.state.lock().events.len()
239 }
240
241 pub fn is_empty(&self) -> bool {
243 self.state.lock().events.is_empty()
244 }
245
246 pub fn events(&self) -> Vec<EventEnvelope> {
248 self.snapshot().recent(usize::MAX).cloned().collect()
249 }
250
251 pub fn snapshot(&self) -> EventBusSnapshot {
253 EventBusSnapshot {
254 events: Arc::clone(&self.state.lock().events),
255 }
256 }
257
258 pub fn stats(&self) -> EventBusStats {
260 self.state.lock().stats()
261 }
262
263 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;