Skip to main content

codoseo_crawler/
politeness.rs

1//! Request pacing for one site: a minimum gap between requests, a cap on parallel
2//! connections, and a back-off when the site answers 429 or 503.
3//!
4//! Callers take a [`Permit`] with [`Limiter::acquire`] before every request and report
5//! the status with [`Limiter::on_response`] afterwards. Time comes from
6//! `tokio::time`, so tests can run on a paused clock.
7
8use 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
16/// A `Retry-After` longer than this is cut to it.
17pub const MAX_RETRY_AFTER: Duration = Duration::from_secs(60);
18
19struct State {
20    interval: Duration,
21    /// The earliest start time still free for the next request.
22    next_slot: Instant,
23    paused_until: Instant,
24}
25
26/// Paces the requests to one site and shares a global in-flight cap with other sites.
27pub struct Limiter {
28    state: Mutex<State>,
29    max_interval: Duration,
30    /// `1 / requests_per_sec`, before any crawl delay or back-off.
31    rate_gap: Duration,
32    site: Arc<Semaphore>,
33    global: Arc<Semaphore>,
34}
35
36/// Holds the site and global connection slots until dropped.
37#[derive(Debug)]
38pub struct Permit {
39    _site: OwnedSemaphorePermit,
40    _global: OwnedSemaphorePermit,
41}
42
43impl Limiter {
44    /// The gap between requests is the slower of `1 / requests_per_sec` and the site's
45    /// `Crawl-delay`, never above `Politeness::max_crawl_delay`.
46    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    /// Waits for a site connection, then for this request's start time, then for a
71    /// global slot. The start time is reserved before waiting, so concurrent callers
72    /// get evenly spaced slots. A pause set while waiting (a 429 or 503 on another
73    /// connection) still applies: the request then reserves a new slot after it.
74    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    /// Reports a response. A 429 or 503 doubles the interval (up to the cap) and holds
106    /// back the next requests for `retry_after`, or two new intervals when the site
107    /// gave none.
108    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    /// Applies a `Crawl-delay` learned after creation (robots.txt is read through this
121    /// limiter). The gap becomes the slower of `1 / requests_per_sec` and `delay`, capped
122    /// as in [`Limiter::new`]; a longer gap from a 429/503 back-off is kept.
123    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    /// The current gap between requests.
133    pub fn interval(&self) -> Duration {
134        self.lock().interval
135    }
136
137    fn lock(&self) -> MutexGuard<'_, State> {
138        // The state is plain data, so a poisoned lock is still consistent.
139        self.state.lock().unwrap_or_else(|e| e.into_inner())
140    }
141}
142
143/// Reads a `Retry-After` header: a number of seconds or an HTTP date. A date in the
144/// past gives zero. The result is not capped; see [`MAX_RETRY_AFTER`].
145pub 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
154/// The `Retry-After` of a response, read against the current wall clock.
155pub(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}