From d3a5ae684a5d5d5e7c869f819299515c1e822de0 Mon Sep 17 00:00:00 2001 From: KKRainbow <5665404+KKRainbow@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:25:08 +0800 Subject: [PATCH] fix(peers): keep TCP hole-punched connections alive with 1s pings (#2632) Disable ping interval backoff for TCP hole-punched connections so idle connections continue to send keepalive traffic every second. Preserve the existing backoff and loss handling for other connections. Keep randomized backoff above zero and cover the one-second schedule and the existing backoff and loss retry behavior with tests. * refactor(peers): model connection origins at admission Record manual, direct, listener, TCP/UDP hole-punch, and attached origins when constructing peer connections. Derive hole-punch state from this origin instead of maintaining separate mutable flags. Choose the one-second TCP hole-punch ping limit in PeerConn and pass only a maximum interval to the pinger. Keep other origins on the existing backoff schedule without inspecting tunnel type strings. Keep origin selection internal and preserve the dedicated attached admission paths. Update public admission callers and cover origin propagation, relay restrictions, and ping interval limits. --- easytier-core/src/connectivity/direct/mod.rs | 11 +- .../src/connectivity/hole_punch/mod.rs | 13 ++- .../connectivity/hole_punch/peer_adapters.rs | 19 ++- .../src/connectivity/hole_punch/tcp.rs | 29 ++++- .../connectivity/hole_punch/udp/runtime.rs | 23 +++- easytier-core/src/connectivity/manual/mod.rs | 4 +- easytier-core/src/gateway/dataplane/tests.rs | 4 +- easytier-core/src/instance/test_utils.rs | 5 +- easytier-core/src/peers/admission.rs | 2 +- easytier-core/src/peers/conn/peer_conn.rs | 64 ++++++++-- .../src/peers/conn/peer_conn_ping.rs | 76 +++++++++++- easytier-core/src/peers/mod.rs | 7 +- easytier-core/src/peers/peer_manager.rs | 109 ++++++------------ easytier-core/src/peers/tests.rs | 1 - easytier-web/src/central_network/gateway.rs | 2 +- easytier/docs/credential_peer.md | 2 +- easytier/src/tests/three_node.rs | 2 +- 17 files changed, 256 insertions(+), 117 deletions(-) diff --git a/easytier-core/src/connectivity/direct/mod.rs b/easytier-core/src/connectivity/direct/mod.rs index 35379bc1..cf6269d7 100644 --- a/easytier-core/src/connectivity/direct/mod.rs +++ b/easytier-core/src/connectivity/direct/mod.rs @@ -28,8 +28,9 @@ use crate::{ foundation::task::{PeerTaskLauncher, PeerTaskManager}, host::dns::DnsResolver, peers::{ - conn::peer_conn::PeerConnId, foreign_network::ForeignNetworkRpcRegistrar, - peer_manager::PeerManagerCore, peer_rpc::PeerRpcManager, + PeerConnectionOrigin, conn::peer_conn::PeerConnId, + foreign_network::ForeignNetworkRpcRegistrar, peer_manager::PeerManagerCore, + peer_rpc::PeerRpcManager, }, process_runtime::ProtectedTcpPortRegistry, proto::{ @@ -870,7 +871,11 @@ where dst_peer_id: PeerId, ) -> anyhow::Result<(PeerId, PeerConnId)> { self.peer_manager - .add_client_tunnel_with_peer_id_hint(tunnel, true, Some(dst_peer_id)) + .add_client_tunnel_with_peer_id_hint( + tunnel, + PeerConnectionOrigin::Direct, + Some(dst_peer_id), + ) .await .map_err(Into::into) } diff --git a/easytier-core/src/connectivity/hole_punch/mod.rs b/easytier-core/src/connectivity/hole_punch/mod.rs index ec1e951d..c48afc59 100644 --- a/easytier-core/src/connectivity/hole_punch/mod.rs +++ b/easytier-core/src/connectivity/hole_punch/mod.rs @@ -1,5 +1,6 @@ use async_trait::async_trait; +use crate::peers::PeerConnectionOrigin; use crate::proto::rpc_types::{controller::BaseController, handler::Handler}; use crate::tunnel::Tunnel; @@ -27,7 +28,15 @@ pub(crate) trait HolePunchRpcRegistry: Send + Sync + 'static { #[async_trait] pub(crate) trait HolePunchTunnelSink: Send + Sync + 'static { - async fn add_client_tunnel(&self, tunnel: Box) -> anyhow::Result<()>; + async fn add_client_tunnel( + &self, + tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()>; - async fn add_server_tunnel(&self, tunnel: Box) -> anyhow::Result<()>; + async fn add_server_tunnel( + &self, + tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()>; } diff --git a/easytier-core/src/connectivity/hole_punch/peer_adapters.rs b/easytier-core/src/connectivity/hole_punch/peer_adapters.rs index 07d0a7c1..bf700a50 100644 --- a/easytier-core/src/connectivity/hole_punch/peer_adapters.rs +++ b/easytier-core/src/connectivity/hole_punch/peer_adapters.rs @@ -9,7 +9,7 @@ use quanta::Instant; use crate::{ config::{P2pPolicyFlags, PeerId}, foundation::task::ExternalTaskSignal, - peers::peer_manager::PeerManagerCore, + peers::{PeerConnectionOrigin, peer_manager::PeerManagerCore}, proto::{ common::NatType, peer_rpc::{ @@ -97,16 +97,25 @@ impl UdpHolePunchRpcSource for PeerManagerCore { #[async_trait] impl HolePunchTunnelSink for PeerManagerCore { - async fn add_client_tunnel(&self, tunnel: Box) -> anyhow::Result<()> { - PeerManagerCore::add_client_tunnel(self, tunnel, false) + async fn add_client_tunnel( + &self, + tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + self.add_client_tunnel_with_peer_id_hint(tunnel, origin, None) .await .map(|_| ()) .map_err(anyhow::Error::from) } - async fn add_server_tunnel(&self, tunnel: Box) -> anyhow::Result<()> { - PeerManagerCore::add_tunnel_as_server(self, tunnel, false) + async fn add_server_tunnel( + &self, + tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + self.add_tunnel_as_server_with_origin(tunnel, origin) .await + .map(|_| ()) .map_err(anyhow::Error::from) } } diff --git a/easytier-core/src/connectivity/hole_punch/tcp.rs b/easytier-core/src/connectivity/hole_punch/tcp.rs index 2aee0718..6a7eda9e 100644 --- a/easytier-core/src/connectivity/hole_punch/tcp.rs +++ b/easytier-core/src/connectivity/hole_punch/tcp.rs @@ -31,6 +31,7 @@ use crate::{ foundation::task::{ ExternalTaskSignal, PeerTaskLauncher, PeerTaskManager, reap_joinset_background, }, + peers::PeerConnectionOrigin, proto::{ common::{NatType, PeerFeatureFlag}, peer_rpc::{ @@ -159,8 +160,16 @@ where .await .map_err(TcpHolePunchTransportError::Upgrade)?; match admission { - TcpHolePunchAdmission::Client => self.tunnel_sink.add_client_tunnel(tunnel).await, - TcpHolePunchAdmission::Server => self.tunnel_sink.add_server_tunnel(tunnel).await, + TcpHolePunchAdmission::Client => { + self.tunnel_sink + .add_client_tunnel(tunnel, PeerConnectionOrigin::TcpHolePunch) + .await + } + TcpHolePunchAdmission::Server => { + self.tunnel_sink + .add_server_tunnel(tunnel, PeerConnectionOrigin::TcpHolePunch) + .await + } } .map_err(TcpHolePunchTransportError::Admission) } @@ -181,7 +190,7 @@ where ))); }; self.tunnel_sink - .add_server_tunnel(tunnel) + .add_server_tunnel(tunnel, PeerConnectionOrigin::TcpHolePunch) .await .map_err(TcpHolePunchTransportError::Admission) } @@ -979,7 +988,12 @@ mod tests { #[async_trait] impl HolePunchTunnelSink for MockTunnelSink { - async fn add_client_tunnel(&self, _tunnel: Box) -> anyhow::Result<()> { + async fn add_client_tunnel( + &self, + _tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + assert_eq!(origin, PeerConnectionOrigin::TcpHolePunch); if self.fail_client_admission.load(Ordering::Relaxed) { anyhow::bail!("mock client admission failure"); } @@ -987,7 +1001,12 @@ mod tests { Ok(()) } - async fn add_server_tunnel(&self, _tunnel: Box) -> anyhow::Result<()> { + async fn add_server_tunnel( + &self, + _tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + assert_eq!(origin, PeerConnectionOrigin::TcpHolePunch); self.servers.fetch_add(1, Ordering::Relaxed); Ok(()) } diff --git a/easytier-core/src/connectivity/hole_punch/udp/runtime.rs b/easytier-core/src/connectivity/hole_punch/udp/runtime.rs index b23f16c1..5be63eb3 100644 --- a/easytier-core/src/connectivity/hole_punch/udp/runtime.rs +++ b/easytier-core/src/connectivity/hole_punch/udp/runtime.rs @@ -14,6 +14,7 @@ use crate::{ transport::{ConnectedTransport, ConnectedUdpSession}, }, foundation::task::ExternalTaskSignal, + peers::PeerConnectionOrigin, socket::{ ListenerConnectionCounter, SocketContext, udp::{UdpBindOptions, UdpSession, VirtualUdpSocket, VirtualUdpSocketFactory}, @@ -271,7 +272,9 @@ where requested_url: url::Url, ) -> anyhow::Result<()> { let tunnel = self.upgrade(connected, requested_url).await?; - self.tunnel_sink.add_client_tunnel(tunnel).await + self.tunnel_sink + .add_client_tunnel(tunnel, PeerConnectionOrigin::UdpHolePunch) + .await } async fn add_server_transport( @@ -280,7 +283,9 @@ where requested_url: url::Url, ) -> anyhow::Result<()> { let tunnel = self.upgrade(connected, requested_url).await?; - self.tunnel_sink.add_server_tunnel(tunnel).await + self.tunnel_sink + .add_server_tunnel(tunnel, PeerConnectionOrigin::UdpHolePunch) + .await } } @@ -428,12 +433,22 @@ mod tests { #[async_trait] impl HolePunchTunnelSink for MockTunnelSink { - async fn add_client_tunnel(&self, _tunnel: Box) -> anyhow::Result<()> { + async fn add_client_tunnel( + &self, + _tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + assert_eq!(origin, PeerConnectionOrigin::UdpHolePunch); self.clients.fetch_add(1, Ordering::Relaxed); Ok(()) } - async fn add_server_tunnel(&self, _tunnel: Box) -> anyhow::Result<()> { + async fn add_server_tunnel( + &self, + _tunnel: Box, + origin: PeerConnectionOrigin, + ) -> anyhow::Result<()> { + assert_eq!(origin, PeerConnectionOrigin::UdpHolePunch); self.servers.fetch_add(1, Ordering::Relaxed); Ok(()) } diff --git a/easytier-core/src/connectivity/manual/mod.rs b/easytier-core/src/connectivity/manual/mod.rs index 74e5d4e8..40602164 100644 --- a/easytier-core/src/connectivity/manual/mod.rs +++ b/easytier-core/src/connectivity/manual/mod.rs @@ -26,7 +26,7 @@ use crate::{ }, events::{CoreEvent, CoreEventSink}, host::dns::{DnsQuery, DnsResolver}, - peers::peer_manager::PeerManagerCore, + peers::{PeerConnectionOrigin, peer_manager::PeerManagerCore}, proto::common::TunnelInfo, socket::{ IpVersion, SocketContext, @@ -891,7 +891,7 @@ where let (peer_id, conn_id) = with_timeout_budget("handshake", started_at, connect_timeout, async move { peer_manager - .add_client_tunnel_with_peer_id_hint(tunnel, true, None) + .add_client_tunnel_with_peer_id_hint(tunnel, PeerConnectionOrigin::Manual, None) .await .map_err(anyhow::Error::from) }) diff --git a/easytier-core/src/gateway/dataplane/tests.rs b/easytier-core/src/gateway/dataplane/tests.rs index 8ac22793..40b01715 100644 --- a/easytier-core/src/gateway/dataplane/tests.rs +++ b/easytier-core/src/gateway/dataplane/tests.rs @@ -113,8 +113,8 @@ async fn setup_data_plane_pair() -> (DataPlaneEndpoint, DataPlaneEndpoint) { let client_tunnel = registry.connect(listener_id).unwrap().into_tunnel(); let server_tunnel = listener.accept().await.unwrap().into_tunnel(); let (client, server) = tokio::join!( - b.peer_manager.add_client_tunnel(client_tunnel, true), - a.peer_manager.add_tunnel_as_server(server_tunnel, true), + b.peer_manager.add_client_tunnel(client_tunnel), + a.peer_manager.add_tunnel_as_server(server_tunnel), ); client.unwrap(); server.unwrap(); diff --git a/easytier-core/src/instance/test_utils.rs b/easytier-core/src/instance/test_utils.rs index ddee05c4..c7eebe48 100644 --- a/easytier-core/src/instance/test_utils.rs +++ b/easytier-core/src/instance/test_utils.rs @@ -46,11 +46,8 @@ where pub async fn admit_client_tunnel_for_test( &self, tunnel: Box, - is_directly_connected: bool, ) -> Result<(crate::config::PeerId, PeerConnId), crate::peers::error::Error> { - self.peer_manager - .add_client_tunnel(tunnel, is_directly_connected) - .await + self.peer_manager.add_client_tunnel(tunnel).await } #[doc(hidden)] diff --git a/easytier-core/src/peers/admission.rs b/easytier-core/src/peers/admission.rs index f082207d..11f7872b 100644 --- a/easytier-core/src/peers/admission.rs +++ b/easytier-core/src/peers/admission.rs @@ -63,7 +63,7 @@ impl AcceptedTunnelHandler for PeerAcceptedTunnelHandler { tracing::error!(error = %error, "handle conn error"); return Err(anyhow::anyhow!(error)); }; - if let Err(error) = peer_manager.add_tunnel_as_server(tunnel, true).await { + if let Err(error) = peer_manager.add_tunnel_as_server(tunnel).await { self.events.emit(CoreEvent::TunnelAdmissionFailed { local_url, remote_url, diff --git a/easytier-core/src/peers/conn/peer_conn.rs b/easytier-core/src/peers/conn/peer_conn.rs index 19c1840a..072e8707 100644 --- a/easytier-core/src/peers/conn/peer_conn.rs +++ b/easytier-core/src/peers/conn/peer_conn.rs @@ -293,9 +293,6 @@ pub struct PeerConn { info: Option, is_client: Option, - // remote or local - is_hole_punched: bool, - close_event_notifier: Arc, ctrl_resp_sender: broadcast::Sender, @@ -333,7 +330,7 @@ impl PeerConn { tunnel, None, peer_session_store, - PeerConnectionOrigin::Network, + PeerConnectionOrigin::Manual, ) } @@ -396,8 +393,6 @@ impl PeerConn { info: None, is_client: None, - is_hole_punched: true, - close_event_notifier: Arc::new(PeerConnCloseNotify::new(conn_id)), ctrl_resp_sender: ctrl_sender, @@ -443,12 +438,23 @@ impl PeerConn { self.origin == PeerConnectionOrigin::Attached } - pub fn set_is_hole_punched(&mut self, is_hole_punched: bool) { - self.is_hole_punched = is_hole_punched; + pub fn is_hole_punched(&self) -> bool { + matches!( + self.origin, + PeerConnectionOrigin::TcpHolePunch | PeerConnectionOrigin::UdpHolePunch + ) } - pub fn is_hole_punched(&self) -> bool { - self.is_hole_punched + fn max_ping_interval(&self) -> Duration { + match self.origin { + // TCP hole-punched connections need frequent traffic to stay alive. + PeerConnectionOrigin::TcpHolePunch => Duration::from_secs(1), + PeerConnectionOrigin::Manual + | PeerConnectionOrigin::Direct + | PeerConnectionOrigin::Listener + | PeerConnectionOrigin::UdpHolePunch + | PeerConnectionOrigin::Attached => Duration::from_secs(32), + } } pub fn is_closed(&self) -> bool { @@ -1377,6 +1383,7 @@ impl PeerConn { self.context.clone(), self.get_conn_info().network_name, self.liveness.clone(), + self.max_ping_interval(), ); let close_event_notifier = self.close_event_notifier.clone(); @@ -1526,3 +1533,40 @@ impl Drop for PeerConn { self.close_event_notifier.notify_close(); } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::{peers::test_support::NoopPeerContext, tunnel::ring::create_ring_tunnel_pair}; + + #[tokio::test] + async fn connection_origin_determines_hole_punch_and_ping_policy() { + for (origin, is_hole_punched, max_interval) in [ + (PeerConnectionOrigin::Manual, false, 32), + (PeerConnectionOrigin::Direct, false, 32), + (PeerConnectionOrigin::Listener, false, 32), + (PeerConnectionOrigin::TcpHolePunch, true, 1), + (PeerConnectionOrigin::UdpHolePunch, true, 32), + (PeerConnectionOrigin::Attached, false, 32), + ] { + // Admission determines the policy even when the transport is a ring. + let (tunnel, _remote_tunnel) = create_ring_tunnel_pair(); + let conn = PeerConn::new_with_peer_id_hint_and_origin( + 1, + Arc::new(NoopPeerContext::default()), + tunnel, + None, + Arc::new(PeerSessionStore::new()), + origin, + ); + + assert_eq!(conn.is_hole_punched(), is_hole_punched, "{origin:?}"); + assert_eq!( + conn.max_ping_interval(), + Duration::from_secs(max_interval), + "{origin:?}", + ); + assert_eq!(conn.is_attached(), origin == PeerConnectionOrigin::Attached); + } + } +} diff --git a/easytier-core/src/peers/conn/peer_conn_ping.rs b/easytier-core/src/peers/conn/peer_conn_ping.rs index 88f9daf4..caa0836d 100644 --- a/easytier-core/src/peers/conn/peer_conn_ping.rs +++ b/easytier-core/src/peers/conn/peer_conn_ping.rs @@ -34,6 +34,7 @@ struct PingIntervalController { loss_counter: Arc, interval: Interval, + max_interval: Duration, logic_time: u64, last_send_logic_time: u64, @@ -53,19 +54,25 @@ impl std::fmt::Debug for PingIntervalController { .field("last_send_logic_time", &self.last_send_logic_time) .field("backoff_idx", &self.backoff_idx) .field("max_backoff_idx", &self.max_backoff_idx) + .field("max_interval", &self.max_interval) .field("last_throughput", &self.last_throughput) .finish() } } impl PingIntervalController { - fn new(throughput: Arc, loss_counter: Arc) -> Self { + fn new( + throughput: Arc, + loss_counter: Arc, + max_interval: Duration, + ) -> Self { let last_throughput = (*throughput).clone(); Self { throughput, loss_counter, interval: interval(Duration::from_secs(1)), + max_interval, logic_time: 0, last_send_logic_time: 0, @@ -99,7 +106,8 @@ impl PingIntervalController { self.last_throughput = (*self.throughput).clone(); - if (self.logic_time - self.last_send_logic_time) < (1 << self.backoff_idx) { + let send_interval = Duration::from_secs(1 << self.backoff_idx).min(self.max_interval); + if Duration::from_secs(self.logic_time - self.last_send_logic_time) < send_interval { return false; } @@ -126,6 +134,7 @@ pub struct PeerConnPinger { context: ArcPeerContext, network_name: String, liveness: PeerConnLiveness, + max_interval: Duration, } impl std::fmt::Debug for PeerConnPinger { @@ -150,6 +159,7 @@ impl PeerConnPinger { context: ArcPeerContext, network_name: String, liveness: PeerConnLiveness, + max_interval: Duration, ) -> Self { Self { my_peer_id, @@ -162,6 +172,7 @@ impl PeerConnPinger { context, network_name, liveness, + max_interval, } } @@ -245,8 +256,10 @@ impl PeerConnPinger { let mut controller_tasks = JoinSet::new(); let throughput = self.throughput_stats.clone(); let controller_loss_counter = loss_counter.clone(); + let max_interval = self.max_interval; controller_tasks.spawn(async move { - let mut controller = PingIntervalController::new(throughput, controller_loss_counter); + let mut controller = + PingIntervalController::new(throughput, controller_loss_counter, max_interval); loop { controller.tick().await; if !controller.should_send_ping() { @@ -337,6 +350,61 @@ mod tests { }, }; + #[cfg(not(target_os = "wasi"))] + #[tokio::test(start_paused = true)] + async fn one_second_limit_disables_ping_backoff() { + let mut controller = PingIntervalController::new( + Arc::new(Throughput::new()), + Arc::new(AtomicU32::new(0)), + Duration::from_secs(1), + ); + + let started_at = tokio::time::Instant::now(); + for second in 0..100 { + controller.tick().await; + assert_eq!(started_at.elapsed(), Duration::from_secs(second)); + assert!(controller.should_send_ping()); + assert!(!controller.should_send_ping()); + } + } + + #[cfg(not(target_os = "wasi"))] + #[tokio::test(start_paused = true)] + async fn ping_backoff_respects_non_power_of_two_limit() { + let mut controller = PingIntervalController::new( + Arc::new(Throughput::new()), + Arc::new(AtomicU32::new(0)), + Duration::from_secs(3), + ); + + for tick in 1..=100 { + controller.tick().await; + assert_eq!(controller.should_send_ping(), tick == 1 || tick % 3 == 0); + } + } + + #[cfg(not(target_os = "wasi"))] + #[tokio::test(start_paused = true)] + async fn default_limit_preserves_backoff_and_loss_retries() { + let loss_counter = Arc::new(AtomicU32::new(0)); + let mut controller = PingIntervalController::new( + Arc::new(Throughput::new()), + loss_counter.clone(), + Duration::from_secs(32), + ); + + for second in 1..=14 { + controller.tick().await; + assert_eq!(controller.should_send_ping(), matches!(second, 1 | 3 | 7)); + } + + loss_counter.store(1, Ordering::Relaxed); + for _ in 0..3 { + controller.tick().await; + assert!(controller.should_send_ping()); + } + } + #[tokio::test(flavor = "current_thread")] async fn ingress_traffic_does_not_mask_failed_round_trips() { let (local_tunnel, _remote_tunnel) = create_ring_tunnel_pair(); @@ -354,6 +422,7 @@ mod tests { Arc::new(NoopPeerContext::default()), "test".to_owned(), PeerConnLiveness::new(), + Duration::from_secs(32), ); let ingress = tokio::spawn(async move { @@ -412,6 +481,7 @@ mod tests { Arc::new(NoopPeerContext::default()), "test".to_owned(), local_liveness, + Duration::from_secs(32), ); let result = timeout(Duration::from_secs(12), pinger.pingpong()).await; diff --git a/easytier-core/src/peers/mod.rs b/easytier-core/src/peers/mod.rs index 4b796380..abd980b7 100644 --- a/easytier-core/src/peers/mod.rs +++ b/easytier-core/src/peers/mod.rs @@ -27,9 +27,14 @@ use tokio::sync::mpsc::error::{SendError, TryRecvError, TrySendError}; use self::conn::peer_conn::PeerConnId; use crate::config::PeerId; +/// The local entry point that created a peer connection. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub(crate) enum PeerConnectionOrigin { - Network, + Manual, + Direct, + Listener, + TcpHolePunch, + UdpHolePunch, Attached, } diff --git a/easytier-core/src/peers/peer_manager.rs b/easytier-core/src/peers/peer_manager.rs index 489bc1e8..e630126f 100644 --- a/easytier-core/src/peers/peer_manager.rs +++ b/easytier-core/src/peers/peer_manager.rs @@ -1565,21 +1565,19 @@ impl PeerManagerCore { pub async fn add_client_tunnel( &self, tunnel: Box, - is_directly_connected: bool, ) -> Result<(PeerId, PeerConnId), Error> { - self.peer_connection_admission - .add_client_tunnel(tunnel, is_directly_connected) + self.add_client_tunnel_with_peer_id_hint(tunnel, PeerConnectionOrigin::Manual, None) .await } - pub async fn add_client_tunnel_with_peer_id_hint( + pub(crate) async fn add_client_tunnel_with_peer_id_hint( &self, tunnel: Box, - is_directly_connected: bool, + origin: PeerConnectionOrigin, peer_id_hint: Option, ) -> Result<(PeerId, PeerConnId), Error> { self.peer_connection_admission - .add_client_tunnel_with_peer_id_hint(tunnel, is_directly_connected, peer_id_hint) + .add_client_tunnel_with_peer_id_hint(tunnel, origin, peer_id_hint) .await } @@ -1587,18 +1585,23 @@ impl PeerManagerCore { &self, tunnel: Box, ) -> Result<(PeerId, PeerConnId), Error> { - self.peer_connection_admission - .add_client_tunnel_with_origin(tunnel, true, None, PeerConnectionOrigin::Attached) + self.add_client_tunnel_with_peer_id_hint(tunnel, PeerConnectionOrigin::Attached, None) .await } - pub async fn add_tunnel_as_server( + pub async fn add_tunnel_as_server(&self, tunnel: Box) -> Result<(), Error> { + self.add_tunnel_as_server_with_origin(tunnel, PeerConnectionOrigin::Listener) + .await + .map(|_| ()) + } + + pub(crate) async fn add_tunnel_as_server_with_origin( &self, tunnel: Box, - is_directly_connected: bool, - ) -> Result<(), Error> { + origin: PeerConnectionOrigin, + ) -> Result<(PeerId, PeerConnId), Error> { self.peer_connection_admission - .add_tunnel_as_server(tunnel, is_directly_connected) + .add_tunnel_as_server_with_origin(tunnel, origin) .await } @@ -1606,8 +1609,7 @@ impl PeerManagerCore { &self, tunnel: Box, ) -> Result<(PeerId, PeerConnId), Error> { - self.peer_connection_admission - .add_tunnel_as_server_with_origin(tunnel, true, PeerConnectionOrigin::Attached) + self.add_tunnel_as_server_with_origin(tunnel, PeerConnectionOrigin::Attached) .await } @@ -1958,36 +1960,11 @@ impl PeerConnectionAdmission { } } - pub async fn add_client_tunnel( - &self, - tunnel: Box, - is_directly_connected: bool, - ) -> Result<(PeerId, PeerConnId), Error> { - self.add_client_tunnel_with_peer_id_hint(tunnel, is_directly_connected, None) - .await - } - pub async fn add_client_tunnel_with_peer_id_hint( &self, tunnel: Box, - is_directly_connected: bool, - peer_id_hint: Option, - ) -> Result<(PeerId, PeerConnId), Error> { - self.add_client_tunnel_with_origin( - tunnel, - is_directly_connected, - peer_id_hint, - PeerConnectionOrigin::Network, - ) - .await - } - - async fn add_client_tunnel_with_origin( - &self, - tunnel: Box, - is_directly_connected: bool, - peer_id_hint: Option, origin: PeerConnectionOrigin, + peer_id_hint: Option, ) -> Result<(PeerId, PeerConnId), Error> { let mut peer = PeerConn::new_with_peer_id_hint_and_origin( self.my_peer_id, @@ -1997,7 +1974,6 @@ impl PeerConnectionAdmission { self.peer_session_store.clone(), origin, ); - peer.set_is_hole_punched(!is_directly_connected); peer.do_handshake_as_client().await?; let conn_id = peer.get_conn_id(); let peer_id = peer.get_peer_id(); @@ -2036,24 +2012,9 @@ impl PeerConnectionAdmission { } #[tracing::instrument(ret, skip(self, tunnel))] - pub async fn add_tunnel_as_server( - &self, - tunnel: Box, - is_directly_connected: bool, - ) -> Result<(), Error> { - self.add_tunnel_as_server_with_origin( - tunnel, - is_directly_connected, - PeerConnectionOrigin::Network, - ) - .await - .map(|_| ()) - } - async fn add_tunnel_as_server_with_origin( &self, tunnel: Box, - is_directly_connected: bool, origin: PeerConnectionOrigin, ) -> Result<(PeerId, PeerConnId), Error> { tracing::info!("add tunnel as server start"); @@ -2138,8 +2099,6 @@ impl PeerConnectionAdmission { )); } - conn.set_is_hole_punched(!is_directly_connected); - let add_peer_ret = if is_local_network { let local_secure_mode = self .context @@ -4424,19 +4383,27 @@ mod tests { #[test] fn forged_attached_source_header_does_not_bypass_relay_disable() { let packet = data_packet(77, 3); - let network_ingress = PeerPacketIngress::Peer { - peer_id: 2, - conn_id: PeerConnId::new_v4(), - origin: PeerConnectionOrigin::Network, - }; + for origin in [ + PeerConnectionOrigin::Manual, + PeerConnectionOrigin::Direct, + PeerConnectionOrigin::Listener, + PeerConnectionOrigin::TcpHolePunch, + PeerConnectionOrigin::UdpHolePunch, + ] { + let network_ingress = PeerPacketIngress::Peer { + peer_id: 2, + conn_id: PeerConnId::new_v4(), + origin, + }; - assert!(should_drop_relay_data( - true, - &packet, - 1, - network_ingress, - false, - )); + assert!(should_drop_relay_data( + true, + &packet, + 1, + network_ingress, + false, + )); + } } #[test] @@ -4450,7 +4417,7 @@ mod tests { let network_ingress = PeerPacketIngress::Peer { peer_id: 2, conn_id: PeerConnId::new_v4(), - origin: PeerConnectionOrigin::Network, + origin: PeerConnectionOrigin::Listener, }; assert!(!should_drop_relay_data( diff --git a/easytier-core/src/peers/tests.rs b/easytier-core/src/peers/tests.rs index bc928a98..ad1b55c0 100644 --- a/easytier-core/src/peers/tests.rs +++ b/easytier-core/src/peers/tests.rs @@ -285,7 +285,6 @@ async fn peer_channel_uses_admission_origin_instead_of_packet_header() { ); client_ret.unwrap(); server_ret.unwrap(); - server_conn.set_is_hole_punched(false); let server_conn_id = server_conn.get_conn_id(); let (client_tx, _client_rx) = create_packet_recv_chan(); diff --git a/easytier-web/src/central_network/gateway.rs b/easytier-web/src/central_network/gateway.rs index 404e51fa..f6467bf0 100644 --- a/easytier-web/src/central_network/gateway.rs +++ b/easytier-web/src/central_network/gateway.rs @@ -379,7 +379,7 @@ impl NetworkInstanceManager { tokio::select! { biased; _ = retiring.cancelled() => None, - result = peer_manager.add_tunnel_as_server(tunnel, true) => Some(result), + result = peer_manager.add_tunnel_as_server(tunnel) => Some(result), } } diff --git a/easytier/docs/credential_peer.md b/easytier/docs/credential_peer.md index bf9837ad..b75bca9a 100644 --- a/easytier/docs/credential_peer.md +++ b/easytier/docs/credential_peer.md @@ -582,7 +582,7 @@ P2P hole punch 的流程: 这个流程不受影响,因为: - 打洞信息交换通过管理节点中继(RPC),不经过临时节点 - P2P tunnel 建立后的握手是直连,不通过临时节点的 listener -- `is_directly_connected=false` 的连接(hole punch 结果)可以被临时节点接受 +- 来源为 `TcpHolePunch` 或 `UdpHolePunch` 的连接可以被临时节点接受 **设计思路**: 将凭据映射为 ACL Group,复用现有的 group-based ACL 规则系统。 diff --git a/easytier/src/tests/three_node.rs b/easytier/src/tests/three_node.rs index dd9fd9b4..1ea68ee9 100644 --- a/easytier/src/tests/three_node.rs +++ b/easytier/src/tests/three_node.rs @@ -3048,7 +3048,7 @@ async fn assert_peer_admission_blocked(inst: &Instance, url: url::Url) { url, ) .await?; - core.admit_client_tunnel_for_test(tunnel, true) + core.admit_client_tunnel_for_test(tunnel) .await .map(|_| ()) .map_err(anyhow::Error::from)