Skip to main content

OutboxDispatcher

Struct OutboxDispatcher 

Source
pub struct OutboxDispatcher { /* private fields */ }
Expand description

Claims due outbox rows, sends them and settles each one.

Implementations§

Source§

impl OutboxDispatcher

Source

pub fn new( outbox: Arc<dyn OutboxStore>, sender: Arc<dyn OutboxSender>, config: DispatchConfig, ) -> Self

Builds a dispatcher over outbox, sending through sender.

Source

pub fn with_observer(self, observer: Arc<dyn Observer>) -> Self

Sends this module’s signals to observer (spec §26.2, §28).

The dispatcher is driven by the application’s own task rather than by the orchestrator, so it is given its observer here rather than inheriting one. Without it the external half of the saga is invisible: ExternalLatency is the only measure of how long the remote system takes, and ExternalReconciled is the other end of ExternalOutcomeUnknown — a rising count of unknowns with no reconciliations behind it is the shape of an operator who has stopped settling them.

Source

pub const fn config(&self) -> &DispatchConfig

The configuration in force.

Source

pub async fn run_once( &self, now: DateTime<Utc>, ) -> Result<DispatchReport, StoreError>

Claims the rows due at now, sends each and settles it.

This is the unit of work an application’s own task calls. It returns when every claimed row has been settled — completed, rescheduled, failed or handed to reconciliation — so a caller that awaits it knows exactly what happened.

§Errors

StoreError when the claim itself could not be made. A row that could not be settled is not an error: it is reported in DispatchReport::unsettled, and the store’s claim timeout will release it for another sweep.

Source

pub async fn reconcile( &self, outbox_id: &OutboxId, reconciler: &dyn OutboxReconciler, now: DateTime<Utc>, ) -> Result<Reconciled, StoreError>

Settles one row whose outcome is unknown, through reconciler (spec §16.5).

The row is read, handed to the reconciler and settled with its answer. A row that is not in OutboxStatus::OutcomeUnknown is left alone and reported as Reconciled::Unresolved: reconciliation is for the rows nobody knows about, and a completed row is not one of them.

§Errors

StoreError when the row could not be read or the settlement was refused.

Source

pub async fn release_expired_claims( &self, claimed_before: DateTime<Utc>, ) -> Result<Vec<OutboxId>, StoreError>

Releases rows a crashed dispatcher left in Dispatching, so another sweep can claim them.

claimed_before is the age at which a claim is considered abandoned; it must be older than the longest send this dispatcher can make, or a slow send is released while it is still running.

§Errors

StoreError when the sweep could not be made.

Trait Implementations§

Source§

impl Clone for OutboxDispatcher

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
Source§

impl Debug for OutboxDispatcher

Source§

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

Formats the value using the given formatter. Read more

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

Source§

fn __clone_box(&self, _: Private) -> *mut ()

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> 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