Skip to main content

dkms_wasm/database/indexed_db/
mod.rs

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
85// SAFETY: In WebAssembly context, there's no true threading, so these are safe
86unsafe 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        // Get IndexedDB factory
134        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            // Open database
144            if let Ok(request) = factory.open(db_name) {
145                // Handle database upgrade needed (first time opening)
146                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                        // Create stores similar to redb tables
153                        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                // Handle successful open
165                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 and kels from IndexedDB
178                            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                            // Set up background flush
185                            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                // Handle errors
194                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
253// Helper function to load key states from IndexedDB
254fn 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
296// Helper function to load KELs from IndexedDB
297fn 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
341// Helper function to load TEL events from IndexedDB
342fn 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
389// Helper function to load management events from IndexedDB
390fn 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 we have a cursor
450                        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
515// Set up background flush for IndexedDB
516fn setup_background_flush(
517    db: web_sys::IdbDatabase, 
518    pending_ops: Rc<RefCell<Vec<PendingDbOperation>>>,
519    flush_flag: Rc<Cell<bool>>
520) {
521    // Create interval function
522    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                    // Get operations to process
533                    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                    // Process operations
543                    if !ops_to_flush.is_empty() {
544                        // Handle key state operations
545                        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                                    // Serialize state to JSON string
562                                    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                        // Handle KEL operations
574                        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                        // Handle TEL event operations
606                        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                                    // Serialize event to JSON string
623                                    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                        // Handle management event operations
635                        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                                    // Serialize event to JSON string
652                                    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                        // Handle identifiers event operations
664                        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                                        // First get the existing record
724                                        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 we have an existing record
732                                                if !result.is_undefined() {
733                                                    if let Ok(item) = result.dyn_into::<js_sys::Object>() {
734                                                        // Update the watcher_url field
735                                                        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                                                        // Put it back in the store
750                                                        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    // Set up interval
777    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                // For simplicity, not implementing range/all queries for receipts
864                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        // Not implemented for WASM yet
897        Ok(())
898    }
899
900    fn get_reply(&self, _id: &IdentifierPrefix, _from_who: &IdentifierPrefix) -> Option<keri_core::query::reply_event::SignedReply> {
901        // Not implemented for WASM yet
902        None
903    }
904}
905
906// Helper methods for IndexedDbDatabase
907impl 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        // Save to in-memory map
916        self.kels.borrow_mut().insert((id.clone(), sn), digest.clone());
917        
918        // Queue for persistence
919        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        // Get current state or default
932        let key_state = self.key_states.borrow()
933            .get(&id)
934            .cloned()
935            .unwrap_or_default();
936        
937        // Apply event to state
938        let updated_state = key_state.apply(event)
939            .map_err(|_| IndexedDbError::DatabaseSaveFailed("Failed to apply event to key state".to_string()))?;
940        
941        // Save updated state
942        self.key_states.borrow_mut().insert(id.clone(), updated_state.clone());
943        
944        // Queue for persistence
945        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        // Find all events for this identifier
986        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        // Get all sequence numbers in the range for this identifier
1009        let mut receipts = Vec::new();
1010        let kels_map = self.kels.borrow();
1011        
1012        // Collect sequence numbers first to make iteration easier
1013        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        // Create receipts for each event in the range
1021        for sn in sequence_numbers {
1022            if let Some(said) = kels_map.get(&(id.to_string(), sn)) {
1023                // Get non-transferable couplets for this digest
1024                if let Ok(nontrans) = self.log_db.get_nontrans_couplets_by_key(said) {
1025                    // Parse identifier
1026                    if let Ok(identifier) = id.parse::<IdentifierPrefix>() {
1027                        // Create receipt
1028                        let rct = Receipt::new(SerializationFormats::JSON, said.clone(), identifier, start);
1029                        
1030                        // Create signed receipt with signatures
1031                        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}