dkms_wasm/database/in_memory/
escrow_database.rs1use std::{
2 collections::HashMap,
3 sync::{Arc, RwLock},
4};
5
6use said::SelfAddressingIdentifier;
7
8use keri_core::{
9 database::EscrowDatabase,
10 database::LogDatabase,
11 event::KeyEvent,
12 event_message::{msg::KeriEvent, signed_event_message::SignedEventMessage},
13 prefix::IdentifierPrefix,
14};
15
16use super::InMemoryDbError;
17
18pub struct InMemoryEscrowDatabase {
19 sn_db: Arc<dyn keri_core::database::SequencedEventDatabase<
20 DatabaseType = (),
21 Error = InMemoryDbError,
22 DigestIter = Box<dyn Iterator<Item = SelfAddressingIdentifier>>,
23 >>,
24 log: Arc<super::logging::InMemoryLogDatabase>,
25 events: RwLock<HashMap<SelfAddressingIdentifier, SignedEventMessage>>,
26}
27
28impl EscrowDatabase for InMemoryEscrowDatabase {
29 type EscrowDatabaseType = ();
30 type LogDatabaseType = super::logging::InMemoryLogDatabase;
31 type Error = InMemoryDbError;
32 type EventIter = Box<dyn Iterator<Item = SignedEventMessage>>;
33
34 fn new(
35 sn_db: Arc<dyn keri_core::database::SequencedEventDatabase<
36 DatabaseType = (),
37 Error = InMemoryDbError,
38 DigestIter = Box<dyn Iterator<Item = SelfAddressingIdentifier>>,
39 >>,
40 log: Arc<super::logging::InMemoryLogDatabase>
41 ) -> Self {
42 Self {
43 sn_db,
44 log,
45 events: RwLock::new(HashMap::new()),
46 }
47 }
48
49 fn save_digest(
50 &self,
51 id: &IdentifierPrefix,
52 sn: u64,
53 event_digest: &SelfAddressingIdentifier,
54 ) -> Result<(), Self::Error> {
55 self.sn_db.insert(id, sn, event_digest)
56 }
57
58 fn insert(&self, event: &SignedEventMessage) -> Result<(), Self::Error> {
59 let said = event.event_message.digest().unwrap();
60 let id = event.event_message.data.get_prefix();
61 let sn = event.event_message.data.sn;
62
63 self.log.log_event_with_new_transaction(event)?;
64 self.events.write().unwrap().insert(said.clone(), event.clone());
65 self.sn_db.insert(&id, sn, &said)
66 }
67
68 fn insert_key_value(
69 &self,
70 id: &IdentifierPrefix,
71 sn: u64,
72 event: &SignedEventMessage,
73 ) -> Result<(), Self::Error> {
74 let said = event.event_message.digest().unwrap();
75
76 self.log.log_event_with_new_transaction(event)?;
77 self.events.write().unwrap().insert(said.clone(), event.clone());
78 self.sn_db.insert(id, sn, &said)
79 }
80
81 fn get(
82 &self,
83 identifier: &IdentifierPrefix,
84 sn: u64,
85 ) -> Result<Self::EventIter, Self::Error> {
86 let digests = self.sn_db.get(identifier, sn)?;
87 let events = self.events.read().unwrap();
88
89 let events_cloned = events.clone();
90 let events_iter = digests.filter_map(move |digest| {
91 events_cloned.get(&digest).cloned()
92 });
93
94 Ok(Box::new(events_iter))
95 }
96
97 fn get_from_sn(
98 &self,
99 identifier: &IdentifierPrefix,
100 sn: u64,
101 ) -> Result<Self::EventIter, Self::Error> {
102 let digests = self.sn_db.get_greater_than(identifier, sn)?;
103 let events = self.events.read().unwrap();
104
105 let events_cloned = events.clone();
106 let events_iter = digests.filter_map(move |digest| {
107 events_cloned.get(&digest).cloned()
108 });
109
110 Ok(Box::new(events_iter))
111 }
112
113 fn remove(&self, event: &KeriEvent<KeyEvent>) {
114 let said = event.digest().unwrap();
115 let id = event.data.get_prefix();
116 let sn = event.data.sn;
117
118 self.sn_db.remove(&id, sn, &said).ok();
119 self.events.write().unwrap().remove(&said);
120 }
121
122 fn contains(
123 &self,
124 id: &IdentifierPrefix,
125 sn: u64,
126 digest: &SelfAddressingIdentifier,
127 ) -> Result<bool, Self::Error> {
128 let mut digests = self.sn_db.get(id, sn)?;
129 let result: bool = digests.any(|d| &d == digest);
130 Ok(result)
131 }
132}