Fix racecondition (save the updated aggregate!), and deadlock (use outer_lock on pub functions).

This commit is contained in:
Tim Bruijnzeels
2019-07-21 15:54:59 -04:00
parent 2e7b6258fe
commit 3c8e9a470c
+32 -15
View File
@@ -63,8 +63,8 @@ pub enum AggregateStoreError {
#[display(fmt = "Event not applicable to aggregate, id or version is off")]
WrongEventForAggregate,
#[display(fmt = "Trying to update outdated aggregate")]
ConcurrentModification,
#[display(fmt = "Trying to update outdated aggregate: {}", _0)]
ConcurrentModification(Handle),
}
impl From<KeyStoreError> for AggregateStoreError {
@@ -76,7 +76,8 @@ pub struct DiskAggregateStore<A: Aggregate> {
store: DiskKeyStore,
cache: RwLock<HashMap<Handle, Arc<A>>>,
use_cache: bool,
listeners: Vec<Arc<EventListener<A>>>
listeners: Vec<Arc<EventListener<A>>>,
outer_lock: RwLock<()>
}
impl<A: Aggregate> DiskAggregateStore<A> {
@@ -85,7 +86,8 @@ impl<A: Aggregate> DiskAggregateStore<A> {
let cache = RwLock::new(HashMap::new());
let use_cache = true;
let listeners = vec![];
Ok(DiskAggregateStore { store, cache, use_cache, listeners })
let lock = RwLock::new(());
Ok(DiskAggregateStore { store, cache, use_cache, listeners, outer_lock: lock })
}
}
@@ -111,10 +113,8 @@ impl<A: Aggregate> DiskAggregateStore<A> {
self.cache.write().unwrap().insert(id.clone(), arc);
}
}
}
impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
fn get_latest(&self, handle: &Handle) -> StoreResult<Arc<A>> {
fn get_latest_no_lock(&self, handle: &Handle) -> StoreResult<Arc<A>> {
debug!("Trying to load aggregate id: {}", handle);
match self.cache_get(handle) {
None => {
@@ -141,11 +141,20 @@ impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
}
}
}
}
impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
fn get_latest(&self, handle: &Handle) -> StoreResult<Arc<A>> {
let _lock = self.outer_lock.read().unwrap();
self.get_latest_no_lock(handle)
}
fn add(
&self,
init: A::InitEvent
) -> StoreResult<Arc<A>> {
let _lock = self.outer_lock.write().unwrap();
self.store.store_event(&init)?;
let handle = init.handle().clone();
@@ -166,12 +175,14 @@ impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
prev: Arc<A>,
events: Vec<A::Event>
) -> StoreResult<Arc<A>> {
let _lock = self.outer_lock.write().unwrap();
// Get the latest arc.
let mut latest = self.get_latest(handle)?;
let mut latest = self.get_latest_no_lock(handle)?;
{
// Verify whether there is a concurrency issue
if prev.version() != latest.version() {
return Err(AggregateStoreError::ConcurrentModification)
return Err(AggregateStoreError::ConcurrentModification(handle.clone()))
}
// forget the previous version
@@ -190,7 +201,7 @@ impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
//
// Also note that we don't need the lock to update the inner arc in the cache. We
// just need it to be in scope until we are done updating.
let _write_lock = self.cache.write().unwrap();
let mut cache = self.cache.write().unwrap();
// There is a possible race condition. We may only have obtained the lock
if self.has_updates(handle, &agg)? {
@@ -211,29 +222,35 @@ impl<A: Aggregate> AggregateStore<A> for DiskAggregateStore<A> {
for event in events {
self.store.store_event(&event)?;
for listener in &self.listeners {
listener.as_ref().listen(agg, &event);
}
agg.apply(event);
agg.apply(event.clone());
if agg.version() % SNAPSHOT_FREQ == 0 {
self.store.store_aggregate(handle, agg)?;
}
for listener in &self.listeners {
listener.as_ref().listen(agg, &event);
}
}
cache.insert(handle.clone(), Arc::new(agg.clone()));
}
Ok(latest)
}
fn has(&self, id: &Handle) -> bool {
let _lock = self.outer_lock.read().unwrap();
self.store.has_aggregate(id)
}
fn list(&self) -> Vec<Handle> {
let _lock = self.outer_lock.read().unwrap();
self.store.aggregates()
}
fn add_listener<L: EventListener<A>>(&mut self, listener: Arc<L>) {
let _lock = self.outer_lock.write().unwrap();
self.listeners.push(listener)
}
}