use crate::{
Result,
config::Config,
constants::{DEFAULT_FETCHER_BROWSER, DEFAULT_TIMEOUT, DEFAULT_USER_AGENT, PAGE_LOAD_WAIT},
converter::{Converter, Format},
error::TarziError,
};
use reqwest::Client;
use tracing::{error, info, warn};
use url::Url;
use super::browser::BrowserManager;
#[derive(Debug)]
pub struct WebFetcher {
http_client: Client,
browser_manager: BrowserManager,
converter: Converter,
browser_enabled: bool,
}
impl WebFetcher {
pub fn new() -> Self {
info!("Initializing WebFetcher");
let http_client = Client::builder()
.timeout(DEFAULT_TIMEOUT)
.user_agent(DEFAULT_USER_AGENT)
.build()
.expect("Failed to create HTTP client");
info!("HTTP client created successfully for WebFetcher");
Self {
http_client,
browser_manager: BrowserManager::new(),
converter: Converter::new(),
browser_enabled: DEFAULT_FETCHER_BROWSER,
}
}
pub fn from_config(config: &Config) -> Self {
info!("Initializing WebFetcher from config");
let mut client_builder = Client::builder()
.timeout(std::time::Duration::from_secs(config.fetcher.timeout))
.user_agent(&config.fetcher.user_agent);
let proxy = crate::config::get_proxy_from_env_or_config(&config.fetcher.proxy);
if let Some(proxy) = proxy
&& !proxy.is_empty()
{
if let Ok(proxy_obj) = reqwest::Proxy::http(&proxy) {
client_builder = client_builder.proxy(proxy_obj);
info!("Using proxy from environment/config: {}", proxy);
} else {
warn!("Invalid proxy configuration: {}", proxy);
}
}
let http_client = client_builder
.build()
.expect("Failed to create HTTP client from config");
Self {
http_client,
browser_manager: BrowserManager::from_config(config),
converter: Converter::new(),
browser_enabled: config.fetcher.browser,
}
}
pub fn browser_enabled(&self) -> bool {
self.browser_enabled
}
pub async fn fetch(&mut self, url: &str, format: Format) -> Result<String> {
let raw_content = self.fetch_raw(url).await?;
let converted_content = self.converter.convert(&raw_content, format).await?;
Ok(converted_content)
}
pub async fn fetch_raw(&mut self, url: &str) -> Result<String> {
match self.fetch_plain(url).await {
Ok(content) if !content.trim().is_empty() => {
info!("Fetch succeeded via plain HTTP for {}", url);
Ok(content)
}
Ok(_) if self.browser_enabled => {
warn!(
"Plain HTTP returned empty content for {}; trying browser",
url
);
self.fetch_browser(url).await
}
Err(e) if self.browser_enabled => {
warn!("Plain HTTP failed for {}: {}; trying browser", url, e);
self.fetch_browser(url).await
}
Ok(content) => Ok(content),
Err(e) => Err(e),
}
}
pub async fn fetch_plain(&self, url: &str) -> Result<String> {
let url = Url::parse(url)?;
let response = self.http_client.get(url).send().await?;
let response = response.error_for_status()?;
let content = response.text().await?;
Ok(content)
}
pub async fn fetch_browser(&mut self, url: &str) -> Result<String> {
self.fetch_with_browser(url).await
}
pub async fn fetch_get_with_headers(
&self,
url: &str,
headers: &[(&str, &str)],
) -> Result<String> {
let url = Url::parse(url)?;
let mut request = self.http_client.get(url);
for (name, value) in headers {
request = request.header(*name, *value);
}
let response = request.send().await?;
let response = response.error_for_status()?;
Ok(response.text().await?)
}
pub async fn fetch_post_json_with_headers(
&self,
url: &str,
headers: &[(&str, &str)],
body: &serde_json::Value,
) -> Result<String> {
let url = Url::parse(url)?;
let mut request = self.http_client.post(url).json(body);
for (name, value) in headers {
request = request.header(*name, *value);
}
let response = request.send().await?;
let response = response.error_for_status()?;
Ok(response.text().await?)
}
async fn fetch_with_browser(&mut self, url: &str) -> Result<String> {
info!("Fetching URL with headless browser: {}", url);
info!("Getting or creating browser instance...");
let browser = self.browser_manager.get_or_create_browser().await?;
info!("Using existing browser instance for fetching");
info!("Navigating to URL: {}", url);
let navigation_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.get(url)).await;
match navigation_result {
Ok(Ok(_)) => {
info!("Successfully navigated to page");
}
Ok(Err(e)) => {
error!("Failed to navigate to URL: {}", e);
let error_msg = if e.to_string().contains("nssFailure")
|| e.to_string().contains("network")
{
format!(
"Network error while navigating to {url}: {e}. This may be due to network connectivity issues, firewall restrictions, or the site being temporarily unavailable."
)
} else {
format!("Failed to navigate to {url}: {e}")
};
return Err(TarziError::Browser(error_msg));
}
Err(_) => {
error!("Timeout while navigating to URL (30 seconds)");
return Err(TarziError::Browser(format!(
"Timeout while navigating to {url} (30 seconds). The page may be slow to load or the site may be experiencing issues."
)));
}
}
info!("Waiting for page to load (2 seconds)...");
tokio::time::sleep(PAGE_LOAD_WAIT).await;
info!("Wait completed");
info!("Extracting page content (dynamic DOM if available)...");
let content = match WebFetcher::get_outer_html_from(browser).await {
Ok(html) => html,
Err(e) => {
warn!(
"Falling back to page source due to error getting dynamic DOM: {}",
e
);
let content_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.source()).await;
match content_result {
Ok(Ok(content)) => content,
Ok(Err(e)) => {
error!("Failed to get page content: {}", e);
return Err(TarziError::Browser(format!("Failed to get content: {e}")));
}
Err(_) => {
error!("Timeout while extracting page content (30 seconds)");
return Err(TarziError::Browser(
"Timeout while extracting page content".to_string(),
));
}
}
}
};
info!(
"Successfully extracted page content ({} characters)",
content.len()
);
Ok(content)
}
pub async fn fetch_with_proxy(
&mut self,
url: &str,
proxy: &str,
format: Format,
) -> Result<String> {
info!("Fetching URL with proxy: {} (proxy: {})", url, proxy);
let raw_content = match self.fetch_plain_with_proxy(url, proxy).await {
Ok(content) if !content.trim().is_empty() => {
info!("Proxy fetch succeeded via plain HTTP for {}", url);
content
}
Ok(_) if self.browser_enabled => {
warn!(
"Plain HTTP via proxy returned empty content for {}; trying browser",
url
);
self.fetch_browser_with_proxy(url, proxy).await?
}
Err(e) if self.browser_enabled => {
warn!(
"Plain HTTP via proxy failed for {}: {}; trying browser",
url, e
);
self.fetch_browser_with_proxy(url, proxy).await?
}
Ok(content) => content,
Err(e) => return Err(e),
};
let converted_content = self.converter.convert(&raw_content, format).await?;
Ok(converted_content)
}
async fn fetch_plain_with_proxy(&self, url: &str, proxy: &str) -> Result<String> {
let proxy_client = match reqwest::Proxy::http(proxy) {
Ok(proxy_config) => match Client::builder()
.timeout(DEFAULT_TIMEOUT)
.user_agent(DEFAULT_USER_AGENT)
.proxy(proxy_config)
.build()
{
Ok(client) => client,
Err(e) => {
warn!(
"Failed to create HTTP client with proxy '{}': {}.",
proxy, e
);
return Err(TarziError::Config(format!(
"Failed to create proxy client: {e}"
)));
}
},
Err(e) => {
warn!("Invalid proxy URL '{}': {}.", proxy, e);
return Err(TarziError::Config(format!("Invalid proxy URL: {e}")));
}
};
let url = Url::parse(url)?;
let response = proxy_client.get(url).send().await?;
let response = response.error_for_status()?;
Ok(response.text().await?)
}
async fn fetch_browser_with_proxy(&mut self, url: &str, proxy: &str) -> Result<String> {
info!(
"Creating headless browser with proxy for fetching: {}",
proxy
);
let instance_id = self
.browser_manager
.create_browser_with_proxy(
None,
Some("proxy_browser".to_string()),
Some(proxy.to_string()),
)
.await?;
let browser = self
.browser_manager
.get_browser(&instance_id)
.ok_or_else(|| {
TarziError::Browser("Failed to get proxy browser instance".to_string())
})?;
let navigation_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.get(url)).await;
match navigation_result {
Ok(Ok(_)) => info!("Successfully navigated to page with proxy"),
Ok(Err(e)) => {
error!("Failed to navigate to URL with proxy: {}", e);
let _ = self.browser_manager.remove_browser(&instance_id).await;
return Err(TarziError::Browser(format!(
"Failed to navigate with proxy: {e}"
)));
}
Err(_) => {
error!("Timeout while navigating to URL with proxy");
let _ = self.browser_manager.remove_browser(&instance_id).await;
return Err(TarziError::Browser(
"Timeout while navigating with proxy".to_string(),
));
}
}
tokio::time::sleep(PAGE_LOAD_WAIT).await;
let content = match WebFetcher::get_outer_html_from(browser).await {
Ok(html) => html,
Err(e) => {
warn!(
"Falling back to page source (proxy) due to error getting dynamic DOM: {}",
e
);
let content_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.source()).await;
match content_result {
Ok(Ok(content)) => content,
Ok(Err(e)) => {
error!("Failed to get page content with proxy: {}", e);
let _ = self.browser_manager.remove_browser(&instance_id).await;
return Err(TarziError::Browser(format!(
"Failed to get content with proxy: {e}"
)));
}
Err(_) => {
error!("Timeout while extracting page content with proxy");
let _ = self.browser_manager.remove_browser(&instance_id).await;
return Err(TarziError::Browser(
"Timeout while extracting content with proxy".to_string(),
));
}
}
}
};
if let Err(e) = self.browser_manager.remove_browser(&instance_id).await {
warn!("Failed to cleanup proxy browser instance: {}", e);
}
Ok(content)
}
pub async fn create_browser_with_user_data(
&mut self,
user_data_dir: Option<std::path::PathBuf>,
instance_id: Option<String>,
) -> Result<String> {
self.browser_manager
.create_browser_with_user_data(user_data_dir, instance_id)
.await
}
pub async fn create_browser_with_proxy(
&mut self,
user_data_dir: Option<std::path::PathBuf>,
instance_id: Option<String>,
proxy: Option<String>,
) -> Result<String> {
self.browser_manager
.create_browser_with_proxy(user_data_dir, instance_id, proxy)
.await
}
pub fn get_browser(&self, instance_id: &str) -> Option<&thirtyfour::WebDriver> {
self.browser_manager.get_browser(instance_id)
}
pub fn get_browser_ids(&self) -> Vec<String> {
self.browser_manager.get_browser_ids()
}
pub async fn remove_browser(&mut self, instance_id: &str) -> Result<()> {
self.browser_manager.remove_browser(instance_id).await?;
Ok(())
}
pub async fn fetch_with_browser_instance(
&mut self,
url: &str,
instance_id: &str,
format: Format,
) -> Result<String> {
info!(
"Fetching URL with browser instance {}: {}",
instance_id, url
);
let browser = self
.browser_manager
.get_browser(instance_id)
.ok_or_else(|| {
TarziError::Browser(format!("Browser instance {instance_id} not found"))
})?;
info!("Using browser instance {} for fetching", instance_id);
info!(
"Navigating to URL in browser instance {}: {}",
instance_id, url
);
let navigation_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.get(url)).await;
match navigation_result {
Ok(Ok(_)) => {
info!(
"Successfully navigated to page in browser instance {}",
instance_id
);
}
Ok(Err(e)) => {
error!(
"Failed to navigate to URL in browser instance {}: {}",
instance_id, e
);
return Err(TarziError::Browser(format!("Failed to navigate: {e}")));
}
Err(_) => {
error!(
"Timeout while navigating to URL in browser instance {} (30 seconds)",
instance_id
);
return Err(TarziError::Browser(
"Timeout while navigating to URL".to_string(),
));
}
}
info!(
"Waiting for page to load in browser instance {} (2 seconds)...",
instance_id
);
tokio::time::sleep(PAGE_LOAD_WAIT).await;
info!("Wait completed for browser instance {}", instance_id);
info!(
"Extracting page content from browser instance {} (dynamic DOM if available)...",
instance_id
);
let content = match WebFetcher::get_outer_html_from(browser).await {
Ok(html) => html,
Err(e) => {
warn!(
"Falling back to page source for instance {} due to error getting dynamic DOM: {}",
instance_id, e
);
let content_result = tokio::time::timeout(DEFAULT_TIMEOUT, browser.source()).await;
match content_result {
Ok(Ok(content)) => content,
Ok(Err(e)) => {
error!(
"Failed to get page content from browser instance {}: {}",
instance_id, e
);
return Err(TarziError::Browser(format!("Failed to get content: {e}")));
}
Err(_) => {
error!(
"Timeout while extracting page content from browser instance {} (30 seconds)",
instance_id
);
return Err(TarziError::Browser(
"Timeout while extracting page content".to_string(),
));
}
}
}
};
let converted_content = self.converter.convert(&content, format).await?;
Ok(converted_content)
}
pub async fn create_multiple_browsers(
&mut self,
count: usize,
base_instance_id: Option<String>,
) -> Result<Vec<String>> {
self.browser_manager
.create_multiple_browsers(count, base_instance_id)
.await
}
pub async fn cleanup_managed_driver(&mut self) -> Result<()> {
self.browser_manager.cleanup_managed_driver().await
}
pub fn has_managed_driver(&self) -> bool {
self.browser_manager.has_managed_driver()
}
pub fn get_managed_driver_info(&self) -> Option<&super::driver::DriverInfo> {
self.browser_manager.get_managed_driver_info()
}
pub async fn shutdown(&mut self) {
self.browser_manager.shutdown().await;
}
async fn get_outer_html_from(browser: &thirtyfour::WebDriver) -> Result<String> {
if let Ok(ret) = browser
.execute(
"return document.readyState;",
Vec::<serde_json::Value>::new(),
)
.await
&& let Ok(state) = ret.convert::<String>()
&& state != "complete"
{
tokio::time::sleep(PAGE_LOAD_WAIT).await;
}
let mut attempts = 0u8;
while attempts < 30 {
if let Ok(ret) = browser
.execute(
"return document.querySelectorAll('a[href]').length;",
Vec::<serde_json::Value>::new(),
)
.await
&& let Ok(count) = ret.convert::<i64>()
&& count >= 20
{
break;
}
attempts += 1;
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
}
for _ in 0..3u8 {
let _ = browser
.execute(
"window.scrollTo({top: document.body.scrollHeight, behavior: 'instant'});",
Vec::<serde_json::Value>::new(),
)
.await;
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
let _ = browser
.execute(
"window.scrollTo({top: 0, behavior: 'instant'});",
Vec::<serde_json::Value>::new(),
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(800)).await;
}
let _ = browser
.execute(
r#"(function(){
try {
const isBraveHost = location.hostname.endsWith('search.brave.com');
if (!isBraveHost) return 0;
const results = [];
const anchors = Array.from(document.querySelectorAll('a[href]'));
const isExternal = (u) => { try { const x=new URL(u, location.origin); return x.hostname!==location.hostname; } catch(_) { return false; } };
for (const a of anchors) {
const href = a.getAttribute('href');
if (!href || !isExternal(href)) continue;
const title = (a.textContent||'').trim();
if (!title || title.length < 5) continue;
const rect = a.getBoundingClientRect();
if (!rect || rect.width === 0 || rect.height === 0) continue;
results.push({ title, url: a.href, snippet: '' });
if (results.length >= 15) break;
}
let s = document.querySelector('#tarzi-brave-results');
if (!s) { s = document.createElement('script'); s.id='tarzi-brave-results'; s.type='application/json'; document.body.appendChild(s); }
s.textContent = JSON.stringify({ results });
return results.length;
} catch(e) { return -1; }
})();"#,
Vec::<serde_json::Value>::new(),
)
.await;
let ret = browser
.execute(
"return document.documentElement.outerHTML;",
Vec::<serde_json::Value>::new(),
)
.await
.map_err(|e| TarziError::Browser(format!("execute() failed: {e}")))?;
let html: String = ret
.convert()
.map_err(|e| TarziError::Browser(format!("failed to convert script return: {e}")))?;
Ok(html)
}
}
impl Default for WebFetcher {
fn default() -> Self {
Self::new()
}
}
impl Drop for WebFetcher {
fn drop(&mut self) {
if self.browser_manager.has_browsers() || self.browser_manager.has_managed_driver() {
tracing::info!(
"WebFetcher dropped without explicit shutdown. Stopping managed driver and dropping sessions."
);
self.browser_manager.stop_managed_driver_sync();
self.browser_manager.clear_browsers();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::Config;
#[test]
fn test_webfetcher_new() {
let fetcher = WebFetcher::new();
assert!(!fetcher.browser_manager.has_browsers());
assert!(!fetcher.browser_manager.has_managed_driver());
}
#[test]
fn test_webfetcher_from_config() {
let config = Config::default();
let fetcher = WebFetcher::from_config(&config);
assert!(!fetcher.browser_manager.has_browsers());
assert!(!fetcher.browser_manager.has_managed_driver());
}
#[test]
fn test_webfetcher_with_proxy_config() {
let mut config = Config::default();
config.fetcher.proxy = Some("http://proxy.example.com:8080".to_string());
let fetcher = WebFetcher::from_config(&config);
assert!(!fetcher.browser_manager.has_browsers());
}
#[test]
fn test_webfetcher_with_custom_timeout() {
let mut config = Config::default();
config.fetcher.timeout = 60; let fetcher = WebFetcher::from_config(&config);
assert!(!fetcher.browser_manager.has_browsers());
}
#[test]
fn test_webfetcher_with_custom_user_agent() {
let mut config = Config::default();
config.fetcher.user_agent = "Test User Agent 1.0".to_string();
let fetcher = WebFetcher::from_config(&config);
assert!(!fetcher.browser_manager.has_browsers());
}
#[test]
fn test_webfetcher_default() {
let fetcher = WebFetcher::default();
assert!(!fetcher.browser_manager.has_browsers());
assert!(!fetcher.browser_manager.has_managed_driver());
}
#[test]
fn test_browser_instance_management() {
let fetcher = WebFetcher::new();
assert!(fetcher.get_browser_ids().is_empty());
assert!(fetcher.get_browser("non-existent").is_none());
}
#[test]
fn test_managed_driver_info() {
let fetcher = WebFetcher::new();
assert!(!fetcher.has_managed_driver());
assert!(fetcher.get_managed_driver_info().is_none());
}
#[tokio::test]
async fn test_invalid_url_handling() {
let mut config = Config::default();
config.fetcher.browser = false;
let mut fetcher = WebFetcher::from_config(&config);
let result = fetcher.fetch_raw("not-a-valid-url").await;
assert!(result.is_err());
if let Err(e) = result {
assert!(e.to_string().contains("relative URL without a base"));
}
}
#[tokio::test]
async fn test_url_validation() {
let mut config = Config::default();
config.fetcher.browser = false;
let mut fetcher = WebFetcher::from_config(&config);
let invalid_urls = vec![
"",
"not-a-url",
"://missing-scheme",
"http://",
"ftp://unsupported-scheme.com",
];
for invalid_url in invalid_urls {
let result = fetcher.fetch_raw(invalid_url).await;
assert!(
result.is_err(),
"Expected error for invalid URL: {invalid_url}"
);
}
}
#[test]
fn test_fetcher_browser_config() {
let fetcher = WebFetcher::new();
assert!(fetcher.browser_enabled());
let mut config = Config::default();
config.fetcher.browser = false;
let fetcher = WebFetcher::from_config(&config);
assert!(!fetcher.browser_enabled());
}
#[tokio::test]
async fn test_multiple_browser_instance_creation() {
let mut fetcher = WebFetcher::new();
let result = fetcher
.create_multiple_browsers(2, Some("test".to_string()))
.await;
match result {
Ok(_) => {
assert!(fetcher.get_browser_ids().len() <= 2);
}
Err(_) => {
let browser_count = fetcher.get_browser_ids().len();
println!("Browser count after failed creation: {browser_count}");
}
}
}
#[tokio::test]
async fn test_browser_with_user_data_dir() {
let mut fetcher = WebFetcher::new();
let temp_dir = tempfile::TempDir::new().unwrap();
let user_data_path = temp_dir.path().to_path_buf();
let result = fetcher
.create_browser_with_user_data(
Some(user_data_path),
Some("test_with_data_dir".to_string()),
)
.await;
match result {
Ok(_) => {
assert!(!fetcher.get_browser_ids().is_empty());
}
Err(_) => {
assert!(fetcher.get_browser_ids().is_empty());
}
}
}
#[tokio::test]
async fn test_browser_with_proxy() {
let mut fetcher = WebFetcher::new();
let result = fetcher
.create_browser_with_proxy(
None,
Some("test_proxy".to_string()),
Some("http://proxy.example.com:8080".to_string()),
)
.await;
match result {
Ok(_) => {
assert!(!fetcher.get_browser_ids().is_empty());
}
Err(_) => {
assert!(fetcher.get_browser_ids().is_empty());
}
}
}
#[test]
fn test_config_merging() {
let mut base_config = Config::default();
base_config.fetcher.timeout = 30;
base_config.fetcher.user_agent = "Base Agent".to_string();
let mut override_config = Config::default();
override_config.fetcher.timeout = 60;
override_config.fetcher.proxy = Some("http://proxy.example.com:8080".to_string());
base_config.merge(&override_config);
let fetcher = WebFetcher::from_config(&base_config);
assert!(!fetcher.browser_manager.has_browsers());
assert_eq!(base_config.fetcher.timeout, 60);
assert_eq!(
base_config.fetcher.proxy,
Some("http://proxy.example.com:8080".to_string())
);
assert_eq!(base_config.fetcher.user_agent, "Base Agent".to_string()); }
#[tokio::test]
async fn test_shutdown() {
let mut fetcher = WebFetcher::new();
fetcher.shutdown().await;
assert!(fetcher.get_browser_ids().is_empty());
assert!(!fetcher.has_managed_driver());
}
#[tokio::test]
async fn test_remove_nonexistent_browser() {
let mut fetcher = WebFetcher::new();
let result = fetcher.remove_browser("non-existent").await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_invalid_proxy_handling() {
let mut config = Config::default();
config.fetcher.browser = false;
let mut fetcher = WebFetcher::from_config(&config);
let test_cases = vec![
("://invalid", "Invalid proxy URL"),
("http://", "Invalid proxy URL"),
("invalid-url", "Invalid proxy URL"),
(
"http://unreachable-proxy-host:9999",
"Network error or timeout",
),
];
for (invalid_proxy, description) in test_cases {
let result = fetcher
.fetch_with_proxy("https://httpbin.org/html", invalid_proxy, Format::Html)
.await;
match result {
Err(TarziError::Config(_)) => {
println!("✓ Test passed for {invalid_proxy}: {description}");
}
Err(TarziError::Http(_)) => {
println!("✓ Test passed for {invalid_proxy}: {description} (HTTP error)");
}
Err(_) => {
println!("✓ Test passed for {invalid_proxy}: {description} (other error)");
}
Ok(_) => {
if invalid_proxy == "://invalid" || invalid_proxy == "http://" {
panic!("Expected error for clearly invalid proxy: {invalid_proxy}");
} else {
println!(
"ℹ Test passed for {invalid_proxy}: {description} (unexpected success, but acceptable)"
);
}
}
}
}
}
#[test]
fn test_drop_warning() {
let _fetcher = WebFetcher::new();
}
}