use distributed::{
AsyncCommitBatch, AsyncGetStream, AsyncStreamWrite, AsyncTransactionalCommit, Entity,
HashMapRepository, StreamIdentity,
};
const AGGREGATE_TYPE: &str = "event_store_test";
fn identity(id: &str) -> StreamIdentity {
StreamIdentity::new(AGGREGATE_TYPE, id).unwrap()
}
async fn get_one(repo: &HashMapRepository, id: &str) -> Option<Entity> {
repo.get_stream(&identity(id)).await.unwrap()
}
async fn commit_one(
repo: &HashMapRepository,
entity: &mut Entity,
) -> Result<(), distributed::RepositoryError> {
let id = entity.id().to_string();
let stream = AsyncStreamWrite::new(identity(&id), entity);
repo.commit_batch_async(AsyncCommitBatch::new(vec![stream]))
.await
}
async fn commit_many(
repo: &HashMapRepository,
entities: &mut [&mut Entity],
) -> Result<(), distributed::RepositoryError> {
let mut streams = Vec::with_capacity(entities.len());
for entity in entities.iter_mut() {
let id = entity.id().to_string();
streams.push(AsyncStreamWrite::new(identity(&id), entity));
}
repo.commit_batch_async(AsyncCommitBatch::new(streams))
.await
}
#[test]
fn digest_adds_events_with_correct_sequences() {
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"data1").unwrap();
entity.digest("Updated", &"data2").unwrap();
entity.digest("Updated", &"data3").unwrap();
assert_eq!(entity.events().len(), 3);
assert_eq!(entity.events()[0].sequence, 1);
assert_eq!(entity.events()[1].sequence, 2);
assert_eq!(entity.events()[2].sequence, 3);
}
#[tokio::test]
async fn multiple_load_modify_commit_cycles_accumulate_all_events() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let mut entity = get_one(&repo, "e1").await.unwrap();
assert_eq!(entity.events().len(), 1);
entity.digest("Updated", &"v2").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let mut entity = get_one(&repo, "e1").await.unwrap();
assert_eq!(entity.events().len(), 2);
entity.digest("Updated", &"v3").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let entity = get_one(&repo, "e1").await.unwrap();
assert_eq!(entity.events().len(), 3);
assert_eq!(entity.events()[0].event_name, "Created");
assert_eq!(entity.events()[1].event_name, "Updated");
assert_eq!(entity.events()[2].event_name, "Updated");
assert_eq!(entity.version(), 3);
}
#[tokio::test]
async fn commit_appends_only_new_events() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let mut entity = get_one(&repo, "e1").await.unwrap();
entity.digest("Updated", &"v2").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let loaded = get_one(&repo, "e1").await.unwrap();
assert_eq!(loaded.events().len(), 2);
assert_eq!(loaded.events()[0].event_name, "Created");
assert_eq!(loaded.events()[1].event_name, "Updated");
}
#[tokio::test]
async fn empty_commit_is_idempotent() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let mut entity = get_one(&repo, "e1").await.unwrap();
assert!(entity.new_events().is_empty());
commit_one(&repo, &mut entity).await.unwrap();
let loaded = get_one(&repo, "e1").await.unwrap();
assert_eq!(loaded.events().len(), 1);
}
#[tokio::test]
async fn events_grow_monotonically() {
let repo = HashMapRepository::new();
for i in 0..5 {
let mut entity = if i == 0 {
Entity::with_id("e1")
} else {
get_one(&repo, "e1").await.unwrap()
};
entity.digest("Event", &format!("v{}", i)).unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let loaded = get_one(&repo, "e1").await.unwrap();
assert_eq!(loaded.events().len(), i + 1);
}
}
#[tokio::test]
async fn concurrent_writes_detected() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let mut reader1 = get_one(&repo, "e1").await.unwrap();
let mut reader2 = get_one(&repo, "e1").await.unwrap();
assert_eq!(reader1.committed_version(), 1);
assert_eq!(reader2.committed_version(), 1);
reader1.digest("UpdatedByR1", &"r1").unwrap();
reader2.digest("UpdatedByR2", &"r2").unwrap();
commit_one(&repo, &mut reader1).await.unwrap();
let err = commit_one(&repo, &mut reader2).await.unwrap_err();
match err {
distributed::RepositoryError::ConcurrentWrite {
id,
expected,
actual,
} => {
assert_eq!(id, format!("{}:e1", AGGREGATE_TYPE));
assert_eq!(expected, 1); assert_eq!(actual, 2); }
other => panic!("expected ConcurrentWrite, got: {:?}", other),
}
}
#[tokio::test]
async fn partial_conflict_rolls_back_entire_commit() {
let repo = HashMapRepository::new();
let mut e1 = Entity::with_id("e1");
e1.digest("Created", &"v1").unwrap();
let mut e2 = Entity::with_id("e2");
e2.digest("Created", &"v1").unwrap();
commit_many(&repo, &mut [&mut e1, &mut e2]).await.unwrap();
let mut e1_a = get_one(&repo, "e1").await.unwrap();
let mut e2_a = get_one(&repo, "e2").await.unwrap();
let mut e2_b = get_one(&repo, "e2").await.unwrap();
e2_b.digest("Conflict", &"b").unwrap();
commit_one(&repo, &mut e2_b).await.unwrap();
e1_a.digest("Update", &"a").unwrap();
e2_a.digest("Update", &"a").unwrap();
let err = commit_many(&repo, &mut [&mut e1_a, &mut e2_a])
.await
.unwrap_err();
match err {
distributed::RepositoryError::ConcurrentWrite { id, .. } => {
assert_eq!(id, format!("{}:e2", AGGREGATE_TYPE));
}
other => panic!("expected ConcurrentWrite, got: {:?}", other),
}
let e1_loaded = get_one(&repo, "e1").await.unwrap();
assert_eq!(e1_loaded.events().len(), 1);
let e2_loaded = get_one(&repo, "e2").await.unwrap();
assert_eq!(e2_loaded.events().len(), 2);
}
#[test]
fn new_entity_has_zero_versions() {
let entity = Entity::with_id("e1");
assert_eq!(entity.committed_version(), 0);
assert_eq!(entity.snapshot_version(), 0);
assert!(entity.new_events().is_empty());
}
#[tokio::test]
async fn load_from_history_sets_committed_version() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
entity.digest("Updated", &"v2").unwrap();
commit_one(&repo, &mut entity).await.unwrap();
let loaded = get_one(&repo, "e1").await.unwrap();
assert_eq!(loaded.committed_version(), 2);
assert_eq!(loaded.snapshot_version(), 0);
assert_eq!(loaded.version(), 2);
assert!(loaded.new_events().is_empty());
}
#[tokio::test]
async fn commit_updates_committed_version() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("Created", &"v1").unwrap();
assert_eq!(entity.committed_version(), 0);
commit_one(&repo, &mut entity).await.unwrap();
assert_eq!(entity.committed_version(), 1);
assert_eq!(entity.snapshot_version(), 0);
assert!(entity.new_events().is_empty());
entity.digest("Updated", &"v2").unwrap();
assert_eq!(entity.new_events().len(), 1);
commit_one(&repo, &mut entity).await.unwrap();
assert_eq!(entity.committed_version(), 2);
assert_eq!(entity.snapshot_version(), 0);
assert!(entity.new_events().is_empty());
}
#[tokio::test]
async fn new_events_returns_only_uncommitted() {
let repo = HashMapRepository::new();
let mut entity = Entity::with_id("e1");
entity.digest("e1", &"a").unwrap();
entity.digest("e2", &"b").unwrap();
assert_eq!(entity.new_events().len(), 2);
commit_one(&repo, &mut entity).await.unwrap();
assert!(entity.new_events().is_empty());
entity.digest("e3", &"c").unwrap();
assert_eq!(entity.new_events().len(), 1);
assert_eq!(entity.new_events()[0].event_name, "e3");
}