Skip to main content

codoseo_crawler/
sitemap.rs

1//! Sitemap parsing (streamed, gzip-aware) and discovery through sitemap indexes.
2//!
3//! Parsing streams through gunzip and a lossy UTF-8 decoder, so memory stays flat
4//! however large a file is. Discovery is bounded three ways: [`MAX_SITEMAP_FILES`]
5//! fetch attempts (failures count), a URL cap, and a deadline.
6
7use std::collections::{HashSet, VecDeque};
8use std::io::{BufReader, Read};
9
10use codoseo_core::crawl::SitemapSummary;
11use codoseo_core::url::{normalize, url_hash};
12use encoding_rs::{Decoder, UTF_8};
13use flate2::read::GzDecoder;
14use quick_xml::Reader;
15use quick_xml::events::Event;
16use tokio::time::{Instant, timeout_at};
17use url::Url;
18use xxhash_rust::xxh3::xxh3_64;
19
20use crate::fetch::Fetcher;
21use crate::politeness::Limiter;
22
23/// The sitemap protocol's size limit, also applied after decompression.
24const MAX_SITEMAP_BYTES: usize = 50 * 1024 * 1024;
25/// Seeds are depth 0; indexes at depth 2 are read but their children are not.
26const MAX_DEPTH: u8 = 2;
27/// Fetch attempts per discovery, successful or not.
28pub const MAX_SITEMAP_FILES: usize = 100;
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
31pub enum SitemapKind {
32    UrlSet,
33    Index,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq)]
37pub struct SitemapDoc {
38    pub kind: SitemapKind,
39    pub locs: Vec<String>,
40}
41
42#[derive(Debug, thiserror::Error)]
43pub enum SitemapError {
44    #[error("not a sitemap")]
45    NotASitemap,
46    #[error("invalid XML: {0}")]
47    Xml(String),
48}
49
50pub struct SitemapDiscovery {
51    pub urls: Vec<Url>,
52    pub summary: SitemapSummary,
53}
54
55pub fn parse_sitemap(bytes: &[u8]) -> Result<SitemapDoc, SitemapError> {
56    let mut locs = Vec::new();
57    let kind = parse_stream(bytes, |_, loc| {
58        locs.push(loc.to_owned());
59        true
60    })?;
61    Ok(SitemapDoc { kind, locs })
62}
63
64/// Streams the `<loc>` values that sit directly under `<url>` or `<sitemap>` (so
65/// `<image:loc>` and friends are skipped) to `on_loc`, which returns false to stop.
66fn parse_stream(
67    bytes: &[u8],
68    mut on_loc: impl FnMut(SitemapKind, &str) -> bool,
69) -> Result<SitemapKind, SitemapError> {
70    let source: Box<dyn Read + '_> = if bytes.starts_with(&[0x1f, 0x8b]) {
71        Box::new(GzDecoder::new(bytes).take(MAX_SITEMAP_BYTES as u64))
72    } else {
73        Box::new(bytes)
74    };
75    let mut reader = Reader::from_reader(BufReader::new(LossyUtf8::new(source)));
76    let mut buf = Vec::new();
77    let mut kind = None;
78    let mut depth = 0usize;
79    let mut in_entry = false;
80    let mut in_loc = false;
81    let mut current = String::new();
82
83    loop {
84        match reader.read_event_into(&mut buf) {
85            Ok(Event::Start(e)) => {
86                depth += 1;
87                let local = e.local_name();
88                match depth {
89                    1 => {
90                        kind = Some(match local.as_ref() {
91                            "urlset" => SitemapKind::UrlSet,
92                            "sitemapindex" => SitemapKind::Index,
93                            _ => return Err(SitemapError::NotASitemap),
94                        })
95                    }
96                    2 => in_entry = matches!(local.as_ref(), "url" | "sitemap"),
97                    3 if in_entry && local.as_ref() == "loc" && e.name().prefix().is_none() => {
98                        in_loc = true;
99                        current.clear();
100                    }
101                    _ => {}
102                }
103            }
104            Ok(Event::Empty(e)) if depth == 0 => {
105                return match e.local_name().as_ref() {
106                    "urlset" => Ok(SitemapKind::UrlSet),
107                    "sitemapindex" => Ok(SitemapKind::Index),
108                    _ => Err(SitemapError::NotASitemap),
109                };
110            }
111            Ok(Event::Text(t)) if in_loc => current.push_str(&t.xml10_content()),
112            Ok(Event::CData(c)) if in_loc => current.push_str(&c.xml10_content()),
113            Ok(Event::GeneralRef(r)) if in_loc => {
114                let ch = match r.resolve_char_ref() {
115                    Ok(Some(ch)) => Some(ch),
116                    _ => match r.as_ref() {
117                        "amp" => Some('&'),
118                        "lt" => Some('<'),
119                        "gt" => Some('>'),
120                        "quot" => Some('"'),
121                        "apos" => Some('\''),
122                        _ => None,
123                    },
124                };
125                current.extend(ch);
126            }
127            Ok(Event::End(_)) => {
128                if in_loc {
129                    in_loc = false;
130                    let loc = current.trim();
131                    if !loc.is_empty()
132                        && let Some(k) = kind
133                        && !on_loc(k, loc)
134                    {
135                        break;
136                    }
137                }
138                depth = depth.saturating_sub(1);
139            }
140            Ok(Event::Eof) => break,
141            // Broken XML after the root element still yields what we read so far.
142            Err(e) if kind.is_none() => return Err(SitemapError::Xml(e.to_string())),
143            Err(_) => break,
144            Ok(_) => {}
145        }
146        buf.clear();
147    }
148    kind.ok_or(SitemapError::NotASitemap)
149}
150
151/// Decodes bytes as UTF-8, replacing invalid sequences, so one bad byte in a
152/// sitemap doesn't end parsing. Also drops a UTF-8 BOM.
153struct LossyUtf8<R: Read> {
154    inner: R,
155    decoder: Decoder,
156    input: Vec<u8>,
157    output: Vec<u8>,
158    pos: usize,
159    done: bool,
160}
161
162impl<R: Read> LossyUtf8<R> {
163    fn new(inner: R) -> Self {
164        LossyUtf8 {
165            inner,
166            decoder: UTF_8.new_decoder_with_bom_removal(),
167            input: vec![0; 16 * 1024],
168            output: Vec::new(),
169            pos: 0,
170            done: false,
171        }
172    }
173}
174
175impl<R: Read> Read for LossyUtf8<R> {
176    fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
177        while self.pos == self.output.len() {
178            if self.done {
179                return Ok(0);
180            }
181            let n = self.inner.read(&mut self.input)?;
182            let last = n == 0;
183            let capacity = self.decoder.max_utf8_buffer_length(n).unwrap_or(n * 3 + 4);
184            self.output.resize(capacity, 0);
185            let (_, _, written, _) =
186                self.decoder
187                    .decode_to_utf8(&self.input[..n], &mut self.output, last);
188            self.output.truncate(written);
189            self.pos = 0;
190            self.done = last;
191        }
192        let n = out.len().min(self.output.len() - self.pos);
193        out[..n].copy_from_slice(&self.output[self.pos..self.pos + n]);
194        self.pos += n;
195        Ok(n)
196    }
197}
198
199/// Fetches the seed sitemaps, follows indexes up to depth 2, and returns up to
200/// `max_urls` unique page URLs. Stops at `deadline` with what it has. With a `limiter`,
201/// every sitemap fetch waits for a permit first.
202pub async fn discover(
203    fetcher: &Fetcher,
204    limiter: Option<&Limiter>,
205    seeds: &[Url],
206    max_urls: u32,
207    deadline: Instant,
208) -> SitemapDiscovery {
209    let max_urls = max_urls as usize;
210    let mut known: HashSet<String> = HashSet::new();
211    let mut queue: VecDeque<(Url, u8)> = VecDeque::new();
212    for seed in seeds {
213        if known.insert(seed.as_str().to_owned()) {
214            queue.push_back((seed.clone(), 0));
215        }
216    }
217    let mut seen: HashSet<u64> = HashSet::new();
218    let mut urls: Vec<Url> = Vec::new();
219    let mut files: Vec<Url> = Vec::new();
220    let mut attempts = 0usize;
221    let mut failed = 0u32;
222    let mut truncated = false;
223    let mut dropped = false;
224    let mut stopped = false;
225
226    while let Some((sitemap_url, depth)) = queue.pop_front() {
227        if truncated {
228            break;
229        }
230        if Instant::now() >= deadline {
231            stopped = true;
232            break;
233        }
234        attempts += 1;
235        let fetch = async {
236            let _permit = match limiter {
237                Some(l) => Some(l.acquire().await),
238                None => None,
239            };
240            fetcher.fetch_raw(&sitemap_url, MAX_SITEMAP_BYTES).await
241        };
242        let res = match timeout_at(deadline, fetch).await {
243            Err(_) => {
244                stopped = true;
245                break;
246            }
247            Ok(Err(_)) => {
248                failed += 1;
249                continue;
250            }
251            Ok(Ok(res)) => res,
252        };
253        let Some(body) = res.body.filter(|_| (200..300).contains(&res.status)) else {
254            failed += 1;
255            continue;
256        };
257
258        let parsed = parse_stream(&body, |kind, loc| match kind {
259            SitemapKind::Index => {
260                if depth >= MAX_DEPTH {
261                    return false;
262                }
263                if attempts + queue.len() >= MAX_SITEMAP_FILES {
264                    dropped = true;
265                    return false;
266                }
267                if let Some(child) = normalize(&sitemap_url, loc)
268                    && known.insert(child.as_str().to_owned())
269                {
270                    queue.push_back((child, depth + 1));
271                }
272                true
273            }
274            SitemapKind::UrlSet => {
275                let Some(url) = normalize(&sitemap_url, loc) else {
276                    return true;
277                };
278                if !seen.insert(url_hash(&url)) {
279                    return true;
280                }
281                if urls.len() >= max_urls {
282                    truncated = true;
283                    return false;
284                }
285                urls.push(url);
286                true
287            }
288        });
289        match parsed {
290            Ok(_) => files.push(sitemap_url),
291            Err(_) => failed += 1,
292        }
293    }
294
295    let complete = !stopped && !dropped && (queue.is_empty() || truncated);
296    let mut hashes: Vec<u64> = urls.iter().map(url_hash).collect();
297    hashes.sort_unstable();
298    let hash_bytes: Vec<u8> = hashes.iter().flat_map(|h| h.to_le_bytes()).collect();
299    SitemapDiscovery {
300        summary: SitemapSummary {
301            files,
302            url_count: urls.len() as u32,
303            hash: xxh3_64(&hash_bytes),
304            truncated,
305            failed_files: failed,
306            complete,
307        },
308        urls,
309    }
310}