From 1df5a1cfb103141b57bd18ec6b9b6154beb48fc8 Mon Sep 17 00:00:00 2001 From: lbl8603 <49143209+lbl8603@users.noreply.github.com> Date: Sun, 30 Jun 2024 23:02:50 +0800 Subject: [PATCH] =?UTF-8?q?=E5=8E=BB=E9=99=A4=E4=B8=8D=E5=AE=89=E5=85=A8?= =?UTF-8?q?=E7=9A=84=E8=AE=A1=E6=95=B0=E5=99=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- vnt/src/core/conn.rs | 13 +-- vnt/src/handle/maintain/up_status.rs | 8 +- vnt/src/handle/tun_tap/tun_handler.rs | 6 +- vnt/src/handle/tun_tap/unix.rs | 6 +- vnt/src/handle/tun_tap/windows.rs | 6 +- vnt/src/tun_tap_device/tun_create_helper.rs | 6 +- vnt/src/util/counter/adder.rs | 122 ++------------------ 7 files changed, 29 insertions(+), 138 deletions(-) diff --git a/vnt/src/core/conn.rs b/vnt/src/core/conn.rs index 97b03fa..7ccdd87 100644 --- a/vnt/src/core/conn.rs +++ b/vnt/src/core/conn.rs @@ -27,9 +27,7 @@ use crate::nat::NatTest; #[cfg(feature = "integrated_tun")] use crate::tun_tap_device::tun_create_helper::{DeviceAdapter, TunDeviceHelper}; use crate::tun_tap_device::vnt_device::DeviceWrite; -use crate::util::{ - Scheduler, SingleU64Adder, StopManager, U64Adder, WatchSingleU64Adder, WatchU64Adder, -}; +use crate::util::{Scheduler, StopManager, U64Adder, WatchU64Adder}; use crate::{nat, VntCallback}; #[derive(Clone)] @@ -68,7 +66,7 @@ pub struct VntInner { context: Arc>>, peer_nat_info_map: Arc>>, down_count_watcher: WatchU64Adder, - up_count_watcher: WatchSingleU64Adder, + up_count_watcher: WatchU64Adder, client_secret_hash: Option<[u8; 16]>, compressor: Compressor, client_cipher: Cipher, @@ -201,14 +199,13 @@ impl VntInner { let (punch_sender, punch_receiver) = maintain::punch_channel(); let peer_nat_info_map: Arc>> = Arc::new(RwLock::new(HashMap::with_capacity(16))); - let down_counter = - U64Adder::with_capacity(config.ports.as_ref().map(|v| v.len()).unwrap_or_default() + 8); + let down_counter = U64Adder::default(); let down_count_watcher = down_counter.watch(); let handshake = Handshake::new( #[cfg(feature = "server_encrypt")] rsa_cipher.clone(), ); - let up_counter = SingleU64Adder::new(); + let up_counter = U64Adder::default(); let up_count_watcher = up_counter.watch(); #[cfg(feature = "integrated_tun")] let tun_device_helper = { @@ -348,7 +345,7 @@ pub fn start( punch: Punch, callback: Call, down_count_watcher: WatchU64Adder, - up_count_watcher: WatchSingleU64Adder, + up_count_watcher: WatchU64Adder, ) { // 定时心跳 maintain::heartbeat( diff --git a/vnt/src/handle/maintain/up_status.rs b/vnt/src/handle/maintain/up_status.rs index 74d554b..1f68ade 100644 --- a/vnt/src/handle/maintain/up_status.rs +++ b/vnt/src/handle/maintain/up_status.rs @@ -3,7 +3,7 @@ use crate::handle::CurrentDeviceInfo; use crate::proto::message::{ClientStatusInfo, PunchNatType, RouteItem}; use crate::protocol::body::ENCRYPTION_RESERVED; use crate::protocol::{service_packet, NetPacket, Protocol, HEAD_LEN, MAX_TTL}; -use crate::util::{Scheduler, WatchSingleU64Adder, WatchU64Adder}; +use crate::util::{Scheduler, WatchU64Adder}; use crossbeam_utils::atomic::AtomicCell; use protobuf::Message; use std::io; @@ -16,7 +16,7 @@ pub fn up_status( context: ChannelContext, current_device_info: Arc>, down_count_watcher: WatchU64Adder, - up_count_watcher: WatchSingleU64Adder, + up_count_watcher: WatchU64Adder, ) { let _ = scheduler.timeout(Duration::from_secs(60), move |x| { up_status0( @@ -34,7 +34,7 @@ fn up_status0( context: ChannelContext, current_device_info: Arc>, down_count_watcher: WatchU64Adder, - up_count_watcher: WatchSingleU64Adder, + up_count_watcher: WatchU64Adder, ) { if let Err(e) = send_up_status_packet( &context, @@ -62,7 +62,7 @@ fn send_up_status_packet( context: &ChannelContext, current_device_info: &AtomicCell, down_count_watcher: &WatchU64Adder, - up_count_watcher: &WatchSingleU64Adder, + up_count_watcher: &WatchU64Adder, ) -> io::Result<()> { let device_info = current_device_info.load(); if device_info.status.offline() { diff --git a/vnt/src/handle/tun_tap/tun_handler.rs b/vnt/src/handle/tun_tap/tun_handler.rs index d8cacfe..52bf5bc 100644 --- a/vnt/src/handle/tun_tap/tun_handler.rs +++ b/vnt/src/handle/tun_tap/tun_handler.rs @@ -26,7 +26,7 @@ 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}; +use crate::util::{StopManager, U64Adder}; fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> anyhow::Result<()> { if ipv4_packet.protocol() == Protocol::Icmp { @@ -53,7 +53,7 @@ pub fn start( #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - mut up_counter: SingleU64Adder, + up_counter: U64Adder, device_list: Arc)>>, compressor: Compressor, device_stop: DeviceStop, @@ -71,7 +71,7 @@ pub fn start( ip_proxy_map, client_cipher, server_cipher, - &mut up_counter, + &up_counter, device_list, compressor, device_stop, diff --git a/vnt/src/handle/tun_tap/unix.rs b/vnt/src/handle/tun_tap/unix.rs index 5c3b1f3..a71a23a 100644 --- a/vnt/src/handle/tun_tap/unix.rs +++ b/vnt/src/handle/tun_tap/unix.rs @@ -7,7 +7,7 @@ use crate::handle::tun_tap::DeviceStop; use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; #[cfg(feature = "ip_proxy")] use crate::ip_proxy::IpProxyMap; -use crate::util::{SingleU64Adder, StopManager}; +use crate::util::{StopManager, U64Adder}; use crossbeam_utils::atomic::AtomicCell; use mio::event::Source; use mio::unix::SourceFd; @@ -30,7 +30,7 @@ pub(crate) fn start_simple( #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: &mut SingleU64Adder, + up_counter: &U64Adder, device_list: Arc)>>, compressor: Compressor, device_stop: DeviceStop, @@ -85,7 +85,7 @@ fn start_simple0( #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: &mut SingleU64Adder, + up_counter: &U64Adder, device_list: Arc)>>, compressor: Compressor, ) -> anyhow::Result<()> { diff --git a/vnt/src/handle/tun_tap/windows.rs b/vnt/src/handle/tun_tap/windows.rs index 9b2d5e7..c462d1f 100644 --- a/vnt/src/handle/tun_tap/windows.rs +++ b/vnt/src/handle/tun_tap/windows.rs @@ -7,7 +7,7 @@ use crate::handle::tun_tap::DeviceStop; use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; #[cfg(feature = "ip_proxy")] use crate::ip_proxy::IpProxyMap; -use crate::util::{SingleU64Adder, StopManager}; +use crate::util::{StopManager, U64Adder}; use crossbeam_utils::atomic::AtomicCell; use parking_lot::Mutex; use std::sync::Arc; @@ -23,7 +23,7 @@ pub(crate) fn start_simple( #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: &mut SingleU64Adder, + up_counter: &U64Adder, device_list: Arc)>>, compressor: Compressor, device_stop: DeviceStop, @@ -76,7 +76,7 @@ fn start_simple0( #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: &mut SingleU64Adder, + up_counter: &U64Adder, device_list: Arc)>>, compressor: Compressor, ) -> anyhow::Result<()> { diff --git a/vnt/src/tun_tap_device/tun_create_helper.rs b/vnt/src/tun_tap_device/tun_create_helper.rs index 99dd74f..09a1a97 100644 --- a/vnt/src/tun_tap_device/tun_create_helper.rs +++ b/vnt/src/tun_tap_device/tun_create_helper.rs @@ -16,7 +16,7 @@ use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; #[cfg(feature = "ip_proxy")] use crate::ip_proxy::IpProxyMap; use crate::tun_tap_device::vnt_device::DeviceWrite; -use crate::util::{SingleU64Adder, StopManager}; +use crate::util::{StopManager, U64Adder}; #[repr(transparent)] #[derive(Clone, Default)] @@ -67,7 +67,7 @@ struct TunDeviceHelperInner { ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: SingleU64Adder, + up_counter: U64Adder, device_list: Arc)>>, compressor: Compressor, } @@ -81,7 +81,7 @@ impl TunDeviceHelper { #[cfg(feature = "ip_proxy")] ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, - up_counter: SingleU64Adder, + up_counter: U64Adder, device_list: Arc)>>, compressor: Compressor, device_adapter: DeviceAdapter, diff --git a/vnt/src/util/counter/adder.rs b/vnt/src/util/counter/adder.rs index f8268ea..4888074 100644 --- a/vnt/src/util/counter/adder.rs +++ b/vnt/src/util/counter/adder.rs @@ -1,139 +1,33 @@ -use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; -/// 不安全的并发计数器,谨慎使用 +use crossbeam_utils::atomic::AtomicCell; +#[derive(Clone, Default)] pub struct U64Adder { - global_index: Arc, - inner: Arc, - index: usize, -} -#[derive(Clone)] -pub struct SingleU64Adder { - inner: Arc, -} -impl SingleU64Adder { - pub fn new() -> Self { - Self { - inner: Arc::new(SingleU64AdderInner::new()), - } - } - pub fn add(&mut self, num: u64) { - self.inner.add(num); - } - pub fn get(&self) -> u64 { - self.inner.get() - } - pub fn watch(&self) -> WatchSingleU64Adder { - WatchSingleU64Adder { - inner: self.inner.clone(), - } - } -} - -struct SingleU64AdderInner { - ptr: *mut u64, -} - -impl SingleU64AdderInner { - fn new() -> Self { - Self { - ptr: Box::into_raw(Box::new(0)), - } - } - #[inline(always)] - fn add(&self, num: u64) { - unsafe { *self.ptr += num } - } - - fn get(&self) -> u64 { - unsafe { *self.ptr } - } -} -impl Drop for SingleU64AdderInner { - fn drop(&mut self) { - unsafe { - let _ = Box::from_raw(self.ptr); - } - } -} - -unsafe impl Send for SingleU64AdderInner {} - -unsafe impl Sync for SingleU64AdderInner {} - -struct U64AdderInner { - base: Vec, -} - -impl U64AdderInner { - pub fn get(&self) -> u64 { - let mut count = 0; - for counter in self.base.iter() { - count += counter.get() - } - count - } + count: Arc>, } impl U64Adder { - /// 计数槽容量 - pub fn with_capacity(capacity: usize) -> Self { - let mut base = Vec::with_capacity(capacity); - for _ in 0..capacity { - base.push(SingleU64AdderInner::new()) - } - let inner = Arc::new(U64AdderInner { base }); - U64Adder { - global_index: Arc::new(AtomicUsize::new(1)), - inner, - index: 0, - } - } pub fn add(&self, num: u64) { - self.inner.base[self.index].add(num); + self.count.fetch_add(num); } pub fn get(&self) -> u64 { - self.inner.get() + self.count.load() } pub fn watch(&self) -> WatchU64Adder { WatchU64Adder { - inner: self.inner.clone(), - } - } -} - -impl Clone for U64Adder { - fn clone(&self) -> Self { - let index = self.global_index.fetch_add(1, Ordering::AcqRel); - if index > self.inner.base.len() { - panic!() - } - - Self { - global_index: self.global_index.clone(), - inner: self.inner.clone(), - index, + count: self.count.clone(), } } } #[derive(Clone)] pub struct WatchU64Adder { - inner: Arc, + count: Arc>, } impl WatchU64Adder { pub fn get(&self) -> u64 { - self.inner.get() - } -} -#[derive(Clone)] -pub struct WatchSingleU64Adder { - inner: Arc, -} -impl WatchSingleU64Adder { - pub fn get(&self) -> u64 { - self.inner.get() + self.count.load() } }