Compare commits

..
Author SHA1 Message Date
fanyang 3b2fd5d471 fix(gateway): route local virtual IP proxy targets to loopback
Avoid proxy loops when the target is this node's virtual IP.

- Normalize local virtual IP destinations to loopback
- Apply the behavior to TCP, KCP, and QUIC proxy paths
- Add macOS utun repro tests and normalization coverage
2026-05-28 23:44:56 +08:00
fanyang 8b1a2b5e80 fix(tunnel): silence macOS BPF unsafe warnings
Make macOS BPF ioctl wrappers explicit about unsafe calls.

- Wrap libc ioctl calls in unsafe blocks
- Keep the existing error handling unchanged
2026-05-28 23:44:56 +08:00
fanyangandGitHub 00957e5f9d feat: add Docker healthcheck (#2279) 2026-05-24 23:12:51 +08:00
fanyangandGitHub bfa3383aaa chore: update kcp-sys (#2277) 2026-05-24 23:11:16 +08:00
10 changed files with 160 additions and 94 deletions
+3
View File
@@ -42,4 +42,7 @@ EXPOSE 11011/tcp
# wss
EXPOSE 11012/tcp
HEALTHCHECK --interval=30s --timeout=10s --start-period=60s --retries=5 \
CMD ["/usr/local/bin/easytier-cli", "--rpc-portal", "127.0.0.1:15888", "--output", "json", "node", "info"]
ENTRYPOINT ["/sbin/tini", "--", "easytier-core"]
Generated
+1 -1
View File
@@ -4623,7 +4623,7 @@ dependencies = [
[[package]]
name = "kcp-sys"
version = "0.1.0"
source = "git+https://github.com/EasyTier/kcp-sys?rev=94964794caaed5d388463137da59b97499619e5f#94964794caaed5d388463137da59b97499619e5f"
source = "git+https://github.com/EasyTier/kcp-sys?rev=d7427c22d764deb1860a7d37acc446ed5033464c#d7427c22d764deb1860a7d37acc446ed5033464c"
dependencies = [
"anyhow",
"auto_impl",
+1 -1
View File
@@ -2595,7 +2595,7 @@ dependencies = [
[[package]]
name = "kcp-sys"
version = "0.1.0"
source = "git+https://github.com/EasyTier/kcp-sys?rev=94964794caaed5d388463137da59b97499619e5f#94964794caaed5d388463137da59b97499619e5f"
source = "git+https://github.com/EasyTier/kcp-sys?rev=d7427c22d764deb1860a7d37acc446ed5033464c#d7427c22d764deb1860a7d37acc446ed5033464c"
dependencies = [
"anyhow",
"auto_impl",
+1 -1
View File
@@ -224,7 +224,7 @@ service-manager = { git = "https://github.com/EasyTier/service-manager-rs.git",
zstd = { version = "0.13", optional = true }
kcp-sys = { git = "https://github.com/EasyTier/kcp-sys", rev = "94964794caaed5d388463137da59b97499619e5f", optional = true }
kcp-sys = { git = "https://github.com/EasyTier/kcp-sys", rev = "d7427c22d764deb1860a7d37acc446ed5033464c", optional = true }
prost-reflect = { version = "0.14.5", default-features = false, features = [
"derive",
-2
View File
@@ -5,12 +5,10 @@ core_clap:
en: |+
config server address, allow format:
full url: --config-server udp://127.0.0.1:22020/admin, 'udp' can be replaced with tcp, ws, wss (when config server ws is proxied to wss)
short link: --config-server https://example.com/easytier/admin, the HTTP(S) response should redirect to the full config server URL
only user name: --config-server admin, will use official server
zh-CN: |+
配置服务器地址。允许格式:
完整URL--config-server udp://127.0.0.1:22020/adminudp可以根据配置服务器替换为 tcp,ws,wss(配置服务器ws被代理为wss时)
短链接:--config-server https://example.com/easytier/adminHTTP(S) 响应应重定向到完整配置服务器 URL
仅用户名:--config-server admin,将使用官方的服务器
machine_id:
en: |+
+4 -4
View File
@@ -19,7 +19,9 @@ use tokio::task::JoinSet;
use super::{
CidrSet,
tcp_proxy::{NatDstConnector, NatDstTcpConnector, TcpProxy},
tcp_proxy::{
NatDstConnector, NatDstTcpConnector, TcpProxy, normalize_dst_for_local_virtual_ip,
},
};
use crate::utils::task::HedgeExt;
use crate::{
@@ -369,9 +371,7 @@ impl KcpProxyDst {
}
let send_to_self = global_ctx.is_ip_local_virtual_ip(&dst_ip);
if send_to_self && global_ctx.no_tun() {
dst_socket = format!("127.0.0.1:{}", dst_socket.port()).parse().unwrap();
}
dst_socket = normalize_dst_for_local_virtual_ip(&global_ctx, dst_socket);
let acl_handler = ProxyAclHandler {
acl_filter: global_ctx.get_acl_filter().clone(),
+2 -4
View File
@@ -2,7 +2,7 @@ use crate::common::PeerId;
use crate::common::acl_processor::PacketInfo;
use crate::common::global_ctx::{ArcGlobalCtx, GlobalCtx};
use crate::gateway::CidrSet;
use crate::gateway::tcp_proxy::{NatDstConnector, TcpProxy};
use crate::gateway::tcp_proxy::{NatDstConnector, TcpProxy, normalize_dst_for_local_virtual_ip};
use crate::gateway::wrapped_proxy::{ProxyAclHandler, TcpProxyForWrappedSrcTrait};
use crate::peers::PeerPacketFilter;
use crate::peers::peer_manager::PeerManager;
@@ -748,9 +748,7 @@ impl QuicStreamReceiver {
}
let send_to_self = global_ctx.is_ip_local_virtual_ip(&dst_ip);
if send_to_self && global_ctx.no_tun() {
dst_socket = format!("127.0.0.1:{}", dst_socket.port()).parse()?;
}
dst_socket = normalize_dst_for_local_virtual_ip(&global_ctx, dst_socket);
let acl_handler = ProxyAclHandler {
acl_filter: global_ctx.get_acl_filter().clone(),
+135 -8
View File
@@ -9,7 +9,7 @@ use pnet::packet::ip::IpNextHeaderProtocols;
use pnet::packet::ipv4::{Ipv4Packet, MutableIpv4Packet};
use pnet::packet::tcp::{MutableTcpPacket, TcpPacket, ipv4_checksum};
use socket2::{SockRef, TcpKeepalive};
use std::net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4};
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, SocketAddrV4};
use std::sync::atomic::{AtomicBool, AtomicU16};
use std::sync::{Arc, Weak};
use std::time::{Duration, Instant};
@@ -40,6 +40,20 @@ use super::CidrSet;
#[cfg(feature = "smoltcp")]
use super::tokio_smoltcp::{self, Net, NetConfig, channel_device};
pub(crate) fn normalize_dst_for_local_virtual_ip(
global_ctx: &GlobalCtx,
dst: SocketAddr,
) -> SocketAddr {
if !global_ctx.is_ip_local_virtual_ip(&dst.ip()) {
return dst;
}
match dst {
SocketAddr::V4(addr) => SocketAddr::new(Ipv4Addr::LOCALHOST.into(), addr.port()),
SocketAddr::V6(addr) => SocketAddr::new(Ipv6Addr::LOCALHOST.into(), addr.port()),
}
}
#[async_trait::async_trait]
pub(crate) trait NatDstConnector: Send + Sync + Clone + 'static {
type DstStream: AsyncRead + AsyncWrite + Unpin + Send;
@@ -762,13 +776,7 @@ impl<C: NatDstConnector> TcpProxy<C> {
return;
}
let nat_dst = if global_ctx.is_ip_local_virtual_ip(&nat_entry.real_dst.ip()) {
format!("127.0.0.1:{}", nat_entry.real_dst.port())
.parse()
.unwrap()
} else {
nat_entry.real_dst
};
let nat_dst = normalize_dst_for_local_virtual_ip(&global_ctx, nat_entry.real_dst);
global_ctx
.stats_manager()
@@ -1033,3 +1041,122 @@ impl<C: NatDstConnector> TcpProxyRpcService<C> {
}
}
}
#[cfg(test)]
mod tests {
use std::net::{Ipv4Addr, SocketAddr};
use super::normalize_dst_for_local_virtual_ip;
#[tokio::test]
async fn normalize_dst_for_local_virtual_ip_maps_to_loopback() {
let global_ctx = crate::common::global_ctx::tests::get_mock_global_ctx();
global_ctx.set_ipv4(Some("10.254.229.6/24".parse().unwrap()));
let local_virtual = SocketAddr::from(([10, 254, 229, 6], 22));
let normalized = normalize_dst_for_local_virtual_ip(&global_ctx, local_virtual);
assert_eq!(normalized, SocketAddr::from((Ipv4Addr::LOCALHOST, 22)));
let remote_virtual = SocketAddr::from(([10, 254, 229, 7], 22));
assert_eq!(
normalize_dst_for_local_virtual_ip(&global_ctx, remote_virtual),
remote_virtual
);
}
}
#[cfg(all(test, target_os = "macos", feature = "tun"))]
mod macos_utun_tests {
use std::{net::SocketAddr, sync::Arc, time::Duration};
use tokio::net::{TcpListener, TcpStream};
use super::{NatDstConnector as _, NatDstTcpConnector, normalize_dst_for_local_virtual_ip};
use crate::{
common::config::{ConfigLoader, TomlConfigLoader},
instance::instance::Instance,
};
async fn run_instance_with_utun(ipv4: &str, enable_kcp: bool, enable_quic: bool) -> Instance {
let config = TomlConfigLoader::default();
config.set_inst_name(format!("macos-utun-tcp-proxy-repro-{ipv4}"));
config.set_ipv4(Some(ipv4.parse().unwrap()));
config.set_ipv6(None);
config.set_listeners(Vec::new());
let mut flags = config.get_flags();
flags.enable_kcp_proxy = enable_kcp;
flags.enable_quic_proxy = enable_quic;
flags.use_smoltcp = false;
config.set_flags(flags);
let mut instance = Instance::new(config);
instance.run().await.expect(
"failed to create macOS utun device; run this ignored reproducer with root privileges",
);
instance
}
async fn assert_wildcard_listener_reachable_via_local_virtual_ip(
ipv4: &str,
enable_kcp: bool,
enable_quic: bool,
) {
let mut instance = run_instance_with_utun(ipv4, enable_kcp, enable_quic).await;
let virtual_ip = instance.get_global_ctx().get_ipv4().unwrap().address();
let listener = Arc::new(TcpListener::bind("0.0.0.0:0").await.unwrap());
let port = listener.local_addr().unwrap().port();
let baseline_listener = listener.clone();
let baseline_accept = tokio::spawn(async move { baseline_listener.accept().await });
let baseline_connect = TcpStream::connect((std::net::Ipv4Addr::LOCALHOST, port));
let (baseline_connect, baseline_accept) = tokio::join!(baseline_connect, baseline_accept);
baseline_connect.unwrap();
baseline_accept.unwrap().unwrap();
let test_listener = listener.clone();
let mut accept = tokio::spawn(async move { test_listener.accept().await });
let dst = SocketAddr::new(virtual_ip.into(), port);
let dst = normalize_dst_for_local_virtual_ip(&instance.get_global_ctx(), dst);
assert_eq!(dst.ip(), std::net::Ipv4Addr::LOCALHOST);
let connect = NatDstTcpConnector {}.connect("0.0.0.0:0".parse().unwrap(), dst);
let connect = tokio::time::timeout(Duration::from_secs(3), connect).await;
let accept_result = tokio::time::timeout(Duration::from_secs(1), &mut accept).await;
if accept_result.is_err() {
accept.abort();
}
instance.clear_resources().await;
match connect {
Ok(Ok(_)) => {}
Ok(Err(error)) => panic!("connect to local EasyTier virtual IP failed: {error:?}"),
Err(_) => panic!("connect to local EasyTier virtual IP timed out"),
}
match accept_result {
Ok(Ok(Ok(_))) => {}
Ok(Ok(Err(error))) => panic!("listener accept failed: {error:?}"),
Ok(Err(error)) => panic!("listener task failed: {error:?}"),
Err(_) => panic!(
"listener did not accept connection to local EasyTier virtual IP {virtual_ip}:{port}"
),
}
}
#[tokio::test]
#[ignore = "requires root and a real macOS utun device; covers issue #2296"]
async fn macos_utun_kcp_proxy_dst_local_virtual_ip_reaches_wildcard_listener() {
assert_wildcard_listener_reachable_via_local_virtual_ip("10.254.229.6/24", true, false)
.await;
}
#[tokio::test]
#[ignore = "requires root and a real macOS utun device; covers issue #2296"]
async fn macos_utun_quic_proxy_dst_local_virtual_ip_reaches_wildcard_listener() {
assert_wildcard_listener_reachable_via_local_virtual_ip("10.254.230.6/24", false, true)
.await;
}
}
@@ -232,7 +232,7 @@ fn iowr<T>(group: u8, num: u8) -> libc::c_ulong {
}
unsafe fn ioctl_ptr<T>(fd: libc::c_int, req: libc::c_ulong, arg: *mut T) -> io::Result<()> {
let ret = libc::ioctl(fd, req, arg);
let ret = unsafe { libc::ioctl(fd, req, arg) };
if ret < 0 {
return Err(io::Error::last_os_error());
}
@@ -240,7 +240,7 @@ unsafe fn ioctl_ptr<T>(fd: libc::c_int, req: libc::c_ulong, arg: *mut T) -> io::
}
unsafe fn ioctl_void(fd: libc::c_int, req: libc::c_ulong) -> io::Result<()> {
let ret = libc::ioctl(fd, req);
let ret = unsafe { libc::ioctl(fd, req) };
if ret < 0 {
return Err(io::Error::last_os_error());
}
+11 -71
View File
@@ -78,22 +78,6 @@ impl TunnelConnector for ConfigServerConnector {
}
}
fn parse_config_server_input(config_server_url_s: &str) -> Result<Url> {
match Url::parse(config_server_url_s) {
Ok(u) => Ok(u),
Err(_) => format!(
"udp://config-server.easytier.cn:22020/{}",
config_server_url_s
)
.parse()
.with_context(|| "failed to parse config server URL"),
}
}
fn is_config_server_http_short_link(url: &Url) -> bool {
matches!(url.scheme(), "http" | "https")
}
impl WebClient {
pub fn new<T: TunnelConnector + 'static, S: ToString, H: ToString>(
connector: T,
@@ -256,7 +240,15 @@ pub async fn run_web_client(
) -> Result<WebClient> {
let machine_id = resolve_machine_id(&machine_id_opts)
.with_context(|| "failed to resolve machine id for web client")?;
let config_server_url = parse_config_server_input(config_server_url_s)?;
let config_server_url = match Url::parse(config_server_url_s) {
Ok(u) => u,
Err(_) => format!(
"udp://config-server.easytier.cn:22020/{}",
config_server_url_s
)
.parse()
.with_context(|| "failed to parse config server URL")?,
};
TunnelScheme::try_from(&config_server_url).map_err(|_| {
anyhow::anyhow!(
@@ -266,8 +258,7 @@ pub async fn run_web_client(
})?;
let mut c_url = config_server_url.clone();
// Keep HTTP(S) paths so HttpTunnelConnector can request short links and handle their redirects.
if !matches!(c_url.scheme(), "ws" | "wss") && !is_config_server_http_short_link(&c_url) {
if !matches!(c_url.scheme(), "ws" | "wss") {
c_url.set_path("");
}
let token = config_server_url
@@ -314,15 +305,7 @@ pub async fn run_web_client(
mod tests {
use std::sync::{Arc, atomic::AtomicBool};
use crate::{
common::{MachineIdOptions, config::TomlConfigLoader, global_ctx::GlobalCtx},
instance_manager::NetworkInstanceManager,
tunnel::TunnelConnector,
};
use tokio::{
io::{AsyncReadExt as _, AsyncWriteExt as _},
net::TcpListener,
};
use crate::{common::MachineIdOptions, instance_manager::NetworkInstanceManager};
#[tokio::test]
async fn test_manager_wait() {
@@ -380,47 +363,4 @@ mod tests {
assert!(!client.is_connected());
drop(client);
}
#[tokio::test]
async fn config_server_short_link_uses_existing_http_redirect_connector() {
let target_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let target_addr = target_listener.local_addr().unwrap();
let target_task = tokio::spawn(async move {
let _ = target_listener.accept().await.unwrap();
});
let http_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let http_addr = http_listener.local_addr().unwrap();
let http_task = tokio::spawn(async move {
let (mut stream, _) = http_listener.accept().await.unwrap();
let mut buf = [0u8; 4096];
let n = stream.read(&mut buf).await.unwrap();
let req = String::from_utf8_lossy(&buf[..n]).to_string();
let resp = format!(
"HTTP/1.1 302 Found\r\nLocation: tcp://{}\r\nContent-Length: 0\r\n\r\n",
target_addr
);
stream.write_all(resp.as_bytes()).await.unwrap();
req
});
let config = TomlConfigLoader::default();
let global_ctx = Arc::new(GlobalCtx::new(config));
let mut flags = global_ctx.get_flags();
flags.bind_device = false;
global_ctx.set_flags(flags);
let url: url::Url = format!("http://{}/short-token", http_addr).parse().unwrap();
let mut connector = super::ConfigServerConnector { url, global_ctx };
let tunnel = connector.connect().await.unwrap();
let req = http_task.await.unwrap();
assert!(req.starts_with("GET /short-token "));
let info = tunnel.info().unwrap();
assert_eq!(
info.resolved_remote_addr.unwrap().url,
format!("tcp://{}", target_addr)
);
target_task.await.unwrap();
}
}