Skip to main content

pardarsh_core/storage/
memory.rs

1use std::collections::{BTreeMap, HashMap};
2
3use super::{IncomingEdge, RecordQuery, Store, check_sequence};
4use crate::error::{Error, Result};
5use crate::event::{Event, EventKind, scrub_event};
6use crate::id::{Origin, RecordId};
7use crate::model::Record;
8
9/// In-memory reference store. Useful for tests, examples and tooling.
10#[derive(Debug)]
11pub struct MemoryStore {
12    origin: Origin,
13    records: BTreeMap<RecordId, (Record, Vec<Event>)>,
14    idempotency: HashMap<String, (RecordId, usize)>,
15}
16
17impl MemoryStore {
18    pub fn new(origin: Origin) -> Self {
19        Self { origin, records: BTreeMap::new(), idempotency: HashMap::new() }
20    }
21}
22
23impl Store for MemoryStore {
24    fn origin(&self) -> &Origin {
25        &self.origin
26    }
27
28    fn get(&self, id: &RecordId) -> Result<Option<Record>> {
29        Ok(self.records.get(id).map(|(r, _)| r.clone()))
30    }
31
32    fn history(&self, id: &RecordId) -> Result<Vec<Event>> {
33        Ok(self.records.get(id).map(|(_, e)| e.clone()).unwrap_or_default())
34    }
35
36    fn event_by_idempotency_key(&self, key: &str) -> Result<Option<Event>> {
37        Ok(self.idempotency.get(key).and_then(|(id, i)| self.records.get(id).map(|(_, events)| events[*i].clone())))
38    }
39
40    fn commit(&mut self, event: Event, snapshot: Record) -> Result<()> {
41        let id = event.record_id.clone();
42        let current = self.records.get(&id).map_or(0, |(r, _)| r.version);
43        check_sequence(&id, current, &event, &snapshot)?;
44        if let Some(key) = &event.idempotency_key {
45            if self.idempotency.contains_key(key) {
46                return Err(Error::IdempotencyConflict(key.clone()));
47            }
48        }
49        let key = event.idempotency_key.clone();
50        let entry = self.records.entry(id.clone()).or_insert_with(|| (snapshot.clone(), Vec::new()));
51        if let EventKind::Redacted { fields } = &event.kind {
52            for e in entry.1.iter_mut() {
53                scrub_event(e, fields);
54            }
55        }
56        entry.0 = snapshot;
57        entry.1.push(event);
58        if let Some(key) = key {
59            self.idempotency.insert(key, (id, entry.1.len() - 1));
60        }
61        Ok(())
62    }
63
64    fn list(&self, query: &RecordQuery) -> Result<Vec<Record>> {
65        Ok(query.page(self.records.values().rev().map(|(r, _)| r.clone())))
66    }
67
68    fn incoming(&self, target: &RecordId) -> Result<Vec<IncomingEdge>> {
69        let mut out = Vec::new();
70        for (record, _) in self.records.values() {
71            for rel in record.relations.iter().filter(|r| &r.target == target) {
72                out.push(IncomingEdge {
73                    source: record.id.clone(),
74                    source_kind: record.kind.clone(),
75                    relation: rel.clone(),
76                });
77            }
78        }
79        Ok(out)
80    }
81}