From c406e1f2e1d4487f153c696bd4cb258e808d414e Mon Sep 17 00:00:00 2001 From: Almaz Date: Thu, 8 May 2025 00:12:30 +0300 Subject: [PATCH] perf: reduce memory usage by ~2-3 times (#604) --- src/checker.rs | 63 +++++++++++++++++++++++++++++++++----------------- src/storage.rs | 4 ++++ 2 files changed, 46 insertions(+), 21 deletions(-) diff --git a/src/checker.rs b/src/checker.rs index ce702a9..e06229e 100644 --- a/src/checker.rs +++ b/src/checker.rs @@ -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, ) -> color_eyre::Result { - 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::>(), )); + + 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) } diff --git a/src/storage.rs b/src/storage.rs index a5d71bc..33aaafb 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -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 {