use crate::actor::{Actor, ActorContext};
use crate::message::{BatchPut, Get, Message, Put};
use crate::types::*;
use crate::utils::{FxHashMap, FxHashSet};
use arena_btreemap::BTreeMap;
use async_trait::async_trait;
use base64::Engine as _;
use log::{error, info, warn};
use std::sync::Arc;
use wasm_bindgen::closure::Closure;
use wasm_bindgen::{JsCast, JsValue};
use web_sys::{IdbDatabase, IdbRequest, IdbTransactionMode, IdbVersionChangeEvent};
pub struct WasmIdbStorage {
db: Option<IdbDatabase>,
cache: Arc<parking_lot::RwLock<FxHashMap<String, Children>>>,
db_name: String,
store_name: String,
db_ready: bool,
}
impl Default for WasmIdbStorage {
fn default() -> Self {
Self::new()
}
}
impl WasmIdbStorage {
pub fn new() -> Self {
Self {
db: None,
cache: Arc::new(parking_lot::RwLock::new(FxHashMap::default())),
db_name: "beam".to_string(),
store_name: "beam_data".to_string(),
db_ready: false,
}
}
pub fn with_name(db_name: &str) -> Self {
Self {
db: None,
cache: Arc::new(parking_lot::RwLock::new(FxHashMap::default())),
db_name: db_name.to_string(),
store_name: "beam_data".to_string(),
db_ready: false,
}
}
fn open_db_async(&mut self, on_ready: Box<dyn FnOnce(IdbDatabase)>) {
let window = match web_sys::window() {
Some(w) => w,
None => {
error!("WasmIdbStorage: no window available");
return;
}
};
let idb = match window.indexed_db() {
Ok(Some(idb)) => idb,
_ => {
error!("WasmIdbStorage: no IndexedDB available");
return;
}
};
let request = match idb.open_with_u32(&self.db_name, 1) {
Ok(r) => r,
Err(e) => {
error!("WasmIdbStorage: open error: {:?}", e);
return;
}
};
let store_name = self.store_name.clone();
let onupgradeneeded: Closure<dyn FnMut(IdbVersionChangeEvent)> =
Closure::new(move |event: IdbVersionChangeEvent| {
let request: IdbRequest = event.target().unwrap().unchecked_into();
if let Ok(db_result) = request.result() {
let db: IdbDatabase = db_result.unchecked_into();
let _ = db.create_object_store(&store_name);
info!("WasmIdbStorage: created object store");
}
});
request.set_onupgradeneeded(Some(onupgradeneeded.as_ref().unchecked_ref()));
onupgradeneeded.forget();
let mut on_ready_opt = Some(on_ready);
let on_ready_wrapped: Closure<dyn FnMut(JsValue)> = Closure::new(move |_event: JsValue| {
let event: web_sys::Event = _event.unchecked_into();
let target = event.target().unwrap();
let request: IdbRequest = target.unchecked_into();
match request.result() {
Ok(db_result) => {
let db: IdbDatabase = db_result.unchecked_into();
info!("WasmIdbStorage: database opened");
if let Some(cb) = on_ready_opt.take() {
cb(db);
}
}
Err(e) => {
error!("WasmIdbStorage: open success callback error: {:?}", e);
}
}
});
request.set_onsuccess(Some(on_ready_wrapped.as_ref().unchecked_ref()));
on_ready_wrapped.forget();
let onerror: Closure<dyn FnMut(JsValue)> = Closure::new(move |event: JsValue| {
error!("WasmIdbStorage: database open error: {:?}", event);
});
request.set_onerror(Some(onerror.as_ref().unchecked_ref()));
onerror.forget();
}
fn serialize_children(children: &Children) -> String {
let map: FxHashMap<String, (Value, f64)> = children
.iter()
.map(|(k, v)| (k.clone(), (v.value.clone(), v.updated_at)))
.collect();
let bytes = postcard::to_allocvec(&map).unwrap_or_default();
base64::engine::general_purpose::STANDARD.encode(&bytes)
}
fn deserialize_children(stored: &str) -> Children {
if stored.is_empty() {
return BTreeMap::default();
}
if let Ok(bytes) = base64::engine::general_purpose::STANDARD.decode(stored) {
if let Ok(map) = postcard::from_bytes::<FxHashMap<String, (Value, f64)>>(&bytes) {
return map
.into_iter()
.map(|(k, (value, updated_at))| (k, NodeData { value, updated_at }))
.collect();
}
}
match serde_json::from_str::<FxHashMap<String, (Value, f64)>>(stored) {
Ok(map) => map
.into_iter()
.map(|(k, (value, updated_at))| (k, NodeData { value, updated_at }))
.collect(),
Err(e) => {
warn!("WasmIdbStorage: deserialize error: {}", e);
BTreeMap::default()
}
}
}
fn handle_get(&self, get: Get, ctx: &ActorContext) {
if let Some(children) = self.cache.read().get(&get.node_id).cloned() {
self.reply_get(&get, children, ctx);
return;
}
let db = match &self.db {
Some(db) => db,
None => {
self.reply_get_empty(&get, ctx);
return;
}
};
let node_id = get.node_id.clone();
let from = get.from.clone();
let get_id = get.id.clone();
let child_key = get.child_key.clone();
let my_addr = ctx.addr.clone();
let cache = self.cache.clone();
let tx = match db
.transaction_with_str_and_mode(&self.store_name, IdbTransactionMode::Readonly)
{
Ok(tx) => tx,
Err(e) => {
error!("WasmIdbStorage: transaction error: {:?}", e);
self.reply_get_empty(&get, ctx);
return;
}
};
let store = match tx.object_store(&self.store_name) {
Ok(s) => s,
Err(e) => {
error!("WasmIdbStorage: object store error: {:?}", e);
self.reply_get_empty(&get, ctx);
return;
}
};
let request = match store.get(&JsValue::from_str(&node_id)) {
Ok(r) => r,
Err(e) => {
error!("WasmIdbStorage: get request error: {:?}", e);
self.reply_get_empty(&get, ctx);
return;
}
};
let onsuccess: Closure<dyn FnMut(JsValue)> = Closure::new(move |event: JsValue| {
let event: web_sys::Event = event.unchecked_into();
let target = event.target().unwrap();
let request: IdbRequest = target.unchecked_into();
match request.result() {
Ok(result) => {
if result.is_null() || result.is_undefined() {
let mut reply_with_nodes = BTreeMap::default();
reply_with_nodes.insert(node_id.clone(), BTreeMap::default());
let put = Put::new(reply_with_nodes, Some(get_id.clone()), my_addr.clone());
let _ = from.send(Message::Put(put));
return;
}
let stored = result.as_string().unwrap_or_default();
let children = WasmIdbStorage::deserialize_children(&stored);
cache.write().insert(node_id.clone(), children.clone());
let reply_with_children = match &child_key {
Some(ck) => match children.get(ck) {
Some(cv) => {
let mut r = BTreeMap::default();
r.insert(ck.clone(), cv.clone());
r
}
None => return,
},
None => children.clone(),
};
let mut reply_with_nodes = BTreeMap::default();
reply_with_nodes.insert(node_id.clone(), reply_with_children);
let put = Put::new(reply_with_nodes, Some(get_id.clone()), my_addr.clone());
let _ = from.send(Message::Put(put));
}
Err(e) => {
error!("WasmIdbStorage: get result error: {:?}", e);
}
}
});
request.set_onsuccess(Some(onsuccess.as_ref().unchecked_ref()));
onsuccess.forget();
let onerror: Closure<dyn FnMut(JsValue)> = Closure::new(move |event: JsValue| {
error!("WasmIdbStorage: get request error: {:?}", event);
});
request.set_onerror(Some(onerror.as_ref().unchecked_ref()));
onerror.forget();
}
fn reply_get(&self, get: &Get, children: Children, ctx: &ActorContext) {
let reply_with_children = match &get.child_key {
Some(ck) => match children.get(ck) {
Some(cv) => {
let mut r = BTreeMap::default();
r.insert(ck.clone(), cv.clone());
r
}
None => return,
},
None => children,
};
let mut reply_with_nodes = BTreeMap::default();
reply_with_nodes.insert(get.node_id.clone(), reply_with_children);
#[allow(clippy::mutable_key_type)]
let mut recipients = FxHashSet::default();
recipients.insert(get.from.clone());
let put = Put::new(reply_with_nodes, Some(get.id.clone()), ctx.addr.clone());
let _ = get.from.send(Message::Put(put));
}
fn reply_get_empty(&self, get: &Get, ctx: &ActorContext) {
let mut reply_with_nodes = BTreeMap::default();
reply_with_nodes.insert(get.node_id.clone(), BTreeMap::default());
let put = Put::new(reply_with_nodes, Some(get.id.clone()), ctx.addr.clone());
let _ = get.from.send(Message::Put(put));
}
fn handle_put(&mut self, put: Put, ctx: &ActorContext) {
for (node_id, update_data) in put.updated_nodes.iter().rev() {
let mut write = self.cache.write();
if let Some(children) = write.get_mut(node_id) {
for (child_id, child_data) in update_data {
if let Some(existing) = children.get(child_id) {
if child_data.updated_at >= existing.updated_at {
children.insert(child_id.clone(), child_data.clone());
}
} else {
children.insert(child_id.clone(), child_data.clone());
}
}
} else {
write.insert(node_id.to_string(), update_data.clone());
}
}
if let Some(db) = &self.db {
for (node_id, update_data) in put.updated_nodes.iter() {
let serialized = WasmIdbStorage::serialize_children(update_data);
let tx = match db
.transaction_with_str_and_mode(&self.store_name, IdbTransactionMode::Readwrite)
{
Ok(tx) => tx,
Err(e) => {
error!("WasmIdbStorage: put transaction error: {:?}", e);
continue;
}
};
let store = match tx.object_store(&self.store_name) {
Ok(s) => s,
Err(e) => {
error!("WasmIdbStorage: put object store error: {:?}", e);
continue;
}
};
let _ = store
.put_with_key(&JsValue::from_str(&serialized), &JsValue::from_str(node_id));
}
}
self.send_put_ack(&put, ctx);
}
fn send_put_ack(&self, put: &Put, ctx: &ActorContext) {
let mut ack_children = BTreeMap::default();
ack_children.insert(
"_ack".to_string(),
NodeData {
value: Value::Text("ok".to_string()),
updated_at: web_time::SystemTime::now()
.duration_since(web_time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as f64,
},
);
let mut nodes = BTreeMap::default();
nodes.insert("_ack".to_string(), ack_children);
let ack = Put::new(nodes, Some(put.id.clone()), ctx.addr.clone());
let _ = put.from.send(Message::Put(ack));
}
fn handle_batch_put(&mut self, batch: BatchPut, ctx: &ActorContext) {
for put in batch.puts.iter() {
for (node_id, update_data) in put.updated_nodes.iter().rev() {
let mut write = self.cache.write();
if let Some(children) = write.get_mut(node_id) {
for (child_id, child_data) in update_data {
if let Some(existing) = children.get(child_id) {
if child_data.updated_at >= existing.updated_at {
children.insert(child_id.clone(), child_data.clone());
}
} else {
children.insert(child_id.clone(), child_data.clone());
}
}
} else {
write.insert(node_id.to_string(), update_data.clone());
}
}
}
if let Some(db) = &self.db {
for put in batch.puts.iter() {
for (node_id, update_data) in put.updated_nodes.iter() {
let serialized = WasmIdbStorage::serialize_children(update_data);
let tx = match db.transaction_with_str_and_mode(
&self.store_name,
IdbTransactionMode::Readwrite,
) {
Ok(tx) => tx,
Err(e) => {
error!("WasmIdbStorage: batch put tx error: {:?}", e);
continue;
}
};
let store = match tx.object_store(&self.store_name) {
Ok(s) => s,
Err(e) => {
error!("WasmIdbStorage: batch put store error: {:?}", e);
continue;
}
};
let _ = store
.put_with_key(&JsValue::from_str(&serialized), &JsValue::from_str(node_id));
}
}
}
let mut ack_children = BTreeMap::default();
ack_children.insert(
"_ack".to_string(),
NodeData {
value: Value::Text("ok".to_string()),
updated_at: web_time::SystemTime::now()
.duration_since(web_time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as f64,
},
);
let mut nodes = BTreeMap::default();
nodes.insert("_ack".to_string(), ack_children);
let ack = Put::new(nodes, Some(batch.id.clone()), ctx.addr.clone());
let _ = batch.from.send(Message::Put(ack));
}
}
#[async_trait]
impl Actor for WasmIdbStorage {
async fn pre_start(&mut self, _ctx: &ActorContext) {
info!("WasmIdbStorage adapter starting");
let db_handle: Arc<parking_lot::Mutex<Option<IdbDatabase>>> =
Arc::new(parking_lot::Mutex::new(None));
let db_handle_for_cb = db_handle.clone();
let _self_cache = self.cache.clone();
self.open_db_async(Box::new(move |db| {
*db_handle_for_cb.lock() = Some(db);
info!("WasmIdbStorage: database ready");
}));
if let Some(db) = db_handle.lock().take() {
self.db = Some(db);
self.db_ready = true;
}
}
async fn handle(&mut self, message: Arc<Message>, ctx: &ActorContext) {
match &*message {
Message::Get(get) => self.handle_get(get.clone(), ctx),
Message::Put(put) => self.handle_put(put.clone(), ctx),
Message::Flush(flush) => {
let mut ack = BTreeMap::default();
ack.insert(
"_flushed".to_string(),
NodeData {
value: Value::Text("true".to_string()),
updated_at: web_time::SystemTime::now()
.duration_since(web_time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as f64,
},
);
let mut nodes = BTreeMap::default();
nodes.insert("_ack".to_string(), ack);
let put = Put::new(nodes, Some(flush.id.clone()), ctx.addr.clone());
let _ = flush.from.send(Message::Put(put));
}
Message::BatchPut(batch) => self.handle_batch_put(batch.clone(), ctx),
_ => {}
}
}
}