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}