dkms_wasm/database/indexed_db/
sn_database.rs1use cesrox::primitives::CesrPrimitive;
2use keri_core::{database::SequencedEventDatabase, prefix::IdentifierPrefix};
3use said::SelfAddressingIdentifier;
4use std::{
5 cell::{Cell, RefCell},
6 collections::HashMap,
7 rc::Rc,
8 sync::{
9 Arc, RwLock,
10 },
11};
12
13use std::str::FromStr;
14
15use super::IndexedDbError;
16
17use wasm_bindgen::prelude::*;
18use wasm_bindgen_futures::spawn_local;
19
20use web_sys::{
21 IdbDatabase, IdbOpenDbRequest,
22};
23
24unsafe impl Send for IndexedDbSnDatabase {}
26unsafe impl Sync for IndexedDbSnDatabase {}
27
28pub struct IndexedDbSnDatabase {
29 data: RwLock<HashMap<(String, u64), Vec<SelfAddressingIdentifier>>>,
30 pending_operations: Rc<RefCell<Vec<PendingOperation>>>,
31}
32
33enum PendingOperation {
34 Insert {
35 id_str: String,
36 sn: u64,
37 digest: SelfAddressingIdentifier,
38 },
39 Remove {
40 id_str: String,
41 sn: u64,
42 },
43}
44
45impl SequencedEventDatabase for IndexedDbSnDatabase {
46 type DatabaseType = ();
47 type Error = IndexedDbError;
48 type DigestIter = Box<dyn Iterator<Item = SelfAddressingIdentifier>>;
49
50 fn new(
51 _db: Arc<Self::DatabaseType>,
52 table_name: &'static str,
53 ) -> Result<Self, Self::Error>
54 where
55 Self: Sized,
56 {
57 let pending_operations = Rc::new(RefCell::new(Vec::new()));
58 let flush_in_progress = Rc::new(Cell::new(false));
59
60 let window =
62 web_sys::window().expect("should have a window in this context");
63 let factory = match window.indexed_db() {
64 Ok(Some(factory)) => factory,
65 Ok(None) => {
66 log::error!("IndexedDB not available");
67 return Ok(Self {
68 data: RwLock::new(HashMap::new()),
69 pending_operations,
70 });
71 }
72 Err(e) => {
73 log::error!("Failed to get IndexedDB: {:?}", e);
74 return Ok(Self {
75 data: RwLock::new(HashMap::new()),
76 pending_operations,
77 });
78 }
79 };
80
81 let db_name = table_name.to_string();
82 let pending_ops_clone = pending_operations.clone();
83 let flush_flag = flush_in_progress.clone();
84
85 match factory.open(&db_name) {
87 Ok(open_request) => {
88 let upgrade_needed_cb = Closure::wrap(Box::new(
89 move |event: web_sys::IdbVersionChangeEvent| {
90 let target = event.target().unwrap();
91 let request =
92 target.dyn_into::<IdbOpenDbRequest>().unwrap();
93 if let Some(db) = request
94 .result()
95 .ok()
96 .and_then(|res| res.dyn_into::<IdbDatabase>().ok())
97 {
98 let mut params =
100 web_sys::IdbObjectStoreParameters::new();
101 params.key_path(Some(&JsValue::from_str("key")));
102
103 if let Err(e) = db
104 .create_object_store_with_optional_parameters(
105 "events", ¶ms,
106 )
107 {
108 log::error!(
109 "Failed to create events store: {:?}",
110 e
111 );
112 }
113 }
114 },
115 )
116 as Box<dyn FnMut(_)>);
117
118 open_request.set_onupgradeneeded(Some(
119 upgrade_needed_cb.as_ref().unchecked_ref(),
120 ));
121 upgrade_needed_cb.forget();
122
123 let data_clone = Rc::new(RefCell::new(HashMap::new()));
124
125 let success_cb = Closure::wrap(Box::new(
127 move |event: web_sys::Event| {
128 let target = event.target().unwrap();
129 let request =
130 target.dyn_into::<IdbOpenDbRequest>().unwrap();
131 if let Ok(db) = request
132 .result()
133 .and_then(|res| res.dyn_into::<IdbDatabase>())
134 {
135 let db_for_load = db.clone();
137 let data_weak = Rc::downgrade(&data_clone);
138 spawn_local(async move {
139 if let Some(data_ref) = data_weak.upgrade() {
140 load_existing_data(&db_for_load, data_ref)
141 .await;
142 }
143 });
144
145 let db_clone = db.clone();
147 let pending_ops = pending_ops_clone.clone();
148 let flush_flag = flush_flag.clone();
149
150 let interval_callback =
152 Closure::wrap(Box::new(move || {
153 let db = db_clone.clone();
154 let pending = pending_ops.clone();
155 let flag = flush_flag.clone();
156
157 if !flag.get() {
158 flag.set(true);
159
160 spawn_local({
161 let pending = pending.clone();
162 let flag = flag.clone();
163
164 async move {
165 let ops_to_flush = {
166 let mut pending_borrow =
167 pending.borrow_mut();
168 if pending_borrow.is_empty()
169 {
170 Vec::new()
171 } else {
172 pending_borrow
173 .drain(..)
174 .collect()
175 }
176 };
177
178 if !ops_to_flush.is_empty() {
179 flush_pending_operations(
180 &db,
181 ops_to_flush,
182 )
183 .await;
184 }
185
186 flag.set(false);
187 }
188 });
189 }
190 })
191 as Box<dyn FnMut()>);
192
193 let interval_fn = interval_callback
195 .as_ref()
196 .unchecked_ref::<js_sys::Function>(
197 );
198 let _interval_id = window.set_interval_with_callback_and_timeout_and_arguments_0(
199 interval_fn,
200 1000
201 ).expect("should be able to set interval");
202 interval_callback.forget(); }
204 },
205 )
206 as Box<dyn FnMut(_)>);
207
208 open_request
209 .set_onsuccess(Some(success_cb.as_ref().unchecked_ref()));
210 success_cb.forget();
211
212 let error_cb =
214 Closure::wrap(Box::new(|event: web_sys::Event| {
215 log::error!("Failed to open IndexedDB: {:?}", event);
216 }) as Box<dyn FnMut(_)>);
217
218 open_request
219 .set_onerror(Some(error_cb.as_ref().unchecked_ref()));
220 error_cb.forget();
221 }
222 Err(e) => {
223 log::error!("Failed to create open request: {:?}", e);
224 }
225 }
226
227 Ok(Self {
228 data: RwLock::new(HashMap::new()),
229 pending_operations,
230 })
231 }
232
233 fn insert(
234 &self,
235 id: &IdentifierPrefix,
236 sn: u64,
237 digest: &SelfAddressingIdentifier,
238 ) -> Result<(), Self::Error> {
239 let mut data = self.data.write().unwrap();
240 let key = (id.to_str(), sn);
241 data.entry(key.clone()).or_default().push(digest.clone());
242
243 let mut pending = self.pending_operations.borrow_mut();
244 pending.push(PendingOperation::Insert {
245 id_str: key.0,
246 sn: key.1,
247 digest: digest.clone(),
248 });
249
250 Ok(())
251 }
252
253 fn get(
254 &self,
255 id: &IdentifierPrefix,
256 sn: u64,
257 ) -> Result<Self::DigestIter, Self::Error> {
258 let data = self.data.read().unwrap();
259 let key = (id.to_str(), sn);
260
261 if let Some(digests) = data.get(&key) {
262 Ok(Box::new(digests.clone().into_iter()))
263 } else {
264 Ok(Box::new(std::iter::empty()))
265 }
266 }
267
268 fn get_greater_than(
269 &self,
270 id: &IdentifierPrefix,
271 sn: u64,
272 ) -> Result<Self::DigestIter, Self::Error> {
273 let data = self.data.read().unwrap();
274 let id_str = id.to_str();
275
276 let mut result = Vec::new();
277 for ((prefix, seq), digests) in data.iter() {
278 if prefix == &id_str && *seq >= sn {
279 result.extend(digests.clone());
280 }
281 }
282
283 Ok(Box::new(result.into_iter()))
284 }
285
286 fn remove(
287 &self,
288 id: &IdentifierPrefix,
289 sn: u64,
290 digest: &SelfAddressingIdentifier,
291 ) -> Result<(), Self::Error> {
292 let mut data = self.data.write().unwrap();
293 let key = (id.to_str(), sn);
294
295 if let Some(digests) = data.get_mut(&key) {
296 digests.retain(|d| d != digest);
297
298 let mut pending = self.pending_operations.borrow_mut();
300 pending.push(PendingOperation::Remove {
301 id_str: key.0.clone(),
302 sn: key.1,
303 });
304
305 if digests.is_empty() {
306 data.remove(&key);
307 }
308 Ok(())
309 } else {
310 Err(IndexedDbError::NotFound(format!("{:?} at sn {}", id, sn)))
311 }
312 }
313}
314
315async fn flush_pending_operations(
317 db: &IdbDatabase,
318 operations: Vec<PendingOperation>,
319) {
320 if let Ok(transaction) = db.transaction_with_str_and_mode(
321 "events",
322 web_sys::IdbTransactionMode::Readwrite,
323 ) {
324 if let Ok(store) = transaction.object_store("events") {
325 for op in operations {
326 match op {
327 PendingOperation::Insert { id_str, sn, digest } => {
328 let key = format!("{}:{}", id_str, sn);
329 let value = js_sys::Object::new();
330 let _ = js_sys::Reflect::set(
331 &value,
332 &"key".into(),
333 &key.into(),
334 );
335 let _ = js_sys::Reflect::set(
336 &value,
337 &"id".into(),
338 &id_str.into(),
339 );
340 let _ = js_sys::Reflect::set(
341 &value,
342 &"sn".into(),
343 &sn.into(),
344 );
345 let _ = js_sys::Reflect::set(
346 &value,
347 &"digest".into(),
348 &digest.to_str().into(),
349 );
350
351 if let Err(e) = store.put(&value) {
352 log::error!(
353 "Failed to store event in IndexedDB: {:?}",
354 e
355 );
356 }
357 }
358 PendingOperation::Remove { id_str, sn } => {
359 let key = format!("{}:{}", id_str, sn);
360 if let Err(e) = store.delete(&key.into()) {
361 log::error!(
362 "Failed to remove event from IndexedDB: {:?}",
363 e
364 );
365 }
366 }
367 }
368 }
369 }
370 }
371}
372
373async fn load_existing_data(
374 db: &IdbDatabase,
375 data: Rc<RefCell<HashMap<(String, u64), Vec<SelfAddressingIdentifier>>>>,
376) {
377 if let Ok(transaction) = db.transaction_with_str_and_mode(
378 "events",
379 web_sys::IdbTransactionMode::Readonly,
380 ) {
381 if let Ok(store) = transaction.object_store("events") {
382 let request = match store.open_cursor() {
384 Ok(request) => request,
385 Err(e) => {
386 log::error!("Failed to open cursor: {:?}", e);
387 return;
388 }
389 };
390
391 let cursor_success = Closure::wrap(Box::new(
393 move |event: web_sys::Event| {
394 let target =
395 event.target().expect("event should have target");
396 let request = target
397 .dyn_into::<web_sys::IdbRequest>()
398 .expect("target should be IdbRequest");
399
400 if let Ok(cursor_val) = request.result() {
401 if !cursor_val.is_null() {
402 if let Ok(cursor) = cursor_val
403 .dyn_into::<web_sys::IdbCursorWithValue>(
404 ) {
405 if let Ok(value) = cursor.value() {
406 if let (
408 Some(id),
409 Some(sn),
410 Some(digest_str),
411 ) = (
412 js_sys::Reflect::get(
413 &value,
414 &"id".into(),
415 )
416 .ok()
417 .and_then(|v| v.as_string()),
418 js_sys::Reflect::get(
419 &value,
420 &"sn".into(),
421 )
422 .ok()
423 .and_then(|v| {
424 if !v.is_bigint() {
425 return None;
426 }
427 let bingint = v
428 .dyn_into::<js_sys::BigInt>()
429 .ok()?;
430 bingint.to_string(10).ok().and_then(
431 |s| {
432 s.as_string()
433 .unwrap()
434 .parse::<u64>()
435 .ok()
436 },
437 )
438 }),
439 js_sys::Reflect::get(
440 &value,
441 &"digest".into(),
442 )
443 .ok()
444 .and_then(|v| v.as_string()),
445 ) {
446 if let Ok(digest) =
447 SelfAddressingIdentifier::from_str(
448 &digest_str,
449 )
450 {
451 let mut data_mut =
452 data.borrow_mut();
453 data_mut
454 .entry((id, sn))
455 .or_default()
456 .push(digest);
457 } else {
458 log::warn!(
459 "Failed to parse digest: {}",
460 digest_str
461 );
462 }
463 } else {
464 log::warn!(
465 "Stored object missing required fields"
466 );
467 }
468 }
469
470 if let Err(e) = cursor.continue_() {
472 log::error!(
473 "Error continuing cursor: {:?}",
474 e
475 );
476 }
477 }
478 }
479 } else {
480 log::error!("Failed to get cursor from result");
481 }
482 },
483 )
484 as Box<dyn FnMut(_)>);
485
486 let cursor_error =
487 Closure::wrap(Box::new(move |event: web_sys::Event| {
488 log::error!("Error in IndexedDB cursor: {:?}", event);
489 }) as Box<dyn FnMut(_)>);
490
491 request
493 .set_onsuccess(Some(cursor_success.as_ref().unchecked_ref()));
494 request.set_onerror(Some(cursor_error.as_ref().unchecked_ref()));
495
496 cursor_success.forget();
498 cursor_error.forget();
499 }
500 }
501}