dkms_wasm/database/in_memory/
mod.rs1use 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
68pub 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 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 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 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 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 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}