Skip to main content

InMemoryExecutor

Struct InMemoryExecutor 

Source
pub struct InMemoryExecutor<W: PureWorkflow> { /* private fields */ }
Expand description

An in-memory WorkflowExecutor for any PureWorkflow.

It does the things an executor must do and nothing else:

  • it refuses a batch whose expected revision is not the current one (ExecutionError::RevisionConflict, I13);
  • it replays known idempotency keys instead of repeating the effect, and refuses a key reused with a different command (ExecutionError::IdempotencyMismatch, I14);
  • it resumes a batch that only half-executed, replaying the prefix and executing the rest;
  • it emits one CommittedEvent per event the pure apply produced.

Event identifiers and timestamps are derived, not random: two runs of the same sequence of batches produce byte-identical commits, which is what makes snapshot and replay assertions possible.

§A batch owns one revision

A batch moves the case from expected_revision to expected_revision + 1 however many envelopes it carries, and a resumed batch lands on the same revision the interrupted one reached. That is what makes execute_prefix followed by execute indistinguishable, from the case’s point of view, from one uninterrupted execute — which is the property crash recovery relies on.

Implementations§

Source§

impl<W: PureWorkflow> InMemoryExecutor<W>

Source

pub fn new(definition: W) -> Self

Builds an empty executor for definition.

Source

pub const fn definition(&self) -> &W

The definition transitions are validated against.

Source

pub fn seed( &self, account: &AccountId, case_id: &CaseId, state: W::State, revision: CaseRevision, )

Installs a case at a chosen revision, bypassing commands. Use it to start a test from a state that would take many turns to reach.

Source

pub fn case_ids(&self, account: &AccountId) -> Vec<CaseId>

The account’s cases, in the order each first appeared: what a directory over this executor lists.

Source

pub fn revision_of(&self, account: &AccountId, case_id: &CaseId) -> CaseRevision

The current revision of a case, or CaseRevision::ZERO when it does not exist for this account.

Source

pub fn state_of( &self, account: &AccountId, case_id: &CaseId, ) -> Option<W::State>

The state of a case, for a test that asserts on what a turn wrote.

Source

pub fn case_count(&self) -> usize

Number of distinct cases stored.

Source

pub fn has_executed(&self, key: &IdempotencyKey) -> bool

Returns true when this idempotency key has already executed.

Source

pub fn replayed_prefix_of(&self, batch: &CommandBatch<W::Command>) -> usize

How many envelopes at the front of batch already executed.

0 means the batch is untouched, batch.envelopes.len() that the whole batch is a replay, anything between that it was interrupted.

Source

pub fn execute_prefix( &self, batch: &CommandBatch<W::Command>, applied: usize, ) -> Result<Commit<W::State, W::Event>, ExecutionError>

Executes only the first applied envelopes of batch, as a process that died mid-batch would have left it.

The case moves to expected_revision + 1 and the applied envelopes are remembered, so executing the whole batch afterwards replays them and runs only the rest. This is how a test creates a partially replayed batch; nothing else in the kit produces one.

§Errors

ExecutionError::ScopeViolation when applied is zero or larger than the batch, plus everything execute can return.

Trait Implementations§

Source§

impl<W: PureWorkflow> Debug for InMemoryExecutor<W>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<W: PureWorkflow + Default> Default for InMemoryExecutor<W>

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl<W: PureWorkflow> WorkflowExecutor<W> for InMemoryExecutor<W>

Source§

fn load<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, account: &'life1 AccountId, case_id: &'life2 CaseId, ) -> Pin<Box<dyn Future<Output = Result<Versioned<Option<W::State>>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Loads a case for an account. None value means the case does not exist; its revision is then crate::ids::CaseRevision::ZERO. An executor must never delete a case to express completion, because a completed case keeps its row and moves to a terminal status, so an absent state always means the case has not been created yet.
Source§

fn execute<'life0, 'async_trait>( &'life0 self, batch: CommandBatch<W::Command>, ) -> Pin<Box<dyn Future<Output = Result<Commit<W::State, W::Event>, ExecutionError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Executes a batch under its atomicity scope with revision and idempotency checks (I13, I14).

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<W, E> CaseLoader<W> for E

Source§

fn load_case<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, account: &'life1 AccountId, case_id: &'life2 CaseId, ) -> Pin<Box<dyn Future<Output = Result<Versioned<Option<<W as WorkflowDefinition>::State>>, StoreError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, E: 'async_trait,

Loads a case for an account. None value means the case does not exist; its revision is then crate::ids::CaseRevision::ZERO. See WorkflowExecutor::load for what an absent state does and does not mean. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more