Skip to content
130 changes: 124 additions & 6 deletions aether/src/cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,13 @@ MASQUE transport:
--fragment-size <n|a-b> fragment chunk size in bytes (default 16-32)
--fragment-delay <n|a-b> delay between fragments in ms (default 2-10)

Health monitoring:
--health-interval <n> background tunnel health check interval in seconds (default 20, 0 to disable)
--health-fails <n> consecutive failed health checks before reconnect (default 2)
--health-timeout <n> timeout for each health check probe in seconds (default 5)
--health-url <url> custom probe URL for health check (default http://www.gstatic.com/generate_204)


WireGuard:
--keepalive <n> persistent keepalive interval in seconds (default 5)
--no-profile-retry don't retry other obfuscation profiles during scan
Expand All @@ -51,16 +58,23 @@ Config files:

Advanced:
--tls-groups <list> TLS key share groups, e.g. \"P-256:X25519:P-384\"
--verbose detailed debug logs: tunnel stages, validation, reconnects, retries
(equivalent to RUST_LOG=info,aether=debug; RUST_LOG overrides this)

-v, --version show version and exit
Logging:
-v, --verbose detailed debug logs (use -vv for trace)
-l, --log-level <level> set log level (error, warn, info, debug, trace)

-V, --version show version and exit
-h, --help show this help and exit
";

pub fn parse_and_apply() -> crate::error::Result<()> {
let args: Vec<String> = env::args().skip(1).collect();
parse_args(&args)
}

pub fn parse_args(args: &[String]) -> crate::error::Result<()> {
let mut i = 0;
let mut verbose_count = 0;

while i < args.len() {
let arg = args[i].as_str();
Expand All @@ -75,7 +89,7 @@ pub fn parse_and_apply() -> crate::error::Result<()> {
}

match arg {
"-v" | "--version" => {
"-V" | "--version" => {
println!("aether {}", env!("CARGO_PKG_VERSION"));
std::process::exit(0);
}
Expand All @@ -85,6 +99,19 @@ pub fn parse_and_apply() -> crate::error::Result<()> {
std::process::exit(0);
}

"-v" | "--verbose" => {
verbose_count += 1;
}
"-vv" => {
verbose_count += 2;
}
"-vvv" => {
verbose_count += 3;
}
"-l" | "--log-level" => {
set("AETHER_LOG", next_value!());
}

"--bind" => set("AETHER_SOCKS", next_value!()),
"--quick-reconnect" => set("AETHER_QUICK_RECONNECT", "1"),
"--no-quick-reconnect" => set("AETHER_QUICK_RECONNECT", "0"),
Expand Down Expand Up @@ -119,11 +146,21 @@ pub fn parse_and_apply() -> crate::error::Result<()> {
set("AETHER_WG_NO_DATA_CHECK", "1");
}
"--validate-secs" => set("AETHER_MASQUE_VALIDATE_SECS", next_value!()),
"--reconnect-secs" => set("AETHER_MASQUE_RECONNECT_SECS", next_value!()),
"--reconnect-secs" => {
let val = next_value!();
set("AETHER_MASQUE_RECONNECT_SECS", val);
set("AETHER_WG_RECONNECT_SECS", val);
}
"--fragment" => set("AETHER_MASQUE_H2_FRAGMENT", "1"),
"--fragment-size" => set("AETHER_MASQUE_H2_FRAGMENT_SIZE", next_value!()),
"--fragment-delay" => set("AETHER_MASQUE_H2_FRAGMENT_DELAY", next_value!()),

"--health-interval" => set("AETHER_HEALTH_INTERVAL", next_value!()),
"--health-fails" => set("AETHER_HEALTH_MAX_FAILS", next_value!()),
"--health-timeout" => set("AETHER_HEALTH_TIMEOUT", next_value!()),
"--health-url" => set("AETHER_HEALTH_PROBE_URL", next_value!()),


"--keepalive" => set("AETHER_WG_KEEPALIVE", next_value!()),
"--no-profile-retry" => set("AETHER_WG_NO_PROFILE_RETRY", "1"),

Expand All @@ -132,7 +169,6 @@ pub fn parse_and_apply() -> crate::error::Result<()> {
"--masque-config" => set("AETHER_MASQUE_CONFIG", next_value!()),

"--tls-groups" => set("AETHER_TLS_GROUPS", next_value!()),
"--verbose" => set("AETHER_VERBOSE", "1"),

other => {
return Err(crate::error::AetherError::Other(format!(
Expand All @@ -144,9 +180,91 @@ pub fn parse_and_apply() -> crate::error::Result<()> {
i += 1;
}

if verbose_count > 0 && std::env::var("AETHER_LOG").is_err() {
if verbose_count == 1 {
set("AETHER_LOG", "info,aether=debug");
} else {
set("AETHER_LOG", "info,aether=trace");
}
}

Ok(())
}

fn set(key: &str, value: &str) {
std::env::set_var(key, value);
}

#[cfg(test)]
pub static ENV_MUTEX: std::sync::Mutex<()> = std::sync::Mutex::new(());

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn test_cli_parse_log_levels() {
let _guard = ENV_MUTEX.lock().unwrap();
std::env::remove_var("AETHER_LOG");
parse_args(&["--verbose".to_string()]).unwrap();
assert_eq!(std::env::var("AETHER_LOG").unwrap(), "info,aether=debug");

std::env::remove_var("AETHER_LOG");
parse_args(&["-v".to_string()]).unwrap();
assert_eq!(std::env::var("AETHER_LOG").unwrap(), "info,aether=debug");

std::env::remove_var("AETHER_LOG");
parse_args(&["-vv".to_string()]).unwrap();
assert_eq!(std::env::var("AETHER_LOG").unwrap(), "info,aether=trace");

std::env::remove_var("AETHER_LOG");
parse_args(&["-l".to_string(), "warn".to_string()]).unwrap();
assert_eq!(std::env::var("AETHER_LOG").unwrap(), "warn");

std::env::remove_var("AETHER_LOG");
parse_args(&["--log-level".to_string(), "trace".to_string()]).unwrap();
assert_eq!(std::env::var("AETHER_LOG").unwrap(), "trace");
std::env::remove_var("AETHER_LOG");
}

#[test]
fn test_cli_parse_health_options() {
let _guard = ENV_MUTEX.lock().unwrap();
parse_args(&[
"--health-interval".to_string(),
"10".to_string(),
"--health-fails".to_string(),
"4".to_string(),
"--health-timeout".to_string(),
"3".to_string(),
"--health-url".to_string(),
"http://cp.cloudflare.com/generate_204".to_string(),
"--reconnect-secs".to_string(),
"5".to_string(),
])
.unwrap();

assert_eq!(std::env::var("AETHER_HEALTH_INTERVAL").unwrap(), "10");
assert_eq!(std::env::var("AETHER_HEALTH_MAX_FAILS").unwrap(), "4");
assert_eq!(std::env::var("AETHER_HEALTH_TIMEOUT").unwrap(), "3");
assert_eq!(
std::env::var("AETHER_HEALTH_PROBE_URL").unwrap(),
"http://cp.cloudflare.com/generate_204"
);
assert_eq!(std::env::var("AETHER_MASQUE_RECONNECT_SECS").unwrap(), "5");
assert_eq!(std::env::var("AETHER_WG_RECONNECT_SECS").unwrap(), "5");

std::env::remove_var("AETHER_HEALTH_INTERVAL");
std::env::remove_var("AETHER_HEALTH_MAX_FAILS");
std::env::remove_var("AETHER_HEALTH_TIMEOUT");
std::env::remove_var("AETHER_HEALTH_PROBE_URL");
std::env::remove_var("AETHER_MASQUE_RECONNECT_SECS");
std::env::remove_var("AETHER_WG_RECONNECT_SECS");
}
}






100 changes: 81 additions & 19 deletions aether/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,12 +41,12 @@ const DEFAULT_CONFIG: &str = "aether.toml";
async fn main() -> Result<()> {
cli::parse_and_apply()?;

let default_filter = if std::env::var("AETHER_VERBOSE").is_ok() {
"info,aether=debug"
let log_env = if std::env::var("AETHER_LOG").is_ok() {
env_logger::Env::default().filter("AETHER_LOG")
} else {
"info"
env_logger::Env::default().default_filter_or("info")
};
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or(default_filter))
env_logger::Builder::from_env(log_env)
.format_timestamp_millis()
.init();

Expand Down Expand Up @@ -557,22 +557,39 @@ async fn run_masque_tunnel(
}
}

let health_cfg = tunnelping::HealthConfig::from_env();
let health_task = tunnelping::spawn_health_monitor(stack.clone(), health_cfg);

let socks_stack = stack.clone();
let socks_task = tokio::spawn(async move {
log::info!("[+] socks5 server listening on {listen}");
socks::serve(listen, socks_stack).await
});

let tunnel_result = tunnel_task.await;
socks_task.abort();
let tunnel_result = tokio::select! {
res = tunnel_task => match res {
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(AetherError::Other(format!("tunnel exited: {e}"))),
Err(e) => Err(AetherError::Other(format!("tunnel task join error: {e}"))),
Comment thread
Vonarian marked this conversation as resolved.
},
health_res = async {
if let Some(h) = health_task {
h.await
} else {
std::future::pending().await
}
} => match health_res {
Ok(Err(e)) => Err(AetherError::Other(format!("health check failure: {e}"))),
Ok(Ok(())) => Err(AetherError::Other("health monitor exited unexpectedly".into())),
Err(e) => Err(AetherError::Other(format!("health monitor join error: {e}"))),
}
};

match tunnel_result {
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(AetherError::Other(format!("tunnel exited: {e}"))),
Err(e) => Err(AetherError::Other(format!("tunnel task join error: {e}"))),
}
socks_task.abort();
tunnel_result
}


fn wg_keepalive_secs() -> u16 {
std::env::var("AETHER_WG_KEEPALIVE")
.ok()
Expand Down Expand Up @@ -851,21 +868,41 @@ async fn run_wireguard_tunnel(

let stack = netstack::spawn(&identity.ipv4, &identity.ipv6, TUNNEL_MTU, inbound_rx, outbound_tx)?;

let health_cfg = tunnelping::HealthConfig::from_env();
let health_task = tunnelping::spawn_health_monitor(stack.clone(), health_cfg);

let socks_stack = stack.clone();
let socks_task = tokio::spawn(async move {
log::info!("[+] socks5 server listening on {listen}");
socks::serve(listen, socks_stack).await
});

let tunnel_result = tunnel.run(outbound_rx).await;
socks_task.abort();
let tunnel_task = tokio::spawn(tunnel.run(outbound_rx));

match tunnel_result {
Ok(()) => Ok(()),
Err(e) => Err(AetherError::Other(format!("wireguard tunnel exited: {e}"))),
}
let res = tokio::select! {
tunnel_res = tunnel_task => match tunnel_res {
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(AetherError::Other(format!("wireguard tunnel exited: {e}"))),
Err(e) => Err(AetherError::Other(format!("wireguard tunnel task join error: {e}"))),
},
health_res = async {
if let Some(h) = health_task {
h.await
} else {
std::future::pending().await
}
} => match health_res {
Ok(Err(e)) => Err(AetherError::Other(format!("health check failure: {e}"))),
Ok(Ok(())) => Err(AetherError::Other("health monitor exited unexpectedly".into())),
Err(e) => Err(AetherError::Other(format!("health monitor join error: {e}"))),
}
};

socks_task.abort();
res
}


async fn establish_wg(
identity: &account::Identity,
peer: SocketAddr,
Expand Down Expand Up @@ -979,11 +1016,36 @@ async fn run_warp_in_warp(

log::info!("[*] establishing inner WARP tunnel (warp-in-warp)...");
let inner_stack = establish_wg(&secondary, forwarder, INNER_MTU, false, 20, "inner").await?;
let health_cfg = tunnelping::HealthConfig::from_env();
let health_task = tunnelping::spawn_health_monitor(inner_stack.clone(), health_cfg);

let socks_task = tokio::spawn(async move {
log::info!("[+] socks5 server listening on {listen}");
socks::serve(listen, inner_stack).await
});

log::info!("[+] socks5 server listening on {listen}");
socks::serve(listen, inner_stack).await
let res = tokio::select! {
socks_res = socks_task => match socks_res {
Ok(res) => res,
Err(e) => Err(AetherError::Other(format!("socks server join error: {e}"))),
},
health_res = async {
if let Some(h) = health_task {
h.await
} else {
std::future::pending().await
}
} => match health_res {
Ok(Err(e)) => Err(AetherError::Other(format!("health check failure: {e}"))),
Ok(Ok(())) => Err(AetherError::Other("health monitor exited unexpectedly".into())),
Err(e) => Err(AetherError::Other(format!("health monitor join error: {e}"))),
}
};

res
}


async fn prompt_line(prompt: &str) -> Option<String> {
use std::io::IsTerminal;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
Expand Down
4 changes: 2 additions & 2 deletions aether/src/masque_h2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -434,7 +434,7 @@ pub async fn run(
Some(ip_packet) => {
let framed = masque::encode_datagram_capsule(&ip_packet);
if let Err(e) = send_capsule(&mut send_stream, Bytes::from(framed)).await {
log::debug!("[h2] send: {e}");
log::trace!("[h2] send: {e}");
return Err(e);
}
}
Expand Down Expand Up @@ -545,7 +545,7 @@ async fn drain_capsules(
Ok(Some(_)) => {}
Ok(None) => break,
Err(e) => {
log::debug!("[h2] capsule parse: {e}");
log::trace!("[h2] capsule parse: {e}");
break;
}
}
Expand Down
Loading