Skip to main content

turnframe_runtime/
dispatch.rs

1//! The outbox dispatcher: the second half of the external-effect saga
2//! (spec §16.4, §16.5, ADR-007).
3//!
4//! [`execute`](crate::execute) enqueues an outbox row inside the turn's one
5//! atomic write, and stops there — deliberately, because the row must be
6//! durable before anything leaves the building. Somebody then has to pick the
7//! row up and call the remote system. That is this module: a reference
8//! dispatcher an adopter can use as it stands, or read and replace.
9//!
10//! # What it is not
11//!
12//! **It is not a background thread.** Nothing here spawns anything. The library
13//! never starts work an application did not ask for: [`OutboxDispatcher`] has
14//! one method that does a unit of work, [`OutboxDispatcher::run_once`], and the
15//! application drives it from its own task, its own scheduler or its own
16//! cron job. A hidden worker would dispatch external effects out of a process
17//! that was only supposed to answer a turn, and would keep doing it while the
18//! operator was shutting the process down.
19//!
20//! ```rust,ignore
21//! // The application's own task. Cancel it, pause it, scale it — it is yours.
22//! let mut ticker = tokio::time::interval(Duration::from_secs(1));
23//! loop {
24//!     tokio::select! {
25//!         _ = shutdown.cancelled() => break,
26//!         _ = ticker.tick() => {
27//!             match dispatcher.run_once(Utc::now()).await {
28//!                 Ok(report) => tracing::debug!(dispatched = report.claimed),
29//!                 Err(error) => tracing::warn!(%error, "the outbox could not be read"),
30//!             }
31//!         }
32//!     }
33//! }
34//! ```
35//!
36//! # The four ways one row ends
37//!
38//! [`OutboxSender::send`] classifies its own outcome, and the classification is
39//! the whole safety contract of the module:
40//!
41//! | [`Dispatched`] | What the row becomes | Why |
42//! | --- | --- | --- |
43//! | [`Completed`](Dispatched::Completed) | `Completed` | the remote confirmed |
44//! | [`Retryable`](Dispatched::Retryable) | `Pending` with a backoff, or `Failed` once the attempts are spent | the request demonstrably did not arrive |
45//! | [`Permanent`](Dispatched::Permanent) | `Failed` | the remote refused, and will refuse again |
46//! | [`Unknown`](Dispatched::Unknown) | `OutcomeUnknown` | the effect **may** exist, and repeating it is the duplicate the library exists to prevent (I15) |
47//!
48//! A send that does not answer within
49//! [`DispatchConfig::send_timeout`] is [`Unknown`](Dispatched::Unknown), never a
50//! retry. That is the same rule the executor applies to a domain timeout, for
51//! the same reason: the request left, so nobody can say it did not land.
52//!
53//! # Exclusivity is the store's, and this module honours it
54//!
55//! [`OutboxWriter::claim_due`](turnframe_store::outbox::OutboxWriter::claim_due) moves rows to `Dispatching` under a worker
56//! identifier with skip-locked semantics, so two dispatchers claim disjoint
57//! sets. This module never bypasses it — it dispatches exactly what a claim
58//! returned — and it never invents an idempotency key: the one the command was
59//! admitted under travels on the row and is handed to the sender, so a remote
60//! that deduplicates can.
61//!
62//! # Unknown outcomes are reconciled, not retried
63//!
64//! A row in `OutcomeUnknown` is out of the dispatcher's hands: only the
65//! application knows how to ask the remote system what happened.
66//! [`OutboxDispatcher::reconcile`] is the hook — it hands the stored record to
67//! an [`OutboxReconciler`] and settles the row with the answer, including
68//! putting it back in the queue when the remote is certain it never arrived.
69
70use std::fmt;
71use std::sync::Arc;
72use std::time::Duration;
73
74use async_trait::async_trait;
75use chrono::{DateTime, Utc};
76use turnframe_core::error::StoreError;
77use turnframe_core::event::{OutboxEntry, OutboxStatus};
78use turnframe_core::ids::OutboxId;
79use turnframe_core::observe::{NoopObserver, Observer, Signal, SignalLabels};
80use turnframe_store::outbox::{OutboxRecord, OutboxStore};
81
82use crate::signals::Stage;
83
84/// Stable codes this module records on a row it settled itself.
85///
86/// There is exactly one: every other settlement carries a code its author
87/// chose, and inventing a second vocabulary next to theirs would only make a
88/// dashboard harder to read.
89pub mod code {
90    /// The retry budget of the row ran out, so a retryable failure became a
91    /// permanent one.
92    pub const ATTEMPTS_EXHAUSTED: &str = "turnframe.dispatch.attempts_exhausted";
93}
94
95/// How one send ended, as the sender classifies it.
96///
97/// There is no `Result` around it on purpose. Every failure mode of an external
98/// call has to land in exactly one of these four, and an author who returns an
99/// error instead of choosing has not answered the only question that matters:
100/// *may this be sent again?*
101#[derive(Debug, Clone, PartialEq, Eq)]
102#[non_exhaustive]
103pub enum Dispatched {
104    /// The remote system accepted it, and said so.
105    Completed {
106        /// Reference the remote gave, when it gave one.
107        remote_ref: Option<String>,
108    },
109    /// It did not arrive, and sending it again is safe.
110    Retryable {
111        /// Stable code, never free text.
112        code: String,
113    },
114    /// The remote refused it, and would refuse it again.
115    Permanent {
116        /// Stable code, never free text.
117        code: String,
118    },
119    /// The request left and no answer came back. The effect may exist (I15).
120    Unknown {
121        /// Reference the remote gave before it went quiet, when it gave one.
122        remote_ref: Option<String>,
123    },
124}
125
126impl Dispatched {
127    /// The remote accepted it, with no reference.
128    #[must_use]
129    pub const fn completed() -> Self {
130        Self::Completed { remote_ref: None }
131    }
132
133    /// Stable snake-case label of the outcome.
134    #[must_use]
135    pub const fn as_str(&self) -> &'static str {
136        match self {
137            Self::Completed { .. } => "completed",
138            Self::Retryable { .. } => "retryable",
139            Self::Permanent { .. } => "permanent",
140            Self::Unknown { .. } => "unknown",
141        }
142    }
143}
144
145/// Calls the external system for one outbox row.
146///
147/// The implementation owns the transport and the credentials; this module owns
148/// the bookkeeping. Two obligations are the application's:
149///
150/// * **forward [`OutboxEntry::idempotency_key`]** to the remote, in whatever
151///   header or field it deduplicates on. It is the same key the command was
152///   admitted under, so a row dispatched twice after a crash is one effect;
153/// * **classify honestly.** A transport error after the bytes were written is
154///   [`Dispatched::Unknown`], not [`Dispatched::Retryable`]. If you cannot tell
155///   the two apart, it is `Unknown`.
156#[async_trait]
157pub trait OutboxSender: Send + Sync {
158    /// Sends one row.
159    async fn send(&self, entry: &OutboxEntry) -> Dispatched;
160}
161
162/// What a reconciler found out about a row whose outcome was unknown.
163#[derive(Debug, Clone, PartialEq, Eq)]
164#[non_exhaustive]
165pub enum Reconciled {
166    /// The remote has it. The row is `Completed`.
167    Completed,
168    /// The remote definitively does not have it, and never will. The row is
169    /// `Failed`.
170    Failed {
171        /// Stable code, never free text.
172        code: String,
173    },
174    /// The remote definitively never received it, so it may be queued again.
175    /// Only answer this when the remote is *certain*: it is the one path that
176    /// can turn an unknown outcome back into a second send.
177    Resend,
178    /// Still unknown. The row is left exactly as it is, for the next sweep.
179    Unresolved,
180}
181
182impl Reconciled {
183    /// Stable snake-case label of the answer, for a log line or a metric
184    /// dimension. Never carries the code an author chose.
185    #[must_use]
186    pub const fn as_str(&self) -> &'static str {
187        match self {
188            Self::Completed => "completed",
189            Self::Failed { .. } => "failed",
190            Self::Resend => "resend",
191            Self::Unresolved => "unresolved",
192        }
193    }
194}
195
196/// Asks the external system what happened to a row whose outcome is unknown
197/// (spec §16.5).
198#[async_trait]
199pub trait OutboxReconciler: Send + Sync {
200    /// Settles one row against the remote system.
201    async fn reconcile(&self, record: &OutboxRecord) -> Reconciled;
202}
203
204/// How the dispatcher works (all of it optional, all of it conservative).
205#[derive(Debug, Clone, PartialEq, Eq)]
206#[non_exhaustive]
207pub struct DispatchConfig {
208    /// Identifier this dispatcher claims rows under. Distinct per process.
209    pub worker_id: String,
210    /// Rows claimed per [`OutboxDispatcher::run_once`].
211    pub batch_size: usize,
212    /// Deadline for one [`OutboxSender::send`]. Exceeding it is an unknown
213    /// outcome, not a retry.
214    pub send_timeout: Duration,
215    /// Delay before the second attempt of a row.
216    pub initial_backoff: Duration,
217    /// Ceiling on the computed delay.
218    pub max_backoff: Duration,
219    /// Multiplier applied per further attempt.
220    pub backoff_multiplier: u32,
221    /// Attempts a row gets before a retryable failure becomes a permanent one.
222    pub max_attempts: u32,
223}
224
225impl DispatchConfig {
226    /// Sixteen rows a sweep, thirty seconds a send, one second of backoff
227    /// growing by four up to a minute, five attempts.
228    #[must_use]
229    pub fn new(worker_id: impl Into<String>) -> Self {
230        Self {
231            worker_id: worker_id.into(),
232            batch_size: 16,
233            send_timeout: Duration::from_secs(30),
234            initial_backoff: Duration::from_secs(1),
235            max_backoff: Duration::from_secs(60),
236            backoff_multiplier: 4,
237            max_attempts: 5,
238        }
239    }
240
241    /// Returns a copy claiming at most `batch_size` rows a sweep.
242    #[must_use]
243    pub fn with_batch_size(mut self, batch_size: usize) -> Self {
244        self.batch_size = batch_size;
245        self
246    }
247
248    /// Returns a copy with another send deadline.
249    #[must_use]
250    pub const fn with_send_timeout(mut self, send_timeout: Duration) -> Self {
251        self.send_timeout = send_timeout;
252        self
253    }
254
255    /// Returns a copy with another attempt budget.
256    #[must_use]
257    pub const fn with_max_attempts(mut self, max_attempts: u32) -> Self {
258        self.max_attempts = max_attempts;
259        self
260    }
261
262    /// Returns a copy with another backoff schedule.
263    #[must_use]
264    pub const fn with_backoff(mut self, initial: Duration, max: Duration) -> Self {
265        self.initial_backoff = initial;
266        self.max_backoff = max;
267        self
268    }
269
270    /// The delay before attempt number `attempt`, 1-based.
271    #[must_use]
272    pub fn backoff_for(&self, attempt: u32) -> Duration {
273        let step = attempt.saturating_sub(1);
274        if step == 0 {
275            return self.initial_backoff;
276        }
277        self.initial_backoff
278            .saturating_mul(self.backoff_multiplier.saturating_pow(step.min(16)))
279            .min(self.max_backoff)
280    }
281}
282
283/// What one sweep did.
284#[derive(Debug, Clone, Default, PartialEq, Eq)]
285#[non_exhaustive]
286pub struct DispatchReport {
287    /// Rows completed by the remote.
288    pub completed: Vec<OutboxId>,
289    /// Rows put back in the queue with a backoff.
290    pub retried: Vec<OutboxId>,
291    /// Rows the remote refused, or whose attempts ran out.
292    pub failed: Vec<OutboxId>,
293    /// Rows whose outcome nobody knows yet; a reconciler must settle them.
294    pub unknown: Vec<OutboxId>,
295    /// Rows the store refused to settle, left as the store has them.
296    pub unsettled: Vec<OutboxId>,
297}
298
299impl DispatchReport {
300    /// How many rows the sweep claimed.
301    #[must_use]
302    pub fn claimed(&self) -> usize {
303        self.completed.len()
304            + self.retried.len()
305            + self.failed.len()
306            + self.unknown.len()
307            + self.unsettled.len()
308    }
309
310    /// Returns `true` when the sweep found nothing due.
311    #[must_use]
312    pub fn is_empty(&self) -> bool {
313        self.claimed() == 0
314    }
315}
316
317/// Claims due outbox rows, sends them and settles each one.
318#[derive(Clone)]
319pub struct OutboxDispatcher {
320    outbox: Arc<dyn OutboxStore>,
321    sender: Arc<dyn OutboxSender>,
322    config: DispatchConfig,
323    observer: Arc<dyn Observer>,
324}
325
326impl fmt::Debug for OutboxDispatcher {
327    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
328        f.debug_struct("OutboxDispatcher")
329            .field("config", &self.config)
330            .finish_non_exhaustive()
331    }
332}
333
334impl OutboxDispatcher {
335    /// Builds a dispatcher over `outbox`, sending through `sender`.
336    #[must_use]
337    pub fn new(
338        outbox: Arc<dyn OutboxStore>,
339        sender: Arc<dyn OutboxSender>,
340        config: DispatchConfig,
341    ) -> Self {
342        Self {
343            outbox,
344            sender,
345            config,
346            observer: Arc::new(NoopObserver),
347        }
348    }
349
350    /// Sends this module's signals to `observer` (spec §26.2, §28).
351    ///
352    /// The dispatcher is driven by the application's own task rather than by
353    /// the orchestrator, so it is given its observer here rather than
354    /// inheriting one. Without it the external half of the saga is invisible:
355    /// [`ExternalLatency`](Signal::ExternalLatency) is the only measure of how
356    /// long the remote system takes, and
357    /// [`ExternalReconciled`](Signal::ExternalReconciled) is the other end of
358    /// [`ExternalOutcomeUnknown`](Signal::ExternalOutcomeUnknown) — a rising
359    /// count of unknowns with no reconciliations behind it is the shape of an
360    /// operator who has stopped settling them.
361    #[must_use]
362    pub fn with_observer(mut self, observer: Arc<dyn Observer>) -> Self {
363        self.observer = observer;
364        self
365    }
366
367    /// The configuration in force.
368    #[must_use]
369    pub const fn config(&self) -> &DispatchConfig {
370        &self.config
371    }
372
373    /// Claims the rows due at `now`, sends each and settles it.
374    ///
375    /// This is the unit of work an application's own task calls. It returns
376    /// when every claimed row has been settled — completed, rescheduled, failed
377    /// or handed to reconciliation — so a caller that awaits it knows exactly
378    /// what happened.
379    ///
380    /// # Errors
381    ///
382    /// [`StoreError`] when the claim itself could not be made. A row that could
383    /// not be *settled* is not an error: it is reported in
384    /// [`DispatchReport::unsettled`], and the store's claim timeout will release
385    /// it for another sweep.
386    pub async fn run_once(&self, now: DateTime<Utc>) -> Result<DispatchReport, StoreError> {
387        let claimed = self
388            .outbox
389            .claim_due(now, self.config.batch_size, &self.config.worker_id)
390            .await?;
391        let mut report = DispatchReport::default();
392        for entry in claimed {
393            self.dispatch_one(&entry, now, &mut report).await;
394        }
395        Ok(report)
396    }
397
398    /// Sends one claimed row and settles it.
399    async fn dispatch_one(
400        &self,
401        entry: &OutboxEntry,
402        now: DateTime<Utc>,
403        report: &mut DispatchReport,
404    ) {
405        let outcome = self.send(entry).await;
406        tracing::debug!(
407            target: "turnframe.dispatch",
408            outbox_id = %entry.outbox_id,
409            destination = entry.destination.as_str(),
410            attempt = entry.attempt_count,
411            outcome = outcome.as_str(),
412            "outbox row dispatched"
413        );
414        let settled = match &outcome {
415            Dispatched::Completed { .. } => self.outbox.mark_completed(&entry.outbox_id).await,
416            Dispatched::Unknown { remote_ref } => {
417                self.outbox
418                    .mark_outcome_unknown(&entry.outbox_id, remote_ref.clone())
419                    .await
420            }
421            Dispatched::Permanent { code } => {
422                self.outbox
423                    .mark_failed(&entry.outbox_id, code.clone(), None)
424                    .await
425            }
426            Dispatched::Retryable { code } => {
427                if entry.attempt_count >= self.config.max_attempts {
428                    self.outbox
429                        .mark_failed(&entry.outbox_id, code::ATTEMPTS_EXHAUSTED.to_owned(), None)
430                        .await
431                } else {
432                    let delay = self.config.backoff_for(entry.attempt_count);
433                    // A schedule so long that chrono refuses it is a
434                    // misconfiguration, not a reason to retry immediately.
435                    let retry_at = now
436                        + chrono::TimeDelta::from_std(delay)
437                            .unwrap_or_else(|_| chrono::TimeDelta::hours(1));
438                    self.outbox
439                        .mark_failed(&entry.outbox_id, code.clone(), Some(retry_at))
440                        .await
441                }
442            }
443        };
444        if let Err(error) = settled {
445            tracing::warn!(
446                target: "turnframe.dispatch",
447                outbox_id = %entry.outbox_id,
448                error = %error,
449                "the outbox row could not be settled; it stays claimed until the claim expires"
450            );
451            report.unsettled.push(entry.outbox_id);
452            return;
453        }
454        match outcome {
455            Dispatched::Completed { .. } => report.completed.push(entry.outbox_id),
456            Dispatched::Unknown { .. } => report.unknown.push(entry.outbox_id),
457            Dispatched::Permanent { .. } => report.failed.push(entry.outbox_id),
458            Dispatched::Retryable { .. } => {
459                if entry.attempt_count >= self.config.max_attempts {
460                    report.failed.push(entry.outbox_id);
461                } else {
462                    report.retried.push(entry.outbox_id);
463                }
464            }
465        }
466    }
467
468    /// Sends one row under the configured deadline.
469    ///
470    /// A send that does not answer in time is an unknown outcome and never a
471    /// retry: the bytes left, and nobody can say they did not land (I15).
472    async fn send(&self, entry: &OutboxEntry) -> Dispatched {
473        // The stage is the call to the remote system and not the bookkeeping
474        // around it (§28), and it is measured whether the call answered,
475        // refused or timed out: a send that hangs for the whole deadline is the
476        // most interesting point in the distribution.
477        let stage = Stage::enter();
478        let outcome =
479            match tokio::time::timeout(self.config.send_timeout, self.sender.send(entry)).await {
480                Ok(outcome) => outcome,
481                Err(_) => {
482                    tracing::warn!(
483                        target: "turnframe.dispatch",
484                        outbox_id = %entry.outbox_id,
485                        "the send did not answer in time; the outcome is unknown, not a failure"
486                    );
487                    Dispatched::Unknown { remote_ref: None }
488                }
489            };
490        stage.observe(
491            self.observer.as_ref(),
492            Signal::ExternalLatency,
493            &SignalLabels::none(),
494        );
495        outcome
496    }
497
498    /// Settles one row whose outcome is unknown, through `reconciler`
499    /// (spec §16.5).
500    ///
501    /// The row is read, handed to the reconciler and settled with its answer.
502    /// A row that is not in [`OutboxStatus::OutcomeUnknown`] is left alone and
503    /// reported as [`Reconciled::Unresolved`]: reconciliation is for the rows
504    /// nobody knows about, and a completed row is not one of them.
505    ///
506    /// # Errors
507    ///
508    /// [`StoreError`] when the row could not be read or the settlement was
509    /// refused.
510    pub async fn reconcile(
511        &self,
512        outbox_id: &OutboxId,
513        reconciler: &dyn OutboxReconciler,
514        now: DateTime<Utc>,
515    ) -> Result<Reconciled, StoreError> {
516        let record = self.outbox.get(outbox_id).await?;
517        if record.entry.status != OutboxStatus::OutcomeUnknown {
518            return Ok(Reconciled::Unresolved);
519        }
520        let answer = reconciler.reconcile(&record).await;
521        tracing::debug!(
522            target: "turnframe.dispatch",
523            outbox_id = %outbox_id,
524            answer = answer.as_str(),
525            "unknown outcome reconciled"
526        );
527        match &answer {
528            Reconciled::Completed => self.outbox.mark_completed(outbox_id).await?,
529            Reconciled::Failed { code } => {
530                self.outbox
531                    .mark_failed(outbox_id, code.clone(), None)
532                    .await?;
533            }
534            Reconciled::Resend => self.outbox.reschedule(outbox_id, now).await?,
535            Reconciled::Unresolved => {}
536        }
537        // An unknown outcome that is now known, and only then: `Unresolved`
538        // settled nothing and the row is still waiting for the next sweep.
539        if !matches!(answer, Reconciled::Unresolved) {
540            self.observer
541                .observe_labeled(&Signal::ExternalReconciled, &SignalLabels::none());
542        }
543        Ok(answer)
544    }
545
546    /// Releases rows a crashed dispatcher left in `Dispatching`, so another
547    /// sweep can claim them.
548    ///
549    /// `claimed_before` is the age at which a claim is considered abandoned;
550    /// it must be older than the longest send this dispatcher can make, or a
551    /// slow send is released while it is still running.
552    ///
553    /// # Errors
554    ///
555    /// [`StoreError`] when the sweep could not be made.
556    pub async fn release_expired_claims(
557        &self,
558        claimed_before: DateTime<Utc>,
559    ) -> Result<Vec<OutboxId>, StoreError> {
560        self.outbox.release_expired_claims(claimed_before).await
561    }
562}
563
564#[cfg(test)]
565mod tests {
566    use super::*;
567
568    #[test]
569    fn the_backoff_grows_and_is_capped() {
570        let config =
571            DispatchConfig::new("w").with_backoff(Duration::from_secs(1), Duration::from_secs(60));
572        assert_eq!(config.backoff_for(0), Duration::from_secs(1));
573        assert_eq!(config.backoff_for(1), Duration::from_secs(1));
574        assert_eq!(config.backoff_for(2), Duration::from_secs(4));
575        assert_eq!(config.backoff_for(3), Duration::from_secs(16));
576        assert_eq!(config.backoff_for(9), Duration::from_secs(60));
577        assert_eq!(config.backoff_for(u32::MAX), Duration::from_secs(60));
578    }
579
580    #[test]
581    fn a_report_counts_every_row_it_settled() {
582        let mut report = DispatchReport::default();
583        assert!(report.is_empty());
584        report.completed.push(OutboxId::nil());
585        report.unknown.push(OutboxId::nil());
586        assert_eq!(report.claimed(), 2);
587        assert!(!report.is_empty());
588    }
589
590    #[test]
591    fn every_outcome_names_itself() {
592        assert_eq!(Dispatched::completed().as_str(), "completed");
593        assert_eq!(
594            Dispatched::Retryable {
595                code: "x".to_owned()
596            }
597            .as_str(),
598            "retryable"
599        );
600        assert_eq!(
601            Dispatched::Permanent {
602                code: "x".to_owned()
603            }
604            .as_str(),
605            "permanent"
606        );
607        assert_eq!(Dispatched::Unknown { remote_ref: None }.as_str(), "unknown");
608        assert!(
609            format!("{:?}", DispatchConfig::new("w")).contains("worker_id"),
610            "the configuration is inspectable"
611        );
612    }
613}