use std::fmt;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use turnframe_core::error::StoreError;
use turnframe_core::event::{OutboxEntry, OutboxStatus};
use turnframe_core::ids::OutboxId;
use turnframe_core::observe::{NoopObserver, Observer, Signal, SignalLabels};
use turnframe_store::outbox::{OutboxRecord, OutboxStore};
use crate::signals::Stage;
pub mod code {
pub const ATTEMPTS_EXHAUSTED: &str = "turnframe.dispatch.attempts_exhausted";
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum Dispatched {
Completed {
remote_ref: Option<String>,
},
Retryable {
code: String,
},
Permanent {
code: String,
},
Unknown {
remote_ref: Option<String>,
},
}
impl Dispatched {
#[must_use]
pub const fn completed() -> Self {
Self::Completed { remote_ref: None }
}
#[must_use]
pub const fn as_str(&self) -> &'static str {
match self {
Self::Completed { .. } => "completed",
Self::Retryable { .. } => "retryable",
Self::Permanent { .. } => "permanent",
Self::Unknown { .. } => "unknown",
}
}
}
#[async_trait]
pub trait OutboxSender: Send + Sync {
async fn send(&self, entry: &OutboxEntry) -> Dispatched;
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum Reconciled {
Completed,
Failed {
code: String,
},
Resend,
Unresolved,
}
impl Reconciled {
#[must_use]
pub const fn as_str(&self) -> &'static str {
match self {
Self::Completed => "completed",
Self::Failed { .. } => "failed",
Self::Resend => "resend",
Self::Unresolved => "unresolved",
}
}
}
#[async_trait]
pub trait OutboxReconciler: Send + Sync {
async fn reconcile(&self, record: &OutboxRecord) -> Reconciled;
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct DispatchConfig {
pub worker_id: String,
pub batch_size: usize,
pub send_timeout: Duration,
pub initial_backoff: Duration,
pub max_backoff: Duration,
pub backoff_multiplier: u32,
pub max_attempts: u32,
}
impl DispatchConfig {
#[must_use]
pub fn new(worker_id: impl Into<String>) -> Self {
Self {
worker_id: worker_id.into(),
batch_size: 16,
send_timeout: Duration::from_secs(30),
initial_backoff: Duration::from_secs(1),
max_backoff: Duration::from_secs(60),
backoff_multiplier: 4,
max_attempts: 5,
}
}
#[must_use]
pub fn with_batch_size(mut self, batch_size: usize) -> Self {
self.batch_size = batch_size;
self
}
#[must_use]
pub const fn with_send_timeout(mut self, send_timeout: Duration) -> Self {
self.send_timeout = send_timeout;
self
}
#[must_use]
pub const fn with_max_attempts(mut self, max_attempts: u32) -> Self {
self.max_attempts = max_attempts;
self
}
#[must_use]
pub const fn with_backoff(mut self, initial: Duration, max: Duration) -> Self {
self.initial_backoff = initial;
self.max_backoff = max;
self
}
#[must_use]
pub fn backoff_for(&self, attempt: u32) -> Duration {
let step = attempt.saturating_sub(1);
if step == 0 {
return self.initial_backoff;
}
self.initial_backoff
.saturating_mul(self.backoff_multiplier.saturating_pow(step.min(16)))
.min(self.max_backoff)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct DispatchReport {
pub completed: Vec<OutboxId>,
pub retried: Vec<OutboxId>,
pub failed: Vec<OutboxId>,
pub unknown: Vec<OutboxId>,
pub unsettled: Vec<OutboxId>,
}
impl DispatchReport {
#[must_use]
pub fn claimed(&self) -> usize {
self.completed.len()
+ self.retried.len()
+ self.failed.len()
+ self.unknown.len()
+ self.unsettled.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.claimed() == 0
}
}
#[derive(Clone)]
pub struct OutboxDispatcher {
outbox: Arc<dyn OutboxStore>,
sender: Arc<dyn OutboxSender>,
config: DispatchConfig,
observer: Arc<dyn Observer>,
}
impl fmt::Debug for OutboxDispatcher {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("OutboxDispatcher")
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl OutboxDispatcher {
#[must_use]
pub fn new(
outbox: Arc<dyn OutboxStore>,
sender: Arc<dyn OutboxSender>,
config: DispatchConfig,
) -> Self {
Self {
outbox,
sender,
config,
observer: Arc::new(NoopObserver),
}
}
#[must_use]
pub fn with_observer(mut self, observer: Arc<dyn Observer>) -> Self {
self.observer = observer;
self
}
#[must_use]
pub const fn config(&self) -> &DispatchConfig {
&self.config
}
pub async fn run_once(&self, now: DateTime<Utc>) -> Result<DispatchReport, StoreError> {
let claimed = self
.outbox
.claim_due(now, self.config.batch_size, &self.config.worker_id)
.await?;
let mut report = DispatchReport::default();
for entry in claimed {
self.dispatch_one(&entry, now, &mut report).await;
}
Ok(report)
}
async fn dispatch_one(
&self,
entry: &OutboxEntry,
now: DateTime<Utc>,
report: &mut DispatchReport,
) {
let outcome = self.send(entry).await;
tracing::debug!(
target: "turnframe.dispatch",
outbox_id = %entry.outbox_id,
destination = entry.destination.as_str(),
attempt = entry.attempt_count,
outcome = outcome.as_str(),
"outbox row dispatched"
);
let settled = match &outcome {
Dispatched::Completed { .. } => self.outbox.mark_completed(&entry.outbox_id).await,
Dispatched::Unknown { remote_ref } => {
self.outbox
.mark_outcome_unknown(&entry.outbox_id, remote_ref.clone())
.await
}
Dispatched::Permanent { code } => {
self.outbox
.mark_failed(&entry.outbox_id, code.clone(), None)
.await
}
Dispatched::Retryable { code } => {
if entry.attempt_count >= self.config.max_attempts {
self.outbox
.mark_failed(&entry.outbox_id, code::ATTEMPTS_EXHAUSTED.to_owned(), None)
.await
} else {
let delay = self.config.backoff_for(entry.attempt_count);
let retry_at = now
+ chrono::TimeDelta::from_std(delay)
.unwrap_or_else(|_| chrono::TimeDelta::hours(1));
self.outbox
.mark_failed(&entry.outbox_id, code.clone(), Some(retry_at))
.await
}
}
};
if let Err(error) = settled {
tracing::warn!(
target: "turnframe.dispatch",
outbox_id = %entry.outbox_id,
error = %error,
"the outbox row could not be settled; it stays claimed until the claim expires"
);
report.unsettled.push(entry.outbox_id);
return;
}
match outcome {
Dispatched::Completed { .. } => report.completed.push(entry.outbox_id),
Dispatched::Unknown { .. } => report.unknown.push(entry.outbox_id),
Dispatched::Permanent { .. } => report.failed.push(entry.outbox_id),
Dispatched::Retryable { .. } => {
if entry.attempt_count >= self.config.max_attempts {
report.failed.push(entry.outbox_id);
} else {
report.retried.push(entry.outbox_id);
}
}
}
}
async fn send(&self, entry: &OutboxEntry) -> Dispatched {
let stage = Stage::enter();
let outcome =
match tokio::time::timeout(self.config.send_timeout, self.sender.send(entry)).await {
Ok(outcome) => outcome,
Err(_) => {
tracing::warn!(
target: "turnframe.dispatch",
outbox_id = %entry.outbox_id,
"the send did not answer in time; the outcome is unknown, not a failure"
);
Dispatched::Unknown { remote_ref: None }
}
};
stage.observe(
self.observer.as_ref(),
Signal::ExternalLatency,
&SignalLabels::none(),
);
outcome
}
pub async fn reconcile(
&self,
outbox_id: &OutboxId,
reconciler: &dyn OutboxReconciler,
now: DateTime<Utc>,
) -> Result<Reconciled, StoreError> {
let record = self.outbox.get(outbox_id).await?;
if record.entry.status != OutboxStatus::OutcomeUnknown {
return Ok(Reconciled::Unresolved);
}
let answer = reconciler.reconcile(&record).await;
tracing::debug!(
target: "turnframe.dispatch",
outbox_id = %outbox_id,
answer = answer.as_str(),
"unknown outcome reconciled"
);
match &answer {
Reconciled::Completed => self.outbox.mark_completed(outbox_id).await?,
Reconciled::Failed { code } => {
self.outbox
.mark_failed(outbox_id, code.clone(), None)
.await?;
}
Reconciled::Resend => self.outbox.reschedule(outbox_id, now).await?,
Reconciled::Unresolved => {}
}
if !matches!(answer, Reconciled::Unresolved) {
self.observer
.observe_labeled(&Signal::ExternalReconciled, &SignalLabels::none());
}
Ok(answer)
}
pub async fn release_expired_claims(
&self,
claimed_before: DateTime<Utc>,
) -> Result<Vec<OutboxId>, StoreError> {
self.outbox.release_expired_claims(claimed_before).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_backoff_grows_and_is_capped() {
let config =
DispatchConfig::new("w").with_backoff(Duration::from_secs(1), Duration::from_secs(60));
assert_eq!(config.backoff_for(0), Duration::from_secs(1));
assert_eq!(config.backoff_for(1), Duration::from_secs(1));
assert_eq!(config.backoff_for(2), Duration::from_secs(4));
assert_eq!(config.backoff_for(3), Duration::from_secs(16));
assert_eq!(config.backoff_for(9), Duration::from_secs(60));
assert_eq!(config.backoff_for(u32::MAX), Duration::from_secs(60));
}
#[test]
fn a_report_counts_every_row_it_settled() {
let mut report = DispatchReport::default();
assert!(report.is_empty());
report.completed.push(OutboxId::nil());
report.unknown.push(OutboxId::nil());
assert_eq!(report.claimed(), 2);
assert!(!report.is_empty());
}
#[test]
fn every_outcome_names_itself() {
assert_eq!(Dispatched::completed().as_str(), "completed");
assert_eq!(
Dispatched::Retryable {
code: "x".to_owned()
}
.as_str(),
"retryable"
);
assert_eq!(
Dispatched::Permanent {
code: "x".to_owned()
}
.as_str(),
"permanent"
);
assert_eq!(Dispatched::Unknown { remote_ref: None }.as_str(), "unknown");
assert!(
format!("{:?}", DispatchConfig::new("w")).contains("worker_id"),
"the configuration is inspectable"
);
}
}