later 0.0.44

Distributed Background jobs manager and runner for Rust
//! Retry policies for failed job handlers.
//!
//! A message type can implement [`crate::retry::JobRetryPolicy`] to replace the server's
//! default policy. The worker calls the policy after the first handler failure,
//! passing the decoded message and application context. The result is then
//! stored with the job, so later retries do not fetch settings again.
//!
//! Jobs written by older Later versions have no resolution marker. Their
//! stored retry settings remain authoritative after an upgrade.

use async_trait::async_trait;
use std::any::Any;
use std::time::Duration;

/// Delay calculation used between failed handler attempts.
///
/// A calculated delay is the earliest retry time. Worker polling and delivery
/// can make the next attempt start later.
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum RetryBackoff {
    /// Makes the retry eligible immediately.
    Immediate,
    /// Waits for the same delay after every failed attempt.
    Fixed {
        /// Delay before each retry.
        delay: Duration,
    },
    /// Doubles the delay after each failed attempt, up to a maximum.
    Exponential {
        /// Delay before the first retry.
        initial_delay: Duration,
        /// Longest delay between attempts.
        max_delay: Duration,
    },
    /// Adds bounded jitter to an exponential delay.
    ///
    /// Jitter spreads retries from many failed jobs across time. The total
    /// delay never exceeds `max_delay`.
    ExponentialWithJitter {
        /// Delay before adding jitter to the first retry.
        initial_delay: Duration,
        /// Longest total delay between attempts.
        max_delay: Duration,
        /// Largest random-looking delay added to the exponential value.
        max_jitter: Duration,
    },
}

impl RetryBackoff {
    pub(crate) fn delay(&self, job_key: &str, retry_number: usize) -> Duration {
        match self {
            Self::Immediate => Duration::ZERO,
            Self::Fixed { delay } => *delay,
            Self::Exponential {
                initial_delay,
                max_delay,
            } => exponential_delay(*initial_delay, *max_delay, retry_number),
            Self::ExponentialWithJitter {
                initial_delay,
                max_delay,
                max_jitter,
            } => exponential_delay(*initial_delay, *max_delay, retry_number)
                .saturating_add(jitter(job_key, retry_number, *max_jitter))
                .min(*max_delay),
        }
    }
}

/// The retry limit and delay calculation assigned to a job.
///
/// The default allows six retries with exponential backoff from one second to
/// one minute and up to 500 milliseconds of jitter.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RetryPolicy {
    max_retries: usize,
    backoff: RetryBackoff,
}

impl RetryPolicy {
    /// Creates a policy with an explicit retry limit and backoff.
    ///
    /// `max_retries` counts attempts after the first handler failure. Set it to
    /// zero to fail a job after its first unsuccessful attempt.
    pub fn new(max_retries: usize, backoff: RetryBackoff) -> Self {
        Self {
            max_retries,
            backoff,
        }
    }

    /// Creates a policy that retries immediately.
    pub fn immediate(max_retries: usize) -> Self {
        Self::new(max_retries, RetryBackoff::Immediate)
    }

    /// Creates a policy that does not retry a failed handler.
    pub fn no_retries() -> Self {
        Self::immediate(0)
    }

    /// Creates a policy with the same delay between every attempt.
    pub fn fixed(max_retries: usize, delay: Duration) -> Self {
        Self::new(max_retries, RetryBackoff::Fixed { delay })
    }

    /// Creates an exponential policy without jitter.
    pub fn exponential(max_retries: usize, initial_delay: Duration, max_delay: Duration) -> Self {
        Self::new(
            max_retries,
            RetryBackoff::Exponential {
                initial_delay,
                max_delay,
            },
        )
    }

    /// Creates an exponential policy with deterministic bounded jitter.
    ///
    /// The jitter differs by job ID and retry number. It does not need a random
    /// number generator and remains stable if another worker reads the job.
    pub fn exponential_with_jitter(
        max_retries: usize,
        initial_delay: Duration,
        max_delay: Duration,
        max_jitter: Duration,
    ) -> Self {
        Self::new(
            max_retries,
            RetryBackoff::ExponentialWithJitter {
                initial_delay,
                max_delay,
                max_jitter,
            },
        )
    }

    /// Returns the number of attempts allowed after the first failure.
    pub fn max_retries(&self) -> usize {
        self.max_retries
    }

    /// Returns the delay calculation used between attempts.
    pub fn backoff(&self) -> &RetryBackoff {
        &self.backoff
    }
}

impl Default for RetryPolicy {
    fn default() -> Self {
        Self::exponential_with_jitter(
            6,
            Duration::from_secs(1),
            Duration::from_secs(60),
            Duration::from_millis(500),
        )
    }
}

