Skip to main content

oxicode_sdk/observability/
event_store.rs

1//! Event sourcing — append-only event store with streams and replay.
2
3use serde::{Deserialize, Serialize};
4use std::collections::HashMap;
5use std::sync::atomic::{AtomicU64, Ordering};
6use tokio::sync::broadcast;
7
8/// Configuration for the event store.
9#[derive(Debug, Clone)]
10pub struct EventStoreConfig {
11    /// Maximum entries per stream before eviction.
12    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/// A stored domain event.
24#[derive(Debug, Clone, Serialize, Deserialize)]
25pub struct StoredEvent {
26    /// Monotonically increasing sequence number.
27    pub sequence: u64,
28    /// Logical stream this event belongs to.
29    pub stream_id: String,
30    /// Event type name.
31    pub event_type: String,
32    /// Arbitrary payload.
33    pub payload: serde_json::Value,
34    /// Wall-clock time (ms since epoch).
35    pub timestamp_ms: u64,
36}
37
38/// Filter for querying events.
39#[derive(Debug, Clone, Default)]
40pub struct EventQuery {
41    /// Filter by stream.
42    pub stream_id: Option<String>,
43    /// Filter by event type.
44    pub event_type: Option<String>,
45    /// Filter by minimum sequence.
46    pub min_sequence: Option<u64>,
47    /// Filter by maximum sequence.
48    pub max_sequence: Option<u64>,
49}
50
51/// Append-only event store with per-stream indexing.
52///
53/// Internally uses a flat sequence (no nested streams), filtered at query time.
54pub struct EventStore {
55    /// Monotonically increasing global sequence counter.
56    sequence: AtomicU64,
57    /// Append-only flat event log.
58    events: parking_lot::RwLock<Vec<StoredEvent>>,
59    /// Per-stream sequence cursor.
60    stream_cursors: parking_lot::RwLock<HashMap<String, u64>>,
61    /// Broadcast channel for subscribers.
62    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    /// Create a new event store with the given configuration.
74    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    /// Append an event to a stream. Returns the assigned sequence number.
86    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            // Simple eviction: keep at most max_entries_per_stream overall
108            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    /// Query events matching the filter.
125    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    /// Replay all events for a particular stream, in order.
150    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    /// Subscribe to new events.
160    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}