#[cfg(not(feature = "std"))]
extern crate alloc;
#[cfg(not(feature = "std"))]
use alloc::boxed::Box;
#[cfg(feature = "std")]
use std::time::SystemTime;
use crate::policy::{GatePolicy, ReceptorPolicy};
use crate::transport::error::TransportFailure;
use crate::{Frame, Message};
#[cfg(feature = "std")]
use crate::utils::jitter::decorrelated_bounds;
pub trait PolicyConf
where
Self: Sized,
{
fn with_restart<P: RestartPolicy + 'static>(self, _: P) -> Self {
self
}
fn with_emitter_gate<G: GatePolicy + 'static>(self, _: G) -> Self {
self
}
fn with_collector_gate<G: GatePolicy + 'static>(self, _: G) -> Self {
self
}
fn with_receptor_gate<T: Message, R: ReceptorPolicy<T> + 'static>(self, _: R) -> Self {
self
}
fn with_timeout(self, _: core::time::Duration) -> Self {
self
}
}
pub trait CoreRetryPolicy: Send + Sync {
fn max_attempts(&self) -> usize;
fn delay_ms(&self, attempt: usize) -> u64;
}
pub trait RestartPolicy: CoreRetryPolicy {
fn evaluate(&self, frame: Box<Frame>, failure: &TransportFailure, attempt: usize) -> RetryAction;
}
#[derive(Debug, Clone, PartialEq)]
pub enum RetryAction {
Retry(Box<Frame>),
NoRetry,
}
pub trait JitterStrategy: Send + Sync {
fn apply(&self, base_delay: u64) -> u64;
}
#[cfg(feature = "std")]
#[derive(Default)]
pub struct DecorrelatedJitter;
#[cfg(feature = "std")]
impl JitterStrategy for DecorrelatedJitter {
fn apply(&self, base_delay: u64) -> u64 {
let seed = SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(base_delay);
let (min, range) = decorrelated_bounds(base_delay);
if range == 0 {
return base_delay;
}
min + (seed % range)
}
}
#[derive(Default)]
pub struct NoRestart;
impl RestartPolicy for NoRestart {
fn evaluate(&self, _frame: Box<Frame>, _failure: &TransportFailure, _attempt: usize) -> RetryAction {
RetryAction::NoRetry
}
}
#[cfg(feature = "std")]
pub struct RestartExponentialBackoff {
pub max_attempts: usize,
pub scale_factor: u64,
pub jitter: Option<Box<dyn JitterStrategy>>,
}
#[cfg(feature = "std")]
impl RestartExponentialBackoff {
pub fn new(max_attempts: usize, scale_factor: u64, jitter: Option<Box<dyn JitterStrategy>>) -> Self {
Self { max_attempts, scale_factor, jitter }
}
}
#[cfg(feature = "std")]
impl Default for RestartExponentialBackoff {
fn default() -> Self {
Self { max_attempts: 5, scale_factor: 1000, jitter: Some(Box::new(DecorrelatedJitter)) }
}
}
#[cfg(feature = "std")]
pub struct RestartLinearBackoff {
pub max_attempts: usize,
pub interval_ms: u64,
pub scale_factor: u64,
pub jitter: Option<Box<dyn JitterStrategy>>,
}
#[cfg(feature = "std")]
impl RestartLinearBackoff {
pub fn new(
max_attempts: usize,
interval_ms: u64,
scale_factor: u64,
jitter: Option<Box<dyn JitterStrategy>>,
) -> Self {
Self { max_attempts, interval_ms, scale_factor, jitter }
}
}
#[cfg(feature = "std")]
impl Default for RestartLinearBackoff {
fn default() -> Self {
Self {
max_attempts: 5,
interval_ms: 1000,
scale_factor: 1,
jitter: Some(Box::new(DecorrelatedJitter)),
}
}
}
#[cfg(feature = "std")]
macro_rules! impl_timed_backoff_policy {
($policy:ident, $delay_calc:expr) => {
impl RestartPolicy for $policy {
fn evaluate(&self, frame: Box<Frame>, _failure: &TransportFailure, attempt: usize) -> RetryAction {
use core::time::Duration;
if attempt >= self.max_attempts {
return RetryAction::NoRetry;
}
match &self.jitter {
Some(jitter_strategy) => {
let delay_ms = $delay_calc(self, attempt);
let delay_ms = jitter_strategy.apply(delay_ms);
std::thread::sleep(Duration::from_millis(delay_ms));
}
None => {
std::thread::sleep(Duration::from_millis($delay_calc(self, attempt)));
}
}
RetryAction::Retry(frame)
}
}
};
}
#[cfg(feature = "std")]
impl_timed_backoff_policy!(
RestartExponentialBackoff,
|policy: &RestartExponentialBackoff, attempt: usize| {
let exp = (attempt as u32).min(63);
policy.scale_factor.saturating_mul(2_u64.saturating_pow(exp))
}
);
#[cfg(feature = "std")]
impl_timed_backoff_policy!(RestartLinearBackoff, |policy: &RestartLinearBackoff, attempt: usize| {
policy
.scale_factor
.saturating_mul(policy.interval_ms)
.saturating_mul(attempt as u64 + 1)
});
#[cfg(feature = "std")]
impl CoreRetryPolicy for RestartExponentialBackoff {
fn max_attempts(&self) -> usize {
self.max_attempts
}
fn delay_ms(&self, attempt: usize) -> u64 {
let exp = (attempt as u32).min(63);
let base_delay = self.scale_factor.saturating_mul(2_u64.saturating_pow(exp));
match &self.jitter {
Some(jitter_strategy) => jitter_strategy.apply(base_delay),
None => base_delay,
}
}
}
#[cfg(feature = "std")]
impl CoreRetryPolicy for RestartLinearBackoff {
fn max_attempts(&self) -> usize {
self.max_attempts
}
fn delay_ms(&self, attempt: usize) -> u64 {
let base_delay = self
.scale_factor
.saturating_mul(self.interval_ms)
.saturating_mul(attempt as u64 + 1);
match &self.jitter {
Some(jitter_strategy) => jitter_strategy.apply(base_delay),
None => base_delay,
}
}
}
impl CoreRetryPolicy for NoRestart {
fn max_attempts(&self) -> usize {
0
}
fn delay_ms(&self, _attempt: usize) -> u64 {
0
}
}