From ccd6e44c2bab0453e7bfab8b9b4fc152d06bb4a9 Mon Sep 17 00:00:00 2001 From: lbl8603 <49143209+lbl8603@users.noreply.github.com> Date: Thu, 9 May 2024 20:14:01 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BD=BF=E7=94=A8=E5=BC=82=E6=AD=A5tun?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- vnt/src/handle/tun_tap/mod.rs | 192 +------------------- vnt/src/handle/tun_tap/tun_handler.rs | 252 +++++++++++++++++--------- vnt/src/handle/tun_tap/unix.rs | 183 +++++++++++++++++++ vnt/src/handle/tun_tap/windows.rs | 124 +++++++++++++ 4 files changed, 486 insertions(+), 265 deletions(-) create mode 100644 vnt/src/handle/tun_tap/unix.rs create mode 100644 vnt/src/handle/tun_tap/windows.rs diff --git a/vnt/src/handle/tun_tap/mod.rs b/vnt/src/handle/tun_tap/mod.rs index 3697320..5d80406 100644 --- a/vnt/src/handle/tun_tap/mod.rs +++ b/vnt/src/handle/tun_tap/mod.rs @@ -1,187 +1,11 @@ -use std::io; -use std::net::Ipv4Addr; - -use parking_lot::Mutex; - -use packet::ip::ipv4::packet::IpV4Packet; -use packet::ip::ipv4::protocol::Protocol; - -use crate::channel::context::ChannelContext; -use crate::cipher::Cipher; -use crate::external_route::ExternalRoute; -use crate::handle::{check_dest, CurrentDeviceInfo, PeerDeviceInfo}; -#[cfg(feature = "ip_proxy")] -use crate::ip_proxy::{IpProxyMap, ProxyHandler}; -use crate::protocol; -use crate::protocol::body::ENCRYPTION_RESERVED; -use crate::protocol::ip_turn_packet::BroadcastPacket; -use crate::protocol::{ip_turn_packet, NetPacket, MAX_TTL}; - mod channel_group; pub mod tun_handler; -fn broadcast( - server_cipher: &Cipher, - sender: &ChannelContext, - net_packet: &mut NetPacket<&mut [u8]>, - current_device: &CurrentDeviceInfo, - device_list: &Mutex<(u16, Vec)>, -) -> io::Result<()> { - let list: Vec = device_list - .lock() - .1 - .iter() - .filter(|info| info.status.is_online()) - .map(|info| info.virtual_ip) - .collect(); - const MAX_COUNT: usize = 8; - let mut p2p_ips = Vec::with_capacity(8); - let mut relay_ips = Vec::with_capacity(8); - let mut overflow = false; - for (index, peer_ip) in list.into_iter().enumerate() { - if index > MAX_COUNT { - overflow = true; - break; - } - if let Some(route) = sender.route_table.route_one_p2p(&peer_ip) { - if sender - .send_by_key(net_packet.buffer(), route.route_key()) - .is_ok() - { - p2p_ips.push(peer_ip); - continue; - } - } - relay_ips.push(peer_ip); - } - if !overflow && relay_ips.is_empty() { - //全部p2p,不需要服务器中转 - return Ok(()); - } - - if p2p_ips.is_empty() { - //都没有p2p则直接由服务器转发 - if current_device.status.online() { - sender.send_default(net_packet.buffer(), current_device.connect_server)?; - } - return Ok(()); - } - if !overflow && relay_ips.len() == 2 { - // 如果转发的ip数不多就直接发 - for peer_ip in relay_ips { - //非直连的广播要改变目的地址,不然服务端收到了会再次广播 - net_packet.set_destination(peer_ip); - sender.send_ipv4_by_id( - net_packet.buffer(), - &peer_ip, - current_device.connect_server, - current_device.status.online(), - )?; - } - return Ok(()); - } - if current_device.status.offline() { - //离线的不再转发 - return Ok(()); - } - let buf = vec![0u8; 12 + 1 + p2p_ips.len() * 4 + net_packet.data_len() + ENCRYPTION_RESERVED]; - //剩余的发送到服务端,需要告知哪些已发送过 - let mut server_packet = NetPacket::new_encrypt(buf)?; - server_packet.set_default_version(); - server_packet.set_gateway_flag(true); - server_packet.first_set_ttl(MAX_TTL); - server_packet.set_source(net_packet.source()); - //使用对应的目的地址 - server_packet.set_destination(net_packet.destination()); - server_packet.set_protocol(protocol::Protocol::IpTurn); - server_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4Broadcast.into()); - - let mut broadcast = BroadcastPacket::unchecked(server_packet.payload_mut()); - broadcast.set_address(&p2p_ips)?; - broadcast.set_data(net_packet.buffer())?; - server_cipher.encrypt_ipv4(&mut server_packet)?; - sender.send_default(server_packet.buffer(), current_device.connect_server) -} - -/// 实现一个原地发送,必须保证是如下结构 -/// |12字节开头|ip报文|至少1024字节结尾| -/// -#[inline] -pub fn base_handle( - context: &ChannelContext, - buf: &mut [u8], - data_len: usize, //数据总长度=12+ip包长度 - current_device: CurrentDeviceInfo, - ip_route: &ExternalRoute, - #[cfg(feature = "ip_proxy")] proxy_map: &Option, - client_cipher: &Cipher, - server_cipher: &Cipher, - device_list: &Mutex<(u16, Vec)>, -) -> io::Result<()> { - let ipv4_packet = IpV4Packet::new(&buf[12..data_len])?; - let protocol = ipv4_packet.protocol(); - let src_ip = ipv4_packet.source_ip(); - let mut dest_ip = ipv4_packet.destination_ip(); - let mut net_packet = NetPacket::new0(data_len, buf)?; - net_packet.set_default_version(); - net_packet.set_protocol(protocol::Protocol::IpTurn); - net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into()); - net_packet.first_set_ttl(6); - net_packet.set_source(src_ip); - net_packet.set_destination(dest_ip); - if dest_ip == current_device.virtual_gateway { - // 发到网关的加密方式不一样,要单独处理 - if protocol == Protocol::Icmp { - net_packet.set_gateway_flag(true); - server_cipher.encrypt_ipv4(&mut net_packet)?; - context.send_default(net_packet.buffer(), current_device.connect_server)?; - } - return Ok(()); - } - if dest_ip.is_multicast() { - //当作广播处理 - dest_ip = Ipv4Addr::BROADCAST; - net_packet.set_destination(Ipv4Addr::BROADCAST); - } - if dest_ip.is_broadcast() || current_device.broadcast_ip == dest_ip { - // 广播 发送到直连目标 - client_cipher.encrypt_ipv4(&mut net_packet)?; - broadcast( - server_cipher, - context, - &mut net_packet, - ¤t_device, - device_list, - )?; - return Ok(()); - } - if !check_dest( - dest_ip, - current_device.virtual_netmask, - current_device.virtual_network, - ) { - if let Some(r_dest_ip) = ip_route.route(&dest_ip) { - //路由的目标不能是自己 - if r_dest_ip == src_ip { - return Ok(()); - } - //需要修改目的地址 - dest_ip = r_dest_ip; - net_packet.set_destination(r_dest_ip); - } else { - return Ok(()); - } - } - #[cfg(feature = "ip_proxy")] - if let Some(proxy_map) = proxy_map { - let mut ipv4_packet = IpV4Packet::new(net_packet.payload_mut())?; - proxy_map.send_handle(&mut ipv4_packet)?; - } - client_cipher.encrypt_ipv4(&mut net_packet)?; - context.send_ipv4_by_id( - net_packet.buffer(), - &dest_ip, - current_device.connect_server, - current_device.status.online(), - ) -} +#[cfg(unix)] +mod unix; +#[cfg(unix)] +pub(crate) use unix::*; +#[cfg(target_os = "windows")] +mod windows; +#[cfg(target_os = "windows")] +pub(crate) use windows::*; diff --git a/vnt/src/handle/tun_tap/tun_handler.rs b/vnt/src/handle/tun_tap/tun_handler.rs index bfa18b9..cda5e3c 100644 --- a/vnt/src/handle/tun_tap/tun_handler.rs +++ b/vnt/src/handle/tun_tap/tun_handler.rs @@ -1,3 +1,4 @@ +use std::net::Ipv4Addr; use std::sync::Arc; use std::{io, thread}; @@ -6,22 +7,27 @@ use parking_lot::Mutex; use packet::icmp::icmp::IcmpPacket; use packet::icmp::Kind; -use packet::ip::ipv4; use packet::ip::ipv4::packet::IpV4Packet; +use packet::ip::ipv4::protocol::Protocol; use tun::device::IFace; use tun::Device; use crate::channel::context::ChannelContext; use crate::cipher::Cipher; use crate::external_route::ExternalRoute; -use crate::handle::tun_tap::channel_group::{channel_group, GroupSyncSender}; -use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; +use crate::handle::tun_tap::channel_group::channel_group; +use crate::handle::{check_dest, CurrentDeviceInfo, PeerDeviceInfo}; #[cfg(feature = "ip_proxy")] use crate::ip_proxy::IpProxyMap; +use crate::ip_proxy::ProxyHandler; +use crate::protocol; +use crate::protocol::body::ENCRYPTION_RESERVED; +use crate::protocol::ip_turn_packet::BroadcastPacket; +use crate::protocol::{ip_turn_packet, NetPacket, MAX_TTL}; use crate::util::{SingleU64Adder, StopManager}; fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> io::Result<()> { - if ipv4_packet.protocol() == ipv4::protocol::Protocol::Icmp { + if ipv4_packet.protocol() == Protocol::Icmp { let mut icmp = IcmpPacket::new(ipv4_packet.payload_mut())?; if icmp.kind() == Kind::EchoRequest { icmp.set_kind(Kind::EchoReply); @@ -37,7 +43,7 @@ fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> io::R } /// 接收tun数据,并且转发到udp上 -fn handle( +pub(crate) fn handle( context: &ChannelContext, data: &mut [u8], len: usize, @@ -59,7 +65,7 @@ fn handle( if src_ip == dest_ip { return icmp(&device_writer, ipv4_packet); } - return crate::handle::tun_tap::base_handle( + return base_handle( context, data, len, @@ -86,24 +92,6 @@ pub fn start( mut up_counter: SingleU64Adder, device_list: Arc)>>, ) -> io::Result<()> { - #[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))] - let worker = { - #[cfg(target_os = "macos")] - let current_device = current_device.clone(); - let device = device.clone(); - stop_manager.add_listener("tun_device".into(), move || { - if let Err(e) = device.shutdown() { - log::warn!("{:?}", e); - } - #[cfg(target_os = "macos")] - { - let ip = current_device.load().virtual_ip; - if let Ok(udp) = std::net::UdpSocket::bind("0.0.0.0:0") { - let _ = udp.send_to(b"stop", format!("{:?}:1234", ip)); - } - } - })? - }; if parallel > 1 { let (sender, receivers) = channel_group::<(Vec, usize)>(parallel, 16); for (index, receiver) in receivers.into_iter().enumerate() { @@ -148,17 +136,20 @@ pub fn start( thread::Builder::new() .name("tunHandlerM".into()) .spawn(move || { - if let Err(e) = start_multi(stop_manager, device, sender, &mut up_counter) { + if let Err(e) = crate::handle::tun_tap::start_multi( + stop_manager, + device, + sender, + &mut up_counter, + ) { log::warn!("stop:{}", e); } - #[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))] - worker.stop_all(); })?; } else { thread::Builder::new() .name("tunHandlerS".into()) .spawn(move || { - if let Err(e) = start_simple( + if let Err(e) = crate::handle::tun_tap::start_simple( stop_manager, &context, device, @@ -173,74 +164,173 @@ pub fn start( ) { log::warn!("stop:{}", e); } - #[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))] - worker.stop_all(); })?; } Ok(()) } -fn start_simple( - stop_manager: StopManager, - context: &ChannelContext, - device: Arc, - current_device: Arc>, - ip_route: ExternalRoute, - #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, - client_cipher: Cipher, - server_cipher: Cipher, - up_counter: &mut SingleU64Adder, - device_list: Arc)>>, +fn broadcast( + server_cipher: &Cipher, + sender: &ChannelContext, + net_packet: &mut NetPacket<&mut [u8]>, + current_device: &CurrentDeviceInfo, + device_list: &Mutex<(u16, Vec)>, ) -> io::Result<()> { - let mut buf = [0; 1024 * 16]; - loop { - if stop_manager.is_stop() { - return Ok(()); + let list: Vec = device_list + .lock() + .1 + .iter() + .filter(|info| info.status.is_online()) + .map(|info| info.virtual_ip) + .collect(); + const MAX_COUNT: usize = 8; + let mut p2p_ips = Vec::with_capacity(8); + let mut relay_ips = Vec::with_capacity(8); + let mut overflow = false; + for (index, peer_ip) in list.into_iter().enumerate() { + if index > MAX_COUNT { + overflow = true; + break; } - let len = device.read(&mut buf[12..])? + 12; - //单线程的 - up_counter.add(len as u64); - #[cfg(any(target_os = "macos"))] - let mut buf = &mut buf[4..]; - // buf是重复利用的,需要重置头部 - buf[..12].fill(0); - match handle( - context, - &mut buf, - len, - &device, - current_device.load(), - &ip_route, - #[cfg(feature = "ip_proxy")] - &ip_proxy_map, - &client_cipher, - &server_cipher, - &device_list, - ) { - Ok(_) => {} - Err(e) => { - log::warn!("{:?}", e) + if let Some(route) = sender.route_table.route_one_p2p(&peer_ip) { + if sender + .send_by_key(net_packet.buffer(), route.route_key()) + .is_ok() + { + p2p_ips.push(peer_ip); + continue; } } + relay_ips.push(peer_ip); } + if !overflow && relay_ips.is_empty() { + //全部p2p,不需要服务器中转 + return Ok(()); + } + + if p2p_ips.is_empty() { + //都没有p2p则直接由服务器转发 + if current_device.status.online() { + sender.send_default(net_packet.buffer(), current_device.connect_server)?; + } + return Ok(()); + } + if !overflow && relay_ips.len() == 2 { + // 如果转发的ip数不多就直接发 + for peer_ip in relay_ips { + //非直连的广播要改变目的地址,不然服务端收到了会再次广播 + net_packet.set_destination(peer_ip); + sender.send_ipv4_by_id( + net_packet.buffer(), + &peer_ip, + current_device.connect_server, + current_device.status.online(), + )?; + } + return Ok(()); + } + if current_device.status.offline() { + //离线的不再转发 + return Ok(()); + } + let buf = vec![0u8; 12 + 1 + p2p_ips.len() * 4 + net_packet.data_len() + ENCRYPTION_RESERVED]; + //剩余的发送到服务端,需要告知哪些已发送过 + let mut server_packet = NetPacket::new_encrypt(buf)?; + server_packet.set_default_version(); + server_packet.set_gateway_flag(true); + server_packet.first_set_ttl(MAX_TTL); + server_packet.set_source(net_packet.source()); + //使用对应的目的地址 + server_packet.set_destination(net_packet.destination()); + server_packet.set_protocol(protocol::Protocol::IpTurn); + server_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4Broadcast.into()); + + let mut broadcast = BroadcastPacket::unchecked(server_packet.payload_mut()); + broadcast.set_address(&p2p_ips)?; + broadcast.set_data(net_packet.buffer())?; + server_cipher.encrypt_ipv4(&mut server_packet)?; + sender.send_default(server_packet.buffer(), current_device.connect_server) } -fn start_multi( - stop_manager: StopManager, - device: Arc, - mut group_sync_sender: GroupSyncSender<(Vec, usize)>, - up_counter: &mut SingleU64Adder, +/// 实现一个原地发送,必须保证是如下结构 +/// |12字节开头|ip报文|至少1024字节结尾| +/// +#[inline] +fn base_handle( + context: &ChannelContext, + buf: &mut [u8], + data_len: usize, //数据总长度=12+ip包长度 + current_device: CurrentDeviceInfo, + ip_route: &ExternalRoute, + #[cfg(feature = "ip_proxy")] proxy_map: &Option, + client_cipher: &Cipher, + server_cipher: &Cipher, + device_list: &Mutex<(u16, Vec)>, ) -> io::Result<()> { - loop { - if stop_manager.is_stop() { - return Ok(()); + let ipv4_packet = IpV4Packet::new(&buf[12..data_len])?; + let protocol = ipv4_packet.protocol(); + let src_ip = ipv4_packet.source_ip(); + let mut dest_ip = ipv4_packet.destination_ip(); + let mut net_packet = NetPacket::new0(data_len, buf)?; + net_packet.set_default_version(); + net_packet.set_protocol(protocol::Protocol::IpTurn); + net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into()); + net_packet.first_set_ttl(6); + net_packet.set_source(src_ip); + net_packet.set_destination(dest_ip); + if dest_ip == current_device.virtual_gateway { + // 发到网关的加密方式不一样,要单独处理 + if protocol == Protocol::Icmp { + net_packet.set_gateway_flag(true); + server_cipher.encrypt_ipv4(&mut net_packet)?; + context.send_default(net_packet.buffer(), current_device.connect_server)?; } - let mut buf = vec![0; 1024 * 16]; - let len = device.read(&mut buf[12..])? + 12; - //单线程的 - up_counter.add(len as u64); - if group_sync_sender.send((buf, len)).is_err() { + return Ok(()); + } + if dest_ip.is_multicast() { + //当作广播处理 + dest_ip = Ipv4Addr::BROADCAST; + net_packet.set_destination(Ipv4Addr::BROADCAST); + } + if dest_ip.is_broadcast() || current_device.broadcast_ip == dest_ip { + // 广播 发送到直连目标 + client_cipher.encrypt_ipv4(&mut net_packet)?; + broadcast( + server_cipher, + context, + &mut net_packet, + ¤t_device, + device_list, + )?; + return Ok(()); + } + if !check_dest( + dest_ip, + current_device.virtual_netmask, + current_device.virtual_network, + ) { + if let Some(r_dest_ip) = ip_route.route(&dest_ip) { + //路由的目标不能是自己 + if r_dest_ip == src_ip { + return Ok(()); + } + //需要修改目的地址 + dest_ip = r_dest_ip; + net_packet.set_destination(r_dest_ip); + } else { return Ok(()); } } + #[cfg(feature = "ip_proxy")] + if let Some(proxy_map) = proxy_map { + let mut ipv4_packet = IpV4Packet::new(net_packet.payload_mut())?; + proxy_map.send_handle(&mut ipv4_packet)?; + } + client_cipher.encrypt_ipv4(&mut net_packet)?; + context.send_ipv4_by_id( + net_packet.buffer(), + &dest_ip, + current_device.connect_server, + current_device.status.online(), + ) } diff --git a/vnt/src/handle/tun_tap/unix.rs b/vnt/src/handle/tun_tap/unix.rs new file mode 100644 index 0000000..0c52966 --- /dev/null +++ b/vnt/src/handle/tun_tap/unix.rs @@ -0,0 +1,183 @@ +use crate::channel::context::ChannelContext; +use crate::cipher::Cipher; +use crate::external_route::ExternalRoute; +use crate::handle::tun_tap::channel_group::GroupSyncSender; +use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; +use crate::ip_proxy::IpProxyMap; +use crate::util::{SingleU64Adder, StopManager}; +use crossbeam_utils::atomic::AtomicCell; +use mio::event::Source; +use mio::unix::SourceFd; +use mio::{Events, Interest, Poll, Token, Waker}; +use parking_lot::Mutex; +use std::io; +use std::os::fd::AsRawFd; +use std::sync::Arc; +use tun::Device; + +const STOP: Token = Token(0); +const FD: Token = Token(1); + +pub(crate) fn start_simple( + stop_manager: StopManager, + context: &ChannelContext, + device: Arc, + current_device: Arc>, + ip_route: ExternalRoute, + #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, + client_cipher: Cipher, + server_cipher: Cipher, + up_counter: &mut SingleU64Adder, + device_list: Arc)>>, +) -> io::Result<()> { + let poll = Poll::new()?; + let waker = Arc::new(Waker::new(poll.registry(), STOP)?); + let _waker = waker.clone(); + let worker = stop_manager.add_listener("tun_device".into(), move || { + let _ = waker.wake(); + })?; + if let Err(e) = start_simple0( + poll, + context, + device, + current_device, + ip_route, + #[cfg(feature = "ip_proxy")] + ip_proxy_map, + client_cipher, + server_cipher, + up_counter, + device_list, + ) { + log::error!("{:?}", e); + }; + worker.stop_all(); + drop(_waker); + Ok(()) +} + +fn start_simple0( + mut poll: Poll, + context: &ChannelContext, + device: Arc, + current_device: Arc>, + ip_route: ExternalRoute, + #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, + client_cipher: Cipher, + server_cipher: Cipher, + up_counter: &mut SingleU64Adder, + device_list: Arc)>>, +) -> io::Result<()> { + let mut buf = [0; 1024 * 16]; + let fd = device.as_tun_fd(); + fd.set_nonblock()?; + SourceFd(&fd.as_raw_fd()).register(poll.registry(), FD, Interest::READABLE)?; + let mut evnets = Events::with_capacity(4); + #[cfg(not(target_os = "macos"))] + let start = 12; + #[cfg(target_os = "macos")] + let start = 12 - 4; + loop { + poll.poll(&mut evnets, None)?; + for event in evnets.iter() { + if event.token() == STOP { + return Ok(()); + } + loop { + let len = match fd.read(&mut buf[start..]) { + Ok(len) => len + start, + Err(e) => { + if e.kind() == io::ErrorKind::WouldBlock { + break; + } + Err(e)? + } + }; + //单线程的 + up_counter.add(len as u64); + // buf是重复利用的,需要重置头部 + buf[..12].fill(0); + match crate::handle::tun_tap::tun_handler::handle( + context, + &mut buf, + len, + &device, + current_device.load(), + &ip_route, + #[cfg(feature = "ip_proxy")] + &ip_proxy_map, + &client_cipher, + &server_cipher, + &device_list, + ) { + Ok(_) => {} + Err(e) => { + log::warn!("{:?}", e) + } + } + } + } + } +} + +pub(crate) fn start_multi( + stop_manager: StopManager, + device: Arc, + group_sync_sender: GroupSyncSender<(Vec, usize)>, + up_counter: &mut SingleU64Adder, +) -> io::Result<()> { + let poll = Poll::new()?; + let waker = Arc::new(Waker::new(poll.registry(), STOP)?); + let _waker = waker.clone(); + let worker = stop_manager.add_listener("tun_device".into(), move || { + let _ = waker.wake(); + })?; + if let Err(e) = start_multi0(poll, device, group_sync_sender, up_counter) { + log::error!("{:?}", e); + }; + worker.stop_all(); + drop(_waker); + Ok(()) +} + +fn start_multi0( + mut poll: Poll, + device: Arc, + mut group_sync_sender: GroupSyncSender<(Vec, usize)>, + up_counter: &mut SingleU64Adder, +) -> io::Result<()> { + let fd = device.as_tun_fd(); + fd.set_nonblock()?; + SourceFd(&fd.as_raw_fd()).register(poll.registry(), FD, Interest::READABLE)?; + let mut evnets = Events::with_capacity(4); + let mut buf = vec![0; 1024 * 16]; + #[cfg(not(target_os = "macos"))] + let start = 12; + #[cfg(target_os = "macos")] + let start = 12 - 4; + loop { + poll.poll(&mut evnets, None)?; + for event in evnets.iter() { + if event.token() == STOP { + return Ok(()); + } + loop { + let len = match fd.read(&mut buf[start..]) { + Ok(len) => len + start, + Err(e) => { + if e.kind() == io::ErrorKind::WouldBlock { + break; + } + Err(e)? + } + }; + //单线程的 + up_counter.add(len as u64); + if group_sync_sender.send((buf, len)).is_err() { + return Ok(()); + } + buf = vec![0; 1024 * 16]; + } + } + } +} diff --git a/vnt/src/handle/tun_tap/windows.rs b/vnt/src/handle/tun_tap/windows.rs new file mode 100644 index 0000000..8546bff --- /dev/null +++ b/vnt/src/handle/tun_tap/windows.rs @@ -0,0 +1,124 @@ +use crate::channel::context::ChannelContext; +use crate::cipher::Cipher; +use crate::external_route::ExternalRoute; +use crate::handle::tun_tap::channel_group::GroupSyncSender; +use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; +use crate::ip_proxy::IpProxyMap; +use crate::util::{SingleU64Adder, StopManager}; +use crossbeam_utils::atomic::AtomicCell; +use parking_lot::Mutex; +use std::io; +use std::sync::Arc; +use tun::device::IFace; +use tun::Device; + +pub(crate) fn start_simple( + stop_manager: StopManager, + context: &ChannelContext, + device: Arc, + current_device: Arc>, + ip_route: ExternalRoute, + #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, + client_cipher: Cipher, + server_cipher: Cipher, + up_counter: &mut SingleU64Adder, + device_list: Arc)>>, +) -> io::Result<()> { + let worker = { + let device = device.clone(); + stop_manager.add_listener("tun_device".into(), move || { + if let Err(e) = device.shutdown() { + log::warn!("{:?}", e); + } + })? + }; + if let Err(e) = start_simple0( + context, + device, + current_device, + ip_route, + #[cfg(feature = "ip_proxy")] + ip_proxy_map, + client_cipher, + server_cipher, + up_counter, + device_list, + ) { + log::error!("{:?}", e); + } + worker.stop_all(); + Ok(()) +} +fn start_simple0( + context: &ChannelContext, + device: Arc, + current_device: Arc>, + ip_route: ExternalRoute, + #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, + client_cipher: Cipher, + server_cipher: Cipher, + up_counter: &mut SingleU64Adder, + device_list: Arc)>>, +) -> io::Result<()> { + let mut buf = [0; 1024 * 16]; + loop { + let len = device.read(&mut buf[12..])? + 12; + //单线程的 + up_counter.add(len as u64); + // buf是重复利用的,需要重置头部 + buf[..12].fill(0); + match crate::handle::tun_tap::tun_handler::handle( + context, + &mut buf, + len, + &device, + current_device.load(), + &ip_route, + #[cfg(feature = "ip_proxy")] + &ip_proxy_map, + &client_cipher, + &server_cipher, + &device_list, + ) { + Ok(_) => {} + Err(e) => { + log::warn!("tun/tap {:?}", e) + } + } + } +} +pub(crate) fn start_multi( + stop_manager: StopManager, + device: Arc, + group_sync_sender: GroupSyncSender<(Vec, usize)>, + up_counter: &mut SingleU64Adder, +) -> io::Result<()> { + let worker = { + let device = device.clone(); + stop_manager.add_listener("tun_device_multi".into(), move || { + if let Err(e) = device.shutdown() { + log::warn!("{:?}", e); + } + })? + }; + if let Err(e) = start_multi0(device, group_sync_sender, up_counter) { + log::error!("{:?}", e); + }; + worker.stop_all(); + Ok(()) +} +fn start_multi0( + device: Arc, + mut group_sync_sender: GroupSyncSender<(Vec, usize)>, + up_counter: &mut SingleU64Adder, +) -> io::Result<()> { + loop { + let mut buf = vec![0; 1024 * 16]; + let len = device.read(&mut buf[12..])? + 12; + //单线程的 + up_counter.add(len as u64); + if group_sync_sender.send((buf, len)).is_err() { + return Ok(()); + } + } +}