codoseo_crawler/
politeness.rs1use std::sync::{Arc, Mutex, MutexGuard};
9use std::time::{Duration, SystemTime};
10
11use codoseo_core::crawl::Politeness;
12use reqwest::header::{HeaderMap, RETRY_AFTER};
13use tokio::sync::{OwnedSemaphorePermit, Semaphore};
14use tokio::time::{Instant, sleep_until};
15
16pub const MAX_RETRY_AFTER: Duration = Duration::from_secs(60);
18
19struct State {
20 interval: Duration,
21 next_slot: Instant,
23 paused_until: Instant,
24}
25
26pub struct Limiter {
28 state: Mutex<State>,
29 max_interval: Duration,
30 rate_gap: Duration,
32 site: Arc<Semaphore>,
33 global: Arc<Semaphore>,
34}
35
36#[derive(Debug)]
38pub struct Permit {
39 _site: OwnedSemaphorePermit,
40 _global: OwnedSemaphorePermit,
41}
42
43impl Limiter {
44 pub fn new(p: &Politeness, crawl_delay: Option<Duration>, global: Arc<Semaphore>) -> Limiter {
47 let max_interval = p.max_crawl_delay;
48 let rate_gap = if p.requests_per_sec.is_finite() && p.requests_per_sec > 0.0 {
49 Duration::from_secs_f64(1.0 / f64::from(p.requests_per_sec))
50 } else {
51 max_interval
52 };
53 let interval = rate_gap
54 .max(crawl_delay.unwrap_or(Duration::ZERO))
55 .min(max_interval);
56 let now = Instant::now();
57 Limiter {
58 state: Mutex::new(State {
59 interval,
60 next_slot: now,
61 paused_until: now,
62 }),
63 max_interval,
64 rate_gap,
65 site: Arc::new(Semaphore::new(p.per_site_connections.max(1) as usize)),
66 global,
67 }
68 }
69
70 pub async fn acquire(&self) -> Permit {
75 let site = self
76 .site
77 .clone()
78 .acquire_owned()
79 .await
80 .expect("the site semaphore is never closed");
81 loop {
82 let slot = {
83 let mut s = self.lock();
84 let slot = Instant::now().max(s.next_slot).max(s.paused_until);
85 s.next_slot = slot + s.interval;
86 slot
87 };
88 sleep_until(slot).await;
89 if self.lock().paused_until <= Instant::now() {
90 break;
91 }
92 }
93 let global = self
94 .global
95 .clone()
96 .acquire_owned()
97 .await
98 .expect("the global semaphore is never closed");
99 Permit {
100 _site: site,
101 _global: global,
102 }
103 }
104
105 pub fn on_response(&self, status: u16, retry_after: Option<Duration>) {
109 if status != 429 && status != 503 {
110 return;
111 }
112 let mut s = self.lock();
113 s.interval = s.interval.saturating_mul(2).min(self.max_interval);
114 let pause = retry_after
115 .unwrap_or_else(|| s.interval.saturating_mul(2))
116 .min(MAX_RETRY_AFTER);
117 s.paused_until = s.paused_until.max(Instant::now() + pause);
118 }
119
120 pub fn set_crawl_delay(&self, delay: Option<Duration>) {
124 let base = self
125 .rate_gap
126 .max(delay.unwrap_or(Duration::ZERO))
127 .min(self.max_interval);
128 let mut s = self.lock();
129 s.interval = s.interval.max(base);
130 }
131
132 pub fn interval(&self) -> Duration {
134 self.lock().interval
135 }
136
137 fn lock(&self) -> MutexGuard<'_, State> {
138 self.state.lock().unwrap_or_else(|e| e.into_inner())
140 }
141}
142
143pub fn parse_retry_after(value: &str, now: SystemTime) -> Option<Duration> {
146 let value = value.trim();
147 if let Ok(secs) = value.parse::<u64>() {
148 return Some(Duration::from_secs(secs));
149 }
150 let at = httpdate::parse_http_date(value).ok()?;
151 Some(at.duration_since(now).unwrap_or(Duration::ZERO))
152}
153
154pub(crate) fn retry_after_of(headers: &HeaderMap) -> Option<Duration> {
156 headers
157 .get(RETRY_AFTER)
158 .and_then(|v| v.to_str().ok())
159 .and_then(|v| parse_retry_after(v, SystemTime::now()))
160}