allsource_core/infrastructure/persistence/
index.rs1use crate::error::Result;
2use dashmap::DashMap;
3use std::sync::Arc;
4use uuid::Uuid;
5
6#[derive(Debug, Clone)]
8pub struct IndexEntry {
9 pub event_id: Uuid,
10 pub offset: usize,
11 pub timestamp: chrono::DateTime<chrono::Utc>,
12}
13
14pub struct EventIndex {
16 entity_index: Arc<DashMap<String, Vec<IndexEntry>>>,
18
19 type_index: Arc<DashMap<String, Vec<IndexEntry>>>,
21
22 id_index: Arc<DashMap<Uuid, usize>>,
24
25 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 #[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 self.entity_index
57 .entry(entity_id.to_string())
58 .or_default()
59 .push(entry.clone());
60
61 self.type_index
63 .entry(event_type.to_string())
64 .or_default()
65 .push(entry.clone());
66
67 self.id_index.insert(event_id, offset);
69
70 let mut total = self.total_events.write();
72 *total += 1;
73
74 Ok(())
75 }
76
77 #[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 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 #[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 #[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 #[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 pub fn get_all_entities(&self) -> Vec<String> {
130 self.entity_index.iter().map(|e| e.key().clone()).collect()
131 }
132
133 pub fn get_all_types(&self) -> Vec<String> {
135 self.type_index.iter().map(|e| e.key().clone()).collect()
136 }
137
138 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 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 let entries = index.get_by_type_prefix("");
245 assert_eq!(entries.len(), 3);
246 }
247}