Skip to main content

allsource_core/infrastructure/persistence/
index.rs

1use crate::error::Result;
2use dashmap::DashMap;
3use std::sync::Arc;
4use uuid::Uuid;
5
6/// Event index entry
7#[derive(Debug, Clone)]
8pub struct IndexEntry {
9    pub event_id: Uuid,
10    pub offset: usize,
11    pub timestamp: chrono::DateTime<chrono::Utc>,
12}
13
14/// High-performance concurrent index for fast event lookups
15pub struct EventIndex {
16    /// Index by entity_id -> list of event entries
17    entity_index: Arc<DashMap<String, Vec<IndexEntry>>>,
18
19    /// Index by event_type -> list of event entries
20    type_index: Arc<DashMap<String, Vec<IndexEntry>>>,
21
22    /// Index by event_id -> offset (for direct lookups)
23    id_index: Arc<DashMap<Uuid, usize>>,
24
25    /// Total indexed events
26    total_events: parking_lot::RwLock<usize>,
27}
28
29impl EventIndex {
30    pub fn new() -> Self {
31        Self {
32            entity_index: Arc::new(DashMap::new()),
33            type_index: Arc::new(DashMap::new()),
34            id_index: Arc::new(DashMap::new()),
35            total_events: parking_lot::RwLock::new(0),
36        }
37    }
38
39    /// Add an event to all relevant indices
40    #[cfg_attr(feature = "hotpath", hotpath::measure)]
41    pub fn index_event(
42        &self,
43        event_id: Uuid,
44        entity_id: &str,
45        event_type: &str,
46        timestamp: chrono::DateTime<chrono::Utc>,
47        offset: usize,
48    ) -> Result<()> {
49        let entry = IndexEntry {
50            event_id,
51            offset,
52            timestamp,
53        };
54
55        // Index by entity_id
56        self.entity_index
57            .entry(entity_id.to_string())
58            .or_default()
59            .push(entry.clone());
60
61        // Index by event_type
62        self.type_index
63            .entry(event_type.to_string())
64            .or_default()
65            .push(entry.clone());
66
67        // Index by event_id
68        self.id_index.insert(event_id, offset);
69
70        // Increment total
71        let mut total = self.total_events.write();
72        *total += 1;
73
74        Ok(())
75    }
76
77    /// Get all event offsets for an entity
78    #[cfg_attr(feature = "hotpath", hotpath::measure)]
79    pub fn get_by_entity(&self, entity_id: &str) -> Option<Vec<IndexEntry>> {
80        self.entity_index
81            .get(entity_id)
82            .map(|entries| entries.clone())
83    }
84
85    /// Refuse an oversized entity before allocating its offset vector.
86    pub(crate) fn get_by_entity_bounded(
87        &self,
88        entity_id: &str,
89        max: usize,
90    ) -> Result<Vec<IndexEntry>> {
91        let Some(entries) = self.entity_index.get(entity_id) else {
92            return Ok(Vec::new());
93        };
94        if entries.len() > max {
95            return Err(crate::error::AllSourceError::StorageError(
96                "Strict retained entity exceeds its event budget".into(),
97            ));
98        }
99        Ok(entries.clone())
100    }
101
102    /// Get all event offsets for an event type
103    #[cfg_attr(feature = "hotpath", hotpath::measure)]
104    pub fn get_by_type(&self, event_type: &str) -> Option<Vec<IndexEntry>> {
105        self.type_index
106            .get(event_type)
107            .map(|entries| entries.clone())
108    }
109
110    /// Get event offset by ID
111    #[cfg_attr(feature = "hotpath", hotpath::measure)]
112    pub fn get_by_id(&self, event_id: &Uuid) -> Option<usize> {
113        self.id_index.get(event_id).map(|offset| *offset)
114    }
115
116    /// Get all event entries matching a type prefix (e.g., "index." matches "index.created")
117    #[cfg_attr(feature = "hotpath", hotpath::measure)]
118    pub fn get_by_type_prefix(&self, prefix: &str) -> Vec<IndexEntry> {
119        let mut entries = Vec::new();
120        for item in self.type_index.iter() {
121            if item.key().starts_with(prefix) {
122                entries.extend(item.value().clone());
123            }
124        }
125        entries
126    }
127
128    /// Get all entities being tracked
129    pub fn get_all_entities(&self) -> Vec<String> {
130        self.entity_index.iter().map(|e| e.key().clone()).collect()
131    }
132
133    /// Get all event types
134    pub fn get_all_types(&self) -> Vec<String> {
135        self.type_index.iter().map(|e| e.key().clone()).collect()
136    }
137
138    /// Get statistics
139    pub fn stats(&self) -> IndexStats {
140        IndexStats {
141            total_events: *self.total_events.read(),
142            total_entities: self.entity_index.len(),
143            total_event_types: self.type_index.len(),
144        }
145    }
146
147    /// Clear all indices (useful for testing)
148    pub fn clear(&self) {
149        self.entity_index.clear();
150        self.type_index.clear();
151        self.id_index.clear();
152        let mut total = self.total_events.write();
153        *total = 0;
154    }
155}
156
157impl Default for EventIndex {
158    fn default() -> Self {
159        Self::new()
160    }
161}
162
163#[derive(Debug, Clone, serde::Serialize)]
164pub struct IndexStats {
165    pub total_events: usize,
166    pub total_entities: usize,
167    pub total_event_types: usize,
168}
169
170#[cfg(test)]
171mod tests {
172    use super::*;
173
174    #[test]
175    fn test_index_event() {
176        let index = EventIndex::new();
177        let event_id = Uuid::new_v4();
178        let timestamp = chrono::Utc::now();
179
180        index
181            .index_event(event_id, "user-123", "user.created", timestamp, 0)
182            .unwrap();
183
184        assert_eq!(index.stats().total_events, 1);
185        assert_eq!(index.stats().total_entities, 1);
186        assert_eq!(index.stats().total_event_types, 1);
187    }
188
189    #[test]
190    fn test_get_by_entity() {
191        let index = EventIndex::new();
192        let event_id = Uuid::new_v4();
193        let timestamp = chrono::Utc::now();
194
195        index
196            .index_event(event_id, "user-123", "user.created", timestamp, 0)
197            .unwrap();
198
199        let entries = index.get_by_entity("user-123").unwrap();
200        assert_eq!(entries.len(), 1);
201        assert_eq!(entries[0].event_id, event_id);
202    }
203
204    #[test]
205    fn test_get_by_type() {
206        let index = EventIndex::new();
207        let event_id = Uuid::new_v4();
208        let timestamp = chrono::Utc::now();
209
210        index
211            .index_event(event_id, "user-123", "user.created", timestamp, 0)
212            .unwrap();
213
214        let entries = index.get_by_type("user.created").unwrap();
215        assert_eq!(entries.len(), 1);
216        assert_eq!(entries[0].event_id, event_id);
217    }
218
219    #[test]
220    fn test_get_by_type_prefix() {
221        let index = EventIndex::new();
222        let ts = chrono::Utc::now();
223
224        index
225            .index_event(Uuid::new_v4(), "e-1", "index.created", ts, 0)
226            .unwrap();
227        index
228            .index_event(Uuid::new_v4(), "e-2", "index.updated", ts, 1)
229            .unwrap();
230        index
231            .index_event(Uuid::new_v4(), "e-3", "trade.created", ts, 2)
232            .unwrap();
233
234        let entries = index.get_by_type_prefix("index.");
235        assert_eq!(entries.len(), 2);
236
237        let entries = index.get_by_type_prefix("trade.");
238        assert_eq!(entries.len(), 1);
239
240        let entries = index.get_by_type_prefix("nonexistent.");
241        assert_eq!(entries.len(), 0);
242
243        // Empty prefix matches all
244        let entries = index.get_by_type_prefix("");
245        assert_eq!(entries.len(), 3);
246    }
247}