Skip to main content

codoseo_crawler/
frontier.rs

1//! The URL queue: depth levels, de-duplication and a hard cap on the total.
2
3use std::collections::{HashSet, VecDeque};
4
5use codoseo_core::url::url_hash;
6use url::Url;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct Queued {
10    pub url: Url,
11    /// Clicks from the start page; `None` for pages found only in sitemaps.
12    pub depth: Option<u16>,
13}
14
15/// Holds the URLs still to fetch. Never admits more than `max_pages` URLs to fetch in
16/// total (fetched plus queued), so memory stays flat on endless URL spaces. Redirect
17/// targets are remembered past the cap but never queued, so records stay at most
18/// `2 × max_pages`.
19///
20/// Levels are a barrier: `pop` only returns URLs of the current level, and the next level
21/// becomes current when the caller asks with `advance`. Sitemap-only URLs come last.
22#[derive(Debug)]
23pub struct Frontier {
24    seen: HashSet<u64>,
25    current: VecDeque<Queued>,
26    next: VecDeque<Queued>,
27    sitemap_only: VecDeque<Queued>,
28    admitted: u32,
29    cap: u32,
30    capped: bool,
31}
32
33impl Frontier {
34    pub fn new(max_pages: u32) -> Frontier {
35        Frontier {
36            seen: HashSet::new(),
37            current: VecDeque::new(),
38            next: VecDeque::new(),
39            sitemap_only: VecDeque::new(),
40            admitted: 0,
41            cap: max_pages,
42            capped: false,
43        }
44    }
45
46    /// Queues the start URL at depth 0. False when it was seen already or the cap is reached.
47    pub fn seed(&mut self, url: Url) -> bool {
48        if !self.try_admit(&url) {
49            return false;
50        }
51        self.current.push_back(Queued {
52            url,
53            depth: Some(0),
54        });
55        true
56    }
57
58    /// Queues a link found on a page at `from_depth` for the next level.
59    pub fn push_link(&mut self, url: Url, from_depth: u16) -> bool {
60        if !self.try_admit(&url) {
61            return false;
62        }
63        self.next.push_back(Queued {
64            url,
65            depth: Some(from_depth.saturating_add(1)),
66        });
67        true
68    }
69
70    /// Marks a redirect target as seen without queueing or counting it: its record is
71    /// built from the response that already arrived, so it costs no fetch. The cap does
72    /// not apply (it bounds fetches), but a URL already seen is refused.
73    pub fn admit_redirect_target(&mut self, url: &Url) -> bool {
74        self.seen.insert(url_hash(url))
75    }
76
77    /// Queues sitemap URLs that nothing linked to; they are fetched after link exploration.
78    pub fn add_sitemap_urls(&mut self, urls: impl IntoIterator<Item = Url>) {
79        for url in urls {
80            if !self.try_admit(&url) {
81                if self.capped {
82                    break;
83                }
84                continue;
85            }
86            self.sitemap_only.push_back(Queued { url, depth: None });
87        }
88    }
89
90    /// The next URL of the current level.
91    pub fn pop(&mut self) -> Option<Queued> {
92        self.current.pop_front()
93    }
94
95    /// Starts the next level. When link exploration is over, starts the sitemap-only URLs
96    /// (once). False when nothing is left.
97    pub fn advance(&mut self) -> bool {
98        if !self.next.is_empty() {
99            self.current.append(&mut self.next);
100        } else if !self.sitemap_only.is_empty() {
101            self.current.append(&mut self.sitemap_only);
102        }
103        !self.current.is_empty()
104    }
105
106    /// Puts a URL back at the front of the current level (a 429/503 retry). It was counted
107    /// when it was first queued and is not counted again.
108    pub fn requeue_front(&mut self, q: Queued) {
109        self.current.push_front(q);
110    }
111
112    pub fn is_seen(&self, url: &Url) -> bool {
113        self.seen.contains(&url_hash(url))
114    }
115
116    /// True once a URL was refused because the cap was reached.
117    pub fn capped(&self) -> bool {
118        self.capped
119    }
120
121    /// URLs waiting to be fetched, in every queue.
122    pub fn queued(&self) -> usize {
123        self.current.len() + self.next.len() + self.sitemap_only.len()
124    }
125
126    /// Remembers and counts a new URL. A refused URL is not remembered.
127    fn try_admit(&mut self, url: &Url) -> bool {
128        let hash = url_hash(url);
129        if self.seen.contains(&hash) {
130            return false;
131        }
132        if self.admitted >= self.cap {
133            self.capped = true;
134            return false;
135        }
136        self.seen.insert(hash);
137        self.admitted += 1;
138        true
139    }
140}
141
142#[cfg(test)]
143mod tests {
144    use codoseo_core::url::normalize;
145
146    use super::*;
147
148    fn u(s: &str) -> Url {
149        Url::parse(s).unwrap()
150    }
151
152    #[test]
153    fn levels_dedup_sitemap_and_cap() {
154        let mut f = Frontier::new(5);
155        assert!(f.seed(u("https://e.com/")));
156        assert_eq!(f.pop().unwrap().depth, Some(0));
157        assert!(f.pop().is_none());
158        assert!(f.push_link(u("https://e.com/a"), 0));
159        assert!(!f.push_link(u("https://e.com/a"), 0)); // dedup
160        assert!(f.pop().is_none()); // next level waits for advance()
161        assert!(f.advance());
162        assert_eq!(f.pop().unwrap().depth, Some(1));
163        f.add_sitemap_urls([u("https://e.com/a"), u("https://e.com/s")]); // /a already seen
164        assert!(f.advance());
165        let s = f.pop().unwrap();
166        assert_eq!((s.url.path(), s.depth), ("/s", None));
167        assert!(!f.advance());
168        for i in 0..10 {
169            f.push_link(u(&format!("https://e.com/p{i}")), 1); // 3 admitted so far, cap 5
170        }
171        assert!(f.capped());
172        assert_eq!(f.queued(), 2);
173    }
174
175    #[test]
176    fn trivial_variants_are_one_page() {
177        let base = u("https://e.com/");
178        let mut f = Frontier::new(10);
179        for href in ["/a", "/a?", "/a#x"] {
180            f.push_link(normalize(&base, href).unwrap(), 0);
181        }
182        assert_eq!(f.queued(), 1);
183        assert!(f.push_link(normalize(&base, "/A").unwrap(), 0));
184    }
185
186    #[test]
187    fn redirect_targets_are_marked_seen_past_the_cap_without_counting() {
188        let mut f = Frontier::new(2);
189        assert!(f.seed(u("https://e.com/")));
190        assert!(f.admit_redirect_target(&u("https://e.com/new")));
191        assert!(f.is_seen(&u("https://e.com/new")));
192        assert!(!f.admit_redirect_target(&u("https://e.com/new"))); // dedup
193        assert!(!f.push_link(u("https://e.com/new"), 0));
194        assert_eq!(f.queued(), 1);
195        // Not counted: one more link still fits.
196        assert!(f.push_link(u("https://e.com/a"), 0));
197        // The cap is full: a link is refused, flagged, and not remembered...
198        assert!(!f.capped());
199        assert!(!f.push_link(u("https://e.com/x"), 0));
200        assert!(f.capped());
201        assert!(!f.is_seen(&u("https://e.com/x")));
202        // ...but a redirect target is still admitted.
203        assert!(f.admit_redirect_target(&u("https://e.com/t")));
204        assert_eq!(f.queued(), 2);
205    }
206
207    #[test]
208    fn seen_url_does_not_trip_the_cap() {
209        let mut f = Frontier::new(1);
210        assert!(f.seed(u("https://e.com/")));
211        assert!(!f.push_link(u("https://e.com/"), 0));
212        assert!(!f.capped());
213    }
214
215    #[test]
216    fn requeue_front_goes_first_and_is_not_counted_again() {
217        let mut f = Frontier::new(2);
218        f.seed(u("https://e.com/"));
219        f.push_link(u("https://e.com/a"), 0);
220        f.pop();
221        f.advance();
222        let a = f.pop().unwrap();
223        f.requeue_front(a.clone());
224        assert_eq!(f.queued(), 1);
225        assert_eq!(f.pop().unwrap(), a);
226        // Still room for nothing new: 2 admitted, cap 2.
227        assert!(!f.push_link(u("https://e.com/b"), 1));
228        assert!(f.capped());
229    }
230
231    #[test]
232    fn sitemap_urls_wait_for_link_levels_and_respect_the_cap() {
233        let mut f = Frontier::new(3);
234        f.seed(u("https://e.com/"));
235        f.add_sitemap_urls((0..10).map(|i| u(&format!("https://e.com/s{i}"))));
236        assert!(f.capped());
237        assert_eq!(f.queued(), 3); // seed + 2 sitemap URLs
238        assert_eq!(f.pop().unwrap().depth, Some(0));
239        assert!(f.pop().is_none()); // sitemap URLs are not in the current level yet
240        assert!(f.advance());
241        assert_eq!(f.pop().unwrap().depth, None);
242        assert_eq!(f.pop().unwrap().depth, None);
243        assert!(!f.advance());
244    }
245
246    #[test]
247    fn link_levels_come_before_sitemap_urls() {
248        let mut f = Frontier::new(10);
249        f.seed(u("https://e.com/"));
250        f.add_sitemap_urls([u("https://e.com/s")]);
251        f.pop();
252        f.push_link(u("https://e.com/a"), 0);
253        assert!(f.advance());
254        assert_eq!(f.pop().unwrap().url.path(), "/a");
255        assert!(f.pop().is_none());
256        assert!(f.advance());
257        assert_eq!(f.pop().unwrap().url.path(), "/s");
258    }
259}