From d1e36c7cefde1cd85c997933c42e9fc926e2fb51 Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Thu, 3 Sep 2026 23:41:37 +0800 Subject: [PATCH 01/14] feat: support trusted IKEv2 relay clients --- README.md | 1 + src/args_config.rs | 32 ++++++++++ vnt-core/proto/control_message.proto | 7 +++ vnt-core/proto/rpc.proto | 8 ++- vnt-core/src/context/config.rs | 6 ++ vnt-core/src/context/mod.rs | 14 ++++- vnt-core/src/core/mod.rs | 6 +- vnt-core/src/enhanced_tunnel/outbound.rs | 27 +++++++- vnt-core/src/protocol/control_message.rs | 45 ++++++++++++++ vnt-core/src/protocol/ip_packet_protocol.rs | 5 +- vnt-core/src/tunnel_core/outbound.rs | 62 +++++++++++++++++++ .../tunnel_core/server/connection_manager.rs | 1 + vnt-core/src/tunnel_core/server/inbound.rs | 20 ++++++ .../tunnel_core/server/transport/config.rs | 3 + vnt-jni/java_example/com/vnt/VntConfig.java | 12 ++++ vnt-jni/src/lib.rs | 20 +++++- vnt-web/src/service_http.rs | 18 +++++- vnt-web/ui/src/utils/configHelp.js | 7 +++ vnt-web/ui/src/utils/toml.js | 11 ++++ vnt-web/ui/src/views/ConfigEditor.vue | 9 +++ vnt-web/ui/src/views/PeersView.vue | 5 +- 21 files changed, 306 insertions(+), 13 deletions(-) diff --git a/README.md b/README.md index a24c2c37..4dc8c091 100644 --- a/README.md +++ b/README.md @@ -72,6 +72,7 @@ IPv4 广播和组播默认开启。可使用 `--no-broadcast`(配置文件中 ## 安全说明 - **建议设置组网密码**。设置密码后(命令行 `-p` / `--password`,或配置文件中的 `password`),节点之间的数据采用端到端加密(ChaCha20-Poly1305),**服务端仅负责转发密文,无法解密通信内容**。即使使用公共服务端,通信内容也不会泄露给服务端。 +- 如需与服务端接入的 IKEv2/IPsec 客户端通信,使用 `--allow-ikev2`(配置文件中为 `allow_ikev2 = true`)。该功能会信任已认证服务端注入的 IKEv2 明文 IPv4 数据,并让发往 IKEv2 类型设备的流量固定走服务端;默认关闭,且该路径不受 VNT 节点间密码的端到端加密保护。 - 同一虚拟网络内的所有设备必须使用**相同的密码**,否则无法互相通信。 - 未设置密码时,节点间数据不加密,经过服务端中继的流量理论上可被服务端查看,请仅在可信网络环境下省略密码。 - 此外,客户端与服务端之间的连接本身支持 tcp-tls、quic、wss 等加密传输协议,并可绑定服务端证书,防止伪造服务端攻击。 diff --git a/src/args_config.rs b/src/args_config.rs index b3f6d867..fef92fa9 100644 --- a/src/args_config.rs +++ b/src/args_config.rs @@ -20,6 +20,7 @@ pub struct FileConfig { pub ip: Option, pub no_punch: Option, pub no_broadcast: Option, + pub allow_ikev2: Option, pub rtx: Option, pub compress: Option, pub fec: Option, @@ -169,6 +170,9 @@ pub struct Args { /// 关闭虚拟网络内的 IPv4 广播和组播转发 #[clap(long)] pub no_broadcast: bool, + /// 允许与 IKEv2/IPsec 客户端通信,并信任服务端注入的 IKEv2 明文 IPv4 包 + #[clap(long)] + pub allow_ikev2: bool, /// 服务端证书验证 #[clap(long)] pub cert_mode: Option, @@ -317,6 +321,7 @@ fn build_from_args_and_file(args: Args, file: FileConfig) -> anyhow::Result<(Con ip: args.ip.or(file.ip), no_punch: args.no_punch || file.no_punch.unwrap_or(false), no_broadcast: args.no_broadcast || file.no_broadcast.unwrap_or(false), + allow_ikev2: args.allow_ikev2 || file.allow_ikev2.unwrap_or(false), rtx: args.rtx || file.rtx.unwrap_or(false), compress: args.compress || file.compress.unwrap_or(false), fec: args.fec || file.fec.unwrap_or(false), @@ -367,6 +372,7 @@ fn build_from_args_only(args: Args) -> anyhow::Result<(Config, CtrlConfig)> { ip: args.ip, no_punch: args.no_punch, no_broadcast: args.no_broadcast, + allow_ikev2: args.allow_ikev2, rtx: args.rtx, input: args.input, subnet_mapping: args.subnet_mapping, @@ -441,6 +447,7 @@ fn build_from_file_only(file: FileConfig) -> anyhow::Result<(Config, CtrlConfig) ip: file.ip, no_punch: file.no_punch.unwrap_or(false), no_broadcast: file.no_broadcast.unwrap_or(false), + allow_ikev2: file.allow_ikev2.unwrap_or(false), rtx: file.rtx.unwrap_or(false), input: file.input.unwrap_or_default(), subnet_mapping: file.subnet_mapping.unwrap_or_default(), @@ -520,6 +527,9 @@ server = ["quic://1.2.3.4:29872"] # 是否关闭 IPv4 广播和组播转发 (默认 false,即开启) # no_broadcast = false +# 是否允许与 IKEv2 客户端通信,并信任服务端注入的 IKEv2 明文 IPv4 包 +# allow_ikev2 = false + # 是否启用 LZ4 压缩 (默认 false,设置为true时开启) # compress = false @@ -836,4 +846,26 @@ mod tests { let (config, _) = build_config_from_args_and_file(Some(args), Some(file)).unwrap(); assert_eq!(config.device_mode, DeviceMode::Tap); } + + #[test] + fn allow_ikev2_is_opt_in_for_cli_and_toml() { + let args = Args::try_parse_from([ + "vnt", + "-s", + "quic://127.0.0.1:29872", + "-n", + "test-net", + "--allow-ikev2", + ]) + .unwrap(); + let (config, _) = build_from_args_only(args).unwrap(); + assert!(config.allow_ikev2); + + let file: FileConfig = toml::from_str( + "server = [\"quic://127.0.0.1:29872\"]\nnetwork_code = \"test\"\nallow_ikev2 = true", + ) + .unwrap(); + let (config, _) = build_config_from_args_and_file(None, Some(file)).unwrap(); + assert!(config.allow_ikev2); + } } diff --git a/vnt-core/proto/control_message.proto b/vnt-core/proto/control_message.proto index 90be3f47..b0f75610 100644 --- a/vnt-core/proto/control_message.proto +++ b/vnt-core/proto/control_message.proto @@ -7,6 +7,11 @@ enum RegistrationMode { PRE_REGISTER = 1; } +enum ClientType { + VNT = 0; + IKEV2 = 1; +} + message RegRequestMsg { string network_code = 1; string device_id = 2; @@ -18,6 +23,7 @@ message RegRequestMsg { fixed32 server_id = 8; RegistrationMode registration_mode = 9; repeated Ipv4Subnet advertised_subnets = 10; + bool allow_ikev2 = 11; } @@ -92,6 +98,7 @@ message SelectiveBroadcast { message ClientSimpleInfo{ fixed32 ip = 1; bool online = 2; + ClientType client_type = 3; } message ClientSimpleInfoList{ diff --git a/vnt-core/proto/rpc.proto b/vnt-core/proto/rpc.proto index e6077e70..ebd30256 100644 --- a/vnt-core/proto/rpc.proto +++ b/vnt-core/proto/rpc.proto @@ -2,6 +2,11 @@ syntax = "proto3"; package protocol.rpc; +enum ClientType { + VNT = 0; + IKEV2 = 1; +} + message RpcMessageRequest{ uint64 id = 1; oneof rpc_req_payload{ @@ -30,8 +35,9 @@ message ClientInfo{ bool online = 5; int64 last_connected_time = 6; string id = 7; + ClientType client_type = 8; } message ClientListResponse{ repeated ClientInfo list = 1; -} \ No newline at end of file +} diff --git a/vnt-core/src/context/config.rs b/vnt-core/src/context/config.rs index dc3c173e..e6821205 100644 --- a/vnt-core/src/context/config.rs +++ b/vnt-core/src/context/config.rs @@ -242,6 +242,8 @@ pub struct Config { pub password: Option, pub no_punch: bool, pub no_broadcast: bool, + /// 允许与由服务端终结的 IKEv2/IPsec 客户端互通。 + pub allow_ikev2: bool, pub compress: bool, pub rtx: bool, pub fec: bool, @@ -294,6 +296,9 @@ impl Config { if self.server_addr.is_empty() { bail!("服务器地址不能为空"); } + if self.allow_ikev2 && self.device_mode == DeviceMode::No { + bail!("allow_ikev2 requires device_mode = tun or tap"); + } if self.server_addr.len() > 1 { let mut set = HashSet::new(); @@ -357,6 +362,7 @@ impl Config { &self.output, &self.subnet_mapping, )), + allow_ikev2: self.allow_ikev2, default_interface, } } diff --git a/vnt-core/src/context/mod.rs b/vnt-core/src/context/mod.rs index 5b13a284..1a52af7e 100644 --- a/vnt-core/src/context/mod.rs +++ b/vnt-core/src/context/mod.rs @@ -3,7 +3,7 @@ use crate::context::nat::{MyNatInfo, PunchBackoff}; use crate::nat::SubnetExternalRoute; use crate::protocol::client_message::PunchInfo; use crate::protocol::control_message::{ - ClientSimpleInfo, ClientSimpleInfoList, SubnetSyncResponse, + ClientSimpleInfo, ClientSimpleInfoList, ClientType, SubnetSyncResponse, }; use crate::tunnel_core::p2p::route_table::RouteTable; use crate::tunnel_core::server::transport::config::ProtocolAddress; @@ -423,7 +423,7 @@ impl ServerInfoCollection { ( v.client_map .iter() - .filter(|(_, v)| v.online) + .filter(|(_, v)| v.online && v.client_type == ClientType::Vnt) .map(|(k, _)| *k) .collect(), v.rtt.unwrap_or(500), @@ -496,13 +496,18 @@ impl ServerInfoCollection { self.client_simple_list .read() .iter() - .filter(|v| v.online) + .filter(|v| v.online && v.client_type == ClientType::Vnt) .map(|c| c.ip) .collect() } pub fn client_ips(&self) -> Vec { self.client_simple_list.read().clone() } + pub fn is_ikev2_client(&self, ip: &Ipv4Addr) -> bool { + self.client_simple_list.read().iter().any(|client| { + client.ip == *ip && client.online && client.client_type == ClientType::Ikev2 + }) + } pub fn data_version(&self, server_id: u32) -> u64 { self.server_node_map .read() @@ -604,6 +609,9 @@ impl ServerInfoCollection { if x.online { v.online = true; } + if x.client_type == ClientType::Ikev2 { + v.client_type = ClientType::Ikev2; + } } else { client_simple_map.insert(x.ip, x.clone()); } diff --git a/vnt-core/src/core/mod.rs b/vnt-core/src/core/mod.rs index 05eac990..7d4bff92 100644 --- a/vnt-core/src/core/mod.rs +++ b/vnt-core/src/core/mod.rs @@ -50,6 +50,7 @@ struct RegistrationContext { fec_decoder: FecDecoder, turn: std::sync::Arc>, auto_sync_subnet: bool, + allow_ikev2: bool, } pub struct NetworkManager { @@ -194,7 +195,8 @@ impl NetworkManager { subnet_packet_mapper.clone(), fec_encoder, ) - .with_no_broadcast(config.no_broadcast); + .with_no_broadcast(config.no_broadcast) + .with_allow_ikev2(config.allow_ikev2); let port_mapping_manager = PortMappingManager::new( config.device_mode == DeviceMode::No, config.allow_port_mapping, @@ -289,6 +291,7 @@ impl NetworkManager { fec_decoder, turn, auto_sync_subnet: config.auto_sync_subnet, + allow_ikev2: config.allow_ikev2, }); app_state.set_config(config.clone()); @@ -403,6 +406,7 @@ impl NetworkManager { fec_decoder: ctx.fec_decoder.clone(), turn: ctx.turn.clone(), auto_sync_subnet: ctx.auto_sync_subnet, + allow_ikev2: ctx.allow_ikev2, }); turn_manager.data_handle_task_connected(task_group, handler_config); } diff --git a/vnt-core/src/enhanced_tunnel/outbound.rs b/vnt-core/src/enhanced_tunnel/outbound.rs index 6c25e3c1..343e30b0 100644 --- a/vnt-core/src/enhanced_tunnel/outbound.rs +++ b/vnt-core/src/enhanced_tunnel/outbound.rs @@ -2,6 +2,7 @@ use crate::context::SharedNetworkAddr; use crate::enhanced_tunnel::quic_over::quic_outbound::EnhancedQuicOutbound; use crate::ethernet::{ MacTable, build_arp_reply, is_broadcast_or_multicast, mac_from_ip, parse_arp_ipv4, parse_frame, + strip_ipv4, }; use crate::nat::SubnetMappingTable; use crate::nat::subnet_packet::SubnetPacketMapper; @@ -86,9 +87,25 @@ impl EnhancedOutbound { if frame.ethertype == EtherTypes::Arp && let Some(arp) = parse_arp_ipv4(data.as_ref()) && arp.operation == ArpOperations::Request - && arp.target_ip == net.gateway + && (arp.target_ip == net.gateway + || self.hybrid_outbound.is_ikev2_client(&arp.target_ip)) { - return Ok(build_arp_reply(data.as_ref(), net.gateway)); + return Ok(build_arp_reply(data.as_ref(), arp.target_ip)); + } + + if frame.ethertype == EtherTypes::Ipv4 { + let Some(ipv4) = Ipv4Packet::new(&data[frame.payload_offset..]) else { + return Ok(None); + }; + let dest = ipv4.get_destination(); + if self.hybrid_outbound.is_ikev2_client(&dest) { + if let Some(ip) = strip_ipv4(data) { + self.hybrid_outbound + .ikev2_relay_outbound(net, ip, dest) + .await?; + } + return Ok(None); + } } // Only frames explicitly addressed to our proxy-ARP gateway enter the @@ -143,6 +160,12 @@ impl EnhancedOutbound { // 发送到网关 return self.hybrid_outbound.ipv4_gateway_outbound(net, data).await; } + if self.hybrid_outbound.is_ikev2_client(&dest) { + return self + .hybrid_outbound + .ikev2_relay_outbound(net, data, dest) + .await; + } if dest.is_multicast() || dest == net.broadcast || dest.is_broadcast() { // 广播或组播 if self.hybrid_outbound.no_broadcast() { diff --git a/vnt-core/src/protocol/control_message.rs b/vnt-core/src/protocol/control_message.rs index 7b000fd9..243a6400 100644 --- a/vnt-core/src/protocol/control_message.rs +++ b/vnt-core/src/protocol/control_message.rs @@ -12,6 +12,8 @@ mod proto { include!(concat!(env!("OUT_DIR"), "/protocol.control_message.rs")); } +pub use proto::ClientType; + #[derive(Debug, Clone, Copy, Eq, PartialEq, Default)] pub enum RegistrationMode { #[default] @@ -47,6 +49,7 @@ pub(crate) struct RegRequestMsg { pub server_id: u32, pub registration_mode: RegistrationMode, pub advertised_subnets: Vec, + pub allow_ikev2: bool, } impl RegRequestMsg { // pub fn check(&self) -> anyhow::Result<()> { @@ -117,6 +120,7 @@ impl RegRequestMsg { .into_iter() .map(ipv4_subnet_to_proto) .collect(), + allow_ikev2: self.allow_ikev2, } } } @@ -333,18 +337,21 @@ impl SelectiveBroadcast { pub struct ClientSimpleInfo { pub ip: Ipv4Addr, pub online: bool, + pub client_type: ClientType, } impl ClientSimpleInfo { pub fn from(msg: proto::ClientSimpleInfo) -> anyhow::Result { Ok(Self { ip: msg.ip.into(), online: msg.online, + client_type: msg.client_type(), }) } pub fn to(self) -> proto::ClientSimpleInfo { proto::ClientSimpleInfo { ip: self.ip.into(), online: self.online, + client_type: self.client_type as i32, } } } @@ -406,6 +413,7 @@ mod tests { server_id: 0, registration_mode: RegistrationMode::Normal, advertised_subnets: vec![advertised], + allow_ikev2: false, }) .encode(); let request = proto::RequestMessage::decode(encoded.as_ref()).unwrap(); @@ -443,4 +451,41 @@ mod tests { assert_eq!(snapshot.snapshot_hash, vec![1, 2, 3]); assert_eq!(snapshot.nodes[0].subnets, vec![advertised]); } + + #[test] + fn ikev2_capability_and_client_type_round_trip() { + let encoded = RequestMessage::Reg(RegRequestMsg { + network_code: "test".to_string(), + device_id: "device".to_string(), + ip: None, + name: "node".to_string(), + version: "1".to_string(), + key_sign: None, + ip_variable: true, + server_id: 0, + registration_mode: RegistrationMode::Normal, + advertised_subnets: Vec::new(), + allow_ikev2: true, + }) + .encode(); + let request = proto::RequestMessage::decode(encoded.as_ref()).unwrap(); + let RequestPayload::Reg(request) = request.request_payload.unwrap() else { + panic!("expected registration request"); + }; + assert!(request.allow_ikev2); + + let list = proto::ClientSimpleInfoList { + data_version: 1, + list: vec![proto::ClientSimpleInfo { + ip: Ipv4Addr::new(10, 26, 0, 8).into(), + online: true, + client_type: proto::ClientType::Ikev2 as i32, + }], + is_all: true, + time: 0, + } + .encode_to_vec(); + let decoded = ClientSimpleInfoList::from_slice(&list).unwrap(); + assert_eq!(decoded.list[0].client_type, ClientType::Ikev2); + } } diff --git a/vnt-core/src/protocol/ip_packet_protocol.rs b/vnt-core/src/protocol/ip_packet_protocol.rs index 67433c5d..b6615339 100644 --- a/vnt-core/src/protocol/ip_packet_protocol.rs +++ b/vnt-core/src/protocol/ip_packet_protocol.rs @@ -117,6 +117,7 @@ pub enum MsgType { FastReg = 22, SubnetSyncReq = 23, SubnetSyncRes = 24, + Ikev2Relay = 25, } impl From for u8 { fn from(val: MsgType) -> Self { @@ -159,6 +160,7 @@ impl TryFrom for MsgType { 22 => MsgType::FastReg, 23 => MsgType::SubnetSyncReq, 24 => MsgType::SubnetSyncRes, + 25 => MsgType::Ikev2Relay, _ => { return Err(io::Error::new( io::ErrorKind::InvalidInput, @@ -363,6 +365,7 @@ mod tests { MsgType::FastReg, MsgType::SubnetSyncReq, MsgType::SubnetSyncRes, + MsgType::Ikev2Relay, ]; for msg_type in all { let byte = u8::from(msg_type); @@ -374,7 +377,7 @@ mod tests { } // 未分配的取值必须报错 assert!(MsgType::try_from(0u8).is_err()); - assert!(MsgType::try_from(25u8).is_err()); + assert!(MsgType::try_from(26u8).is_err()); } #[test] diff --git a/vnt-core/src/tunnel_core/outbound.rs b/vnt-core/src/tunnel_core/outbound.rs index b502e8b4..702a374b 100644 --- a/vnt-core/src/tunnel_core/outbound.rs +++ b/vnt-core/src/tunnel_core/outbound.rs @@ -115,6 +115,16 @@ impl BasicOutbound { .await } + pub async fn send_server_raw( + &self, + dest: Ipv4Addr, + packet: NetPacket, + ) -> anyhow::Result<()> { + self.server_outbound + .send_raw(dest, packet.into_bytes()) + .await + } + /// 广播发送 pub async fn send_raw_broadcast( &self, @@ -218,6 +228,7 @@ pub(crate) struct HybridOutbound { subnet_packet_mapper: SubnetPacketMapper, fec_encoder: Option, no_broadcast: bool, + allow_ikev2: bool, } impl HybridOutbound { #[allow(clippy::too_many_arguments)] @@ -243,6 +254,7 @@ impl HybridOutbound { subnet_packet_mapper, fec_encoder, no_broadcast: false, + allow_ikev2: false, } } @@ -250,6 +262,50 @@ impl HybridOutbound { self.no_broadcast = no_broadcast; self } + pub fn with_allow_ikev2(mut self, allow_ikev2: bool) -> Self { + self.allow_ikev2 = allow_ikev2; + self + } + pub fn is_ikev2_client(&self, ip: &Ipv4Addr) -> bool { + self.allow_ikev2 && self.server_info.is_ikev2_client(ip) + } + pub async fn ikev2_relay_outbound( + &self, + net: NetworkAddr, + mut data: TransmissionBytes, + dest: Ipv4Addr, + ) -> anyhow::Result<()> { + if !self.is_ikev2_client(&dest) { + return Ok(()); + } + let Some(ipv4) = Ipv4Packet::new(data.as_ref()) else { + return Ok(()); + }; + let header_length = ipv4.get_header_length() as usize * 4; + let total_length = ipv4.get_total_length() as usize; + if header_length < Ipv4Packet::minimum_packet_size() + || total_length < header_length + || total_length > data.len() + || ipv4.get_source() != net.ip + || ipv4.get_destination() != dest + { + return Ok(()); + } + if total_length < data.len() { + let trailing = data.len() - total_length; + data.shrink_end(trailing); + } + let len = data.len() as u64; + data.retreat_head(HEAD_LENGTH)?; + let mut packet = NetPacket::new(data)?; + packet.set_msg_type(MsgType::Ikev2Relay); + packet.set_src_id(net.ip.into()); + packet.set_dest_id(dest.into()); + packet.set_ttl(5); + self.basic_outbound.send_server_raw(dest, packet).await?; + self.traffic_stats.record_tx(dest, len); + Ok(()) + } pub async fn outbound_raw( &self, dest: Ipv4Addr, @@ -357,6 +413,12 @@ impl HybridOutbound { data: TransmissionBytes, mut dest: Ipv4Addr, ) -> anyhow::Result<()> { + if self.is_ikev2_client(&dest) { + let Some(ip) = crate::ethernet::strip_ipv4(data) else { + return Ok(()); + }; + return self.ikev2_relay_outbound(net, ip, dest).await; + } if dest == net.gateway { let Some(ip) = crate::ethernet::strip_ipv4(data) else { return Ok(()); diff --git a/vnt-core/src/tunnel_core/server/connection_manager.rs b/vnt-core/src/tunnel_core/server/connection_manager.rs index 1c89e743..fdecb304 100644 --- a/vnt-core/src/tunnel_core/server/connection_manager.rs +++ b/vnt-core/src/tunnel_core/server/connection_manager.rs @@ -40,6 +40,7 @@ pub struct InboundHandlerConfig { pub fec_decoder: FecDecoder, pub turn: Arc>, pub auto_sync_subnet: bool, + pub allow_ikev2: bool, } pub struct ServerTurnManager { diff --git a/vnt-core/src/tunnel_core/server/inbound.rs b/vnt-core/src/tunnel_core/server/inbound.rs index 5bd7d070..9dce80c0 100644 --- a/vnt-core/src/tunnel_core/server/inbound.rs +++ b/vnt-core/src/tunnel_core/server/inbound.rs @@ -361,6 +361,7 @@ pub(crate) struct ServerTurnInboundHandler { fec_decoder: FecDecoder, turn: Arc>, auto_sync_subnet: bool, + allow_ikev2: bool, } impl ServerTurnInboundHandler { pub fn new( @@ -383,6 +384,7 @@ impl ServerTurnInboundHandler { fec_decoder: config.fec_decoder, turn: config.turn, auto_sync_subnet: config.auto_sync_subnet, + allow_ikev2: config.allow_ikev2, } } fn network_contains(&self, ip: &Ipv4Addr) -> bool { @@ -421,6 +423,24 @@ impl ServerTurnInboundHandler { let mut net_packet = self.packet_compression.decompress(net_packet)?; match msg_type { + MsgType::Ikev2Relay if self.allow_ikev2 => { + let Some(ipv4) = Ipv4Packet::new(net_packet.payload()) else { + return Ok(()); + }; + let header_length = ipv4.get_header_length() as usize * 4; + if ipv4.get_version() != 4 + || header_length < Ipv4Packet::minimum_packet_size() + || ipv4.get_total_length() as usize != net_packet.payload().len() + || ipv4.get_source() != src + || ipv4.get_destination() != network_addr.ip + || Ipv4Addr::from(net_packet.dest_id()) != network_addr.ip + { + return Ok(()); + } + self.enhanced_inbound + .inbound(&network_addr, MsgType::Turn, src, net_packet) + .await?; + } MsgType::Turn => { // 只允许icmp EchoReply let Some(ipv4) = Ipv4Packet::new(net_packet.payload()) else { diff --git a/vnt-core/src/tunnel_core/server/transport/config.rs b/vnt-core/src/tunnel_core/server/transport/config.rs index e834abe0..c805debe 100644 --- a/vnt-core/src/tunnel_core/server/transport/config.rs +++ b/vnt-core/src/tunnel_core/server/transport/config.rs @@ -39,6 +39,7 @@ pub(crate) struct ConnectRegConfig { pub key_sign: Option, pub ip_variable: bool, pub advertised_subnets: Arc>, + pub allow_ikev2: bool, pub default_interface: Option, } #[derive(Debug, Clone)] @@ -131,6 +132,7 @@ impl ConnectRegConfig { server_id, registration_mode, advertised_subnets: self.advertised_subnets.as_ref().clone(), + allow_ikev2: self.allow_ikev2, } } /// 解析出全部候选服务器地址:动态地址(DNS TXT 记录或 http(s) 接口返回的 @@ -263,6 +265,7 @@ mod tests { key_sign: None, ip_variable: true, advertised_subnets: Arc::new(Vec::new()), + allow_ikev2: false, default_interface: None, }; diff --git a/vnt-jni/java_example/com/vnt/VntConfig.java b/vnt-jni/java_example/com/vnt/VntConfig.java index bcca0499..de849c84 100644 --- a/vnt-jni/java_example/com/vnt/VntConfig.java +++ b/vnt-jni/java_example/com/vnt/VntConfig.java @@ -29,6 +29,7 @@ public class VntConfig { private final String certMode; private final boolean noPunch; private final boolean noBroadcast; + private final boolean allowIkev2; private final boolean compress; private final boolean rtx; private final boolean fec; @@ -58,6 +59,7 @@ private VntConfig(Builder builder) { this.certMode = builder.certMode; this.noPunch = builder.noPunch; this.noBroadcast = builder.noBroadcast; + this.allowIkev2 = builder.allowIkev2; this.compress = builder.compress; this.rtx = builder.rtx; this.fec = builder.fec; @@ -117,6 +119,7 @@ String toJson() { // 布尔值 json.put("no_punch", noPunch); json.put("no_broadcast", noBroadcast); + json.put("allow_ikev2", allowIkev2); json.put("compress", compress); json.put("rtx", rtx); json.put("fec", fec); @@ -181,6 +184,7 @@ public static class Builder { private String certMode; private boolean noPunch = false; private boolean noBroadcast = false; + private boolean allowIkev2 = false; private boolean compress = false; private boolean rtx = false; private boolean fec = false; @@ -325,6 +329,14 @@ public Builder setNoBroadcast(boolean noBroadcast) { return this; } + /** + * 允许与服务端接入的 IKEv2/IPsec 客户端通信(默认false) + */ + public Builder setAllowIkev2(boolean allowIkev2) { + this.allowIkev2 = allowIkev2; + return this; + } + /** * 启用压缩(默认false) */ diff --git a/vnt-jni/src/lib.rs b/vnt-jni/src/lib.rs index 9d51cd7f..676e18c8 100644 --- a/vnt-jni/src/lib.rs +++ b/vnt-jni/src/lib.rs @@ -565,7 +565,7 @@ pub extern "system" fn Java_com_vnt_VntApi_nativeGetClientList<'local>( let local_clients: HashMap<_, _> = api .client_ips() .into_iter() - .map(|client| (client.ip, client.online)) + .map(|client| (client.ip, (client.online, client.client_type))) .collect(); let server_clients: HashMap<_, _> = runtime .block_on(api.server_rpc().client_list()) @@ -601,7 +601,7 @@ pub extern "system" fn Java_com_vnt_VntApi_nativeGetClientList<'local>( .map(|route| route.route_key().protocol().to_string()); let route_metric = route.as_ref().map(|route| route.metric()); let rtt = route.as_ref().map(|route| route.rtt()); - let online = local_clients.get(&ip).copied().unwrap_or(false) + let online = local_clients.get(&ip).map(|value| value.0).unwrap_or(false) || server_client.map(|client| client.online).unwrap_or(false) || has_route; let packet_loss = api.packet_loss_info(&ip).map(|info| { @@ -621,13 +621,24 @@ pub extern "system" fn Java_com_vnt_VntApi_nativeGetClientList<'local>( "ip": ip.to_string(), "name": server_client.map(|client| client.name.as_str()).unwrap_or(""), "version": server_client.map(|client| client.version.as_str()).unwrap_or(""), + "client_type": server_client + .map(|client| if client.client_type == 1 { "IKEV2" } else { "VNT" }) + .or_else(|| local_clients.get(&ip).map(|value| match value.1 { + vnt_core::protocol::control_message::ClientType::Ikev2 => "IKEV2", + vnt_core::protocol::control_message::ClientType::Vnt => "VNT", + })) + .unwrap_or("VNT"), "online": online, "direct": direct, "route_protocol": route_protocol, "route_metric": route_metric, "rtt": rtt, "key_equal": server_client - .map(|client| encryption_state(local_key.as_deref(), client.key_sign.as_deref())) + .map(|client| if client.client_type == 1 { + 0 + } else { + encryption_state(local_key.as_deref(), client.key_sign.as_deref()) + }) .unwrap_or(0), "packet_loss": packet_loss, "traffic": traffic, @@ -1047,6 +1058,8 @@ fn parse_config_from_json(json_str: &str) -> anyhow::Result { #[serde(default)] no_broadcast: bool, #[serde(default)] + allow_ikev2: bool, + #[serde(default)] compress: bool, #[serde(default)] rtx: bool, @@ -1162,6 +1175,7 @@ fn parse_config_from_json(json_str: &str) -> anyhow::Result { ip: cfg.ip, no_punch: cfg.no_punch, no_broadcast: cfg.no_broadcast, + allow_ikev2: cfg.allow_ikev2, rtx: cfg.rtx, compress: cfg.compress, device_id, diff --git a/vnt-web/src/service_http.rs b/vnt-web/src/service_http.rs index 3c5d62d7..308b2e02 100644 --- a/vnt-web/src/service_http.rs +++ b/vnt-web/src/service_http.rs @@ -270,6 +270,8 @@ pub struct StartConfig { #[serde(default)] pub no_broadcast: bool, #[serde(default)] + pub allow_ikev2: bool, + #[serde(default)] pub compress: bool, #[serde(default)] pub rtx: bool, @@ -375,6 +377,7 @@ struct HttpClientItem { online: bool, route: Option, version: String, + client_type: String, last_connected_time: i64, key_equal: i32, nat_info: Option, @@ -1468,6 +1471,7 @@ fn convert_config(cfg: StartConfig) -> anyhow::Result { ip: cfg.ip, no_punch: cfg.no_punch, no_broadcast: cfg.no_broadcast, + allow_ikev2: cfg.allow_ikev2, rtx: cfg.rtx, compress: cfg.compress, device_id, @@ -1596,6 +1600,11 @@ async fn get_peers( online: v.online || has_route, route, version: String::new(), + client_type: match v.client_type { + vnt_core::protocol::control_message::ClientType::Ikev2 => "IKEV2", + vnt_core::protocol::control_message::ClientType::Vnt => "VNT", + } + .to_string(), last_connected_time: 0, key_equal: 0, nat_info: build_nat_info(&ip), @@ -1613,6 +1622,7 @@ async fn get_peers( let route = build_route(&ip); // 如果有路由,说明设备在线(可以直接通信) let has_route = route.is_some(); + let client_type = if v.client_type == 1 { "IKEV2" } else { "VNT" }; merged.insert( ip, HttpClientItem { @@ -1621,8 +1631,13 @@ async fn get_peers( online: v.online || has_route, route, version: v.version, + client_type: client_type.to_string(), last_connected_time: v.last_connected_time, - key_equal: calc_key_equal(&v.key_sign), + key_equal: if v.client_type == 1 { + 0 + } else { + calc_key_equal(&v.key_sign) + }, nat_info: build_nat_info(&ip), packet_loss: build_packet_loss(&ip), traffic: build_traffic(&ip), @@ -1807,6 +1822,7 @@ mod tests { password: None, no_punch: false, no_broadcast: false, + allow_ikev2: false, compress: false, rtx: false, fec: false, diff --git a/vnt-web/ui/src/utils/configHelp.js b/vnt-web/ui/src/utils/configHelp.js index f58a62ad..64a78bef 100644 --- a/vnt-web/ui/src/utils/configHelp.js +++ b/vnt-web/ui/src/utils/configHelp.js @@ -89,6 +89,13 @@ export const configHelp = { usage: "可减少发现协议和广播风暴产生的流量;依赖局域网发现、组播服务或某些游戏联机时不要开启。", format: "开关;默认关闭,即允许转发。", }, + allow_ikev2: { + param: "allow_ikev2 / --allow-ikev2", + summary: "允许本节点与接入同一虚拟网络的 IKEv2/IPsec 客户端互通。", + usage: "开启后,本节点信任已认证 VNT 服务端注入的 IKEv2 明文 IPv4 包,并把发往 IKEv2 类型地址的流量固定交给服务端中继。", + format: "开关;默认关闭。", + notes: ["此路径不使用 VNT 节点间 password 加密,应只连接受信任的服务端。"], + }, password: { param: "password", summary: "对虚拟网络中的节点间数据进行加密。", diff --git a/vnt-web/ui/src/utils/toml.js b/vnt-web/ui/src/utils/toml.js index cdd7b42e..d1f5fef4 100644 --- a/vnt-web/ui/src/utils/toml.js +++ b/vnt-web/ui/src/utils/toml.js @@ -13,6 +13,7 @@ export const emptyFormData = () => ({ compress: false, no_punch: false, no_broadcast: false, + allow_ikev2: false, input: [], subnet_mapping: [], output: [], @@ -83,6 +84,8 @@ export const parseTomlToForm = (toml) => { data.no_punch = trimmed.includes("true"); } else if (trimmed.match(/^no_broadcast\s*=/)) { data.no_broadcast = trimmed.includes("true"); + } else if (trimmed.match(/^allow_ikev2\s*=/)) { + data.allow_ikev2 = trimmed.includes("true"); } else if (trimmed.startsWith("input")) { const match = trimmed.match(/input\s*=\s*\[(.*)\]/); if (match) { @@ -228,6 +231,11 @@ export const formToToml = (formData) => { toml += "no_broadcast = true\n"; } + if (formData.allow_ikev2) { + toml += "\n# 允许与 IKEv2 客户端通信,并信任服务端注入的数据\n"; + toml += "allow_ikev2 = true\n"; + } + if (formData.compress) { toml += "\n# 是否启用 LZ4 压缩 (默认 false)\n"; toml += "compress = true\n"; @@ -377,6 +385,9 @@ server = ["quic://1.2.3.4:29872"] # 是否关闭 IPv4 广播和组播转发 (默认 false,即开启) # no_broadcast = false +# 是否允许与 IKEv2 客户端通信,并信任服务端注入的 IKEv2 明文 IPv4 包 +# allow_ikev2 = false + # 是否启用 LZ4 压缩 (默认 false,设置为true时开启) # compress = false diff --git a/vnt-web/ui/src/views/ConfigEditor.vue b/vnt-web/ui/src/views/ConfigEditor.vue index acf8a32f..d8b3fc09 100644 --- a/vnt-web/ui/src/views/ConfigEditor.vue +++ b/vnt-web/ui/src/views/ConfigEditor.vue @@ -435,6 +435,15 @@ const sectionTitleClass = "text-md mb-4 flex items-center font-bold text-slate-9 + diff --git a/vnt-web/ui/src/views/PeersView.vue b/vnt-web/ui/src/views/PeersView.vue index 407db70c..52f93338 100644 --- a/vnt-web/ui/src/views/PeersView.vue +++ b/vnt-web/ui/src/views/PeersView.vue @@ -216,7 +216,10 @@ const switcherClass = (fileName) => {{ peer.ip }} - {{ peer.name || "-" }} + + {{ peer.name || "-" }} + {{ peer.client_type || "VNT" }} + {{ peer.version || "-" }}
From fc6320374fb4aa740be4f17d301919b01648cac0 Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Fri, 4 Sep 2026 22:07:43 +0800 Subject: [PATCH 02/14] fix(p2p): improve punch scheduling and backoff --- vnt-core/src/context/mod.rs | 81 ++++++++++++++++++- vnt-core/src/context/nat.rs | 80 ++++++++++++++++-- vnt-core/src/core/mod.rs | 1 + vnt-core/src/tunnel_core/p2p/inbound.rs | 20 ++++- vnt-core/src/tunnel_core/p2p/route_table.rs | 6 +- .../src/tunnel_core/p2p/transport/punch.rs | 49 +++++------ vnt-core/src/tunnel_core/server/inbound.rs | 6 +- 7 files changed, 203 insertions(+), 40 deletions(-) diff --git a/vnt-core/src/context/mod.rs b/vnt-core/src/context/mod.rs index 1a52af7e..1ac3ebe5 100644 --- a/vnt-core/src/context/mod.rs +++ b/vnt-core/src/context/mod.rs @@ -583,7 +583,7 @@ impl ServerInfoCollection { self_ip: Ipv4Addr, client_simple_list: ClientSimpleInfoList, now: i64, - ) { + ) -> Vec { let mut guard = self.server_node_map.write(); let server_node = guard.entry(server_id).or_default(); if now > client_simple_list.time { @@ -619,7 +619,21 @@ impl ServerInfoCollection { } let mut guard = self.client_simple_list.write(); - *guard = client_simple_map.into_values().collect() + let previous_online: HashMap = + guard.iter().map(|info| (info.ip, info.online)).collect(); + let mut changed = Vec::new(); + for (ip, info) in &client_simple_map { + if previous_online.get(ip).copied().unwrap_or(false) != info.online { + changed.push(*ip); + } + } + for (ip, was_online) in previous_online { + if was_online && !client_simple_map.contains_key(&ip) { + changed.push(ip); + } + } + *guard = client_simple_map.into_values().collect(); + changed } pub fn set_server_connected(&self, server_id: u32, val: bool) -> bool { let mut mutex_guard = self.server_node_map.write(); @@ -808,6 +822,69 @@ mod subnet_sync_tests { ); } } + +#[cfg(test)] +mod client_status_tests { + use super::ServerInfoCollection; + use crate::protocol::control_message::{ClientSimpleInfo, ClientSimpleInfoList, ClientType}; + use crate::tunnel_core::server::transport::config::ProtocolAddress; + use std::net::Ipv4Addr; + + fn update(servers: &ServerInfoCollection, peer: Ipv4Addr, online: bool) -> Vec { + servers.update_client_simple_list( + 0, + Ipv4Addr::new(10, 26, 0, 1), + ClientSimpleInfoList { + data_version: 1, + list: vec![ClientSimpleInfo { + ip: peer, + online, + client_type: ClientType::Vnt, + }], + is_all: true, + time: 0, + }, + 0, + ) + } + + #[test] + fn online_transitions_are_reported_once() { + let servers = ServerInfoCollection::default(); + servers.update_server(vec![(0, ProtocolAddress::default())]); + let peer = Ipv4Addr::new(10, 26, 0, 2); + + assert_eq!(update(&servers, peer, true), vec![peer]); + assert!(update(&servers, peer, true).is_empty()); + assert_eq!(update(&servers, peer, false), vec![peer]); + assert!(update(&servers, peer, false).is_empty()); + assert_eq!(update(&servers, peer, true), vec![peer]); + } + + #[test] + fn disappearing_from_full_snapshot_is_one_offline_transition() { + let servers = ServerInfoCollection::default(); + servers.update_server(vec![(0, ProtocolAddress::default())]); + let peer = Ipv4Addr::new(10, 26, 0, 2); + assert_eq!(update(&servers, peer, true), vec![peer]); + + let empty_snapshot = || ClientSimpleInfoList { + data_version: 2, + list: Vec::new(), + is_all: true, + time: 0, + }; + assert_eq!( + servers.update_client_simple_list(0, Ipv4Addr::new(10, 26, 0, 1), empty_snapshot(), 0,), + vec![peer] + ); + assert!( + servers + .update_client_simple_list(0, Ipv4Addr::new(10, 26, 0, 1), empty_snapshot(), 0,) + .is_empty() + ); + } +} #[derive(Copy, Clone, Debug)] pub struct NetworkAddr { pub gateway: Ipv4Addr, diff --git a/vnt-core/src/context/nat.rs b/vnt-core/src/context/nat.rs index 6a96bff4..79d6774c 100644 --- a/vnt-core/src/context/nat.rs +++ b/vnt-core/src/context/nat.rs @@ -136,7 +136,7 @@ fn mapping_addr(addr: SocketAddr) -> Option<(Ipv4Addr, u16)> { #[derive(Copy, Clone, Debug)] pub struct PunchState { - /// 打洞交互累计次数,退避时长按「BASE × count」线性增长,封顶 MAX_BACKOFF。 + /// 连续失败的打洞轮次,成功或节点上下线时清零。 pub count: u32, /// 该目标下次允许打洞的时刻,到点之前不应再次打洞。 /// 用单调时钟 Instant,不受系统时间调整影响。 @@ -149,20 +149,50 @@ pub struct PunchBackoff { } impl PunchBackoff { - const MAX_BACKOFF: Duration = Duration::from_secs(3_600); // 常规退避上限 1h const NAT_CHANGE_CAP: Duration = Duration::from_secs(600); // NAT 变化后最多等待 10 分钟 - const BASE: Duration = Duration::from_secs(3); + const RETRY_BACKOFF: [Duration; 7] = [ + Duration::from_secs(3), + Duration::from_secs(10), + Duration::from_secs(30), + Duration::from_secs(60), + Duration::from_secs(300), + Duration::from_secs(900), + Duration::from_secs(3_600), + ]; - /// 记录一次打洞交互:退避时长 = BASE × 累计次数,封顶 MAX_BACKOFF。 - pub fn record(&self, ip: Ipv4Addr) { + #[cfg(test)] + fn record(&self, ip: Ipv4Addr) { let now = Instant::now(); let mut map = self.inner.write(); let entry = map.entry(ip).or_insert(PunchState { count: 0, backoff_until: now, }); - entry.count += 1; - entry.backoff_until = now + (Self::BASE * entry.count).min(Self::MAX_BACKOFF); + Self::record_entry(entry, now); + } + + /// 原子地检查并占用一轮打洞机会。 + /// + /// 返回 `true` 时已经立即更新退避期,即使后续协商包丢失, + /// 也不会在下一调度周期中无限重发。 + pub fn try_begin(&self, ip: Ipv4Addr) -> bool { + let now = Instant::now(); + let mut map = self.inner.write(); + let entry = map.entry(ip).or_insert(PunchState { + count: 0, + backoff_until: now, + }); + if now < entry.backoff_until { + return false; + } + Self::record_entry(entry, now); + true + } + + fn record_entry(entry: &mut PunchState, now: Instant) { + entry.count = entry.count.saturating_add(1); + let index = (entry.count.saturating_sub(1) as usize).min(Self::RETRY_BACKOFF.len() - 1); + entry.backoff_until = now + Self::RETRY_BACKOFF[index]; } pub fn should_punch(&self, ip: Ipv4Addr) -> bool { @@ -171,6 +201,11 @@ impl PunchBackoff { .is_none_or(|state| Instant::now() >= state.backoff_until) } + /// 直连建立成功或节点在线状态变化后,忘记该节点的历史失败。 + pub fn reset(&self, ip: Ipv4Addr) { + self.inner.write().remove(&ip); + } + /// 对端 NAT 变化:不删除退避记录,把该目标的退避截止时刻压缩到 /// 「当前 + NAT_CHANGE_CAP」以内,并归零增长指数,允许尽快重新尝试打洞。 pub fn cap(&self, ip: Ipv4Addr) { @@ -295,6 +330,37 @@ mod tests { assert!(backoff.should_punch(ip)); } + #[test] + fn try_begin_atomically_reserves_attempt_and_advances_schedule() { + let backoff = PunchBackoff::default(); + let ip = Ipv4Addr::new(10, 26, 0, 2); + + assert!(backoff.try_begin(ip)); + assert!(!backoff.try_begin(ip), "同一退避期只能占用一次"); + + backoff.inner.write().get_mut(&ip).unwrap().backoff_until = + Instant::now() - Duration::from_secs(1); + let now = Instant::now(); + assert!(backoff.try_begin(ip)); + let state = backoff.inner.read()[&ip]; + assert_eq!(state.count, 2); + assert!(state.backoff_until - now >= Duration::from_secs(9)); + } + + #[test] + fn reset_forgets_previous_failures() { + let backoff = PunchBackoff::default(); + let ip = Ipv4Addr::new(10, 26, 0, 2); + assert!(backoff.try_begin(ip)); + assert!(!backoff.should_punch(ip)); + + backoff.reset(ip); + + assert!(backoff.should_punch(ip)); + assert!(backoff.try_begin(ip)); + assert_eq!(backoff.inner.read()[&ip].count, 1); + } + #[test] fn record_after_expiry_keeps_growth() { let backoff = PunchBackoff::default(); diff --git a/vnt-core/src/core/mod.rs b/vnt-core/src/core/mod.rs index 7d4bff92..c965fd3c 100644 --- a/vnt-core/src/core/mod.rs +++ b/vnt-core/src/core/mod.rs @@ -276,6 +276,7 @@ impl NetworkManager { fec_decoder: fec_decoder.clone(), turn: turn.clone(), basic_outbound, + punch_backoff: app_state.punch_backoff.clone(), }); p2p_task.start(handler); } diff --git a/vnt-core/src/tunnel_core/p2p/inbound.rs b/vnt-core/src/tunnel_core/p2p/inbound.rs index 79fbf43b..60d22f8b 100644 --- a/vnt-core/src/tunnel_core/p2p/inbound.rs +++ b/vnt-core/src/tunnel_core/p2p/inbound.rs @@ -1,5 +1,6 @@ use crate::compression::PacketCompression; use crate::context::config::{TurnRule, allow_punch}; +use crate::context::nat::PunchBackoff; use crate::context::{NetworkAddr, NetworkRoute, PacketLossStats}; use crate::crypto::PacketCrypto; use crate::enhanced_tunnel::inbound::EnhancedInbound; @@ -59,6 +60,7 @@ pub(crate) struct P2pInboundConfig { pub fec_decoder: FecDecoder, pub turn: Arc>, pub basic_outbound: BasicOutbound, + pub punch_backoff: PunchBackoff, } #[derive(Clone)] @@ -72,6 +74,7 @@ pub(crate) struct P2pInboundHandler { fec_decoder: FecDecoder, turn: Arc>, basic_outbound: BasicOutbound, + punch_backoff: PunchBackoff, } impl P2pInboundHandler { @@ -86,6 +89,7 @@ impl P2pInboundHandler { fec_decoder: config.fec_decoder, turn: config.turn, basic_outbound: config.basic_outbound, + punch_backoff: config.punch_backoff, } } fn network_contains(&self, ip: &Ipv4Addr) -> bool { @@ -283,7 +287,9 @@ impl P2pInboundHandler { ctx.src_ip, ctx.dest_ip ); - self.route_table.add_owner_route(ctx.src_ip, route_key); + if self.route_table.add_owner_route(ctx.src_ip, route_key) { + self.punch_backoff.reset(ctx.src_ip); + } let mut packet = build_handshake_response( MsgType::PunchRes, net.ip, @@ -317,7 +323,9 @@ impl P2pInboundHandler { ctx.src_ip, ctx.dest_ip ); - self.route_table.add_owner_route(ctx.src_ip, route_key); + if self.route_table.add_owner_route(ctx.src_ip, route_key) { + self.punch_backoff.reset(ctx.src_ip); + } } MsgType::DirectConnectReq => { if !valid_punch_source(net, ctx.src_ip) { @@ -332,7 +340,9 @@ impl P2pInboundHandler { ctx.src_ip, net.ip ); - self.route_table.add_owner_route(ctx.src_ip, route_key); + if self.route_table.add_owner_route(ctx.src_ip, route_key) { + self.punch_backoff.reset(ctx.src_ip); + } let mut packet = build_handshake_response( MsgType::DirectConnectRes, net.ip, @@ -355,7 +365,9 @@ impl P2pInboundHandler { ctx.src_ip, ctx.dest_ip ); - self.route_table.add_owner_route(ctx.src_ip, route_key); + if self.route_table.add_owner_route(ctx.src_ip, route_key) { + self.punch_backoff.reset(ctx.src_ip); + } } MsgType::PingTurn => {} MsgType::PongTurn => {} diff --git a/vnt-core/src/tunnel_core/p2p/route_table.rs b/vnt-core/src/tunnel_core/p2p/route_table.rs index d3e0de25..0ebe25f2 100644 --- a/vnt-core/src/tunnel_core/p2p/route_table.rs +++ b/vnt-core/src/tunnel_core/p2p/route_table.rs @@ -185,10 +185,12 @@ impl RouteTable { } /// 添加 owner 路由(打洞请求响应时调用) - pub fn add_owner_route(&self, id: Ipv4Addr, key: RouteKey) { - if self.inner.add_owner_route(id, key) { + pub fn add_owner_route(&self, id: Ipv4Addr, key: RouteKey) -> bool { + let first_direct = self.inner.add_owner_route(id, key); + if first_direct { self.inner.first_direct_route_notify.notify_one(); } + first_direct } /// 添加路由(心跳时调用,用于更新路由时间和添加跨节点转发路由) diff --git a/vnt-core/src/tunnel_core/p2p/transport/punch.rs b/vnt-core/src/tunnel_core/p2p/transport/punch.rs index f049c246..e54c0f0e 100644 --- a/vnt-core/src/tunnel_core/p2p/transport/punch.rs +++ b/vnt-core/src/tunnel_core/p2p/transport/punch.rs @@ -38,35 +38,38 @@ pub async fn punch_task( let Some(punch_info) = (ctx.punch_info_getter)() else { continue; }; + if !ctx.server_info.is_any_server_connected(None) { + continue; + } let mut list = ctx.server_info.client_online_ips(); + list.retain(|dest_ip| { + *dest_ip > src_ip + && allow_punch(&ctx.turn, dest_ip) + && route_table.need_punch(dest_ip) + && ctx.punch_backoff.should_punch(*dest_ip) + }); list.shuffle(&mut rand::rng()); list.truncate(5); for dest_ip in list { - if dest_ip <= src_ip { - continue; - } - if !allow_punch(&ctx.turn, &dest_ip) { + // should_punch() 只用于候选过滤;最终必须原子占用, + // 防止并发调度或重复报文同时启动同一目标。 + if !ctx.punch_backoff.try_begin(dest_ip) { continue; } - if ctx.server_info.is_any_server_connected(None) && route_table.need_punch(&dest_ip) { - if !ctx.punch_backoff.should_punch(dest_ip) { - continue; - } - log::info!("punching {dest_ip}"); + log::info!("punching {dest_ip}"); - let data = punch_info.encode(); - let mut net_packet = NetPacket::new(TransmissionBytes::zeroed_size( - HEAD_LENGTH + data.len(), - tunnel_to_server.encrypt_reserve(), - ))?; - net_packet.set_msg_type(MsgType::PunchStart1); - net_packet.set_ttl(2); - net_packet.set_src_id(src_ip.into()); - net_packet.set_dest_id(dest_ip.into()); - net_packet.set_payload(data.as_ref())?; - if let Err(e) = tunnel_to_server.send(dest_ip, net_packet).await { - error!("punch send error {:?}", e); - } + let data = punch_info.encode(); + let mut net_packet = NetPacket::new(TransmissionBytes::zeroed_size( + HEAD_LENGTH + data.len(), + tunnel_to_server.encrypt_reserve(), + ))?; + net_packet.set_msg_type(MsgType::PunchStart1); + net_packet.set_ttl(2); + net_packet.set_src_id(src_ip.into()); + net_packet.set_dest_id(dest_ip.into()); + net_packet.set_payload(data.as_ref())?; + if let Err(e) = tunnel_to_server.send(dest_ip, net_packet).await { + error!("punch send error {:?}", e); } } } @@ -97,7 +100,7 @@ impl NatPuncher { if self.puncher.is_none() { return Ok(false); } - if !self.punch_backoff.should_punch(dest_ip) { + if !self.punch_backoff.try_begin(dest_ip) { return Ok(false); } self.punch_uncheck_delay(dest_ip, punch_info, Some(Duration::from_millis(50)))?; diff --git a/vnt-core/src/tunnel_core/server/inbound.rs b/vnt-core/src/tunnel_core/server/inbound.rs index 9dce80c0..8cd8696f 100644 --- a/vnt-core/src/tunnel_core/server/inbound.rs +++ b/vnt-core/src/tunnel_core/server/inbound.rs @@ -482,12 +482,15 @@ impl ServerTurnInboundHandler { } MsgType::PushClientIps => { let list = ClientSimpleInfoList::from_slice(net_packet.payload())?; - self.server_info.update_client_simple_list( + let changed = self.server_info.update_client_simple_list( self.server_id, network_addr.ip, list, now, ); + for ip in changed { + self.punch_backoff.reset(ip); + } } MsgType::RpcRes => { // 设置rpc响应 @@ -662,7 +665,6 @@ impl ServerTurnInboundHandler { log::debug!("ignore configured turn target PunchStart2 from {src}"); return Ok(()); } - self.punch_backoff.record(src); // 对方回复开始打洞 let peer_punch_info = PunchInfo::from_slice(net_packet.payload())?; self.update_peer_nat_info(src, peer_punch_info.nat_info.clone()); From 751db88ae4e4e1d492cc4956cc36df4370c758ac Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Fri, 4 Sep 2026 22:32:29 +0800 Subject: [PATCH 03/14] fix(p2p): limit concurrent punch tasks --- .../src/tunnel_core/p2p/transport/punch.rs | 84 ++++++++++++++++++- 1 file changed, 80 insertions(+), 4 deletions(-) diff --git a/vnt-core/src/tunnel_core/p2p/transport/punch.rs b/vnt-core/src/tunnel_core/p2p/transport/punch.rs index e54c0f0e..c7f34d98 100644 --- a/vnt-core/src/tunnel_core/p2p/transport/punch.rs +++ b/vnt-core/src/tunnel_core/p2p/transport/punch.rs @@ -14,6 +14,28 @@ use rustp2p_core::punch::{PunchModel, Puncher}; use std::net::Ipv4Addr; use std::sync::Arc; use std::time::Duration; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; + +const MAX_CONCURRENT_PUNCHES: usize = 4; + +#[derive(Clone)] +struct PunchLimiter { + semaphore: Arc, +} + +impl Default for PunchLimiter { + fn default() -> Self { + Self { + semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_PUNCHES)), + } + } +} + +impl PunchLimiter { + fn try_acquire(&self) -> Option { + self.semaphore.clone().try_acquire_owned().ok() + } +} pub struct PunchTaskContext { pub network: SharedNetworkAddr, @@ -80,6 +102,7 @@ pub struct NatPuncher { punch_backoff: PunchBackoff, puncher: Option, packet_crypto: PacketCrypto, + limiter: PunchLimiter, } impl NatPuncher { @@ -94,22 +117,33 @@ impl NatPuncher { punch_backoff, puncher, packet_crypto, + limiter: PunchLimiter::default(), } } pub fn punch(&self, dest_ip: Ipv4Addr, punch_info: PunchInfo) -> anyhow::Result { - if self.puncher.is_none() { + let Some(puncher) = self.puncher.clone() else { return Ok(false); - } + }; + let Some(permit) = self.limiter.try_acquire() else { + log::debug!("skip punch to {dest_ip}: concurrent punch limit reached"); + return Ok(false); + }; if !self.punch_backoff.try_begin(dest_ip) { return Ok(false); } - self.punch_uncheck_delay(dest_ip, punch_info, Some(Duration::from_millis(50)))?; + self.spawn_punch( + puncher, + dest_ip, + punch_info, + Some(Duration::from_millis(50)), + permit, + )?; Ok(true) } pub fn punch_uncheck(&self, dest_ip: Ipv4Addr, punch_info: PunchInfo) -> anyhow::Result<()> { self.punch_uncheck_delay(dest_ip, punch_info, None) } - pub fn punch_uncheck_delay( + fn punch_uncheck_delay( &self, dest_ip: Ipv4Addr, punch_info: PunchInfo, @@ -118,11 +152,26 @@ impl NatPuncher { let Some(puncher) = self.puncher.clone() else { return Ok(()); }; + let Some(permit) = self.limiter.try_acquire() else { + log::debug!("skip punch to {dest_ip}: concurrent punch limit reached"); + return Ok(()); + }; + self.spawn_punch(puncher, dest_ip, punch_info, time, permit) + } + fn spawn_punch( + &self, + puncher: Puncher, + dest_ip: Ipv4Addr, + punch_info: PunchInfo, + time: Option, + permit: OwnedSemaphorePermit, + ) -> anyhow::Result<()> { let Some(src_ip) = self.network.ip() else { bail!("not ip"); }; let packet_crypto = self.packet_crypto.clone(); tokio::spawn(async move { + let _permit = permit; if let Some(time) = time { tokio::time::sleep(time).await; } @@ -157,3 +206,30 @@ async fn punch_now( .await?; Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn punch_limiter_is_shared_and_releases_capacity() { + let limiter = PunchLimiter::default(); + let clone = limiter.clone(); + let mut permits = Vec::new(); + + for _ in 0..MAX_CONCURRENT_PUNCHES { + permits.push(limiter.try_acquire().expect("前四个任务应获得许可")); + } + assert!( + clone.try_acquire().is_none(), + "克隆必须共享同一个全局并发上限" + ); + + permits.pop(); + + assert!( + clone.try_acquire().is_some(), + "任务结束释放许可后应能立即开始下一轮" + ); + } +} From b956b41faa78cae2cd74083eda438dc8669360f6 Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Fri, 4 Sep 2026 22:57:15 +0800 Subject: [PATCH 04/14] feat(p2p): configure punch modes per target --- README.md | 6 + src/args_config.rs | 62 ++++++- vnt-core/proto/client.proto | 4 +- vnt-core/src/context/config.rs | 166 ++++++++++++++++++ vnt-core/src/context/mod.rs | 11 +- vnt-core/src/core/mod.rs | 5 + vnt-core/src/protocol/client_message.rs | 109 +++++++++++- .../src/tunnel_core/p2p/transport/punch.rs | 83 ++++++++- .../src/tunnel_core/p2p/transport/task.rs | 2 +- .../tunnel_core/server/connection_manager.rs | 3 +- vnt-core/src/tunnel_core/server/inbound.rs | 9 +- vnt-jni/java_example/com/vnt/VntConfig.java | 13 ++ vnt-jni/src/lib.rs | 33 +++- vnt-web/src/service_http.rs | 34 +++- vnt-web/ui/src/utils/configHelp.js | 8 + vnt-web/ui/src/utils/toml.js | 16 ++ vnt-web/ui/src/views/ConfigEditor.vue | 34 ++++ 17 files changed, 577 insertions(+), 21 deletions(-) diff --git a/README.md b/README.md index 4dc8c091..37237120 100644 --- a/README.md +++ b/README.md @@ -58,6 +58,12 @@ P2P 打洞;中转节点已有直连路由时,数据包会优先发给中转节点继续转发。中转 IP 也可以填写当前虚拟网络的网关 IP,此时强制通过服务器中继。 +可重复使用 `--punch-model <目标IP或CIDR,打洞方式...>`(配置文件中为 +`punch_model = ["10.26.0.2,IPv4Udp", "10.26.1.0/24,IPv4Tcp,IPv4Udp"]`) +限制指定目标可使用的打洞方式。可选方式为 `IPv4Tcp`、`IPv4Udp`、`IPv6Tcp`、 +`IPv6Udp`;双方实际采用各自允许集合的交集。重叠规则按最长前缀匹配,同一目标的 +多条规则会合并;未命中规则时默认允许全部方式。 + IPv4 广播和组播默认开启。可使用 `--no-broadcast`(配置文件中为 `no_broadcast = true`)停止转发本机发出的 IPv4 广播和组播;ARP、非 IPv4 二层广播 以及单播流量不受影响。 diff --git a/src/args_config.rs b/src/args_config.rs index fef92fa9..14e8ce47 100644 --- a/src/args_config.rs +++ b/src/args_config.rs @@ -4,7 +4,7 @@ use ipnet::Ipv4Net; use serde::{Deserialize, Serialize}; use std::net::Ipv4Addr; use std::path::{Path, PathBuf}; -use vnt_core::context::config::{Config, DeviceMode, PeerAddress, TurnRule}; +use vnt_core::context::config::{Config, DeviceMode, PeerAddress, PunchRule, TurnRule}; use vnt_core::nat::{NetInput, SubnetMapping}; use vnt_core::tls::verifier::CertValidationMode; use vnt_core::tunnel_core::server::transport::config::ProtocolAddress; @@ -16,6 +16,7 @@ pub struct FileConfig { pub server: Option>, pub peer_address: Option>, pub turn: Option>, + pub punch_model: Option>, pub network_code: Option, pub ip: Option, pub no_punch: Option, @@ -98,6 +99,18 @@ impl FileConfig { }) .collect() } + pub fn to_punch_model(&self) -> anyhow::Result> { + self.punch_model + .as_deref() + .unwrap_or_default() + .iter() + .map(|value| { + value + .parse::() + .map_err(|error| anyhow!("invalid punch model rule '{}': {}", value, error)) + }) + .collect() + } pub fn to_port_mapping(&self) -> anyhow::Result> { if let Some(port_mapping_raw) = &self.port_mapping { let mut port_mapping = Vec::with_capacity(port_mapping_raw.len()); @@ -126,6 +139,9 @@ pub struct Args { /// 指定目标 IP/网段的优先中转节点,可重复指定,格式为 target,turn_ip #[clap(long)] pub turn: Vec, + /// 指定目标 IP/网段允许的打洞方式,可重复指定,格式为 target,mode[,mode...] + #[clap(long)] + pub punch_model: Vec, /// 网络编号,相同编号的会组同一个局域网 #[clap(short, long)] pub network_code: Option, @@ -263,6 +279,11 @@ fn build_from_args_and_file(args: Args, file: FileConfig) -> anyhow::Result<(Con } else { args.turn }; + let punch_model = if args.punch_model.is_empty() { + file.to_punch_model()? + } else { + args.punch_model + }; let port_mapping = if args.port_mapping.is_empty() { file.to_port_mapping()? } else { @@ -317,6 +338,7 @@ fn build_from_args_and_file(args: Args, file: FileConfig) -> anyhow::Result<(Con server_addr, peer_address, turn, + punch_model, network_code, ip: args.ip.or(file.ip), no_punch: args.no_punch || file.no_punch.unwrap_or(false), @@ -366,6 +388,7 @@ fn build_from_args_only(args: Args) -> anyhow::Result<(Config, CtrlConfig)> { server_addr: args.server, peer_address: args.peer_address, turn: args.turn, + punch_model: args.punch_model, network_code: args .network_code .ok_or_else(|| anyhow!("network_code is required"))?, @@ -412,6 +435,7 @@ fn build_from_file_only(file: FileConfig) -> anyhow::Result<(Config, CtrlConfig) let server_addr = file.to_server_addr()?; let peer_address = file.to_peer_address()?; let turn = file.to_turn()?; + let punch_model = file.to_punch_model()?; let port_mapping = file.to_port_mapping()?; let cert_mode = file @@ -441,6 +465,7 @@ fn build_from_file_only(file: FileConfig) -> anyhow::Result<(Config, CtrlConfig) server_addr, peer_address, turn, + punch_model, network_code: file .network_code .ok_or_else(|| anyhow!("network_code is required"))?, @@ -510,6 +535,10 @@ server = ["quic://1.2.3.4:29872"] # 命中目标不参与 P2P 打洞 # turn = ["10.26.0.0/24,10.26.0.2", "10.26.1.9,10.26.0.3"] +# 指定目标虚拟 IP 或网段允许的 P2P 打洞方式,同一目标的多条规则会合并 +# 可选模式:IPv4Tcp、IPv4Udp、IPv6Tcp、IPv6Udp +# punch_model = ["10.26.0.2,IPv4Udp", "10.26.1.0/24,IPv4Tcp,IPv4Udp"] + # ===简单使用以下参数可以不动=== # 自定义虚拟 IP (可选) @@ -758,6 +787,37 @@ mod tests { assert_eq!(config.turn[0].to_string(), "10.26.0.0/16,10.26.0.2"); } + #[test] + fn test_punch_model_cli_and_file_precedence() { + let file: FileConfig = toml::from_str( + "punch_model = [\"10.26.0.0/16,IPv4Udp\"]\nnetwork_code = \"test-net\"\nserver = [\"quic://127.0.0.1:29872\"]", + ) + .unwrap(); + let args = Args::try_parse_from([ + "vnt", + "-s", + "quic://127.0.0.1:29872", + "-n", + "test-net", + "--punch-model", + "10.26.1.9,IPv4Tcp,IPv6Udp", + ]) + .unwrap(); + let (config, _) = build_config_from_args_and_file(Some(args), Some(file)).unwrap(); + assert_eq!(config.punch_model.len(), 1); + assert_eq!( + config.punch_model[0].to_string(), + "10.26.1.9,IPv4Tcp,IPv6Udp" + ); + + let file: FileConfig = toml::from_str( + "punch_model = [\"10.26.0.0/16,IPv4Udp\"]\nnetwork_code = \"test-net\"\nserver = [\"quic://127.0.0.1:29872\"]", + ) + .unwrap(); + let (config, _) = build_config_from_args_and_file(None, Some(file)).unwrap(); + assert_eq!(config.punch_model[0].to_string(), "10.26.0.0/16,IPv4Udp"); + } + #[test] fn test_subnet_mapping_cli_and_file_precedence() { let file: FileConfig = toml::from_str( diff --git a/vnt-core/proto/client.proto b/vnt-core/proto/client.proto index f5d95158..be20eac0 100644 --- a/vnt-core/proto/client.proto +++ b/vnt-core/proto/client.proto @@ -59,4 +59,6 @@ message NatInfo { message PunchInfo{ NatInfo nat_info = 1; -} \ No newline at end of file + // Bit 0..3: IPv4Tcp, IPv4Udp, IPv6Tcp, IPv6Udp. Zero means legacy/all. + uint32 punch_model = 2; +} diff --git a/vnt-core/src/context/config.rs b/vnt-core/src/context/config.rs index e6821205..4170e7f3 100644 --- a/vnt-core/src/context/config.rs +++ b/vnt-core/src/context/config.rs @@ -5,6 +5,7 @@ use crate::tls::verifier::CertValidationMode; use crate::tunnel_core::server::transport::config::{ConnectRegConfig, ProtocolAddress}; use anyhow::bail; use ipnet::Ipv4Net; +use rustp2p_core::punch::{PunchPolicy, PunchPolicySet}; use rustp2p_core::route_table::Protocol; use rustp2p_core::socket::LocalInterface; use std::collections::HashSet; @@ -18,6 +19,114 @@ pub const MAX_NAME_LEN: usize = 128; pub const MAX_VERSION_LEN: usize = 32; pub const MAX_MTU: u16 = 1500; +const PUNCH_POLICIES: [PunchPolicy; 4] = [ + PunchPolicy::IPv4Tcp, + PunchPolicy::IPv4Udp, + PunchPolicy::IPv6Tcp, + PunchPolicy::IPv6Udp, +]; + +#[derive(Debug, Clone)] +pub struct PunchRule { + target: Ipv4Net, + policies: PunchPolicySet, +} + +impl PunchRule { + pub fn target(&self) -> Ipv4Net { + self.target + } + + pub fn policies(&self) -> PunchPolicySet { + self.policies.clone() + } + + pub fn matches(&self, ip: &Ipv4Addr) -> bool { + self.target.contains(ip) + } + + fn merge(&mut self, other: &Self) { + for policy in PUNCH_POLICIES { + if other.policies.is_match(policy) { + self.policies.or(policy); + } + } + } +} + +impl FromStr for PunchRule { + type Err = anyhow::Error; + + fn from_str(value: &str) -> Result { + let parts = value.split(',').map(str::trim).collect::>(); + if parts.len() < 2 || parts.iter().any(|part| part.is_empty()) { + bail!("invalid punch model rule '{value}', expected target_ip_or_cidr,mode[,mode...]") + } + let target = if parts[0].contains('/') { + parts[0] + .parse::() + .map_err(|error| anyhow::anyhow!("invalid punch target '{}': {error}", parts[0]))? + } else { + let ip = parts[0] + .parse::() + .map_err(|error| anyhow::anyhow!("invalid punch target '{}': {error}", parts[0]))?; + Ipv4Net::new(ip, 32)? + }; + let mut policies = PunchPolicySet::empty(); + for value in &parts[1..] { + policies.or(parse_punch_policy(value)?); + } + Ok(Self { target, policies }) + } +} + +impl Display for PunchRule { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + if self.target.prefix_len() == 32 { + write!(f, "{}", self.target.addr())?; + } else { + write!(f, "{}", self.target)?; + } + for policy in PUNCH_POLICIES { + if self.policies.is_match(policy) { + write!(f, ",{}", punch_policy_name(policy))?; + } + } + Ok(()) + } +} + +fn parse_punch_policy(value: &str) -> anyhow::Result { + let normalized = value.trim().to_ascii_lowercase().replace(['-', '_'], ""); + match normalized.as_str() { + "ipv4tcp" => Ok(PunchPolicy::IPv4Tcp), + "ipv4udp" => Ok(PunchPolicy::IPv4Udp), + "ipv6tcp" => Ok(PunchPolicy::IPv6Tcp), + "ipv6udp" => Ok(PunchPolicy::IPv6Udp), + _ => bail!( + "invalid punch mode '{value}', expected one of: IPv4Tcp, IPv4Udp, IPv6Tcp, IPv6Udp" + ), + } +} + +fn punch_policy_name(policy: PunchPolicy) -> &'static str { + match policy { + PunchPolicy::IPv4Tcp => "IPv4Tcp", + PunchPolicy::IPv4Udp => "IPv4Udp", + PunchPolicy::IPv6Tcp => "IPv6Tcp", + PunchPolicy::IPv6Udp => "IPv6Udp", + } +} + +pub fn punch_model_for(rules: &[PunchRule], target: &Ipv4Addr) -> PunchPolicySet { + rules + .iter() + .filter(|rule| rule.matches(target)) + .max_by_key(|rule| rule.target.prefix_len()) + .map(PunchRule::policies) + .unwrap_or_else(PunchPolicySet::all) +} + #[derive(Debug, Clone, Eq, PartialEq, Hash)] pub struct TurnRule { target: Ipv4Net, @@ -231,6 +340,7 @@ pub struct Config { pub server_addr: Vec, pub peer_address: Vec, pub turn: Vec, + pub punch_model: Vec, pub cert_mode: CertValidationMode, pub network_code: String, pub device_id: String, @@ -267,6 +377,15 @@ impl Config { self.check_turn_rules()?; let mut seen = HashSet::new(); self.turn.retain(|rule| seen.insert(rule.clone())); + let mut merged = Vec::::new(); + for rule in self.punch_model.drain(..) { + if let Some(existing) = merged.iter_mut().find(|item| item.target == rule.target) { + existing.merge(&rule); + } else { + merged.push(rule); + } + } + self.punch_model = merged; crate::nat::subnet_mapping::normalize_and_validate(&mut self.subnet_mapping, &self.output)?; Ok(()) } @@ -435,6 +554,53 @@ mod tests { } } + #[test] + fn punch_rule_parses_aliases_and_rejects_invalid_values() { + let host = "10.26.0.2,ipv4-udp".parse::().unwrap(); + assert_eq!(host.target().prefix_len(), 32); + assert!(host.policies().is_match(PunchPolicy::IPv4Udp)); + assert_eq!(host.to_string(), "10.26.0.2,IPv4Udp"); + + let cidr = "10.26.1.0/24,IPv4Tcp,ipv6_udp" + .parse::() + .unwrap(); + assert!(cidr.policies().is_match(PunchPolicy::IPv4Tcp)); + assert!(cidr.policies().is_match(PunchPolicy::IPv6Udp)); + assert_eq!(cidr.to_string(), "10.26.1.0/24,IPv4Tcp,IPv6Udp"); + + for value in ["10.26.0.2", "10.26.0.2,", "bad,IPv4Udp", "10.26.0.2,quic"] { + assert!(value.parse::().is_err(), "{value} must fail"); + } + } + + #[test] + fn punch_rules_merge_and_use_longest_prefix() { + let mut config = Config { + punch_model: vec![ + "10.26.0.0/16,IPv4Tcp".parse().unwrap(), + "10.26.1.0/24,IPv4Udp".parse().unwrap(), + "10.26.1.0/24,IPv6Tcp,IPv4Udp".parse().unwrap(), + ], + ..Default::default() + }; + config.normalize().unwrap(); + assert_eq!(config.punch_model.len(), 2); + + let narrow = punch_model_for(&config.punch_model, &Ipv4Addr::new(10, 26, 1, 9)); + assert!(narrow.is_match(PunchPolicy::IPv4Udp)); + assert!(narrow.is_match(PunchPolicy::IPv6Tcp)); + assert!(!narrow.is_match(PunchPolicy::IPv4Tcp)); + + let broad = punch_model_for(&config.punch_model, &Ipv4Addr::new(10, 26, 2, 9)); + assert!(broad.is_match(PunchPolicy::IPv4Tcp)); + assert!(!broad.is_match(PunchPolicy::IPv4Udp)); + + let unmatched = punch_model_for(&config.punch_model, &Ipv4Addr::new(10, 27, 0, 1)); + for policy in PUNCH_POLICIES { + assert!(unmatched.is_match(policy)); + } + } + #[test] fn turn_rule_parses_ip_and_cidr_and_uses_longest_prefix() { let host = "10.26.1.9,10.26.0.3".parse::().unwrap(); diff --git a/vnt-core/src/context/mod.rs b/vnt-core/src/context/mod.rs index 1ac3ebe5..78fa8474 100644 --- a/vnt-core/src/context/mod.rs +++ b/vnt-core/src/context/mod.rs @@ -1,4 +1,4 @@ -use crate::context::config::Config; +use crate::context::config::{Config, punch_model_for}; use crate::context::nat::{MyNatInfo, PunchBackoff}; use crate::nat::SubnetExternalRoute; use crate::protocol::client_message::PunchInfo; @@ -955,9 +955,16 @@ impl AppState { } } impl AppState { - pub fn get_punch_info(&self) -> Option { + pub fn get_punch_info(&self, target: Ipv4Addr) -> Option { + let punch_model = self + .config + .lock() + .as_ref() + .map(|config| punch_model_for(&config.punch_model, &target)) + .unwrap_or_else(rustp2p_core::punch::PunchPolicySet::all); self.nat_info.get().map(|info| PunchInfo { nat_info: self.filter_ip(info), + punch_model, }) } pub fn get_nat_info(&self) -> Option { diff --git a/vnt-core/src/core/mod.rs b/vnt-core/src/core/mod.rs index c965fd3c..c441f023 100644 --- a/vnt-core/src/core/mod.rs +++ b/vnt-core/src/core/mod.rs @@ -49,6 +49,7 @@ struct RegistrationContext { enhanced_inbound: EnhancedInbound, fec_decoder: FecDecoder, turn: std::sync::Arc>, + punch_model: std::sync::Arc>, auto_sync_subnet: bool, allow_ikev2: bool, } @@ -87,6 +88,7 @@ impl NetworkManager { config.normalize()?; config.check()?; let turn = std::sync::Arc::new(config.turn.clone()); + let punch_model = std::sync::Arc::new(config.punch_model.clone()); let outbound_interface_name = config .outbound_interface .as_deref() @@ -161,6 +163,7 @@ impl NetworkManager { app_state.punch_backoff.clone(), puncher, packet_crypto.clone(), + punch_model.clone(), ); let subnet_external_route = app_state.subnet_route.clone(); subnet_external_route.set_route_table(config.input.clone()); @@ -291,6 +294,7 @@ impl NetworkManager { enhanced_inbound, fec_decoder, turn, + punch_model, auto_sync_subnet: config.auto_sync_subnet, allow_ikev2: config.allow_ikev2, }); @@ -401,6 +405,7 @@ impl NetworkManager { peer_map: app_state.peer_map.clone(), punch_backoff: app_state.punch_backoff.clone(), puncher: ctx.puncher.clone(), + punch_model: ctx.punch_model.clone(), packet_crypto: ctx.packet_crypto.clone(), packet_compression: ctx.packet_compression.clone(), enhanced_inbound: ctx.enhanced_inbound.clone(), diff --git a/vnt-core/src/protocol/client_message.rs b/vnt-core/src/protocol/client_message.rs index cd82938b..c6b1024c 100644 --- a/vnt-core/src/protocol/client_message.rs +++ b/vnt-core/src/protocol/client_message.rs @@ -5,6 +5,7 @@ mod proto { use anyhow::bail; use bytes::BytesMut; use prost::Message; +use rustp2p_core::punch::{PunchPolicy, PunchPolicySet}; use std::net::{Ipv4Addr, Ipv6Addr}; use crate::protocol::ProtoToBytesMut; @@ -79,6 +80,7 @@ pub fn decode_nat_info(msg: proto::NatInfo) -> anyhow::Result BytesMut { let message = proto::PunchInfo { nat_info: Some(encode_nat_info(&self.nat_info)), + punch_model: encode_punch_model(&self.punch_model), }; message.encode_bytes_mut() } } + +const IPV4_TCP: u32 = 1 << 0; +const IPV4_UDP: u32 = 1 << 1; +const IPV6_TCP: u32 = 1 << 2; +const IPV6_UDP: u32 = 1 << 3; + +fn encode_punch_model(model: &PunchPolicySet) -> u32 { + let mut bits = 0; + for (policy, bit) in [ + (PunchPolicy::IPv4Tcp, IPV4_TCP), + (PunchPolicy::IPv4Udp, IPV4_UDP), + (PunchPolicy::IPv6Tcp, IPV6_TCP), + (PunchPolicy::IPv6Udp, IPV6_UDP), + ] { + if model.is_match(policy) { + bits |= bit; + } + } + bits +} + +fn decode_punch_model(bits: u32) -> PunchPolicySet { + if bits == 0 { + return PunchPolicySet::all(); + } + let mut model = PunchPolicySet::empty(); + for (policy, bit) in [ + (PunchPolicy::IPv4Tcp, IPV4_TCP), + (PunchPolicy::IPv4Udp, IPV4_UDP), + (PunchPolicy::IPv6Tcp, IPV6_TCP), + (PunchPolicy::IPv6Udp, IPV6_UDP), + ] { + if bits & bit != 0 { + model.or(policy); + } + } + model +} + +#[cfg(test)] +mod tests { + use super::*; + + fn encoded_punch_info(punch_model: u32) -> Vec { + proto::PunchInfo { + nat_info: Some(proto::NatInfo { + nat_type: proto::NatType::Cone.into(), + public_ips: Vec::new(), + public_udp_ports: Vec::new(), + public_port_range: 0, + local_ipv4s: Vec::new(), + ipv6: None, + local_udp_ports: Vec::new(), + local_tcp_port: 0, + public_tcp_port: 0, + }), + punch_model, + } + .encode_to_vec() + } + + #[test] + fn punch_model_bits_round_trip() { + for bits in 1..=(IPV4_TCP | IPV4_UDP | IPV6_TCP | IPV6_UDP) { + let decoded = PunchInfo::from_slice(&encoded_punch_info(bits)).unwrap(); + assert_eq!(encode_punch_model(&decoded.punch_model), bits); + + let encoded = decoded.encode(); + let wire = proto::PunchInfo::decode(encoded.as_ref()).unwrap(); + assert_eq!(wire.punch_model, bits); + } + } + + #[test] + fn legacy_zero_means_all_and_unknown_bits_do_not_open_modes() { + let legacy = PunchInfo::from_slice(&encoded_punch_info(0)) + .unwrap() + .punch_model; + for policy in [ + PunchPolicy::IPv4Tcp, + PunchPolicy::IPv4Udp, + PunchPolicy::IPv6Tcp, + PunchPolicy::IPv6Udp, + ] { + assert!(legacy.is_match(policy)); + } + + let unknown_only = PunchInfo::from_slice(&encoded_punch_info(1 << 20)) + .unwrap() + .punch_model; + for policy in [ + PunchPolicy::IPv4Tcp, + PunchPolicy::IPv4Udp, + PunchPolicy::IPv6Tcp, + PunchPolicy::IPv6Udp, + ] { + assert!(!unknown_only.is_match(policy)); + } + } +} diff --git a/vnt-core/src/tunnel_core/p2p/transport/punch.rs b/vnt-core/src/tunnel_core/p2p/transport/punch.rs index c7f34d98..6315a0c2 100644 --- a/vnt-core/src/tunnel_core/p2p/transport/punch.rs +++ b/vnt-core/src/tunnel_core/p2p/transport/punch.rs @@ -1,4 +1,4 @@ -use crate::context::config::{TurnRule, allow_punch}; +use crate::context::config::{PunchRule, TurnRule, allow_punch, punch_model_for}; use crate::context::nat::PunchBackoff; use crate::context::{ServerInfoCollection, SharedNetworkAddr}; use crate::crypto::PacketCrypto; @@ -10,7 +10,7 @@ use crate::tunnel_core::server::outbound::ServerOutbound; use anyhow::bail; use log::error; use rand::seq::SliceRandom; -use rustp2p_core::punch::{PunchModel, Puncher}; +use rustp2p_core::punch::{PunchModel, PunchPolicy, PunchPolicySet, Puncher}; use std::net::Ipv4Addr; use std::sync::Arc; use std::time::Duration; @@ -45,7 +45,7 @@ pub struct PunchTaskContext { pub turn: Arc>, } -pub type PunchInfoGetter = std::sync::Arc Option + Send + Sync>; +pub type PunchInfoGetter = std::sync::Arc Option + Send + Sync>; pub async fn punch_task( tunnel_to_server: ServerOutbound, @@ -57,9 +57,6 @@ pub async fn punch_task( let Some(src_ip) = ctx.network.ip() else { continue; }; - let Some(punch_info) = (ctx.punch_info_getter)() else { - continue; - }; if !ctx.server_info.is_any_server_connected(None) { continue; } @@ -73,6 +70,9 @@ pub async fn punch_task( list.shuffle(&mut rand::rng()); list.truncate(5); for dest_ip in list { + let Some(punch_info) = (ctx.punch_info_getter)(dest_ip) else { + continue; + }; // should_punch() 只用于候选过滤;最终必须原子占用, // 防止并发调度或重复报文同时启动同一目标。 if !ctx.punch_backoff.try_begin(dest_ip) { @@ -103,6 +103,7 @@ pub struct NatPuncher { puncher: Option, packet_crypto: PacketCrypto, limiter: PunchLimiter, + punch_rules: Arc>, } impl NatPuncher { @@ -111,6 +112,7 @@ impl NatPuncher { punch_backoff: PunchBackoff, puncher: Option, packet_crypto: PacketCrypto, + punch_rules: Arc>, ) -> Self { Self { network, @@ -118,12 +120,17 @@ impl NatPuncher { puncher, packet_crypto, limiter: PunchLimiter::default(), + punch_rules, } } pub fn punch(&self, dest_ip: Ipv4Addr, punch_info: PunchInfo) -> anyhow::Result { let Some(puncher) = self.puncher.clone() else { return Ok(false); }; + let Some(punch_model) = self.effective_punch_model(dest_ip, &punch_info) else { + log::debug!("skip punch to {dest_ip}: punch model intersection is empty"); + return Ok(false); + }; let Some(permit) = self.limiter.try_acquire() else { log::debug!("skip punch to {dest_ip}: concurrent punch limit reached"); return Ok(false); @@ -135,6 +142,7 @@ impl NatPuncher { puncher, dest_ip, punch_info, + punch_model, Some(Duration::from_millis(50)), permit, )?; @@ -152,17 +160,29 @@ impl NatPuncher { let Some(puncher) = self.puncher.clone() else { return Ok(()); }; + let Some(punch_model) = self.effective_punch_model(dest_ip, &punch_info) else { + log::debug!("skip punch to {dest_ip}: punch model intersection is empty"); + return Ok(()); + }; let Some(permit) = self.limiter.try_acquire() else { log::debug!("skip punch to {dest_ip}: concurrent punch limit reached"); return Ok(()); }; - self.spawn_punch(puncher, dest_ip, punch_info, time, permit) + self.spawn_punch(puncher, dest_ip, punch_info, punch_model, time, permit) + } + fn effective_punch_model( + &self, + dest_ip: Ipv4Addr, + punch_info: &PunchInfo, + ) -> Option { + effective_punch_model(&self.punch_rules, dest_ip, punch_info.punch_model.clone()) } fn spawn_punch( &self, puncher: Puncher, dest_ip: Ipv4Addr, punch_info: PunchInfo, + punch_model: PunchModel, time: Option, permit: OwnedSemaphorePermit, ) -> anyhow::Result<()> { @@ -175,7 +195,16 @@ impl NatPuncher { if let Some(time) = time { tokio::time::sleep(time).await; } - if let Err(e) = punch_now(puncher, src_ip, dest_ip, punch_info, packet_crypto).await { + if let Err(e) = punch_now( + puncher, + src_ip, + dest_ip, + punch_info, + punch_model, + packet_crypto, + ) + .await + { log::warn!("punch send error {:?}", e); } }); @@ -187,6 +216,7 @@ async fn punch_now( src_ip: Ipv4Addr, dest_ip: Ipv4Addr, nat_info: PunchInfo, + punch_model: PunchModel, packet_crypto: PacketCrypto, ) -> anyhow::Result<()> { let mut packet = NetPacket::new(TransmissionBytes::zeroed_size( @@ -200,13 +230,30 @@ async fn punch_now( packet.set_payload(&crate::utils::time::now_ts_ms().to_be_bytes())?; packet_crypto.encrypt_in_place(&mut packet)?; let buf = packet.into_buffer().into_bytes().freeze(); - let punch_info = rustp2p_core::punch::PunchInfo::new(PunchModel::all(), nat_info.nat_info); + let punch_info = rustp2p_core::punch::PunchInfo::new(punch_model, nat_info.nat_info); puncher .punch_now(Some(buf.clone()), buf, punch_info) .await?; Ok(()) } +fn effective_punch_model( + rules: &[PunchRule], + dest_ip: Ipv4Addr, + peer_model: PunchPolicySet, +) -> Option { + let model = punch_model_for(rules, &dest_ip) & peer_model; + [ + PunchPolicy::IPv4Tcp, + PunchPolicy::IPv4Udp, + PunchPolicy::IPv6Tcp, + PunchPolicy::IPv6Udp, + ] + .into_iter() + .any(|policy| model.is_match(policy)) + .then_some(model) +} + #[cfg(test)] mod tests { use super::*; @@ -232,4 +279,22 @@ mod tests { "任务结束释放许可后应能立即开始下一轮" ); } + + #[test] + fn punch_model_uses_local_and_peer_intersection() { + let rules = vec!["10.26.0.2,IPv4Tcp,IPv6Udp".parse::().unwrap()]; + let mut peer = PunchPolicySet::empty(); + peer.or(PunchPolicy::IPv4Tcp); + peer.or(PunchPolicy::IPv4Udp); + + let effective = effective_punch_model(&rules, Ipv4Addr::new(10, 26, 0, 2), peer) + .expect("IPv4Tcp is allowed by both peers"); + assert!(effective.is_match(PunchPolicy::IPv4Tcp)); + assert!(!effective.is_match(PunchPolicy::IPv4Udp)); + assert!(!effective.is_match(PunchPolicy::IPv6Udp)); + + let mut incompatible = PunchPolicySet::empty(); + incompatible.or(PunchPolicy::IPv6Tcp); + assert!(effective_punch_model(&rules, Ipv4Addr::new(10, 26, 0, 2), incompatible).is_none()); + } } diff --git a/vnt-core/src/tunnel_core/p2p/transport/task.rs b/vnt-core/src/tunnel_core/p2p/transport/task.rs index e14fcdf0..971f3c73 100644 --- a/vnt-core/src/tunnel_core/p2p/transport/task.rs +++ b/vnt-core/src/tunnel_core/p2p/transport/task.rs @@ -82,7 +82,7 @@ pub async fn init_tunnel( network: app_state.network.clone(), server_info: app_state.server_info_collection.clone(), punch_backoff: app_state.punch_backoff.clone(), - punch_info_getter: Arc::new(move || app_state_for_punch.get_punch_info()), + punch_info_getter: Arc::new(move |target| app_state_for_punch.get_punch_info(target)), turn: config.turn.clone(), }; task_group.spawn(punch_task(tunnel_to_server, route_table.clone(), punch_ctx)); diff --git a/vnt-core/src/tunnel_core/server/connection_manager.rs b/vnt-core/src/tunnel_core/server/connection_manager.rs index fdecb304..2be595cc 100644 --- a/vnt-core/src/tunnel_core/server/connection_manager.rs +++ b/vnt-core/src/tunnel_core/server/connection_manager.rs @@ -1,5 +1,5 @@ use crate::compression::PacketCompression; -use crate::context::config::{Config, TurnRule}; +use crate::context::config::{Config, PunchRule, TurnRule}; use crate::context::nat::{MyNatInfo, PunchBackoff}; use crate::context::{AppState, NetworkRoute, PeerInfoMap, ServerInfoCollection}; use crate::crypto::PacketCrypto; @@ -34,6 +34,7 @@ pub struct InboundHandlerConfig { pub peer_map: PeerInfoMap, pub punch_backoff: PunchBackoff, pub puncher: NatPuncher, + pub punch_model: Arc>, pub packet_crypto: PacketCrypto, pub packet_compression: PacketCompression, pub enhanced_inbound: EnhancedInbound, diff --git a/vnt-core/src/tunnel_core/server/inbound.rs b/vnt-core/src/tunnel_core/server/inbound.rs index 8cd8696f..07181ae5 100644 --- a/vnt-core/src/tunnel_core/server/inbound.rs +++ b/vnt-core/src/tunnel_core/server/inbound.rs @@ -1,5 +1,5 @@ use crate::compression::PacketCompression; -use crate::context::config::{DeviceMode, TurnRule, allow_punch}; +use crate::context::config::{DeviceMode, PunchRule, TurnRule, allow_punch, punch_model_for}; use crate::context::nat::{MyNatInfo, PunchBackoff}; use crate::context::{ NetworkAddr, NetworkRoute, PeerInfoMap, ServerInfoCollection, SharedNetworkAddr, @@ -355,6 +355,7 @@ pub(crate) struct ServerTurnInboundHandler { peer_map: PeerInfoMap, punch_backoff: PunchBackoff, puncher: NatPuncher, + punch_model: Arc>, packet_crypto: PacketCrypto, packet_compression: PacketCompression, enhanced_inbound: EnhancedInbound, @@ -378,6 +379,7 @@ impl ServerTurnInboundHandler { peer_map: config.peer_map, punch_backoff: config.punch_backoff, puncher: config.puncher, + punch_model: config.punch_model, packet_crypto: config.packet_crypto, packet_compression: config.packet_compression, enhanced_inbound: config.enhanced_inbound, @@ -397,9 +399,10 @@ impl ServerTurnInboundHandler { info.local_ipv4s.retain(|ip| !self.network_contains(ip)); info } - fn get_punch_info(&self) -> Option { + fn get_punch_info(&self, target: Ipv4Addr) -> Option { self.nat_info.get().map(|info| PunchInfo { nat_info: self.filter_ip(info), + punch_model: punch_model_for(&self.punch_model, &target), }) } fn update_peer_nat_info(&self, ip: Ipv4Addr, nat_info: NatInfo) { @@ -635,7 +638,7 @@ impl ServerTurnInboundHandler { } // 对方发起打洞 let peer_punch_info = PunchInfo::from_slice(net_packet.payload())?; - let Some(self_punch_info) = self.get_punch_info() else { + let Some(self_punch_info) = self.get_punch_info(src) else { return Ok(()); }; log::info!( diff --git a/vnt-jni/java_example/com/vnt/VntConfig.java b/vnt-jni/java_example/com/vnt/VntConfig.java index de849c84..d02fc9a9 100644 --- a/vnt-jni/java_example/com/vnt/VntConfig.java +++ b/vnt-jni/java_example/com/vnt/VntConfig.java @@ -15,6 +15,7 @@ public class VntConfig { private final List servers; private final List peerAddresses; private final List turnRules; + private final List punchModelRules; private final List inputRoutes; private final List subnetMappings; private final List outputRoutes; @@ -45,6 +46,7 @@ private VntConfig(Builder builder) { this.servers = builder.servers; this.peerAddresses = builder.peerAddresses; this.turnRules = builder.turnRules; + this.punchModelRules = builder.punchModelRules; this.inputRoutes = builder.inputRoutes; this.subnetMappings = builder.subnetMappings; this.outputRoutes = builder.outputRoutes; @@ -101,6 +103,7 @@ String toJson() { } json.put("turn", turnArray); } + putStringArray(json, "punch_model", punchModelRules); putStringArray(json, "input", inputRoutes); putStringArray(json, "subnet_mapping", subnetMappings); @@ -170,6 +173,7 @@ public static class Builder { private List servers = new ArrayList<>(); private List peerAddresses = new ArrayList<>(); private List turnRules = new ArrayList<>(); + private List punchModelRules = new ArrayList<>(); private List inputRoutes = new ArrayList<>(); private List subnetMappings = new ArrayList<>(); private List outputRoutes = new ArrayList<>(); @@ -223,6 +227,15 @@ public Builder addTurnRule(String turnRule) { return this; } + /** + * 添加目标节点打洞方式规则(可选)。 + * @param rule 格式:目标虚拟IP或CIDR,IPv4Tcp[,IPv4Udp,IPv6Tcp,IPv6Udp] + */ + public Builder addPunchModelRule(String rule) { + this.punchModelRules.add(rule); + return this; + } + /** 添加点对网入口路由,格式:映射CIDR,出口节点虚拟IP。 */ public Builder addInputRoute(String route) { this.inputRoutes.add(route); diff --git a/vnt-jni/src/lib.rs b/vnt-jni/src/lib.rs index 676e18c8..f51966e0 100644 --- a/vnt-jni/src/lib.rs +++ b/vnt-jni/src/lib.rs @@ -12,7 +12,7 @@ use std::os::fd::{FromRawFd, OwnedFd}; use std::sync::Arc; use tokio::runtime::Runtime; use vnt_core::api::VntApi; -use vnt_core::context::config::{Config, DeviceMode, PeerAddress, TurnRule}; +use vnt_core::context::config::{Config, DeviceMode, PeerAddress, PunchRule, TurnRule}; use vnt_core::core::{NetworkManager, RegisterResponse}; use vnt_core::nat::{NetInput, SubnetMapping}; use vnt_core::port_mapping::PortMapping; @@ -1038,6 +1038,8 @@ fn parse_config_from_json(json_str: &str) -> anyhow::Result { peer_address: Vec, #[serde(default)] turn: Vec, + #[serde(default)] + punch_model: Vec, network_code: String, #[serde(default)] device_id: Option, @@ -1124,6 +1126,16 @@ fn parse_config_from_json(json_str: &str) -> anyhow::Result { }) .collect::>()?; + let punch_model: Vec = cfg + .punch_model + .iter() + .map(|value| { + value + .parse() + .map_err(|error| anyhow::anyhow!("invalid punch_model rule '{}': {}", value, error)) + }) + .collect::>()?; + let port_mapping: Vec = cfg .port_mapping .iter() @@ -1171,6 +1183,7 @@ fn parse_config_from_json(json_str: &str) -> anyhow::Result { server_addr: server_addrs, peer_address, turn, + punch_model, network_code: cfg.network_code, ip: cfg.ip, no_punch: cfg.no_punch, @@ -1491,6 +1504,24 @@ mod tests { assert_eq!(config.turn[1].to_string(), "10.26.1.9,10.26.0.3"); } + #[test] + fn parses_punch_model_rules_from_json() { + let config = parse_config_from_json( + r#"{ + "server":["quic://127.0.0.1:29872"], + "network_code":"test-net", + "punch_model":["10.26.0.2,IPv4Udp","10.26.1.0/24,IPv4Tcp,IPv6Udp"] + }"#, + ) + .unwrap(); + assert_eq!(config.punch_model.len(), 2); + assert_eq!(config.punch_model[0].to_string(), "10.26.0.2,IPv4Udp"); + assert_eq!( + config.punch_model[1].to_string(), + "10.26.1.0/24,IPv4Tcp,IPv6Udp" + ); + } + #[test] fn parses_exit_subnet_mapping_from_json() { let mut config = parse_config_from_json( diff --git a/vnt-web/src/service_http.rs b/vnt-web/src/service_http.rs index 308b2e02..f177fbc9 100644 --- a/vnt-web/src/service_http.rs +++ b/vnt-web/src/service_http.rs @@ -28,7 +28,9 @@ use tokio_util::sync::CancellationToken; use tower::ServiceExt; use tower_http::cors::{Any, CorsLayer}; use vnt_core::api::VntApi; -use vnt_core::context::config::{Config as CoreConfig, DeviceMode, PeerAddress, TurnRule}; +use vnt_core::context::config::{ + Config as CoreConfig, DeviceMode, PeerAddress, PunchRule, TurnRule, +}; use vnt_core::core::{DEFAULT_MTU, NetworkManager, RegisterResponse}; use vnt_core::nat::{NetInput, SubnetMapping}; use vnt_core::port_mapping::PortMapping; @@ -257,6 +259,8 @@ pub struct StartConfig { pub peer_address: Vec, #[serde(default)] pub turn: Vec, + #[serde(default)] + pub punch_model: Vec, pub cert_mode: Option, pub network_code: String, pub device_id: Option, @@ -1423,6 +1427,16 @@ fn convert_config(cfg: StartConfig) -> anyhow::Result { }) .collect::>()?; + let punch_model: Vec = cfg + .punch_model + .iter() + .map(|value| { + value + .parse() + .map_err(|error| anyhow!("invalid punch_model rule '{}': {}", value, error)) + }) + .collect::>()?; + let port_mapping: Vec = cfg .port_mapping .iter() @@ -1467,6 +1481,7 @@ fn convert_config(cfg: StartConfig) -> anyhow::Result { server_addr: server_addrs, peer_address, turn, + punch_model, network_code: cfg.network_code, ip: cfg.ip, no_punch: cfg.no_punch, @@ -1812,6 +1827,7 @@ mod tests { server: Vec::new(), peer_address: Vec::new(), turn: Vec::new(), + punch_model: Vec::new(), cert_mode: None, network_code: "test".to_string(), device_id: Some("device-a".to_string()), @@ -1890,6 +1906,22 @@ network_code = "test" assert_eq!(core.turn[1].to_string(), "10.26.1.9,10.26.0.3"); } + #[test] + fn test_convert_config_keeps_punch_model_rules() { + let mut config = new_test_config(); + config.punch_model = vec![ + "10.26.0.2,IPv4Udp".to_string(), + "10.26.1.0/24,IPv4Tcp,IPv6Udp".to_string(), + ]; + let core = convert_config(config).unwrap(); + assert_eq!(core.punch_model.len(), 2); + assert_eq!(core.punch_model[0].to_string(), "10.26.0.2,IPv4Udp"); + assert_eq!( + core.punch_model[1].to_string(), + "10.26.1.0/24,IPv4Tcp,IPv6Udp" + ); + } + #[test] fn test_convert_config_keeps_exit_subnet_mapping() { let mut config = new_test_config(); diff --git a/vnt-web/ui/src/utils/configHelp.js b/vnt-web/ui/src/utils/configHelp.js index 64a78bef..d5cd4fb9 100644 --- a/vnt-web/ui/src/utils/configHelp.js +++ b/vnt-web/ui/src/utils/configHelp.js @@ -38,6 +38,14 @@ export const configHelp = { format: "目标虚拟IP或CIDR,中转节点虚拟IP", example: "10.26.0.0/24,10.26.0.2", }, + punch_model: { + param: "punch_model", + summary: "按目标虚拟 IP 或网段限制允许使用的 P2P 打洞方式。", + usage: "双方会交换各自规则,实际只尝试双方允许集合的交集。重叠规则按最长前缀匹配,相同目标的多条规则会合并;未命中时允许全部方式。", + format: "目标虚拟IP或CIDR,IPv4Tcp|IPv4Udp|IPv6Tcp|IPv6Udp(可填写多种,以逗号分隔)", + example: "10.26.1.0/24,IPv4Tcp,IPv4Udp", + notes: ["模式名称不区分大小写,也支持 ipv4-tcp 等连字符写法。", "指定中转规则和禁用打洞仍具有更高优先级。"], + }, ip: { param: "ip", summary: "请求一个固定的本机虚拟 IPv4 地址。", diff --git a/vnt-web/ui/src/utils/toml.js b/vnt-web/ui/src/utils/toml.js index d1f5fef4..983c3032 100644 --- a/vnt-web/ui/src/utils/toml.js +++ b/vnt-web/ui/src/utils/toml.js @@ -6,6 +6,7 @@ export const emptyFormData = () => ({ server: [""], peer_address: [], turn: [], + punch_model: [], ip: "", mtu: null, rtx: false, @@ -68,6 +69,12 @@ export const parseTomlToForm = (toml) => { const items = match[1].match(/"([^"]*)"/g); if (items) data.turn = items.map((s) => s.replace(/"/g, "")); } + } else if (trimmed.match(/^punch_model\s*=/)) { + const match = trimmed.match(/punch_model\s*=\s*\[(.*)\]/); + if (match) { + const items = match[1].match(/"([^"]*)"/g); + if (items) data.punch_model = items.map((s) => s.replace(/"/g, "")); + } } else if (trimmed.includes("ip =")) { const match = trimmed.match(/ip\s*=\s*"([^"]*)"/); if (match) data.ip = match[1]; @@ -204,6 +211,12 @@ export const formToToml = (formData) => { toml += `turn = [${turnRules.map((s) => `"${s}"`).join(", ")}]\n`; } + const punchModelRules = formData.punch_model.filter((s) => s.trim()); + if (punchModelRules.length > 0) { + toml += "\n# 按目标虚拟 IP 或网段限制 P2P 打洞方式;双方实际使用允许集合的交集\n"; + toml += `punch_model = [${punchModelRules.map((s) => `"${s}"`).join(", ")}]\n`; + } + if (formData.ip) { toml += "\n# 自定义虚拟 IP (可选)\n"; toml += `ip = "${formData.ip}"\n`; @@ -366,6 +379,9 @@ server = ["quic://1.2.3.4:29872"] # 命中目标不参与 P2P 打洞 # turn = ["10.26.0.0/24,10.26.0.2", "10.26.1.9,10.26.0.3"] +# 按目标虚拟 IP 或网段限制 P2P 打洞方式;可选 IPv4Tcp、IPv4Udp、IPv6Tcp、IPv6Udp +# punch_model = ["10.26.0.2,IPv4Udp", "10.26.1.0/24,IPv4Tcp,IPv4Udp"] + # ===简单使用以下参数可以不动=== # 自定义虚拟 IP (可选) diff --git a/vnt-web/ui/src/views/ConfigEditor.vue b/vnt-web/ui/src/views/ConfigEditor.vue index d8b3fc09..2366cc9a 100644 --- a/vnt-web/ui/src/views/ConfigEditor.vue +++ b/vnt-web/ui/src/views/ConfigEditor.vue @@ -342,6 +342,40 @@ const sectionTitleClass = "text-md mb-4 flex items-center font-bold text-slate-9

命中目标不参与 P2P 打洞;中转节点已直连时优先经其转发,填写网关 IP 时强制走服务器。

+
+ +
+
+ + +
+ +
+

可选 IPv4Tcp、IPv4Udp、IPv6Tcp、IPv6Udp;双方实际使用允许集合的交集。

+
From 88c279a4ba8f40efd8f3373a6ce6faf5c3efdadd Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Fri, 4 Sep 2026 23:22:02 +0800 Subject: [PATCH 05/14] =?UTF-8?q?refactor(ui):=20=E9=87=8D=E7=BB=84?= =?UTF-8?q?=E9=85=8D=E7=BD=AE=E7=BC=96=E8=BE=91=E5=99=A8=E5=88=86=E7=BB=84?= =?UTF-8?q?=E5=B9=B6=E6=94=AF=E6=8C=81=E5=88=86=E5=8C=BA=E6=8A=98=E5=8F=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 基础配置只保留配置名称、网络编号、服务器地址;新增连接与打洞 分区(可直连节点、中转规则、打洞方式规则、隧道端口、关闭P2P打洞); 虚拟网卡模式移入网络设置;NAT与路由改名子网路由并移出网卡模式。 高级分区(子网路由、端口映射、STUN)默认折叠,已有数据时自动展开。 --- vnt-web/ui/src/views/ConfigEditor.vue | 1123 ++++++++++++++----------- 1 file changed, 635 insertions(+), 488 deletions(-) diff --git a/vnt-web/ui/src/views/ConfigEditor.vue b/vnt-web/ui/src/views/ConfigEditor.vue index 2366cc9a..34740522 100644 --- a/vnt-web/ui/src/views/ConfigEditor.vue +++ b/vnt-web/ui/src/views/ConfigEditor.vue @@ -39,6 +39,42 @@ const deviceModeOptions = [ ]; const isWindows = /Windows/i.test(globalThis.navigator?.userAgent || ""); +// 分区折叠状态:子网路由、端口映射、STUN 属于高级配置,默认折叠 +const DEFAULT_SECTIONS = { + basic: true, + connect: true, + network: true, + transport: true, + security: true, + subnet: false, + portmap: false, + device: true, + stun: false, +}; +const sectionExpanded = ref({ ...DEFAULT_SECTIONS }); +const toggleSection = (key) => { + sectionExpanded.value[key] = !sectionExpanded.value[key]; +}; +// 打开编辑器时重置折叠状态;已有数据的分区自动展开,避免配置被藏起来 +const resetSections = (data) => { + sectionExpanded.value = { ...DEFAULT_SECTIONS }; + if ( + data.input.length || + data.output.length || + data.subnet_mapping.length || + data.auto_sync_subnet || + data.no_nat + ) { + sectionExpanded.value.subnet = true; + } + if (data.port_mapping.length || data.allow_mapping) { + sectionExpanded.value.portmap = true; + } + if (data.udp_stun.length || data.tcp_stun.length) { + sectionExpanded.value.stun = true; + } +}; + // 打开时加载内容 watch( () => props.show, @@ -57,6 +93,7 @@ watch( originalToml.value = data; // 保存原始TOML isParsingToml.value = true; formData.value = parseTomlToForm(data); + resetSections(formData.value); nextTick(() => { isParsingToml.value = false; }); @@ -68,6 +105,7 @@ watch( // 新建配置,初始化表单 originalToml.value = ""; formData.value = emptyFormData(); + resetSections(formData.value); editorContent.value = NEW_CONFIG_TEMPLATE; } }, @@ -163,7 +201,9 @@ const removeBtnClass = "shrink-0 rounded-lg bg-red-50 px-3 py-2 text-red-500 transition-colors hover:bg-red-100 dark:bg-red-900/20 dark:text-red-500 dark:hover:bg-red-900/40"; const addBtnClass = "flex w-full items-center justify-center gap-1 rounded-lg bg-slate-100 px-3 py-2 text-sm text-slate-600 transition-colors hover:bg-slate-200 dark:bg-slate-700/50 dark:text-slate-300 dark:hover:bg-slate-700"; -const sectionTitleClass = "text-md mb-4 flex items-center font-bold text-slate-900 dark:text-white"; +const sectionTitleClass = "text-md flex items-center font-bold text-slate-900 dark:text-white"; +const sectionChevronClass = (expanded) => + `h-5 w-5 shrink-0 text-slate-400 transition-transform duration-200 ${expanded ? "rotate-180" : ""}`;