1pub mod escrow_database;
2pub mod sn_database;
3pub mod logging;
4
5use std::sync::Arc;
6use std::rc::Rc;
7use std::cell::{RefCell, Cell};
8use keri_core::event::receipt::Receipt;
9use keri_core::oobi::{LocationScheme, Scheme};
10use keri_core::prefix::SeedPrefix;
11use keri_sdk::TelEventDatabase;
12use said::sad::SerializationFormats;
13use teliox::event::verifiable_event::VerifiableEvent;
14use teliox::event::Event;
15use wasm_bindgen::prelude::*;
16use wasm_bindgen_futures::spawn_local;
17
18use keri_core::{
19 event::KeyEvent,
20 event_message::{
21 msg::KeriEvent,
22 signature::Transferable,
23 signed_event_message::{
24 SignedEventMessage, SignedNontransferableReceipt, SignedTransferableReceipt,
25 },
26 },
27 prefix::IdentifierPrefix,
28 state::IdentifierState,
29};
30
31use keri_core::database::{timestamped, EscrowCreator, EscrowDatabase, EventDatabase, QueryParameters, LogDatabase, SequencedEventDatabase};
32use escrow_database::IndexedDbEscrowDatabase;
33use logging::IndexedDbLogDatabase;
34use sn_database::IndexedDbSnDatabase;
35use said::SelfAddressingIdentifier;
36
37#[derive(Debug, thiserror::Error)]
38pub enum IndexedDbError {
39 #[error("Failed to save to database")]
40 DatabaseSaveFailed(String),
41 #[error("Not found: {0}")]
42 NotFound(String),
43 #[error("Event not found")]
44 MissingDigest,
45 #[error("Lock error")]
46 LockError,
47 #[error("Invalid signature")]
48 InvalidSignature,
49 #[error("Failed to encode")]
50 EncodingFailed(String),
51 #[error("Failed to decode")]
52 DecodingFailed(String),
53 #[error("Key format error")]
54 KeyFormatError,
55}
56
57#[derive(Debug, Clone)]
58pub struct IdentifierRecord {
59 pub said: IdentifierPrefix,
60 pub seed: SeedPrefix,
61 pub watcher_oobi: Option<LocationScheme>,
62}
63
64pub struct IndexedDbDatabase {
65 log_db: Arc<IndexedDbLogDatabase>,
66 key_states: Rc<RefCell<std::collections::HashMap<String, IdentifierState>>>,
67 kels: Rc<RefCell<std::collections::HashMap<(String, u64), SelfAddressingIdentifier>>>,
68 pending_operations: Rc<RefCell<Vec<PendingDbOperation>>>,
69 flush_in_progress: Rc<Cell<bool>>,
70 tel_events: Rc<RefCell<std::collections::HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>>,
71 management_events: Rc<RefCell<std::collections::HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>>,
72 identifiers: Rc<RefCell<std::collections::HashMap<String, IdentifierRecord>>>,
73}
74
75enum PendingDbOperation {
76 SaveKeyState { id: String, state: IdentifierState },
77 SaveKel { id: String, sn: u64, digest: SelfAddressingIdentifier },
78 SaveTelEvent { id: IdentifierPrefix, event: VerifiableEvent },
79 SaveManagementEvent { id: IdentifierPrefix, event: VerifiableEvent },
80 SaveIdentifier { alias: String, said: IdentifierPrefix, seed: SeedPrefix },
81 RemoveIdentifier { alias: String },
82 AddWatcher { alias: String, watcher_oobi: LocationScheme },
83}
84
85unsafe impl Send for IndexedDbDatabase {}
87unsafe impl Sync for IndexedDbDatabase {}
88
89impl Default for IndexedDbDatabase {
90 fn default() -> Self {
91 Self::new()
92 }
93}
94
95impl IndexedDbDatabase {
96 pub fn get_identifiers(&self) -> Vec<(String, IdentifierRecord)> {
97 self.identifiers.borrow().iter().map(|(k, v)| (k.clone(), v.clone())).collect()
98 }
99
100 pub fn get_identifier(&self, alias: &str) -> Option<IdentifierRecord> {
101 self.identifiers.borrow().get(alias).cloned()
102 }
103
104 pub fn new() -> Self {
105 let log_db = Arc::new(IndexedDbLogDatabase::new(Arc::new(())).unwrap());
106 let key_states = Rc::new(RefCell::new(std::collections::HashMap::new()));
107 let kels = Rc::new(RefCell::new(std::collections::HashMap::new()));
108 let pending_operations = Rc::new(RefCell::new(Vec::new()));
109 let flush_in_progress = Rc::new(Cell::new(false));
110 let tel_events = Rc::new(RefCell::new(std::collections::HashMap::new()));
111 let management_events = Rc::new(RefCell::new(std::collections::HashMap::new()));
112 let identifiers = Rc::new(RefCell::new(std::collections::HashMap::new()));
113
114 let mut db = Self {
115 log_db,
116 key_states,
117 kels,
118 pending_operations,
119 flush_in_progress,
120 tel_events,
121 management_events,
122 identifiers,
123 };
124
125 db.init_db("keri_indexed_db");
126
127 db
128 }
129
130 fn init_db(&mut self, db_name: &str) {
131 let window = web_sys::window().expect("should have a window");
132
133 if let Ok(Some(factory)) = window.indexed_db() {
135 let pending_ops = self.pending_operations.clone();
136 let flush_flag = self.flush_in_progress.clone();
137 let key_states = self.key_states.clone();
138 let kels = self.kels.clone();
139 let tel_events = self.tel_events.clone();
140 let management_events = self.management_events.clone();
141 let identifiers = self.identifiers.clone();
142
143 if let Ok(request) = factory.open(db_name) {
145 let upgrade_needed_cb = Closure::wrap(Box::new(move |event: web_sys::IdbVersionChangeEvent| {
147 if let Some(db) = event.target()
148 .and_then(|t| t.dyn_into::<web_sys::IdbOpenDbRequest>().ok())
149 .and_then(|r| r.result().ok())
150 .and_then(|r| r.dyn_into::<web_sys::IdbDatabase>().ok())
151 {
152 let _ = db.create_object_store("kels");
154 let _ = db.create_object_store("key_states");
155 let _ = db.create_object_store("tel_events");
156 let _ = db.create_object_store("management_events");
157 let _ = db.create_object_store("identifiers");
158 }
159 }) as Box<dyn FnMut(_)>);
160
161 request.set_onupgradeneeded(Some(upgrade_needed_cb.as_ref().unchecked_ref()));
162 upgrade_needed_cb.forget();
163
164 let success_cb = {
166 let this_db = Rc::new(RefCell::new(None::<web_sys::IdbDatabase>));
167 let this_db_clone = this_db.clone();
168
169 Closure::wrap(Box::new(move |event: web_sys::Event| {
170 if let Some(db) = event.target()
171 .and_then(|t| t.dyn_into::<web_sys::IdbOpenDbRequest>().ok())
172 .and_then(|r| r.result().ok())
173 .and_then(|r| r.dyn_into::<web_sys::IdbDatabase>().ok())
174 {
175 *this_db_clone.borrow_mut() = Some(db.clone());
176
177 load_key_states(&db, key_states.clone());
179 load_kels(&db, kels.clone());
180 load_tel_events(&db, tel_events.clone());
181 load_management_events(&db, management_events.clone());
182 load_identifiers(&db, identifiers.clone());
183
184 setup_background_flush(db, pending_ops.clone(), flush_flag.clone());
186 }
187 }) as Box<dyn FnMut(_)>)
188 };
189
190 request.set_onsuccess(Some(success_cb.as_ref().unchecked_ref()));
191 success_cb.forget();
192
193 let error_cb = Closure::wrap(Box::new(|event: web_sys::Event| {
195 log::error!("Failed to open IndexedDB: {:?}", event);
196 }) as Box<dyn FnMut(_)>);
197
198 request.set_onerror(Some(error_cb.as_ref().unchecked_ref()));
199 error_cb.forget();
200 }
201 }
202 }
203
204 pub fn add_identifier(&self, alias: &str, said: &IdentifierPrefix, seed: &SeedPrefix) -> Result<(), IndexedDbError> {
205 let record = IdentifierRecord {
206 said: said.clone(),
207 seed: seed.clone(),
208 watcher_oobi: None,
209 };
210
211 self.identifiers.borrow_mut().insert(alias.to_string(), record);
212 self.pending_operations.borrow_mut().push(PendingDbOperation::SaveIdentifier {
213 alias: alias.to_string(),
214 said: said.clone(),
215 seed: seed.clone()
216 });
217
218 Ok(())
219 }
220
221 pub fn update_identifier_watcher(&self, alias: &str, watcher_oobi: LocationScheme) -> Result<(), IndexedDbError> {
222 if let Some(record) = self.identifiers.borrow_mut().get_mut(alias) {
223 record.watcher_oobi = Some(watcher_oobi.clone());
224 self.pending_operations.borrow_mut().push(PendingDbOperation::AddWatcher {
225 alias: alias.to_string(),
226 watcher_oobi: watcher_oobi.clone()
227 });
228 Ok(())
229 } else {
230 Err(IndexedDbError::NotFound(format!("Identifier with alias {} not found", alias)))
231 }
232 }
233
234 pub fn update_identifier_alias(&self, old_alias: &str, new_alias: &str) -> Result<(), IndexedDbError> {
235 let mut identifiers = self.identifiers.borrow_mut();
236 if let Some(record) = identifiers.remove(old_alias) {
237 identifiers.insert(new_alias.to_string(), record.clone());
238 self.pending_operations.borrow_mut().push(PendingDbOperation::SaveIdentifier {
239 alias: new_alias.to_string(),
240 said: record.said.clone(),
241 seed: record.seed.clone()
242 });
243 self.pending_operations.borrow_mut().push(PendingDbOperation::RemoveIdentifier {
244 alias: old_alias.to_string(),
245 });
246 Ok(())
247 } else {
248 Err(IndexedDbError::NotFound(format!("Identifier with alias {} not found", old_alias)))
249 }
250 }
251}
252
253fn load_key_states(db: &web_sys::IdbDatabase, key_states: Rc<RefCell<std::collections::HashMap<String, IdentifierState>>>) {
255 if let Ok(transaction) = db.transaction_with_str_and_mode(
256 "key_states",
257 web_sys::IdbTransactionMode::Readwrite,
258 ) {
259 if let Ok(store) = transaction.object_store("key_states") {
260 if let Ok(request) = store.get_all() {
261 let callback = Closure::wrap(Box::new(move |event: web_sys::Event| {
262 if let Some(result) = event.target()
263 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
264 .and_then(|r| r.result().ok())
265 {
266 if let Ok(array) = result.dyn_into::<js_sys::Array>() {
267 let mut states = key_states.borrow_mut();
268 for i in 0..array.length() {
269 if let Ok(item) = array.get(i).dyn_into::<js_sys::Object>() {
270 if let (Ok(id), Ok(state_bytes)) = (
271 js_sys::Reflect::get(&item, &"id".into()),
272 js_sys::Reflect::get(&item, &"state".into())
273 ) {
274 if let (Some(id_str), Some(state_str)) = (
275 id.as_string(),
276 state_bytes.as_string()
277 ) {
278 if let Ok(state) = serde_json::from_str::<IdentifierState>(&state_str) {
279 states.insert(id_str, state);
280 }
281 }
282 }
283 }
284 }
285 }
286 }
287 }) as Box<dyn FnMut(_)>);
288
289 request.set_onsuccess(Some(callback.as_ref().unchecked_ref()));
290 callback.forget();
291 }
292 }
293 }
294}
295
296fn load_kels(db: &web_sys::IdbDatabase, kels: Rc<RefCell<std::collections::HashMap<(String, u64), SelfAddressingIdentifier>>>) {
298 if let Ok(transaction) = db.transaction_with_str_and_mode(
299 "kels",
300 web_sys::IdbTransactionMode::Readwrite,
301 ) {
302 if let Ok(store) = transaction.object_store("kels") {
303 if let Ok(request) = store.get_all() {
304 let callback = Closure::wrap(Box::new(move |event: web_sys::Event| {
305 if let Some(result) = event.target()
306 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
307 .and_then(|r| r.result().ok())
308 {
309 if let Ok(array) = result.dyn_into::<js_sys::Array>() {
310 let mut kel_map = kels.borrow_mut();
311 for i in 0..array.length() {
312 if let Ok(item) = array.get(i).dyn_into::<js_sys::Object>() {
313 if let (Ok(id), Ok(sn), Ok(digest)) = (
314 js_sys::Reflect::get(&item, &"id".into()),
315 js_sys::Reflect::get(&item, &"sn".into()),
316 js_sys::Reflect::get(&item, &"digest".into())
317 ) {
318 if let (Some(id_str), Some(sn_num), Some(digest_str)) = (
319 id.as_string(),
320 sn.as_f64().map(|n| n as u64),
321 digest.as_string()
322 ) {
323 if let Ok(said) = digest_str.parse::<SelfAddressingIdentifier>() {
324 kel_map.insert((id_str, sn_num), said);
325 }
326 }
327 }
328 }
329 }
330 }
331 }
332 }) as Box<dyn FnMut(_)>);
333
334 request.set_onsuccess(Some(callback.as_ref().unchecked_ref()));
335 callback.forget();
336 }
337 }
338 }
339}
340
341fn load_tel_events(db: &web_sys::IdbDatabase, tel_events: Rc<RefCell<std::collections::HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>>) {
343 if let Ok(transaction) = db.transaction_with_str_and_mode(
344 "tel_events",
345 web_sys::IdbTransactionMode::Readwrite,
346 ) {
347 if let Ok(store) = transaction.object_store("tel_events") {
348 if let Ok(request) = store.get_all() {
349 let callback = Closure::wrap(Box::new(move |event: web_sys::Event| {
350 if let Some(result) = event.target()
351 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
352 .and_then(|r| r.result().ok())
353 {
354 if let Ok(array) = result.dyn_into::<js_sys::Array>() {
355 let mut events_map = tel_events.borrow_mut();
356 for i in 0..array.length() {
357 if let Ok(item) = array.get(i).dyn_into::<js_sys::Object>() {
358 if let (Ok(id), Ok(event_json)) = (
359 js_sys::Reflect::get(&item, &"id".into()),
360 js_sys::Reflect::get(&item, &"event".into())
361 ) {
362 if let (Some(id_str), Some(event_str)) = (
363 id.as_string(),
364 event_json.as_string()
365 ) {
366 if let (Ok(prefix), Ok(event)) = (
367 id_str.parse::<IdentifierPrefix>(),
368 serde_json::from_str::<VerifiableEvent>(&event_str)
369 ) {
370 events_map.entry(prefix)
371 .or_default()
372 .push(event);
373 }
374 }
375 }
376 }
377 }
378 }
379 }
380 }) as Box<dyn FnMut(_)>);
381
382 request.set_onsuccess(Some(callback.as_ref().unchecked_ref()));
383 callback.forget();
384 }
385 }
386 }
387}
388
389fn load_management_events(db: &web_sys::IdbDatabase, management_events: Rc<RefCell<std::collections::HashMap<IdentifierPrefix, Vec<VerifiableEvent>>>>) {
391 if let Ok(transaction) = db.transaction_with_str_and_mode(
392 "management_events",
393 web_sys::IdbTransactionMode::Readwrite,
394 ) {
395 if let Ok(store) = transaction.object_store("management_events") {
396 if let Ok(request) = store.get_all() {
397 let callback = Closure::wrap(Box::new(move |event: web_sys::Event| {
398 if let Some(result) = event.target()
399 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
400 .and_then(|r| r.result().ok())
401 {
402 if let Ok(array) = result.dyn_into::<js_sys::Array>() {
403 let mut events_map = management_events.borrow_mut();
404 for i in 0..array.length() {
405 if let Ok(item) = array.get(i).dyn_into::<js_sys::Object>() {
406 if let (Ok(id), Ok(event_json)) = (
407 js_sys::Reflect::get(&item, &"id".into()),
408 js_sys::Reflect::get(&item, &"event".into())
409 ) {
410 if let (Some(id_str), Some(event_str)) = (
411 id.as_string(),
412 event_json.as_string()
413 ) {
414 if let (Ok(prefix), Ok(event)) = (
415 id_str.parse::<IdentifierPrefix>(),
416 serde_json::from_str::<VerifiableEvent>(&event_str)
417 ) {
418 events_map.entry(prefix)
419 .or_default()
420 .push(event);
421 }
422 }
423 }
424 }
425 }
426 }
427 }
428 }) as Box<dyn FnMut(_)>);
429
430 request.set_onsuccess(Some(callback.as_ref().unchecked_ref()));
431 callback.forget();
432 }
433 }
434 }
435}
436
437fn load_identifiers(db: &web_sys::IdbDatabase, identifiers: Rc<RefCell<std::collections::HashMap<String, IdentifierRecord>>>) {
438 if let Ok(transaction) = db.transaction_with_str_and_mode(
439 "identifiers",
440 web_sys::IdbTransactionMode::Readwrite,
441 ) {
442 if let Ok(store) = transaction.object_store("identifiers") {
443 if let Ok(request) = store.open_cursor() {
444 let callback = Closure::wrap(Box::new(move |event: web_sys::Event| {
445 if let Some(cursor_result) = event.target()
446 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
447 .and_then(|r| r.result().ok())
448 {
449 if !cursor_result.is_undefined() {
451 if let Ok(cursor) = cursor_result.dyn_into::<web_sys::IdbCursorWithValue>() {
452 let key = cursor.key().ok()
453 .and_then(|k| k.as_string())
454 .unwrap_or_default();
455 let item = cursor.value().unwrap();
456
457 if let (Ok(said), Ok(seed), watcher_url, watcher_eid, watcher_scheme) = (
458 js_sys::Reflect::get(&item, &"said".into()),
459 js_sys::Reflect::get(&item, &"seed".into()),
460 js_sys::Reflect::get(&item, &"watcher_url".into()),
461 js_sys::Reflect::get(&item, &"watcher_eid".into()),
462 js_sys::Reflect::get(&item, &"watcher_scheme".into())
463 ) {
464 if let (Some(said_str), Some(seed_str)) = (
465 said.as_string(),
466 seed.as_string()
467 ) {
468 if let (Ok(prefix), Ok(seed)) = (
469 said_str.parse::<IdentifierPrefix>(),
470 seed_str.parse::<SeedPrefix>(),
471 ) {
472 let mut watcher_oobi = None;
473 if let (Ok(watcher_url), Ok(watcher_eid), Ok(watcher_scheme)) = (watcher_url, watcher_eid, watcher_scheme) {
474 if let (Some(url_str), Some(eid_str), Some(scheme_str)) = (
475 watcher_url.as_string(),
476 watcher_eid.as_string(),
477 watcher_scheme.as_string()
478 ) {
479 if let (Ok(url), Ok(eid), Ok(scheme)) = (
480 url::Url::parse(&url_str),
481 eid_str.parse::<IdentifierPrefix>(),
482 scheme_str.parse::<Scheme>()
483 ) {
484 watcher_oobi = Some(LocationScheme::new(eid, scheme, url));
485 }
486 }
487 }
488
489 let id = IdentifierRecord {
490 said: prefix,
491 seed,
492 watcher_oobi,
493 };
494
495 identifiers.borrow_mut().insert(key, id);
496 }
497 }
498 }
499
500 if let Err(e) = cursor.continue_() {
501 log::error!("Error continuing cursor: {:?}", e);
502 }
503 }
504 }
505 }
506 }) as Box<dyn FnMut(_)>);
507
508 request.set_onsuccess(Some(callback.as_ref().unchecked_ref()));
509 callback.forget();
510 }
511 }
512 }
513}
514
515fn setup_background_flush(
517 db: web_sys::IdbDatabase,
518 pending_ops: Rc<RefCell<Vec<PendingDbOperation>>>,
519 flush_flag: Rc<Cell<bool>>
520) {
521 let interval_callback = Closure::wrap(Box::new(move || {
523 if !flush_flag.get() {
524 flush_flag.set(true);
525
526 spawn_local({
527 let db = db.clone();
528 let pending = pending_ops.clone();
529 let flag = flush_flag.clone();
530
531 async move {
532 let ops_to_flush = {
534 let mut pending_borrow = pending.borrow_mut();
535 if pending_borrow.is_empty() {
536 Vec::new()
537 } else {
538 pending_borrow.drain(..).collect::<Vec<_>>()
539 }
540 };
541
542 if !ops_to_flush.is_empty() {
544 let key_state_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
546 if let PendingDbOperation::SaveKeyState { id, state } = op {
547 Some((id.clone(), state.clone()))
548 } else {
549 None
550 }
551 }).collect();
552
553 if !key_state_ops.is_empty() && db.transaction_with_str_and_mode("key_states", web_sys::IdbTransactionMode::Readwrite).is_ok() {
554 let transaction = db.transaction_with_str_and_mode("key_states", web_sys::IdbTransactionMode::Readwrite).unwrap();
555 if let Ok(store) = transaction.object_store("key_states") {
556 for (id, state) in key_state_ops {
557 let value = js_sys::Object::new();
558 js_sys::Reflect::set(&value, &"id".into(), &id.clone().into()).unwrap();
559
560
561 if let Ok(state_json) = serde_json::to_string(&state) {
563 js_sys::Reflect::set(&value, &"state".into(), &state_json.into()).unwrap();
564
565 if let Err(e) = store.put_with_key(&value, &id.into()) {
566 log::error!("Failed to store key state in IndexedDB: {:?}", e);
567 }
568 }
569 }
570 }
571 }
572
573 let kel_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
575 match op {
576 PendingDbOperation::SaveKel { id, sn, digest } => {
577 Some(("save", id.clone(), *sn, digest.clone()))
578 }
579 _ => None
580 }
581 }).collect();
582
583 if !kel_ops.is_empty() && db.transaction_with_str_and_mode("kels", web_sys::IdbTransactionMode::Readwrite).is_ok() {
584 let transaction = db.transaction_with_str_and_mode("kels", web_sys::IdbTransactionMode::Readwrite).unwrap();
585 if let Ok(store) = transaction.object_store("kels") {
586 for (op_type, id, sn, digest) in kel_ops {
587 let key = format!("{}:{}", id, sn);
588
589 if op_type == "save" {
590 let value = js_sys::Object::new();
591 js_sys::Reflect::set(&value, &"id".into(), &id.into()).unwrap();
592 js_sys::Reflect::set(&value, &"sn".into(), &(sn as f64).into()).unwrap();
593 js_sys::Reflect::set(&value, &"digest".into(), &digest.to_string().into()).unwrap();
594
595 if let Err(e) = store.put_with_key(&value, &key.into()) {
596 log::error!("Failed to store KEL in IndexedDB: {:?}", e);
597 }
598 } else if let Err(e) = store.delete(&key.into()) {
599 log::error!("Failed to remove KEL from IndexedDB: {:?}", e);
600 }
601 }
602 }
603 }
604
605 let tel_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
607 if let PendingDbOperation::SaveTelEvent { id, event } = op {
608 Some((id.clone(), event.clone()))
609 } else {
610 None
611 }
612 }).collect();
613
614 if !tel_ops.is_empty() && db.transaction_with_str_and_mode("tel_events", web_sys::IdbTransactionMode::Readwrite).is_ok() {
615 let transaction = db.transaction_with_str_and_mode("tel_events", web_sys::IdbTransactionMode::Readwrite).unwrap();
616 if let Ok(store) = transaction.object_store("tel_events") {
617 for (id, event) in tel_ops {
618 let key = id.to_string();
619 let value = js_sys::Object::new();
620 js_sys::Reflect::set(&value, &"id".into(), &key.clone().into()).unwrap();
621
622 if let Ok(event_json) = serde_json::to_string(&event) {
624 js_sys::Reflect::set(&value, &"event".into(), &event_json.into()).unwrap();
625
626 if let Err(e) = store.put_with_key(&value, &key.into()) {
627 log::error!("Failed to store TEL event in IndexedDB: {:?}", e);
628 }
629 }
630 }
631 }
632 }
633
634 let mgmt_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
636 if let PendingDbOperation::SaveManagementEvent { id, event } = op {
637 Some((id.clone(), event.clone()))
638 } else {
639 None
640 }
641 }).collect();
642
643 if !mgmt_ops.is_empty() && db.transaction_with_str_and_mode("management_events", web_sys::IdbTransactionMode::Readwrite).is_ok() {
644 let transaction = db.transaction_with_str_and_mode("management_events", web_sys::IdbTransactionMode::Readwrite).unwrap();
645 if let Ok(store) = transaction.object_store("management_events") {
646 for (id, event) in mgmt_ops {
647 let key = id.to_string();
648 let value = js_sys::Object::new();
649 js_sys::Reflect::set(&value, &"id".into(), &key.clone().into()).unwrap();
650
651 if let Ok(event_json) = serde_json::to_string(&event) {
653 js_sys::Reflect::set(&value, &"event".into(), &event_json.into()).unwrap();
654
655 if let Err(e) = store.put_with_key(&value, &key.into()) {
656 log::error!("Failed to store management event in IndexedDB: {:?}", e);
657 }
658 }
659 }
660 }
661 }
662
663 let id_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
665 if let PendingDbOperation::SaveIdentifier { alias, said, seed } = op {
666 Some((alias.clone(), said.clone(), seed.clone()))
667 } else {
668 None
669 }
670 }).collect();
671
672 if !id_ops.is_empty() && db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite).is_ok() {
673 let transaction = db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite).unwrap();
674 if let Ok(store) = transaction.object_store("identifiers") {
675 for (alias, said, seed) in id_ops {
676 let value = js_sys::Object::new();
677
678 let key = alias.to_string();
679 let said = said.to_string();
680 js_sys::Reflect::set(&value, &"said".into(), &said.clone().into()).unwrap();
681 let seed_str = serde_json::to_string(&seed).unwrap();
682 if let serde_json::Value::String(seed) = serde_json::from_str(&seed_str).unwrap() {
683 js_sys::Reflect::set(&value, &"seed".into(), &seed.clone().into()).unwrap();
684 }
685 if let Err(e) = store.put_with_key(&value, &key.into()) {
686 log::error!("Failed to store management event in IndexedDB: {:?}", e);
687 }
688 }
689 }
690 }
691
692 let id_alias_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
693 if let PendingDbOperation::RemoveIdentifier { alias } = op {
694 Some(alias.clone())
695 } else {
696 None
697 }
698 }).collect();
699
700 if !id_alias_ops.is_empty() && db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite).is_ok() {
701 let transaction = db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite).unwrap();
702 if let Ok(store) = transaction.object_store("identifiers") {
703 for alias in id_alias_ops {
704 if let Err(e) = store.delete(&alias.into()) {
705 log::error!("Failed to delete old identifier alias in IndexedDB: {:?}", e);
706 }
707 }
708 }
709 }
710
711 let id_watcher_ops: Vec<_> = ops_to_flush.iter().filter_map(|op| {
712 if let PendingDbOperation::AddWatcher { alias, watcher_oobi } = op {
713 Some((alias.clone(), watcher_oobi.clone()))
714 } else {
715 None
716 }
717 }).collect();
718
719 if !id_watcher_ops.is_empty() && db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite).is_ok() {
720 if let Ok(transaction) = db.transaction_with_str_and_mode("identifiers", web_sys::IdbTransactionMode::Readwrite) {
721 for (alias, watcher_oobi) in id_watcher_ops {
722 if let Ok(store) = transaction.object_store("identifiers") {
723 let request = store.get(&alias.clone().into());
725 let alias_clone = alias.clone();
726 let on_success = Closure::wrap(Box::new(move |event: web_sys::Event| {
727 if let Some(result) = event.target()
728 .and_then(|t| t.dyn_into::<web_sys::IdbRequest>().ok())
729 .and_then(|r| r.result().ok())
730 {
731 if !result.is_undefined() {
733 if let Ok(item) = result.dyn_into::<js_sys::Object>() {
734 js_sys::Reflect::set(&item, &"watcher_url".into(), &watcher_oobi.url.to_string().clone().into()).unwrap_or_else(|_| {
736 log::error!("Failed to set watcher_url property");
737 false
738 });
739 js_sys::Reflect::set(&item, &"watcher_eid".into(), &watcher_oobi.eid.to_string().clone().into()).unwrap_or_else(|_| {
740 log::error!("Failed to set watcher_eid property");
741 false
742 });
743 if let serde_json::Value::String(scheme_str) = serde_json::to_value(&watcher_oobi.scheme).unwrap() {
744 js_sys::Reflect::set(&item, &"watcher_scheme".into(), &scheme_str.clone().into()).unwrap_or_else(|_| {
745 log::error!("Failed to set watcher_scheme property");
746 false
747 });
748 }
749 let transaction = store.transaction();
751 if let Ok(store_again) = transaction.object_store("identifiers") {
752 let key_for_put = alias_clone.clone();
753 if let Err(e) = store_again.put_with_key(&item, &key_for_put.into()) {
754 log::error!("Failed to update identifier watcher URL in IndexedDB: {:?}", e);
755 }
756 }
757 }
758 }
759 }
760 }) as Box<dyn FnMut(_)>);
761
762 request.unwrap().set_onsuccess(Some(on_success.as_ref().unchecked_ref()));
763 on_success.forget();
764 }
765 }
766 }
767 }
768 }
769
770 flag.set(false);
771 }
772 });
773 }
774 }) as Box<dyn FnMut()>);
775
776 let window = web_sys::window().unwrap();
778 let _ = window.set_interval_with_callback_and_timeout_and_arguments_0(
779 interval_callback.as_ref().unchecked_ref(),
780 1000
781 );
782
783 interval_callback.forget();
784}
785
786impl EventDatabase for IndexedDbDatabase {
787 type Error = IndexedDbError;
788 type LogDatabaseType = IndexedDbLogDatabase;
789
790 fn get_log_db(&self) -> Arc<Self::LogDatabaseType> {
791 self.log_db.clone()
792 }
793
794 fn add_kel_finalized_event(
795 &self,
796 signed_event: SignedEventMessage,
797 _id: &IdentifierPrefix,
798 ) -> Result<(), Self::Error> {
799 self.update_key_state(&signed_event.event_message)?;
800 self.log_db.log_event_with_new_transaction(&signed_event)?;
801 self.save_to_kel(&signed_event.event_message)?;
802
803 Ok(())
804 }
805
806 fn add_receipt_t(
807 &self,
808 receipt: SignedTransferableReceipt,
809 _id: &IdentifierPrefix,
810 ) -> Result<(), Self::Error> {
811 let digest = receipt.body.receipted_event_digest;
812 let transferable = Transferable::Seal(receipt.validator_seal, receipt.signatures);
813 self.log_db.insert_trans_receipt(&digest, &[transferable])
814 }
815
816 fn add_receipt_nt(
817 &self,
818 receipt: SignedNontransferableReceipt,
819 _id: &IdentifierPrefix,
820 ) -> Result<(), Self::Error> {
821 let receipted_event_digest = receipt.body.receipted_event_digest;
822 let receipts = receipt.signatures;
823 self.log_db.insert_nontrans_receipt(&receipted_event_digest, &receipts)
824 }
825
826 fn get_key_state(&self, id: &IdentifierPrefix) -> Option<IdentifierState> {
827 let key = id.to_string();
828 self.key_states.borrow().get(&key).cloned()
829 }
830
831 fn get_kel_finalized_events(
832 &self,
833 params: QueryParameters,
834 ) -> Option<impl DoubleEndedIterator<Item = timestamped::TimestampedSignedEventMessage>> {
835 match params {
836 QueryParameters::BySn { id, sn } => {
837 self.get_kel(&id, sn, 1)
838 .map(|events| events.into_iter())
839 }
840 QueryParameters::Range { id, start, limit } => {
841 self.get_kel(&id, start, limit)
842 .map(|events| events.into_iter())
843 }
844 QueryParameters::All { id } => {
845 self.get_full_kel(id)
846 .map(|events| events.into_iter())
847 }
848 }
849 }
850
851 fn get_receipts_t(
852 &self,
853 params: QueryParameters,
854 ) -> Option<impl DoubleEndedIterator<Item = Transferable>> {
855 match params {
856 QueryParameters::BySn { id, sn } => {
857 let key = id.to_string();
858 let digest = self.kels.borrow().get(&(key, sn)).cloned()?;
859 let receipts = self.log_db.get_trans_receipts(&digest).ok()?;
860 Some(receipts.collect::<Vec<_>>().into_iter())
861 }
862 QueryParameters::Range {..} | QueryParameters::All {..} => {
863 None
865 }
866 }
867 }
868
869 fn get_receipts_nt(
870 &self,
871 params: QueryParameters,
872 ) -> Option<impl DoubleEndedIterator<Item = SignedNontransferableReceipt>> {
873 match params {
874 QueryParameters::BySn { id, sn } => self
875 .get_nontrans_receipts_range(&id.to_string(), sn, 1)
876 .ok()
877 .map(|e| e.into_iter()),
878 QueryParameters::Range { id, start, limit } => self
879 .get_nontrans_receipts_range(&id.to_string(), start, limit)
880 .ok()
881 .map(|e| e.into_iter()),
882 QueryParameters::All { id } => self
883 .get_nontrans_receipts_range(&id.to_string(), 0, u64::MAX)
884 .ok()
885 .map(|e| e.into_iter()),
886 }
887 }
888
889 fn accept_to_kel(&self, event: &KeriEvent<KeyEvent>) -> Result<(), Self::Error> {
890 self.save_to_kel(event)?;
891 self.update_key_state(event)?;
892 Ok(())
893 }
894
895 fn save_reply(&self, _reply: keri_core::query::reply_event::SignedReply) -> Result<(), Self::Error> {
896 Ok(())
898 }
899
900 fn get_reply(&self, _id: &IdentifierPrefix, _from_who: &IdentifierPrefix) -> Option<keri_core::query::reply_event::SignedReply> {
901 None
903 }
904}
905
906impl IndexedDbDatabase {
908 fn save_to_kel(&self, event: &KeriEvent<KeyEvent>) -> Result<(), IndexedDbError> {
909 let digest = event.digest()
910 .map_err(|_| IndexedDbError::EncodingFailed("Could not get event digest".to_string()))?;
911
912 let id = event.data.prefix.to_string();
913 let sn = event.data.sn;
914
915 self.kels.borrow_mut().insert((id.clone(), sn), digest.clone());
917
918 self.pending_operations.borrow_mut().push(PendingDbOperation::SaveKel {
920 id,
921 sn,
922 digest
923 });
924
925 Ok(())
926 }
927
928 fn update_key_state(&self, event: &KeriEvent<KeyEvent>) -> Result<(), IndexedDbError> {
929 let id = event.data.prefix.to_string();
930
931 let key_state = self.key_states.borrow()
933 .get(&id)
934 .cloned()
935 .unwrap_or_default();
936
937 let updated_state = key_state.apply(event)
939 .map_err(|_| IndexedDbError::DatabaseSaveFailed("Failed to apply event to key state".to_string()))?;
940
941 self.key_states.borrow_mut().insert(id.clone(), updated_state.clone());
943
944 self.pending_operations.borrow_mut().push(PendingDbOperation::SaveKeyState {
946 id,
947 state: updated_state
948 });
949
950 Ok(())
951 }
952
953 #[allow(dead_code)]
954 fn get_event_digest(&self, id: &IdentifierPrefix, sn: u64) -> Option<SelfAddressingIdentifier> {
955 let key = id.to_string();
956 self.kels.borrow().get(&(key, sn)).cloned()
957 }
958
959 fn get_kel(&self, id: &IdentifierPrefix, from: u64, limit: u64) -> Option<Vec<timestamped::TimestampedSignedEventMessage>> {
960 let id_str = id.to_string();
961 let mut events = Vec::new();
962
963 for sn in from..(from + limit) {
964 if let Some(digest) = self.kels.borrow().get(&(id_str.clone(), sn)) {
965 if let Ok(Some(event)) = self.log_db.get_signed_event(digest) {
966 events.push(event);
967 }
968 } else {
969 break;
970 }
971 }
972
973 if events.is_empty() {
974 None
975 } else {
976 Some(events)
977 }
978 }
979
980 fn get_full_kel(&self, id: &IdentifierPrefix) -> Option<Vec<timestamped::TimestampedSignedEventMessage>> {
981 let id_str = id.to_string();
982 let mut events = Vec::new();
983 let mut sn = 0;
984
985 while let Some(digest) = self.kels.borrow().get(&(id_str.clone(), sn)) {
987 if let Ok(Some(event)) = self.log_db.get_signed_event(digest) {
988 events.push(event);
989 sn += 1;
990 } else {
991 break;
992 }
993 };
994
995 if events.is_empty() {
996 None
997 } else {
998 Some(events)
999 }
1000 }
1001
1002 fn get_nontrans_receipts_range(
1003 &self,
1004 id: &str,
1005 start: u64,
1006 limit: u64,
1007 ) -> Result<Vec<SignedNontransferableReceipt>, IndexedDbError> {
1008 let mut receipts = Vec::new();
1010 let kels_map = self.kels.borrow();
1011
1012 let mut sequence_numbers = Vec::new();
1014 for sn in start..(start + limit) {
1015 if kels_map.contains_key(&(id.to_string(), sn)) {
1016 sequence_numbers.push(sn);
1017 }
1018 }
1019
1020 for sn in sequence_numbers {
1022 if let Some(said) = kels_map.get(&(id.to_string(), sn)) {
1023 if let Ok(nontrans) = self.log_db.get_nontrans_couplets_by_key(said) {
1025 if let Ok(identifier) = id.parse::<IdentifierPrefix>() {
1027 let rct = Receipt::new(SerializationFormats::JSON, said.clone(), identifier, start);
1029
1030 let signatures = nontrans
1032 .unwrap()
1033 .collect();
1034
1035 let signed_receipt = SignedNontransferableReceipt {
1036 body: rct,
1037 signatures,
1038 };
1039
1040 receipts.push(signed_receipt);
1041 }
1042 }
1043 }
1044 }
1045
1046 Ok(receipts)
1047 }
1048}
1049
1050impl EscrowCreator for IndexedDbDatabase {
1051 type EscrowDatabaseType = IndexedDbEscrowDatabase;
1052
1053 fn create_escrow_db(&self, table_name: &'static str) -> Self::EscrowDatabaseType {
1054
1055 IndexedDbEscrowDatabase::new(
1056 Arc::new(IndexedDbSnDatabase::new(Arc::new(()), table_name).unwrap()),
1057 self.log_db.clone(),
1058 )
1059 }
1060}
1061
1062impl TelEventDatabase for IndexedDbDatabase {
1063 fn new(_path: impl AsRef<std::path::Path>) -> Result<Self, teliox::error::Error>
1064 where
1065 Self: Sized,
1066 {
1067 Ok(Self::new())
1068 }
1069
1070 fn add_new_event(&self, event: VerifiableEvent, id: &IdentifierPrefix) -> Result<(), teliox::error::Error> {
1071 match event.event {
1072 Event::Vc(_) => {
1073 self.tel_events.borrow_mut()
1074 .entry(id.clone())
1075 .or_default()
1076 .push(event.clone());
1077
1078 self.pending_operations.borrow_mut()
1079 .push(PendingDbOperation::SaveTelEvent {
1080 id: id.clone(),
1081 event: event.clone()
1082 });
1083 },
1084 Event::Management(_) => {
1085 self.management_events.borrow_mut()
1086 .entry(id.clone())
1087 .or_default()
1088 .push(event.clone());
1089
1090 self.pending_operations.borrow_mut()
1091 .push(PendingDbOperation::SaveManagementEvent {
1092 id: id.clone(),
1093 event: event.clone()
1094 });
1095 },
1096 }
1097
1098 Ok(())
1099 }
1100
1101 fn get_events(
1102 &self,
1103 id: &IdentifierPrefix,
1104 ) -> Option<impl DoubleEndedIterator<Item = VerifiableEvent>> {
1105 if let Some(events) = self.tel_events.borrow().get(id) {
1106 return Some(events.clone().into_iter());
1107 }
1108
1109 None
1110 }
1111
1112 fn get_management_events(
1113 &self,
1114 id: &IdentifierPrefix,
1115 ) -> Option<impl DoubleEndedIterator<Item = VerifiableEvent>> {
1116 if let Some(events) = self.management_events.borrow().get(id) {
1117 return Some(events.clone().into_iter());
1118 }
1119
1120 None
1121 }
1122}