Skip to main content

subscription_proxy_pool/
health.rs

1use std::{
2    num::{NonZeroU32, NonZeroUsize},
3    time::{Duration, Instant},
4};
5
6use crate::{Error, Result};
7
8/// Passive health feedback is enabled by default. Active probes require an
9/// explicit URL, so a reusable library never depends on a particular website.
10#[derive(Clone, Debug)]
11pub struct HealthPolicy {
12    /// Optional URL fetched through each proxy; 2xx/3xx passes the probe.
13    pub check_url: Option<String>,
14    /// Maximum duration of an active probe.
15    pub timeout: Duration,
16    /// Maximum simultaneous active probes.
17    pub concurrency: NonZeroUsize,
18    /// Consecutive transport failures that open a node's circuit.
19    pub failure_threshold: NonZeroU32,
20    /// First cooldown. Repeated failed recoveries double this delay.
21    pub cooldown: Duration,
22    /// Cap for exponential cooldowns.
23    pub max_cooldown: Duration,
24    /// Random delay variation in [0, 0.5]; zero gives deterministic delays.
25    pub cooldown_jitter: f64,
26    /// Probe new nodes before normal use. Failed nodes enter cooldown and
27    /// can receive one recovery attempt after it expires.
28    /// Requires `check_url` to be set.
29    pub check_on_build: bool,
30}
31
32impl Default for HealthPolicy {
33    fn default() -> Self {
34        Self {
35            check_url: None,
36            timeout: Duration::from_secs(5),
37            concurrency: NonZeroUsize::new(16).unwrap(),
38            failure_threshold: NonZeroU32::new(2).unwrap(),
39            cooldown: Duration::from_secs(30),
40            max_cooldown: Duration::from_secs(300),
41            cooldown_jitter: 0.2,
42            check_on_build: false,
43        }
44    }
45}
46
47impl HealthPolicy {
48    /// Enable initial and periodic probes against a caller-chosen HTTP(S) URL.
49    pub fn active(url: impl Into<String>) -> Result<Self> {
50        let policy = Self {
51            check_url: Some(url.into()),
52            check_on_build: true,
53            ..Self::default()
54        };
55        policy.validate()?;
56        Ok(policy)
57    }
58
59    pub(crate) fn validate(&self) -> Result<()> {
60        if self.check_on_build && self.check_url.is_none() {
61            return Err(Error::Config("initial health checks require a check URL"));
62        }
63        if let Some(value) = &self.check_url {
64            let url = reqwest::Url::parse(value)
65                .map_err(|_| Error::Config("invalid health check URL"))?;
66            if !matches!(url.scheme(), "http" | "https") || url.host_str().is_none() {
67                return Err(Error::Config("health check URL must use HTTP or HTTPS"));
68            }
69        }
70        if self.cooldown.is_zero()
71            || self.max_cooldown < self.cooldown
72            || Instant::now().checked_add(self.max_cooldown).is_none()
73        {
74            return Err(Error::Config(
75                "cooldown must be positive and its maximum must be representable and at least the initial delay",
76            ));
77        }
78        if !self.cooldown_jitter.is_finite() || !(0.0..=0.5).contains(&self.cooldown_jitter) {
79            return Err(Error::Config("cooldown jitter must be between 0 and 0.5"));
80        }
81        Ok(())
82    }
83
84    pub(crate) fn cooldown_for(&self, failures: u32) -> Duration {
85        let nanos = 1_u128
86            .checked_shl(failures.saturating_sub(1))
87            .and_then(|factor| self.cooldown.as_nanos().checked_mul(factor))
88            .unwrap_or(self.max_cooldown.as_nanos())
89            .min(self.max_cooldown.as_nanos());
90        let base = Duration::new(
91            (nanos / 1_000_000_000) as u64,
92            (nanos % 1_000_000_000) as u32,
93        );
94        if self.cooldown_jitter == 0.0 {
95            return base;
96        }
97        let factor =
98            rand::random_range((1.0 - self.cooldown_jitter)..=(1.0 + self.cooldown_jitter));
99        base.mul_f64(factor)
100            .min(self.max_cooldown)
101            .max(Duration::from_nanos(1))
102    }
103}
104
105#[cfg(test)]
106mod tests {
107    use super::*;
108    #[test]
109    fn repeated_failures_back_off_and_stop_at_the_cap() {
110        let policy = HealthPolicy {
111            cooldown: Duration::from_secs(2),
112            max_cooldown: Duration::from_secs(10),
113            cooldown_jitter: 0.0,
114            ..HealthPolicy::default()
115        };
116        let delays: Vec<_> = [1, 2, 3, 4, u32::MAX]
117            .map(|n| policy.cooldown_for(n).as_secs())
118            .into();
119        assert_eq!(delays, [2, 4, 8, 10, 10]);
120        let small = HealthPolicy {
121            cooldown: Duration::from_nanos(1),
122            ..policy
123        };
124        assert_eq!(small.cooldown_for(33), Duration::from_nanos(1_u64 << 32));
125        assert_eq!(small.cooldown_for(u32::MAX), small.max_cooldown);
126    }
127    #[test]
128    fn defaults_need_no_external_probe_and_bad_policies_are_rejected() {
129        assert!(HealthPolicy::default().check_url.is_none());
130        assert!(HealthPolicy::default().validate().is_ok());
131        assert!(
132            HealthPolicy {
133                check_on_build: true,
134                ..HealthPolicy::default()
135            }
136            .validate()
137            .is_err()
138        );
139        assert!(
140            HealthPolicy {
141                cooldown_jitter: f64::NAN,
142                ..HealthPolicy::default()
143            }
144            .validate()
145            .is_err()
146        );
147    }
148}