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)