Skip to main content

dkms_wasm/database/in_memory/
mod.rs

1use std::{
2    collections::HashMap,
3    sync::{Arc, RwLock},
4};
5
6use keri_sdk::TelEventDatabase;
7use said::SelfAddressingIdentifier;
8use teliox::event::{verifiable_event::VerifiableEvent, Event};
9
10use keri_core::{
11    database::SequencedEventDatabase,
12    event::KeyEvent,
13    event_message::{
14        msg::KeriEvent,
15        signed_event_message::{
16            SignedEventMessage, SignedNontransferableReceipt,
17            SignedTransferableReceipt,
18        },
19    },
20    prefix::IdentifierPrefix,
21    state::IdentifierState,
22};
23
24use keri_core::database::{
25    timestamped, EscrowCreator, EscrowDatabase, EventDatabase, QueryParameters,
26};
27
28use cesrox::primitives::CesrPrimitive;
29
30pub mod logging;
31use keri_core::database::LogDatabase;
32use logging::InMemoryLogDatabase;
33
34pub mod escrow_database;
35pub mod sn_database;
36
37use escrow_database::InMemoryEscrowDatabase;
38use sn_database::InMemorySnDatabase;
39
40impl EscrowCreator for InMemoryDatabase {
41    type EscrowDatabaseType = InMemoryEscrowDatabase;
42
43    fn create_escrow_db(
44        &self,
45        _table_name: &'static str,
46    ) -> Self::EscrowDatabaseType {
47        InMemoryEscrowDatabase::new(
48            Arc::new(
49                InMemorySnDatabase::new(Arc::new(()), _table_name).unwrap(),
50            ),
51            self.log_db.clone(),
52        )
53    }
54}
55
56#[derive(Debug, thiserror::Error)]
57pub enum InMemoryDbError {
58    #[error("Not found: {0}")]
59    NotFound(String),
60    #[error("Already saved: {0}")]
61    AlreadySaved(SelfAddressingIdentifier),
62    #[error("Event not found")]
63    MissingDigest,
64    #[error("Lock error")]
65    LockError,
66}
67
68/// In-memory implementation of the event database
69pub struct InMemoryDatabase {
70    kels: RwLock<HashMap<(String, u64), SelfAddressingIdentifier>>,
71    key_states: RwLock<HashMap<String, IdentifierState>>,
72    events: RwLock<
73        HashMap<
74            SelfAddressingIdentifier,
75            timestamped::TimestampedSignedEventMessage,
76        >,
77    >,
78    trans_receipts: RwLock<
79        HashMap<
80            SelfAddressingIdentifier,
81            Vec<keri_core::event_message::signature::Transferable>,
82        >,
83    >,
84    nontrans_receipts: RwLock<
85        HashMap<
86            SelfAddressingIdentifier,
87            Vec<keri_core::event_message::signature::Nontransferable>,
88        >,
89    >,
90    log_db: Arc<InMemoryLogDatabase>,
91    tel_events: RwLock<HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>,
92    management_events: RwLock<HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>,
93}
94
95impl InMemoryDatabase {
96    pub fn new() -> Self {
97        Self {
98            kels: RwLock::new(HashMap::new()),
99            key_states: RwLock::new(HashMap::new()),
100            events: RwLock::new(HashMap::new()),
101            trans_receipts: RwLock::new(HashMap::new()),
102            nontrans_receipts: RwLock::new(HashMap::new()),
103            log_db: Arc::new(InMemoryLogDatabase::new(Arc::new(())).unwrap()),
104            tel_events: RwLock::new(HashMap::new()),
105            management_events: RwLock::new(HashMap::new()),
106        }
107    }
108}
109
110impl Default for InMemoryDatabase {
111    fn default() -> Self {
112        Self::new()
113    }
114}
115
116impl EventDatabase for InMemoryDatabase {
117    type Error = InMemoryDbError;
118    type LogDatabaseType = InMemoryLogDatabase;
119
120    fn get_log_db(&self) -> Arc<Self::LogDatabaseType> {
121        self.log_db.clone()
122    }
123
124    fn add_kel_finalized_event(
125        &self,
126        signed_event: SignedEventMessage,
127        id: &IdentifierPrefix,
128    ) -> Result<(), Self::Error> {
129        let event = &signed_event.event_message;
130        let digest = event.digest().unwrap();
131        let id_str = id.to_str();
132        let sn = event.data.sn;
133
134        // Update key state
135        let mut key_states = self.key_states.write().unwrap();
136        let mut key_state =
137            key_states.get(&id_str).cloned().unwrap_or_default();
138        key_state = key_state
139            .apply(event)
140            .map_err(|_| InMemoryDbError::AlreadySaved(digest.clone()))?;
141        key_states.insert(id_str.clone(), key_state);
142
143        // Save to KEL
144        self.kels
145            .write()
146            .unwrap()
147            .insert((id_str, sn), digest.clone());
148
149        self.log_db.log_event_with_new_transaction(&signed_event)?;
150
151        // Save event
152        self.events.write().unwrap().insert(
153            digest,
154            timestamped::TimestampedSignedEventMessage::new(signed_event),
155        );
156
157        Ok(())
158    }
159
160    fn add_receipt_t(
161        &self,
162        receipt: SignedTransferableReceipt,
163        _id: &IdentifierPrefix,
164    ) -> Result<(), Self::Error> {
165        let digest = receipt.body.receipted_event_digest;
166        let transferable =
167            keri_core::event_message::signature::Transferable::Seal(
168                receipt.validator_seal,
169                receipt.signatures,
170            );
171
172        let mut receipts = self.trans_receipts.write().unwrap();
173        receipts.entry(digest).or_default().push(transferable);
174
175        Ok(())
176    }
177
178    fn add_receipt_nt(
179        &self,
180        receipt: SignedNontransferableReceipt,
181        _id: &IdentifierPrefix,
182    ) -> Result<(), Self::Error> {
183        let digest = receipt.body.receipted_event_digest;
184
185        let mut receipts = self.nontrans_receipts.write().unwrap();
186        receipts
187            .entry(digest)
188            .or_default()
189            .extend(receipt.signatures);
190
191        Ok(())
192    }
193
194    fn get_key_state(&self, id: &IdentifierPrefix) -> Option<IdentifierState> {
195        self.key_states.read().unwrap().get(&id.to_str()).cloned()
196    }
197
198    fn get_kel_finalized_events(
199        &self,
200        params: QueryParameters,
201    ) -> Option<
202        impl DoubleEndedIterator<Item = timestamped::TimestampedSignedEventMessage>,
203    > {
204        match params {
205            QueryParameters::BySn { id, sn } => {
206                let key = (id.to_str(), sn);
207                let kels = self.kels.read().unwrap();
208                let events = self.events.read().unwrap();
209
210                kels.get(&key)
211                    .and_then(|digest| events.get(digest))
212                    .cloned()
213                    .map(|event| vec![event].into_iter())
214            }
215            QueryParameters::Range { id, start, limit } => {
216                let id_str = id.to_str();
217                let kels = self.kels.read().unwrap();
218                let events = self.events.read().unwrap();
219
220                let mut result = Vec::new();
221                for sn in start..(start + limit) {
222                    if let Some(digest) = kels.get(&(id_str.clone(), sn)) {
223                        if let Some(event) = events.get(digest) {
224                            result.push(event.clone());
225                        }
226                    }
227                }
228
229                if result.is_empty() {
230                    None
231                } else {
232                    Some(result.into_iter())
233                }
234            }
235            QueryParameters::All { id } => {
236                let id_str = id.to_str();
237                let kels = self.kels.read().unwrap();
238                let events = self.events.read().unwrap();
239
240                let mut result = Vec::new();
241                for ((prefix, _), digest) in kels.iter() {
242                    if prefix == &id_str {
243                        if let Some(event) = events.get(digest) {
244                            result.push(event.clone());
245                        }
246                    }
247                }
248
249                if result.is_empty() {
250                    None
251                } else {
252                    Some(result.into_iter())
253                }
254            }
255        }
256    }
257
258    fn get_receipts_t(
259        &self,
260        params: QueryParameters,
261    ) -> Option<
262        impl DoubleEndedIterator<
263            Item = keri_core::event_message::signature::Transferable,
264        >,
265    > {
266        match params {
267            QueryParameters::BySn { id, sn } => {
268                let key = (id.to_str(), sn);
269                let kels = self.kels.read().unwrap();
270                let receipts = self.trans_receipts.read().unwrap();
271
272                kels.get(&key)
273                    .and_then(|digest| receipts.get(digest))
274                    .cloned()
275                    .map(|r| r.into_iter())
276            }
277            _ => None,
278        }
279    }
280
281    fn get_receipts_nt(
282        &self,
283        params: QueryParameters,
284    ) -> Option<impl DoubleEndedIterator<Item = SignedNontransferableReceipt>>
285    {
286        match params {
287            QueryParameters::BySn { id: _, sn: _ } => Some(vec![].into_iter()),
288            _ => None,
289        }
290    }
291
292    fn accept_to_kel(
293        &self,
294        event: &KeriEvent<KeyEvent>,
295    ) -> Result<(), Self::Error> {
296        let digest = event.digest().unwrap();
297        let id_str = event.data.get_prefix().to_str();
298        let sn = event.data.sn;
299
300        // Update key state
301        let mut key_states = self.key_states.write().unwrap();
302        let mut key_state =
303            key_states.get(&id_str).cloned().unwrap_or_default();
304        key_state = key_state
305            .apply(event)
306            .map_err(|_| InMemoryDbError::AlreadySaved(digest.clone()))?;
307        key_states.insert(id_str.clone(), key_state);
308
309        // Save to KEL
310        self.kels.write().unwrap().insert((id_str, sn), digest);
311
312        Ok(())
313    }
314
315    fn save_reply(
316        &self,
317        _reply: keri_core::query::reply_event::SignedReply,
318    ) -> Result<(), Self::Error> {
319        Ok(())
320    }
321    fn get_reply(
322        &self,
323        _id: &IdentifierPrefix,
324        _from_who: &IdentifierPrefix,
325    ) -> Option<keri_core::query::reply_event::SignedReply> {
326        None
327    }
328}
329
330impl TelEventDatabase for InMemoryDatabase {
331    fn new(
332        _path: impl AsRef<std::path::Path>,
333    ) -> Result<Self, teliox::error::Error>
334    where
335        Self: Sized,
336    {
337        Ok(Self::new())
338    }
339
340    fn add_new_event(
341        &self,
342        event: VerifiableEvent,
343        id: &IdentifierPrefix,
344    ) -> Result<(), teliox::error::Error> {
345        match event.event {
346            Event::Vc(_) => {
347                let mut events_map = self
348                    .tel_events
349                    .write()
350                    .map_err(|_| teliox::error::Error::RwLockingError)?;
351                events_map
352                    .entry(id.clone())
353                    .or_insert_with(Vec::new)
354                    .push(event);
355            }
356            Event::Management(_) => {
357                let mut events_map = self
358                    .management_events
359                    .write()
360                    .map_err(|_| teliox::error::Error::RwLockingError)?;
361                events_map
362                    .entry(id.clone())
363                    .or_insert_with(Vec::new)
364                    .push(event);
365            }
366        }
367        Ok(())
368    }
369
370    fn get_events(
371        &self,
372        id: &IdentifierPrefix,
373    ) -> Option<impl DoubleEndedIterator<Item = VerifiableEvent>> {
374        let events_map = self.tel_events.read().ok()?;
375        events_map.get(id).map(|events| events.clone().into_iter())
376    }
377
378    fn get_management_events(
379        &self,
380        id: &IdentifierPrefix,
381    ) -> Option<impl DoubleEndedIterator<Item = VerifiableEvent>> {
382        let events_map = self.management_events.read().ok()?;
383        events_map.get(id).map(|events| events.clone().into_iter())
384    }
385}