perf: reduce memory usage by ~2-3 times (#604)
This commit is contained in:
+42
-21
@@ -1,4 +1,4 @@
|
||||
use std::sync::Arc;
|
||||
use std::{collections::VecDeque, sync::Arc};
|
||||
|
||||
use color_eyre::eyre::WrapErr as _;
|
||||
|
||||
@@ -41,36 +41,57 @@ pub async fn check_all(
|
||||
storage: ProxyStorage,
|
||||
#[cfg(feature = "tui")] tx: tokio::sync::mpsc::UnboundedSender<Event>,
|
||||
) -> color_eyre::Result<ProxyStorage> {
|
||||
let semaphore = Arc::new(tokio::sync::Semaphore::new(
|
||||
config.max_concurrent_checks.min(tokio::sync::Semaphore::MAX_PERMITS),
|
||||
let workers_count = config.max_concurrent_checks.min(storage.len());
|
||||
if workers_count == 0 {
|
||||
return Ok(ProxyStorage::new(config.sources.keys().cloned().collect()));
|
||||
}
|
||||
|
||||
let queue = Arc::new(tokio::sync::Mutex::new(
|
||||
storage.into_iter().collect::<VecDeque<_>>(),
|
||||
));
|
||||
|
||||
let (result_tx, mut result_rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
|
||||
let mut join_set = tokio::task::JoinSet::new();
|
||||
for proxy in storage {
|
||||
|
||||
for _ in 0..workers_count {
|
||||
let queue = Arc::clone(&queue);
|
||||
let config = Arc::clone(&config);
|
||||
#[cfg(feature = "tui")]
|
||||
let tx = tx.clone();
|
||||
let permit = Arc::clone(&semaphore)
|
||||
.acquire_owned()
|
||||
.await
|
||||
.wrap_err("failed to acquire semaphore")?;
|
||||
let result_tx = result_tx.clone();
|
||||
join_set.spawn(async move {
|
||||
let result = check_one(
|
||||
config,
|
||||
proxy,
|
||||
#[cfg(feature = "tui")]
|
||||
tx,
|
||||
)
|
||||
.await;
|
||||
drop(permit);
|
||||
result
|
||||
loop {
|
||||
let Some(proxy) = queue.lock().await.pop_front() else {
|
||||
break Ok(());
|
||||
};
|
||||
if let Ok(proxy) = check_one(
|
||||
Arc::clone(&config),
|
||||
proxy,
|
||||
#[cfg(feature = "tui")]
|
||||
tx.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
if let Err(e) = result_tx.send(proxy) {
|
||||
break Err(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
drop(result_tx);
|
||||
|
||||
while let Some(res) = join_set.join_next().await {
|
||||
res.wrap_err("failed to join proxy checking task")?
|
||||
.wrap_err("proxy checking worker failed")?;
|
||||
}
|
||||
|
||||
let mut new_storage =
|
||||
ProxyStorage::new(config.sources.keys().cloned().collect());
|
||||
while let Some(res) = join_set.join_next().await {
|
||||
if let Ok(proxy) = res.wrap_err("failed to join proxy checking task")? {
|
||||
new_storage.insert(proxy);
|
||||
}
|
||||
while let Some(proxy) = result_rx.recv().await {
|
||||
new_storage.insert(proxy);
|
||||
}
|
||||
Ok(new_storage)
|
||||
}
|
||||
|
||||
@@ -35,6 +35,10 @@ impl ProxyStorage {
|
||||
pub fn iter(&self) -> hash_set::Iter<'_, Proxy> {
|
||||
self.proxies.iter()
|
||||
}
|
||||
|
||||
pub fn len(&self) -> usize {
|
||||
self.proxies.len()
|
||||
}
|
||||
}
|
||||
|
||||
impl IntoIterator for ProxyStorage {
|
||||
|
||||
Reference in New Issue
Block a user