pub mod claim;
pub mod traveler;
pub mod trip;
use std::collections::HashMap;
use std::fmt;
use std::sync::Mutex;
use chrono::{DateTime, Utc};
use turnframe_core::case::Versioned;
use turnframe_core::command::{AtomicityScope, CommandBatch, IdempotencyKey};
use turnframe_core::error::{DomainRejection, ExecutionError, RevisionConflict, StoreError};
use turnframe_core::event::{Commit, CommittedEvent};
use turnframe_core::flow::{WorkflowDefinition, WorkflowExecutor};
use turnframe_core::hash::{Digest, canonical_digest, derive_uuid};
use turnframe_core::ids::{AccountId, CaseId, CaseRevision, EventId};
use crate::explore::SimulatedTransition;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Applied<S, E> {
pub state: S,
pub events: Vec<E>,
}
impl<S, E> Applied<S, E> {
#[must_use]
pub const fn new(state: S, events: Vec<E>) -> Self {
Self { state, events }
}
}
pub trait PureWorkflow: WorkflowDefinition {
fn apply(
&self,
state: Option<&Self::State>,
command: &Self::Command,
) -> Result<Applied<Self::State, Self::Event>, DomainRejection>;
fn event_type(&self, event: &Self::Event) -> String;
}
pub(crate) fn in_italian(
specs: Vec<turnframe_core::operation::OperationSpec>,
summaries: &[(&str, &str)],
) -> Vec<turnframe_core::operation::OperationSpec> {
specs
.into_iter()
.map(
|spec| match summaries.iter().find(|(key, _)| spec.key.as_str() == *key) {
Some((_, italian)) => spec.summary_in("it-IT", *italian),
None => spec,
},
)
.collect()
}
pub fn simulate<W: PureWorkflow>(
definition: &W,
state: Option<&W::State>,
command: &W::Command,
) -> SimulatedTransition<W::State, W::Event> {
match definition.apply(state, command) {
Ok(applied) => SimulatedTransition::applied(applied.state, applied.events),
Err(rejection) => SimulatedTransition::rejected(rejection),
}
}
const EVENT_ID_DOMAIN: &str = "turnframe.test.event.v1";
const CLOCK_EPOCH_SECONDS: i64 = 1_700_000_000;
#[derive(Clone)]
struct Replayable<S, E> {
command_digest: Digest,
revision_before: CaseRevision,
events: Vec<CommittedEvent<E>>,
state_after: S,
}
struct Store<S, E> {
cases: HashMap<(AccountId, CaseId), Versioned<S>>,
created: Vec<(AccountId, CaseId)>,
replays: HashMap<IdempotencyKey, Replayable<S, E>>,
sequence: u64,
}
pub struct InMemoryExecutor<W: PureWorkflow> {
definition: W,
store: Mutex<Store<W::State, W::Event>>,
}
impl<W: PureWorkflow> fmt::Debug for InMemoryExecutor<W> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("InMemoryExecutor")
.field("workflow", &self.definition.key())
.field("version", &self.definition.version())
.finish_non_exhaustive()
}
}
impl<W: PureWorkflow + Default> Default for InMemoryExecutor<W> {
fn default() -> Self {
Self::new(W::default())
}
}
impl<W: PureWorkflow> InMemoryExecutor<W> {
#[must_use]
pub fn new(definition: W) -> Self {
Self {
definition,
store: Mutex::new(Store {
cases: HashMap::new(),
created: Vec::new(),
replays: HashMap::new(),
sequence: 0,
}),
}
}
#[must_use]
pub const fn definition(&self) -> &W {
&self.definition
}
fn store(&self) -> std::sync::MutexGuard<'_, Store<W::State, W::Event>> {
self.store
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn seed(
&self,
account: &AccountId,
case_id: &CaseId,
state: W::State,
revision: CaseRevision,
) {
let key = (account.clone(), case_id.clone());
let mut store = self.store();
store.remember(&key);
store.cases.insert(key, Versioned::new(state, revision));
}
#[must_use]
pub fn case_ids(&self, account: &AccountId) -> Vec<CaseId> {
self.store()
.created
.iter()
.filter(|(owner, _)| owner == account)
.map(|(_, case_id)| case_id.clone())
.collect()
}
#[must_use]
pub fn revision_of(&self, account: &AccountId, case_id: &CaseId) -> CaseRevision {
self.store()
.cases
.get(&(account.clone(), case_id.clone()))
.map_or(CaseRevision::ZERO, |case| case.revision)
}
#[must_use]
pub fn state_of(&self, account: &AccountId, case_id: &CaseId) -> Option<W::State> {
self.store()
.cases
.get(&(account.clone(), case_id.clone()))
.map(|case| case.value.clone())
}
#[must_use]
pub fn case_count(&self) -> usize {
self.store().cases.len()
}
#[must_use]
pub fn has_executed(&self, key: &IdempotencyKey) -> bool {
self.store().replays.contains_key(key)
}
#[must_use]
pub fn replayed_prefix_of(&self, batch: &CommandBatch<W::Command>) -> usize {
let store = self.store();
batch
.envelopes
.iter()
.take_while(|envelope| store.replays.contains_key(&envelope.idempotency_key))
.count()
}
pub fn execute_prefix(
&self,
batch: &CommandBatch<W::Command>,
applied: usize,
) -> Result<Commit<W::State, W::Event>, ExecutionError> {
if applied == 0 || applied > batch.envelopes.len() {
return Err(ExecutionError::ScopeViolation);
}
self.run(batch, applied)
}
fn run(
&self,
batch: &CommandBatch<W::Command>,
limit: usize,
) -> Result<Commit<W::State, W::Event>, ExecutionError> {
let first = batch
.envelopes
.first()
.ok_or(ExecutionError::ScopeViolation)?;
if matches!(batch.scope, AtomicityScope::PerCase) && !batch.is_single_case() {
return Err(ExecutionError::ScopeViolation);
}
let key = (first.account_id().clone(), first.case_ref.case_id.clone());
let mut store = self.store();
let envelopes = &batch.envelopes[..limit.min(batch.envelopes.len())];
let mut recorded: Vec<Option<Replayable<W::State, W::Event>>> =
Vec::with_capacity(envelopes.len());
for envelope in envelopes {
let digest = canonical_digest(&envelope.command)
.map_err(|_| ExecutionError::Store(StoreError::Serialization))?;
match store.replays.get(&envelope.idempotency_key) {
None => recorded.push(None),
Some(entry) if entry.command_digest == digest => recorded.push(Some(entry.clone())),
Some(_) => {
return Err(ExecutionError::IdempotencyMismatch {
command_id: envelope.command_id,
});
}
}
}
let replayed = recorded.iter().take_while(|entry| entry.is_some()).count();
if let Some(position) = recorded[replayed..].iter().position(Option::is_some) {
return Err(ExecutionError::IdempotencyMismatch {
command_id: envelopes[replayed + position].command_id,
});
}
let current = store.cases.get(&key).cloned();
let current_revision = current.as_ref().map_or(CaseRevision::ZERO, |c| c.revision);
let conflict = |current_revision| {
ExecutionError::RevisionConflict(RevisionConflict {
expected: first.case_ref.clone(),
current_revision,
})
};
let last_replayed = recorded[..replayed].last().and_then(Option::as_ref);
let (mut state, revision_before, mut events) = match last_replayed {
None => {
if current_revision != first.case_ref.expected_revision {
return Err(conflict(current_revision));
}
(current.map(|c| c.value), current_revision, Vec::new())
}
Some(entry) => {
if entry.revision_before != first.case_ref.expected_revision
|| current_revision != entry.revision_before.next()
{
return Err(conflict(current_revision));
}
let events = recorded[..replayed]
.iter()
.flatten()
.flat_map(|entry| entry.events.clone())
.collect();
(
Some(entry.state_after.clone()),
entry.revision_before,
events,
)
}
};
if replayed == envelopes.len() {
let Some(committed) = state else {
return Err(ExecutionError::ScopeViolation);
};
return Ok(Commit {
state: Some(committed),
new_revision: revision_before.next(),
events,
idempotency_replay: true,
});
}
let new_revision = revision_before.next();
let mut fresh = Vec::new();
for envelope in &envelopes[replayed..] {
self.definition
.validate_command(state.as_ref(), &envelope.command)
.map_err(ExecutionError::Rejected)?;
let applied = self
.definition
.apply(state.as_ref(), &envelope.command)
.map_err(ExecutionError::Rejected)?;
let mut committed_events = Vec::with_capacity(applied.events.len());
for payload in applied.events {
let event_type = self.definition.event_type(&payload);
let (sequence, occurred_at) = store.tick();
committed_events.push(CommittedEvent {
event_id: EventId::from(derive_uuid(
EVENT_ID_DOMAIN,
&[
key.0.as_str(),
key.1.as_str(),
&sequence.to_string(),
&event_type,
],
)),
event_type,
occurred_at,
payload,
});
}
let digest = canonical_digest(&envelope.command)
.map_err(|_| ExecutionError::Store(StoreError::Serialization))?;
state = Some(applied.state.clone());
fresh.push((
envelope.idempotency_key.clone(),
Replayable {
command_digest: digest,
revision_before,
events: committed_events.clone(),
state_after: applied.state,
},
));
events.extend(committed_events);
}
let Some(committed) = state else {
return Err(ExecutionError::ScopeViolation);
};
store.remember(&key);
store
.cases
.insert(key, Versioned::new(committed.clone(), new_revision));
for (idempotency_key, entry) in fresh {
store.replays.insert(idempotency_key, entry);
}
Ok(Commit {
state: Some(committed),
new_revision,
events,
idempotency_replay: replayed > 0,
})
}
}
impl<S, E> Store<S, E> {
fn remember(&mut self, key: &(AccountId, CaseId)) {
if !self.cases.contains_key(key) {
self.created.push(key.clone());
}
}
fn tick(&mut self) -> (u64, DateTime<Utc>) {
self.sequence += 1;
let seconds = CLOCK_EPOCH_SECONDS.saturating_add(self.sequence as i64);
(
self.sequence,
DateTime::from_timestamp(seconds, 0).unwrap_or_default(),
)
}
}
#[async_trait::async_trait]
impl<W: PureWorkflow> WorkflowExecutor<W> for InMemoryExecutor<W> {
async fn load(
&self,
account: &AccountId,
case_id: &CaseId,
) -> Result<Versioned<Option<W::State>>, StoreError> {
Ok(self
.store()
.cases
.get(&(account.clone(), case_id.clone()))
.map_or_else(
|| Versioned::new(None, CaseRevision::ZERO),
|case| Versioned::new(Some(case.value.clone()), case.revision),
))
}
async fn execute(
&self,
batch: CommandBatch<W::Command>,
) -> Result<Commit<W::State, W::Event>, ExecutionError> {
self.run(&batch, batch.envelopes.len())
}
}