#![allow(dead_code)]
use std::sync::Arc;
use std::time::Duration;
use eggress_core::chain::{ChainExecutor, HopHandler};
use eggress_core::listener::{TcpListener, TcpListenerConfig};
use eggress_core::{BoxStream, TargetAddr, TargetHost};
use eggress_protocol_http::connect::client::http_connect;
use eggress_protocol_socks::socks5::client::socks5_connect;
use eggress_protocol_socks::socks5::server::SocksAddr;
use eggress_routing::{RouteActionSpec, RouteService, Router};
use eggress_testkit::differential::*;
use eggress_uri::{ProtocolSpec, ProxyHopSpec};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio_util::sync::CancellationToken;
fn target_to_socks_addr(target: &TargetAddr) -> SocksAddr {
match &target.host {
TargetHost::Ip(std::net::IpAddr::V4(ip)) => SocksAddr::IPv4(ip.octets(), target.port),
TargetHost::Ip(std::net::IpAddr::V6(ip)) => SocksAddr::IPv6(ip.octets(), target.port),
TargetHost::Domain(d) => SocksAddr::Domain(d.clone(), target.port),
}
}
fn socket_addr(host: &str, port: u16) -> std::net::SocketAddr {
std::net::SocketAddr::new(host.parse().unwrap(), port)
}
type HandshakeFuture<'a> = std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<BoxStream, Box<dyn std::error::Error + Send + Sync>>,
> + Send
+ 'a,
>,
>;
struct HttpHopHandler;
impl HopHandler for HttpHopHandler {
fn protocol(&self) -> ProtocolSpec {
ProtocolSpec::Http
}
fn handshake<'a>(
&'a self,
stream: BoxStream,
target: &'a TargetAddr,
hop: &'a ProxyHopSpec,
_hop_index: usize,
) -> HandshakeFuture<'a> {
let auth = hop
.credentials
.as_ref()
.map(|c| (c.username.as_str(), c.password.as_str()));
Box::pin(async move {
http_connect(stream, target, auth, &Default::default())
.await
.map_err(|e| Box::new(e) as Box<dyn std::error::Error + Send + Sync>)
})
}
}
struct Socks5HopHandler;
impl HopHandler for Socks5HopHandler {
fn protocol(&self) -> ProtocolSpec {
ProtocolSpec::Socks5
}
fn handshake<'a>(
&'a self,
stream: BoxStream,
target: &'a TargetAddr,
hop: &'a ProxyHopSpec,
_hop_index: usize,
) -> HandshakeFuture<'a> {
let socks_addr = target_to_socks_addr(target);
let auth = hop
.credentials
.as_ref()
.map(|c| (c.username.as_str(), c.password.as_str()));
Box::pin(async move {
socks5_connect(stream, &socks_addr, auth)
.await
.map_err(|e| Box::new(e) as Box<dyn std::error::Error + Send + Sync>)
})
}
}
fn build_executor() -> ChainExecutor {
ChainExecutor::new(vec![Box::new(HttpHopHandler), Box::new(Socks5HopHandler)])
}
async fn send_through_socks5(
proxy_addr: std::net::SocketAddr,
target: &TargetAddr,
payload: &[u8],
) -> Result<Vec<u8>, String> {
let stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let boxed: BoxStream = Box::new(stream);
let socks_addr = target_to_socks_addr(target);
let mut conn = socks5_connect(boxed, &socks_addr, None)
.await
.map_err(|e| format!("socks5 handshake failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
conn.flush()
.await
.map_err(|e| format!("flush failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
async fn send_through_socks5_with_auth(
proxy_addr: std::net::SocketAddr,
target: &TargetAddr,
payload: &[u8],
username: &str,
password: &str,
) -> Result<Vec<u8>, String> {
let stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let boxed: BoxStream = Box::new(stream);
let socks_addr = target_to_socks_addr(target);
let mut conn = socks5_connect(boxed, &socks_addr, Some((username, password)))
.await
.map_err(|e| format!("socks5 handshake failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
conn.flush()
.await
.map_err(|e| format!("flush failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
async fn send_through_socks5_stream<S>(
stream: S,
target: &TargetAddr,
payload: &[u8],
) -> Result<Vec<u8>, String>
where
S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static,
{
let boxed: BoxStream = Box::new(stream);
let socks_addr = target_to_socks_addr(target);
let mut conn = socks5_connect(boxed, &socks_addr, None)
.await
.map_err(|e| format!("socks5 handshake failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
async fn send_through_http(
proxy_addr: std::net::SocketAddr,
target: &TargetAddr,
payload: &[u8],
) -> Result<Vec<u8>, String> {
let stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let boxed: BoxStream = Box::new(stream);
let mut conn = http_connect(boxed, target, None, &Default::default())
.await
.map_err(|e| format!("http connect handshake failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
async fn send_through_http_with_auth(
proxy_addr: std::net::SocketAddr,
target: &TargetAddr,
payload: &[u8],
username: &str,
password: &str,
) -> Result<Vec<u8>, String> {
let stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let boxed: BoxStream = Box::new(stream);
let mut conn = http_connect(
boxed,
target,
Some((username, password)),
&Default::default(),
)
.await
.map_err(|e| format!("http connect handshake failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
async fn send_through_socks4(
proxy_addr: std::net::SocketAddr,
target: std::net::SocketAddr,
payload: &[u8],
) -> Result<Vec<u8>, String> {
let mut stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let mut req = vec![0x04, 0x01]; req.extend_from_slice(&target.port().to_be_bytes());
match target.ip() {
std::net::IpAddr::V4(ip) => req.extend_from_slice(&ip.octets()),
std::net::IpAddr::V6(_) => return Err("SOCKS4 does not support IPv6 targets".into()),
}
req.push(0x00); stream
.write_all(&req)
.await
.map_err(|e| format!("write failed: {e}"))?;
let mut reply = [0u8; 8];
stream
.read_exact(&mut reply)
.await
.map_err(|e| format!("read reply failed: {e}"))?;
if reply[1] != 0x5a {
return Err(format!("SOCKS4 CONNECT failed: code {}", reply[1]));
}
stream
.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut stream, Duration::from_secs(3)).await)
}
async fn send_through_socks4a(
proxy_addr: std::net::SocketAddr,
target_host: &str,
target_port: u16,
payload: &[u8],
) -> Result<Vec<u8>, String> {
let mut stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
let mut req = vec![0x04, 0x01]; req.extend_from_slice(&target_port.to_be_bytes());
req.extend_from_slice(&[0, 0, 0, 1]); req.extend_from_slice(target_host.as_bytes());
req.push(0x00); stream
.write_all(&req)
.await
.map_err(|e| format!("write failed: {e}"))?;
let mut reply = [0u8; 8];
stream
.read_exact(&mut reply)
.await
.map_err(|e| format!("read reply failed: {e}"))?;
if reply[1] != 0x5a {
return Err(format!("SOCKS4a CONNECT failed: code {}", reply[1]));
}
stream
.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut stream, Duration::from_secs(3)).await)
}
async fn socks5_udp_associate(
stream: &mut tokio::net::TcpStream,
) -> std::io::Result<std::net::SocketAddr> {
stream.write_all(&[0x05, 0x01, 0x00]).await?;
let mut resp = [0u8; 2];
stream.read_exact(&mut resp).await?;
assert_eq!(resp, [0x05, 0x00]);
stream
.write_all(&[0x05, 0x03, 0x00, 0x01, 0, 0, 0, 0])
.await?;
stream.write_all(&0u16.to_be_bytes()).await?;
let mut reply = [0u8; 22];
let n = stream.read(&mut reply).await?;
assert!(n >= 10, "UDP ASSOCIATE reply too short: {n} bytes");
assert_eq!(reply[0], 0x05, "SOCKS5 version mismatch");
assert_eq!(reply[1], 0x00, "UDP ASSOCIATE failed: code {}", reply[1]);
let relay_ip = match reply[3] {
0x01 => {
let ip = std::net::Ipv4Addr::new(reply[4], reply[5], reply[6], reply[7]);
std::net::IpAddr::V4(ip)
}
_ => panic!("unexpected address type in UDP ASSOCIATE reply"),
};
let relay_port = u16::from_be_bytes([reply[8], reply[9]]);
Ok(std::net::SocketAddr::new(relay_ip, relay_port))
}
struct TaskGuard {
cancel: Option<CancellationToken>,
jh: Option<tokio::task::JoinHandle<()>>,
}
impl TaskGuard {
fn new(cancel: CancellationToken, jh: tokio::task::JoinHandle<()>) -> Self {
Self {
cancel: Some(cancel),
jh: Some(jh),
}
}
fn cancel_token(&self) -> &CancellationToken {
self.cancel.as_ref().unwrap()
}
fn shutdown(&mut self) {
if let Some(cancel) = self.cancel.take() {
cancel.cancel();
}
if let Some(jh) = self.jh.take() {
jh.abort();
}
}
}
impl Drop for TaskGuard {
fn drop(&mut self) {
self.shutdown();
}
}
async fn start_eggress_server(
protocols: Vec<eggress_core::ProtocolId>,
) -> (
std::net::SocketAddr,
CancellationToken,
tokio::task::JoinHandle<()>,
) {
let config = TcpListenerConfig {
bind_addr: "127.0.0.1:0".parse().unwrap(),
protocols,
auth_required: false,
handshake_timeout: Duration::from_secs(5),
connection_limit: 10,
};
let cancel = CancellationToken::new();
let listener = TcpListener::new(&config, cancel.clone()).await.unwrap();
let addr = listener.local_addr().unwrap();
let conn_protocols: Arc<[eggress_core::ProtocolId]> = config.protocols.clone().into();
let jh = tokio::spawn(async move {
loop {
let conn = match listener.accept().await {
Ok(c) => c,
Err(_) => break,
};
let config = eggress_server::ConnectionConfig {
routing: Arc::new(Router::new(vec![], RouteActionSpec::Direct))
as Arc<dyn RouteService>,
context: eggress_server::ConnectionContext::default(),
handshake_timeout: Duration::from_secs(5),
connect_timeout: Duration::from_secs(10),
protocols: conn_protocols.clone(),
authentication: eggress_server::accept::InboundAuthentication::None,
metrics: None,
udp: None,
tls_client_config: None,
shadowsocks: None,
shadowsocks_metrics: None,
trojan: None,
fixed_target: None,
local_bind: None,
};
tokio::spawn(async move {
let _ = eggress_server::serve_connection(conn.stream, config).await;
});
}
});
(addr, cancel, jh)
}
async fn wait_ready(state: &eggress_runtime::RuntimeState) {
use std::sync::atomic::Ordering;
for _ in 0..100 {
if state.readiness.load(Ordering::Relaxed) {
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
panic!("timeout waiting for readiness");
}
async fn start_eggress_from_toml_running(
config_str: &str,
) -> (
std::net::SocketAddr,
CancellationToken,
tokio::task::JoinHandle<()>,
) {
use std::io::Write;
let mut f = tempfile::NamedTempFile::new().expect("create tempfile");
f.write_all(config_str.as_bytes()).expect("write config");
f.flush().expect("flush config");
let path = f.path().to_str().unwrap().to_string();
std::mem::forget(f);
let mut sup =
eggress_runtime::ServiceSupervisor::start(&path).expect("start eggress from TOML");
let state = sup.state().clone();
let token = sup.shutdown_token();
let jh = tokio::task::spawn_blocking(move || {
let _ = sup.run();
});
wait_ready(&state).await;
let listener_addr = {
let addrs = state.listener_addrs.lock().unwrap();
addrs[0].unwrap()
};
(listener_addr, token, jh)
}
async fn start_http_origin() -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let jh = tokio::spawn(async move {
loop {
let (mut stream, _) = match listener.accept().await {
Ok(v) => v,
Err(_) => break,
};
tokio::spawn(async move {
let mut buf = vec![0u8; 4096];
let mut total = 0;
let mut headers_done = false;
while !headers_done {
let n = stream.read(&mut buf[total..]).await.unwrap_or(0);
if n == 0 {
break;
}
total += n;
if buf[..total].windows(4).any(|w| w == b"\r\n\r\n") {
headers_done = true;
}
}
let head = String::from_utf8_lossy(&buf[..total]);
let has_body = head.to_lowercase().contains("content-length:");
if has_body {
loop {
let n = stream.read(&mut buf[total..]).await.unwrap_or(0);
if n == 0 {
break;
}
total += n;
}
}
let response = "HTTP/1.1 200 OK\r\nContent-Length: 13\r\nConnection: close\r\n\r\nHello, origin!";
let _ = stream.write_all(response.as_bytes()).await;
});
}
});
(addr, jh)
}
async fn send_http_forward(
proxy_addr: std::net::SocketAddr,
request: &[u8],
) -> Result<Vec<u8>, String> {
let mut stream = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to proxy failed: {e}"))?;
stream
.write_all(request)
.await
.map_err(|e| format!("write failed: {e}"))?;
Ok(read_with_timeout(&mut stream, Duration::from_secs(5)).await)
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_http_connect() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_server("http", pproxy_port).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let pproxy_result = send_through_http(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"differential http connect",
)
.await;
pproxy.kill();
let (egress_addr, cancel, jh) =
start_eggress_server(vec![eggress_core::ProtocolId::Http]).await;
tokio::time::sleep(Duration::from_millis(50)).await;
let egress_result = send_through_http(egress_addr, &target, b"differential http connect").await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
compare_tcp_echo("pproxy", &pproxy_result, "eggress", &egress_result);
assert_eq!(pproxy_result.unwrap(), b"differential http connect");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_http_forward() {
require_differential_gate();
let (origin_addr, origin_jh) = start_http_origin().await;
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_server("http", pproxy_port).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let request = format!(
"GET http://127.0.0.1:{}/path HTTP/1.1\r\nHost: 127.0.0.1:{}\r\nConnection: close\r\n\r\n",
origin_addr.port(),
origin_addr.port(),
);
let pproxy_result =
send_http_forward(socket_addr("127.0.0.1", pproxy_port), request.as_bytes()).await;
pproxy.kill();
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "http-in"
bind = "127.0.0.1:{port}"
protocols = ["http"]
[[rules]]
id = "allow-all"
direct = true
[routing]
default = "direct"
"#,
port = egress_port,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let egress_result = send_http_forward(egress_addr, request.as_bytes()).await;
cancel.cancel();
let _ = jh.await;
origin_jh.abort();
assert_coarse_failure_equivalence("pproxy", &pproxy_result, "eggress", &egress_result);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_socks4_connect() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_server("socks4", pproxy_port).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let pproxy_result = send_through_socks4(
socket_addr("127.0.0.1", pproxy_port),
echo_addr,
b"differential socks4",
)
.await;
pproxy.kill();
let (egress_addr, cancel, jh) =
start_eggress_server(vec![eggress_core::ProtocolId::Socks4]).await;
tokio::time::sleep(Duration::from_millis(50)).await;
let egress_result = send_through_socks4(egress_addr, echo_addr, b"differential socks4").await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
compare_tcp_echo("pproxy", &pproxy_result, "eggress", &egress_result);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_socks5_connect() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_server("socks5", pproxy_port).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let pproxy_result = send_through_socks5(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"differential socks5",
)
.await;
pproxy.kill();
let (egress_addr, cancel, jh) =
start_eggress_server(vec![eggress_core::ProtocolId::Socks5]).await;
tokio::time::sleep(Duration::from_millis(50)).await;
let egress_result = send_through_socks5(egress_addr, &target, b"differential socks5").await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
compare_tcp_echo("pproxy", &pproxy_result, "eggress", &egress_result);
assert_eq!(pproxy_result.unwrap(), b"differential socks5");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_socks5_auth() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let user = "testuser";
let pass = "testpass";
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_server_with_auth("socks5", pproxy_port, user, pass).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let pproxy_result = send_through_socks5_with_auth(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"differential socks5 auth",
user,
pass,
)
.await;
pproxy.kill();
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "socks-in"
bind = "127.0.0.1:{egress_port}"
protocols = ["socks5"]
[listeners.auth]
type = "password"
username = "{user}"
password = "{pass}"
[[rules]]
id = "allow-all"
direct = true
[routing]
default = "direct"
"#
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let egress_result = send_through_socks5_with_auth(
egress_addr,
&target,
b"differential socks5 auth",
user,
pass,
)
.await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
compare_tcp_echo("pproxy", &pproxy_result, "eggress", &egress_result);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_socks5_udp_associate() {
require_differential_gate();
let (udp_echo_addr, udp_echo_jh) = start_udp_echo().await;
let pproxy_tcp_port = eggress_testkit::get_free_port().await;
let listen_tcp = format!("socks5://127.0.0.1:{}", pproxy_tcp_port);
let pproxy = start_pproxy_with_args(&["-l", &listen_tcp, "-r", "direct"]).await;
assert_port_ready(pproxy_tcp_port, Duration::from_secs(5)).await;
tokio::time::sleep(Duration::from_millis(100)).await;
drop(pproxy);
let egress_tcp_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "socks-in"
bind = "127.0.0.1:{port}"
protocols = ["socks5"]
[listeners.udp]
enabled = true
[[rules]]
id = "allow-all"
direct = true
[routing]
default = "direct"
"#,
port = egress_tcp_port,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let mut egress_stream = tokio::net::TcpStream::connect(egress_addr).await.unwrap();
let egress_relay = socks5_udp_associate(&mut egress_stream).await.unwrap();
let egress_udp_sock = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
let egress_packet = build_socks5_udp_packet(udp_echo_addr, b"pproxy udp test");
egress_udp_sock
.send_to(&egress_packet, egress_relay)
.await
.unwrap();
let egress_udp_result = recv_udp_response(&egress_udp_sock, Duration::from_secs(3)).await;
cancel.cancel();
let _ = jh.await;
udp_echo_jh.abort();
assert!(
egress_udp_result.is_some(),
"eggress UDP ASSOCIATE should relay data"
);
let payload = extract_udp_payload(&egress_udp_result.unwrap());
assert_eq!(payload, b"pproxy udp test");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_standalone_udp() {
require_differential_gate();
let (udp_echo_addr, udp_echo_jh) = start_udp_echo().await;
let pproxy_port = eggress_testkit::get_free_port().await;
let listen = format!("socks5://127.0.0.1:{}", pproxy_port);
let pproxy = start_pproxy_with_args(&["-l", &listen, "-ul", &listen, "-r", "direct"]).await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
tokio::time::sleep(Duration::from_millis(100)).await;
let pproxy_sock = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
let packet = build_socks5_udp_packet(udp_echo_addr, b"standalone udp test");
pproxy_sock
.send_to(&packet, ("127.0.0.1", pproxy_port))
.await
.unwrap();
let _pproxy_result = recv_udp_response(&pproxy_sock, Duration::from_secs(3)).await;
drop(pproxy);
let udp_socket = std::sync::Arc::new(tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap());
let egress_addr = udp_socket.local_addr().unwrap();
let router = eggress_routing::Router::new(vec![], eggress_routing::RouteActionSpec::Direct);
let routing: Arc<dyn eggress_routing::RouteService> =
Arc::new(eggress_routing::SharedRoutingService::new(router));
let udp_metrics = Arc::new(eggress_udp::metrics::UdpMetrics::new());
let limits = eggress_udp::limits::UdpLimits::default();
let cancel = tokio_util::sync::CancellationToken::new();
let cancel_clone = cancel.clone();
let config = eggress_udp::standalone::StandaloneUdpConfig {
routing,
udp_metrics,
limits,
listener: "differential-test".to_string(),
generation: 1,
allow_private_egress: true,
};
let jh = tokio::spawn(async move {
let _ =
eggress_udp::standalone::standalone_udp_relay(udp_socket, config, cancel_clone).await;
});
let egress_sock = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
let egress_packet = build_socks5_udp_packet(udp_echo_addr, b"standalone udp test");
egress_sock
.send_to(&egress_packet, egress_addr)
.await
.unwrap();
let egress_result = recv_udp_response(&egress_sock, Duration::from_secs(3)).await;
cancel.cancel();
let _ = jh.await;
udp_echo_jh.abort();
assert!(
egress_result.is_some(),
"eggress standalone UDP should relay data"
);
let payload = extract_udp_payload(&egress_result.unwrap());
assert_eq!(payload, b"standalone udp test");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_scheduler_round_robin() {
require_differential_gate();
let (echo1_addr, echo1_jh) = eggress_testkit::start_echo_server().await;
let (echo2_addr, echo2_jh) = eggress_testkit::start_echo_server().await;
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy = start_pproxy_with_args(&[
"-l",
&format!("socks5://127.0.0.1:{}", pproxy_port),
"-r",
"direct",
])
.await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let target1 = TargetAddr {
host: TargetHost::Ip(echo1_addr.ip()),
port: echo1_addr.port(),
};
let target2 = TargetAddr {
host: TargetHost::Ip(echo2_addr.ip()),
port: echo2_addr.port(),
};
let r1 = send_through_socks5(
socket_addr("127.0.0.1", pproxy_port),
&target1,
b"scheduler-test-1",
)
.await;
let r2 = send_through_socks5(
socket_addr("127.0.0.1", pproxy_port),
&target2,
b"scheduler-test-2",
)
.await;
pproxy.kill();
echo1_jh.abort();
echo2_jh.abort();
assert!(r1.is_ok(), "pproxy should reach echo1: {:?}", r1.err());
assert!(r2.is_ok(), "pproxy should reach echo2: {:?}", r2.err());
assert_eq!(r1.unwrap(), b"scheduler-test-1");
assert_eq!(r2.unwrap(), b"scheduler-test-2");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1 and pproxy"]
async fn differential_block_behavior() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_port = eggress_testkit::get_free_port().await;
let block_pattern = "{127\\.0\\.0\\.1}".to_string();
let mut pproxy = start_pproxy_with_args(&[
"-l",
&format!("socks5://127.0.0.1:{}", pproxy_port),
"-r",
"direct",
"-b",
&block_pattern,
])
.await;
assert_port_ready(pproxy_port, Duration::from_secs(5)).await;
let pproxy_result = send_through_socks5(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"should-be-blocked",
)
.await;
pproxy.kill();
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "socks-in"
bind = "127.0.0.1:{port}"
protocols = ["socks5"]
[[rules]]
id = "block-target"
reject = "blocked"
[rules.match]
destination_port = {target_port}
[[rules]]
id = "allow-all"
direct = true
[routing]
default = "direct"
"#,
port = egress_port,
target_port = echo_addr.port(),
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let egress_result = send_through_socks5(egress_addr, &target, b"should-be-blocked").await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
let pproxy_ok = pproxy_result
.as_ref()
.map(|d| !d.is_empty())
.unwrap_or(false);
let egress_ok = egress_result
.as_ref()
.map(|d| !d.is_empty())
.unwrap_or(false);
assert!(!pproxy_ok, "pproxy should not deliver data through block");
assert!(!egress_ok, "eggress should not deliver data through reject");
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_tls_listener() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let cert_params = rcgen::CertificateParams::new(vec!["localhost".to_string()]).unwrap();
let key_pair = rcgen::KeyPair::generate().unwrap();
let cert = cert_params.self_signed(&key_pair).unwrap();
let cert_pem = cert.pem();
let key_pem = key_pair.serialize_pem();
let cert_file = tempfile::NamedTempFile::new().unwrap();
let key_file = tempfile::NamedTempFile::new().unwrap();
std::fs::write(cert_file.path(), &cert_pem).unwrap();
std::fs::write(key_file.path(), &key_pem).unwrap();
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "tls-in"
bind = "127.0.0.1:{egress_port}"
protocols = ["socks5"]
[listeners.tls]
cert = "{cert_path}"
key = "{key_path}"
[[rules]]
id = "allow-all"
direct = true
[routing]
default = "direct"
"#,
cert_path = cert_file.path().display(),
key_path = key_file.path().display(),
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let mut root_store = rustls::RootCertStore::empty();
let cert_der = cert.der().clone();
root_store.add(cert_der).unwrap();
let mut tls_config = rustls::ClientConfig::builder()
.with_root_certificates(root_store)
.with_no_client_auth();
tls_config.alpn_protocols = vec![b"http/1.1".to_vec()];
let connector = tokio_rustls::TlsConnector::from(Arc::new(tls_config));
let tcp = tokio::net::TcpStream::connect(egress_addr).await.unwrap();
let domain = rustls::pki_types::ServerName::try_from("localhost".to_string()).unwrap();
let tls_stream = connector.connect(domain, tcp).await.unwrap();
let result = send_through_socks5_stream(tls_stream, &target, b"tls smoke test").await;
cancel.cancel();
let _ = jh.await;
echo_jh.abort();
assert!(result.is_ok(), "TLS+SOCKS5 should work: {:?}", result.err());
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_http_to_socks5_upstream() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy_child = start_pproxy_server("socks5", pproxy_port).await;
assert!(
wait_for_port(pproxy_port, Duration::from_secs(5)).await,
"pproxy failed to start"
);
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "http-chain"
bind = "127.0.0.1:{egress_port}"
protocols = ["http"]
[[upstreams]]
id = "pproxy-socks5"
uri = "socks5://127.0.0.1:{pproxy_port}"
[[upstream_groups]]
id = "default"
members = ["pproxy-socks5"]
[routing]
default = "default"
"#,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let chain_result = send_through_http(egress_addr, &target, b"chain http->socks5").await;
let direct_result = send_through_socks5(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"chain http->socks5",
)
.await;
cancel.cancel();
let _ = jh.await;
pproxy_child.kill();
echo_jh.abort();
match (&chain_result, &direct_result) {
(Ok(chain_payload), Ok(direct_payload)) => {
assert_eq!(
chain_payload, direct_payload,
"chain payload mismatch with direct pproxy"
);
assert_eq!(*chain_payload, b"chain http->socks5");
}
(Err(e), _) => panic!("chain through eggress HTTP -> pproxy SOCKS5 failed: {e}"),
(_, Err(e)) => panic!("direct pproxy SOCKS5 failed: {e}"),
}
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_http_to_http_upstream() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_port = eggress_testkit::get_free_port().await;
let mut pproxy_child = start_pproxy_server("http", pproxy_port).await;
assert!(
wait_for_port(pproxy_port, Duration::from_secs(5)).await,
"pproxy failed to start"
);
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "http-chain"
bind = "127.0.0.1:{egress_port}"
protocols = ["http"]
[[upstreams]]
id = "pproxy-http"
uri = "http://127.0.0.1:{pproxy_port}"
[[upstream_groups]]
id = "default"
members = ["pproxy-http"]
[routing]
default = "default"
"#,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let chain_result = send_through_http(egress_addr, &target, b"chain http->http").await;
let direct_result = send_through_http(
socket_addr("127.0.0.1", pproxy_port),
&target,
b"chain http->http",
)
.await;
cancel.cancel();
let _ = jh.await;
pproxy_child.kill();
echo_jh.abort();
match (&chain_result, &direct_result) {
(Ok(chain_payload), Ok(direct_payload)) => {
assert_eq!(
chain_payload, direct_payload,
"chain payload mismatch with direct pproxy"
);
assert_eq!(*chain_payload, b"chain http->http");
}
(Err(e), _) => panic!("chain through eggress HTTP -> pproxy HTTP failed: {e}"),
(_, Err(e)) => panic!("direct pproxy HTTP failed: {e}"),
}
}
async fn send_through_trojan(
proxy_addr: std::net::SocketAddr,
target: &TargetAddr,
payload: &[u8],
password: &str,
) -> Result<Vec<u8>, String> {
use eggress_protocol_trojan::tcp::trojan_connect;
use eggress_transport_tls::TlsClientConfigBuilder;
let tcp = tokio::net::TcpStream::connect(proxy_addr)
.await
.map_err(|e| format!("connect to trojan proxy failed: {e}"))?;
let boxed: BoxStream = Box::new(tcp);
let builder = TlsClientConfigBuilder::new().with_insecure();
let tls_config = builder
.build()
.map_err(|e| format!("TLS config build failed: {e}"))?;
let mut conn = trojan_connect(boxed, target, password, "localhost", Some(tls_config))
.await
.map_err(|e| format!("trojan connect failed: {e}"))?;
conn.write_all(payload)
.await
.map_err(|e| format!("write failed: {e}"))?;
conn.flush()
.await
.map_err(|e| format!("flush failed: {e}"))?;
Ok(read_with_timeout(&mut conn, Duration::from_secs(3)).await)
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_trojan_upstream() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let cert_params = rcgen::CertificateParams::new(vec!["localhost".to_string()]).unwrap();
let key_pair = rcgen::KeyPair::generate().unwrap();
let cert = cert_params.self_signed(&key_pair).unwrap();
let cert_pem = cert.pem();
let key_pem = key_pair.serialize_pem();
let cert_file = tempfile::NamedTempFile::new().unwrap();
let key_file = tempfile::NamedTempFile::new().unwrap();
std::fs::write(cert_file.path(), &cert_pem).unwrap();
std::fs::write(key_file.path(), &key_pem).unwrap();
let password = "test-trojan-password";
let pproxy_port = eggress_testkit::get_free_port().await;
let cert_path = cert_file.path().to_str().unwrap();
let key_path = key_file.path().to_str().unwrap();
let listen = format!("trojan+ssl://127.0.0.1:{}#{}", pproxy_port, password);
let mut pproxy = start_pproxy_with_args(&[
"-l",
&listen,
"--ssl",
&format!("{cert_path},{key_path}"),
"-r",
"direct",
])
.await;
assert!(
wait_for_port(pproxy_port, Duration::from_secs(5)).await,
"pproxy failed to start"
);
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "trojan-in"
bind = "127.0.0.1:{egress_port}"
protocols = ["trojan"]
[listeners.tls]
cert = "{cert_path}"
key = "{key_path}"
[listeners.trojan]
password = "{password}"
[routing]
default = "direct"
"#,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
tokio::time::sleep(Duration::from_millis(100)).await;
let pproxy_result = send_through_trojan(
std::net::SocketAddr::new(
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
pproxy_port,
),
&target,
b"trojan differential test",
password,
)
.await;
let egress_result =
send_through_trojan(egress_addr, &target, b"trojan differential test", password).await;
cancel.cancel();
let _ = jh.await;
pproxy.kill();
echo_jh.abort();
compare_tcp_echo("pproxy", &pproxy_result, "eggress", &egress_result);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_trojan_auth_failure() {
require_differential_gate();
let (echo_addr, echo_jh) = eggress_testkit::start_echo_server().await;
let cert_params = rcgen::CertificateParams::new(vec!["localhost".to_string()]).unwrap();
let key_pair = rcgen::KeyPair::generate().unwrap();
let cert = cert_params.self_signed(&key_pair).unwrap();
let cert_pem = cert.pem();
let key_pem = key_pair.serialize_pem();
let cert_file = tempfile::NamedTempFile::new().unwrap();
let key_file = tempfile::NamedTempFile::new().unwrap();
std::fs::write(cert_file.path(), &cert_pem).unwrap();
std::fs::write(key_file.path(), &key_pem).unwrap();
let password = "correct-password";
let wrong_password = "wrong-password";
let pproxy_port = eggress_testkit::get_free_port().await;
let cert_path = cert_file.path().to_str().unwrap();
let key_path = key_file.path().to_str().unwrap();
let listen = format!("trojan+ssl://127.0.0.1:{}#{}", pproxy_port, password);
let mut pproxy = start_pproxy_with_args(&[
"-l",
&listen,
"--ssl",
&format!("{cert_path},{key_path}"),
"-r",
"direct",
])
.await;
assert!(
wait_for_port(pproxy_port, Duration::from_secs(5)).await,
"pproxy failed to start"
);
let egress_port = eggress_testkit::get_free_port().await;
let toml = format!(
r#"version = 1
[[listeners]]
name = "trojan-in"
bind = "127.0.0.1:{egress_port}"
protocols = ["trojan"]
[listeners.tls]
cert = "{cert_path}"
key = "{key_path}"
[listeners.trojan]
password = "{password}"
[routing]
default = "direct"
"#,
);
let (egress_addr, cancel, jh) = start_eggress_from_toml_running(&toml).await;
let target = TargetAddr {
host: TargetHost::Ip(echo_addr.ip()),
port: echo_addr.port(),
};
let pproxy_result = send_through_trojan(
std::net::SocketAddr::new(
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
pproxy_port,
),
&target,
b"should fail",
wrong_password,
)
.await;
let egress_result =
send_through_trojan(egress_addr, &target, b"should fail", wrong_password).await;
cancel.cancel();
let _ = jh.await;
pproxy.kill();
echo_jh.abort();
assert_coarse_failure_equivalence("pproxy", &pproxy_result, "eggress", &egress_result);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_help_output() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args(["run", "--bin", "eggress", "--", "--help"])
.output()
.expect("failed to run eggress --help");
let help = String::from_utf8_lossy(&output.stdout);
assert!(help.contains("eggress"), "help should mention program name");
assert!(
help.contains("--config") || help.contains("-c"),
"help should mention config flag"
);
assert!(
help.contains("--listen") || help.contains("-l"),
"help should mention listen flag"
);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_version_output() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args(["run", "--bin", "eggress", "--", "--version"])
.output()
.expect("failed to run eggress --version");
let version = String::from_utf8_lossy(&output.stdout);
assert!(
version.contains("eggress"),
"version output should mention program name"
);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_pproxy_translate() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args([
"run",
"--bin",
"eggress",
"--",
"pproxy",
"translate",
"--",
"-l",
"socks5://:1080",
"-r",
"socks5://127.0.0.1:8080",
])
.output()
.expect("failed to run eggress pproxy translate");
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
output.status.success(),
"pproxy translate should succeed: {stderr}"
);
assert!(
stdout.contains("[[listeners]]"),
"output should contain TOML listeners section"
);
assert!(
stdout.contains("1080"),
"output should contain the listen port"
);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_pproxy_check() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args([
"run",
"--bin",
"eggress",
"--",
"pproxy",
"check",
"--",
"-l",
"socks5://:1080",
"-r",
"socks5://127.0.0.1:8080",
])
.output()
.expect("failed to run eggress pproxy check");
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
output.status.success(),
"pproxy check should succeed: {}",
String::from_utf8_lossy(&output.stderr)
);
assert!(
stdout.contains("parity tier:"),
"check should report compatibility status: {stdout}"
);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_invalid_uri_diagnostic() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args([
"run",
"--bin",
"eggress",
"--",
"pproxy",
"translate",
"--",
"-l",
"not_a_valid_uri",
])
.output()
.expect("failed to run eggress with invalid URI");
assert!(
!output.status.success() || {
let stderr = String::from_utf8_lossy(&output.stderr);
stderr.contains("error") || stderr.contains("warning") || stderr.contains("diagnostic")
},
"invalid URI should produce an error or diagnostic"
);
}
#[tokio::test]
#[ignore = "requires EGRESS_RUN_PPROXY_DIFFERENTIAL=1"]
async fn differential_cli_unsupported_uri_diagnostic() {
require_differential_gate();
let output = std::process::Command::new("cargo")
.args([
"run",
"--bin",
"eggress",
"--",
"pproxy",
"translate",
"--",
"-l",
"ssh://:22",
])
.output()
.expect("failed to run eggress with unsupported URI");
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
!output.status.success()
|| stderr.contains("unsupported")
|| stderr.contains("diagnostic")
|| stderr.contains("warning"),
"unsupported URI should produce a diagnostic"
);
}