Compare commits

...
28 Commits
Author SHA1 Message Date
lubeilin 707b07b8d3 ipv6改为完整地址 2023-09-24 15:42:19 +08:00
lubeilin de5a6971f0 更新读取时间不需要再插入 2023-09-24 14:58:03 +08:00
lubeilin 056036c4d2 减少注册和探测nat的频率 2023-09-24 14:53:37 +08:00
lubeilin 1cfb188845 去除多余依赖 2023-09-24 13:45:47 +08:00
lubeilin 73a2c31854 连接通过关闭同时关闭tap 2023-09-24 13:42:09 +08:00
lubeilin 92eea536f8 删除多余依赖 2023-09-24 12:57:49 +08:00
lubeilin 00936a923e 避免直接关闭网卡 2023-09-24 12:55:42 +08:00
lubeilin ba69ba78af 增加线程名称 2023-09-23 23:09:37 +08:00
lubeilin 17b206bace 1.2.4.3 2023-09-23 21:38:31 +08:00
lubeilin 58d5a4f5da 增加wintun日志 2023-09-23 21:38:20 +08:00
lubeilin 2438d14175 避免短时间重复上传服务端密钥 2023-09-23 21:33:00 +08:00
lubeilin 99f8526799 去除tap广播路由 2023-09-22 22:49:05 +08:00
lubeilin 56fcbd64ed 增加日志 2023-09-22 22:43:23 +08:00
lubeilin 3766b2b7c1 修改命令超时时间 2023-09-22 22:43:02 +08:00
lubeilin baf0698fe4 去除广播路由 2023-09-22 22:17:13 +08:00
lubeilin 301938b9fc 增加小版本 2023-09-22 18:19:13 +08:00
lubeilin 16a37c713a 增加日志 2023-09-22 18:18:29 +08:00
lubeilin d412a769dd fmt 2023-09-22 18:18:10 +08:00
lubeilin c8eecc87fd 调整心跳间隔,服务端和客户端心跳分离 2023-09-22 18:17:12 +08:00
lubeilin 6a11db70c8 调整代理超时时间 2023-09-22 18:16:06 +08:00
lubeilin d7c121a756 commit:
1.去除缓冲池
2.数据处理改为同步方法
3.fmt
2023-09-20 19:54:49 +08:00
lubeilin 3429ee8bd6 增加提示 2023-09-20 16:03:25 +08:00
lubeilin 4422f9f8b7 Merge remote-tracking branch 'origin/main'
# Conflicts:
#	vnt/src/ip_proxy/tcp_proxy.rs
2023-09-20 15:39:54 +08:00
lubeilin 57ed454c93 修复内网ip断线问题 2023-09-20 11:11:31 +08:00
lubeilin 236205c0f3 修复内网ip断线问题 2023-09-19 18:25:14 +08:00
lubeilin 99b4bf0041 Merge remote-tracking branch 'origin/main' 2023-09-18 21:51:17 +08:00
lubeilin 9495e39700 优化nat校验 2023-09-18 21:51:08 +08:00
lbl8603 75e244e3a8 Update README.md 2023-09-18 11:17:21 +08:00
30 changed files with 384 additions and 256 deletions
+1 -1
View File
@@ -68,7 +68,7 @@ A virtual network tool (VPN)
- Mac
- Linux
- 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)
- Android
- [VntApp](https://github.com/lbl8603/VntApp)
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "common"
version = "1.2.3"
version = "1.2.4"
edition = "2021"
# 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]
name = "vnt-cli"
version = "1.2.3"
version = "1.2.4"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+8
View File
@@ -36,6 +36,11 @@
### -W
开启和服务端通信的数据加密,采用rsa+aes256gcm加密客户端和服务端之间通信的数据,可以避免token泄漏、中间人攻击
注意:
1. -w `<password>`是用于客户端-客户端之间的加密,password不会传递到服务端,只添加这个参数不会加密客户端-服务端通信的数据
2. -W 用于开启客户端-服务端之间的加密
### -m
模拟组播,高频使用组播通信时,可以尝试开启此参数,默认情况下会把组播当作广播发给所有节点
@@ -65,8 +70,11 @@
| 1~8位 | aes_ecb | AES128-ECB |
| `>=`8 | aes_ecb | AES256-ECB |
### --finger
开启数据指纹校验,可增加安全性,如果服务端开启指纹校验,则客户端也必须开启,开启会损耗一部分性能
注意:默认情况下服务端不会对中转的数据做校验,如果要对中转的数据做校验,则需要客户端、服务端都开启此参数
### --relay
禁用p2p,在网络环境很差时,只使用服务器中转效果可能更好(可以配合--tcp参数一起使用)
### --list
+1 -1
View File
@@ -26,7 +26,7 @@ impl CommandClient {
}
};
let udp = UdpSocket::bind("127.0.0.1:0")?;
udp.set_read_timeout(Some(Duration::from_secs(2)))?;
udp.set_read_timeout(Some(Duration::from_secs(5)))?;
udp.connect(SocketAddr::V4(SocketAddrV4::new(
Ipv4Addr::new(127, 0, 0, 1),
port,
+1 -1
View File
@@ -17,7 +17,7 @@ pub enum CommandEnum {
pub fn command(cmd: CommandEnum) {
if let Err(e) = command_(cmd) {
println!("cmd: {}", e);
println!("cmd: {:?}", e);
}
}
+7 -2
View File
@@ -17,15 +17,20 @@ impl CommandServer {
let udp = UdpSocket::bind("127.0.0.1:0").await?;
let path_buf = crate::app_home()?.join("command-port");
let mut file = std::fs::File::create(path_buf)?;
file.write_all(udp.local_addr()?.port().to_string().as_bytes())?;
let addr = udp.local_addr()?;
file.write_all(addr.port().to_string().as_bytes())?;
file.sync_all()?;
log::info!("启动后台cmd:{:?}", addr);
let mut buf = [0u8; 64];
loop {
let (len, addr) = udp.recv_from(&mut buf).await?;
match std::str::from_utf8(&buf[..len]) {
Ok(cmd) => {
log::info!("收到cmd={:?}", cmd);
if let Ok(out) = command(cmd, &vnt) {
let _ = udp.send_to(out.as_bytes(), addr).await;
if let Err(e) = udp.send_to(out.as_bytes(), addr).await {
log::warn!("cmd={},err={:?}", cmd, e);
}
if "stopped" == &out {
break;
}
+9 -4
View File
@@ -214,10 +214,13 @@ fn main() {
return;
}
let cipher_model = matches
.opt_get::<CipherModel>("model")
.unwrap()
.unwrap_or(CipherModel::AesGcm);
let cipher_model = match matches.opt_get::<CipherModel>("model") {
Ok(model) => model.unwrap_or(CipherModel::AesGcm),
Err(e) => {
println!("'--model ' invalid,{}", e);
return;
}
};
let finger = matches.opt_present("finger");
let punch_model = matches
@@ -250,6 +253,7 @@ fn main() {
main0(config, !unused_cmd);
std::process::exit(0);
}
#[tokio::main]
async fn main0(config: Config, show_cmd: bool) {
let server_encrypt = config.server_encrypt;
@@ -358,6 +362,7 @@ async fn main0(config: Config, show_cmd: bool) {
let vnt_c = vnt.clone();
tokio::spawn(async {
if let Err(e) = command::server::CommandServer::new().start(vnt_c).await {
log::warn!("cmd:{:?}", e);
println!("command error :{}", e);
}
});
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "vnt-jni"
version = "1.2.3"
version = "1.2.4"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+1 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "vnt"
version = "1.2.3"
version = "1.2.4"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
@@ -14,8 +14,6 @@ crossbeam-utils = "0.8"
crossbeam-epoch = "0.9.15"
dashmap = "5.5.1"
parking_lot = "0.12.1"
byte-pool = "0.2.4"
lazy_static = "1.4.0"
rand = "0.8.5"
sha2 = { version = "0.10.6", features = ["oid"] }
thiserror = "1.0.37"
+54 -68
View File
@@ -6,7 +6,6 @@ use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::time::{Duration, Instant};
use byte_pool::{Block, BytePool};
use crossbeam_epoch::{Atomic, Owned};
use crossbeam_utils::atomic::AtomicCell;
use dashmap::DashMap;
@@ -23,9 +22,6 @@ use crate::handle::recv_handler::ChannelDataHandler;
use crate::handle::CurrentDeviceInfo;
use crate::ip_proxy::DashMapNew;
lazy_static::lazy_static! {
static ref POOL:BytePool = BytePool::new();
}
pub struct ContextInner {
//udp用于打洞、服务端通信(可选)
pub(crate) main_channel: Arc<StdUdpSocket>,
@@ -544,13 +540,13 @@ impl Channel {
#[derive(Clone)]
struct BufSenderGroup(
usize,
Vec<std::sync::mpsc::SyncSender<(Block<'static>, usize, usize, RouteKey)>>,
Vec<std::sync::mpsc::SyncSender<(Vec<u8>, usize, usize, RouteKey)>>,
);
struct BufReceiverGroup(Vec<std::sync::mpsc::Receiver<(Block<'static>, usize, usize, RouteKey)>>);
struct BufReceiverGroup(Vec<std::sync::mpsc::Receiver<(Vec<u8>, usize, usize, RouteKey)>>);
impl BufSenderGroup {
pub fn send(&mut self, val: (Block<'static>, usize, usize, RouteKey)) -> bool {
pub fn send(&mut self, val: (Vec<u8>, usize, usize, RouteKey)) -> bool {
let index = self.0 % self.1.len();
self.0 = self.0.wrapping_add(1);
self.1[index].send(val).is_ok()
@@ -562,7 +558,7 @@ fn buf_channel_group(size: usize) -> (BufSenderGroup, BufReceiverGroup) {
let mut buf_receiver_group = Vec::with_capacity(size);
for _ in 0..size {
let (buf_sender, buf_receiver) =
std::sync::mpsc::sync_channel::<(Block<'static, Vec<u8>>, usize, usize, RouteKey)>(1);
std::sync::mpsc::sync_channel::<(Vec<u8>, usize, usize, RouteKey)>(1);
buf_sender_group.push(buf_sender);
buf_receiver_group.push(buf_receiver);
}
@@ -595,9 +591,7 @@ impl Channel {
tcp_r
.read_exact(&mut buf[head_reserve..head_reserve + len])
.await?;
handler
.handle(&mut buf, head_reserve, head_reserve + len, key, &context)
.await;
handler.handle(&mut buf, head_reserve, head_reserve + len, key, &context);
}
}
async fn start_tcp(
@@ -680,24 +674,20 @@ impl Channel {
let main_channel = context.inner.main_channel.clone();
let buf_sender = if parallel > 1 {
let (buf_sender, buf_receiver) = buf_channel_group(parallel);
let mut num = 0;
for buf_receiver in buf_receiver.0 {
let context = context.clone();
let handler = handler.clone();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
log::info!("启动异步处理");
runtime.block_on(async move {
std::thread::Builder::new()
.name(format!("recv-handler-{}", num))
.spawn(move || {
while let Ok((mut buf, start, end, route_key)) = buf_receiver.recv() {
handler
.handle(&mut buf, start, end, route_key, &context)
.await;
handler.handle(&mut buf, start, end, route_key, &context);
}
log::warn!("异步处理停止");
});
});
})
.unwrap();
num += 1;
}
Some(buf_sender)
} else {
@@ -720,22 +710,21 @@ impl Channel {
let main_channel_ipv6 = main_channel_ipv6.clone();
let handler = handler.clone();
let buf_sender = buf_sender.clone();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
log::info!("启动udp v6");
runtime.block_on(Self::main_start_(
worker,
context,
UDP_V6_ID,
main_channel_ipv6,
handler,
buf_sender,
head_reserve,
));
});
std::thread::Builder::new()
.name("ipv6-recv".into())
.spawn(move || {
log::info!("启动udp v6");
Self::main_start_(
worker,
context,
UDP_V6_ID,
main_channel_ipv6,
handler,
buf_sender,
head_reserve,
)
})
.unwrap();
}
{
let worker = worker.worker("main_channel_1");
@@ -743,22 +732,21 @@ impl Channel {
let main_channel = main_channel.clone();
let handler = handler.clone();
let buf_sender = buf_sender.clone();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
log::info!("启动udp v4");
runtime.block_on(Self::main_start_(
worker,
context,
UDP_ID,
main_channel,
handler,
buf_sender,
head_reserve,
));
});
std::thread::Builder::new()
.name("ipv4-recv".into())
.spawn(move || {
log::info!("启动udp v4");
Self::main_start_(
worker,
context,
UDP_ID,
main_channel,
handler,
buf_sender,
head_reserve,
)
})
.unwrap();
}
if relay {
worker.stop_wait().await;
@@ -811,7 +799,7 @@ impl Channel {
}
worker.stop_all();
}
async fn main_start_(
fn main_start_(
worker: VntWorker,
context: Context,
id: usize,
@@ -832,15 +820,13 @@ impl Channel {
break;
}
}
handler
.handle(
&mut buf,
head_reserve,
end,
RouteKey::new(id, addr),
&context,
)
.await;
handler.handle(
&mut buf,
head_reserve,
end,
RouteKey::new(id, addr),
&context,
);
}
Err(e) => {
log::error!("udp :{:?}", e);
@@ -849,7 +835,7 @@ impl Channel {
}
}
Some(mut buf_sender) => loop {
let mut buf = POOL.alloc(4096);
let mut buf = vec![0; 4096];
match udp.recv_from(&mut buf[head_reserve..]) {
Ok((len, addr)) => {
let end = head_reserve + len;
@@ -897,7 +883,7 @@ impl Channel {
rs=udp.recv_from(&mut buf[head_reserve..])=>{
match rs {
Ok((len, addr)) => {
handler.handle(&mut buf, head_reserve, head_reserve + len, RouteKey::new(id, addr), &context).await;
handler.handle(&mut buf, head_reserve, head_reserve + len, RouteKey::new(id, addr), &context);
}
Err(e) => {
log::error!("{:?}",e)
@@ -931,7 +917,7 @@ impl Channel {
}
}
Some(mut buf_sender) => loop {
let mut buf = POOL.alloc(4096);
let mut buf = vec![0; 4096];
tokio::select! {
rs=udp.recv_from(&mut buf[head_reserve..])=>{
match rs {
+5 -2
View File
@@ -50,9 +50,12 @@ impl NatInfo {
public_port_range: u16,
local_ipv4_addr: SocketAddrV4,
ipv6_addr: SocketAddrV6,
nat_type: NatType,
mut nat_type: NatType,
) -> Self {
public_ips.retain(|ip| !ip.is_loopback() && !ip.is_private());
public_ips.retain(|ip| !ip.is_loopback() && !ip.is_private() && !ip.is_unspecified());
if public_ips.len() > 1 {
nat_type = NatType::Symmetric;
}
Self {
public_ips,
public_port,
+1 -1
View File
@@ -28,7 +28,7 @@ impl FromStr for CipherModel {
"aes_gcm" => Ok(CipherModel::AesGcm),
"aes_cbc" => Ok(CipherModel::AesCbc),
"aes_ecb" => Ok(CipherModel::AesEcb),
_ => Err(format!("not match '{}'", s)),
_ => Err(format!("not match '{}', enum:aes_gcm/aes_cbc/aes_ecb", s)),
}
}
}
+9 -1
View File
@@ -403,12 +403,20 @@ impl VntUtil {
let device_list = device_list.clone();
let current_device = current_device.clone();
// 定时心跳
heartbeat_handler::start_heartbeat_main(
vnt_status_manager.worker("main-heartbeat"),
channel_sender.clone(),
device_list.clone(),
current_device.clone(),
config.server_address_str,
client_cipher.clone(),
self.server_cipher.clone(),
);
heartbeat_handler::start_heartbeat(
vnt_status_manager.worker("heartbeat"),
channel_sender.clone(),
device_list.clone(),
current_device.clone(),
config.server_address_str,
client_cipher.clone(),
self.server_cipher.clone(),
);
+65 -21
View File
@@ -43,6 +43,28 @@ async fn start_idle_(idle: Idle, sender: ChannelSender) -> io::Result<()> {
}
pub fn start_heartbeat(
mut worker: VntWorker,
sender: ChannelSender,
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
client_cipher: Cipher,
server_cipher: Cipher,
) {
tokio::spawn(async move {
tokio::select! {
_=worker.stop_wait()=>{
return;
}
rs=start_heartbeat_(sender, device_list, current_device,client_cipher,server_cipher)=>{
if let Err(e) = rs {
log::warn!("心跳任务停止:{:?}", e);
}
}
}
worker.stop_all();
});
}
pub fn start_heartbeat_main(
mut worker: VntWorker,
sender: ChannelSender,
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
@@ -56,11 +78,11 @@ pub fn start_heartbeat(
_=worker.stop_wait()=>{
return;
}
rs=start_heartbeat_(sender, device_list, current_device,server_address_str,client_cipher,server_cipher)=>{
rs=start_heartbeat_main_(sender, device_list, current_device,server_address_str,client_cipher,server_cipher)=>{
if let Err(e) = rs {
log::warn!("心跳任务停止:{:?}", e);
log::warn!("心跳任务停止:{:?}", e);
}
}
}
}
worker.stop_all();
});
@@ -97,7 +119,7 @@ fn heartbeat_packet(
net_packet
}
async fn start_heartbeat_(
async fn start_heartbeat_main_(
sender: ChannelSender,
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
@@ -106,26 +128,14 @@ async fn start_heartbeat_(
server_cipher: Cipher,
) -> io::Result<()> {
let mut count = 0;
log::info!("启动心跳任务");
log::info!("启动心跳任务");
loop {
if sender.is_close() {
return Ok(());
}
let mut current_dev = current_device.load();
//如果和服务端使用tcp连接,则维持udp洞的频率要更高些
if (sender.is_main_tcp() && count % 2 == 0) || (!sender.is_main_tcp() && count % 20 == 1) {
let mut packet = NetPacket::new_encrypt([0; 12 + ENCRYPTION_RESERVED])?;
packet.set_version(Version::V1);
packet.set_gateway_flag(true);
packet.set_protocol(Protocol::Control);
packet.set_transport_protocol(control_packet::Protocol::AddrRequest.into());
packet.first_set_ttl(MAX_TTL);
packet.set_source(current_dev.virtual_ip());
packet.set_destination(current_dev.virtual_gateway);
server_cipher.encrypt_ipv4(&mut packet)?;
let _ = sender.send_main_udp(packet.buffer(), current_dev.connect_server);
}
if count % 20 == 19 {
let src = current_dev.virtual_ip();
if count % 40 == 19 {
if let Ok(mut addr) = server_address_str.to_socket_addrs() {
if let Some(addr) = addr.next() {
if addr != current_dev.connect_server {
@@ -143,7 +153,6 @@ async fn start_heartbeat_(
}
}
}
let src = current_dev.virtual_ip();
let server_packet = heartbeat_packet(
MAX_TTL,
&device_list,
@@ -156,6 +165,41 @@ async fn start_heartbeat_(
if let Err(e) = sender.send_main(server_packet.buffer(), current_dev.connect_server) {
log::warn!("connect_server:{:?},e:{:?}", current_dev.connect_server, e);
}
count += 1;
tokio::time::sleep(Duration::from_millis(3000)).await;
}
}
async fn start_heartbeat_(
sender: ChannelSender,
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
client_cipher: Cipher,
server_cipher: Cipher,
) -> io::Result<()> {
let mut count = 0;
log::info!("启动心跳任务");
loop {
if sender.is_close() {
return Ok(());
}
let current_dev = current_device.load();
//如果和服务端使用tcp连接,则维持udp洞的频率要更高些
if (sender.is_main_tcp() && count % 4 == 0) || (!sender.is_main_tcp() && count % 40 == 1) {
let mut packet = NetPacket::new_encrypt([0; 12 + ENCRYPTION_RESERVED])?;
packet.set_version(Version::V1);
packet.set_gateway_flag(true);
packet.set_protocol(Protocol::Control);
packet.set_transport_protocol(control_packet::Protocol::AddrRequest.into());
packet.first_set_ttl(MAX_TTL);
packet.set_source(current_dev.virtual_ip());
packet.set_destination(current_dev.virtual_gateway);
server_cipher.encrypt_ipv4(&mut packet)?;
let _ = sender.send_main_udp(packet.buffer(), current_dev.connect_server);
}
let src = current_dev.virtual_ip();
if count < 7 || count % 7 == 0 {
let mut route_list: Option<Vec<(Ipv4Addr, Vec<Route>)>> = None;
let peer_list = { device_list.lock().1.clone() };
@@ -251,6 +295,6 @@ async fn start_heartbeat_(
}
count += 1;
tokio::time::sleep(Duration::from_millis(5000)).await;
tokio::time::sleep(Duration::from_millis(3000)).await;
}
}
+85 -73
View File
@@ -1,5 +1,6 @@
use std::net::{Ipv4Addr, Ipv6Addr, SocketAddrV4, SocketAddrV6};
use std::sync::Arc;
use std::time::{Duration, Instant};
use crossbeam_utils::atomic::AtomicCell;
use dashmap::DashMap;
@@ -53,6 +54,7 @@ pub struct ChannelDataHandler {
rsa_cipher: Option<RsaCipher>,
relay: bool,
token: String,
time: Arc<AtomicCell<Instant>>,
}
impl ChannelDataHandler {
@@ -93,12 +95,13 @@ impl ChannelDataHandler {
rsa_cipher,
relay,
token,
time: Arc::new(AtomicCell::new(Instant::now())),
}
}
}
impl ChannelDataHandler {
pub async fn handle(
pub fn handle(
&self,
buf: &mut [u8],
start: usize,
@@ -107,14 +110,14 @@ impl ChannelDataHandler {
context: &Context,
) {
assert_eq!(start, 14);
match self.handle0(&mut buf[..end], &route_key, context).await {
match self.handle0(&mut buf[..end], &route_key, context) {
Ok(_) => {}
Err(e) => {
log::warn!("{:?}", e);
}
}
}
async fn handle0(
fn handle0(
&self,
buf: &mut [u8],
route_key: &RouteKey,
@@ -132,6 +135,7 @@ impl ChannelDataHandler {
&& !destination.is_multicast()
&& destination != current_device.broadcast_address;
if current_device.virtual_ip() != destination
&& !net_packet.is_gateway()
&& not_broadcast
&& !destination.is_unspecified()
{
@@ -160,6 +164,14 @@ impl ChannelDataHandler {
== crate::protocol::error_packet::Protocol::NoKey.into()
{
if let Some(rsa_cipher) = &self.rsa_cipher {
let last = self.time.load();
if last.elapsed() < Duration::from_secs(3)
|| self.time.compare_exchange(last, Instant::now()).is_err()
{
//短时间不重复上传服务端密钥
return Ok(());
}
log::warn!("上传服务端密钥");
secret_handshake_req(
context,
current_device.connect_server,
@@ -173,8 +185,7 @@ impl ChannelDataHandler {
//服务端解密
self.server_cipher.decrypt_ipv4(&mut net_packet)?;
let data_len = net_packet.data_len();
self.server_packet_handle(context, current_device, buf, data_len, route_key)
.await?;
self.server_packet_handle(context, current_device, buf, data_len, route_key)?;
}
return Ok(());
}
@@ -274,8 +285,10 @@ impl ChannelDataHandler {
}
_ => {
log::warn!(
"不支持的ip代理Icmp协议:{}",
destination
"不支持的ip代理Icmp协议:{}->{}->{}",
source,
destination,
dest_ip
);
return Err(Error::Warn(
"不支持的ip代理Icmp协议".to_string(),
@@ -284,18 +297,36 @@ impl ChannelDataHandler {
}
}
_ => {
log::warn!("不支持的ip代理ipv4协议:{}", destination);
log::warn!(
"不支持的ip代理ipv4协议{:?}:{}->{}->{}",
ipv4.protocol(),
source,
destination,
ipv4.destination_ip()
);
return Err(Error::Warn(
"不支持的ip代理ipv4协议".to_string(),
));
}
}
} else {
log::warn!("没有ip代理规则:{}", destination);
log::warn!(
"没有ip代理规则{:?}:{}->{}->{}",
ipv4.protocol(),
source,
destination,
ipv4.destination_ip()
);
return Err(Error::Warn("没有ip代理规则".to_string()));
}
} else {
log::warn!("不支持ip代理:{}", destination);
log::warn!(
"不支持ip代理{:?}:{}->{}->{}",
ipv4.protocol(),
source,
destination,
ipv4.destination_ip()
);
return Err(Error::Warn("不支持ip代理".to_string()));
}
}
@@ -313,12 +344,10 @@ impl ChannelDataHandler {
Protocol::Service => {}
Protocol::Error => {}
Protocol::Control => {
self.control(context, current_device, source, net_packet, route_key)
.await?;
self.control(context, current_device, source, net_packet, route_key)?;
}
Protocol::OtherTurn => {
self.other_turn(context, current_device, source, net_packet, route_key)
.await?;
self.other_turn(context, current_device, source, net_packet, route_key)?;
}
Protocol::UnKnow(e) => {
log::info!("不支持的协议:{}", e);
@@ -327,7 +356,7 @@ impl ChannelDataHandler {
Ok(())
}
async fn pong_packet(
fn pong_packet(
&self,
gateway: bool,
metric: u8,
@@ -361,7 +390,7 @@ impl ChannelDataHandler {
}
Ok(())
}
async fn control(
fn control(
&self,
context: &Context,
current_device: CurrentDeviceInfo,
@@ -390,8 +419,7 @@ impl ChannelDataHandler {
source,
pong_packet,
route_key,
)
.await?;
)?;
}
ControlPacket::PunchRequest => {
if self.relay {
@@ -431,22 +459,13 @@ impl ChannelDataHandler {
}
std::net::IpAddr::V6(_) => {}
},
ControlPacket::AddrResponse(addr_packet) => {
if !addr_packet.ipv4().is_multicast()
&& !addr_packet.ipv4().is_broadcast()
&& !addr_packet.ipv4().is_unspecified()
&& !addr_packet.ipv4().is_loopback()
&& !addr_packet.ipv4().is_private()
&& addr_packet.port() != 0
{
self.nat_test
.update_addr(addr_packet.ipv4(), addr_packet.port())
}
}
ControlPacket::AddrResponse(addr_packet) => self
.nat_test
.update_addr(addr_packet.ipv4(), addr_packet.port()),
}
Ok(())
}
async fn other_turn(
fn other_turn(
&self,
context: &Context,
current_device: CurrentDeviceInfo,
@@ -527,12 +546,12 @@ impl ChannelDataHandler {
// let _ = context.try_send_main_udp(packet.buffer(),
// SocketAddr::V4(SocketAddrV4::new(peer_nat_info.local_ip, peer_nat_info.local_port)));
// }
if self.punch(source, peer_nat_info).await {
if self.punch(source, peer_nat_info) {
self.client_cipher.encrypt_ipv4(&mut punch_packet)?;
context.try_send_by_key(punch_packet.buffer(), route_key)?;
}
} else {
self.punch(source, peer_nat_info).await;
self.punch(source, peer_nat_info);
}
}
other_turn_packet::Protocol::Unknown(e) => {
@@ -541,7 +560,7 @@ impl ChannelDataHandler {
}
Ok(())
}
async fn punch(&self, peer_ip: Ipv4Addr, peer_nat_info: NatInfo) -> bool {
fn punch(&self, peer_ip: Ipv4Addr, peer_nat_info: NatInfo) -> bool {
match peer_nat_info.nat_type {
NatType::Symmetric => self
.symmetric_sender
@@ -554,7 +573,7 @@ impl ChannelDataHandler {
/// 处理服务端数据
impl ChannelDataHandler {
async fn server_packet_handle(
fn server_packet_handle(
&self,
context: &Context,
current_device: CurrentDeviceInfo,
@@ -566,16 +585,13 @@ impl ChannelDataHandler {
let source = net_packet.source();
match net_packet.protocol() {
Protocol::Service => {
self.service(context, current_device, net_packet, route_key)
.await?;
self.service(context, current_device, net_packet, route_key)?;
}
Protocol::Error => {
self.error(context, current_device, source, net_packet, route_key)
.await?;
self.error(context, current_device, source, net_packet, route_key)?;
}
Protocol::Control => {
self.control_gateway(context, current_device, net_packet, route_key)
.await?;
self.control_gateway(context, current_device, net_packet, route_key)?;
}
Protocol::IpTurn => {
match ip_turn_packet::Protocol::from(net_packet.transport_protocol()) {
@@ -609,7 +625,7 @@ impl ChannelDataHandler {
}
return Ok(());
}
async fn control_gateway(
fn control_gateway(
&self,
context: &Context,
current_device: CurrentDeviceInfo,
@@ -627,26 +643,16 @@ impl ChannelDataHandler {
net_packet.source(),
pong_packet,
route_key,
)
.await?;
}
ControlPacket::AddrResponse(addr_packet) => {
if addr_packet.port() != 0
&& !addr_packet.ipv4().is_multicast()
&& !addr_packet.ipv4().is_broadcast()
&& !addr_packet.ipv4().is_unspecified()
&& !addr_packet.ipv4().is_loopback()
&& !addr_packet.ipv4().is_private()
{
self.nat_test
.update_addr(addr_packet.ipv4(), addr_packet.port())
}
)?;
}
ControlPacket::AddrResponse(addr_packet) => self
.nat_test
.update_addr(addr_packet.ipv4(), addr_packet.port()),
_ => {}
}
Ok(())
}
async fn service(
fn service(
&self,
context: &Context,
current_device: CurrentDeviceInfo,
@@ -658,23 +664,29 @@ impl ChannelDataHandler {
service_packet::Protocol::RegistrationResponse => {
let response = RegistrationResponse::parse_from_bytes(net_packet.payload())?;
{
if self.nat_test.can_update() {
let context = context.clone();
let nat_test = self.nat_test.clone();
tokio::spawn(async move {
let local_port = context.main_local_ipv4_port().unwrap_or(0);
let local_ipv4_addr = nat::local_ipv4_addr(local_port);
let local_port = context.main_local_ipv6_port().unwrap_or(0);
let ipv6_addr = nat::local_ipv6_addr(local_port);
let nat_info = nat_test
.re_test(
Ipv4Addr::from(response.public_ip),
response.public_port as u16,
local_ipv4_addr,
ipv6_addr,
)
.await;
context.switch(nat_info.nat_type);
std::thread::spawn(move || {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async move {
let local_port = context.main_local_ipv4_port().unwrap_or(0);
let local_ipv4_addr = nat::local_ipv4_addr(local_port);
let local_port = context.main_local_ipv6_port().unwrap_or(0);
let ipv6_addr = nat::local_ipv6_addr(local_port);
let nat_info = nat_test
.re_test(
Ipv4Addr::from(response.public_ip),
response.public_port as u16,
local_ipv4_addr,
ipv6_addr,
)
.await;
context.switch(nat_info.nat_type);
})
});
}
let new_ip = Ipv4Addr::from(response.virtual_ip);
@@ -746,7 +758,7 @@ impl ChannelDataHandler {
}
Ok(())
}
async fn error(
fn error(
&self,
_context: &Context,
current_device: CurrentDeviceInfo,
+1 -1
View File
@@ -229,7 +229,7 @@ impl Register {
}
pub fn fast_register(&self, ip: Ipv4Addr) -> crate::Result<()> {
let last = self.time.load();
if last.elapsed() < Duration::from_secs(2)
if last.elapsed() < Duration::from_secs(3)
|| self.time.compare_exchange(last, Instant::now()).is_err()
{
//短时间不重复注册
+4 -6
View File
@@ -1,15 +1,13 @@
use byte_pool::Block;
#[derive(Clone)]
pub struct BufSenderGroup(
usize,
Vec<std::sync::mpsc::SyncSender<(Block<'static>, usize, usize)>>,
Vec<std::sync::mpsc::SyncSender<(Vec<u8>, usize, usize)>>,
);
pub struct BufReceiverGroup(pub Vec<std::sync::mpsc::Receiver<(Block<'static>, usize, usize)>>);
pub struct BufReceiverGroup(pub Vec<std::sync::mpsc::Receiver<(Vec<u8>, usize, usize)>>);
impl BufSenderGroup {
pub fn send(&mut self, val: (Block<'static>, usize, usize)) -> bool {
pub fn send(&mut self, val: (Vec<u8>, usize, usize)) -> bool {
let index = self.0 % self.1.len();
self.0 = self.0.wrapping_add(1);
self.1[index].send(val).is_ok()
@@ -21,7 +19,7 @@ pub fn buf_channel_group(size: usize) -> (BufSenderGroup, BufReceiverGroup) {
let mut buf_receiver_group = Vec::with_capacity(size);
for _ in 0..size {
let (buf_sender, buf_receiver) =
std::sync::mpsc::sync_channel::<(Block<'static>, usize, usize)>(1);
std::sync::mpsc::sync_channel::<(Vec<u8>, usize, usize)>(1);
buf_sender_group.push(buf_sender);
buf_receiver_group.push(buf_receiver);
}
+4 -6
View File
@@ -1,9 +1,7 @@
use byte_pool::BytePool;
use std::sync::Arc;
use std::{io, thread};
use crossbeam_utils::atomic::AtomicCell;
use lazy_static::lazy_static;
use packet::arp::arp::ArpPacket;
use packet::ethernet;
@@ -22,9 +20,6 @@ use crate::handle::CurrentDeviceInfo;
use crate::igmp_server::IgmpServer;
use crate::ip_proxy::IpProxyMap;
use crate::tun_tap_device::{DeviceReader, DeviceWriter};
lazy_static! {
static ref POOL: BytePool<Vec<u8>> = BytePool::<Vec<u8>>::new();
}
pub fn start(
worker: VntWorker,
@@ -116,7 +111,7 @@ fn start_(
mut buf_sender: BufSenderGroup,
) -> io::Result<()> {
loop {
let mut buf = POOL.alloc(4096);
let mut buf = vec![0; 4096];
if sender.is_close() {
return Ok(());
}
@@ -144,6 +139,9 @@ fn start_simple(
) -> io::Result<()> {
let mut buf = [0; 4096];
loop {
if sender.is_close() {
return Ok(());
}
let len = device_reader.read(&mut buf)?;
if let Err(e) = handle(
&mut buf,
+1 -5
View File
@@ -1,4 +1,3 @@
use byte_pool::BytePool;
use std::sync::Arc;
use std::{io, thread};
@@ -19,9 +18,6 @@ use crate::handle::CurrentDeviceInfo;
use crate::igmp_server::IgmpServer;
use crate::ip_proxy::IpProxyMap;
use crate::tun_tap_device::{DeviceReader, DeviceWriter};
lazy_static::lazy_static! {
static ref POOL:BytePool<Vec<u8>> = BytePool::<Vec<u8>>::new();
}
fn icmp(device_writer: &DeviceWriter, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> Result<()> {
if ipv4_packet.protocol() == ipv4::protocol::Protocol::Icmp {
let mut icmp = IcmpPacket::new(ipv4_packet.payload_mut())?;
@@ -169,7 +165,7 @@ fn start_(
mut buf_sender: BufSenderGroup,
) -> io::Result<()> {
loop {
let mut buf = POOL.alloc(4096);
let mut buf = vec![0; 4096];
buf[..12].fill(0);
if sender.is_close() {
return Ok(());
+41 -11
View File
@@ -1,9 +1,12 @@
use crossbeam_utils::atomic::AtomicCell;
use dashmap::DashMap;
use std::io;
use std::net::{SocketAddr, SocketAddrV4};
use std::sync::Arc;
use std::time::Duration;
use std::time::{Duration, Instant};
use tokio::io::AsyncReadExt;
use tokio::io::AsyncWriteExt;
use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf};
use tokio::net::{TcpListener, TcpStream};
pub struct TcpProxy {
@@ -37,7 +40,7 @@ impl TcpProxy {
Duration::from_secs(5),
TcpStream::connect(dest_addr),
)
.await
.await
{
Ok(peer_tcp_stream) => match peer_tcp_stream {
Ok(peer_tcp_stream) => peer_tcp_stream,
@@ -79,15 +82,42 @@ impl TcpProxy {
}
}
async fn proxy(mut client: TcpStream, mut server: TcpStream) -> io::Result<()> {
let (mut client_reader, mut client_writer) = client.split();
let (mut server_reader, mut server_writer) = server.split();
async fn proxy(client: TcpStream, server: TcpStream) -> io::Result<()> {
let (client_read, client_write) = client.into_split();
let (server_read, server_write) = server.into_split();
let time = Arc::new(AtomicCell::new(Instant::now()));
let time1 = time.clone();
tokio::spawn(async move {
if let Err(e) = copy(client_read, server_write, &time1).await {
log::warn!("{:?}", e);
}
});
copy(server_read, client_write, &time).await
}
let client_to_server = tokio::io::copy(&mut client_reader, &mut server_writer);
let server_to_client = tokio::io::copy(&mut server_reader, &mut client_writer);
tokio::select! {
_ = tokio::time::timeout(Duration::from_secs(10), client_to_server) =>{},
_ = tokio::time::timeout(Duration::from_secs(10), server_to_client) =>{},
async fn copy(
mut read: OwnedReadHalf,
mut write: OwnedWriteHalf,
time: &AtomicCell<Instant>,
) -> io::Result<()> {
let mut buf = [0; 10240];
loop {
tokio::select! {
result = read.read(&mut buf) =>{
let len = result?;
if len==0{
break;
}
write.write_all(&buf[..len]).await?;
time.store(Instant::now());
}
_ = tokio::time::sleep(Duration::from_secs(600)) =>{
if time.load().elapsed()>=Duration::from_secs(580){
//读写均超时再退出
break;
}
}
}
}
Ok(())
}
+21 -9
View File
@@ -1,10 +1,12 @@
use crate::ip_proxy::DashMapNew;
use crossbeam_utils::atomic::AtomicCell;
use dashmap::DashMap;
use std::io;
use std::net::{SocketAddr, SocketAddrV4};
use std::sync::Arc;
use std::time::Duration;
use tokio::net::UdpSocket;
use tokio::time::Instant;
/// 一个udp代理,作用是利用系统协议栈,将udp数据报解析出来再转发到目的地址
pub struct UdpProxy {
@@ -22,7 +24,8 @@ impl UdpProxy {
let udp_socket = self.udp_socket;
let mut buf = [0u8; 65536];
let inner_map: Arc<DashMap<SocketAddrV4, Arc<UdpSocket>>> = Arc::new(DashMap::new0());
let inner_map: Arc<DashMap<SocketAddrV4, (Arc<UdpSocket>, Arc<AtomicCell<Instant>>)>> =
Arc::new(DashMap::new0());
loop {
match udp_socket.recv_from(&mut buf).await {
@@ -49,29 +52,36 @@ impl UdpProxy {
async fn start0(
buf: &[u8],
sender_addr: SocketAddrV4,
inner_map: &Arc<DashMap<SocketAddrV4, Arc<UdpSocket>>>,
inner_map: &Arc<DashMap<SocketAddrV4, (Arc<UdpSocket>, Arc<AtomicCell<Instant>>)>>,
map: &Arc<DashMap<SocketAddrV4, SocketAddrV4>>,
udp_socket: &Arc<UdpSocket>,
) -> io::Result<()> {
if let Some(entry) = inner_map.get(&sender_addr) {
let udp = entry.value().clone();
entry.value().1.store(Instant::now());
let udp = entry.value().0.clone();
drop(entry);
udp.send(buf).await?;
} else if let Some(entry) = map.get(&sender_addr) {
let dest_addr = *entry.value();
drop(entry);
let peer_udp_socket = UdpSocket::bind("0.0.0.0:0").await?;
//先使用相同的端口,冲突了再随机端口
let peer_udp_socket = match UdpSocket::bind(format!("0.0.0.0:{}", sender_addr.port())).await
{
Ok(udp) => udp,
Err(_) => UdpSocket::bind("0.0.0.0:0").await?,
};
peer_udp_socket.connect(dest_addr).await?;
peer_udp_socket.send(buf).await?;
let peer_udp_socket = Arc::new(peer_udp_socket);
let inner_map = inner_map.clone();
inner_map.insert(sender_addr, peer_udp_socket.clone());
let time = Arc::new(AtomicCell::new(Instant::now()));
inner_map.insert(sender_addr, (peer_udp_socket.clone(), time.clone()));
let udp_socket = udp_socket.clone();
let map = map.clone();
tokio::spawn(async move {
let mut buf = [0u8; 65536];
loop {
match tokio::time::timeout(Duration::from_secs(300), peer_udp_socket.recv(&mut buf))
match tokio::time::timeout(Duration::from_secs(600), peer_udp_socket.recv(&mut buf))
.await
{
Ok(rs) => match rs {
@@ -98,9 +108,11 @@ async fn start0(
}
},
Err(_) => {
//超时关闭
log::warn!("udp代理超时关闭,来源:{},目标:{}", sender_addr, dest_addr);
break;
if time.load().elapsed() > Duration::from_secs(580) {
//超时关闭
log::warn!("udp代理超时关闭,来源:{},目标:{}", sender_addr, dest_addr);
break;
}
}
}
}
+1 -1
View File
@@ -1,5 +1,5 @@
use crate::error::Error;
pub const VNT_VERSION: &'static str = "1.2.3";
pub const VNT_VERSION: &'static str = "1.2.4.4";
pub type Result<T> = std::result::Result<T, Error>;
pub mod channel;
+30 -13
View File
@@ -1,7 +1,9 @@
use crossbeam_utils::atomic::AtomicCell;
use std::io;
use std::net::UdpSocket;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddrV4, SocketAddrV6};
use std::sync::Arc;
use std::time::{Duration, Instant};
use parking_lot::Mutex;
@@ -22,7 +24,7 @@ pub fn local_ipv4() -> io::Result<Ipv4Addr> {
pub fn local_ipv6() -> io::Result<Ipv6Addr> {
let socket = UdpSocket::bind("[::]:0")?;
socket.connect("[2001:4860:4860::8888]:80")?;
socket.connect("[2001:4860:4860:0000:0000:0000:0000:8888]:80")?;
let addr = socket.local_addr()?;
match addr.ip() {
IpAddr::V4(_) => Ok(Ipv6Addr::UNSPECIFIED),
@@ -54,6 +56,7 @@ pub fn local_ipv6_addr(port: u16) -> SocketAddrV6 {
pub struct NatTest {
stun_server: Vec<String>,
info: Arc<Mutex<NatInfo>>,
time: Arc<AtomicCell<Instant>>,
}
impl From<NatType> for PunchNatType {
@@ -93,16 +96,33 @@ impl NatTest {
NatType::Cone,
);
let info = Arc::new(Mutex::new(nat_info));
NatTest { stun_server, info }
NatTest {
stun_server,
info,
time: Arc::new(AtomicCell::new(Instant::now())),
}
}
pub fn can_update(&self) -> bool {
let last = self.time.load();
last.elapsed() > Duration::from_secs(10)
&& self.time.compare_exchange(last, Instant::now()).is_ok()
}
pub fn nat_info(&self) -> NatInfo {
self.info.lock().clone()
}
pub fn update_addr(&self, ip: Ipv4Addr, port: u16) {
let mut guard = self.info.lock();
guard.public_port = port;
if !guard.public_ips.contains(&ip) {
guard.public_ips.push(ip);
if !ip.is_multicast()
&& !ip.is_broadcast()
&& !ip.is_unspecified()
&& !ip.is_loopback()
&& !ip.is_private()
&& port != 0
{
let mut guard = self.info.lock();
guard.public_port = port;
if !guard.public_ips.contains(&ip) {
guard.public_ips.push(ip);
}
}
}
pub async fn re_test(
@@ -120,6 +140,7 @@ impl NatTest {
ipv6_addr,
)
.await;
log::info!("探测nat类型={:?}", info);
*self.info.lock() = info.clone();
info
}
@@ -131,13 +152,9 @@ impl NatTest {
ipv6_addr: SocketAddrV6,
) -> NatInfo {
return match stun_test::stun_test_nat(stun_server.clone()).await {
Ok((nat_type, ips, port_range)) => {
let mut public_ips = Vec::new();
public_ips.push(Ipv4Addr::from(public_ip));
for ip in ips {
if ip != public_ip {
public_ips.push(ip);
}
Ok((nat_type, mut public_ips, port_range)) => {
if !public_ips.contains(&public_ip) {
public_ips.push(public_ip)
}
NatInfo::new(
public_ips,
+11 -11
View File
@@ -5,7 +5,6 @@ use bytes::BufMut;
use packet::ethernet;
use parking_lot::Mutex;
use std::net::Ipv4Addr;
use std::os::unix::io::AsRawFd;
#[cfg(any(target_os = "linux"))]
use tun::platform::linux::Device;
#[cfg(any(target_os = "macos"))]
@@ -112,16 +111,17 @@ impl DeviceWriter {
}
}
pub fn close(&self) -> io::Result<()> {
unsafe {
match &self.writer {
DeviceW::Tun(writer) => {
libc::close(writer.as_raw_fd());
}
DeviceW::Tap((writer, _)) => {
libc::close(writer.as_raw_fd());
}
}
}
//早期使用close直接切断网卡,现在并不需要这么做也能正常关闭
// unsafe {
// match &self.writer {
// DeviceW::Tun(writer) => {
// libc::close(writer.as_raw_fd());
// }
// DeviceW::Tap((writer, _)) => {
// libc::close(writer.as_raw_fd());
// }
// }
// }
Ok(())
}
pub fn is_tun(&self) -> bool {
+4 -4
View File
@@ -117,7 +117,7 @@ impl DeviceWriter {
// 当前网段路由
dev.add_route(address, netmask, gateway, 1)?;
// 广播和组播路由
dev.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
// dev.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
dev.add_route(
Ipv4Addr::from([224, 0, 0, 0]),
Ipv4Addr::from([240, 0, 0, 0]),
@@ -230,8 +230,8 @@ fn create_tun(
}
// 当前网段路由
tun_device.add_route(address, netmask, gateway, 1)?;
// 广播和组播路由
tun_device.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
// 广播和组播路由 修改了广播路由会导致发不出广播
// tun_device.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
tun_device.add_route(
Ipv4Addr::from([224, 0, 0, 0]),
Ipv4Addr::from([240, 0, 0, 0]),
@@ -309,7 +309,7 @@ fn create_tap(
tap_device.add_route(*address, *netmask, gateway, 1)?;
}
// 广播和组播路由
tap_device.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
// tap_device.add_route(Ipv4Addr::BROADCAST, Ipv4Addr::BROADCAST, gateway, 1)?;
tap_device.add_route(
Ipv4Addr::from([224, 0, 0, 0]),
Ipv4Addr::from([240, 0, 0, 0]),
+2 -1
View File
@@ -29,5 +29,6 @@ features = [
"winerror",
"ipexport",
"iphlpapi",
"handleapi"
"handleapi",
"ifdef"
]
+2 -1
View File
@@ -119,7 +119,8 @@ impl TapDevice {
}
pub fn delete(self) -> io::Result<()> {
iface::delete_interface(&self.luid)
// iface::delete_interface(&self.luid)
Ok(())
}
}
+12 -6
View File
@@ -1,12 +1,12 @@
use std::io;
use std::net::Ipv4Addr;
use winapi::um::{handleapi, synchapi, winbase, winnt};
use winapi::um::{synchapi, winbase, winnt};
use crate::{decode_utf16, encode_utf16, ffi, netsh, route, IFace};
use rand::Rng;
mod log;
pub mod packet;
mod wintun_log;
mod wintun_raw;
/// The maximum size of wintun's internal ring buffer (in bytes)
@@ -80,7 +80,7 @@ impl TunDevice {
let guid_struct: wintun_raw::GUID = unsafe { std::mem::transmute(guid) };
let guid_ptr = &guid_struct as *const wintun_raw::GUID;
log::set_default_logger_if_unset(&win_tun);
wintun_log::set_default_logger_if_unset(&win_tun);
//SAFETY: the function is loaded from the wintun dll properly, we are providing valid
//pointers, and all the strings are correct null terminated UTF-16. This safety rationale
@@ -88,6 +88,7 @@ impl TunDevice {
let adapter =
win_tun.WintunCreateAdapter(pool_utf16.as_ptr(), name_utf16.as_ptr(), guid_ptr);
if adapter.is_null() {
log::error!("adapter.is_null {:?}", io::Error::last_os_error());
return Err(io::Error::new(
io::ErrorKind::Other,
"Failed to crate adapter",
@@ -102,6 +103,7 @@ impl TunDevice {
// 开启session
let session = win_tun.WintunStartSession(adapter, 128 * 1024);
if session.is_null() {
log::error!("session.is_null {:?}", io::Error::last_os_error());
return Err(io::Error::new(
io::ErrorKind::Other,
"WintunStartSession failed",
@@ -138,10 +140,14 @@ impl TunDevice {
));
}
};
log::set_default_logger_if_unset(&win_tun);
wintun_log::set_default_logger_if_unset(&win_tun);
let name_utf16 = encode_utf16(name);
let adapter = win_tun.WintunOpenAdapter(name_utf16.as_ptr());
if adapter.is_null() {
log::error!(
"delete_for_name adapter.is_null {:?}",
io::Error::last_os_error()
);
return Err(io::Error::new(
io::ErrorKind::Other,
"Failed to open adapter",
@@ -187,8 +193,8 @@ pub struct Version {
impl IFace for TunDevice {
fn shutdown(&self) -> io::Result<()> {
let _ = unsafe { synchapi::SetEvent(self.shutdown_event) };
let _ = unsafe { handleapi::CloseHandle(self.shutdown_event) };
// let _ = unsafe { synchapi::SetEvent(self.shutdown_event) };
// let _ = unsafe { handleapi::CloseHandle(self.shutdown_event) };
Ok(())
}