Compare commits

...
10 Commits
Author SHA1 Message Date
lubeilin 26d68ac059 版本改为1.2.7 2023-10-31 20:49:25 +08:00
lubeilin 6292c1c381 去除溢出检查 2023-10-31 20:49:16 +08:00
lubeilin b0c3f25a29 使用读写锁简化操作,增加延迟优先选项 2023-10-31 20:48:38 +08:00
lubeilin d2e09d3da5 增加配置文件的参数说明 2023-10-11 20:46:39 +08:00
lbl8603 293c5b90a4 Merge pull request #22 from taotieren/contrib
Update README.md
2023-10-11 09:19:37 +08:00
taotieren e2323361f9 Update README.md 2023-10-10 23:23:25 +08:00
lbl8603 d05cff99ee Merge pull request #21 from taotieren/aur
Update README.md
2023-10-10 22:52:30 +08:00
taotieren e9b1b2ef3b Update README.md 2023-10-10 22:31:38 +08:00
lbl8603 a8ea2c14fc Merge pull request #20 from taotieren/aur
Add AUR vnt-git
2023-10-10 22:22:10 +08:00
taotieren 0688cb4515 Add AUR vnt-git 2023-10-10 22:18:37 +08:00
17 changed files with 142 additions and 277 deletions
-1
View File
@@ -6,7 +6,6 @@ opt-level = 'z'
debug = 0 debug = 0
debug-assertions = false debug-assertions = false
strip= "debuginfo" strip= "debuginfo"
overflow-checks = true
lto = true lto = true
panic = 'abort' panic = 'abort'
incremental = false incremental = false
+39 -1
View File
@@ -111,6 +111,37 @@ sudo iptables -t nat -A POSTROUTING -o eth0 -s 10.26.0.0/24 -j MASQUERADE
# 查看设置 # 查看设置
iptables -vnL -t nat iptables -vnL -t nat
``` ```
### Arch Linux
[![Packaging status](https://repology.org/badge/vertical-allrepos/vnt.svg)](https://repology.org/project/vnt/versions)
- 通过 AUR 安装 [vnt-git](https://aur.archlinux.org/packages/vnt-git)
```bash
yay -Syu vnt
```
- 通过 `systemd` 设置开机自启及配置
```bash
sudo systemctl enable --now vnt-cli@
sudo systemctl status vnt-cli@
```
- 启用内置 `IPv4` 转发规则
```bash
sudo sysctl --system
```
- 通过内置防火墙文件配置防火墙转发规则
```bash
sudo cat /etc/vnt/iptables-vnt.rules >> /etc/iptables/iptables.rules
sudo iptables-restore iptables.rules
```
### macos ### macos
```shell ```shell
# 开启ip转发 # 开启ip转发
@@ -127,6 +158,7 @@ sudo pfctl -f /etc/pf.conf -e
- Mac - Mac
- Linux - Linux
- Arch Linux `yay -Syu vnt`
- Windows - Windows
- 默认使用tun网卡 依赖wintun.dll([win-tun](https://www.wintun.net/))(将dll放到同目录下,建议使用版本0.14.1) - 默认使用tun网卡 依赖wintun.dll([win-tun](https://www.wintun.net/))(将dll放到同目录下,建议使用版本0.14.1)
- 使用tap网卡 依赖tap-windows([win-tap](https://build.openvpn.net/downloads/releases/))(建议使用版本9.24.7) - 使用tap网卡 依赖tap-windows([win-tap](https://build.openvpn.net/downloads/releases/))(建议使用版本9.24.7)
@@ -226,10 +258,16 @@ vnt默认使用10.26.0.0/24网段,和本地网络适配器的ip冲突
### 交流群 ### 交流群
QQ:1034868233 QQ: 1034868233
### 其他 ### 其他
可使用社区小伙伴搭建的中继服务器 可使用社区小伙伴搭建的中继服务器
1. -s vnt.8443.eu.org:29871 1. -s vnt.8443.eu.org:29871
### 参与贡献
<a href="https://github.com/lbl8603/vnt/graphs/contributors">
<img src="https://contrib.rocks/image?repo=lbl8603/vnt" />
</a>
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "common" name = "common"
version = "1.2.6" version = "1.2.7"
edition = "2021" edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "vnt-cli" name = "vnt-cli"
version = "1.2.6" version = "1.2.7"
edition = "2021" edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+5 -1
View File
@@ -82,6 +82,8 @@
取值0~65535,指定本地监听的端口,默认取随机端口 取值0~65535,指定本地监听的端口,默认取随机端口
### --cmd ### --cmd
开启交互式命令,开启后可以直接在窗口下输入命令,如需后台运行请勿开启 开启交互式命令,开启后可以直接在窗口下输入命令,如需后台运行请勿开启
### --first_latency
优先使用低延迟通道,默认情况下优先使用p2p通道,某些情况下可能p2p比客户端中继延迟更高,可使用此参数进行优化传输
### --no-proxy ### --no-proxy
关闭内置的ip代理,内置的代理较为简单,而且一般来说直接使用网卡NAT转发性能会更高, 关闭内置的ip代理,内置的代理较为简单,而且一般来说直接使用网卡NAT转发性能会更高,
有需要可以自行配置NAT转发,[可参考‘编译’小节中的NAT配置](https://github.com/lbl8603/vnt#%E7%BC%96%E8%AF%91) 有需要可以自行配置NAT转发,[可参考‘编译’小节中的NAT配置](https://github.com/lbl8603/vnt#%E7%BC%96%E8%AF%91)
@@ -112,9 +114,11 @@ server_encrypt: true #服务端加密
parallel: 1 #任务并行度 parallel: 1 #任务并行度
cipher_model: aes_gcm #客户端加密算法 cipher_model: aes_gcm #客户端加密算法
finger: false #关闭数据指纹 finger: false #关闭数据指纹
punch_model: ipv4 #打洞模式 punch_model: ipv4 #打洞模式
port: 0 #使用随机端口 port: 0 #使用随机端口
cmd: false #关闭控制台输入 cmd: false #关闭控制台输入
no_proxy: false #是否关闭内置代理,true为关闭
first_latency: false #是否优先低延迟通道,默认为false,表示优先使用p2p通道
``` ```
或者需要哪个配置就加哪个,当然token是必须的 或者需要哪个配置就加哪个,当然token是必须的
+3
View File
@@ -33,6 +33,7 @@ pub struct FileConfig {
pub punch_model: String, pub punch_model: String,
pub port: u16, pub port: u16,
pub cmd: bool, pub cmd: bool,
pub first_latency: bool,
} }
impl Default for FileConfig { impl Default for FileConfig {
@@ -64,6 +65,7 @@ impl Default for FileConfig {
punch_model: "".to_string(), punch_model: "".to_string(),
port: 0, port: 0,
cmd: false, cmd: false,
first_latency: false,
} }
} }
} }
@@ -155,6 +157,7 @@ pub fn read_config(file_path: &str) -> io::Result<(Config, bool)> {
file_conf.finger, file_conf.finger,
punch_model, punch_model,
file_conf.port, file_conf.port,
file_conf.first_latency,
); );
Ok((config, file_conf.cmd)) Ok((config, file_conf.cmd))
} }
+4
View File
@@ -59,6 +59,7 @@ fn main() {
opts.optopt("", "port", "监听的端口", "<port>"); opts.optopt("", "port", "监听的端口", "<port>");
opts.optflag("", "cmd", "开启窗口输入"); opts.optflag("", "cmd", "开启窗口输入");
opts.optflag("", "no-proxy", "关闭内置代理"); opts.optflag("", "no-proxy", "关闭内置代理");
opts.optflag("", "first-latency", "优先延迟");
opts.optopt("f", "", "配置文件", "<conf>"); opts.optopt("f", "", "配置文件", "<conf>");
//"后台运行时,查看其他设备列表" //"后台运行时,查看其他设备列表"
opts.optflag("", "list", "后台运行时,查看其他设备列表"); opts.optflag("", "list", "后台运行时,查看其他设备列表");
@@ -261,6 +262,7 @@ fn main() {
let cmd = matches.opt_present("cmd"); let cmd = matches.opt_present("cmd");
#[cfg(feature = "ip_proxy")] #[cfg(feature = "ip_proxy")]
let no_proxy = matches.opt_present("no-proxy"); let no_proxy = matches.opt_present("no-proxy");
let first_latency = matches.opt_present("first-latency");
let config = Config::new( let config = Config::new(
tap, tap,
token, token,
@@ -285,6 +287,7 @@ fn main() {
finger, finger,
punch_model, punch_model,
port, port,
first_latency,
); );
(config, cmd) (config, cmd)
}; };
@@ -541,6 +544,7 @@ fn print_usage(program: &str, _opts: Options) {
println!(" --cmd 开启交互式命令,使用此参数开启控制台输入"); println!(" --cmd 开启交互式命令,使用此参数开启控制台输入");
#[cfg(feature = "ip_proxy")] #[cfg(feature = "ip_proxy")]
println!(" --no-proxy 关闭内置代理,如需点对网则需要配置网卡NAT转发"); println!(" --no-proxy 关闭内置代理,如需点对网则需要配置网卡NAT转发");
println!(" --first-latency 优先低延迟的通道,默认情况优先使用p2p通道");
println!(); println!();
println!( println!(
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "vnt-jni" name = "vnt-jni"
version = "1.2.6" version = "1.2.7"
edition = "2021" edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+2
View File
@@ -66,6 +66,7 @@ fn new_sync(env: &mut JNIEnv, config: JObject) -> Result<VntUtilSync, Error> {
let cipher_model = to_string_not_null(env, &config, "cipherModel")?; let cipher_model = to_string_not_null(env, &config, "cipherModel")?;
let tcp = env.get_field(&config, "tcp", "Z")?.z()?; let tcp = env.get_field(&config, "tcp", "Z")?.z()?;
let finger = env.get_field(&config, "finger", "Z")?.z()?; let finger = env.get_field(&config, "finger", "Z")?.z()?;
let first_latency = env.get_field(&config, "firstLatency", "Z")?.z()?;
let in_ips = to_string(env, &config, "inIps")?; let in_ips = to_string(env, &config, "inIps")?;
let out_ips = to_string(env, &config, "outIps")?; let out_ips = to_string(env, &config, "outIps")?;
let port = env.get_field(&config, "port", "I")?.i()? as u16; let port = env.get_field(&config, "port", "I")?.i()? as u16;
@@ -152,6 +153,7 @@ fn new_sync(env: &mut JNIEnv, config: JObject) -> Result<VntUtilSync, Error> {
finger, finger,
PunchModel::All, PunchModel::All,
port, port,
first_latency,
); );
match VntUtilSync::new(config) { match VntUtilSync::new(config) {
Ok(vnt_util) => Ok(vnt_util), Ok(vnt_util) => Ok(vnt_util),
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "vnt" name = "vnt"
version = "1.2.6" version = "1.2.7"
edition = "2021" edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+57 -189
View File
@@ -3,13 +3,12 @@ use std::io::{Read, Write};
use std::net::TcpStream; use std::net::TcpStream;
use std::net::UdpSocket as StdUdpSocket; use std::net::UdpSocket as StdUdpSocket;
use std::net::{Ipv4Addr, Ipv6Addr, Shutdown, SocketAddr}; use std::net::{Ipv4Addr, Ipv6Addr, Shutdown, SocketAddr};
use std::sync::atomic::Ordering;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use std::{io, thread}; use std::{io, thread};
use crossbeam_epoch::{Atomic, Owned};
use crossbeam_utils::atomic::AtomicCell; use crossbeam_utils::atomic::AtomicCell;
use parking_lot::RwLock;
use tokio::net::UdpSocket; use tokio::net::UdpSocket;
use tokio::sync::watch::{channel, Receiver, Sender}; use tokio::sync::watch::{channel, Receiver, Sender};
@@ -25,12 +24,13 @@ pub struct ContextInner {
pub(crate) main_channel_ipv6: Option<Arc<StdUdpSocket>>, pub(crate) main_channel_ipv6: Option<Arc<StdUdpSocket>>,
//在udp的基础上,可以选择使用tcp和服务端通信 //在udp的基础上,可以选择使用tcp和服务端通信
pub(crate) main_tcp_channel: Option<std::sync::mpsc::SyncSender<Vec<u8>>>, pub(crate) main_tcp_channel: Option<std::sync::mpsc::SyncSender<Vec<u8>>>,
pub(crate) route_table: Atomic<HashMap<Ipv4Addr, Vec<(Route, Arc<AtomicCell<Instant>>)>>>, pub(crate) route_table: RwLock<HashMap<Ipv4Addr, Vec<(Route, AtomicCell<Instant>)>>>,
pub(crate) status_receiver: Receiver<Status>, pub(crate) status_receiver: Receiver<Status>,
pub(crate) status_sender: Sender<Status>, pub(crate) status_sender: Sender<Status>,
pub(crate) udp_map: Atomic<HashMap<usize, Arc<UdpSocket>>>, pub(crate) udp_map: RwLock<HashMap<usize, Arc<UdpSocket>>>,
pub(crate) channel_num: usize, pub(crate) channel_num: usize,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>, current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
first_latency: bool,
} }
#[derive(Clone)] #[derive(Clone)]
@@ -45,6 +45,7 @@ impl Context {
main_tcp_channel: Option<std::sync::mpsc::SyncSender<Vec<u8>>>, main_tcp_channel: Option<std::sync::mpsc::SyncSender<Vec<u8>>>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>, current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
_channel_num: usize, _channel_num: usize,
first_latency: bool,
) -> Self { ) -> Self {
//当前版本只支持一个通道 //当前版本只支持一个通道
let channel_num = 1; let channel_num = 1;
@@ -53,12 +54,13 @@ impl Context {
main_channel, main_channel,
main_channel_ipv6, main_channel_ipv6,
main_tcp_channel, main_tcp_channel,
route_table: Atomic::new(HashMap::with_capacity(16)), route_table: RwLock::new(HashMap::with_capacity(16)),
status_receiver, status_receiver,
status_sender, status_sender,
udp_map: Atomic::new(HashMap::with_capacity(16)), udp_map: RwLock::new(HashMap::with_capacity(16)),
channel_num, channel_num,
current_device, current_device,
first_latency,
}); });
Self { inner } Self { inner }
} }
@@ -120,41 +122,10 @@ impl Context {
} }
} }
fn insert_udp(&self, id: usize, udp: Arc<UdpSocket>) { fn insert_udp(&self, id: usize, udp: Arc<UdpSocket>) {
self.insert_udp_(id, Some(udp)) self.inner.udp_map.write().insert(id, udp);
} }
fn remove_udp(&self, id: usize) { fn remove_udp(&self, id: usize) {
self.insert_udp_(id, None) self.inner.udp_map.write().remove(&id);
}
fn insert_udp_(&self, id: usize, udp: Option<Arc<UdpSocket>>) {
let guard = &crossbeam_epoch::pin();
let udp_map = &self.inner.udp_map;
let mut udp_map_shared = udp_map.load(Ordering::Acquire, guard);
loop {
let mut map = unsafe { udp_map_shared.deref().clone() };
match udp.clone() {
None => {
map.remove(&id);
}
Some(udp) => {
map.insert(id, udp);
}
}
match udp_map.compare_exchange(
udp_map_shared,
Owned::new(map),
Ordering::AcqRel,
Ordering::Relaxed,
guard,
) {
Ok(_p) => unsafe {
guard.defer_destroy(udp_map_shared);
return;
},
Err(e) => {
udp_map_shared = e.current;
}
}
}
} }
pub fn send_main_udp(&self, buf: &[u8], addr: SocketAddr) -> io::Result<usize> { pub fn send_main_udp(&self, buf: &[u8], addr: SocketAddr) -> io::Result<usize> {
if addr.is_ipv6() { if addr.is_ipv6() {
@@ -181,19 +152,12 @@ impl Context {
} }
pub(crate) fn try_send_all(&self, buf: &[u8], addr: SocketAddr) -> io::Result<()> { pub(crate) fn try_send_all(&self, buf: &[u8], addr: SocketAddr) -> io::Result<()> {
let table = unsafe { let table = self.inner.udp_map.read();
let guard = &crossbeam_epoch::pin();
self.inner
.udp_map
.load(Ordering::Relaxed, guard)
.deref()
.clone()
};
if table.is_empty() { if table.is_empty() {
log::error!("udp列表为空,addr={}", addr); log::error!("udp列表为空,addr={}", addr);
return Ok(()); return Ok(());
} }
for (_, udp) in table { for (_, udp) in table.iter() {
//使用ipv6的udp发送ipv4报文会出错 //使用ipv6的udp发送ipv4报文会出错
if let Err(e) = udp.try_send_to(buf, addr) { if let Err(e) = udp.try_send_to(buf, addr) {
log::error!("{:?}", e); log::error!("{:?}", e);
@@ -211,14 +175,7 @@ impl Context {
self.try_send_by_key(buf, &route.route_key()) self.try_send_by_key(buf, &route.route_key())
} }
fn get_route_by_id(&self, id: &Ipv4Addr) -> io::Result<Route> { fn get_route_by_id(&self, id: &Ipv4Addr) -> io::Result<Route> {
let guard = &crossbeam_epoch::pin(); if let Some(v) = self.inner.route_table.read().get(id) {
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
if let Some(v) = table.get(id) {
if v.is_empty() { if v.is_empty() {
return Err(io::Error::new(io::ErrorKind::NotFound, "route not found")); return Err(io::Error::new(io::ErrorKind::NotFound, "route not found"));
} }
@@ -297,9 +254,7 @@ impl Context {
} }
} }
fn get_udp_by_route(&self, route_key: &RouteKey) -> Option<Arc<UdpSocket>> { fn get_udp_by_route(&self, route_key: &RouteKey) -> Option<Arc<UdpSocket>> {
let guard = &crossbeam_epoch::pin(); self.inner.udp_map.read().get(&route_key.index).cloned()
let udp_map = unsafe { self.inner.udp_map.load(Ordering::Relaxed, guard).deref() };
udp_map.get(&route_key.index).cloned()
} }
pub fn add_route_if_absent(&self, id: Ipv4Addr, route: Route) { pub fn add_route_if_absent(&self, id: Ipv4Addr, route: Route) {
@@ -310,96 +265,58 @@ impl Context {
} }
fn add_route_(&self, id: Ipv4Addr, route: Route, only_if_absent: bool) { fn add_route_(&self, id: Ipv4Addr, route: Route, only_if_absent: bool) {
let key = route.route_key(); let key = route.route_key();
let guard = &crossbeam_epoch::pin(); let mut route_table = self.inner.route_table.write();
let route_table = &self.inner.route_table; let list = route_table
let mut table_share = route_table.load(Ordering::Acquire, guard); .entry(id)
loop { .or_insert_with(|| Vec::with_capacity(4));
let mut table = unsafe { table_share.deref().clone() }; let mut exist = false;
let list = table.entry(id).or_insert_with(|| Vec::with_capacity(4)); for (x, time) in list.iter_mut() {
let mut exist = false; if x.metric < route.metric {
for (x, time) in list.iter_mut() { //不能比当前的路径更长
if x.metric < route.metric { return;
//不能比当前的路径更长 }
if x.route_key() == key {
if only_if_absent {
return; return;
} }
if x.route_key() == key { x.metric = route.metric;
if only_if_absent { x.rt = route.rt;
return; exist = true;
} time.store(Instant::now());
x.metric = route.metric; break;
x.rt = route.rt;
exist = true;
time.store(Instant::now());
break;
}
} }
if exist { }
list.sort_by_key(|(k, _)| k.sort_key()); if exist {
} else { list.sort_by_key(|(k, _)| k.rt);
if route.metric == 1 { } else {
//添加了直连的则排除非直连的 if route.metric == 1 && !self.inner.first_latency {
list.retain(|(k, _)| k.metric == 1); //非优先延迟的情况下 添加了直连的则排除非直连的
} list.retain(|(k, _)| k.metric == 1);
list.push((route, Arc::new(AtomicCell::new(Instant::now()))));
list.sort_by_key(|(k, _)| k.sort_key());
let max_len = self.inner.channel_num + 1;
if list.len() > max_len {
list.truncate(max_len);
}
} }
match route_table.compare_exchange( list.sort_by_key(|(k, _)| k.rt);
table_share, let max_len = self.inner.channel_num;
Owned::new(table), if list.len() > max_len {
Ordering::AcqRel, list.truncate(max_len);
Ordering::Relaxed,
guard,
) {
Ok(_p) => unsafe {
guard.defer_destroy(table_share);
break;
},
Err(e) => {
table_share = e.current;
}
} }
list.push((route, AtomicCell::new(Instant::now())));
} }
} }
pub fn route(&self, id: &Ipv4Addr) -> Option<Vec<Route>> { pub fn route(&self, id: &Ipv4Addr) -> Option<Vec<Route>> {
let guard = &crossbeam_epoch::pin(); if let Some(v) = self.inner.route_table.read().get(id) {
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
if let Some(v) = table.get(id) {
Some(v.iter().map(|(i, _)| *i).collect()) Some(v.iter().map(|(i, _)| *i).collect())
} else { } else {
None None
} }
} }
pub fn route_one(&self, id: &Ipv4Addr) -> Option<Route> { pub fn route_one(&self, id: &Ipv4Addr) -> Option<Route> {
let guard = &crossbeam_epoch::pin(); if let Some(v) = self.inner.route_table.read().get(id) {
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
if let Some(v) = table.get(id) {
v.first().map(|(i, _)| *i) v.first().map(|(i, _)| *i)
} else { } else {
None None
} }
} }
pub fn route_to_id(&self, route_key: &RouteKey) -> Option<Ipv4Addr> { pub fn route_to_id(&self, route_key: &RouteKey) -> Option<Ipv4Addr> {
let guard = &crossbeam_epoch::pin(); let table = self.inner.route_table.read();
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
for (k, v) in table.iter() { for (k, v) in table.iter() {
for (route, _) in v { for (route, _) in v {
if &route.route_key() == route_key && route.is_p2p() { if &route.route_key() == route_key && route.is_p2p() {
@@ -410,14 +327,7 @@ impl Context {
None None
} }
pub fn need_punch(&self, id: &Ipv4Addr) -> bool { pub fn need_punch(&self, id: &Ipv4Addr) -> bool {
let guard = &crossbeam_epoch::pin(); if let Some(v) = self.inner.route_table.read().get(id) {
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
if let Some(v) = table.get(id) {
if v.iter().filter(|(k, _)| k.is_p2p()).count() >= self.inner.channel_num { if v.iter().filter(|(k, _)| k.is_p2p()).count() >= self.inner.channel_num {
return false; return false;
} }
@@ -425,13 +335,7 @@ impl Context {
true true
} }
pub fn route_table(&self) -> Vec<(Ipv4Addr, Vec<Route>)> { pub fn route_table(&self) -> Vec<(Ipv4Addr, Vec<Route>)> {
let guard = &crossbeam_epoch::pin(); let table = self.inner.route_table.read();
let table = unsafe {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
table table
.iter() .iter()
.map(|(k, v)| (k.clone(), v.iter().map(|(i, _)| *i).collect())) .map(|(k, v)| (k.clone(), v.iter().map(|(i, _)| *i).collect()))
@@ -439,14 +343,8 @@ impl Context {
} }
pub fn route_table_one(&self) -> Vec<(Ipv4Addr, Route)> { pub fn route_table_one(&self) -> Vec<(Ipv4Addr, Route)> {
let mut list = Vec::with_capacity(8); let mut list = Vec::with_capacity(8);
let guard = &crossbeam_epoch::pin(); let table = self.inner.route_table.read();
let table = unsafe { for (k, v) in table.iter() {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
for (k, v) in table {
if let Some((route, _)) = v.first() { if let Some((route, _)) = v.first() {
list.push((*k, *route)); list.push((*k, *route));
} }
@@ -455,14 +353,8 @@ impl Context {
} }
pub fn direct_route_table_one(&self) -> Vec<(Ipv4Addr, Route)> { pub fn direct_route_table_one(&self) -> Vec<(Ipv4Addr, Route)> {
let mut list = Vec::with_capacity(8); let mut list = Vec::with_capacity(8);
let guard = &crossbeam_epoch::pin(); let table = self.inner.route_table.read();
let table = unsafe { for (k, v) in table.iter() {
self.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
for (k, v) in table {
if let Some((route, _)) = v.first() { if let Some((route, _)) = v.first() {
if route.metric == 1 { if route.metric == 1 {
list.push((*k, *route)); list.push((*k, *route));
@@ -473,38 +365,14 @@ impl Context {
} }
pub fn remove_route(&self, id: &Ipv4Addr, route_key: RouteKey) { pub fn remove_route(&self, id: &Ipv4Addr, route_key: RouteKey) {
let guard = &crossbeam_epoch::pin(); if let Some(routes) = self.inner.route_table.write().get_mut(id) {
let route_table = &self.inner.route_table; routes.retain(|(x, _)| x.route_key() != route_key);
let mut table_share = route_table.load(Ordering::Acquire, guard); } else {
loop { return;
let mut table = unsafe { table_share.deref().clone() };
if let Some(routes) = table.get_mut(id) {
routes.retain(|(x, _)| x.route_key() != route_key);
match route_table.compare_exchange(
table_share,
Owned::new(table),
Ordering::AcqRel,
Ordering::Relaxed,
guard,
) {
Ok(_p) => unsafe {
guard.defer_destroy(table_share);
return;
},
Err(e) => {
table_share = e.current;
}
}
} else {
return;
}
} }
} }
pub fn update_read_time(&self, id: &Ipv4Addr, route_key: &RouteKey) { pub fn update_read_time(&self, id: &Ipv4Addr, route_key: &RouteKey) {
let guard = &crossbeam_epoch::pin(); if let Some(routes) = self.inner.route_table.read().get(id) {
let table_share = self.inner.route_table.load(Ordering::Relaxed, guard);
let table = unsafe { table_share.deref() };
if let Some(routes) = table.get(id) {
for (route, time) in routes { for (route, time) in routes {
if &route.route_key() == route_key { if &route.route_key() == route_key {
time.store(Instant::now()); time.store(Instant::now());
+4 -12
View File
@@ -1,11 +1,11 @@
use crate::channel::channel::Context;
use crate::channel::RouteKey;
use std::io; use std::io;
use std::io::{Error, ErrorKind}; use std::io::{Error, ErrorKind};
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::sync::atomic::Ordering;
use std::time::Duration; use std::time::Duration;
use crate::channel::channel::Context;
use crate::channel::RouteKey;
pub struct Idle { pub struct Idle {
read_idle: Duration, read_idle: Duration,
context: Context, context: Context,
@@ -23,15 +23,7 @@ impl Idle {
loop { loop {
let mut max = Duration::from_secs(0); let mut max = Duration::from_secs(0);
{ {
let guard = &crossbeam_epoch::pin(); for (ip, routes) in self.context.inner.route_table.read().iter() {
let table = unsafe {
self.context
.inner
.route_table
.load(Ordering::Relaxed, guard)
.deref()
};
for (ip, routes) in table.iter() {
for (route, time) in routes { for (route, time) in routes {
let last_read = time.load().elapsed(); let last_read = time.load().elapsed();
if last_read >= self.read_idle { if last_read >= self.read_idle {
+9 -10
View File
@@ -3,13 +3,11 @@ use std::io;
use std::net::TcpStream; use std::net::TcpStream;
use std::net::UdpSocket; use std::net::UdpSocket;
use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4}; use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4};
use std::sync::atomic::Ordering;
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
use crossbeam_epoch::Atomic;
use crossbeam_utils::atomic::AtomicCell; use crossbeam_utils::atomic::AtomicCell;
use parking_lot::Mutex; use parking_lot::{Mutex, RwLock};
use rand::Rng; use rand::Rng;
use tokio::sync::mpsc::channel; use tokio::sync::mpsc::channel;
@@ -53,7 +51,7 @@ pub struct Vnt {
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>, device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
nat_test: NatTest, nat_test: NatTest,
connect_status: Arc<AtomicCell<ConnectStatus>>, connect_status: Arc<AtomicCell<ConnectStatus>>,
peer_nat_info_map: Arc<Atomic<HashMap<Ipv4Addr, NatInfo>>>, peer_nat_info_map: Arc<RwLock<HashMap<Ipv4Addr, NatInfo>>>,
} }
pub struct VntUtil { pub struct VntUtil {
@@ -262,6 +260,7 @@ impl VntUtil {
tcp_sender, tcp_sender,
current_device.clone(), current_device.clone(),
1, 1,
config.first_latency,
); );
let punch = Punch::new(context.clone(), config.punch_model); let punch = Punch::new(context.clone(), config.punch_model);
let idle = Idle::new(Duration::from_secs(16), context.clone()); let idle = Idle::new(Duration::from_secs(16), context.clone());
@@ -278,8 +277,8 @@ impl VntUtil {
)); ));
let device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>> = let device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>> =
Arc::new(Mutex::new((response.epoch, response.device_info_list))); Arc::new(Mutex::new((response.epoch, response.device_info_list)));
let peer_nat_info_map: Arc<Atomic<HashMap<Ipv4Addr, NatInfo>>> = let peer_nat_info_map: Arc<RwLock<HashMap<Ipv4Addr, NatInfo>>> =
Arc::new(Atomic::new(HashMap::new())); Arc::new(RwLock::new(HashMap::with_capacity(16)));
let connect_status = Arc::new(AtomicCell::new(ConnectStatus::Connected)); let connect_status = Arc::new(AtomicCell::new(ConnectStatus::Connected));
let public_ip = response.public_ip; let public_ip = response.public_ip;
let public_port = response.public_port; let public_port = response.public_port;
@@ -504,10 +503,7 @@ impl Vnt {
self.current_device.load() self.current_device.load()
} }
pub fn peer_nat_info(&self, ip: &Ipv4Addr) -> Option<NatInfo> { pub fn peer_nat_info(&self, ip: &Ipv4Addr) -> Option<NatInfo> {
let guard = &crossbeam_epoch::pin(); self.peer_nat_info_map.read().get(ip).cloned()
let shared = self.peer_nat_info_map.load(Ordering::Acquire, guard);
let map = unsafe { shared.deref() };
map.get(ip).map(|e| e.clone())
} }
pub fn connection_status(&self) -> ConnectStatus { pub fn connection_status(&self) -> ConnectStatus {
self.connect_status.load() self.connect_status.load()
@@ -590,6 +586,7 @@ pub struct Config {
pub finger: bool, pub finger: bool,
pub punch_model: PunchModel, pub punch_model: PunchModel,
pub port: u16, pub port: u16,
pub first_latency: bool,
} }
impl Config { impl Config {
@@ -616,6 +613,7 @@ impl Config {
finger: bool, finger: bool,
punch_model: PunchModel, punch_model: PunchModel,
port: u16, port: u16,
first_latency: bool,
) -> Self { ) -> Self {
for x in stun_server.iter_mut() { for x in stun_server.iter_mut() {
if !x.contains(":") { if !x.contains(":") {
@@ -646,6 +644,7 @@ impl Config {
finger, finger,
punch_model, punch_model,
port, port,
first_latency,
} }
} }
} }
+5 -14
View File
@@ -1,12 +1,10 @@
use std::collections::HashMap; use std::collections::HashMap;
use std::net::{Ipv4Addr, Ipv6Addr, SocketAddrV4, SocketAddrV6}; use std::net::{Ipv4Addr, Ipv6Addr, SocketAddrV4, SocketAddrV6};
use std::sync::atomic::Ordering;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use crossbeam_epoch::{Atomic, Owned};
use crossbeam_utils::atomic::AtomicCell; use crossbeam_utils::atomic::AtomicCell;
use parking_lot::Mutex; use parking_lot::{Mutex, RwLock};
use protobuf::Message; use protobuf::Message;
use tokio::sync::mpsc::Sender; use tokio::sync::mpsc::Sender;
@@ -47,7 +45,7 @@ pub struct ChannelDataHandler {
igmp_server: Option<IgmpServer>, igmp_server: Option<IgmpServer>,
device_writer: DeviceWriter, device_writer: DeviceWriter,
connect_status: Arc<AtomicCell<ConnectStatus>>, connect_status: Arc<AtomicCell<ConnectStatus>>,
peer_nat_info_map: Arc<Atomic<HashMap<Ipv4Addr, NatInfo>>>, peer_nat_info_map: Arc<RwLock<HashMap<Ipv4Addr, NatInfo>>>,
#[cfg(feature = "ip_proxy")] #[cfg(feature = "ip_proxy")]
ip_proxy_map: Option<IpProxyMap>, ip_proxy_map: Option<IpProxyMap>,
out_external_route: AllowExternalRoute, out_external_route: AllowExternalRoute,
@@ -70,7 +68,7 @@ impl ChannelDataHandler {
igmp_server: Option<IgmpServer>, igmp_server: Option<IgmpServer>,
device_writer: DeviceWriter, device_writer: DeviceWriter,
connect_status: Arc<AtomicCell<ConnectStatus>>, connect_status: Arc<AtomicCell<ConnectStatus>>,
peer_nat_info_map: Arc<Atomic<HashMap<Ipv4Addr, NatInfo>>>, peer_nat_info_map: Arc<RwLock<HashMap<Ipv4Addr, NatInfo>>>,
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>, #[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
out_external_route: AllowExternalRoute, out_external_route: AllowExternalRoute,
cone_sender: Sender<(Ipv4Addr, NatInfo)>, cone_sender: Sender<(Ipv4Addr, NatInfo)>,
@@ -422,15 +420,8 @@ impl ChannelDataHandler {
punch_info.nat_type.enum_value_or_default().into(), punch_info.nat_type.enum_value_or_default().into(),
); );
{ {
let guard = &crossbeam_epoch::pin(); let peer_nat_info = peer_nat_info.clone();
let nat_map = &self.peer_nat_info_map; self.peer_nat_info_map.write().insert(source, peer_nat_info);
let nat_map_shared = nat_map.load(Ordering::Acquire, guard);
let mut map = unsafe { nat_map_shared.deref().clone() };
map.insert(source, peer_nat_info.clone());
nat_map.store(Owned::new(map), Ordering::Release);
unsafe {
guard.defer_destroy(nat_map_shared);
}
} }
if !punch_info.reply { if !punch_info.reply {
let mut punch_reply = PunchInfo::new(); let mut punch_reply = PunchInfo::new();
+8 -43
View File
@@ -1,10 +1,8 @@
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::sync::atomic::Ordering;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use crossbeam_epoch::{Atomic, Owned};
use parking_lot::RwLock; use parking_lot::RwLock;
use packet::igmp::igmp_v2::IgmpV2Packet; use packet::igmp::igmp_v2::IgmpV2Packet;
@@ -51,13 +49,13 @@ impl Multicast {
#[derive(Clone)] #[derive(Clone)]
pub struct IgmpServer { pub struct IgmpServer {
multicast: Arc<Atomic<HashMap<Ipv4Addr, Arc<RwLock<Multicast>>>>>, multicast: Arc<RwLock<HashMap<Ipv4Addr, Arc<RwLock<Multicast>>>>>,
} }
impl IgmpServer { impl IgmpServer {
pub fn new(device_writer: DeviceWriter) -> Self { pub fn new(device_writer: DeviceWriter) -> Self {
let multicast: Arc<Atomic<HashMap<Ipv4Addr, Arc<RwLock<Multicast>>>>> = let multicast: Arc<RwLock<HashMap<Ipv4Addr, Arc<RwLock<Multicast>>>>> =
Arc::new(Atomic::new(HashMap::with_capacity(16))); Arc::new(RwLock::new(HashMap::with_capacity(16)));
std::thread::spawn(move || { std::thread::spawn(move || {
//预留以太网帧头和ip头 //预留以太网帧头和ip头
let mut buf = [0; 14 + 24 + 12]; let mut buf = [0; 14 + 24 + 12];
@@ -98,17 +96,10 @@ impl IgmpServer {
Self { multicast } Self { multicast }
} }
pub fn load(&self, multicast_addr: &Ipv4Addr) -> Option<Arc<RwLock<Multicast>>> { pub fn load(&self, multicast_addr: &Ipv4Addr) -> Option<Arc<RwLock<Multicast>>> {
let guard = &crossbeam_epoch::pin(); self.multicast.read().get(multicast_addr).cloned()
let multicast = unsafe { self.multicast.load(Ordering::Relaxed, guard).deref() };
if let Some(entry) = multicast.get(multicast_addr) {
Some(entry.clone())
} else {
None
}
} }
pub fn handle(&self, buf: &[u8], source: Ipv4Addr) -> crate::Result<()> { pub fn handle(&self, buf: &[u8], source: Ipv4Addr) -> crate::Result<()> {
let guard = &crossbeam_epoch::pin(); let multicast = self.multicast.read();
let multicast = unsafe { self.multicast.load(Ordering::Relaxed, guard).deref() };
for (_, v) in multicast.iter() { for (_, v) in multicast.iter() {
let mut list = Vec::new(); let mut list = Vec::new();
let mut write_guard = v.write(); let mut write_guard = v.write();
@@ -142,7 +133,7 @@ impl IgmpServer {
if !multicast_addr.is_multicast() { if !multicast_addr.is_multicast() {
return Ok(()); return Ok(());
} }
if let Some(entry) = self.get_multicast(&multicast_addr) { if let Some(entry) = self.load(&multicast_addr) {
let mut guard = entry.write(); let mut guard = entry.write();
guard.map.remove(&source); guard.map.remove(&source);
guard.members.remove(&source); guard.members.remove(&source);
@@ -234,35 +225,9 @@ impl IgmpServer {
} }
Ok(()) Ok(())
} }
fn get_multicast(&self, multicast_addr: &Ipv4Addr) -> Option<Arc<RwLock<Multicast>>> {
let guard = &crossbeam_epoch::pin();
let multicast = &self.multicast;
let table_share = multicast.load(Ordering::Acquire, guard);
unsafe { table_share.deref().get(multicast_addr).map(|v| v.clone()) }
}
fn add_multicast(&self, multicast_addr: Ipv4Addr) -> Arc<RwLock<Multicast>> { fn add_multicast(&self, multicast_addr: Ipv4Addr) -> Arc<RwLock<Multicast>> {
let guard = &crossbeam_epoch::pin();
let multicast = &self.multicast;
let mut table_share = multicast.load(Ordering::Acquire, guard);
let value = Arc::new(RwLock::new(Multicast::new())); let value = Arc::new(RwLock::new(Multicast::new()));
loop { self.multicast.write().insert(multicast_addr, value.clone());
let mut table = unsafe { table_share.deref().clone() }; value
table.insert(multicast_addr, value.clone());
match self.multicast.compare_exchange(
table_share,
Owned::new(table.clone()),
Ordering::AcqRel,
Ordering::Relaxed,
guard,
) {
Ok(_p) => unsafe {
guard.defer_destroy(table_share);
return value;
},
Err(e) => {
table_share = e.current;
}
}
}
} }
} }
+1 -1
View File
@@ -1,4 +1,3 @@
#[cfg(not(target_os = "android"))]
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::net::SocketAddrV4; use std::net::SocketAddrV4;
use std::sync::Arc; use std::sync::Arc;
@@ -14,6 +13,7 @@ use packet::ip::ipv4::packet::IpV4Packet;
use crate::channel::sender::ChannelSender; use crate::channel::sender::ChannelSender;
use crate::cipher::Cipher; use crate::cipher::Cipher;
use crate::handle::CurrentDeviceInfo; use crate::handle::CurrentDeviceInfo;
#[cfg(not(target_os = "android"))]
use crate::ip_proxy::icmp_proxy::IcmpHandler; use crate::ip_proxy::icmp_proxy::IcmpHandler;
use crate::ip_proxy::tcp_proxy::{TcpHandler, TcpProxy}; use crate::ip_proxy::tcp_proxy::{TcpHandler, TcpProxy};
use crate::ip_proxy::udp_proxy::{UdpHandler, UdpProxy}; use crate::ip_proxy::udp_proxy::{UdpHandler, UdpProxy};
+1 -1
View File
@@ -1,5 +1,5 @@
use crate::error::Error; use crate::error::Error;
pub const VNT_VERSION: &'static str = "1.2.6"; pub const VNT_VERSION: &'static str = "1.2.7";
pub type Result<T> = std::result::Result<T, Error>; pub type Result<T> = std::result::Result<T, Error>;
pub mod channel; pub mod channel;