#![allow(clippy::mutable_key_type)]
use crate::actor::{Actor, ActorContext};
use crate::message::{BatchPut, Get, Message, Put};
use crate::types::*;
use crate::utils::FxHashMap;
use arena_btreemap::BTreeMap;
use async_trait::async_trait;
use log::{error, info, warn};
use parking_lot::RwLock;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use wasm_bindgen::prelude::*;
use wasm_bindgen_futures::spawn_local;
#[wasm_bindgen(module = "node:fs/promises")]
extern "C" {
fn readFile(path: &str) -> js_sys::Promise;
fn writeFile(path: &str, data: &js_sys::Uint8Array) -> js_sys::Promise;
fn mkdir(path: &str, options: &JsValue) -> js_sys::Promise;
}
#[wasm_bindgen(module = "node:path")]
extern "C" {
fn join(base: &str, name: &str) -> String;
}
pub struct WasmNodeFsStorage {
cache: Arc<RwLock<FxHashMap<String, Children>>>,
base_dir: String,
db_ready: Arc<AtomicBool>,
}
#[cfg(test)]
impl WasmNodeFsStorage {
pub(crate) fn base_dir_str(&self) -> &str {
&self.base_dir
}
pub(crate) fn is_ready(&self) -> bool {
self.db_ready.load(Ordering::SeqCst)
}
}
impl Default for WasmNodeFsStorage {
fn default() -> Self {
Self::new()
}
}
impl WasmNodeFsStorage {
pub fn new() -> Self {
Self {
cache: Arc::new(RwLock::new(FxHashMap::default())),
base_dir: "beam_data".to_string(),
db_ready: Arc::new(AtomicBool::new(false)),
}
}
pub fn with_dir(dir: &str) -> Self {
Self {
cache: Arc::new(RwLock::new(FxHashMap::default())),
base_dir: dir.to_string(),
db_ready: Arc::new(AtomicBool::new(false)),
}
}
fn node_path(&self, node_id: &str) -> String {
let safe_name = if node_id.is_empty() { "_root" } else { node_id };
join(&self.base_dir, safe_name)
}
fn serialize_children(children: &Children) -> Vec<u8> {
postcard::to_allocvec(children).unwrap_or_default()
}
fn deserialize_children(bytes: &[u8]) -> Children {
if bytes.is_empty() {
return BTreeMap::default();
}
match postcard::from_bytes::<Children>(bytes) {
Ok(children) => children,
Err(e) => {
warn!("WasmNodeFsStorage: 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 path = self.node_path(&get.node_id);
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();
spawn_local(async move {
let promise = readFile(&path);
match wasm_bindgen_futures::JsFuture::from(promise).await {
Ok(result) => {
let bytes = js_sys::Uint8Array::from(result);
let raw = bytes.to_vec();
let children = WasmNodeFsStorage::deserialize_children(&raw);
cache.write().insert(node_id.clone(), children.clone());
let reply_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,
};
let mut reply_nodes = BTreeMap::default();
reply_nodes.insert(node_id, reply_children);
let put = Put::new(reply_nodes, Some(get_id), my_addr);
let _ = from.send(Message::Put(put));
}
Err(_e) => {
let mut reply_nodes = BTreeMap::default();
reply_nodes.insert(node_id, BTreeMap::default());
let put = Put::new(reply_nodes, Some(get_id), my_addr);
let _ = from.send(Message::Put(put));
}
}
});
}
fn reply_get(&self, get: &Get, children: Children, ctx: &ActorContext) {
let reply_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_nodes = BTreeMap::default();
reply_nodes.insert(get.node_id.clone(), reply_children);
let put = Put::new(reply_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 self.db_ready.load(Ordering::SeqCst) {
for (node_id, _update_data) in put.updated_nodes.iter() {
let path = self.node_path(node_id);
let children = self.cache.read().get(node_id).cloned().unwrap_or_default();
let bytes = WasmNodeFsStorage::serialize_children(&children);
let js_bytes = js_sys::Uint8Array::from(&bytes[..]);
spawn_local(async move {
let promise = writeFile(&path, &js_bytes);
if let Err(e) = wasm_bindgen_futures::JsFuture::from(promise).await {
error!("WasmNodeFsStorage: writeFile error: {:?}", e);
}
});
}
}
self.send_put_ack(&put, ctx);
}
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 self.db_ready.load(Ordering::SeqCst) {
for put in batch.puts.iter() {
for (node_id, _update_data) in put.updated_nodes.iter() {
let path = self.node_path(node_id);
let children = self.cache.read().get(node_id).cloned().unwrap_or_default();
let bytes = WasmNodeFsStorage::serialize_children(&children);
let js_bytes = js_sys::Uint8Array::from(&bytes[..]);
spawn_local(async move {
let promise = writeFile(&path, &js_bytes);
if let Err(e) = wasm_bindgen_futures::JsFuture::from(promise).await {
error!("WasmNodeFsStorage: batch writeFile error: {:?}", e);
}
});
}
}
}
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));
}
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));
}
}
#[async_trait]
impl Actor for WasmNodeFsStorage {
async fn pre_start(&mut self, _ctx: &ActorContext) {
info!(
"WasmNodeFsStorage adapter starting (dir: {})",
self.base_dir
);
let base_dir = self.base_dir.clone();
let db_ready = self.db_ready.clone();
spawn_local(async move {
let opts = js_sys::Object::new();
js_sys::Reflect::set(&opts, &"recursive".into(), &true.into()).ok();
let promise = mkdir(&base_dir, &opts);
match wasm_bindgen_futures::JsFuture::from(promise).await {
Ok(_) => {
db_ready.store(true, Ordering::SeqCst);
info!("WasmNodeFsStorage: directory ready: {}", base_dir);
}
Err(e) => {
error!("WasmNodeFsStorage: mkdir error: {:?}", e);
}
}
});
}
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),
_ => {}
}
}
fn try_clone_storage(&self) -> Option<Box<dyn Actor>> {
None
}
}