107 lines
3.5 KiB
Rust
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())
|
|
}
|