From 82c66acc233e0d5047c083b64f07c9155e5aac6b Mon Sep 17 00:00:00 2001 From: fanyang Date: Thu, 18 Jun 2026 09:27:13 +0800 Subject: [PATCH] fix: bound peer rpc packet queues --- easytier/src/common/stats_manager.rs | 3 + easytier/src/peers/foreign_network_manager.rs | 108 +++++++++++++++-- easytier/src/peers/peer_manager.rs | 114 ++++++++++++++++-- 3 files changed, 209 insertions(+), 16 deletions(-) diff --git a/easytier/src/common/stats_manager.rs b/easytier/src/common/stats_manager.rs index 0db8edbe..36952def 100644 --- a/easytier/src/common/stats_manager.rs +++ b/easytier/src/common/stats_manager.rs @@ -23,6 +23,8 @@ pub enum MetricName { PeerRpcDuration, /// RPC errors PeerRpcErrors, + /// RPC/control packets dropped because the peer RPC queue is unavailable + PeerRpcPacketQueueDrops, /// Data-plane traffic bytes sent TrafficBytesTx, @@ -116,6 +118,7 @@ impl fmt::Display for MetricName { MetricName::PeerRpcServerRx => write!(f, "peer_rpc_server_rx"), MetricName::PeerRpcDuration => write!(f, "peer_rpc_duration_ms"), MetricName::PeerRpcErrors => write!(f, "peer_rpc_errors"), + MetricName::PeerRpcPacketQueueDrops => write!(f, "peer_rpc_packet_queue_drops"), MetricName::TrafficBytesTx => write!(f, "traffic_bytes_tx"), MetricName::TrafficBytesTxByInstance => write!(f, "traffic_bytes_tx_by_instance"), diff --git a/easytier/src/peers/foreign_network_manager.rs b/easytier/src/peers/foreign_network_manager.rs index 9a9fb641..8d5dd9a9 100644 --- a/easytier/src/peers/foreign_network_manager.rs +++ b/easytier/src/peers/foreign_network_manager.rs @@ -18,7 +18,7 @@ use guarden::{Guard, defer}; use tokio::{ sync::{ Mutex, - mpsc::{self, UnboundedReceiver, UnboundedSender}, + mpsc::{self, Receiver, Sender, error::TrySendError}, }, task::JoinSet, }; @@ -30,7 +30,7 @@ use crate::{ error::Error, global_ctx::{ArcGlobalCtx, GlobalCtx, GlobalCtxEvent, NetworkIdentity, TrustedKeySource}, join_joinset_background, shrink_dashmap, - stats_manager::{LabelSet, LabelType, MetricName, StatsManager}, + stats_manager::{CounterHandle, LabelSet, LabelType, MetricName, StatsManager}, token_bucket::TokenBucket, }, peer_center::instance::{PeerCenterInstance, PeerMapWithPeerRpcManager}, @@ -64,6 +64,35 @@ use super::{ }, }; +const PEER_RPC_PACKET_QUEUE_CAPACITY: usize = 1024; + +fn try_enqueue_peer_rpc_packet( + sender: &Sender, + packet: ZCPacket, + dropped_packets: &CounterHandle, + queue_name: &'static str, +) -> bool { + match sender.try_send(packet) { + Ok(()) => true, + Err(TrySendError::Full(_)) => { + dropped_packets.inc(); + tracing::warn!( + queue = queue_name, + "drop peer rpc/control packet because queue is full" + ); + false + } + Err(TrySendError::Closed(_)) => { + dropped_packets.inc(); + tracing::warn!( + queue = queue_name, + "drop peer rpc/control packet because receiver is closed" + ); + false + } + } +} + #[async_trait::async_trait] #[auto_impl::auto_impl(&, Box, Arc)] pub trait GlobalForeignNetworkAccessor: Send + Sync + 'static { @@ -87,7 +116,7 @@ struct ForeignNetworkEntry { pm_packet_sender: Mutex>, peer_rpc: Arc, - rpc_sender: UnboundedSender, + rpc_sender: Sender, packet_recv: Mutex>, @@ -312,12 +341,12 @@ impl ForeignNetworkEntry { fn build_rpc_tspt( my_peer_id: PeerId, peer_map: Arc, - ) -> (Arc, UnboundedSender) { + ) -> (Arc, Sender) { struct RpcTransport { my_peer_id: PeerId, peer_map: Weak, - packet_recv: Mutex>, + packet_recv: Mutex>, } #[async_trait::async_trait] @@ -359,7 +388,8 @@ impl ForeignNetworkEntry { } } - let (rpc_transport_sender, peer_rpc_tspt_recv) = mpsc::unbounded_channel(); + let (rpc_transport_sender, peer_rpc_tspt_recv) = + mpsc::channel(PEER_RPC_PACKET_QUEUE_CAPACITY); let tspt = RpcTransport { my_peer_id, peer_map: Arc::downgrade(&peer_map), @@ -478,6 +508,9 @@ impl ForeignNetworkEntry { let rx_packets = self .stats_mgr .get_counter(MetricName::TrafficPacketsRx, label_set.clone()); + let rpc_queue_drops = self + .stats_mgr + .get_counter(MetricName::PeerRpcPacketQueueDrops, label_set.clone()); self.tasks.lock().await.spawn(async move { while let Ok(mut zc_packet) = recv_packet_from_chan(&mut recv).await { @@ -526,7 +559,12 @@ impl ForeignNetworkEntry { { rx_bytes.add(buf_len as u64); rx_packets.inc(); - rpc_sender.send(zc_packet).unwrap(); + try_enqueue_peer_rpc_packet( + &rpc_sender, + zc_packet, + &rpc_queue_drops, + "foreign_network_peer_rpc", + ); continue; } tracing::trace!( @@ -1236,7 +1274,7 @@ impl Drop for ForeignNetworkManager { pub mod tests { use crate::{ common::global_ctx::tests::get_mock_global_ctx_with_network, - common::stats_manager::{LabelSet, LabelType, MetricName}, + common::stats_manager::{LabelSet, LabelType, MetricName, StatsManager}, connector::udp_hole_punch::tests::{ create_mock_peer_manager_with_mock_stun, replace_stun_info_collector, }, @@ -1253,6 +1291,7 @@ pub mod tests { }, }; use std::{collections::HashMap, time::Duration}; + use tokio::sync::mpsc; use super::*; @@ -1265,6 +1304,59 @@ pub mod tests { .unwrap_or(0) } + #[tokio::test] + async fn peer_rpc_queue_helper_enqueues_when_available() { + let (sender, mut receiver) = mpsc::channel(1); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"rpc"), + &dropped_packets, + "test_foreign_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 0); + assert!(receiver.try_recv().is_ok()); + } + + #[tokio::test] + async fn peer_rpc_queue_helper_drops_when_full() { + let (sender, _receiver) = mpsc::channel(1); + sender + .try_send(ZCPacket::new_with_payload(b"existing")) + .unwrap(); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(!try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"overflow"), + &dropped_packets, + "test_foreign_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 1); + } + + #[tokio::test] + async fn peer_rpc_queue_helper_drops_when_closed() { + let (sender, receiver) = mpsc::channel(1); + drop(receiver); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(!try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"closed"), + &dropped_packets, + "test_foreign_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 1); + } + async fn create_mock_peer_manager_for_foreign_network_ext( network: &str, secret: &str, diff --git a/easytier/src/peers/peer_manager.rs b/easytier/src/peers/peer_manager.rs index 17144d4c..a9de53c1 100644 --- a/easytier/src/peers/peer_manager.rs +++ b/easytier/src/peers/peer_manager.rs @@ -13,7 +13,7 @@ use std::{ use tokio::{ sync::{ Mutex, RwLock, - mpsc::{self, UnboundedReceiver, UnboundedSender}, + mpsc::{self, Receiver, Sender, error::TrySendError}, }, task::JoinSet, }; @@ -72,14 +72,43 @@ use super::{ route_trait::{ArcRoute, Route}, }; +const PEER_RPC_PACKET_QUEUE_CAPACITY: usize = 1024; + +fn try_enqueue_peer_rpc_packet( + sender: &Sender, + packet: ZCPacket, + dropped_packets: &CounterHandle, + queue_name: &'static str, +) -> bool { + match sender.try_send(packet) { + Ok(()) => true, + Err(TrySendError::Full(_)) => { + dropped_packets.inc(); + tracing::warn!( + queue = queue_name, + "drop peer rpc/control packet because queue is full" + ); + false + } + Err(TrySendError::Closed(_)) => { + dropped_packets.inc(); + tracing::warn!( + queue = queue_name, + "drop peer rpc/control packet because receiver is closed" + ); + false + } + } +} + struct RpcTransport { my_peer_id: PeerId, peers: Weak, // TODO: this seems can be removed foreign_peers: Mutex>>, - packet_recv: Mutex>, - peer_rpc_tspt_sender: UnboundedSender, + packet_recv: Mutex>, + peer_rpc_tspt_sender: Sender, encryptor: Arc, is_secure_mode_enabled: bool, @@ -273,7 +302,8 @@ impl PeerManager { .unwrap_or(false); // TODO: remove these because we have impl pipeline processor. - let (peer_rpc_tspt_sender, peer_rpc_tspt_recv) = mpsc::unbounded_channel(); + let (peer_rpc_tspt_sender, peer_rpc_tspt_recv) = + mpsc::channel(PEER_RPC_PACKET_QUEUE_CAPACITY); let rpc_tspt = Arc::new(RpcTransport { my_peer_id, peers: Arc::downgrade(&peers), @@ -1243,7 +1273,8 @@ impl PeerManager { // for peer rpc packet struct PeerRpcPacketProcessor { - peer_rpc_tspt_sender: UnboundedSender, + peer_rpc_tspt_sender: Sender, + dropped_packets: CounterHandle, } #[async_trait::async_trait] @@ -1254,15 +1285,27 @@ impl PeerManager { || hdr.packet_type == PacketType::RpcReq as u8 || hdr.packet_type == PacketType::RpcResp as u8 { - self.peer_rpc_tspt_sender.send(packet).unwrap(); + try_enqueue_peer_rpc_packet( + &self.peer_rpc_tspt_sender, + packet, + &self.dropped_packets, + "local_peer_rpc", + ); None } else { Some(packet) } } } + let peer_rpc_queue_drops = self.global_ctx.stats_manager().get_counter( + MetricName::PeerRpcPacketQueueDrops, + LabelSet::new().with_label_type(LabelType::NetworkName( + self.global_ctx.get_network_name().to_string(), + )), + ); self.add_packet_process_pipeline(Box::new(PeerRpcPacketProcessor { peer_rpc_tspt_sender: self.peer_rpc_tspt.peer_rpc_tspt_sender.clone(), + dropped_packets: peer_rpc_queue_drops, })) .await; } @@ -2209,7 +2252,7 @@ mod tests { PeerId, config::Flags, global_ctx::{NetworkIdentity, tests::get_mock_global_ctx}, - stats_manager::{LabelSet, LabelType, MetricName}, + stats_manager::{LabelSet, LabelType, MetricName, StatsManager}, }, connector::{ create_connector_by_url, direct::PeerManagerForDirectConnector, @@ -2240,7 +2283,9 @@ mod tests { }, }; - use super::PeerManager; + use tokio::sync::mpsc; + + use super::{PeerManager, try_enqueue_peer_rpc_packet}; async fn create_lazy_peer_manager() -> Arc { let peer_mgr = create_mock_peer_manager_with_mock_stun(NatType::Unknown).await; @@ -2265,6 +2310,59 @@ mod tests { )) } + #[tokio::test] + async fn peer_rpc_queue_helper_enqueues_when_available() { + let (sender, mut receiver) = mpsc::channel(1); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"rpc"), + &dropped_packets, + "test_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 0); + assert!(receiver.try_recv().is_ok()); + } + + #[tokio::test] + async fn peer_rpc_queue_helper_drops_when_full() { + let (sender, _receiver) = mpsc::channel(1); + sender + .try_send(ZCPacket::new_with_payload(b"existing")) + .unwrap(); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(!try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"overflow"), + &dropped_packets, + "test_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 1); + } + + #[tokio::test] + async fn peer_rpc_queue_helper_drops_when_closed() { + let (sender, receiver) = mpsc::channel(1); + drop(receiver); + let stats_manager = StatsManager::new(); + let dropped_packets = stats_manager.get_simple_counter(MetricName::PeerRpcPacketQueueDrops); + + assert!(!try_enqueue_peer_rpc_packet( + &sender, + ZCPacket::new_with_payload(b"closed"), + &dropped_packets, + "test_peer_rpc", + )); + + assert_eq!(dropped_packets.get(), 1); + } + struct TestCostCalculator { costs: HashMap<(PeerId, PeerId), i32>, }