Skip to main content

Runtime

Struct Runtime 

Source
pub struct Runtime<S> { /* private fields */ }
Expand description

Executes effects against a store. Cheap to clone; clones share state.

Implementations§

Source§

impl<S: EffectStore> Runtime<S>

Source

pub fn compensation( &self, name: impl Into<String>, key: impl Display, ) -> CompensationBuilder<'_, S>

Starts undoing the closure effect (name, key). See CompensationBuilder::run.

Source§

impl<S: EffectStore> Runtime<S>

Source

pub async fn recover(&self) -> Result<RecoveryReport, RuntimeError>

One recovery pass, in two steps.

  1. Every effect whose worker’s lease expired mid-attempt or mid-verification becomes Unknown, never Failed: it may have changed the outside world.
  2. Every unsettled effect nobody holds (see Self::pending) whose name has a registered handler is finished from its stored input, exactly as if its caller had called again: verified, re-run if that is safe, or escalated. This includes effects left Pending by a crash. Effects waiting for an operator, and retries scheduled for later, are left alone. Unsettled closure effects are reported as unhandled: only a caller can re-run them.

Effects are resumed one at a time, so a pass lasts as long as their retries and verifications take. Safe to run from several workers at once: every change happens under a lease.

§Errors

RuntimeError::Store if the store fails. Progress made before the failure is kept. A handler that cannot run is reported in resume_errors instead.

Source

pub async fn run_recovery(&self, interval: Duration)

Runs Self::recover every interval, forever, followed by Self::prune when the runtime has a retention policy. Spawn it:

ⓘ
tokio::spawn({
    let runtime = runtime.clone();
    async move { runtime.run_recovery(Duration::from_secs(30)).await }
});

A failed pass is logged and retried at the next tick.

Source

pub async fn pending( &self, after: Option<EffectId>, limit: usize, ) -> Result<Vec<EffectRecord>, RuntimeError>

Effects waiting for a caller to re-run them or an operator to decide, ordered by creation, limit at a time after after.

These are the effects that are unsettled with nobody working on them: Pending (for example a worker died during a backoff wait), Executing or Verifying with an expired lease (not yet recovered), Unknown, NeedsIntervention, AwaitingApproval, an interrupted Compensating, and CompensationFailed.

§Errors

RuntimeError::Store if the store fails.

Source

pub async fn resolve( &self, id: EffectId, resolution: Resolution, actor: impl Into<String>, note: impl Into<String>, ) -> Result<EffectRecord, RuntimeError>

Records an operator’s decision about an effect that is Unknown or NeedsIntervention. actor identifies the operator (e.g. operator:alice); note says why, for the audit trail.

§Errors

RuntimeError::Store wrapping:

Source

pub async fn approve( &self, id: EffectId, actor: impl Into<String>, note: impl Into<String>, ) -> Result<EffectRecord, RuntimeError>

Approves an effect waiting in AwaitingApproval. It becomes Pending: a caller’s next call runs it, and so does recovery for a registered handler. actor is recorded as the approver.

§Errors

As for Self::resolve; InvalidTransition if it is not awaiting approval.

Source

pub async fn deny( &self, id: EffectId, actor: impl Into<String>, reason: impl Into<String>, ) -> Result<EffectRecord, RuntimeError>

Denies an effect waiting in AwaitingApproval. It ends Rejected, with reason as its error.

§Errors

As for Self::approve.

Source§

impl<S: EffectStore> Runtime<S>

Source

pub async fn prune(&self) -> Result<PruneReport, RuntimeError>

Deletes the settled records the retention policy says are old enough, in batches, with their audit trails. Does nothing without a policy. Self::run_recovery calls it every round; call it yourself to prune on another schedule.

§Errors

RuntimeError::Store if the store fails. Batches deleted before the failure stay deleted.

Source§

impl<S: EffectStore> Runtime<S>

Source

pub fn new(store: S) -> Self

A runtime with default settings.

Source

pub fn builder(store: S) -> RuntimeBuilder<S>

Starts configuring a runtime.

Source

pub fn effect( &self, name: impl Into<String>, key: impl Display, ) -> EffectBuilder<S>

Describes an effect: name is its type (e.g. payment.charge), key identifies this occurrence (e.g. the order id). Every call with the same name and key refers to the same effect.

Source

pub fn submit<H: EffectHandler>( &self, key: impl Display, input: H::Input, ) -> Submission<'_, S, H>

Runs a registered handler’s effect with input, or attaches to an earlier run with the same key. Await the returned Submission.

Behaves like EffectBuilder::run, with the handler’s properties. The input is stored in full, so recovery can finish the effect if this process dies.

Source

pub fn compensate<H: EffectHandler>( &self, key: impl Display, ) -> CompensationSubmission<'_, S, H>

Undoes a registered, compensable handler’s effect. Await the returned submission. See compensation.

Source

pub fn store(&self) -> &S

The underlying store.

Source

pub fn worker_id(&self) -> &WorkerId

This runtime’s lease-holder identity.

Source

pub async fn wait<T: DeserializeOwned>( &self, id: EffectId, timeout: Duration, ) -> Result<EffectOutcome<T>, RuntimeError>

Waits until nobody is working on effect id, or timeout passes, and reports where it stands. For a caller that got EffectOutcome::InProgress.

Returns InProgress if the effect is still being worked on at the deadline, or if it is unsettled and nobody holds it (for example a worker crashed while waiting to retry). Running the effect again then takes it over.

§Errors

RuntimeError::Store if the store fails or has no such effect, and RuntimeError::Output if a committed output does not deserialize into T.

Trait Implementations§

Source§

impl<S> Clone for Runtime<S>

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more

Auto Trait Implementations§

§

impl<S> !RefUnwindSafe for Runtime<S>

§

impl<S> !UnwindSafe for Runtime<S>

§

impl<S> Freeze for Runtime<S>
where Arc<Inner<S>>: Freeze,

§

impl<S> Send for Runtime<S>
where Arc<Inner<S>>: Send,

§

impl<S> Sync for Runtime<S>
where Arc<Inner<S>>: Sync,

§

impl<S> Unpin for Runtime<S>
where Arc<Inner<S>>: Unpin,

§

impl<S> UnsafeUnpin for Runtime<S>
where Arc<Inner<S>>: UnsafeUnpin,

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<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
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> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
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<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