From 6c5565959579e19bd58c42a97c4b84bb53cc9b9f Mon Sep 17 00:00:00 2001 From: Danlian Akhmedzianov Date: Wed, 10 Jun 2026 17:49:49 +0300 Subject: [PATCH 1/2] Remove complex perf testing --- Cargo.lock | 39 ------ Cargo.toml | 1 - src/cmd/commands.rs | 4 - src/cmd/mod.rs | 1 - src/cmd/perf.rs | 182 -------------------------- src/cmd/select.rs | 93 +------------ src/entities/mod.rs | 3 - src/entities/perf_result.rs | 97 -------------- src/entities/proxied_client.rs | 123 ----------------- src/entities/xray_server.rs | 65 +-------- src/entities/xray_server_with_perf.rs | 35 ----- src/main.rs | 53 +------- src/services/dns_client.rs | 91 ------------- src/services/measures.rs | 55 -------- src/services/mod.rs | 3 - src/services/portman.rs | 84 ------------ src/utils/mod.rs | 1 - src/utils/urldecode.rs | 17 --- 18 files changed, 10 insertions(+), 937 deletions(-) delete mode 100644 src/cmd/perf.rs delete mode 100644 src/entities/perf_result.rs delete mode 100644 src/entities/proxied_client.rs delete mode 100644 src/entities/xray_server_with_perf.rs delete mode 100644 src/services/dns_client.rs delete mode 100644 src/services/measures.rs delete mode 100644 src/services/portman.rs delete mode 100644 src/utils/urldecode.rs diff --git a/Cargo.lock b/Cargo.lock index 9b95d6b..ed9ea2c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -177,17 +177,6 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" -[[package]] -name = "chacha20" -version = "0.10.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6f8d983286843e49675a4b7a2d174efe136dc93a18d69130dd18198a6c167601" -dependencies = [ - "cfg-if 1.0.4", - "cpufeatures", - "rand_core 0.10.1", -] - [[package]] name = "chomp1" version = "0.3.4" @@ -321,15 +310,6 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" -[[package]] -name = "cpufeatures" -version = "0.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" -dependencies = [ - "libc", -] - [[package]] name = "custom_derive" version = "0.1.7" @@ -626,7 +606,6 @@ dependencies = [ "cfg-if 1.0.4", "libc", "r-efi 6.0.0", - "rand_core 0.10.1", "wasip2", "wasip3", ] @@ -1264,17 +1243,6 @@ dependencies = [ "rand_core 0.9.5", ] -[[package]] -name = "rand" -version = "0.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d2e8e8bcc7961af1fdac401278c6a831614941f6164ee3bf4ce61b7edb162207" -dependencies = [ - "chacha20", - "getrandom 0.4.2", - "rand_core 0.10.1", -] - [[package]] name = "rand_chacha" version = "0.3.1" @@ -1313,12 +1281,6 @@ dependencies = [ "getrandom 0.3.4", ] -[[package]] -name = "rand_core" -version = "0.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" - [[package]] name = "rayconf" version = "0.5.1" @@ -1333,7 +1295,6 @@ dependencies = [ "http", "percent-encoding", "querystring", - "rand 0.10.1", "reqwest", "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 66de2d8..b41b776 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -19,7 +19,6 @@ querystring = "1.1.0" tokio = { version = "1.48.0", features = ["full"] } dnsclient = "0.2.0" futures = "0.3.32" -rand = "0.10.1" clap_complete = "4.6.5" [profile.release] diff --git a/src/cmd/commands.rs b/src/cmd/commands.rs index 4c98812..21999d6 100644 --- a/src/cmd/commands.rs +++ b/src/cmd/commands.rs @@ -1,6 +1,5 @@ use crate::cmd::add::AddParams; use crate::cmd::completion::CompletionParams; -use crate::cmd::perf::PerfParams; use crate::cmd::remote::RemoteParams; use crate::cmd::remove::RemoveParams; use crate::cmd::select::SelectParams; @@ -20,9 +19,6 @@ pub(crate) enum Commands { /// Remote URLs (subscription link) Remote(RemoteParams), - #[clap(visible_aliases = ["p, performance"], about = "Performance analyzis")] - Perf(PerfParams), - #[clap(about = "Prepare and print shell completion script")] Completion(CompletionParams), } diff --git a/src/cmd/mod.rs b/src/cmd/mod.rs index 1cabc3f..dbf68fe 100644 --- a/src/cmd/mod.rs +++ b/src/cmd/mod.rs @@ -2,7 +2,6 @@ pub(crate) mod add; pub(crate) mod cli; pub(crate) mod commands; pub(crate) mod completion; -pub(crate) mod perf; pub(crate) mod remote; pub(crate) mod remove; pub(crate) mod select; diff --git a/src/cmd/perf.rs b/src/cmd/perf.rs deleted file mode 100644 index d39cfb3..0000000 --- a/src/cmd/perf.rs +++ /dev/null @@ -1,182 +0,0 @@ -use std::{process::Stdio, time::Duration}; - -use clap::Args; -use reqwest::Client; -use tokio::{io::AsyncWriteExt, process::Child, sync::OnceCell}; - -use crate::entities::perf_result::{PerfData, PerfResult}; -use crate::entities::proxied_client::ProxiedClient; -use crate::entities::remote::Remote; -use crate::entities::xray_server::XRayServer; -use crate::services::config::Config; -use crate::services::measures::Measures; -use crate::services::portman::Portman; -use crate::utils::tap::Tap; -use crate::utils::urldecode::URLDecode; -use crate::v2parser::parser::create_json_config; - -static FIFTY_MB_IN_BYTES: u128 = 50 * 1024 * 1024; -static UNKNOWN_SERVER_NAME: &'static str = "Unknown"; -static PROXY_HOST: &'static str = "127.0.0.1"; -static PROXY_TEST_TIMEOUT: Duration = Duration::from_mins(5); -static DELAY_BETWEEN_ATTEMPS: Duration = Duration::from_millis(200); -static PROXY_TEST_ATTEMPTS: usize = 5; - -#[derive(Args)] -pub(crate) struct PerfParams { - #[clap(skip)] - portman: OnceCell, - - #[clap(skip)] - measures: OnceCell, -} - -impl PerfParams { - async fn select_remote(&self) -> Vec { - let config = Config::read_or_default().await; - let remotes: Vec = config - .remotes() - .into_iter() - .map(ToOwned::to_owned) - .collect(); - - if remotes.is_empty() { - return vec![]; - } - - let remote_idx_result = dialoguer::FuzzySelect::new() - .with_prompt("Select remote subscription URL") - .items(&remotes) - .tap(|s| match &remotes.is_empty() { - true => s, - false => s.default(0), - }) - .interact(); - - let Ok(remote_idx) = remote_idx_result else { - return vec![]; - }; - - let selected_remote = remotes - .get(remote_idx) - .ok_or_else(|| anyhow::Error::msg("No remote URL")); - - let Ok(selected_remote) = selected_remote else { - return vec![]; - }; - - let Ok(client) = Client::builder().build() else { - return vec![]; - }; - let Ok(response) = client.get(selected_remote.url()).send().await else { - return vec![]; - }; - let Ok(content) = response.text().await else { - return vec![]; - }; - - selected_remote - .decoder() - .decode(content) - .iter() - .filter_map(|u| XRayServer::try_from(u.to_string()).ok()) - .collect::>() - } - - async fn get_measures(&self) -> &Measures { - self.measures - .get_or_init(async || Measures::read_or_default().await) - .await - } - - async fn get_portman(&self) -> &Portman { - self.portman.get_or_init(async || Portman::new()).await - } - - async fn measure_server_speed(&self, server: &XRayServer) -> anyhow::Result { - server - .measure_rtt(Duration::from_secs(5)) - .await - .ok_or_else(|| anyhow::anyhow!("Cannot connect to a server"))?; - let portman = self.get_portman().await; - let url = server.url(); - - let result = portman - .lease_port(async move |port| -> anyhow::Result { - let config_json = create_json_config(url.as_str(), Some(port), None, None)?; - let mut command = self.run_xray(config_json).await?; - let measure_result = ProxiedClient::try_new(PROXY_HOST, port, PROXY_TEST_TIMEOUT)? - .measure_first_successful( - FIFTY_MB_IN_BYTES, - DELAY_BETWEEN_ATTEMPS, - PROXY_TEST_ATTEMPTS, - ) - .await; - command.kill().await?; - measure_result - }) - .await; - - if let Ok(result) = result { - let perf_data = PerfData::from_raw_data(FIFTY_MB_IN_BYTES as u64, result); - let perf_result = PerfResult::now(&url, Some(perf_data)); - self.get_measures() - .await - .add_measure(&url, perf_result) - .await; - } - - result - } - - pub async fn start(&self) -> anyhow::Result<()> { - let servers = self.select_remote().await; - - for server in servers.iter() { - let result = self.measure_server_speed(server).await; - let server_display_name = server - .url() - .fragment() - .unwrap_or_else(|| UNKNOWN_SERVER_NAME) - .decode_as_urlencoded() - .unwrap_or_else(|_| UNKNOWN_SERVER_NAME.to_string()); - let output = match result { - Ok(result) => format!( - "Server {server_display_name} responded with time: {} ms", - result.as_millis() - ), - Err(error) => format!( - "Server {server_display_name} did not responded, error: {}", - error.to_string() - ), - }; - println!("{output}"); - } - - self.get_measures().await.dump().await?; - - Ok(()) - } - - async fn run_xray(&self, config: impl AsRef) -> anyhow::Result { - static BIN_NAME: &'static str = "xray"; - let mut command = tokio::process::Command::new(BIN_NAME) - .kill_on_drop(true) - .stderr(Stdio::null()) - .stdout(Stdio::null()) - .stdin(Stdio::piped()) - .spawn() - .map_err(|_| anyhow::anyhow!("Failed to run XRay, please make sure It is installed"))?; - - let mut child_stdin = command - .stdin - .take() - .ok_or_else(|| anyhow::anyhow!("No child stdin"))?; - - child_stdin.write_all(config.as_ref().as_bytes()).await?; - - drop(child_stdin); - - Ok(command) - } -} diff --git a/src/cmd/select.rs b/src/cmd/select.rs index 22c6eaf..8d0b2a6 100644 --- a/src/cmd/select.rs +++ b/src/cmd/select.rs @@ -1,17 +1,11 @@ -use std::process::Stdio; -use std::{cmp::Ordering, collections::HashMap}; - use crate::entities::remote::Remote; use crate::entities::xray_server::XRayServer; -use crate::entities::xray_server_with_perf::XRayServerWithPerf; use crate::services::config::Config; -use crate::services::measures::Measures; use crate::utils::tap::Tap; use crate::v2parser::entities::log::{Log, LogLevel}; use crate::v2parser::parser::create_json_config; use clap::Args; use reqwest::Client; -use tokio::io::AsyncWriteExt; static DEFAULT_HTTP_PORT: u16 = 8080; static DEFAULT_SOCKS_PORT: u16 = 1080; @@ -27,14 +21,6 @@ pub(crate) struct SelectParams { )] remote: bool, - #[arg( - long, - short, - default_value_t = false, - help = "Do not run XRay, just print config" - )] - dry_run: bool, - #[arg(long, conflicts_with = "http_port", help = format!("SOCKS5 port, default: {DEFAULT_SOCKS_PORT}"))] socks_port: Option, @@ -49,51 +35,6 @@ pub(crate) struct SelectParams { } impl SelectParams { - async fn sort_servers(&self, servers: Vec) -> Vec { - let measures = Measures::read_or_default().await; - - let mut measure_results = Vec::with_capacity(servers.len()); - for server in servers.iter() { - let measure_result = measures.get_latest_measure(server.url()).await; - measure_results.push(measure_result); - } - - let mut measure_results_index = HashMap::with_capacity(measure_results.len()); - for measure_result in measure_results { - if let Some(measure_result) = measure_result { - measure_results_index.insert(measure_result.url().to_owned(), measure_result); - } - } - - let mut servers = servers.clone(); - - servers.sort_by(|server_a, server_b| { - let (speed_a, speed_b) = ( - measure_results_index.get(server_a.url()), - measure_results_index.get(server_b.url()), - ); - - match (speed_a, speed_b) { - (None, None) => Ordering::Equal, - (None, Some(_)) => Ordering::Greater, - (Some(_), None) => Ordering::Less, - (Some(a), Some(b)) => match (a.perf_data(), b.perf_data()) { - (None, None) => Ordering::Equal, - (None, Some(_)) => Ordering::Greater, - (Some(_), None) => Ordering::Less, - (Some(a), Some(b)) => b.download_speed_mbps().cmp(&a.download_speed_mbps()), - }, - } - }); - - servers - .into_iter() - .map(|ref s| { - XRayServerWithPerf::new(s.clone(), measure_results_index.get(s.url()).cloned()) - }) - .collect() - } - async fn select_local(&self) -> anyhow::Result { let config = Config::read_or_default().await; let items: Vec = config @@ -163,8 +104,6 @@ impl SelectParams { .filter_map(|u| XRayServer::try_from(u.to_string()).ok()) .collect::>(); - let xray_servers = self.sort_servers(xray_servers).await; - if xray_servers.is_empty() { return Err(anyhow::Error::msg("This URL has note XRay servers")); } @@ -182,7 +121,7 @@ impl SelectParams { .get(server_idx) .ok_or_else(|| anyhow::Error::msg("No XRay server selected"))?; - Ok(selected_xray_server.xray_server().url().to_string()) + Ok(selected_xray_server.url().to_string()) } fn get_log(&self) -> Log { @@ -217,35 +156,7 @@ impl SelectParams { let config_json = create_json_config(url.as_str(), socks_port, http_port, Some(self.get_log()))?; - match self.dry_run { - true => { - println!("{}", config_json); - Ok(()) - } - false => self.run_xray(config_json).await, - } - } - - async fn run_xray(&self, config: impl AsRef) -> anyhow::Result<()> { - static BIN_NAME: &'static str = "xray"; - let mut command = tokio::process::Command::new(BIN_NAME) - .kill_on_drop(true) - .stderr(Stdio::inherit()) - .stdout(Stdio::inherit()) - .stdin(Stdio::piped()) - .spawn() - .map_err(|_| anyhow::anyhow!("Failed to run XRay, please make sure It is installed"))?; - - let mut child_stdin = command - .stdin - .take() - .ok_or_else(|| anyhow::anyhow!("No child stdin"))?; - - child_stdin.write_all(config.as_ref().as_bytes()).await?; - - drop(child_stdin); - - command.wait().await?; + println!("{}", config_json); Ok(()) } diff --git a/src/entities/mod.rs b/src/entities/mod.rs index 45cd889..ac921cf 100644 --- a/src/entities/mod.rs +++ b/src/entities/mod.rs @@ -1,6 +1,3 @@ -pub(crate) mod perf_result; -pub(crate) mod proxied_client; pub(crate) mod remote; pub(crate) mod remote_decoder; pub(crate) mod xray_server; -pub(crate) mod xray_server_with_perf; diff --git a/src/entities/perf_result.rs b/src/entities/perf_result.rs deleted file mode 100644 index 755f788..0000000 --- a/src/entities/perf_result.rs +++ /dev/null @@ -1,97 +0,0 @@ -use std::time::Duration; - -use serde::{Deserialize, Serialize}; -use tokio::time::Instant; -use url::Url; - -#[derive(Clone, Copy, Debug, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct PerfData { - download_speed_mbps: u64, -} - -impl PerfData { - pub fn from_raw_data(bytes_downloaded: u64, time_taken: Duration) -> Self { - let download_speed_mbps = bytes_downloaded * 8 / time_taken.as_secs() / 1_000_000; - Self { - download_speed_mbps, - } - } - - #[allow(unused)] - pub fn new(download_speed_mbps: u64) -> Self { - Self { - download_speed_mbps, - } - } - - #[allow(unused)] - pub fn download_speed_mbps(&self) -> u64 { - return self.download_speed_mbps; - } -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct PerfResult { - url: url::Url, - - #[serde(with = "approx_instant")] - date: Instant, - - perf_data: Option, -} - -impl PerfResult { - pub fn now(url: &Url, perf_data: Option) -> Self { - Self { - url: url.to_owned(), - date: Instant::now(), - perf_data, - } - } - - #[allow(unused)] - pub fn url(&self) -> &url::Url { - &self.url - } - - #[allow(unused)] - pub fn date(&self) -> Instant { - self.date - } - - #[allow(unused)] - pub fn perf_data(&self) -> Option { - return self.perf_data; - } -} - -mod approx_instant { - use std::time::SystemTime; - - use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error}; - use tokio::time::Instant; - - pub fn serialize(instant: &Instant, serializer: S) -> Result - where - S: Serializer, - { - let system_now = SystemTime::now(); - let instant_now = Instant::now(); - let approx = system_now - (instant_now - *instant); - approx.serialize(serializer) - } - - pub fn deserialize<'de, D>(deserializer: D) -> Result - where - D: Deserializer<'de>, - { - let de = SystemTime::deserialize(deserializer)?; - let system_now = SystemTime::now(); - let instant_now = Instant::now(); - let duration = system_now.duration_since(de).map_err(Error::custom)?; - let approx = instant_now - duration; - Ok(approx) - } -} diff --git a/src/entities/proxied_client.rs b/src/entities/proxied_client.rs deleted file mode 100644 index 257c3a0..0000000 --- a/src/entities/proxied_client.rs +++ /dev/null @@ -1,123 +0,0 @@ -use std::time::Duration; - -use reqwest::{Client, ClientBuilder, Proxy}; -use tokio::time::Instant; -use url::Url; - -pub(crate) struct ProxiedClient(Client); - -impl ProxiedClient { - #[inline] - fn proxy_url(proxy_host: impl AsRef, proxy_port: u16) -> String { - let proxy_host = proxy_host.as_ref(); - format!("socks5h://{proxy_host}:{proxy_port}") - } - - #[inline] - fn download_url(bytes_len: u128) -> anyhow::Result { - static DOWNLOAD_URL_BASE: &'static str = "https://speed.cloudflare.com/__down"; - let mut download_url = Url::parse(DOWNLOAD_URL_BASE)?; - - let bytes_kv = format!("bytes={bytes_len}"); - download_url.set_query(Some(bytes_kv.as_str())); - - Ok(download_url) - } - - pub fn try_new( - proxy_host: impl AsRef, - proxy_port: u16, - timeout: Duration, - ) -> anyhow::Result { - let proxy_url = Self::proxy_url(proxy_host.as_ref(), proxy_port); - let proxy = Proxy::all(proxy_url)?; - let client = ClientBuilder::new() - .proxy(proxy) - .no_gzip() - .no_brotli() - .no_deflate() - .timeout(timeout) - .build()?; - - let instance = Self(client); - - Ok(instance) - } - - #[must_use] - #[allow(unused)] - pub async fn measure_first_successful( - &self, - bytes_len: u128, - delay: Duration, - attempts: usize, - ) -> anyhow::Result { - for _ in 0..attempts { - tokio::time::sleep(delay).await; - match self.measure_attempt(bytes_len).await { - Ok(duration) => return Ok(duration), - _ => {} - } - } - - anyhow::bail!("Cannot estimate") - } - - #[must_use] - #[allow(unused)] - pub async fn measure_median( - &self, - bytes_len: u128, - delay: Duration, - attempts: usize, - ) -> anyhow::Result { - let mut results = Vec::with_capacity(attempts); - for _ in 0..attempts { - tokio::time::sleep(delay).await; - match self.measure_attempt(bytes_len).await { - Ok(duration) => results.push(duration), - _ => {} - } - } - - if results.len() == 0 { - anyhow::bail!("Cannot estimate, no successful measures") - } - - results.sort_by_key(|v| v.as_millis()); - - results - .get(results.len() / 2) - .map(|v| *v) - .ok_or_else(|| anyhow::anyhow!("Cannot estimate")) - } - - async fn measure_attempt(&self, bytes_len: u128) -> anyhow::Result { - let download_url = Self::download_url(bytes_len)?; - - let request = self - .0 - .get(download_url) - .header(reqwest::header::ACCEPT_ENCODING, "identity") - .build()?; - - let start = Instant::now(); - let mut response = self.0.execute(request).await?; - - let mut bytes_read = 0u128; - - while let Some(chunk) = response.chunk().await? { - bytes_read += chunk.len() as u128; - } - - let elapsed = start.elapsed(); - - if bytes_len != bytes_read { - let err_msg = format!("requested bytes: {bytes_len}, but downloaded: {bytes_read}"); - let err = anyhow::anyhow!(err_msg); - return Err(err); - } - - Ok(elapsed) - } -} diff --git a/src/entities/xray_server.rs b/src/entities/xray_server.rs index 7d9248f..2c8bc73 100644 --- a/src/entities/xray_server.rs +++ b/src/entities/xray_server.rs @@ -1,12 +1,7 @@ -use crate::services::dns_client; use percent_encoding::percent_decode_str; use serde::{Deserialize, Serialize}; use std::fmt::Display; -use std::net::{IpAddr, SocketAddr}; -use std::sync::Arc; -use std::time::Duration; -use tokio::time::Instant; -use url::{Host, Url}; +use url::Url; const UNNAMED: &'static str = "Unnamed XRay Server"; @@ -18,64 +13,6 @@ impl XRayServer { pub fn url(&self) -> &Url { &self.0 } - - async fn measure_rtt_for_addr( - &self, - timeout: Duration, - addr: &SocketAddr, - ) -> anyhow::Result { - let start = Instant::now(); - match tokio::time::timeout(timeout, tokio::net::TcpStream::connect(addr)).await { - Ok(Ok(stream)) => { - let result = start.elapsed(); - drop(stream); - Ok(result) - } - Ok(Err(err)) => Err(err.into()), - _ => anyhow::bail!("Timeout"), - } - } - - async fn measure_rtt_for_addrs( - &self, - timeout: Duration, - addrs: impl Iterator, - ) -> anyhow::Result { - let mut durations = Vec::new(); - for addr in addrs { - let duration = self.measure_rtt_for_addr(timeout, &addr).await; - durations.push(duration); - } - if durations.is_empty() { - anyhow::bail!("Empty addrs list"); - } - let mut ok_durations = durations - .into_iter() - .filter_map(|r| r.ok()) - .collect::>(); - ok_durations.sort_by_key(|v| v.as_millis()); - - ok_durations - .get((ok_durations.len() / 2) as usize) - .cloned() - .ok_or_else(|| anyhow::anyhow!("Cannot estimate")) - } - - pub async fn measure_rtt(&self, timeout: Duration) -> Option { - let Some((host, port)) = self.url().host().zip(self.0.port()) else { - return None; - }; - let ip_addresses = match host { - Host::Domain(dns_name) => dns_client::resolve(dns_name).await, - Host::Ipv4(ipv4) => Arc::new(vec![IpAddr::V4(ipv4)]), - Host::Ipv6(ipv6) => Arc::new(vec![IpAddr::V6(ipv6)]), - }; - let socket_addrs = ip_addresses - .iter() - .map(|v| SocketAddr::new(v.to_owned(), port)); - - self.measure_rtt_for_addrs(timeout, socket_addrs).await.ok() - } } impl Display for XRayServer { diff --git a/src/entities/xray_server_with_perf.rs b/src/entities/xray_server_with_perf.rs deleted file mode 100644 index 53ca18e..0000000 --- a/src/entities/xray_server_with_perf.rs +++ /dev/null @@ -1,35 +0,0 @@ -use std::fmt::Display; - -use crate::entities::{perf_result::PerfResult, xray_server::XRayServer}; - -pub(crate) struct XRayServerWithPerf(XRayServer, Option); - -impl XRayServerWithPerf { - pub fn new(xray_server: XRayServer, duration: Option) -> Self { - Self(xray_server, duration) - } - - #[allow(unused)] - pub fn xray_server(&self) -> &XRayServer { - &self.0 - } - - #[allow(unused)] - pub fn duration(&self) -> Option { - return self.1.clone(); - } -} - -impl Display for XRayServerWithPerf { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let server_name = self.0.to_string(); - let duration = self - .1 - .clone() - .and_then(|v| v.perf_data()) - .map(|v| format!("{} mbps", v.download_speed_mbps())) - .unwrap_or_else(|| "n/a".to_string()); - let title = format!("{}, {}", server_name, duration); - f.write_str(title.as_str()) - } -} diff --git a/src/main.rs b/src/main.rs index 298f901..4f468fb 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,5 +1,3 @@ -use std::process::ExitCode; - use crate::cmd::cli::Cli; use crate::cmd::commands::Commands; use clap::Parser; @@ -11,51 +9,14 @@ mod utils; mod v2parser; #[tokio::main] -async fn main() -> ExitCode { +async fn main() -> anyhow::Result<()> { let cli = Cli::parse(); - let command = async move { - match cli.command { - Commands::Add(add_params) => add_params.add().await, - Commands::Delete(delete_params) => delete_params.remove().await, - Commands::Select(select_params) => select_params.select().await, - Commands::Remote(remote_params) => remote_params.handle_action().await, - Commands::Perf(perf_params) => perf_params.start().await, - Commands::Completion(completion_params) => completion_params.generate().await, - } - }; - - wait_for_exit(command).await -} - -async fn wait_for_exit( - fut: impl Future>, -) -> ExitCode { - tokio::select! { - result = fut => { - match result { - Ok(_) => { - ExitCode::SUCCESS - }, - Err(error) => { - eprintln!("\nRayconf finished with error: {error}"); - ExitCode::FAILURE - } - } - } - - signal = tokio::signal::ctrl_c() => { - match signal { - Ok(()) => { - eprintln!("\nGraceful shutdown"); - ExitCode::from(130) - } - - Err(error) => { - eprintln!("\nfailed to listen for Ctrl+C: {error}"); - ExitCode::FAILURE - } - } - } + match cli.command { + Commands::Add(add_params) => add_params.add().await, + Commands::Delete(delete_params) => delete_params.remove().await, + Commands::Select(select_params) => select_params.select().await, + Commands::Remote(remote_params) => remote_params.handle_action().await, + Commands::Completion(completion_params) => completion_params.generate().await, } } diff --git a/src/services/dns_client.rs b/src/services/dns_client.rs deleted file mode 100644 index 3731d97..0000000 --- a/src/services/dns_client.rs +++ /dev/null @@ -1,91 +0,0 @@ -use dnsclient::UpstreamServer; -use dnsclient::r#async::DNSClient as AsyncDNSClient; -use std::collections::HashMap; -use std::net::{IpAddr, SocketAddr}; -use std::sync::{Arc, LazyLock}; -use tokio::sync::{Mutex, OnceCell}; - -static DNS_CLIENTS_RAW: &'static str = include_str!("./dns_servers.txt"); -static SEP: &'static str = "\n"; - -struct DNSClient { - client: AsyncDNSClient, - cache: Arc>>>>>>, -} - -impl DNSClient { - fn try_new() -> anyhow::Result { - let dns_servers: Vec = DNS_CLIENTS_RAW - .split(SEP) - .filter_map(|v| -> Option { v.parse().ok() }) - .map(|v| UpstreamServer::new(v)) - .collect(); - - if dns_servers.is_empty() { - return Err(anyhow::Error::msg("No DNS servers found")); - } - - let client = AsyncDNSClient::new(dns_servers); - - Ok(Self { - client, - cache: Arc::default(), - }) - } - - async fn get_all_addrs(&self, dns_name: impl AsRef) -> Vec { - let ipv4 = self - .client - .query_a(dns_name.as_ref()) - .await - .unwrap_or_else(|_| vec![]) - .into_iter() - .map(|v| IpAddr::V4(v)) - .collect::>(); - let ipv6 = self - .client - .query_aaaa(dns_name.as_ref()) - .await - .unwrap_or_else(|_| vec![]) - .into_iter() - .map(|v| IpAddr::V6(v)) - .collect::>(); - let mut result = Vec::with_capacity(ipv4.len() + ipv6.len()); - - result.extend(ipv4); - result.extend(ipv6); - - result - } - - async fn resolve_or_insert( - &self, - dns_name: impl AsRef, - ) -> Arc>>> { - let mut cache = self.cache.lock().await; - let result = cache - .entry(dns_name.as_ref().to_owned()) - .or_insert_with(|| Arc::new(OnceCell::new())); - - result.clone() - } - - pub async fn resolve(&self, dns_name: impl AsRef) -> Arc> { - let cell = self.resolve_or_insert(dns_name.as_ref()).await; - cell.get_or_init(async move || { - let addrs = self.get_all_addrs(dns_name.as_ref()).await; - Arc::new(addrs) - }) - .await - .clone() - } -} - -static DEFAULT_DNS_CLIENT: LazyLock> = LazyLock::new(DNSClient::try_new); - -pub(crate) async fn resolve(dns_name: impl AsRef) -> Arc> { - match DEFAULT_DNS_CLIENT.as_ref() { - Err(_) => Arc::default(), - Ok(client) => client.resolve(dns_name).await, - } -} diff --git a/src/services/measures.rs b/src/services/measures.rs deleted file mode 100644 index 94890de..0000000 --- a/src/services/measures.rs +++ /dev/null @@ -1,55 +0,0 @@ -use std::collections::{HashMap, LinkedList}; -use std::sync::Arc; - -use serde::{Deserialize, Serialize}; -use tokio::sync::RwLock; -use url::Url; - -use crate::entities::perf_result::PerfResult; -use crate::services::fileman::FileMan; - -static DEFAULT_MEASURES_FILENAME: &'static str = "measures.json"; - -#[derive(Clone, Debug, Serialize, Deserialize, Default)] -#[serde(rename_all = "camelCase")] -struct MeasuresData { - measures: HashMap>, -} - -pub struct Measures { - data: Arc>, -} - -impl Measures { - pub async fn read_or_default() -> Self { - let measures_data = FileMan::read_data(DEFAULT_MEASURES_FILENAME) - .await - .ok() - .and_then(|v| serde_json::from_slice::(v.as_ref()).ok()) - .unwrap_or_default(); - let data = Arc::new(RwLock::new(measures_data)); - Self { data } - } - - pub async fn add_measure(&self, url: &Url, measure: PerfResult) { - let mut data = self.data.write().await; - data.measures - .entry(url.to_owned()) - .or_insert_with(Default::default) - .push_front(measure); - } - - pub async fn get_latest_measure(&self, url: &Url) -> Option { - let data = self.data.read().await; - data.measures.get(url)?.front().cloned() - } - - pub async fn dump(&self) -> anyhow::Result<()> { - let data = self.data.read().await; - let measures_data = (*data).clone(); - let bytes = serde_json::to_string_pretty(&measures_data)?; - FileMan::save_data(DEFAULT_MEASURES_FILENAME, bytes).await?; - - Ok(()) - } -} diff --git a/src/services/mod.rs b/src/services/mod.rs index 074382d..ec94232 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,5 +1,2 @@ pub(crate) mod config; -pub(crate) mod dns_client; pub(crate) mod fileman; -pub(crate) mod measures; -pub(crate) mod portman; diff --git a/src/services/portman.rs b/src/services/portman.rs deleted file mode 100644 index 93fe063..0000000 --- a/src/services/portman.rs +++ /dev/null @@ -1,84 +0,0 @@ -use std::collections::HashSet; -use std::sync::{Arc, Mutex}; - -#[must_use] -pub(crate) struct PortLease(u16, F) -where - F: FnMut(); - -impl PortLease -where - F: FnMut(), -{ - fn new(port: u16, on_drop: F) -> Self { - PortLease(port, on_drop) - } - - pub fn port(&self) -> u16 { - self.0 - } -} - -impl Drop for PortLease -where - F: FnMut(), -{ - fn drop(&mut self) { - self.1(); - } -} - -pub(crate) struct Portman { - ports: Arc>>, -} - -impl Portman { - pub fn new() -> Portman { - Self { - ports: Arc::default(), - } - } - - fn hold(&self, port: u16) -> anyhow::Result<()> { - let mut ports = match self.ports.lock() { - Ok(ports) => ports, - Err(e) => e.into_inner(), - }; - if ports.contains(&port) { - anyhow::bail!("Port already taken") - } - ports.insert(port); - Ok(()) - } - - fn get_random_port() -> u16 { - rand::random_range(10000..65535) - } - - fn hold_random(&self) -> PortLease { - let mut random = Self::get_random_port(); - while self.hold(random).is_err() { - random = Self::get_random_port(); - } - - let ports = self.ports.clone(); - PortLease::new(random, move || { - let mut ports = match ports.lock() { - Ok(ports) => ports, - Err(e) => e.into_inner(), - }; - ports.remove(&random); - }) - } - - pub async fn lease_port(&self, doer: F) -> R - where - Fut: Future, - F: FnOnce(u16) -> Fut, - { - let lease = self.hold_random(); - let result = doer(lease.port()).await; - - result - } -} diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 270cbca..aa8c44b 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -1,3 +1,2 @@ pub(crate) mod tap; pub(crate) mod to_err; -pub(crate) mod urldecode; diff --git a/src/utils/urldecode.rs b/src/utils/urldecode.rs deleted file mode 100644 index d0167e1..0000000 --- a/src/utils/urldecode.rs +++ /dev/null @@ -1,17 +0,0 @@ -pub(crate) trait URLDecode { - fn decode_as_urlencoded(&self) -> anyhow::Result; -} - -impl URLDecode for String { - fn decode_as_urlencoded(&self) -> anyhow::Result { - let result = urlencoding::decode(self).map(|v| v.to_owned().to_string())?; - Ok(result) - } -} - -impl URLDecode for &str { - fn decode_as_urlencoded(&self) -> anyhow::Result { - let result = urlencoding::decode(self).map(|v| v.to_owned().to_string())?; - Ok(result) - } -} From 280f497d9bf0f50b0dd786053c408037c2822378 Mon Sep 17 00:00:00 2001 From: Danlian Akhmedzianov Date: Wed, 10 Jun 2026 19:45:56 +0300 Subject: [PATCH 2/2] Measure RTT and sort XRay servers for selection - Add async `measure_rtt` to `XRayServer` and helper functions to probe TCP latency. - Introduce `XRayServerWithDuration` to display server name with measured RTT. - Implement `sort_servers` in the select command to order servers by latency. - Create a DNS client service with caching for domain resolution. - Update imports, module exports, and related code to use the new functionality. --- src/cmd/select.rs | 45 ++++++++++++++++- src/entities/xray_server.rs | 96 +++++++++++++++++++++++++++++++++++- src/services/dns_client.rs | 91 ++++++++++++++++++++++++++++++++++ src/services/dns_servers.txt | 2 +- src/services/mod.rs | 1 + 5 files changed, 230 insertions(+), 5 deletions(-) create mode 100644 src/services/dns_client.rs diff --git a/src/cmd/select.rs b/src/cmd/select.rs index 8d0b2a6..2218f13 100644 --- a/src/cmd/select.rs +++ b/src/cmd/select.rs @@ -1,10 +1,15 @@ +use std::cmp::Ordering; +use std::collections::HashMap; +use std::time::Duration; + use crate::entities::remote::Remote; -use crate::entities::xray_server::XRayServer; +use crate::entities::xray_server::{XRayServer, XRayServerWithDuration}; use crate::services::config::Config; use crate::utils::tap::Tap; use crate::v2parser::entities::log::{Log, LogLevel}; use crate::v2parser::parser::create_json_config; use clap::Args; +use futures::{StreamExt, stream}; use reqwest::Client; static DEFAULT_HTTP_PORT: u16 = 8080; @@ -35,6 +40,40 @@ pub(crate) struct SelectParams { } impl SelectParams { + async fn sort_servers<'a>(&self, servers: &'a [XRayServer]) -> Vec { + let mut sorted_map = HashMap::new(); + let measure_results = stream::iter(servers.iter().enumerate()) + .map(async move |(idx, s)| (idx, s.measure_rtt(Duration::from_secs(5)).await)) + .buffer_unordered(5) + .collect::>() + .await; + for (idx, duration) in measure_results.into_iter() { + let Some(server) = servers.get(idx) else { + continue; + }; + sorted_map.insert(server.to_owned(), duration); + } + + let mut servers_sorted = sorted_map + .into_iter() + .collect::)>>(); + servers_sorted.sort_by( + |(_, duration_a), (_, duration_b)| match (duration_a, duration_b) { + (None, None) => Ordering::Equal, + (None, Some(_)) => Ordering::Greater, + (Some(_), None) => Ordering::Less, + (Some(duration_a), Some(duration_b)) => { + duration_a.as_millis().cmp(&duration_b.as_millis()) + } + }, + ); + + servers_sorted + .into_iter() + .map(|(server, duration)| XRayServerWithDuration::new(server, duration)) + .collect() + } + async fn select_local(&self) -> anyhow::Result { let config = Config::read_or_default().await; let items: Vec = config @@ -104,13 +143,15 @@ impl SelectParams { .filter_map(|u| XRayServer::try_from(u.to_string()).ok()) .collect::>(); + let sorted = self.sort_servers(xray_servers.as_slice()).await; + if xray_servers.is_empty() { return Err(anyhow::Error::msg("This URL has note XRay servers")); } let server_idx = dialoguer::FuzzySelect::new() .with_prompt("Select one XRay server") - .items(&xray_servers) + .items(&sorted) .tap(|d| match &xray_servers.is_empty() { true => d, false => d.default(0), diff --git a/src/entities/xray_server.rs b/src/entities/xray_server.rs index 2c8bc73..74b27df 100644 --- a/src/entities/xray_server.rs +++ b/src/entities/xray_server.rs @@ -1,7 +1,13 @@ +use crate::services::dns_client; use percent_encoding::percent_decode_str; use serde::{Deserialize, Serialize}; -use std::fmt::Display; -use url::Url; +use std::{ + fmt::Display, + net::{IpAddr, SocketAddr}, + sync::Arc, + time::{Duration, Instant}, +}; +use url::{Host, Url}; const UNNAMED: &'static str = "Unnamed XRay Server"; @@ -13,6 +19,64 @@ impl XRayServer { pub fn url(&self) -> &Url { &self.0 } + + async fn measure_rtt_for_addr( + &self, + timeout: Duration, + addr: &SocketAddr, + ) -> anyhow::Result { + let start = Instant::now(); + match tokio::time::timeout(timeout, tokio::net::TcpStream::connect(addr)).await { + Ok(Ok(stream)) => { + let result = start.elapsed(); + drop(stream); + Ok(result) + } + Ok(Err(err)) => Err(err.into()), + _ => anyhow::bail!("Timeout"), + } + } + + async fn measure_rtt_for_addrs( + &self, + timeout: Duration, + addrs: impl Iterator, + ) -> anyhow::Result { + let mut durations = Vec::new(); + for addr in addrs { + let duration = self.measure_rtt_for_addr(timeout, &addr).await; + durations.push(duration); + } + if durations.is_empty() { + anyhow::bail!("Empty addrs list"); // unused, почему? + } + let mut ok_durations = durations + .into_iter() + .filter_map(|r| r.ok()) + .collect::>(); + ok_durations.sort_by_key(|v| v.as_millis()); + + ok_durations + .get((ok_durations.len() / 2) as usize) + .cloned() + .ok_or_else(|| anyhow::anyhow!("Cannot estimate")) + } + + pub async fn measure_rtt(&self, timeout: Duration) -> Option { + let Some((host, port)) = self.url().host().zip(self.0.port()) else { + return None; + }; + let ip_addresses = match host { + Host::Domain(dns_name) => dns_client::resolve(dns_name).await, + Host::Ipv4(ipv4) => Arc::new(vec![IpAddr::V4(ipv4)]), + Host::Ipv6(ipv6) => Arc::new(vec![IpAddr::V6(ipv6)]), + }; + let socket_addrs = ip_addresses + .iter() + .map(|v| SocketAddr::new(v.to_owned(), port)); + + self.measure_rtt_for_addrs(timeout, socket_addrs).await.ok() + } } impl Display for XRayServer { @@ -37,3 +101,31 @@ impl TryFrom for XRayServer { Ok(url) } } + +pub(crate) struct XRayServerWithDuration { + xray_server: XRayServer, + duration: Option, +} + +impl XRayServerWithDuration { + pub fn new(xray_server: XRayServer, duration: Option) -> Self { + Self { + xray_server, + duration, + } + } +} + +impl Display for XRayServerWithDuration { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let server_name = self.xray_server.to_string(); + let text = self + .duration + .map(|d| { + let ms = d.as_millis(); + format!("{server_name}, {ms} ms") + }) + .unwrap_or_else(|| format!("{server_name}, not accessible")); + f.write_str(text.as_str()) + } +} diff --git a/src/services/dns_client.rs b/src/services/dns_client.rs new file mode 100644 index 0000000..3731d97 --- /dev/null +++ b/src/services/dns_client.rs @@ -0,0 +1,91 @@ +use dnsclient::UpstreamServer; +use dnsclient::r#async::DNSClient as AsyncDNSClient; +use std::collections::HashMap; +use std::net::{IpAddr, SocketAddr}; +use std::sync::{Arc, LazyLock}; +use tokio::sync::{Mutex, OnceCell}; + +static DNS_CLIENTS_RAW: &'static str = include_str!("./dns_servers.txt"); +static SEP: &'static str = "\n"; + +struct DNSClient { + client: AsyncDNSClient, + cache: Arc>>>>>>, +} + +impl DNSClient { + fn try_new() -> anyhow::Result { + let dns_servers: Vec = DNS_CLIENTS_RAW + .split(SEP) + .filter_map(|v| -> Option { v.parse().ok() }) + .map(|v| UpstreamServer::new(v)) + .collect(); + + if dns_servers.is_empty() { + return Err(anyhow::Error::msg("No DNS servers found")); + } + + let client = AsyncDNSClient::new(dns_servers); + + Ok(Self { + client, + cache: Arc::default(), + }) + } + + async fn get_all_addrs(&self, dns_name: impl AsRef) -> Vec { + let ipv4 = self + .client + .query_a(dns_name.as_ref()) + .await + .unwrap_or_else(|_| vec![]) + .into_iter() + .map(|v| IpAddr::V4(v)) + .collect::>(); + let ipv6 = self + .client + .query_aaaa(dns_name.as_ref()) + .await + .unwrap_or_else(|_| vec![]) + .into_iter() + .map(|v| IpAddr::V6(v)) + .collect::>(); + let mut result = Vec::with_capacity(ipv4.len() + ipv6.len()); + + result.extend(ipv4); + result.extend(ipv6); + + result + } + + async fn resolve_or_insert( + &self, + dns_name: impl AsRef, + ) -> Arc>>> { + let mut cache = self.cache.lock().await; + let result = cache + .entry(dns_name.as_ref().to_owned()) + .or_insert_with(|| Arc::new(OnceCell::new())); + + result.clone() + } + + pub async fn resolve(&self, dns_name: impl AsRef) -> Arc> { + let cell = self.resolve_or_insert(dns_name.as_ref()).await; + cell.get_or_init(async move || { + let addrs = self.get_all_addrs(dns_name.as_ref()).await; + Arc::new(addrs) + }) + .await + .clone() + } +} + +static DEFAULT_DNS_CLIENT: LazyLock> = LazyLock::new(DNSClient::try_new); + +pub(crate) async fn resolve(dns_name: impl AsRef) -> Arc> { + match DEFAULT_DNS_CLIENT.as_ref() { + Err(_) => Arc::default(), + Ok(client) => client.resolve(dns_name).await, + } +} diff --git a/src/services/dns_servers.txt b/src/services/dns_servers.txt index 39cbb37..3ae0cba 100644 --- a/src/services/dns_servers.txt +++ b/src/services/dns_servers.txt @@ -1,4 +1,4 @@ 1.1.1.1:53 1.0.0.1:53 8.8.8.8:53 -8.8.4.4:53 \ No newline at end of file +8.8.4.4:53 diff --git a/src/services/mod.rs b/src/services/mod.rs index ec94232..15fe709 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,2 +1,3 @@ pub(crate) mod config; +pub(crate) mod dns_client; pub(crate) mod fileman;