Skip to main content

turnframe_test/workflows/
mod.rs

1//! Complete sample domains, and the pieces they share: a travel-disruption desk.
2//!
3//! Three workflows ship with the kit and are meant to be reused by every other
4//! crate, by the examples and by integration tests:
5//!
6//! * [`trip`]: the disruption case of one booking, with obligations open at once,
7//!   one of them per extra; a rebooking card bound to the case revision and to the
8//!   quote it shows; an airline that may answer late or never, whose states are
9//!   never collapsed into "done"; and a leg the domain refuses to change once the
10//!   traveler asked to keep it;
11//! * [`traveler`]: the flat onboarding slice, where deleting is destructive and a
12//!   new email is a sensitive change behind a review card, and where one field is
13//!   three-valued (untouched, answered, declined) because a two-valued field cannot
14//!   tell "we have not asked" from "they said no";
15//! * [`claim`]: a receipt arrives, values are read from it and held as a proposal,
16//!   and a review card turns the proposal into state: the worked recipe for
17//!   *proposed values awaiting review*, which is domain state rather than framework
18//!   vocabulary.
19//!
20//! All three render a receipt for a
21//! [`ReceiptEvent::Redacted`](turnframe_core::event::ReceiptEvent), because a turn
22//! that quietly drops one reads as a turn in which nothing happened. The copy says
23//! the least any of them could honestly say: the step is on record, its detail was
24//! erased. The receipt still cites the event, so it still passes the claim guard,
25//! which `tests/claim_guard.rs` pins.
26//!
27//! Each implements [`PureWorkflow`], so the same pure `apply` drives the
28//! [`InMemoryExecutor`] of the integration tests and the
29//! [`WorkflowModel`](crate::explore::WorkflowModel) of exploration.
30
31pub 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/// The state and the events one command produced.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct Applied<S, E> {
53    /// State after the command.
54    pub state: S,
55    /// Events the executor commits, in order.
56    pub events: Vec<E>,
57}
58
59impl<S, E> Applied<S, E> {
60    /// Pairs a state with its events.
61    #[must_use]
62    pub const fn new(state: S, events: Vec<E>) -> Self {
63        Self { state, events }
64    }
65}
66
67/// A workflow whose transitions are one pure function.
68///
69/// Splitting this out of [`WorkflowDefinition`] keeps the definition free of
70/// execution concerns while letting the kit derive both an executor and an
71/// exploration model from a single description of what a command does.
72pub trait PureWorkflow: WorkflowDefinition {
73    /// Applies one command to one state, refusing it the way the real domain
74    /// would. Must be deterministic and must not mutate `state`.
75    fn apply(
76        &self,
77        state: Option<&Self::State>,
78        command: &Self::Command,
79    ) -> Result<Applied<Self::State, Self::Event>, DomainRejection>;
80
81    /// Stable event type label of an event, e.g. `"trip.extra_added"`.
82    fn event_type(&self, event: &Self::Event) -> String;
83}
84
85/// The operations with the Italian summary `summaries` gives each by key, for a turn in
86/// Italian: a deployment writes its summaries in every language it serves.
87pub(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
102/// Turns [`PureWorkflow::apply`] into a [`SimulatedTransition`], so a model
103/// only has to describe its initial states and candidate commands.
104pub 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
115/// Domain separation of the derived event identifiers.
116const EVENT_ID_DOMAIN: &str = "turnframe.test.event.v1";
117
118/// First second of the executor's synthetic clock (2023-11-14T22:13:20Z).
119const CLOCK_EPOCH_SECONDS: i64 = 1_700_000_000;
120
121/// What one already-executed envelope remembers.
122///
123/// The unit is the **envelope**, not the batch. Keying the memory on the batch
124/// would make a batch that half-executed unrepresentable: its first command
125/// really did commit, and pretending otherwise is exactly the crash-recovery
126/// bug the executor exists to model (spec §23.1).
127#[derive(Clone)]
128struct Replayable<S, E> {
129    /// Digest of *this envelope's* command, so a key reused with a different
130    /// command is a mismatch (I14).
131    command_digest: Digest,
132    /// Revision the case was at when the batch that contains this envelope
133    /// started.
134    revision_before: CaseRevision,
135    /// Events this envelope committed.
136    events: Vec<CommittedEvent<E>>,
137    /// State after this envelope.
138    state_after: S,
139}
140
141struct Store<S, E> {
142    cases: HashMap<(AccountId, CaseId), Versioned<S>>,
143    /// Every case, in the order it first appeared.
144    created: Vec<(AccountId, CaseId)>,
145    replays: HashMap<IdempotencyKey, Replayable<S, E>>,
146    sequence: u64,
147}
148
149/// An in-memory [`WorkflowExecutor`] for any [`PureWorkflow`].
150///
151/// It does the things an executor must do and nothing else:
152///
153/// * it refuses a batch whose expected revision is not the current one
154///   ([`ExecutionError::RevisionConflict`], I13);
155/// * it replays known idempotency keys instead of repeating the effect, and
156///   refuses a key reused with a different command
157///   ([`ExecutionError::IdempotencyMismatch`], I14);
158/// * it resumes a batch that only half-executed, replaying the prefix and
159///   executing the rest;
160/// * it emits one [`CommittedEvent`] per event the pure `apply` produced.
161///
162/// Event identifiers and timestamps are derived, not random: two runs of the
163/// same sequence of batches produce byte-identical commits, which is what makes
164/// snapshot and replay assertions possible.
165///
166/// # A batch owns one revision
167///
168/// A batch moves the case from `expected_revision` to `expected_revision + 1`
169/// however many envelopes it carries, and a resumed batch lands on the *same*
170/// revision the interrupted one reached. That is what makes
171/// [`execute_prefix`](Self::execute_prefix) followed by
172/// [`execute`](WorkflowExecutor::execute) indistinguishable, from the case's
173/// point of view, from one uninterrupted `execute` — which is the property
174/// crash recovery relies on.
175pub 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    /// Builds an empty executor for `definition`.
197    #[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    /// The definition transitions are validated against.
211    #[must_use]
212    pub const fn definition(&self) -> &W {
213        &self.definition
214    }
215
216    /// Locks the store, recovering from a poisoned mutex rather than panicking:
217    /// a test that already failed must not cascade into unrelated failures.
218    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    /// Installs a case at a chosen revision, bypassing commands. Use it to
225    /// start a test from a state that would take many turns to reach.
226    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    /// The account's cases, in the order each first appeared: what a directory over
240    /// this executor lists.
241    #[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    /// The current revision of a case, or [`CaseRevision::ZERO`] when it does
252    /// not exist for this account.
253    #[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    /// The state of a case, for a test that asserts on what a turn wrote.
262    #[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    /// Number of distinct cases stored.
271    #[must_use]
272    pub fn case_count(&self) -> usize {
273        self.store().cases.len()
274    }
275
276    /// Returns `true` when this idempotency key has already executed.
277    #[must_use]
278    pub fn has_executed(&self, key: &IdempotencyKey) -> bool {
279        self.store().replays.contains_key(key)
280    }
281
282    /// How many envelopes at the front of `batch` already executed.
283    ///
284    /// `0` means the batch is untouched, `batch.envelopes.len()` that the whole
285    /// batch is a replay, anything between that it was interrupted.
286    #[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    /// Executes only the first `applied` envelopes of `batch`, as a process
297    /// that died mid-batch would have left it.
298    ///
299    /// The case moves to `expected_revision + 1` and the applied envelopes are
300    /// remembered, so executing the whole batch afterwards replays them and
301    /// runs only the rest. This is how a test *creates* a partially replayed
302    /// batch; nothing else in the kit produces one.
303    ///
304    /// # Errors
305    ///
306    /// [`ExecutionError::ScopeViolation`] when `applied` is zero or larger than
307    /// the batch, plus everything [`execute`](WorkflowExecutor::execute) can
308    /// return.
309    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    /// The whole batch, or as much of it as `limit` allows.
321    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        // What of this batch already executed, and does it still mean the same?
338        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            // A key from the middle of the batch executed while an earlier one
356            // did not: these envelopes never travelled together.
357            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            // Nothing of this batch ran yet: the ordinary optimistic check.
374            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            // The batch was interrupted. Its own revision is the one to check
381            // against, and the case must still be where the interrupted batch
382            // left it — otherwise something else wrote in between and resuming
383            // would silently overwrite it.
384            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            // Every envelope is a replay: return the original outcome without
405            // repeating a single effect (I14).
406            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            // Part of the batch came back from the journal rather than from the
476            // domain, which is what a caller has to know before it renders a
477            // receipt for "what just happened".
478            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    /// The next synthetic instant, one second after the previous one.
491    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}