subscription_proxy_pool/
health.rs1use std::{
2 num::{NonZeroU32, NonZeroUsize},
3 time::{Duration, Instant},
4};
5
6use crate::{Error, Result};
7
8#[derive(Clone, Debug)]
11pub struct HealthPolicy {
12 pub check_url: Option<String>,
14 pub timeout: Duration,
16 pub concurrency: NonZeroUsize,
18 pub failure_threshold: NonZeroU32,
20 pub cooldown: Duration,
22 pub max_cooldown: Duration,
24 pub cooldown_jitter: f64,
26 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 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}