mod namespaced;
pub use namespaced::NamespacedStore;
use std::collections::HashMap;
use std::sync::mpsc::{channel, Sender};
use std::sync::{Arc, Mutex, RwLock};
use arora_types::data::{DataError, DataStore, Key, Slot, State, StateChange, Subscription};
use arora_types::value::Value;
type Cell = Arc<RwLock<Option<Value>>>;
#[derive(Default)]
struct Inner {
cells: RwLock<HashMap<Key, Cell>>,
subscribers: Mutex<Vec<Sender<StateChange>>>,
}
impl Inner {
fn notify(&self, change: StateChange) {
if change.is_empty() {
return;
}
let mut subs = self.subscribers.lock().unwrap();
subs.retain(|tx| tx.send(change.clone()).is_ok());
}
}
#[derive(Clone, Default)]
pub struct SimpleDataStore {
inner: Arc<Inner>,
}
impl SimpleDataStore {
pub fn new() -> Self {
Self::default()
}
fn cell(&self, key: &Key) -> Cell {
if let Some(cell) = self.inner.cells.read().unwrap().get(key) {
return cell.clone();
}
self.inner
.cells
.write()
.unwrap()
.entry(key.clone())
.or_insert_with(|| Arc::new(RwLock::new(None)))
.clone()
}
}
impl DataStore for SimpleDataStore {
fn read(&self, keys: &[Key]) -> Vec<Option<Value>> {
let cells = self.inner.cells.read().unwrap();
keys.iter()
.map(|k| cells.get(k).and_then(|c| c.read().unwrap().clone()))
.collect()
}
fn write(&self, changes: StateChange) -> Result<(), DataError> {
let mut effective = StateChange::new();
for (key, value) in &changes.set {
let cell = self.cell(key);
let mut current = cell.write().unwrap();
if current.as_ref() != value.as_ref() {
*current = value.clone();
effective.set.insert(key.clone(), value.clone());
}
}
for key in &changes.unset {
let cell = self.cell(key);
let mut current = cell.write().unwrap();
if current.is_some() {
*current = None;
effective.unset.insert(key.clone());
}
}
if !effective.is_empty() {
self.inner.notify(effective);
}
Ok(())
}
fn snapshot(&self) -> State {
let cells = self.inner.cells.read().unwrap();
let storage = cells
.iter()
.map(|(k, c)| (k.clone(), c.read().unwrap().clone()))
.collect();
State { storage }
}
fn slot(&self, key: &Key) -> Box<dyn Slot> {
Box::new(SimpleSlot {
cell: self.cell(key),
key: key.clone(),
inner: self.inner.clone(),
})
}
fn clone_box(&self) -> Box<dyn DataStore> {
Box::new(self.clone())
}
fn subscribe(&self) -> Subscription {
let (tx, rx) = channel();
let mut subscribers = self.inner.subscribers.lock().unwrap();
let mut initial = StateChange::new();
for (key, value) in self.snapshot().storage {
initial.set.insert(key, value);
}
let _ = tx.send(initial);
subscribers.push(tx);
Subscription::new(rx)
}
}
struct SimpleSlot {
cell: Cell,
key: Key,
inner: Arc<Inner>,
}
impl Slot for SimpleSlot {
fn get(&self) -> Option<Value> {
self.cell.read().unwrap().clone()
}
fn set(&self, value: Option<Value>) -> Result<(), DataError> {
{
let mut current = self.cell.write().unwrap();
if *current == value {
return Ok(());
}
*current = value.clone();
}
self.inner.notify(StateChange {
set: HashMap::from([(self.key.clone(), value)]),
unset: Default::default(),
});
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn write_then_read() {
let store = SimpleDataStore::new();
store
.write(StateChange::set("a/b", Value::Boolean(true)))
.unwrap();
assert_eq!(
store.read(&[Key::from("a/b")]),
vec![Some(Value::Boolean(true))]
);
assert_eq!(store.read(&[Key::from("missing")]), vec![None]);
}
#[test]
fn slot_and_store_coincide() {
let store = SimpleDataStore::new();
let slot = store.slot(&Key::from("x"));
slot.set(Some(Value::Boolean(true))).unwrap();
assert_eq!(
store.read(&[Key::from("x")]),
vec![Some(Value::Boolean(true))]
);
store
.write(StateChange::set("x", Value::Boolean(false)))
.unwrap();
assert_eq!(slot.get(), Some(Value::Boolean(false)));
}
#[test]
fn subscribe_delivers_changes_to_all() {
let store = SimpleDataStore::new();
let s1 = store.subscribe();
let s2 = store.subscribe();
s1.try_recv().expect("s1 opening state");
s2.try_recv().expect("s2 opening state");
store
.write(StateChange::set("k", Value::Boolean(true)))
.unwrap();
assert!(s1.try_recv().expect("s1 change").contains(&Key::from("k")));
assert!(s2.try_recv().expect("s2 change").contains(&Key::from("k")));
}
#[test]
fn subscribe_opens_on_the_current_state() {
let store = SimpleDataStore::new();
store
.write(StateChange::set("already", Value::Boolean(true)))
.unwrap();
let sub = store.subscribe();
let opening = sub.try_recv().expect("opening state");
assert!(opening.contains(&Key::from("already")));
store
.write(StateChange::set("later", Value::Boolean(false)))
.unwrap();
assert!(sub
.try_recv()
.expect("change")
.contains(&Key::from("later")));
}
#[test]
fn slot_set_notifies_subscribers() {
let store = SimpleDataStore::new();
let sub = store.subscribe();
sub.try_recv().expect("opening state");
store
.slot(&Key::from("y"))
.set(Some(Value::Boolean(true)))
.unwrap();
assert!(sub.try_recv().expect("change").contains(&Key::from("y")));
}
#[test]
fn snapshot_returns_all() {
let store = SimpleDataStore::new();
store
.write(StateChange::set("a", Value::Boolean(true)))
.unwrap();
store
.write(StateChange::set("b", Value::Boolean(false)))
.unwrap();
assert_eq!(store.snapshot().storage.len(), 2);
}
#[test]
fn clones_share_storage() {
let store = SimpleDataStore::new();
let other = store.clone();
store
.write(StateChange::set("shared", Value::Boolean(true)))
.unwrap();
assert_eq!(
other.read(&[Key::from("shared")]),
vec![Some(Value::Boolean(true))]
);
}
}