Files
proxy-scraper-checker/src/checker.rs
T
2026-01-16 16:14:21 +03:00

107 lines
3.5 KiB
Rust

use std::sync::Arc;
use color_eyre::eyre::OptionExt as _;
#[cfg(feature = "tui")]
use crate::event::{AppEvent, Event};
use crate::{config::Config, proxy::Proxy, utils::pretty_error};
pub async fn check_all<R: reqwest::dns::Resolve + Clone + 'static>(
config: Arc<Config>,
dns_resolver: R,
proxies: Vec<Proxy>,
mut tls_backend: rustls::ClientConfig,
token: tokio_util::sync::CancellationToken,
#[cfg(feature = "tui")] tx: tokio::sync::mpsc::UnboundedSender<Event>,
) -> crate::Result<Vec<Proxy>> {
if config.checking.check_url.is_none() {
return Ok(proxies);
}
let workers_count =
config.checking.max_concurrent_checks.min(proxies.len());
if workers_count == 0 {
return Ok(Vec::new());
}
#[cfg(not(feature = "tui"))]
tracing::info!("Started checking {} proxies", proxies.len());
let queue = Arc::new(parking_lot::Mutex::new(proxies));
let checked_proxies = Arc::new(parking_lot::Mutex::new(Vec::new()));
tls_backend.alpn_protocols = vec![b"http/1.1".to_vec()];
let mut join_set = tokio::task::JoinSet::<()>::new();
for _ in 0..workers_count {
let queue = Arc::clone(&queue);
let config = Arc::clone(&config);
let dns_resolver = dns_resolver.clone();
let tls_backend = tls_backend.clone();
let checked_proxies = Arc::clone(&checked_proxies);
let token = token.clone();
#[cfg(feature = "tui")]
let tx = tx.clone();
join_set.spawn(async move {
tokio::select! {
biased;
res = async move {
loop {
let Some(mut proxy) = queue.lock().pop() else {
break;
};
let check_result = proxy.check(&config, dns_resolver.clone(), tls_backend.clone()).await;
#[cfg(feature = "tui")]
drop(tx.send(Event::App(AppEvent::ProxyChecked(proxy.protocol))));
match check_result {
Ok(()) => {
#[cfg(feature = "tui")]
drop(tx.send(Event::App(AppEvent::ProxyWorking(proxy.protocol))));
checked_proxies.lock().push(proxy);
}
Err(e) if tracing::event_enabled!(tracing::Level::DEBUG) => {
tracing::debug!(
"{}: {}",
proxy.to_string(true),
pretty_error(&e)
);
}
Err(_) => {}
}
}
} => res,
() = token.cancelled() => (),
}
});
}
drop(config);
drop(dns_resolver);
drop(queue);
drop(tls_backend);
drop(token);
#[cfg(feature = "tui")]
drop(tx);
while let Some(res) = join_set.join_next().await {
match res {
Ok(()) => {}
Err(e) if e.is_panic() => {
tracing::error!(
"Proxy checking task panicked: {}",
pretty_error(&e.into())
);
}
Err(e) => {
return Err(e.into());
}
}
}
drop(join_set);
Ok(Arc::into_inner(checked_proxies)
.ok_or_eyre("failed to unwrap Arc")?
.into_inner())
}