mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-04 01:55:41 +00:00
feat: support disabling relay data forwarding (#2188)
- add a disable_relay_data runtime/config patch option - reuse the existing avoid_relay_data feature flag when relay data forwarding is disabled
This commit is contained in:
@@ -56,7 +56,7 @@ use super::{
|
||||
route_trait::NextHopPolicy,
|
||||
traffic_metrics::{
|
||||
InstanceLabelKind, LogicalTrafficMetrics, TrafficKind, TrafficMetricRecorder,
|
||||
route_peer_info_instance_id, traffic_kind,
|
||||
is_relay_data_packet_type, route_peer_info_instance_id, traffic_kind,
|
||||
},
|
||||
};
|
||||
|
||||
@@ -69,11 +69,16 @@ pub trait GlobalForeignNetworkAccessor: Send + Sync + 'static {
|
||||
struct ForeignNetworkEntry {
|
||||
my_peer_id: PeerId,
|
||||
|
||||
// Node-global runtime flags, such as disable_relay_data, live on the parent
|
||||
// context. The foreign context is scoped to the foreign network's OSPF view.
|
||||
parent_global_ctx: ArcGlobalCtx,
|
||||
global_ctx: ArcGlobalCtx,
|
||||
network: NetworkIdentity,
|
||||
peer_map: Arc<PeerMap>,
|
||||
relay_peer_map: Arc<RelayPeerMap>,
|
||||
peer_session_store: Arc<PeerSessionStore>,
|
||||
// Static per-network permission from the whitelist check. disable_relay_data
|
||||
// is the node-wide runtime override layered on top of this value.
|
||||
relay_data: bool,
|
||||
pm_packet_sender: Mutex<Option<PacketRecvChan>>,
|
||||
|
||||
@@ -205,6 +210,7 @@ impl ForeignNetworkEntry {
|
||||
Self {
|
||||
my_peer_id,
|
||||
|
||||
parent_global_ctx: global_ctx.clone(),
|
||||
global_ctx: foreign_global_ctx,
|
||||
network,
|
||||
peer_map,
|
||||
@@ -231,6 +237,27 @@ impl ForeignNetworkEntry {
|
||||
}
|
||||
}
|
||||
|
||||
fn desired_avoid_relay_data_feature_flag(
|
||||
parent_global_ctx: &ArcGlobalCtx,
|
||||
relay_data: bool,
|
||||
) -> bool {
|
||||
!relay_data || parent_global_ctx.get_feature_flags().avoid_relay_data
|
||||
}
|
||||
|
||||
fn sync_parent_relay_data_feature_flag(
|
||||
parent_global_ctx: &ArcGlobalCtx,
|
||||
global_ctx: &ArcGlobalCtx,
|
||||
relay_data: bool,
|
||||
) -> bool {
|
||||
let avoid_relay_data =
|
||||
Self::desired_avoid_relay_data_feature_flag(parent_global_ctx, relay_data);
|
||||
if global_ctx.get_feature_flags().avoid_relay_data == avoid_relay_data {
|
||||
return false;
|
||||
}
|
||||
|
||||
global_ctx.set_avoid_relay_data_preference(avoid_relay_data)
|
||||
}
|
||||
|
||||
fn build_foreign_global_ctx(
|
||||
network: &NetworkIdentity,
|
||||
global_ctx: ArcGlobalCtx,
|
||||
@@ -258,10 +285,9 @@ impl ForeignNetworkEntry {
|
||||
|
||||
let mut feature_flag = global_ctx.get_feature_flags();
|
||||
feature_flag.is_public_server = true;
|
||||
if !relay_data {
|
||||
feature_flag.avoid_relay_data = true;
|
||||
}
|
||||
foreign_global_ctx.set_feature_flags(feature_flag);
|
||||
feature_flag.avoid_relay_data =
|
||||
Self::desired_avoid_relay_data_feature_flag(&global_ctx, relay_data);
|
||||
foreign_global_ctx.set_base_advertised_feature_flags(feature_flag);
|
||||
|
||||
for u in global_ctx.get_running_listeners().into_iter() {
|
||||
foreign_global_ctx.add_running_listener(u);
|
||||
@@ -412,6 +438,7 @@ impl ForeignNetworkEntry {
|
||||
let peer_map = self.peer_map.clone();
|
||||
let relay_peer_map = self.relay_peer_map.clone();
|
||||
let traffic_metrics = self.traffic_metrics.clone();
|
||||
let parent_global_ctx = self.parent_global_ctx.clone();
|
||||
let relay_data = self.relay_data;
|
||||
let pm_sender = self.pm_packet_sender.lock().await.take().unwrap();
|
||||
let network_name = self.network.network_name.clone();
|
||||
@@ -497,11 +524,16 @@ impl ForeignNetworkEntry {
|
||||
"ignore packet in foreign network"
|
||||
);
|
||||
} else {
|
||||
if packet_type == PacketType::Data as u8
|
||||
|| packet_type == PacketType::KcpSrc as u8
|
||||
|| packet_type == PacketType::KcpDst as u8
|
||||
{
|
||||
if !relay_data {
|
||||
if is_relay_data_packet_type(packet_type) {
|
||||
let disable_relay_data = parent_global_ctx.flags_arc().disable_relay_data;
|
||||
if !relay_data || disable_relay_data {
|
||||
tracing::debug!(
|
||||
?from_peer_id,
|
||||
?to_peer_id,
|
||||
packet_type,
|
||||
disable_relay_data,
|
||||
"drop foreign network relay data"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
if !bps_limiter.try_consume(len.into()) {
|
||||
@@ -589,10 +621,31 @@ impl ForeignNetworkEntry {
|
||||
});
|
||||
}
|
||||
|
||||
async fn run_parent_feature_flag_sync_routine(&self) {
|
||||
let parent_global_ctx = self.parent_global_ctx.clone();
|
||||
let global_ctx = self.global_ctx.clone();
|
||||
let relay_data = self.relay_data;
|
||||
self.tasks.lock().await.spawn(async move {
|
||||
let mut parent_events = parent_global_ctx.subscribe();
|
||||
loop {
|
||||
ForeignNetworkEntry::sync_parent_relay_data_feature_flag(
|
||||
&parent_global_ctx,
|
||||
&global_ctx,
|
||||
relay_data,
|
||||
);
|
||||
|
||||
if parent_events.recv().await.is_err() {
|
||||
parent_events = parent_global_ctx.subscribe();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
async fn prepare(&self, accessor: Box<dyn GlobalForeignNetworkAccessor>) {
|
||||
self.prepare_route(accessor).await;
|
||||
self.start_packet_recv().await;
|
||||
self.run_relay_session_gc_routine().await;
|
||||
self.run_parent_feature_flag_sync_routine().await;
|
||||
self.peer_rpc.run();
|
||||
self.peer_center.init().await;
|
||||
}
|
||||
@@ -1397,6 +1450,92 @@ pub mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn disable_relay_data_blocks_foreign_network_transit_data() {
|
||||
let pm_center = create_mock_peer_manager_with_mock_stun(NatType::Unknown).await;
|
||||
let pma_net1 = create_mock_peer_manager_for_foreign_network("net1").await;
|
||||
let pmb_net1 = create_mock_peer_manager_for_foreign_network("net1").await;
|
||||
|
||||
connect_peer_manager(pma_net1.clone(), pm_center.clone()).await;
|
||||
connect_peer_manager(pmb_net1.clone(), pm_center.clone()).await;
|
||||
wait_route_appear(pma_net1.clone(), pmb_net1.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let mut flags = pm_center.get_global_ctx().get_flags();
|
||||
flags.disable_relay_data = true;
|
||||
pm_center.get_global_ctx().set_flags(flags);
|
||||
pm_center
|
||||
.get_global_ctx()
|
||||
.issue_event(GlobalCtxEvent::ConfigPatched(Default::default()));
|
||||
|
||||
let center_peer_id = pm_center
|
||||
.get_foreign_network_manager()
|
||||
.get_network_peer_id("net1")
|
||||
.unwrap();
|
||||
wait_for_condition(
|
||||
|| {
|
||||
let pma_net1 = pma_net1.clone();
|
||||
async move {
|
||||
pma_net1.list_routes().await.iter().any(|route| {
|
||||
route.peer_id == center_peer_id
|
||||
&& route
|
||||
.feature_flag
|
||||
.as_ref()
|
||||
.map(|flag| flag.avoid_relay_data)
|
||||
.unwrap_or(false)
|
||||
})
|
||||
}
|
||||
},
|
||||
Duration::from_secs(5),
|
||||
)
|
||||
.await;
|
||||
|
||||
let network_labels =
|
||||
LabelSet::new().with_label_type(LabelType::NetworkName("net1".to_string()));
|
||||
let forwarded_bytes_before = metric_value(
|
||||
&pm_center,
|
||||
MetricName::TrafficBytesForwarded,
|
||||
network_labels.clone(),
|
||||
);
|
||||
let forwarded_packets_before = metric_value(
|
||||
&pm_center,
|
||||
MetricName::TrafficPacketsForwarded,
|
||||
network_labels.clone(),
|
||||
);
|
||||
|
||||
let mut transit_pkt = ZCPacket::new_with_payload(b"foreign-transit-disabled");
|
||||
transit_pkt.fill_peer_manager_hdr(
|
||||
pma_net1.my_peer_id(),
|
||||
pmb_net1.my_peer_id(),
|
||||
PacketType::Data as u8,
|
||||
);
|
||||
pma_net1
|
||||
.get_foreign_network_client()
|
||||
.send_msg(transit_pkt, center_peer_id)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
|
||||
assert_eq!(
|
||||
metric_value(
|
||||
&pm_center,
|
||||
MetricName::TrafficBytesForwarded,
|
||||
network_labels.clone()
|
||||
),
|
||||
forwarded_bytes_before
|
||||
);
|
||||
assert_eq!(
|
||||
metric_value(
|
||||
&pm_center,
|
||||
MetricName::TrafficPacketsForwarded,
|
||||
network_labels
|
||||
),
|
||||
forwarded_packets_before
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn foreign_network_transit_control_forwarding_records_control_forwarded_metrics() {
|
||||
let pm_center = create_mock_peer_manager_with_mock_stun(NatType::Unknown).await;
|
||||
@@ -1409,6 +1548,10 @@ pub mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let mut flags = pm_center.get_global_ctx().get_flags();
|
||||
flags.disable_relay_data = true;
|
||||
pm_center.get_global_ctx().set_flags(flags);
|
||||
|
||||
let center_peer_id = pm_center
|
||||
.get_foreign_network_manager()
|
||||
.get_network_peer_id("net1")
|
||||
@@ -1657,6 +1800,81 @@ pub mod tests {
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn foreign_entry_feature_flag_tracks_parent_disable_relay_data_toggle() {
|
||||
let global_ctx = get_mock_global_ctx_with_network(Some(NetworkIdentity::new(
|
||||
"__access__".to_string(),
|
||||
"access_secret".to_string(),
|
||||
)));
|
||||
let foreign_network = NetworkIdentity::new("net1".to_string(), "net1_secret".to_string());
|
||||
let (pm_packet_sender, _pm_packet_recv) = create_packet_recv_chan();
|
||||
let entry = ForeignNetworkEntry::new(
|
||||
foreign_network,
|
||||
1,
|
||||
global_ctx.clone(),
|
||||
true,
|
||||
Arc::new(PeerSessionStore::new()),
|
||||
pm_packet_sender,
|
||||
);
|
||||
assert!(!entry.global_ctx.get_feature_flags().avoid_relay_data);
|
||||
|
||||
entry.run_parent_feature_flag_sync_routine().await;
|
||||
|
||||
let mut flags = global_ctx.get_flags();
|
||||
flags.disable_relay_data = true;
|
||||
global_ctx.set_flags(flags);
|
||||
global_ctx.issue_event(GlobalCtxEvent::ConfigPatched(Default::default()));
|
||||
|
||||
wait_for_condition(
|
||||
|| async { entry.global_ctx.get_feature_flags().avoid_relay_data },
|
||||
Duration::from_secs(2),
|
||||
)
|
||||
.await;
|
||||
|
||||
let mut flags = global_ctx.get_flags();
|
||||
flags.disable_relay_data = false;
|
||||
global_ctx.set_flags(flags);
|
||||
global_ctx.issue_event(GlobalCtxEvent::ConfigPatched(Default::default()));
|
||||
|
||||
wait_for_condition(
|
||||
|| async { !entry.global_ctx.get_feature_flags().avoid_relay_data },
|
||||
Duration::from_secs(2),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn foreign_entry_without_relay_data_keeps_avoid_feature_flag() {
|
||||
let global_ctx = get_mock_global_ctx_with_network(Some(NetworkIdentity::new(
|
||||
"__access__".to_string(),
|
||||
"access_secret".to_string(),
|
||||
)));
|
||||
let foreign_network = NetworkIdentity::new("net1".to_string(), "net1_secret".to_string());
|
||||
let (pm_packet_sender, _pm_packet_recv) = create_packet_recv_chan();
|
||||
let entry = ForeignNetworkEntry::new(
|
||||
foreign_network,
|
||||
1,
|
||||
global_ctx.clone(),
|
||||
false,
|
||||
Arc::new(PeerSessionStore::new()),
|
||||
pm_packet_sender,
|
||||
);
|
||||
|
||||
assert!(entry.global_ctx.get_feature_flags().avoid_relay_data);
|
||||
|
||||
let mut flags = global_ctx.get_flags();
|
||||
flags.disable_relay_data = false;
|
||||
global_ctx.set_flags(flags);
|
||||
|
||||
ForeignNetworkEntry::sync_parent_relay_data_feature_flag(
|
||||
&global_ctx,
|
||||
&entry.global_ctx,
|
||||
entry.relay_data,
|
||||
);
|
||||
|
||||
assert!(entry.global_ctx.get_feature_flags().avoid_relay_data);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn credential_trust_path_rejects_admin_identity() {
|
||||
assert!(ForeignNetworkManager::should_reject_credential_trust_path(
|
||||
|
||||
Reference in New Issue
Block a user