pub struct OutboxDispatcher { /* private fields */ }Expand description
Claims due outbox rows, sends them and settles each one.
Implementations§
Source§impl OutboxDispatcher
impl OutboxDispatcher
Sourcepub fn new(
outbox: Arc<dyn OutboxStore>,
sender: Arc<dyn OutboxSender>,
config: DispatchConfig,
) -> Self
pub fn new( outbox: Arc<dyn OutboxStore>, sender: Arc<dyn OutboxSender>, config: DispatchConfig, ) -> Self
Builds a dispatcher over outbox, sending through sender.
Sourcepub fn with_observer(self, observer: Arc<dyn Observer>) -> Self
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.
Sourcepub const fn config(&self) -> &DispatchConfig
pub const fn config(&self) -> &DispatchConfig
The configuration in force.
Sourcepub async fn run_once(
&self,
now: DateTime<Utc>,
) -> Result<DispatchReport, StoreError>
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.
Sourcepub async fn reconcile(
&self,
outbox_id: &OutboxId,
reconciler: &dyn OutboxReconciler,
now: DateTime<Utc>,
) -> Result<Reconciled, StoreError>
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.
Sourcepub async fn release_expired_claims(
&self,
claimed_before: DateTime<Utc>,
) -> Result<Vec<OutboxId>, StoreError>
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.