use std::collections::{HashSet, VecDeque};
use std::hash::{DefaultHasher, Hash, Hasher};
use std::time::{Duration, SystemTime};
use tokio::task::JoinSet;
use tokio::time::{MissedTickBehavior, interval};
use url::Url;
use crate::bridge::{self, PageFetcher};
use crate::net;
use crate::robots::RobotsPolicy;
use crate::scope::{is_same_site, matches_scope, normalize_url};
const MAX_HTML_BYTES: usize = 2 * 1024 * 1024;
#[must_use = "options do nothing until passed to crawl() or crawl_each()"]
#[derive(Debug, Clone)]
pub struct CrawlOptions {
pub(crate) url: String,
pub(crate) limit: usize,
pub(crate) max_depth: usize,
pub(crate) timeout: Duration,
pub(crate) settle: Duration,
pub(crate) include: Vec<String>,
pub(crate) exclude: Vec<String>,
pub(crate) selector: Option<String>,
pub(crate) json: bool,
pub(crate) user_agent: Option<String>,
pub(crate) robots_user_agent: Option<String>,
pub(crate) concurrency: usize,
pub(crate) delay: Option<Duration>,
pub(crate) cookies: Vec<crate::cookies::CookieSpec>,
pub(crate) headers: http::HeaderMap,
}
impl CrawlOptions {
pub fn new(url: &str) -> Self {
Self {
url: url.into(),
limit: 50,
max_depth: 3,
timeout: Duration::from_secs(30),
settle: Duration::ZERO,
include: Vec::new(),
exclude: Vec::new(),
selector: None,
json: false,
user_agent: None,
robots_user_agent: None,
concurrency: 1,
delay: Some(Duration::from_millis(500)),
cookies: Vec::new(),
headers: http::HeaderMap::new(),
}
}
pub(crate) fn validate(&self) -> crate::error::Result<()> {
net::validate_url(&self.url)?;
if !self.include.is_empty() {
crate::scope::build_globset(&self.include)?;
}
if !self.exclude.is_empty() {
crate::scope::build_globset(&self.exclude)?;
}
Ok(())
}
pub fn limit(mut self, n: usize) -> Self {
self.limit = n;
self
}
pub fn max_depth(mut self, n: usize) -> Self {
self.max_depth = n;
self
}
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
pub fn settle(mut self, settle: Duration) -> Self {
self.settle = settle;
self
}
pub fn include(mut self, patterns: &[&str]) -> Self {
self.include = patterns.iter().map(|s| (*s).to_string()).collect();
self
}
pub fn exclude(mut self, patterns: &[&str]) -> Self {
self.exclude = patterns.iter().map(|s| (*s).to_string()).collect();
self
}
pub fn json(mut self, json: bool) -> Self {
self.json = json;
self
}
pub fn selector(mut self, selector: impl Into<String>) -> Self {
self.selector = Some(selector.into());
self
}
pub fn user_agent(mut self, ua: impl Into<String>) -> Self {
self.user_agent = Some(net::sanitize_user_agent(ua.into()));
self
}
pub fn concurrency(mut self, n: usize) -> Self {
self.concurrency = n.max(1);
self
}
pub fn delay(mut self, delay: Option<Duration>) -> Self {
self.delay = delay;
self
}
pub fn cookies(mut self, cookies: Vec<crate::cookies::CookieSpec>) -> Self {
self.cookies = cookies;
self
}
pub fn headers(mut self, headers: http::HeaderMap) -> Self {
self.headers = headers;
self
}
}
#[derive(Debug)]
#[non_exhaustive]
pub struct CrawlResult {
pub url: String,
pub depth: usize,
pub fetched_at: SystemTime,
pub outcome: Result<CrawlPage, crate::error::Error>,
}
#[derive(Debug, Clone)]
pub struct CrawlPage {
pub title: Option<String>,
pub content: String,
pub links_found: usize,
}
impl serde::Serialize for CrawlResult {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
use serde::ser::SerializeMap;
let fetched_at = humantime::format_rfc3339_millis(self.fetched_at).to_string();
match &self.outcome {
Ok(page) => {
let mut map = serializer.serialize_map(None)?;
map.serialize_entry("type", "page")?;
map.serialize_entry("url", &self.url)?;
map.serialize_entry("depth", &self.depth)?;
map.serialize_entry("fetchedAt", &fetched_at)?;
if let Some(t) = &page.title {
map.serialize_entry("title", t)?;
}
map.serialize_entry("content", &page.content)?;
map.serialize_entry("linksFound", &page.links_found)?;
map.end()
}
Err(e) => {
let mut map = serializer.serialize_map(None)?;
map.serialize_entry("type", "error")?;
map.serialize_entry("url", &self.url)?;
map.serialize_entry("depth", &self.depth)?;
map.serialize_entry("fetchedAt", &fetched_at)?;
map.serialize_entry("error", &e.report())?;
map.end()
}
}
}
}
impl CrawlResult {
fn from_internal(r: CrawlPageResult) -> Self {
let outcome = match r.status {
CrawlStatus::Ok => Ok(CrawlPage {
title: r.title,
content: r.content.unwrap_or_default(),
links_found: r.links_found,
}),
CrawlStatus::Error => Err(r
.error
.unwrap_or_else(|| crate::error::Error::engine("unknown crawl error", None))),
};
Self {
url: r.url,
depth: r.depth,
fetched_at: r.fetched_at,
outcome,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SuppressionReason {
DuplicateContent,
}
#[derive(Debug)]
pub(crate) struct SuppressedPage {
pub url: String,
pub depth: usize,
pub reason: SuppressionReason,
}
pub(crate) enum CrawlSessionEvent {
Result(CrawlResult),
Suppressed(SuppressedPage),
}
pub(crate) fn crawl_each_in_process_blocking_with_events<F>(
opts: &CrawlOptions,
on_event: F,
) -> crate::error::Result<()>
where
F: FnMut(CrawlSessionEvent) -> crate::error::Result<()>,
{
crate::runtime::block_on(crawl_each_in_process_with_events(opts, on_event))
.map_err(|error| crate::error::Error::engine(error, None))?
}
async fn crawl_each_in_process_with_events<F>(opts: &CrawlOptions, mut on_event: F) -> crate::error::Result<()>
where
F: FnMut(CrawlSessionEvent) -> crate::error::Result<()>,
{
net::ensure_crypto_provider();
let plan = build_crawl_plan(opts)?;
let client = crate::transfer::client()?;
let headers = crate::transfer::Headers::new(&plan.headers, plan.robots_user_agent.as_deref());
let robots = crate::robots::fetch(&client, &plan.seed, &headers, Duration::from_secs(plan.timeout_secs)).await;
run(plan, robots, &bridge::ServoFetcher, |event| {
on_event(match event {
CrawlRunEvent::Result(result) => CrawlSessionEvent::Result(CrawlResult::from_internal(result)),
CrawlRunEvent::Suppressed(page) => CrawlSessionEvent::Suppressed(page),
})
})
.await
}
fn emit_crawl_result(event: CrawlSessionEvent, on_page: &mut impl FnMut(CrawlResult)) {
match event {
CrawlSessionEvent::Result(result) => on_page(result),
CrawlSessionEvent::Suppressed(page) => {
tracing::debug!(
url = %page.url,
depth = page.depth,
reason = ?page.reason,
"crawl result suppressed"
);
}
}
}
pub fn crawl_each_blocking<F>(opts: &CrawlOptions, mut on_page: F) -> crate::error::Result<()>
where
F: FnMut(CrawlResult) + Send,
{
crawl_each_in_process_blocking_with_events(opts, |event| {
emit_crawl_result(event, &mut on_page);
Ok(())
})
}
pub async fn crawl_each<F>(opts: &CrawlOptions, mut on_page: F) -> crate::error::Result<()>
where
F: FnMut(CrawlResult) + Send,
{
crawl_each_in_process_with_events(opts, |event| {
emit_crawl_result(event, &mut on_page);
Ok(())
})
.await
}
pub fn crawl_blocking(opts: &CrawlOptions) -> crate::error::Result<Vec<CrawlResult>> {
let mut results = Vec::new();
crawl_each_blocking(opts, |r| results.push(r))?;
Ok(results)
}
pub async fn crawl(opts: &CrawlOptions) -> crate::error::Result<Vec<CrawlResult>> {
let mut results = Vec::new();
crawl_each(opts, |r| results.push(r)).await?;
Ok(results)
}
fn build_crawl_plan(opts: &CrawlOptions) -> crate::error::Result<CrawlPlan> {
let seed = net::validate_url(&opts.url)?;
let include = if opts.include.is_empty() {
None
} else {
Some(crate::scope::build_globset(&opts.include)?)
};
let exclude = if opts.exclude.is_empty() {
None
} else {
Some(crate::scope::build_globset(&opts.exclude)?)
};
Ok(CrawlPlan {
seed,
limit: opts.limit,
max_depth: opts.max_depth,
timeout_secs: opts.timeout.as_secs().max(1),
settle_ms: u64::try_from(opts.settle.as_millis()).unwrap_or(u64::MAX),
include,
exclude,
selector: opts.selector.clone(),
json: opts.json,
user_agent: opts.user_agent.clone(),
robots_user_agent: opts.robots_user_agent.clone().or_else(|| opts.user_agent.clone()),
concurrency: opts.concurrency,
delay: opts.delay,
cookies: opts.cookies.clone(),
headers: opts.headers.clone(),
})
}
pub(crate) struct CrawlPlan {
pub seed: Url,
pub limit: usize,
pub max_depth: usize,
pub timeout_secs: u64,
pub settle_ms: u64,
pub include: Option<globset::GlobSet>,
pub exclude: Option<globset::GlobSet>,
pub selector: Option<String>,
pub json: bool,
pub user_agent: Option<String>,
pub robots_user_agent: Option<String>,
pub concurrency: usize,
pub delay: Option<Duration>,
pub cookies: Vec<crate::cookies::CookieSpec>,
pub headers: http::HeaderMap,
}
pub(crate) struct CrawlPageResult {
pub url: String,
pub depth: usize,
pub status: CrawlStatus,
pub title: Option<String>,
pub content: Option<String>,
pub error: Option<crate::error::Error>,
pub links_found: usize,
pub fetched_at: SystemTime,
}
pub(crate) enum CrawlStatus {
Ok,
Error,
}
struct Frontier {
queue: VecDeque<(Url, usize)>,
visited: HashSet<String>,
content_hashes: HashSet<u64>,
}
impl Frontier {
fn new(seed: &Url) -> Self {
Self {
queue: VecDeque::from([(seed.clone(), 0)]),
visited: HashSet::from([normalize_url(seed)]),
content_hashes: HashSet::new(),
}
}
fn try_enqueue(&mut self, url: Url, depth: usize) -> bool {
if self.visited.insert(normalize_url(&url)) {
self.queue.push_back((url, depth));
true
} else {
false
}
}
fn pop(&mut self) -> Option<(Url, usize)> {
self.queue.pop_front()
}
fn is_duplicate_content(&mut self, content: &str) -> bool {
let mut h = DefaultHasher::new();
content.hash(&mut h);
!self.content_hashes.insert(h.finish())
}
fn pending(&self) -> usize {
self.queue.len()
}
}
fn extract_links_from_html(html: &str, base: &Url) -> Vec<Url> {
dom_query::Document::from(html)
.select("a[href]")
.iter()
.filter_map(|el| {
let href = el.attr("href")?;
let href = href.trim();
if href.is_empty() {
return None;
}
let resolved = base.join(href).ok()?;
matches!(resolved.scheme(), "http" | "https").then_some(resolved)
})
.collect()
}
pub(crate) enum CrawlRunEvent {
Result(CrawlPageResult),
Suppressed(SuppressedPage),
}
pub(crate) async fn run(
opts: CrawlPlan,
robots: RobotsPolicy,
fetcher: &(impl PageFetcher + Clone),
mut on_event: impl FnMut(CrawlRunEvent) -> crate::error::Result<()>,
) -> crate::error::Result<()> {
let mut frontier = Frontier::new(&opts.seed);
let mut completed: usize = 0;
let mut in_flight: JoinSet<FetchOutcome> = JoinSet::new();
let mut ticker = opts.delay.map(|period| {
let mut t = interval(period);
t.set_missed_tick_behavior(MissedTickBehavior::Delay);
t
});
let concurrency = opts.concurrency.max(1);
let failure = 'crawl: loop {
while in_flight.len() < concurrency && completed + in_flight.len() < opts.limit {
let Some((url, depth)) = frontier.pop() else {
break;
};
if let Some(t) = ticker.as_mut() {
t.tick().await;
}
spawn_fetch(&mut in_flight, fetcher, &opts, url, depth);
}
let outcome = match in_flight.join_next().await {
None => break None,
Some(Ok(outcome)) => outcome,
Some(Err(error)) => break Some(crate::worker::worker_error(error)),
};
let FetchOutcome {
url,
depth,
result,
fetched_at,
} = outcome;
let page = match result {
Ok(p) => p,
Err(err) => {
if let Err(error) = on_event(CrawlRunEvent::Result(error_result(&url, depth, err, fetched_at))) {
break 'crawl Some(error);
}
completed += 1;
continue;
}
};
let budget_used = completed + in_flight.len() + 1;
let mut ctx = CrawlContext {
frontier: &mut frontier,
robots: &robots,
opts: &opts,
};
match process_ok_fetch(&mut ctx, &url, depth, &page, budget_used, fetched_at) {
event @ CrawlRunEvent::Result(_) => {
if let Err(error) = on_event(event) {
break 'crawl Some(error);
}
completed += 1;
}
event @ CrawlRunEvent::Suppressed(_) => {
if let Err(error) = on_event(event) {
break 'crawl Some(error);
}
}
}
};
if let Some(error) = failure {
in_flight.shutdown().await;
Err(error)
} else {
Ok(())
}
}
fn spawn_fetch(
in_flight: &mut JoinSet<FetchOutcome>,
fetcher: &(impl PageFetcher + Clone),
opts: &CrawlPlan,
url: Url,
depth: usize,
) {
let url_str = url.to_string();
let timeout = opts.timeout_secs;
let settle = opts.settle_ms;
let user_agent = opts.user_agent.clone();
let cookies = opts.cookies.clone();
let headers = opts.headers.clone();
let f = fetcher.clone();
in_flight.spawn_blocking(move || {
let result = f
.fetch_page(bridge::PageOptions {
url: &url_str,
timeout_secs: timeout,
settle_ms: settle,
mode: bridge::FetchMode::Content { include_a11y: false },
user_agent: user_agent.as_deref(),
cookies: &cookies,
headers: &headers,
})
.map_err(|error| match error {
bridge::EngineError::Timeout(timeout_secs) => crate::error::Error::Timeout {
url: url_str.clone(),
timeout: Duration::from_secs(timeout_secs),
},
bridge::EngineError::Other(error) => crate::error::Error::engine(error, Some(url_str.clone())),
});
FetchOutcome {
url,
depth,
result,
fetched_at: SystemTime::now(),
}
});
}
struct CrawlContext<'a> {
frontier: &'a mut Frontier,
robots: &'a RobotsPolicy,
opts: &'a CrawlPlan,
}
fn process_ok_fetch(
ctx: &mut CrawlContext<'_>,
url: &Url,
depth: usize,
page: &bridge::ServoPage,
budget_used: usize,
fetched_at: SystemTime,
) -> CrawlRunEvent {
let html = if page.html.len() > MAX_HTML_BYTES {
&page.html[..crate::sanitize::floor_char_boundary(&page.html, MAX_HTML_BYTES)]
} else {
&page.html
};
let document_url = match net::validate_url(&page.url) {
Ok(document_url) => document_url,
Err(error) => return CrawlRunEvent::Result(error_result(url, depth, error, fetched_at)),
};
let input = crate::extract::ExtractInput::new(html, document_url.as_str())
.with_layout_json(page.layout_json.as_deref())
.with_inner_text(page.inner_text.as_deref())
.with_selector(ctx.opts.selector.as_deref());
let content = if ctx.opts.json {
crate::extract::extract_article(&input).and_then(|article| serde_json::to_string(&article).map_err(Into::into))
} else {
crate::extract::extract_text(&input)
};
let content = match content {
Ok(content) => content,
Err(error) => {
return CrawlRunEvent::Result(error_result(url, depth, error.into(), fetched_at));
}
};
if ctx.frontier.is_duplicate_content(&content) {
return CrawlRunEvent::Suppressed(SuppressedPage {
url: url.to_string(),
depth,
reason: SuppressionReason::DuplicateContent,
});
}
let links = extract_links_from_html(html, &document_url);
let links_found = links.len();
if depth < ctx.opts.max_depth {
for link in &links {
if budget_used + ctx.frontier.pending() >= ctx.opts.limit {
break;
}
if !is_same_site(&ctx.opts.seed, link)
|| net::validate_url_with_policy(link.as_str(), bridge::engine_policy()).is_err()
|| !ctx.robots.is_allowed(link)
|| !matches_scope(link, ctx.opts.include.as_ref(), ctx.opts.exclude.as_ref())
{
continue;
}
ctx.frontier.try_enqueue(link.clone(), depth + 1);
}
}
let title = {
let doc = dom_query::Document::from(html);
let t = doc.select("title").text().to_string();
(!t.is_empty()).then_some(t)
};
CrawlRunEvent::Result(CrawlPageResult {
url: url.to_string(),
depth,
status: CrawlStatus::Ok,
title,
content: Some(crate::sanitize::sanitize(&content).into_owned()),
error: None,
links_found,
fetched_at,
})
}
struct FetchOutcome {
url: Url,
depth: usize,
result: Result<bridge::ServoPage, crate::error::Error>,
fetched_at: SystemTime,
}
fn error_result(url: &Url, depth: usize, error: crate::error::Error, fetched_at: SystemTime) -> CrawlPageResult {
CrawlPageResult {
url: url.to_string(),
depth,
status: CrawlStatus::Error,
title: None,
content: None,
error: Some(error),
links_found: 0,
fetched_at,
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Barrier, Mutex};
use tokio::sync::oneshot;
use super::*;
#[test]
fn crawl_options_defaults() {
let opts = CrawlOptions::new("https://example.com");
assert_eq!(opts.url, "https://example.com");
assert_eq!(opts.limit, 50);
assert_eq!(opts.max_depth, 3);
assert_eq!(opts.timeout, Duration::from_secs(30));
assert!(opts.include.is_empty());
assert!(opts.exclude.is_empty());
assert_eq!(opts.concurrency, 1);
assert_eq!(opts.delay, Some(Duration::from_millis(500)));
}
#[test]
fn crawl_options_chaining() {
let opts = CrawlOptions::new("https://example.com")
.limit(100)
.max_depth(5)
.timeout(Duration::from_secs(60))
.include(&["/docs/**"])
.exclude(&["/docs/archive/**"])
.concurrency(4)
.delay(None);
assert_eq!(opts.limit, 100);
assert_eq!(opts.max_depth, 5);
assert_eq!(opts.include, vec!["/docs/**"]);
assert_eq!(opts.exclude, vec!["/docs/archive/**"]);
assert_eq!(opts.concurrency, 4);
assert_eq!(opts.delay, None);
}
#[test]
fn crawl_options_concurrency_clamps_below_one() {
let opts = CrawlOptions::new("https://example.com").concurrency(0);
assert_eq!(opts.concurrency, 1);
}
#[test]
fn crawl_options_delay_custom_value() {
let opts = CrawlOptions::new("https://example.com").delay(Some(Duration::from_secs(2)));
assert_eq!(opts.delay, Some(Duration::from_secs(2)));
}
#[test]
fn crawl_user_agent_sanitizes_crlf() {
let opts = CrawlOptions::new("https://example.com").user_agent("Crawler\r\n/2.0");
assert_eq!(opts.user_agent.as_deref(), Some("Crawler /2.0"));
}
#[derive(Clone)]
struct MockFetcher(Arc<HashMap<String, (String, String)>>);
impl MockFetcher {
fn new(pages: &[(&str, &str)]) -> Self {
Self(Arc::new(
pages
.iter()
.map(|(url, html)| (url.to_string(), (url.to_string(), html.to_string())))
.collect(),
))
}
fn with_document_urls(pages: &[(&str, &str, &str)]) -> Self {
Self(Arc::new(
pages
.iter()
.map(|(requested_url, document_url, html)| {
(requested_url.to_string(), (document_url.to_string(), html.to_string()))
})
.collect(),
))
}
}
impl PageFetcher for MockFetcher {
fn fetch_page(&self, opts: bridge::PageOptions<'_>) -> Result<bridge::ServoPage, bridge::EngineError> {
self.0
.get(opts.url)
.map(|(url, html)| bridge::ServoPage {
url: url.clone(),
html: html.clone(),
..Default::default()
})
.ok_or_else(|| bridge::EngineError::Other(anyhow::anyhow!("not found: {}", opts.url)))
}
}
#[derive(Clone)]
struct TimeoutFetcher;
impl PageFetcher for TimeoutFetcher {
fn fetch_page(&self, opts: bridge::PageOptions<'_>) -> Result<bridge::ServoPage, bridge::EngineError> {
Err(bridge::EngineError::Timeout(opts.timeout_secs))
}
}
struct PanicSignal(Option<oneshot::Sender<()>>);
impl Drop for PanicSignal {
fn drop(&mut self) {
if let Some(sender) = self.0.take() {
let _ = sender.send(());
}
}
}
struct ControlledFailureState {
start: Barrier,
panicked: Mutex<Option<oneshot::Sender<()>>>,
release: Mutex<Option<oneshot::Receiver<()>>>,
blocking_completed: AtomicBool,
}
#[derive(Clone)]
struct ControlledFailureFetcher(Arc<ControlledFailureState>);
impl PageFetcher for ControlledFailureFetcher {
fn fetch_page(&self, opts: bridge::PageOptions<'_>) -> Result<bridge::ServoPage, bridge::EngineError> {
match Url::parse(opts.url).expect("test URL").path() {
"/" => Ok(bridge::ServoPage {
url: opts.url.to_string(),
html: page(&["/controlled", "/block"]),
..Default::default()
}),
"/block" => {
self.0.start.wait();
self.0
.release
.lock()
.unwrap()
.take()
.expect("blocking task waits once")
.blocking_recv()
.expect("test releases blocking task");
self.0.blocking_completed.store(true, Ordering::SeqCst);
Ok(bridge::ServoPage {
url: opts.url.to_string(),
html: distinct_page("block"),
..Default::default()
})
}
"/controlled" => {
self.0.start.wait();
let panicked = self.0.panicked.lock().unwrap().take();
if panicked.is_some() {
let _signal = PanicSignal(panicked);
panic!("intentional crawl fetch panic");
}
Ok(bridge::ServoPage {
url: opts.url.to_string(),
html: distinct_page("controlled"),
..Default::default()
})
}
path => panic!("unexpected test path {path}"),
}
}
}
fn controlled_fetcher(panicked: Option<oneshot::Sender<()>>) -> (ControlledFailureFetcher, oneshot::Sender<()>) {
let (release_tx, release) = oneshot::channel();
let fetcher = ControlledFailureFetcher(Arc::new(ControlledFailureState {
start: Barrier::new(2),
panicked: Mutex::new(panicked),
release: Mutex::new(Some(release)),
blocking_completed: AtomicBool::new(false),
}));
(fetcher, release_tx)
}
fn controlled_plan() -> CrawlPlan {
let mut plan = test_plan("https://example.com/");
plan.concurrency = 2;
plan
}
async fn assert_crawl_pending(crawling: &mut tokio::task::JoinHandle<crate::error::Result<()>>) {
assert!(
tokio::time::timeout(Duration::from_millis(100), crawling)
.await
.is_err(),
"crawl returned while a started blocking sibling was still running"
);
}
fn page(links: &[&str]) -> String {
use std::fmt::Write as _;
let mut anchors = String::new();
for l in links {
write!(anchors, r#"<a href="{l}">link</a>"#).unwrap();
}
format!("<html><head><title>Test</title></head><body>{anchors}</body></html>")
}
fn distinct_page(tag: &str) -> String {
format!("<html><head><title>{tag}</title></head><body>page {tag}</body></html>")
}
fn test_plan(seed: &str) -> CrawlPlan {
CrawlPlan {
seed: Url::parse(seed).unwrap(),
limit: 50,
max_depth: 3,
timeout_secs: 30,
settle_ms: 0,
include: None,
exclude: None,
selector: None,
json: false,
user_agent: None,
robots_user_agent: None,
concurrency: 1,
delay: None,
cookies: Vec::new(),
headers: http::HeaderMap::new(),
}
}
async fn check(
pages: &[(&str, &str)],
configure: impl FnOnce(&mut CrawlPlan),
assert: impl FnOnce(&[CrawlPageResult]),
) {
let fetcher = MockFetcher::new(pages);
let seed = pages[0].0;
let mut opts = test_plan(seed);
configure(&mut opts);
let mut results = Vec::new();
run(opts, RobotsPolicy::Unavailable, &fetcher, |progress| {
if let CrawlRunEvent::Result(result) = progress {
results.push(result);
}
Ok(())
})
.await
.unwrap();
assert(&results);
}
#[tokio::test]
async fn crawl_single_page() {
check(
&[("https://example.com/", &page(&[]))],
|_| {},
|r| {
assert_eq!(r.len(), 1);
assert_eq!(r[0].url, "https://example.com/");
},
)
.await;
}
#[tokio::test]
async fn crawl_event_sink_failure_waits_for_started_blocking_sibling() {
let (fetcher, release_tx) = controlled_fetcher(None);
let state = Arc::clone(&fetcher.0);
let (sink_tx, sink_rx) = oneshot::channel();
let mut sink_tx = Some(sink_tx);
let mut crawling = tokio::spawn(async move {
run(controlled_plan(), RobotsPolicy::Unavailable, &fetcher, |event| {
if matches!(event, CrawlRunEvent::Result(ref result) if result.url.ends_with("/controlled")) {
sink_tx
.take()
.expect("one controlled event")
.send(())
.expect("test waits");
Err(crate::error::Error::SessionCancelled)
} else {
Ok(())
}
})
.await
});
sink_rx.await.expect("sink failed after blocking sibling started");
assert_crawl_pending(&mut crawling).await;
release_tx.send(()).expect("blocking sibling still running");
let error = crawling.await.expect("crawl joins").expect_err("sink failure returned");
assert!(matches!(error, crate::error::Error::SessionCancelled));
assert!(state.blocking_completed.load(Ordering::SeqCst));
}
#[tokio::test]
async fn crawl_preserves_page_load_timeout_kind() {
let mut opts = test_plan("https://example.com/");
opts.limit = 1;
opts.max_depth = 0;
opts.timeout_secs = 7;
let mut results = Vec::new();
run(opts, RobotsPolicy::Unavailable, &TimeoutFetcher, |progress| {
if let CrawlRunEvent::Result(result) = progress {
results.push(result);
}
Ok(())
})
.await
.unwrap();
assert_eq!(results.len(), 1);
assert!(matches!(
results[0].error,
Some(crate::error::Error::Timeout { ref url, timeout })
if url == "https://example.com/" && timeout == Duration::from_secs(7)
));
}
#[tokio::test]
async fn crawl_preserves_invalid_selector_as_page_error() {
check(
&[("https://example.com/", &page(&[]))],
|o| o.selector = Some("[".to_string()),
|r| {
assert_eq!(r.len(), 1);
assert!(matches!(
r[0].error,
Some(crate::error::Error::Extract(
crate::extract::ExtractError::InvalidSelector
))
));
},
)
.await;
}
#[tokio::test]
async fn crawl_fetch_task_join_failure_waits_for_started_blocking_sibling() {
let (panicked_tx, panicked_rx) = oneshot::channel();
let (fetcher, release_tx) = controlled_fetcher(Some(panicked_tx));
let state = Arc::clone(&fetcher.0);
let mut crawling =
tokio::spawn(async move { run(controlled_plan(), RobotsPolicy::Unavailable, &fetcher, |_| Ok(())).await });
panicked_rx.await.expect("failing task panicked");
assert_crawl_pending(&mut crawling).await;
release_tx.send(()).expect("blocking sibling still running");
let error = crawling.await.expect("crawl joins").expect_err("join failure returned");
assert!(matches!(error, crate::error::Error::WorkerUnavailable { .. }));
assert!(state.blocking_completed.load(Ordering::SeqCst));
}
#[tokio::test]
async fn crawl_follows_links() {
check(
&[
("https://example.com/", &page(&["/a", "/b"])),
(
"https://example.com/a",
"<html><head><title>A</title></head><body>page a</body></html>",
),
(
"https://example.com/b",
"<html><head><title>B</title></head><body>page b</body></html>",
),
],
|_| {},
|r| assert_eq!(r.len(), 3),
)
.await;
}
#[tokio::test]
async fn crawl_uses_document_url_for_extraction_and_link_resolution() {
let fetcher = MockFetcher::with_document_urls(&[
(
"https://request.example.com/start",
"https://final.example.com/dir/",
"<!doctype html><html><body><main><a href=\"next\">Next page</a></main></body></html>",
),
(
"https://final.example.com/dir/next",
"https://final.example.com/dir/next",
"<!doctype html><html><body><main>Redirect target child</main></body></html>",
),
]);
let mut opts = test_plan("https://request.example.com/start");
opts.json = true;
opts.limit = 2;
let mut results = Vec::new();
run(opts, RobotsPolicy::Unavailable, &fetcher, |event| {
if let CrawlRunEvent::Result(result) = event {
results.push(result);
}
Ok(())
})
.await
.unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[0].url, "https://request.example.com/start");
let article: serde_json::Value = serde_json::from_str(results[0].content.as_deref().unwrap()).unwrap();
assert_eq!(article["url"], "https://final.example.com/dir/");
assert!(
article["textContent"]
.as_str()
.is_some_and(|markdown| markdown.contains("https://final.example.com/dir/next"))
);
assert_eq!(results[1].url, "https://final.example.com/dir/next");
assert!(matches!(results[1].status, CrawlStatus::Ok));
}
#[tokio::test]
async fn crawl_respects_depth_limit() {
check(
&[
("https://example.com/", &page(&["/a"])),
("https://example.com/a", &page(&["/b"])),
("https://example.com/b", &page(&["/c"])),
("https://example.com/c", &page(&[])),
],
|o| o.max_depth = 1,
|r| assert_eq!(r.len(), 2),
)
.await;
}
#[tokio::test]
async fn crawl_respects_limit() {
check(
&[
("https://example.com/", &page(&["/a", "/b", "/c"])),
("https://example.com/a", &page(&[])),
("https://example.com/b", &page(&[])),
("https://example.com/c", &page(&[])),
],
|o| o.limit = 2,
|r| assert_eq!(r.len(), 2),
)
.await;
}
#[tokio::test]
async fn crawl_skips_cross_site_links() {
check(
&[
("https://example.com/", &page(&["https://other.com/x"])),
("https://other.com/x", &page(&[])),
],
|_| {},
|r| assert_eq!(r.len(), 1),
)
.await;
}
#[tokio::test]
async fn crawl_deduplicates_urls() {
check(
&[
("https://example.com/", &page(&["/a", "/a", "/a"])),
("https://example.com/a", &page(&["/"])),
],
|_| {},
|r| assert_eq!(r.len(), 2),
)
.await;
}
#[tokio::test]
async fn crawl_handles_fetch_errors() {
check(
&[("https://example.com/", &page(&["/missing"]))],
|_| {},
|r| {
assert_eq!(r.len(), 2);
assert!(matches!(r[1].status, CrawlStatus::Error));
assert!(r[1].error.is_some());
},
)
.await;
}
#[tokio::test]
async fn crawl_applies_include_glob() {
check(
&[
("https://example.com/", &page(&["/docs/a", "/blog/b"])),
("https://example.com/docs/a", &page(&[])),
("https://example.com/blog/b", &page(&[])),
],
|o| o.include = Some(crate::scope::build_globset(&["/docs/**".into()]).unwrap()),
|r| {
assert_eq!(r.len(), 2);
assert!(r.iter().any(|p| p.url == "https://example.com/docs/a"));
assert!(!r.iter().any(|p| p.url == "https://example.com/blog/b"));
},
)
.await;
}
#[tokio::test]
async fn crawl_applies_exclude_glob() {
check(
&[
("https://example.com/", &page(&["/public", "/secret/data"])),
("https://example.com/public", &page(&[])),
("https://example.com/secret/data", &page(&[])),
],
|o| o.exclude = Some(crate::scope::build_globset(&["/secret/**".into()]).unwrap()),
|r| {
assert_eq!(r.len(), 2);
assert!(!r.iter().any(|p| p.url == "https://example.com/secret/data"));
},
)
.await;
}
#[tokio::test]
async fn suppressed_page_does_not_consume_result_limit() {
let root = page(&["/a", "/b", "/c"]);
let duplicate = "<html><body>same</body></html>";
let distinct = r#"<html><body>distinct<a href="/d"></a></body></html>"#;
let final_page = "<html><body>final</body></html>";
let fetcher = MockFetcher::new(&[
("https://example.com/", &root),
("https://example.com/a", duplicate),
("https://example.com/b", duplicate),
("https://example.com/c", distinct),
("https://example.com/d", final_page),
]);
let opts = CrawlPlan {
seed: Url::parse("https://example.com/").unwrap(),
limit: 4,
max_depth: 3,
timeout_secs: 30,
settle_ms: 0,
include: None,
exclude: None,
selector: None,
json: false,
user_agent: None,
robots_user_agent: None,
concurrency: 1,
delay: None,
cookies: Vec::new(),
headers: http::HeaderMap::new(),
};
let mut results = Vec::new();
let mut suppressed = 0;
run(opts, RobotsPolicy::Unavailable, &fetcher, |event| {
match event {
CrawlRunEvent::Result(result) => results.push(result.url),
CrawlRunEvent::Suppressed(_) => suppressed += 1,
}
Ok(())
})
.await
.unwrap();
assert_eq!(suppressed, 1);
assert_eq!(
results,
[
"https://example.com/",
"https://example.com/a",
"https://example.com/c",
"https://example.com/d",
]
);
}
#[tokio::test]
async fn crawl_reports_duplicate_content_suppression() {
let same = "<html><head><title>Same</title></head><body>identical</body></html>";
let fetcher = MockFetcher::new(&[
("https://example.com/", &page(&["/a", "/b"])),
("https://example.com/a", same),
("https://example.com/b", same),
]);
let opts = CrawlPlan {
seed: Url::parse("https://example.com/").unwrap(),
limit: 50,
max_depth: 3,
timeout_secs: 30,
settle_ms: 0,
include: None,
exclude: None,
selector: None,
json: false,
user_agent: None,
robots_user_agent: None,
concurrency: 1,
delay: None,
cookies: Vec::new(),
headers: http::HeaderMap::new(),
};
let mut results = 0;
let mut suppressed = Vec::new();
run(opts, RobotsPolicy::Unavailable, &fetcher, |event| {
match event {
CrawlRunEvent::Result(_) => results += 1,
CrawlRunEvent::Suppressed(page) => suppressed.push(page),
}
Ok(())
})
.await
.unwrap();
assert_eq!(results, 2);
assert_eq!(suppressed.len(), 1);
assert_eq!(suppressed[0].depth, 1);
assert_eq!(suppressed[0].reason, SuppressionReason::DuplicateContent);
assert!(matches!(
suppressed[0].url.as_str(),
"https://example.com/a" | "https://example.com/b"
));
}
#[tokio::test]
async fn crawl_concurrency_visits_all_pages() {
check(
&[
("https://example.com/", &page(&["/a", "/b", "/c", "/d"])),
("https://example.com/a", &distinct_page("a")),
("https://example.com/b", &distinct_page("b")),
("https://example.com/c", &distinct_page("c")),
("https://example.com/d", &distinct_page("d")),
],
|o| o.concurrency = 4,
|r| {
assert_eq!(r.len(), 5);
let urls: HashSet<&str> = r.iter().map(|p| p.url.as_str()).collect();
for u in [
"https://example.com/",
"https://example.com/a",
"https://example.com/b",
"https://example.com/c",
"https://example.com/d",
] {
assert!(urls.contains(u), "missing {u}");
}
},
)
.await;
}
#[tokio::test]
async fn crawl_concurrency_respects_limit() {
check(
&[
("https://example.com/", &page(&["/a", "/b", "/c", "/d"])),
("https://example.com/a", &distinct_page("a")),
("https://example.com/b", &distinct_page("b")),
("https://example.com/c", &distinct_page("c")),
("https://example.com/d", &distinct_page("d")),
],
|o| {
o.concurrency = 4;
o.limit = 3;
},
|r| assert_eq!(r.len(), 3),
)
.await;
}
#[tokio::test]
async fn crawl_concurrency_one_preserves_bfs_order() {
check(
&[
("https://example.com/", &page(&["/a", "/b"])),
("https://example.com/a", &distinct_page("a")),
("https://example.com/b", &distinct_page("b")),
],
|o| o.concurrency = 1,
|r| {
assert_eq!(r.len(), 3);
assert_eq!(r[0].url, "https://example.com/");
assert_eq!(r[1].url, "https://example.com/a");
assert_eq!(r[2].url, "https://example.com/b");
},
)
.await;
}
#[tokio::test(start_paused = true)]
async fn crawl_delay_enforces_minimum_interval() {
let start = tokio::time::Instant::now();
check(
&[
("https://example.com/", &page(&["/a", "/b"])),
("https://example.com/a", &distinct_page("a")),
("https://example.com/b", &distinct_page("b")),
],
|o| {
o.concurrency = 1;
o.delay = Some(Duration::from_millis(500));
},
|r| assert_eq!(r.len(), 3),
)
.await;
let elapsed = start.elapsed();
assert!(
elapsed >= Duration::from_secs(1),
"expected >= 1s for 3 pages with 500ms delay, got {elapsed:?}"
);
}
#[test]
fn frontier_dedup() {
let seed = Url::parse("https://example.com/").unwrap();
let mut f = Frontier::new(&seed);
assert!(!f.try_enqueue(seed, 0));
let other = Url::parse("https://example.com/page").unwrap();
assert!(f.try_enqueue(other.clone(), 1));
assert!(!f.try_enqueue(other, 1));
}
#[test]
fn frontier_pop_and_pending() {
let seed = Url::parse("https://example.com/").unwrap();
let mut f = Frontier::new(&seed);
assert_eq!(f.pending(), 1);
let (url, depth) = f.pop().unwrap();
assert_eq!(url.as_str(), "https://example.com/");
assert_eq!(depth, 0);
assert_eq!(f.pending(), 0);
assert!(f.pop().is_none());
}
#[test]
fn extract_links_filters_dangerous_schemes() {
let html = r#"<a href="https://example.com/a">A</a>
<a href="javascript:void(0)">JS</a>
<a href="JAVASCRIPT:alert(1)">JS upper</a>
<a href="data:text/html,<h1>hi</h1>">Data</a>
<a href="mailto:x@y.com">Mail</a>
<a href="/relative">Rel</a>"#;
let base = Url::parse("https://example.com/").unwrap();
let links = extract_links_from_html(html, &base);
assert_eq!(links.len(), 2);
assert_eq!(links[0].as_str(), "https://example.com/a");
assert_eq!(links[1].as_str(), "https://example.com/relative");
}
#[test]
fn error_result_fields() {
let url = Url::parse("https://example.com/fail").unwrap();
let r = error_result(&url, 2, crate::error::Error::engine("timeout", None), SystemTime::now());
assert!(matches!(r.status, CrawlStatus::Error));
assert!(r.error.as_ref().is_some_and(|e| e.report().contains("timeout")));
assert!(r.content.is_none());
}
#[test]
fn content_hash_dedup() {
let seed = Url::parse("https://example.com/").unwrap();
let mut f = Frontier::new(&seed);
assert!(!f.is_duplicate_content("unique content"));
assert!(f.is_duplicate_content("unique content"));
assert!(!f.is_duplicate_content("different content"));
}
}