oxicode_sdk/observability/
event_store.rs1use serde::{Deserialize, Serialize};
4use std::collections::HashMap;
5use std::sync::atomic::{AtomicU64, Ordering};
6use tokio::sync::broadcast;
7
8#[derive(Debug, Clone)]
10pub struct EventStoreConfig {
11 pub max_entries_per_stream: usize,
13}
14
15impl Default for EventStoreConfig {
16 fn default() -> Self {
17 Self {
18 max_entries_per_stream: 10_000,
19 }
20 }
21}
22
23#[derive(Debug, Clone, Serialize, Deserialize)]
25pub struct StoredEvent {
26 pub sequence: u64,
28 pub stream_id: String,
30 pub event_type: String,
32 pub payload: serde_json::Value,
34 pub timestamp_ms: u64,
36}
37
38#[derive(Debug, Clone, Default)]
40pub struct EventQuery {
41 pub stream_id: Option<String>,
43 pub event_type: Option<String>,
45 pub min_sequence: Option<u64>,
47 pub max_sequence: Option<u64>,
49}
50
51pub struct EventStore {
55 sequence: AtomicU64,
57 events: parking_lot::RwLock<Vec<StoredEvent>>,
59 stream_cursors: parking_lot::RwLock<HashMap<String, u64>>,
61 tx: broadcast::Sender<StoredEvent>,
63 config: EventStoreConfig,
64}
65
66impl Default for EventStore {
67 fn default() -> Self {
68 Self::new(EventStoreConfig::default())
69 }
70}
71
72impl EventStore {
73 pub fn new(config: EventStoreConfig) -> Self {
75 let (tx, _) = broadcast::channel(256);
76 Self {
77 sequence: AtomicU64::new(1),
78 events: parking_lot::RwLock::new(Vec::new()),
79 stream_cursors: parking_lot::RwLock::new(HashMap::new()),
80 tx,
81 config,
82 }
83 }
84
85 pub fn append(
87 &self,
88 stream_id: impl Into<String>,
89 event_type: impl Into<String>,
90 payload: serde_json::Value,
91 ) -> u64 {
92 let sequence = self.sequence.fetch_add(1, Ordering::SeqCst);
93 let stream_id = stream_id.into();
94
95 let event = StoredEvent {
96 sequence,
97 stream_id: stream_id.clone(),
98 event_type: event_type.into(),
99 payload,
100 timestamp_ms: now_ms(),
101 };
102
103 {
104 let mut events = self.events.write();
105 events.push(event.clone());
106
107 if events.len() > self.config.max_entries_per_stream {
109 let drain_count = events.len() - self.config.max_entries_per_stream;
110 events.drain(0..drain_count);
111 }
112 }
113
114 {
115 let mut cursors = self.stream_cursors.write();
116 let _current = cursors.get(&stream_id).copied().unwrap_or(0);
117 cursors.insert(stream_id, sequence);
118 }
119
120 let _ = self.tx.send(event);
121 sequence
122 }
123
124 pub fn query(&self, query: EventQuery) -> Vec<StoredEvent> {
126 self.events
127 .read()
128 .iter()
129 .filter(|e| {
130 if let Some(ref sid) = query.stream_id {
131 &e.stream_id == sid
132 } else {
133 true
134 }
135 })
136 .filter(|e| {
137 if let Some(ref et) = query.event_type {
138 &e.event_type == et
139 } else {
140 true
141 }
142 })
143 .filter(|e| query.min_sequence.map(|m| e.sequence >= m).unwrap_or(true))
144 .filter(|e| query.max_sequence.map(|m| e.sequence <= m).unwrap_or(true))
145 .cloned()
146 .collect()
147 }
148
149 pub fn replay(&self, stream_id: &str) -> Vec<StoredEvent> {
151 self.events
152 .read()
153 .iter()
154 .filter(|e| e.stream_id == stream_id)
155 .cloned()
156 .collect()
157 }
158
159 pub fn subscribe(&self) -> broadcast::Receiver<StoredEvent> {
161 self.tx.subscribe()
162 }
163}
164
165fn now_ms() -> u64 {
166 std::time::SystemTime::now()
167 .duration_since(std::time::UNIX_EPOCH)
168 .map(|d| d.as_millis() as u64)
169 .unwrap_or(0)
170}
171
172#[cfg(test)]
173mod tests {
174 use super::*;
175
176 #[test]
177 fn event_store_append_returns_sequence() {
178 let store = EventStore::default();
179 let seq = store.append("s1", "Created", serde_json::json!({}));
180 assert_eq!(seq, 1);
181
182 let seq2 = store.append("s1", "Updated", serde_json::json!({}));
183 assert_eq!(seq2, 2);
184 }
185
186 #[test]
187 fn event_store_query_by_stream() {
188 let store = EventStore::default();
189 store.append("s1", "A", serde_json::json!({"n": 1}));
190 store.append("s2", "B", serde_json::json!({"n": 2}));
191 store.append("s1", "C", serde_json::json!({"n": 3}));
192
193 let results = store
194 .query(EventQuery {
195 stream_id: Some("s1".into()),
196 ..Default::default()
197 })
198 .into_iter()
199 .map(|e| e.event_type)
200 .collect::<Vec<_>>();
201
202 assert_eq!(results.len(), 2);
203 assert!(results.contains(&"A".into()));
204 assert!(results.contains(&"C".into()));
205 }
206
207 #[test]
208 fn event_store_query_by_event_type() {
209 let store = EventStore::default();
210 store.append("s1", "Click", serde_json::json!({}));
211 store.append("s1", "Hover", serde_json::json!({}));
212 store.append("s1", "Click", serde_json::json!({}));
213
214 let results = store
215 .query(EventQuery {
216 event_type: Some("Click".into()),
217 ..Default::default()
218 })
219 .len();
220
221 assert_eq!(results, 2);
222 }
223
224 #[test]
225 fn event_store_replay() {
226 let store = EventStore::default();
227 store.append("order-1", "Created", serde_json::json!({"id": 1}));
228 store.append("order-1", "Paid", serde_json::json!({"amount": 100}));
229 store.append("order-2", "Created", serde_json::json!({"id": 2}));
230
231 let events = store.replay("order-1");
232 assert_eq!(events.len(), 2);
233 assert_eq!(events[0].event_type, "Created");
234 assert_eq!(events[1].event_type, "Paid");
235 }
236}