#![allow(clippy::mutable_key_type)]
use crate::utils::{FxHashMap, FxHashSet};
use arena_btreemap::BTreeMap;
use crate::actor::{Actor, ActorContext};
use crate::message::{BatchPut, Get, Message, Put};
use crate::types::*;
use async_trait::async_trait;
use log::{debug, info};
use parking_lot::RwLock;
use std::sync::Arc;
#[derive(Clone)]
pub struct MemoryStorage {
store: Arc<RwLock<FxHashMap<String, Children>>>,
}
impl Default for MemoryStorage {
fn default() -> Self {
Self::new()
}
}
impl MemoryStorage {
pub fn new() -> Self {
MemoryStorage {
store: Arc::new(RwLock::new(FxHashMap::default())),
}
}
fn handle_get(&self, get: &Get, ctx: &ActorContext) {
if let Some(children) = self.store.read().get(&get.node_id).cloned() {
debug!("have {}: {:?}", get.node_id, children);
let reply_with_children = match &get.child_key {
Some(child_key) => {
match children.get(child_key) {
Some(child_val) => {
let mut r = BTreeMap::default();
r.insert(child_key.clone(), child_val.clone());
r
}
None => {
return;
}
}
}
None => children.clone(), };
let mut reply_with_nodes = BTreeMap::default();
reply_with_nodes.insert(get.node_id.clone(), reply_with_children);
let mut recipients = FxHashSet::default();
recipients.insert(get.from.clone());
let my_addr = ctx.addr.clone();
let put = Put::new(reply_with_nodes, Some(get.id.clone()), my_addr);
let _ = get.from.send(Message::Put(put));
} else {
debug!("have not {}", get.node_id);
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(&self, put: &Put, ctx: &ActorContext) {
let put_result = self.apply_put(put);
self.send_put_ack(put, &put_result, ctx);
}
fn apply_put(&self, put: &Put) -> Result<(), String> {
for (node_id, update_data) in put.updated_nodes.iter().rev() {
debug!("saving k-v {}: {:?}", node_id, update_data);
let mut write = self.store.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());
}
}
Ok(())
}
fn send_put_ack(&self, put: &Put, result: &Result<(), String>, ctx: &ActorContext) {
let mut ack_children = BTreeMap::default();
match result {
Ok(()) => {
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,
},
);
}
Err(msg) => {
ack_children.insert(
"_err".to_string(),
NodeData {
value: Value::Text(msg.clone()),
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(&self, batch: &BatchPut, ctx: &ActorContext) {
let mut last_err: Option<String> = None;
for put in batch.puts.iter() {
if let Err(e) = self.apply_put(put) {
last_err = Some(e);
}
}
let result = last_err.map(Err).unwrap_or(Ok(()));
self.send_batch_put_ack(batch, &result, ctx);
}
fn send_batch_put_ack(
&self,
batch: &BatchPut,
result: &Result<(), String>,
ctx: &ActorContext,
) {
let mut ack_children = BTreeMap::default();
match result {
Ok(()) => {
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,
},
);
}
Err(msg) => {
ack_children.insert(
"_err".to_string(),
NodeData {
value: Value::Text(msg.clone()),
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 MemoryStorage {
async fn pre_start(&mut self, _ctx: &ActorContext) {
info!("MemoryStorage adapter starting");
}
async fn handle(&mut self, message: Arc<Message>, ctx: &ActorContext) {
match &*message {
Message::Get(get) => self.handle_get(get, ctx),
Message::Put(put) => {
self.handle_put(put, 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());
put.to_string(); let _ = flush.from.send(Message::Put(put));
}
Message::BatchPut(batch) => self.handle_batch_put(batch, ctx),
_ => {}
}
}
fn try_clone_storage(&self) -> Option<Box<dyn Actor>> {
None
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_memory_storage_new() {
let storage = MemoryStorage::new();
assert!(storage.store.read().is_empty());
}
#[tokio::test]
async fn test_memory_storage_default() {
let storage = MemoryStorage::default();
assert!(storage.store.read().is_empty());
}
}