1pub mod claim;
32pub mod traveler;
33pub mod trip;
34
35use std::collections::HashMap;
36use std::fmt;
37use std::sync::Mutex;
38
39use chrono::{DateTime, Utc};
40use turnframe_core::case::Versioned;
41use turnframe_core::command::{AtomicityScope, CommandBatch, IdempotencyKey};
42use turnframe_core::error::{DomainRejection, ExecutionError, RevisionConflict, StoreError};
43use turnframe_core::event::{Commit, CommittedEvent};
44use turnframe_core::flow::{WorkflowDefinition, WorkflowExecutor};
45use turnframe_core::hash::{Digest, canonical_digest, derive_uuid};
46use turnframe_core::ids::{AccountId, CaseId, CaseRevision, EventId};
47
48use crate::explore::SimulatedTransition;
49
50#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct Applied<S, E> {
53 pub state: S,
55 pub events: Vec<E>,
57}
58
59impl<S, E> Applied<S, E> {
60 #[must_use]
62 pub const fn new(state: S, events: Vec<E>) -> Self {
63 Self { state, events }
64 }
65}
66
67pub trait PureWorkflow: WorkflowDefinition {
73 fn apply(
76 &self,
77 state: Option<&Self::State>,
78 command: &Self::Command,
79 ) -> Result<Applied<Self::State, Self::Event>, DomainRejection>;
80
81 fn event_type(&self, event: &Self::Event) -> String;
83}
84
85pub(crate) fn in_italian(
88 specs: Vec<turnframe_core::operation::OperationSpec>,
89 summaries: &[(&str, &str)],
90) -> Vec<turnframe_core::operation::OperationSpec> {
91 specs
92 .into_iter()
93 .map(
94 |spec| match summaries.iter().find(|(key, _)| spec.key.as_str() == *key) {
95 Some((_, italian)) => spec.summary_in("it-IT", *italian),
96 None => spec,
97 },
98 )
99 .collect()
100}
101
102pub fn simulate<W: PureWorkflow>(
105 definition: &W,
106 state: Option<&W::State>,
107 command: &W::Command,
108) -> SimulatedTransition<W::State, W::Event> {
109 match definition.apply(state, command) {
110 Ok(applied) => SimulatedTransition::applied(applied.state, applied.events),
111 Err(rejection) => SimulatedTransition::rejected(rejection),
112 }
113}
114
115const EVENT_ID_DOMAIN: &str = "turnframe.test.event.v1";
117
118const CLOCK_EPOCH_SECONDS: i64 = 1_700_000_000;
120
121#[derive(Clone)]
128struct Replayable<S, E> {
129 command_digest: Digest,
132 revision_before: CaseRevision,
135 events: Vec<CommittedEvent<E>>,
137 state_after: S,
139}
140
141struct Store<S, E> {
142 cases: HashMap<(AccountId, CaseId), Versioned<S>>,
143 created: Vec<(AccountId, CaseId)>,
145 replays: HashMap<IdempotencyKey, Replayable<S, E>>,
146 sequence: u64,
147}
148
149pub struct InMemoryExecutor<W: PureWorkflow> {
176 definition: W,
177 store: Mutex<Store<W::State, W::Event>>,
178}
179
180impl<W: PureWorkflow> fmt::Debug for InMemoryExecutor<W> {
181 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
182 f.debug_struct("InMemoryExecutor")
183 .field("workflow", &self.definition.key())
184 .field("version", &self.definition.version())
185 .finish_non_exhaustive()
186 }
187}
188
189impl<W: PureWorkflow + Default> Default for InMemoryExecutor<W> {
190 fn default() -> Self {
191 Self::new(W::default())
192 }
193}
194
195impl<W: PureWorkflow> InMemoryExecutor<W> {
196 #[must_use]
198 pub fn new(definition: W) -> Self {
199 Self {
200 definition,
201 store: Mutex::new(Store {
202 cases: HashMap::new(),
203 created: Vec::new(),
204 replays: HashMap::new(),
205 sequence: 0,
206 }),
207 }
208 }
209
210 #[must_use]
212 pub const fn definition(&self) -> &W {
213 &self.definition
214 }
215
216 fn store(&self) -> std::sync::MutexGuard<'_, Store<W::State, W::Event>> {
219 self.store
220 .lock()
221 .unwrap_or_else(|poisoned| poisoned.into_inner())
222 }
223
224 pub fn seed(
227 &self,
228 account: &AccountId,
229 case_id: &CaseId,
230 state: W::State,
231 revision: CaseRevision,
232 ) {
233 let key = (account.clone(), case_id.clone());
234 let mut store = self.store();
235 store.remember(&key);
236 store.cases.insert(key, Versioned::new(state, revision));
237 }
238
239 #[must_use]
242 pub fn case_ids(&self, account: &AccountId) -> Vec<CaseId> {
243 self.store()
244 .created
245 .iter()
246 .filter(|(owner, _)| owner == account)
247 .map(|(_, case_id)| case_id.clone())
248 .collect()
249 }
250
251 #[must_use]
254 pub fn revision_of(&self, account: &AccountId, case_id: &CaseId) -> CaseRevision {
255 self.store()
256 .cases
257 .get(&(account.clone(), case_id.clone()))
258 .map_or(CaseRevision::ZERO, |case| case.revision)
259 }
260
261 #[must_use]
263 pub fn state_of(&self, account: &AccountId, case_id: &CaseId) -> Option<W::State> {
264 self.store()
265 .cases
266 .get(&(account.clone(), case_id.clone()))
267 .map(|case| case.value.clone())
268 }
269
270 #[must_use]
272 pub fn case_count(&self) -> usize {
273 self.store().cases.len()
274 }
275
276 #[must_use]
278 pub fn has_executed(&self, key: &IdempotencyKey) -> bool {
279 self.store().replays.contains_key(key)
280 }
281
282 #[must_use]
287 pub fn replayed_prefix_of(&self, batch: &CommandBatch<W::Command>) -> usize {
288 let store = self.store();
289 batch
290 .envelopes
291 .iter()
292 .take_while(|envelope| store.replays.contains_key(&envelope.idempotency_key))
293 .count()
294 }
295
296 pub fn execute_prefix(
310 &self,
311 batch: &CommandBatch<W::Command>,
312 applied: usize,
313 ) -> Result<Commit<W::State, W::Event>, ExecutionError> {
314 if applied == 0 || applied > batch.envelopes.len() {
315 return Err(ExecutionError::ScopeViolation);
316 }
317 self.run(batch, applied)
318 }
319
320 fn run(
322 &self,
323 batch: &CommandBatch<W::Command>,
324 limit: usize,
325 ) -> Result<Commit<W::State, W::Event>, ExecutionError> {
326 let first = batch
327 .envelopes
328 .first()
329 .ok_or(ExecutionError::ScopeViolation)?;
330 if matches!(batch.scope, AtomicityScope::PerCase) && !batch.is_single_case() {
331 return Err(ExecutionError::ScopeViolation);
332 }
333 let key = (first.account_id().clone(), first.case_ref.case_id.clone());
334 let mut store = self.store();
335 let envelopes = &batch.envelopes[..limit.min(batch.envelopes.len())];
336
337 let mut recorded: Vec<Option<Replayable<W::State, W::Event>>> =
339 Vec::with_capacity(envelopes.len());
340 for envelope in envelopes {
341 let digest = canonical_digest(&envelope.command)
342 .map_err(|_| ExecutionError::Store(StoreError::Serialization))?;
343 match store.replays.get(&envelope.idempotency_key) {
344 None => recorded.push(None),
345 Some(entry) if entry.command_digest == digest => recorded.push(Some(entry.clone())),
346 Some(_) => {
347 return Err(ExecutionError::IdempotencyMismatch {
348 command_id: envelope.command_id,
349 });
350 }
351 }
352 }
353 let replayed = recorded.iter().take_while(|entry| entry.is_some()).count();
354 if let Some(position) = recorded[replayed..].iter().position(Option::is_some) {
355 return Err(ExecutionError::IdempotencyMismatch {
358 command_id: envelopes[replayed + position].command_id,
359 });
360 }
361
362 let current = store.cases.get(&key).cloned();
363 let current_revision = current.as_ref().map_or(CaseRevision::ZERO, |c| c.revision);
364 let conflict = |current_revision| {
365 ExecutionError::RevisionConflict(RevisionConflict {
366 expected: first.case_ref.clone(),
367 current_revision,
368 })
369 };
370
371 let last_replayed = recorded[..replayed].last().and_then(Option::as_ref);
372 let (mut state, revision_before, mut events) = match last_replayed {
373 None => {
375 if current_revision != first.case_ref.expected_revision {
376 return Err(conflict(current_revision));
377 }
378 (current.map(|c| c.value), current_revision, Vec::new())
379 }
380 Some(entry) => {
385 if entry.revision_before != first.case_ref.expected_revision
386 || current_revision != entry.revision_before.next()
387 {
388 return Err(conflict(current_revision));
389 }
390 let events = recorded[..replayed]
391 .iter()
392 .flatten()
393 .flat_map(|entry| entry.events.clone())
394 .collect();
395 (
396 Some(entry.state_after.clone()),
397 entry.revision_before,
398 events,
399 )
400 }
401 };
402
403 if replayed == envelopes.len() {
404 let Some(committed) = state else {
407 return Err(ExecutionError::ScopeViolation);
408 };
409 return Ok(Commit {
410 state: Some(committed),
411 new_revision: revision_before.next(),
412 events,
413 idempotency_replay: true,
414 });
415 }
416
417 let new_revision = revision_before.next();
418 let mut fresh = Vec::new();
419 for envelope in &envelopes[replayed..] {
420 self.definition
421 .validate_command(state.as_ref(), &envelope.command)
422 .map_err(ExecutionError::Rejected)?;
423 let applied = self
424 .definition
425 .apply(state.as_ref(), &envelope.command)
426 .map_err(ExecutionError::Rejected)?;
427 let mut committed_events = Vec::with_capacity(applied.events.len());
428 for payload in applied.events {
429 let event_type = self.definition.event_type(&payload);
430 let (sequence, occurred_at) = store.tick();
431 committed_events.push(CommittedEvent {
432 event_id: EventId::from(derive_uuid(
433 EVENT_ID_DOMAIN,
434 &[
435 key.0.as_str(),
436 key.1.as_str(),
437 &sequence.to_string(),
438 &event_type,
439 ],
440 )),
441 event_type,
442 occurred_at,
443 payload,
444 });
445 }
446 let digest = canonical_digest(&envelope.command)
447 .map_err(|_| ExecutionError::Store(StoreError::Serialization))?;
448 state = Some(applied.state.clone());
449 fresh.push((
450 envelope.idempotency_key.clone(),
451 Replayable {
452 command_digest: digest,
453 revision_before,
454 events: committed_events.clone(),
455 state_after: applied.state,
456 },
457 ));
458 events.extend(committed_events);
459 }
460
461 let Some(committed) = state else {
462 return Err(ExecutionError::ScopeViolation);
463 };
464 store.remember(&key);
465 store
466 .cases
467 .insert(key, Versioned::new(committed.clone(), new_revision));
468 for (idempotency_key, entry) in fresh {
469 store.replays.insert(idempotency_key, entry);
470 }
471 Ok(Commit {
472 state: Some(committed),
473 new_revision,
474 events,
475 idempotency_replay: replayed > 0,
479 })
480 }
481}
482
483impl<S, E> Store<S, E> {
484 fn remember(&mut self, key: &(AccountId, CaseId)) {
485 if !self.cases.contains_key(key) {
486 self.created.push(key.clone());
487 }
488 }
489
490 fn tick(&mut self) -> (u64, DateTime<Utc>) {
492 self.sequence += 1;
493 let seconds = CLOCK_EPOCH_SECONDS.saturating_add(self.sequence as i64);
494 (
495 self.sequence,
496 DateTime::from_timestamp(seconds, 0).unwrap_or_default(),
497 )
498 }
499}
500
501#[async_trait::async_trait]
502impl<W: PureWorkflow> WorkflowExecutor<W> for InMemoryExecutor<W> {
503 async fn load(
504 &self,
505 account: &AccountId,
506 case_id: &CaseId,
507 ) -> Result<Versioned<Option<W::State>>, StoreError> {
508 Ok(self
509 .store()
510 .cases
511 .get(&(account.clone(), case_id.clone()))
512 .map_or_else(
513 || Versioned::new(None, CaseRevision::ZERO),
514 |case| Versioned::new(Some(case.value.clone()), case.revision),
515 ))
516 }
517
518 async fn execute(
519 &self,
520 batch: CommandBatch<W::Command>,
521 ) -> Result<Commit<W::State, W::Event>, ExecutionError> {
522 self.run(&batch, batch.envelopes.len())
523 }
524}