/// Supplies a retry policy for a message.
///
/// Implementing this trait overrides the server default for that message type.
/// The async method receives the message and application context. It can read
/// local state or fetch settings from another service. Fetch errors can be
/// returned to the worker. Types without an implementation continue to use
/// the server default.
///
/// ```
/// use later::retry::{JobRetryPolicy, RetryPolicy};
/// use std::time::Duration;
///
/// struct AppContext {
///     retry_limit: usize,
/// }
/// struct SendEmail;
///
/// #[later::async_trait::async_trait]
/// impl JobRetryPolicy for SendEmail {
///     type Context = AppContext;
///
///     async fn retry_policy(&self, context: &AppContext) -> anyhow::Result<RetryPolicy> {
///         Ok(RetryPolicy::exponential_with_jitter(
///             context.retry_limit,
///             Duration::from_secs(1),
///             Duration::from_secs(60),
///             Duration::from_millis(500),
///         ))
///     }
/// }
/// ```
#[async_trait]
pub trait JobRetryPolicy: Sync {
    /// Application context expected by this message policy.
    type Context: Sync + 'static;

    /// Returns the policy stored after this job's first handler failure.
    async fn retry_policy(&self, context: &Self::Context) -> anyhow::Result<RetryPolicy>;
}

/// Supports optional [`JobRetryPolicy`] implementations in generated code.
#[doc(hidden)]
pub struct RetryPolicyResolver<'a, T, C> {
    message: &'a T,
    context: &'a C,
}

impl<'a, T, C> RetryPolicyResolver<'a, T, C> {
    /// Creates a resolver for generated code.
    #[doc(hidden)]
    pub fn new(message: &'a T, context: &'a C) -> Self {
        Self { message, context }
    }
}

/// Resolves an optional message policy without requiring every type to implement it.
#[doc(hidden)]
#[async_trait]
pub trait ResolveRetryPolicy {
    /// Returns the message override, or `None` when the server default applies.
    #[doc(hidden)]
    async fn resolve_retry_policy(self) -> anyhow::Result<Option<RetryPolicy>>;
}

#[async_trait]
impl<T: Sync, C: Sync> ResolveRetryPolicy for &&RetryPolicyResolver<'_, T, C> {
    async fn resolve_retry_policy(self) -> anyhow::Result<Option<RetryPolicy>> {
        Ok(None)
    }
}

#[async_trait]
impl<T, C> ResolveRetryPolicy for &RetryPolicyResolver<'_, T, C>
where
    T: JobRetryPolicy + Sync,
    C: Any + Sync,
{
    async fn resolve_retry_policy(self) -> anyhow::Result<Option<RetryPolicy>> {
        let context = (self.context as &dyn Any)
            .downcast_ref::<T::Context>()
            .ok_or_else(|| {
                anyhow::anyhow!(
                    "retry policy for {} expects application context {}",
                    std::any::type_name::<T>(),
                    std::any::type_name::<T::Context>(),
                )
            })?;
        JobRetryPolicy::retry_policy(self.message, context)
            .await
            .map(Some)
    }
}

fn exponential_delay(initial: Duration, maximum: Duration, retry_number: usize) -> Duration {
    let exponent = u32::try_from(retry_number.saturating_sub(1)).unwrap_or(u32::MAX);
    let multiplier = 2_u32.checked_pow(exponent).unwrap_or(u32::MAX);
    initial.saturating_mul(multiplier).min(maximum)
}

fn jitter(job_key: &str, retry_number: usize, maximum: Duration) -> Duration {
    let mut input = job_key.as_bytes().to_vec();
    input.extend_from_slice(
        &u64::try_from(retry_number)
            .unwrap_or(u64::MAX)
            .to_le_bytes(),
    );
    let hash = blake3::hash(&input);
    let mut bytes = [0_u8; 8];
    bytes.copy_from_slice(&hash.as_bytes()[..8]);
    let value = u64::from_le_bytes(bytes);
    let maximum_nanos = u64::try_from(maximum.as_nanos()).unwrap_or(u64::MAX);
    let nanos = value % maximum_nanos.saturating_add(1);
    Duration::from_nanos(nanos)
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn backoff_presets_apply_limits_and_bounded_jitter() {
        let exponential =
            RetryPolicy::exponential(4, Duration::from_secs(2), Duration::from_secs(5));
        let jittered = RetryPolicy::exponential_with_jitter(
            4,
            Duration::from_secs(2),
            Duration::from_secs(5),
            Duration::from_secs(1),
        );

        assert_eq!(
            vec![
                exponential.backoff.delay("job", 1),
                exponential.backoff.delay("job", 2),
                exponential.backoff.delay("job", 3),
                RetryPolicy::fixed(1, Duration::from_secs(3))
                    .backoff
                    .delay("job", 1),
                RetryPolicy::immediate(1).backoff.delay("job", 1),
            ],
            vec![
                Duration::from_secs(2),
                Duration::from_secs(4),
                Duration::from_secs(5),
                Duration::from_secs(3),
                Duration::ZERO,
            ]
        );
        let first = jittered.backoff.delay("job-a", 1);
        assert!(first >= Duration::from_secs(2));
        assert!(first <= Duration::from_secs(3));
        assert_eq!(first, jittered.backoff.delay("job-a", 1));
        assert_ne!(first, jittered.backoff.delay("job-b", 1));
        assert!(jittered.backoff.delay("job-a", 20) <= Duration::from_secs(5));
        assert_eq!(
            exponential.backoff.delay("job", usize::MAX),
            Duration::from_secs(5)
        );
        assert_eq!(RetryPolicy::no_retries().max_retries(), 0);
    }
}