pjson_rs/infrastructure/adapters/
generic_store.rs1use dashmap::DashMap;
32use std::{hash::Hash, sync::Arc};
33
34const MAX_PREALLOC_SIZE: usize = 1024;
39
40#[derive(Debug)]
60pub struct InMemoryStore<K, V>
61where
62 K: Eq + Hash + Clone + Send + Sync,
63 V: Clone + Send + Sync,
64{
65 data: Arc<DashMap<K, V>>,
66}
67
68impl<K, V> InMemoryStore<K, V>
69where
70 K: Eq + Hash + Clone + Send + Sync,
71 V: Clone + Send + Sync,
72{
73 pub fn new() -> Self {
75 Self {
76 data: Arc::new(DashMap::new()),
77 }
78 }
79
80 pub fn count(&self) -> usize {
82 self.data.len()
83 }
84
85 pub fn clear(&self) {
87 self.data.clear();
88 }
89
90 pub fn all_keys(&self) -> Vec<K> {
97 self.data.iter().map(|entry| entry.key().clone()).collect()
98 }
99
100 pub fn get(&self, key: &K) -> Option<V> {
107 self.data.get(key).map(|entry| entry.value().clone())
108 }
109
110 pub fn insert(&self, key: K, value: V) -> Option<V> {
112 self.data.insert(key, value)
113 }
114
115 pub fn remove(&self, key: &K) -> Option<V> {
117 self.data.remove(key).map(|(_k, v)| v)
118 }
119
120 pub fn filter<P>(&self, predicate: P) -> Vec<V>
128 where
129 P: Fn(&V) -> bool,
130 {
131 self.data
132 .iter()
133 .filter(|entry| predicate(entry.value()))
134 .map(|entry| entry.value().clone())
135 .collect()
136 }
137
138 pub fn filter_limited<P>(
173 &self,
174 predicate: P,
175 result_limit: usize,
176 scan_limit: usize,
177 ) -> (Vec<V>, bool)
178 where
179 P: Fn(&V) -> bool,
180 {
181 let mut results = Vec::with_capacity(result_limit.min(MAX_PREALLOC_SIZE));
182 let mut limit_reached = false;
183
184 for (scanned, entry) in self.data.iter().enumerate() {
185 if scanned >= scan_limit {
187 limit_reached = true;
188 break;
189 }
190
191 if predicate(entry.value()) {
192 results.push(entry.value().clone());
193 if results.len() >= result_limit {
194 limit_reached = true;
195 break;
196 }
197 }
198 }
199
200 (results, limit_reached)
201 }
202
203 pub fn contains_key(&self, key: &K) -> bool {
210 self.data.contains_key(key)
211 }
212
213 pub fn is_empty(&self) -> bool {
215 self.data.is_empty()
216 }
217
218 pub fn update_with<F, R>(&self, key: &K, f: F) -> Option<R>
236 where
237 F: FnOnce(&mut V) -> R,
238 {
239 self.data.get_mut(key).map(|mut entry| f(entry.value_mut()))
240 }
241
242 pub fn iter(&self) -> impl Iterator<Item = dashmap::mapref::multiple::RefMulti<'_, K, V>> {
253 self.data.iter()
254 }
255}
256
257impl<K, V> Default for InMemoryStore<K, V>
258where
259 K: Eq + Hash + Clone + Send + Sync,
260 V: Clone + Send + Sync,
261{
262 fn default() -> Self {
263 Self::new()
264 }
265}
266
267impl<K, V> Clone for InMemoryStore<K, V>
268where
269 K: Eq + Hash + Clone + Send + Sync,
270 V: Clone + Send + Sync,
271{
272 fn clone(&self) -> Self {
273 Self {
274 data: Arc::clone(&self.data),
275 }
276 }
277}
278
279use crate::domain::{
281 aggregates::StreamSession,
282 entities::Stream,
283 value_objects::{SessionId, StreamId},
284};
285
286pub type SessionStore = InMemoryStore<SessionId, StreamSession>;
288
289pub type StreamStore = InMemoryStore<StreamId, Stream>;
291
292#[cfg(test)]
293mod tests {
294 use super::*;
295
296 #[test]
297 fn test_basic_operations() {
298 let store: InMemoryStore<String, i32> = InMemoryStore::new();
299
300 assert!(store.is_empty());
301 assert_eq!(store.count(), 0);
302
303 store.insert("a".to_string(), 1);
304 store.insert("b".to_string(), 2);
305
306 assert_eq!(store.count(), 2);
307 assert_eq!(store.get(&"a".to_string()), Some(1));
308 assert_eq!(store.get(&"c".to_string()), None);
309
310 let keys = store.all_keys();
311 assert_eq!(keys.len(), 2);
312
313 store.remove(&"a".to_string());
314 assert_eq!(store.count(), 1);
315
316 store.clear();
317 assert!(store.is_empty());
318 }
319
320 #[test]
321 fn test_filter() {
322 let store: InMemoryStore<String, i32> = InMemoryStore::new();
323
324 store.insert("a".to_string(), 1);
325 store.insert("b".to_string(), 2);
326 store.insert("c".to_string(), 3);
327
328 let evens = store.filter(|v| v % 2 == 0);
329 assert_eq!(evens, vec![2]);
330 }
331
332 #[test]
333 fn test_filter_limited_returns_at_most_limit_items() {
334 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
335
336 for i in 0..100 {
337 store.insert(i, i);
338 }
339
340 let (results, limit_reached) = store.filter_limited(|_| true, 10, 1000);
341
342 assert_eq!(results.len(), 10);
343 assert!(limit_reached);
344 }
345
346 #[test]
347 fn test_filter_limited_sets_limit_reached_when_scan_exceeded() {
348 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
349
350 for i in 0..100 {
351 store.insert(i, i);
352 }
353
354 let (results, limit_reached) = store.filter_limited(|v| *v > 1000, 100, 50);
355
356 assert!(results.is_empty());
357 assert!(limit_reached);
358 }
359
360 #[test]
361 fn test_filter_limited_sets_limit_reached_when_results_exceeded() {
362 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
363
364 for i in 0..100 {
365 store.insert(i, i);
366 }
367
368 let (results, limit_reached) = store.filter_limited(|_| true, 5, 1000);
369
370 assert_eq!(results.len(), 5);
371 assert!(limit_reached);
372 }
373
374 #[test]
375 fn test_filter_limited_empty_store() {
376 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
377
378 let (results, limit_reached) = store.filter_limited(|_| true, 10, 100);
379
380 assert!(results.is_empty());
381 assert!(!limit_reached);
382 }
383
384 #[test]
385 fn test_filter_limited_no_matches() {
386 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
387
388 for i in 0..10 {
389 store.insert(i, i);
390 }
391
392 let (results, limit_reached) = store.filter_limited(|v| *v > 100, 10, 100);
393
394 assert!(results.is_empty());
395 assert!(!limit_reached);
396 }
397
398 #[test]
399 fn test_filter_limited_partial_match_within_limits() {
400 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
401
402 for i in 0..10 {
403 store.insert(i, i);
404 }
405
406 let (results, limit_reached) = store.filter_limited(|v| v % 2 == 0, 100, 100);
407
408 assert_eq!(results.len(), 5);
409 assert!(!limit_reached);
410 }
411
412 #[test]
413 fn test_clone_shares_data() {
414 let store1: InMemoryStore<String, i32> = InMemoryStore::new();
415 store1.insert("key".to_string(), 42);
416
417 let store2 = store1.clone();
418 assert_eq!(store2.get(&"key".to_string()), Some(42));
419
420 store2.insert("another".to_string(), 100);
421 assert_eq!(store1.get(&"another".to_string()), Some(100));
422 }
423
424 #[test]
425 fn test_contains_key() {
426 let store: InMemoryStore<String, i32> = InMemoryStore::new();
427
428 assert!(!store.contains_key(&"key".to_string()));
429 store.insert("key".to_string(), 42);
430 assert!(store.contains_key(&"key".to_string()));
431 }
432
433 #[test]
437 fn test_concurrent_access() {
438 use std::thread;
439
440 let store: InMemoryStore<i32, String> = InMemoryStore::new();
441 let store_clone = store.clone();
442
443 let write_handle = thread::spawn(move || {
445 for i in 0..100 {
446 store_clone.insert(i, format!("thread1-{}", i));
447 }
448 });
449
450 for i in 100..200 {
452 store.insert(i, format!("thread2-{}", i));
453 }
454
455 write_handle.join().unwrap();
456
457 assert_eq!(store.count(), 200);
459 assert_eq!(store.get(&50), Some("thread1-50".to_string()));
460 assert_eq!(store.get(&150), Some("thread2-150".to_string()));
461
462 let read_store = store.clone();
464 let read_handle = thread::spawn(move || {
465 for i in 0..200 {
466 read_store.get(&i); }
468 });
469
470 for i in 0..200 {
472 store.get(&i);
473 }
474
475 read_handle.join().unwrap();
476 }
477
478 #[test]
479 fn test_iter() {
480 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
481
482 store.insert(1, 10);
483 store.insert(2, 20);
484 store.insert(3, 30);
485
486 let mut count = 0;
487 for entry in store.iter() {
488 assert!(entry.value() == &10 || entry.value() == &20 || entry.value() == &30);
489 count += 1;
490 }
491
492 assert_eq!(count, 3);
493 }
494
495 #[test]
496 fn test_max_prealloc_size_limits_allocation() {
497 let store: InMemoryStore<i32, i32> = InMemoryStore::new();
499 store.insert(1, 1);
500
501 let (results, _) = store.filter_limited(|_| true, 1_000_000, 1_000_000);
503
504 assert_eq!(results.len(), 1);
506 }
507}