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..2218f13 100644 --- a/src/cmd/select.rs +++ b/src/cmd/select.rs @@ -1,17 +1,16 @@ -use std::process::Stdio; -use std::{cmp::Ordering, collections::HashMap}; +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_with_perf::XRayServerWithPerf; +use crate::entities::xray_server::{XRayServer, XRayServerWithDuration}; 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 futures::{StreamExt, stream}; use reqwest::Client; -use tokio::io::AsyncWriteExt; static DEFAULT_HTTP_PORT: u16 = 8080; static DEFAULT_SOCKS_PORT: u16 = 1080; @@ -27,14 +26,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,48 +40,37 @@ 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); - } + 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 = 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) { + 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(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 + (Some(duration_a), Some(duration_b)) => { + duration_a.as_millis().cmp(&duration_b.as_millis()) + } + }, + ); + + servers_sorted .into_iter() - .map(|ref s| { - XRayServerWithPerf::new(s.clone(), measure_results_index.get(s.url()).cloned()) - }) + .map(|(server, duration)| XRayServerWithDuration::new(server, duration)) .collect() } @@ -163,7 +143,7 @@ impl SelectParams { .filter_map(|u| XRayServer::try_from(u.to_string()).ok()) .collect::>(); - let xray_servers = self.sort_servers(xray_servers).await; + 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")); @@ -171,7 +151,7 @@ impl SelectParams { 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), @@ -182,7 +162,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 +197,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..74b27df 100644 --- a/src/entities/xray_server.rs +++ b/src/entities/xray_server.rs @@ -1,11 +1,12 @@ 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 std::{ + fmt::Display, + net::{IpAddr, SocketAddr}, + sync::Arc, + time::{Duration, Instant}, +}; use url::{Host, Url}; const UNNAMED: &'static str = "Unnamed XRay Server"; @@ -47,7 +48,7 @@ impl XRayServer { durations.push(duration); } if durations.is_empty() { - anyhow::bail!("Empty addrs list"); + anyhow::bail!("Empty addrs list"); // unused, почему? } let mut ok_durations = durations .into_iter() @@ -100,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/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_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/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..15fe709 100644 --- a/src/services/mod.rs +++ b/src/services/mod.rs @@ -1,5 +1,3 @@ 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) - } -}