use std::{
collections::{HashMap, HashSet},
marker::PhantomData,
sync::Mutex,
time::SystemTime,
};
use soaprs_core::{BoxFuture, Command, MessageId, SoapError, SoapResult};
use soaprs_cqrs::{
AppendResult, ClaimedOutboxRecord, ClaimedSagaAction, DeadLetter, DeliveryClaimId,
EncodedEvent, EventStore, ExpectedSagaVersion, ExpectedVersion, GlobalPosition, InboxClaim,
InboxStore, OutboxRecord, OutboxStore, ProjectionApplyOutcome, ProjectionCheckpointStore,
ProjectionId, RecordedEvent, SagaAction, SagaActionId, SagaActionRecord, SagaActionStore,
SagaCommitResult, SagaId, SagaStore, SagaVersion, Snapshot, SnapshotStore, StoredEvent,
StoredEventStore, StoredSaga, StreamId, StreamVersion, TransactionalEventOutbox,
TransactionalProjection, TransactionalSagaOutbox,
};
use soaprs_events::{DomainEvent, IntegrationEvent};
#[derive(Debug)]
struct EventStoreState<E> {
streams: HashMap<StreamId, Vec<RecordedEvent<E>>>,
global: Vec<RecordedEvent<E>>,
next_global_position: u64,
}
#[derive(Debug)]
struct StoredEventState<P> {
streams: HashMap<StreamId, Vec<StoredEvent<P>>>,
global: Vec<StoredEvent<P>>,
next_global_position: u64,
}
impl<P> Default for StoredEventState<P> {
fn default() -> Self {
Self {
streams: HashMap::new(),
global: Vec::new(),
next_global_position: 1,
}
}
}
impl<E> Default for EventStoreState<E> {
fn default() -> Self {
Self {
streams: HashMap::new(),
global: Vec::new(),
next_global_position: 1,
}
}
}
#[derive(Debug)]
pub struct MemoryEventStore<E> {
state: Mutex<EventStoreState<E>>,
}
impl<E> MemoryEventStore<E> {
pub fn new() -> Self {
Self {
state: Mutex::new(EventStoreState::default()),
}
}
}
impl<E> Default for MemoryEventStore<E> {
fn default() -> Self {
Self::new()
}
}
impl<E> EventStore<E> for MemoryEventStore<E>
where
E: DomainEvent + Clone,
{
fn append<'a>(
&'a self,
stream_id: &'a StreamId,
expected: ExpectedVersion,
events: Vec<soaprs_events::EventEnvelope<E>>,
) -> BoxFuture<'a, SoapResult<AppendResult>> {
Box::pin(async move {
if events.is_empty() {
return Err(SoapError::validation(
"event store append requires at least one event",
));
}
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("in-memory event store lock poisoned"))?;
let current = state
.streams
.get(stream_id)
.and_then(|stream| stream.last())
.map(|event| event.stream_version);
validate_expected_version(expected, current)?;
let mut stream_version = current
.map_or(Some(StreamVersion::FIRST), StreamVersion::checked_next)
.ok_or_else(|| SoapError::infrastructure("event stream version overflow"))?;
let event_count = events.len();
let mut records = Vec::with_capacity(event_count);
for (index, envelope) in events.into_iter().enumerate() {
let global_position = GlobalPosition::new(state.next_global_position);
state.next_global_position = state
.next_global_position
.checked_add(1)
.ok_or_else(|| SoapError::infrastructure("global event position overflow"))?;
records.push(RecordedEvent {
stream_id: stream_id.clone(),
stream_version,
global_position,
envelope,
});
if index + 1 < event_count {
stream_version = stream_version.checked_next().ok_or_else(|| {
SoapError::infrastructure("event stream version overflow")
})?;
}
}
let Some(last) = records.last() else {
return Err(SoapError::infrastructure(
"event append lost its non-empty record set",
));
};
let result = AppendResult {
stream_version: last.stream_version,
global_position: last.global_position,
};
state
.streams
.entry(stream_id.clone())
.or_default()
.extend(records.iter().cloned());
state.global.extend(records);
Ok(result)
})
}
fn load<'a>(
&'a self,
stream_id: &'a StreamId,
after: Option<StreamVersion>,
) -> BoxFuture<'a, SoapResult<Vec<RecordedEvent<E>>>> {
Box::pin(async move {
let state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("in-memory event store lock poisoned"))?;
Ok(state
.streams
.get(stream_id)
.into_iter()
.flatten()
.filter(|event| after.is_none_or(|version| event.stream_version > version))
.cloned()
.collect())
})
}
fn read_all(
&self,
after: Option<GlobalPosition>,
limit: usize,
) -> BoxFuture<'_, SoapResult<Vec<RecordedEvent<E>>>> {
Box::pin(async move {
if limit == 0 {
return Err(SoapError::validation(
"global event read limit must be greater than zero",
));
}
let state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("in-memory event store lock poisoned"))?;
Ok(state
.global
.iter()
.filter(|event| after.is_none_or(|position| event.global_position > position))
.take(limit)
.cloned()
.collect())
})
}
}
#[derive(Debug)]
pub struct MemoryStoredEventStore<P> {
state: Mutex<StoredEventState<P>>,
}
impl<P> MemoryStoredEventStore<P> {
pub fn new() -> Self {
Self {
state: Mutex::new(StoredEventState::default()),
}
}
}
impl<P> Default for MemoryStoredEventStore<P> {
fn default() -> Self {
Self::new()
}
}
impl<P> StoredEventStore<P> for MemoryStoredEventStore<P>
where
P: Clone + Send,
{
fn append<'a>(
&'a self,
stream_id: &'a StreamId,
expected: ExpectedVersion,
events: Vec<EncodedEvent<P>>,
) -> BoxFuture<'a, SoapResult<AppendResult>> {
Box::pin(async move {
if events.is_empty() {
return Err(SoapError::validation(
"stored event append requires at least one event",
));
}
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("stored event store lock poisoned"))?;
let current = state
.streams
.get(stream_id)
.and_then(|stream| stream.last())
.map(|event| event.stream_version);
validate_expected_version(expected, current)?;
let mut stream_version = current
.map_or(Some(StreamVersion::FIRST), StreamVersion::checked_next)
.ok_or_else(|| SoapError::infrastructure("stored event stream version overflow"))?;
let event_count = events.len();
let mut records = Vec::with_capacity(event_count);
for (index, encoded) in events.into_iter().enumerate() {
let global_position = GlobalPosition::new(state.next_global_position);
state.next_global_position = state
.next_global_position
.checked_add(1)
.ok_or_else(|| SoapError::infrastructure("global event position overflow"))?;
records.push(StoredEvent {
stream_id: stream_id.clone(),
stream_version,
global_position,
encoded,
});
if index + 1 < event_count {
stream_version = stream_version.checked_next().ok_or_else(|| {
SoapError::infrastructure("stored event stream version overflow")
})?;
}
}
let Some(last) = records.last() else {
return Err(SoapError::infrastructure(
"stored event append lost its non-empty record set",
));
};
let result = AppendResult {
stream_version: last.stream_version,
global_position: last.global_position,
};
state
.streams
.entry(stream_id.clone())
.or_default()
.extend(records.iter().cloned());
state.global.extend(records);
Ok(result)
})
}
fn load<'a>(
&'a self,
stream_id: &'a StreamId,
after: Option<StreamVersion>,
) -> BoxFuture<'a, SoapResult<Vec<StoredEvent<P>>>> {
Box::pin(async move {
let state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("stored event store lock poisoned"))?;
Ok(state
.streams
.get(stream_id)
.into_iter()
.flatten()
.filter(|event| after.is_none_or(|version| event.stream_version > version))
.cloned()
.collect())
})
}
fn read_all(
&self,
after: Option<GlobalPosition>,
limit: usize,
) -> BoxFuture<'_, SoapResult<Vec<StoredEvent<P>>>> {
Box::pin(async move {
if limit == 0 {
return Err(SoapError::validation(
"stored global event read limit must be greater than zero",
));
}
let state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("stored event store lock poisoned"))?;
Ok(state
.global
.iter()
.filter(|event| after.is_none_or(|position| event.global_position > position))
.take(limit)
.cloned()
.collect())
})
}
}
fn validate_expected_version(
expected: ExpectedVersion,
current: Option<StreamVersion>,
) -> SoapResult<()> {
let matches = match expected {
ExpectedVersion::Any => true,
ExpectedVersion::NoStream => current.is_none(),
ExpectedVersion::Exact(version) => current == Some(version),
};
if matches {
Ok(())
} else {
Err(SoapError::conflict("event stream version conflict"))
}
}
#[derive(Debug, Default)]
pub struct MemoryProjectionCheckpointStore {
checkpoints: Mutex<HashMap<ProjectionId, GlobalPosition>>,
}
impl MemoryProjectionCheckpointStore {
pub fn new() -> Self {
Self::default()
}
}
impl ProjectionCheckpointStore for MemoryProjectionCheckpointStore {
fn load<'a>(
&'a self,
projection: &'a ProjectionId,
) -> BoxFuture<'a, SoapResult<Option<GlobalPosition>>> {
Box::pin(async move {
Ok(self
.checkpoints
.lock()
.map_err(|_| SoapError::infrastructure("checkpoint store lock poisoned"))?
.get(projection)
.copied())
})
}
fn store<'a>(
&'a self,
projection: &'a ProjectionId,
position: GlobalPosition,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut checkpoints = self
.checkpoints
.lock()
.map_err(|_| SoapError::infrastructure("checkpoint store lock poisoned"))?;
if checkpoints
.get(projection)
.is_some_and(|current| position < *current)
{
return Err(SoapError::conflict(
"projection checkpoint cannot move backwards",
));
}
checkpoints.insert(projection.clone(), position);
Ok(())
})
}
}
#[derive(Debug, Clone)]
struct TransactionalProjectionState<S> {
read_model: S,
checkpoint: Option<GlobalPosition>,
}
pub struct MemoryTransactionalProjection<E, S, F> {
id: ProjectionId,
state: Mutex<TransactionalProjectionState<S>>,
apply: F,
event: PhantomData<fn(E)>,
}
impl<E, S, F> MemoryTransactionalProjection<E, S, F> {
pub fn new(id: ProjectionId, read_model: S, apply: F) -> Self {
Self {
id,
state: Mutex::new(TransactionalProjectionState {
read_model,
checkpoint: None,
}),
apply,
event: PhantomData,
}
}
}
impl<E, S, F> MemoryTransactionalProjection<E, S, F>
where
S: Clone,
{
pub fn snapshot(&self) -> SoapResult<(S, Option<GlobalPosition>)> {
let state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("transactional projection lock poisoned"))?;
Ok((state.read_model.clone(), state.checkpoint))
}
}
impl<E, S, F> std::fmt::Debug for MemoryTransactionalProjection<E, S, F> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("MemoryTransactionalProjection")
.field("id", &self.id)
.finish_non_exhaustive()
}
}
impl<E, S, F> TransactionalProjection<E> for MemoryTransactionalProjection<E, S, F>
where
E: DomainEvent,
S: Clone + Send,
F: Fn(&mut S, &RecordedEvent<E>) -> SoapResult<()> + Send + Sync,
{
fn id(&self) -> &ProjectionId {
&self.id
}
fn checkpoint(&self) -> BoxFuture<'_, SoapResult<Option<GlobalPosition>>> {
Box::pin(async move {
Ok(self
.state
.lock()
.map_err(|_| SoapError::infrastructure("transactional projection lock poisoned"))?
.checkpoint)
})
}
fn project_and_checkpoint<'a>(
&'a self,
event: &'a RecordedEvent<E>,
) -> BoxFuture<'a, SoapResult<ProjectionApplyOutcome>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("transactional projection lock poisoned"))?;
if state
.checkpoint
.is_some_and(|checkpoint| event.global_position <= checkpoint)
{
return Ok(ProjectionApplyOutcome::AlreadyApplied);
}
let mut candidate = state.read_model.clone();
(self.apply)(&mut candidate, event)?;
state.read_model = candidate;
state.checkpoint = Some(event.global_position);
Ok(ProjectionApplyOutcome::Applied)
})
}
}
#[derive(Debug, Default)]
pub struct MemorySnapshotStore<S> {
snapshots: Mutex<HashMap<StreamId, Vec<Snapshot<S>>>>,
}
impl<S> MemorySnapshotStore<S> {
pub fn new() -> Self {
Self {
snapshots: Mutex::new(HashMap::new()),
}
}
}
impl<S> SnapshotStore<S> for MemorySnapshotStore<S>
where
S: Clone + Send,
{
fn load_latest<'a>(
&'a self,
stream_id: &'a StreamId,
) -> BoxFuture<'a, SoapResult<Option<Snapshot<S>>>> {
Box::pin(async move {
Ok(self
.snapshots
.lock()
.map_err(|_| SoapError::infrastructure("snapshot store lock poisoned"))?
.get(stream_id)
.and_then(|items| items.last())
.cloned())
})
}
fn save(&self, snapshot: Snapshot<S>) -> BoxFuture<'_, SoapResult<()>> {
Box::pin(async move {
let mut snapshots = self
.snapshots
.lock()
.map_err(|_| SoapError::infrastructure("snapshot store lock poisoned"))?;
let stream = snapshots.entry(snapshot.stream_id.clone()).or_default();
if stream
.last()
.is_some_and(|current| snapshot.stream_version < current.stream_version)
{
return Err(SoapError::conflict(
"snapshot cannot replace a newer stream version",
));
}
if stream
.last()
.is_some_and(|current| snapshot.stream_version == current.stream_version)
{
stream.pop();
}
stream.push(snapshot);
Ok(())
})
}
fn prune<'a>(
&'a self,
stream_id: &'a StreamId,
retain_from: StreamVersion,
) -> BoxFuture<'a, SoapResult<u64>> {
Box::pin(async move {
let mut snapshots = self
.snapshots
.lock()
.map_err(|_| SoapError::infrastructure("snapshot store lock poisoned"))?;
let Some(stream) = snapshots.get_mut(stream_id) else {
return Ok(0);
};
let before = stream.len();
stream.retain(|snapshot| snapshot.stream_version >= retain_from);
Ok(before.saturating_sub(stream.len()) as u64)
})
}
}
#[derive(Debug, Default)]
pub struct MemorySagaStore<S> {
sagas: Mutex<HashMap<SagaId, StoredSaga<S>>>,
}
impl<S> MemorySagaStore<S> {
pub fn new() -> Self {
Self {
sagas: Mutex::new(HashMap::new()),
}
}
}
impl<S> SagaStore<S> for MemorySagaStore<S>
where
S: Clone + Send,
{
fn load<'a>(&'a self, id: &'a SagaId) -> BoxFuture<'a, SoapResult<Option<StoredSaga<S>>>> {
Box::pin(async move {
Ok(self
.sagas
.lock()
.map_err(|_| SoapError::infrastructure("saga store lock poisoned"))?
.get(id)
.cloned())
})
}
fn save<'a>(
&'a self,
id: &'a SagaId,
expected: ExpectedSagaVersion,
saga: S,
) -> BoxFuture<'a, SoapResult<SagaVersion>> {
Box::pin(async move {
let mut sagas = self
.sagas
.lock()
.map_err(|_| SoapError::infrastructure("saga store lock poisoned"))?;
let current = sagas.get(id).map(|stored| stored.version);
let version = next_saga_version(expected, current)?;
sagas.insert(
id.clone(),
StoredSaga {
version,
state: saga,
},
);
Ok(version)
})
}
}
fn next_saga_version(
expected: ExpectedSagaVersion,
current: Option<SagaVersion>,
) -> SoapResult<SagaVersion> {
let matches = match expected {
ExpectedSagaVersion::Any => true,
ExpectedSagaVersion::NoSaga => current.is_none(),
ExpectedSagaVersion::Exact(version) => current == Some(version),
};
if !matches {
return Err(SoapError::conflict("saga version conflict"));
}
current
.map_or(Some(SagaVersion::FIRST), SagaVersion::checked_next)
.ok_or_else(|| SoapError::infrastructure("saga version overflow"))
}
#[derive(Debug, Clone)]
struct MemorySagaActionEntry<C>
where
C: Command<Output = ()>,
{
record: SagaActionRecord<C>,
claim: Option<(DeliveryClaimId, SystemTime)>,
}
#[derive(Debug)]
struct SagaActionState<C>
where
C: Command<Output = ()>,
{
order: Vec<SagaActionId>,
entries: HashMap<SagaActionId, MemorySagaActionEntry<C>>,
}
impl<C> Default for SagaActionState<C>
where
C: Command<Output = ()>,
{
fn default() -> Self {
Self {
order: Vec::new(),
entries: HashMap::new(),
}
}
}
#[derive(Debug)]
pub struct MemorySagaActionStore<C>
where
C: Command<Output = ()>,
{
state: Mutex<SagaActionState<C>>,
}
impl<C> MemorySagaActionStore<C>
where
C: Command<Output = ()>,
{
pub fn new() -> Self {
Self {
state: Mutex::new(SagaActionState::default()),
}
}
}
impl<C> Default for MemorySagaActionStore<C>
where
C: Command<Output = ()>,
{
fn default() -> Self {
Self::new()
}
}
impl<C> SagaActionStore<C> for MemorySagaActionStore<C>
where
C: Command<Output = ()> + Clone,
{
fn claim_pending(
&self,
claim_id: DeliveryClaimId,
now: SystemTime,
lease_until: SystemTime,
limit: usize,
) -> BoxFuture<'_, SoapResult<Vec<ClaimedSagaAction<C>>>> {
Box::pin(async move {
if limit == 0 {
return Err(SoapError::validation(
"saga action claim limit must be greater than zero",
));
}
if lease_until <= now {
return Err(SoapError::validation(
"saga action claim lease must end after the claim time",
));
}
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("saga action store lock poisoned"))?;
let order = state.order.clone();
let mut claimed = Vec::new();
for id in order {
if claimed.len() == limit {
break;
}
let Some(entry) = state.entries.get_mut(&id) else {
continue;
};
if entry
.claim
.as_ref()
.is_none_or(|(_, current_lease)| *current_lease <= now)
&& entry.record.completed_at.is_none()
&& entry.record.dead_letter.is_none()
&& entry.record.available_at <= now
{
entry.claim = Some((claim_id.clone(), lease_until));
claimed.push(ClaimedSagaAction {
record: entry.record.clone(),
claim_id: claim_id.clone(),
lease_until,
});
}
}
Ok(claimed)
})
}
fn mark_completed<'a>(
&'a self,
id: &'a SagaActionId,
claim_id: &'a DeliveryClaimId,
completed_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("saga action store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("saga action"));
};
validate_saga_action_claim(entry.claim.as_ref(), claim_id, completed_at)?;
entry.record.completed_at = Some(completed_at);
entry.claim = None;
Ok(())
})
}
fn mark_failed<'a>(
&'a self,
id: &'a SagaActionId,
claim_id: &'a DeliveryClaimId,
safe_error: &'a str,
failed_at: SystemTime,
available_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("saga action store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("saga action"));
};
validate_saga_action_claim(entry.claim.as_ref(), claim_id, failed_at)?;
entry.record.attempts = entry.record.attempts.saturating_add(1);
entry.record.available_at = available_at;
entry.record.last_error = Some(safe_error.to_owned());
entry.claim = None;
Ok(())
})
}
fn mark_dead_lettered<'a>(
&'a self,
id: &'a SagaActionId,
claim_id: &'a DeliveryClaimId,
safe_error: &'a str,
failed_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("saga action store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("saga action"));
};
validate_saga_action_claim(entry.claim.as_ref(), claim_id, failed_at)?;
entry.record.attempts = entry.record.attempts.saturating_add(1);
entry.record.last_error = Some(safe_error.to_owned());
entry.record.dead_letter = Some(DeadLetter {
failed_at,
safe_error: safe_error.to_owned(),
});
entry.claim = None;
Ok(())
})
}
}
fn validate_saga_action_claim(
claim: Option<&(DeliveryClaimId, SystemTime)>,
claim_id: &DeliveryClaimId,
outcome_at: SystemTime,
) -> SoapResult<()> {
let Some((current, lease_until)) = claim else {
return Err(SoapError::conflict("saga action has no active claim"));
};
if current != claim_id {
return Err(SoapError::conflict(
"saga action claim is not owned by this worker",
));
}
if *lease_until <= outcome_at {
return Err(SoapError::conflict("saga action claim lease has expired"));
}
Ok(())
}
#[derive(Debug)]
pub struct MemorySagaOutbox<S, C>
where
C: Command<Output = ()>,
{
sagas: MemorySagaStore<S>,
actions: MemorySagaActionStore<C>,
}
impl<S, C> MemorySagaOutbox<S, C>
where
C: Command<Output = ()>,
{
pub fn new() -> Self {
Self {
sagas: MemorySagaStore::new(),
actions: MemorySagaActionStore::new(),
}
}
pub const fn saga_store(&self) -> &MemorySagaStore<S> {
&self.sagas
}
pub const fn action_store(&self) -> &MemorySagaActionStore<C> {
&self.actions
}
}
impl<S, C> Default for MemorySagaOutbox<S, C>
where
C: Command<Output = ()>,
{
fn default() -> Self {
Self::new()
}
}
impl<S, C> TransactionalSagaOutbox<S, C> for MemorySagaOutbox<S, C>
where
S: Send,
C: Command<Output = ()>,
{
fn commit_transition<'a>(
&'a self,
saga_id: &'a SagaId,
expected: ExpectedSagaVersion,
saga: S,
transition: soaprs_core::MessageMetadata,
actions: Vec<SagaAction<C>>,
) -> BoxFuture<'a, SoapResult<SagaCommitResult>> {
Box::pin(async move {
let created_at = transition.created_at();
let mut records = Vec::with_capacity(actions.len());
for (index, action) in actions.into_iter().enumerate() {
let sequence = u32::try_from(index).map_err(|_| {
SoapError::validation("saga transition contains too many actions")
})?;
let available_at = match &action {
SagaAction::Schedule { delay, .. } => {
created_at.checked_add(*delay).ok_or_else(|| {
SoapError::validation("scheduled saga action time overflow")
})?
}
SagaAction::Dispatch(_)
| SagaAction::Compensate(_)
| SagaAction::Complete
| SagaAction::Fail(_) => created_at,
};
records.push(SagaActionRecord {
id: SagaActionId::new(transition.id().clone(), sequence),
saga_id: saga_id.clone(),
action,
attempts: 0,
available_at,
completed_at: None,
dead_letter: None,
last_error: None,
transition: transition.clone(),
});
}
let mut sagas = self
.sagas
.sagas
.lock()
.map_err(|_| SoapError::infrastructure("saga store lock poisoned"))?;
let mut action_state = self
.actions
.state
.lock()
.map_err(|_| SoapError::infrastructure("saga action store lock poisoned"))?;
let current = sagas.get(saga_id).map(|stored| stored.version);
let version = next_saga_version(expected, current)?;
if records
.iter()
.any(|record| action_state.entries.contains_key(&record.id))
{
return Err(SoapError::conflict(
"saga transition action identity already exists",
));
}
let queued_actions = records.len();
for record in records {
action_state.order.push(record.id.clone());
action_state.entries.insert(
record.id.clone(),
MemorySagaActionEntry {
record,
claim: None,
},
);
}
sagas.insert(
saga_id.clone(),
StoredSaga {
version,
state: saga,
},
);
Ok(SagaCommitResult {
version,
queued_actions,
})
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum InboxState {
Claimed {
claim_id: DeliveryClaimId,
lease_until: SystemTime,
},
Completed,
}
#[derive(Debug, Default)]
pub struct MemoryInboxStore {
messages: Mutex<HashMap<(String, MessageId), InboxState>>,
}
impl MemoryInboxStore {
pub fn new() -> Self {
Self::default()
}
}
impl InboxStore for MemoryInboxStore {
fn claim<'a>(
&'a self,
consumer: &'a str,
id: &'a MessageId,
claim_id: DeliveryClaimId,
received_at: SystemTime,
lease_until: SystemTime,
) -> BoxFuture<'a, SoapResult<InboxClaim>> {
Box::pin(async move {
let mut messages = self
.messages
.lock()
.map_err(|_| SoapError::infrastructure("inbox store lock poisoned"))?;
if lease_until <= received_at {
return Err(SoapError::validation(
"inbox claim lease must end after the received time",
));
}
let key = (consumer.to_owned(), id.clone());
match messages.entry(key) {
std::collections::hash_map::Entry::Vacant(entry) => {
entry.insert(InboxState::Claimed {
claim_id,
lease_until,
});
Ok(InboxClaim::Acquired)
}
std::collections::hash_map::Entry::Occupied(mut entry) => match entry.get() {
InboxState::Claimed {
lease_until: current_lease,
..
} if *current_lease <= received_at => {
entry.insert(InboxState::Claimed {
claim_id,
lease_until,
});
Ok(InboxClaim::Acquired)
}
InboxState::Claimed { .. } => Ok(InboxClaim::InProgress),
InboxState::Completed => Ok(InboxClaim::Completed),
},
}
})
}
fn complete<'a>(
&'a self,
consumer: &'a str,
id: &'a MessageId,
claim_id: &'a DeliveryClaimId,
completed_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut messages = self
.messages
.lock()
.map_err(|_| SoapError::infrastructure("inbox store lock poisoned"))?;
let key = (consumer.to_owned(), id.clone());
let Some(state) = messages.get_mut(&key) else {
return Err(SoapError::not_found("inbox message claim"));
};
if !matches!(state, InboxState::Claimed { claim_id: current, .. } if current == claim_id)
{
return Err(SoapError::conflict(
"inbox message claim is not owned by this worker",
));
}
if matches!(state, InboxState::Claimed { lease_until, .. } if *lease_until <= completed_at)
{
return Err(SoapError::conflict("inbox message claim lease has expired"));
}
*state = InboxState::Completed;
Ok(())
})
}
fn release<'a>(
&'a self,
consumer: &'a str,
id: &'a MessageId,
claim_id: &'a DeliveryClaimId,
released_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut messages = self
.messages
.lock()
.map_err(|_| SoapError::infrastructure("inbox store lock poisoned"))?;
let key = (consumer.to_owned(), id.clone());
match messages.get(&key) {
Some(InboxState::Claimed {
claim_id: current, ..
}) if current == claim_id => {
if matches!(messages.get(&key), Some(InboxState::Claimed { lease_until, .. }) if *lease_until <= released_at)
{
return Err(SoapError::conflict("inbox message claim lease has expired"));
}
}
Some(InboxState::Claimed { .. } | InboxState::Completed) => {
return Err(SoapError::conflict(
"inbox message claim is not owned by this worker",
));
}
None => return Err(SoapError::not_found("inbox message claim")),
}
messages.remove(&key);
Ok(())
})
}
}
#[derive(Debug, Clone)]
struct MemoryOutboxEntry<M> {
record: OutboxRecord<M>,
claim: Option<(DeliveryClaimId, SystemTime)>,
}
#[derive(Debug)]
struct OutboxState<M> {
order: Vec<MessageId>,
entries: HashMap<MessageId, MemoryOutboxEntry<M>>,
}
impl<M> Default for OutboxState<M> {
fn default() -> Self {
Self {
order: Vec::new(),
entries: HashMap::new(),
}
}
}
#[derive(Debug)]
pub struct MemoryOutboxStore<M> {
state: Mutex<OutboxState<M>>,
}
impl<M> MemoryOutboxStore<M> {
pub fn new() -> Self {
Self {
state: Mutex::new(OutboxState::default()),
}
}
}
impl<M> Default for MemoryOutboxStore<M> {
fn default() -> Self {
Self::new()
}
}
impl<M> OutboxStore<M> for MemoryOutboxStore<M>
where
M: IntegrationEvent + Clone,
{
fn enqueue(&self, messages: Vec<OutboxRecord<M>>) -> BoxFuture<'_, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
for message in &messages {
if state.entries.contains_key(message.envelope.metadata.id()) {
return Err(SoapError::conflict("duplicate outbox message identity"));
}
}
for message in messages {
let id = message.envelope.metadata.id().clone();
state.order.push(id.clone());
state.entries.insert(
id,
MemoryOutboxEntry {
record: message,
claim: None,
},
);
}
Ok(())
})
}
fn claim_pending(
&self,
claim_id: DeliveryClaimId,
now: SystemTime,
lease_until: SystemTime,
limit: usize,
) -> BoxFuture<'_, SoapResult<Vec<ClaimedOutboxRecord<M>>>> {
Box::pin(async move {
if limit == 0 {
return Err(SoapError::validation(
"outbox claim limit must be greater than zero",
));
}
if lease_until <= now {
return Err(SoapError::validation(
"outbox claim lease must end after the claim time",
));
}
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
let order = state.order.clone();
let mut claimed = Vec::new();
for id in order {
if claimed.len() == limit {
break;
}
let Some(entry) = state.entries.get_mut(&id) else {
continue;
};
if entry
.claim
.as_ref()
.is_none_or(|(_, current_lease)| *current_lease <= now)
&& entry.record.delivered_at.is_none()
&& entry.record.dead_letter.is_none()
&& entry.record.available_at <= now
{
entry.claim = Some((claim_id.clone(), lease_until));
claimed.push(ClaimedOutboxRecord {
record: entry.record.clone(),
claim_id: claim_id.clone(),
lease_until,
});
}
}
Ok(claimed)
})
}
fn mark_delivered<'a>(
&'a self,
id: &'a MessageId,
claim_id: &'a DeliveryClaimId,
delivered_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("outbox message"));
};
if !matches!(&entry.claim, Some((current, _)) if current == claim_id) {
return Err(SoapError::conflict(
"outbox message claim is not owned by this worker",
));
}
if matches!(&entry.claim, Some((_, lease_until)) if *lease_until <= delivered_at) {
return Err(SoapError::conflict(
"outbox message claim lease has expired",
));
}
entry.record.delivered_at = Some(delivered_at);
entry.claim = None;
Ok(())
})
}
fn mark_failed<'a>(
&'a self,
id: &'a MessageId,
claim_id: &'a DeliveryClaimId,
safe_error: &'a str,
failed_at: SystemTime,
available_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("outbox message"));
};
if !matches!(&entry.claim, Some((current, _)) if current == claim_id) {
return Err(SoapError::conflict(
"outbox message claim is not owned by this worker",
));
}
if matches!(&entry.claim, Some((_, lease_until)) if *lease_until <= failed_at) {
return Err(SoapError::conflict(
"outbox message claim lease has expired",
));
}
entry.record.attempts = entry.record.attempts.saturating_add(1);
entry.record.available_at = available_at;
entry.record.last_error = Some(safe_error.to_owned());
entry.claim = None;
Ok(())
})
}
fn mark_dead_lettered<'a>(
&'a self,
id: &'a MessageId,
claim_id: &'a DeliveryClaimId,
safe_error: &'a str,
failed_at: SystemTime,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut state = self
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
let Some(entry) = state.entries.get_mut(id) else {
return Err(SoapError::not_found("outbox message"));
};
if !matches!(&entry.claim, Some((current, _)) if current == claim_id) {
return Err(SoapError::conflict(
"outbox message claim is not owned by this worker",
));
}
if matches!(&entry.claim, Some((_, lease_until)) if *lease_until <= failed_at) {
return Err(SoapError::conflict(
"outbox message claim lease has expired",
));
}
entry.record.attempts = entry.record.attempts.saturating_add(1);
entry.record.last_error = Some(safe_error.to_owned());
entry.record.dead_letter = Some(DeadLetter {
failed_at,
safe_error: safe_error.to_owned(),
});
entry.claim = None;
Ok(())
})
}
}
#[derive(Debug)]
pub struct MemoryEventOutbox<E, M> {
events: MemoryEventStore<E>,
outbox: MemoryOutboxStore<M>,
}
impl<E, M> MemoryEventOutbox<E, M> {
pub fn new() -> Self {
Self {
events: MemoryEventStore::new(),
outbox: MemoryOutboxStore::new(),
}
}
pub const fn event_store(&self) -> &MemoryEventStore<E> {
&self.events
}
pub const fn outbox(&self) -> &MemoryOutboxStore<M> {
&self.outbox
}
}
impl<E, M> Default for MemoryEventOutbox<E, M> {
fn default() -> Self {
Self::new()
}
}
impl<E, M> TransactionalEventOutbox<E, M> for MemoryEventOutbox<E, M>
where
E: DomainEvent + Clone,
M: IntegrationEvent,
{
fn append_with_outbox<'a>(
&'a self,
stream_id: &'a StreamId,
expected: ExpectedVersion,
events: Vec<soaprs_events::EventEnvelope<E>>,
messages: Vec<OutboxRecord<M>>,
) -> BoxFuture<'a, SoapResult<AppendResult>> {
Box::pin(async move {
if events.is_empty() {
return Err(SoapError::validation(
"event store append requires at least one event",
));
}
let mut event_state =
self.events.state.lock().map_err(|_| {
SoapError::infrastructure("in-memory event store lock poisoned")
})?;
let mut outbox_state = self
.outbox
.state
.lock()
.map_err(|_| SoapError::infrastructure("outbox store lock poisoned"))?;
let current = event_state
.streams
.get(stream_id)
.and_then(|stream| stream.last())
.map(|event| event.stream_version);
validate_expected_version(expected, current)?;
let mut new_ids = HashSet::new();
for message in &messages {
let id = message.envelope.metadata.id();
if outbox_state.entries.contains_key(id) || !new_ids.insert(id.clone()) {
return Err(SoapError::conflict("duplicate outbox message identity"));
}
}
let mut next_global_position = event_state.next_global_position;
let mut stream_version = current
.map_or(Some(StreamVersion::FIRST), StreamVersion::checked_next)
.ok_or_else(|| SoapError::infrastructure("event stream version overflow"))?;
let event_count = events.len();
let mut records = Vec::with_capacity(event_count);
for (index, envelope) in events.into_iter().enumerate() {
let global_position = GlobalPosition::new(next_global_position);
next_global_position = next_global_position
.checked_add(1)
.ok_or_else(|| SoapError::infrastructure("global event position overflow"))?;
records.push(RecordedEvent {
stream_id: stream_id.clone(),
stream_version,
global_position,
envelope,
});
if index + 1 < event_count {
stream_version = stream_version.checked_next().ok_or_else(|| {
SoapError::infrastructure("event stream version overflow")
})?;
}
}
let Some(last) = records.last() else {
return Err(SoapError::infrastructure(
"event append lost its non-empty record set",
));
};
let result = AppendResult {
stream_version: last.stream_version,
global_position: last.global_position,
};
for message in messages {
let id = message.envelope.metadata.id().clone();
outbox_state.order.push(id.clone());
outbox_state.entries.insert(
id,
MemoryOutboxEntry {
record: message,
claim: None,
},
);
}
event_state.next_global_position = next_global_position;
event_state
.streams
.entry(stream_id.clone())
.or_default()
.extend(records.iter().cloned());
event_state.global.extend(records);
Ok(result)
})
}
}
#[cfg(test)]
mod tests {
use std::{
sync::{Arc, Mutex},
time::{Duration, UNIX_EPOCH},
};
use soaprs_contract_tests::{
TransactionalEventOutboxFixture, TransactionalSagaOutboxFixture, block_on,
verify_event_store_contract, verify_inbox_contract, verify_outbox_contract,
verify_projection_checkpoint_contract, verify_saga_store_contract,
verify_snapshot_contract, verify_stored_event_store_contract,
verify_transactional_event_outbox_contract, verify_transactional_projection_contract,
verify_transactional_saga_outbox_contract,
};
use soaprs_core::{
BoxFuture, Command, CommandHandler, MessageEnvelope, MessageMetadata, SoapError, SoapResult,
};
use soaprs_cqrs::{
AggregateCommitResult, CodecEventStore, DeliveryClaimId, EncodedEvent, EventPayloadCodec,
EventReplayHandler, EventReplayer, EventSourcedAggregate, EventSourcedRepository,
EventStore, EventUpcaster, ExpectedSagaVersion, ExpectedVersion, GlobalPosition,
IdempotentProjection, InboxClaim, InboxProcessingOutcome, InboxProcessor, InboxStore,
OutboxProcessingReport, OutboxProcessor, OutboxRecord, OutboxStore, Projection,
ProjectionId, ProjectionRunner, RecordedEvent, ReplayOptions, Saga, SagaAction,
SagaActionProcessingReport, SagaActionProcessor, SagaActionStore, SagaCoordinator, SagaId,
SagaStatus, Snapshot, StoredEventStore, StreamId, StreamVersion,
TransactionalProjectionRunner, TransactionalSagaOutbox, TransientRetryPolicy,
UpcasterChain, VersionedPayload,
};
use soaprs_events::{DomainEvent, Event, EventHandler, EventPublisher, IntegrationEvent};
use super::{
MemoryEventOutbox, MemoryEventStore, MemoryInboxStore, MemoryOutboxStore,
MemoryProjectionCheckpointStore, MemorySagaOutbox, MemorySnapshotStore,
MemoryStoredEventStore, MemoryTransactionalProjection,
};
#[derive(Debug, Clone, PartialEq, Eq)]
enum TestEvent {
Opened,
Renamed(String),
Closed,
}
impl Event for TestEvent {
fn event_type(&self) -> &'static str {
match self {
Self::Opened => "contract.opened",
Self::Renamed(_) => "contract.renamed",
Self::Closed => "contract.closed",
}
}
}
impl DomainEvent for TestEvent {}
#[derive(Debug, Clone, PartialEq, Eq)]
enum CounterEvent {
Added(i32),
}
impl Event for CounterEvent {
fn event_type(&self) -> &'static str {
"contract.counter-added"
}
}
impl DomainEvent for CounterEvent {}
#[derive(Debug, Clone, PartialEq, Eq)]
struct CounterAggregate {
stream_id: StreamId,
version: Option<StreamVersion>,
value: i32,
uncommitted: Vec<MessageEnvelope<CounterEvent>>,
}
impl CounterAggregate {
fn pristine(stream_id: StreamId) -> Self {
Self {
stream_id,
version: None,
value: 0,
uncommitted: Vec::new(),
}
}
}
impl EventSourcedAggregate for CounterAggregate {
type Event = CounterEvent;
fn stream_id(&self) -> &StreamId {
&self.stream_id
}
fn version(&self) -> Option<StreamVersion> {
self.version
}
fn apply(&mut self, event: &Self::Event) -> SoapResult<()> {
match event {
CounterEvent::Added(amount) => {
self.value = self
.value
.checked_add(*amount)
.ok_or_else(|| SoapError::domain("counter overflow"))?;
}
}
Ok(())
}
fn uncommitted_events(&self) -> &[MessageEnvelope<Self::Event>] {
&self.uncommitted
}
fn record(&mut self, event: MessageEnvelope<Self::Event>) -> SoapResult<()> {
self.apply(&event.message)?;
self.uncommitted.push(event);
Ok(())
}
fn mark_committed(&mut self, version: StreamVersion) {
self.version = Some(version);
self.uncommitted.clear();
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum VersionedTestEvent {
Created { name: String, email: String },
}
impl Event for VersionedTestEvent {
fn event_type(&self) -> &'static str {
"contract.user-created"
}
fn schema_version(&self) -> u32 {
2
}
}
impl DomainEvent for VersionedTestEvent {}
struct VersionedTestCodec;
impl EventPayloadCodec<VersionedTestEvent, String> for VersionedTestCodec {
fn encode(&self, event: &VersionedTestEvent) -> SoapResult<String> {
match event {
VersionedTestEvent::Created { name, email } => Ok(format!("{name}|{email}")),
}
}
fn decode(&self, event: VersionedPayload<String>) -> SoapResult<VersionedTestEvent> {
if event.event_type != "contract.user-created" || event.schema_version != 2 {
return Err(SoapError::validation(
"versioned test codec received an unsupported schema",
));
}
let Some((name, email)) = event.payload.split_once('|') else {
return Err(SoapError::validation(
"versioned test payload is missing email",
));
};
Ok(VersionedTestEvent::Created {
name: name.to_owned(),
email: email.to_owned(),
})
}
fn current_schema_version(&self, event_type: &str) -> SoapResult<u32> {
if event_type == "contract.user-created" {
Ok(2)
} else {
Err(SoapError::unsupported("unknown versioned test event type"))
}
}
}
struct AddEmailUpcaster;
impl EventUpcaster<String> for AddEmailUpcaster {
fn event_type(&self) -> &str {
"contract.user-created"
}
fn source_version(&self) -> u32 {
1
}
fn target_version(&self) -> u32 {
2
}
fn upcast(&self, payload: String) -> SoapResult<String> {
Ok(format!("{payload}|unknown@example.test"))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Outgoing(String);
impl Event for Outgoing {
fn event_type(&self) -> &'static str {
"contract.outgoing"
}
}
impl IntegrationEvent for Outgoing {}
struct SelectivePublisher;
impl EventPublisher<Outgoing> for SelectivePublisher {
fn publish(&self, event: MessageEnvelope<Outgoing>) -> BoxFuture<'_, SoapResult<()>> {
Box::pin(async move {
match event.message.0.as_str() {
"retry" => Err(SoapError::timeout("publisher timed out")),
"dead-letter" => Err(SoapError::validation("invalid integration event")),
_ => Ok(()),
}
})
}
}
#[derive(Default)]
struct SelectiveHandler {
seen: Mutex<Vec<String>>,
}
impl EventHandler<Outgoing> for SelectiveHandler {
fn handle<'a>(
&'a self,
event: &'a MessageEnvelope<Outgoing>,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
self.seen
.lock()
.map_err(|_| SoapError::infrastructure("test inbox handler lock poisoned"))?
.push(event.message.0.clone());
if event.message.0 == "fail" {
Err(SoapError::domain("incoming handler failed"))
} else {
Ok(())
}
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum TestCommand {
Continue,
Retry,
Reject,
}
impl Command for TestCommand {
type Output = ();
}
#[derive(Default)]
struct TestCommandHandler {
seen: Mutex<Vec<TestCommand>>,
}
impl CommandHandler<TestCommand> for TestCommandHandler {
fn command(&self, command: TestCommand) -> BoxFuture<'_, SoapResult<()>> {
Box::pin(async move {
self.seen
.lock()
.map_err(|_| SoapError::infrastructure("test command handler lock poisoned"))?
.push(command.clone());
match command {
TestCommand::Continue => Ok(()),
TestCommand::Retry => Err(SoapError::timeout("command timed out")),
TestCommand::Reject => Err(SoapError::validation("command was rejected")),
}
})
}
}
#[derive(Debug, Clone)]
struct TestSaga {
id: SagaId,
status: SagaStatus,
}
impl Saga<TestEvent> for TestSaga {
type Command = TestCommand;
fn id(&self) -> &SagaId {
&self.id
}
fn status(&self) -> SagaStatus {
self.status
}
fn react(
&mut self,
_event: &MessageEnvelope<TestEvent>,
) -> SoapResult<Vec<SagaAction<Self::Command>>> {
self.status = SagaStatus::Completed;
Ok(vec![
SagaAction::Dispatch(TestCommand::Continue),
SagaAction::Complete,
])
}
}
fn event(id: &str, message: TestEvent) -> MessageEnvelope<TestEvent> {
MessageEnvelope::new(message, MessageMetadata::new(id, UNIX_EPOCH))
}
fn recorded_event(position: u64, message: TestEvent) -> RecordedEvent<TestEvent> {
RecordedEvent {
stream_id: StreamId::new("transactional-projection-contract"),
stream_version: StreamVersion::new(position),
global_position: GlobalPosition::new(position),
envelope: event(&format!("projection-contract-{position}"), message),
}
}
fn encoded(id: &str, payload: &str) -> EncodedEvent<String> {
EncodedEvent {
event: VersionedPayload {
event_type: "contract.raw-event".to_owned(),
schema_version: 1,
payload: payload.to_owned(),
},
metadata: MessageMetadata::new(id, UNIX_EPOCH),
}
}
#[test]
fn event_store_passes_shared_contract() {
let store = MemoryEventStore::new();
let result = block_on(verify_event_store_contract(
&store,
StreamId::new("order-1"),
StreamId::new("order-2"),
event("event-1", TestEvent::Opened),
event("event-2", TestEvent::Renamed("new".into())),
event("event-3", TestEvent::Closed),
));
assert!(result.is_ok(), "{result:?}");
}
#[test]
fn event_sourced_repository_rehydrates_and_preserves_conflicting_changes() {
let store = MemoryEventStore::new();
let factory = |stream_id| Ok(CounterAggregate::pristine(stream_id));
let repository = EventSourcedRepository::new(&store, &factory);
let stream = StreamId::new("counter-1");
let aggregate: SoapResult<CounterAggregate> = repository.new_aggregate(stream.clone());
let Some(mut aggregate) = aggregate.ok() else {
panic!("counter aggregate factory was rejected");
};
let recorded = aggregate.record(MessageEnvelope::new(
CounterEvent::Added(1),
MessageMetadata::new("counter-event-1", UNIX_EPOCH),
));
assert!(recorded.is_ok(), "{recorded:?}");
let first_commit = block_on(repository.commit(&mut aggregate));
assert!(matches!(
first_commit,
Ok(AggregateCommitResult::Appended { event_count: 1, .. })
));
assert_eq!(aggregate.version, Some(StreamVersion::FIRST));
assert!(aggregate.uncommitted.is_empty());
let unchanged = block_on(repository.commit(&mut aggregate));
assert_eq!(
unchanged.ok(),
Some(AggregateCommitResult::Unchanged {
stream_version: Some(StreamVersion::FIRST),
})
);
let first: SoapResult<Option<CounterAggregate>> = block_on(repository.load(&stream));
let second: SoapResult<Option<CounterAggregate>> = block_on(repository.load(&stream));
let Some(mut first) = first.ok().flatten() else {
panic!("first counter copy was not rehydrated");
};
let Some(mut second) = second.ok().flatten() else {
panic!("second counter copy was not rehydrated");
};
assert_eq!(first.value, 1);
assert_eq!(first.version, Some(StreamVersion::FIRST));
let first_recorded = first.record(MessageEnvelope::new(
CounterEvent::Added(2),
MessageMetadata::new("counter-event-2", UNIX_EPOCH),
));
let second_recorded = second.record(MessageEnvelope::new(
CounterEvent::Added(4),
MessageMetadata::new("counter-event-conflict", UNIX_EPOCH),
));
assert!(first_recorded.is_ok(), "{first_recorded:?}");
assert!(second_recorded.is_ok(), "{second_recorded:?}");
let winning_commit = block_on(repository.commit(&mut first));
assert!(winning_commit.is_ok(), "{winning_commit:?}");
let conflicting_commit = block_on(repository.commit(&mut second));
assert!(conflicting_commit.is_err());
assert_eq!(second.version, Some(StreamVersion::FIRST));
assert_eq!(second.uncommitted.len(), 1);
let reloaded: SoapResult<Option<CounterAggregate>> = block_on(repository.load(&stream));
assert!(matches!(reloaded, Ok(Some(current)) if current.value == 3
&& current.version == Some(StreamVersion::new(2))
&& current.uncommitted.is_empty()));
}
#[test]
fn stored_event_store_passes_shared_contract() {
let store = MemoryStoredEventStore::new();
let result = block_on(verify_stored_event_store_contract(
&store,
StreamId::new("raw-stream-1"),
StreamId::new("raw-stream-2"),
encoded("raw-event-1", "first"),
encoded("raw-event-2", "second"),
encoded("raw-event-3", "other"),
));
assert!(result.is_ok(), "{result:?}");
}
#[test]
fn codec_event_store_upcasts_old_payloads_and_preserves_metadata() {
let raw = MemoryStoredEventStore::new();
let stream = StreamId::new("versioned-user-1");
let seeded = block_on(raw.append(
&stream,
ExpectedVersion::NoStream,
vec![EncodedEvent {
event: VersionedPayload {
event_type: "contract.user-created".to_owned(),
schema_version: 1,
payload: "Ada".to_owned(),
},
metadata: MessageMetadata::new("versioned-event-1", UNIX_EPOCH)
.with_correlation_id("registration-1"),
}],
));
assert!(seeded.is_ok(), "{seeded:?}");
let mut upcasters = UpcasterChain::new();
let registered = upcasters.register(Arc::new(AddEmailUpcaster));
assert!(registered.is_ok(), "{registered:?}");
let store = CodecEventStore::new(raw, VersionedTestCodec, upcasters);
let loaded: SoapResult<Vec<RecordedEvent<VersionedTestEvent>>> =
block_on(store.load(&stream, None));
assert!(matches!(loaded, Ok(events) if events.len() == 1
&& events[0].envelope.message == VersionedTestEvent::Created {
name: "Ada".to_owned(),
email: "unknown@example.test".to_owned(),
}
&& events[0]
.envelope
.metadata
.correlation_id()
.is_some_and(|id| id.as_str() == "registration-1")));
let appended = block_on(store.append(
&stream,
ExpectedVersion::Exact(StreamVersion::FIRST),
vec![MessageEnvelope::new(
VersionedTestEvent::Created {
name: "Grace".to_owned(),
email: "grace@example.test".to_owned(),
},
MessageMetadata::new("versioned-event-2", UNIX_EPOCH),
)],
));
assert!(appended.is_ok(), "{appended:?}");
let raw_events = block_on(store.store().load(&stream, Some(StreamVersion::FIRST)));
assert!(matches!(raw_events, Ok(events) if events.len() == 1
&& events[0].encoded.event.schema_version == 2
&& events[0].encoded.event.payload == "Grace|grace@example.test"));
}
#[test]
fn checkpoint_inbox_and_snapshot_pass_shared_contracts() {
let checkpoint_result = block_on(verify_projection_checkpoint_contract(
&MemoryProjectionCheckpointStore::new(),
));
assert!(checkpoint_result.is_ok(), "{checkpoint_result:?}");
let inbox_result = block_on(verify_inbox_contract(
&MemoryInboxStore::new(),
MessageMetadata::new("incoming-1", UNIX_EPOCH).id(),
UNIX_EPOCH,
));
assert!(inbox_result.is_ok(), "{inbox_result:?}");
let snapshots = MemorySnapshotStore::new();
let first = Snapshot {
stream_id: StreamId::new("order-1"),
stream_version: StreamVersion::new(1),
state: "first".to_owned(),
captured_at: UNIX_EPOCH,
};
let second = Snapshot {
stream_id: StreamId::new("order-1"),
stream_version: StreamVersion::new(2),
state: "second".to_owned(),
captured_at: UNIX_EPOCH + Duration::from_secs(1),
};
let snapshot_result = block_on(verify_snapshot_contract(&snapshots, first, second));
assert!(snapshot_result.is_ok(), "{snapshot_result:?}");
}
#[test]
fn outbox_claims_once_until_failure_releases_the_message() {
let store = MemoryOutboxStore::new();
let record = OutboxRecord::pending(
MessageEnvelope::new(
Outgoing("payload".to_owned()),
MessageMetadata::new("outgoing-1", UNIX_EPOCH),
),
UNIX_EPOCH,
);
let dead_letter_record = OutboxRecord::pending(
MessageEnvelope::new(
Outgoing("dead-letter".to_owned()),
MessageMetadata::new("outgoing-dead-letter", UNIX_EPOCH),
),
UNIX_EPOCH,
);
let result = block_on(verify_outbox_contract(
&store,
record,
dead_letter_record,
UNIX_EPOCH,
));
assert!(result.is_ok(), "{result:?}");
}
#[test]
fn outbox_processor_delivers_retries_and_dead_letters_with_one_policy() {
let store = MemoryOutboxStore::new();
let now = UNIX_EPOCH + Duration::from_secs(10);
let record = |id: &str, payload: &str| {
OutboxRecord::pending(
MessageEnvelope::new(Outgoing(payload.to_owned()), MessageMetadata::new(id, now)),
now,
)
};
let enqueued = block_on(store.enqueue(vec![
record("delivery-success", "success"),
record("delivery-retry", "retry"),
record("delivery-dead-letter", "dead-letter"),
]));
assert!(enqueued.is_ok(), "{enqueued:?}");
let policy = TransientRetryPolicy::new(3, Duration::from_secs(5), Duration::from_secs(30));
let Some(policy) = policy.ok() else {
panic!("valid retry policy was rejected");
};
let clock = move || now;
let publisher = SelectivePublisher;
let processor = OutboxProcessor::new(&store, &publisher, &policy, &clock);
let report = block_on(processor.process_batch::<Outgoing>(
DeliveryClaimId::new("processor-claim"),
Duration::from_secs(60),
10,
));
assert_eq!(
report.ok(),
Some(OutboxProcessingReport {
claimed: 3,
delivered: 1,
retry_scheduled: 1,
dead_lettered: 1,
})
);
let retry = block_on(store.claim_pending(
DeliveryClaimId::new("retry-claim"),
now + Duration::from_secs(5),
now + Duration::from_secs(65),
10,
));
assert!(matches!(retry, Ok(records) if records.len() == 1
&& records[0].record.envelope.metadata.id().as_str() == "delivery-retry"
&& records[0].record.attempts == 1));
}
#[test]
fn inbox_processor_completes_duplicates_and_releases_failures() {
let store = MemoryInboxStore::new();
let handler = SelectiveHandler::default();
let now = UNIX_EPOCH + Duration::from_secs(10);
let clock = move || now;
let processor = InboxProcessor::new(&store, &handler, &clock);
let successful = MessageEnvelope::new(
Outgoing("success".to_owned()),
MessageMetadata::new("incoming-success", now),
);
let first = block_on(processor.process(
"orders",
successful.clone(),
DeliveryClaimId::new("inbox-success-1"),
Duration::from_secs(60),
));
assert_eq!(first.ok(), Some(InboxProcessingOutcome::Processed));
let duplicate = block_on(processor.process(
"orders",
successful,
DeliveryClaimId::new("inbox-success-2"),
Duration::from_secs(60),
));
assert_eq!(
duplicate.ok(),
Some(InboxProcessingOutcome::AlreadyProcessed)
);
let active = MessageEnvelope::new(
Outgoing("active".to_owned()),
MessageMetadata::new("incoming-active", now),
);
let active_claim = block_on(store.claim(
"orders",
active.metadata.id(),
DeliveryClaimId::new("inbox-active-owner"),
now,
now + Duration::from_secs(60),
));
assert!(matches!(active_claim, Ok(InboxClaim::Acquired)));
let in_progress = block_on(processor.process(
"orders",
active,
DeliveryClaimId::new("inbox-active-duplicate"),
Duration::from_secs(60),
));
assert_eq!(in_progress.ok(), Some(InboxProcessingOutcome::InProgress));
let failing = || {
MessageEnvelope::new(
Outgoing("fail".to_owned()),
MessageMetadata::new("incoming-failure", now),
)
};
let first_failure = block_on(processor.process(
"orders",
failing(),
DeliveryClaimId::new("inbox-failure-1"),
Duration::from_secs(60),
));
let second_failure = block_on(processor.process(
"orders",
failing(),
DeliveryClaimId::new("inbox-failure-2"),
Duration::from_secs(60),
));
assert!(first_failure.is_err());
assert!(second_failure.is_err());
assert_eq!(
handler.seen.lock().ok().map(|seen| seen.clone()),
Some(vec![
"success".to_owned(),
"fail".to_owned(),
"fail".to_owned(),
])
);
}
#[test]
fn transactional_outbox_does_not_commit_messages_after_version_conflict() {
let store = MemoryEventOutbox::new();
let stream = StreamId::new("order-atomic");
let outgoing = |id: &str| {
OutboxRecord::pending(
MessageEnvelope::new(
Outgoing(id.to_owned()),
MessageMetadata::new(id, UNIX_EPOCH),
),
UNIX_EPOCH,
)
};
let result = block_on(verify_transactional_event_outbox_contract(
&store,
store.event_store(),
store.outbox(),
TransactionalEventOutboxFixture {
stream,
first_event: event("atomic-event-1", TestEvent::Opened),
conflicting_event: event("atomic-event-2", TestEvent::Closed),
first_message: outgoing("atomic-message-1"),
conflicting_message: outgoing("atomic-message-2"),
now: UNIX_EPOCH,
},
));
assert!(result.is_ok(), "{result:?}");
}
struct RecordingReplay {
seen: Mutex<Vec<u64>>,
}
impl EventReplayHandler<TestEvent> for RecordingReplay {
fn handle_batch<'a>(
&'a self,
events: &'a [RecordedEvent<TestEvent>],
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
let mut seen = self
.seen
.lock()
.map_err(|_| SoapError::infrastructure("replay test lock poisoned"))?;
seen.extend(events.iter().map(|item| item.global_position.get()));
Ok(())
})
}
}
struct RecordingProjection {
id: ProjectionId,
seen: Mutex<Vec<u64>>,
}
impl Projection<TestEvent> for RecordingProjection {
fn id(&self) -> &ProjectionId {
&self.id
}
fn project<'a>(
&'a self,
event: &'a RecordedEvent<TestEvent>,
) -> BoxFuture<'a, SoapResult<()>> {
Box::pin(async move {
self.seen
.lock()
.map_err(|_| SoapError::infrastructure("projection test lock poisoned"))?
.push(event.global_position.get());
Ok(())
})
}
}
impl IdempotentProjection<TestEvent> for RecordingProjection {}
#[test]
fn replay_and_projection_resume_from_durable_positions() {
let store = MemoryEventStore::new();
let stream = StreamId::new("replay-stream");
let append = block_on(store.append(
&stream,
ExpectedVersion::NoStream,
vec![
event("replay-1", TestEvent::Opened),
event("replay-2", TestEvent::Closed),
],
));
assert!(append.is_ok(), "{append:?}");
let handler = RecordingReplay {
seen: Mutex::new(Vec::new()),
};
let replayer = EventReplayer::new(&store);
let Some(options) = ReplayOptions::new(1).ok() else {
panic!("valid replay options were rejected");
};
let replay = block_on(replayer.replay::<TestEvent, _>(options, &handler));
assert_eq!(replay.ok().map(|report| report.processed), Some(2));
assert_eq!(
handler.seen.lock().ok().map(|seen| seen.clone()),
Some(vec![1, 2])
);
let checkpoints = MemoryProjectionCheckpointStore::new();
let projection = RecordingProjection {
id: ProjectionId::new("orders"),
seen: Mutex::new(Vec::new()),
};
let runner = ProjectionRunner::new(&store, &checkpoints);
assert_eq!(
block_on(runner.run_once::<TestEvent, _>(&projection, 10)).ok(),
Some(2)
);
assert_eq!(
block_on(runner.run_once::<TestEvent, _>(&projection, 10)).ok(),
Some(0)
);
}
#[test]
fn transactional_projection_passes_shared_contract() {
let projection = MemoryTransactionalProjection::new(
ProjectionId::new("transactional-contract"),
Vec::<u64>::new(),
|model: &mut Vec<u64>, event: &RecordedEvent<TestEvent>| {
model.push(event.global_position.get());
Ok(())
},
);
let first = recorded_event(1, TestEvent::Opened);
let second = recorded_event(2, TestEvent::Closed);
let result = block_on(verify_transactional_projection_contract(
&projection,
&first,
&second,
));
assert!(result.is_ok(), "{result:?}");
assert_eq!(
projection.snapshot().ok(),
Some((vec![1, 2], Some(GlobalPosition::new(2))))
);
}
#[test]
fn transactional_projection_rolls_back_model_and_checkpoint_on_failure() {
let store = MemoryEventStore::new();
let stream = StreamId::new("transactional-projection-stream");
let append = block_on(store.append(
&stream,
ExpectedVersion::NoStream,
vec![
event("projection-1", TestEvent::Opened),
event("projection-2", TestEvent::Renamed("fail".to_owned())),
event("projection-3", TestEvent::Closed),
],
));
assert!(append.is_ok(), "{append:?}");
let projection = MemoryTransactionalProjection::new(
ProjectionId::new("transactional-rollback"),
Vec::<String>::new(),
|model: &mut Vec<String>, event: &RecordedEvent<TestEvent>| {
match &event.envelope.message {
TestEvent::Opened => model.push("opened".to_owned()),
TestEvent::Renamed(name) => {
model.push(name.clone());
return Err(SoapError::domain("projection update failed"));
}
TestEvent::Closed => model.push("closed".to_owned()),
}
Ok(())
},
);
let runner = TransactionalProjectionRunner::new(&store);
let result = block_on(runner.run_once::<TestEvent, _>(&projection, 10));
assert_eq!(
result.as_ref().map_err(SoapError::kind),
Err(soaprs_core::SoapErrorKind::Domain)
);
assert_eq!(
projection.snapshot().ok(),
Some((vec!["opened".to_owned()], Some(GlobalPosition::new(1))))
);
}
#[test]
fn saga_store_round_trips_application_state() {
let store = super::MemorySagaStore::new();
let id = SagaId::new("checkout-1");
let result = block_on(verify_saga_store_contract(
&store,
&id,
"waiting-for-payment".to_owned(),
"completed".to_owned(),
));
assert!(result.is_ok(), "{result:?}");
}
#[test]
fn saga_outbox_passes_atomic_state_action_and_timer_contract() {
let store = MemorySagaOutbox::new();
let later = UNIX_EPOCH + Duration::from_secs(60);
let result = block_on(verify_transactional_saga_outbox_contract(
&store,
store.saga_store(),
store.action_store(),
TransactionalSagaOutboxFixture {
saga_id: SagaId::new("contract-durable-saga"),
first_state: "running".to_owned(),
conflicting_state: "stale".to_owned(),
first_transition: MessageMetadata::new("saga-transition-1", UNIX_EPOCH),
first_actions: vec![
SagaAction::Dispatch(TestCommand::Continue),
SagaAction::Schedule {
delay: Duration::from_secs(60),
command: TestCommand::Continue,
},
SagaAction::Complete,
],
conflicting_transition: MessageMetadata::new(
"saga-transition-conflict",
UNIX_EPOCH,
),
conflicting_actions: vec![SagaAction::Dispatch(TestCommand::Reject)],
now: UNIX_EPOCH,
later,
expected_due_now: 2,
expected_due_later: 1,
},
));
assert!(result.is_ok(), "{result:?}");
}
#[test]
fn saga_coordinator_commits_actions_instead_of_returning_them() {
let store = MemorySagaOutbox::new();
let coordinator = SagaCoordinator::new(&store);
let saga = TestSaga {
id: SagaId::new("durable-coordinated-saga"),
status: SagaStatus::Pending,
};
let triggering_event = event("durable-saga-event", TestEvent::Opened);
let transition = block_on(coordinator.transition(
saga,
ExpectedSagaVersion::NoSaga,
&triggering_event,
MessageMetadata::new("durable-saga-transition", UNIX_EPOCH),
));
assert_eq!(
transition
.as_ref()
.ok()
.map(|result| result.commit.queued_actions),
Some(2)
);
let claimed = block_on(store.action_store().claim_pending(
DeliveryClaimId::new("durable-saga-worker"),
UNIX_EPOCH,
UNIX_EPOCH + Duration::from_secs(60),
10,
));
assert!(matches!(claimed, Ok(actions) if actions.len() == 2));
}
#[test]
fn saga_action_processor_dispatches_retries_dead_letters_and_waits_for_timers() {
let store = MemorySagaOutbox::new();
let committed = block_on(store.commit_transition(
&SagaId::new("processed-saga"),
ExpectedSagaVersion::NoSaga,
"running".to_owned(),
MessageMetadata::new("processed-saga-transition", UNIX_EPOCH),
vec![
SagaAction::Dispatch(TestCommand::Continue),
SagaAction::Dispatch(TestCommand::Retry),
SagaAction::Compensate(TestCommand::Reject),
SagaAction::Complete,
SagaAction::Schedule {
delay: Duration::from_secs(60),
command: TestCommand::Continue,
},
],
));
assert!(committed.is_ok(), "{committed:?}");
let policy = TransientRetryPolicy::new(3, Duration::from_secs(5), Duration::from_secs(30));
let Some(policy) = policy.ok() else {
panic!("valid retry policy was rejected");
};
let handler = TestCommandHandler::default();
let clock = || UNIX_EPOCH;
let processor = SagaActionProcessor::new(store.action_store(), &handler, &policy, &clock);
let report = block_on(processor.process_batch::<TestCommand>(
DeliveryClaimId::new("saga-action-processor"),
Duration::from_secs(30),
10,
));
assert_eq!(
report.ok(),
Some(SagaActionProcessingReport {
claimed: 4,
dispatched: 1,
markers_completed: 1,
retry_scheduled: 1,
dead_lettered: 1,
})
);
let before_timer = block_on(store.action_store().claim_pending(
DeliveryClaimId::new("saga-before-timer"),
UNIX_EPOCH + Duration::from_secs(4),
UNIX_EPOCH + Duration::from_secs(30),
10,
));
assert!(matches!(before_timer, Ok(actions) if actions.is_empty()));
let retry_claim = DeliveryClaimId::new("saga-retry");
let retry = block_on(store.action_store().claim_pending(
retry_claim.clone(),
UNIX_EPOCH + Duration::from_secs(5),
UNIX_EPOCH + Duration::from_secs(30),
10,
));
let Ok(retry) = retry else {
panic!("retryable saga action could not be claimed");
};
assert_eq!(retry.len(), 1);
assert_eq!(retry[0].record.attempts, 1);
let retry_completed = block_on(store.action_store().mark_completed(
&retry[0].record.id,
&retry_claim,
UNIX_EPOCH + Duration::from_secs(5),
));
assert!(retry_completed.is_ok(), "{retry_completed:?}");
let timer = block_on(store.action_store().claim_pending(
DeliveryClaimId::new("saga-timer"),
UNIX_EPOCH + Duration::from_secs(60),
UNIX_EPOCH + Duration::from_secs(90),
10,
));
assert!(matches!(timer, Ok(actions) if actions.len() == 1
&& matches!(actions[0].record.action, SagaAction::Schedule { .. })));
}
}