Compare commits

...
11 Commits
Author SHA1 Message Date
lubeilin 9badbe180c 不转发来源和目的相同的数据 2024-03-20 12:21:19 +08:00
lubeilin 30b1e71aa1 fmt 2024-03-19 23:54:33 +08:00
lubeilin b36cc352d5 兼容纯ipv4 2024-03-19 23:53:16 +08:00
lubeilin 12d888cefc 忽略跃点设置失败的异常 2024-03-18 21:35:22 +08:00
lubeilin ad1df41029 调整打洞 2024-03-17 15:41:25 +08:00
lubeilin 67498dfc82 增加序列号 2024-03-17 14:28:03 +08:00
lubeilin cf52fdde57 修改线程名称 2024-03-17 14:27:49 +08:00
lubeilin a20082d40b 已支持ipv6 2024-03-14 23:31:40 +08:00
lubeilin 8e29b020e0 Merge branch 'main' into mio 2024-03-14 12:06:47 +08:00
lubeilin 20132861e9 完善日志,修复已知问题 2024-03-14 12:05:35 +08:00
lubeilin a1e9b3c133 增加参数说明 2024-03-13 23:24:36 +08:00
27 changed files with 243 additions and 142 deletions
-1
View File
@@ -213,7 +213,6 @@ sudo pfctl -f /etc/pf.conf -e
### Todo ### Todo
- 桌面UI(测试中) - 桌面UI(测试中)
- 支持Ipv6(1.2.2已支持客户端之间的ipv6,待支持客户端和服务端之间的ipv6通信)
### 常见问题 ### 常见问题
+2 -1
View File
@@ -40,4 +40,5 @@ aes_gcm=["vnt/aes_gcm"]
server_encrypt=["vnt/server_encrypt"] server_encrypt=["vnt/server_encrypt"]
ip_proxy=["vnt/ip_proxy"] ip_proxy=["vnt/ip_proxy"]
[build-dependencies] [build-dependencies]
embed-manifest = "1.4.0" embed-manifest = "1.4.0"
rand = "0.9.0-alpha.0"
+3
View File
@@ -118,6 +118,9 @@ ports:
cmd: false #关闭控制台输入 cmd: false #关闭控制台输入
no_proxy: false #是否关闭内置代理,true为关闭 no_proxy: false #是否关闭内置代理,true为关闭
first_latency: false #是否优先低延迟通道,默认为false,表示优先使用p2p通道 first_latency: false #是否优先低延迟通道,默认为false,表示优先使用p2p通道
device_name: vnt-tun #网卡名称
packet_loss: 0 #指定丢包率 取值0~1之间的数 用于模拟弱网
packet_delay: 0 #指定延迟 单位毫秒 用于模拟弱网
``` ```
或者需要哪个配置就加哪个,当然token是必须的 或者需要哪个配置就加哪个,当然token是必须的
+14 -7
View File
@@ -1,10 +1,17 @@
// use embed_manifest::{embed_manifest, new_manifest}; use rand::Rng;
// use embed_manifest::manifest::ExecutionLevel; use std::fs::File;
use std::io::Write;
fn main() { fn main() {
////强制用管理员运行貌似体验更差了 // 生成随机序列号
// if std::env::var_os("CARGO_CFG_WINDOWS").is_some() { let serial_number = format!(
// embed_manifest(new_manifest("vnt") "{}-{}-{}",
// .requested_execution_level(ExecutionLevel::RequireAdministrator)).expect("unable to embed manifest file"); rand::thread_rng().gen_range(100..1000),
// } rand::thread_rng().gen_range(100..1000),
rand::thread_rng().gen_range(100..1000)
);
let generated_code = format!(r#"pub const SERIAL_NUMBER: &str = "{}";"#, serial_number);
let dest_path = "src/generated_serial_number.rs";
let mut file = File::create(&dest_path).unwrap();
file.write_all(generated_code.as_bytes()).unwrap();
} }
+1
View File
@@ -31,6 +31,7 @@ impl VntCallback for VntHandler {
} }
fn error(&self, info: ErrorInfo) { fn error(&self, info: ErrorInfo) {
log::error!("error {:?}", info);
println!("{}", style(format!("error {}", info)).red()); println!("{}", style(format!("error {}", info)).red());
match info.code { match info.code {
ErrorType::TokenError ErrorType::TokenError
+11 -6
View File
@@ -15,6 +15,7 @@ use vnt::core::{Config, Vnt};
mod command; mod command;
mod config; mod config;
mod console_out; mod console_out;
mod generated_serial_number;
mod root_check; mod root_check;
pub fn app_home() -> io::Result<PathBuf> { pub fn app_home() -> io::Result<PathBuf> {
@@ -336,7 +337,7 @@ fn main() {
(config, cmd) (config, cmd)
}; };
println!("version {}", vnt::VNT_VERSION); println!("version {}", vnt::VNT_VERSION);
println!("Serial:{}", generated_serial_number::SERIAL_NUMBER);
main0(config, cmd); main0(config, cmd);
std::process::exit(0); std::process::exit(0);
} }
@@ -346,11 +347,14 @@ mod callback;
fn main0(config: Config, show_cmd: bool) { fn main0(config: Config, show_cmd: bool) {
let vnt_util = Vnt::new(config, callback::VntHandler {}).unwrap(); let vnt_util = Vnt::new(config, callback::VntHandler {}).unwrap();
let vnt_c = vnt_util.clone(); let vnt_c = vnt_util.clone();
thread::spawn(move || { thread::Builder::new()
if let Err(e) = command::server::CommandServer::new().start(vnt_c) { .name("CommandServer".into())
log::warn!("cmd:{:?}", e); .spawn(move || {
} if let Err(e) = command::server::CommandServer::new().start(vnt_c) {
}); log::warn!("cmd:{:?}", e);
}
})
.expect("CommandServer");
if show_cmd { if show_cmd {
let mut cmd = String::new(); let mut cmd = String::new();
loop { loop {
@@ -406,6 +410,7 @@ fn command(cmd: &str, vnt: &Vnt) -> bool {
fn print_usage(program: &str, _opts: Options) { fn print_usage(program: &str, _opts: Options) {
println!("Usage: {} [options]", program); println!("Usage: {} [options]", program);
println!("version:{}", vnt::VNT_VERSION); println!("version:{}", vnt::VNT_VERSION);
println!("Serial:{}", generated_serial_number::SERIAL_NUMBER);
println!("Options:"); println!("Options:");
println!( println!(
" -k <token> {}", " -k <token> {}",
+17 -21
View File
@@ -28,6 +28,7 @@ impl Context {
is_tcp: bool, is_tcp: bool,
packet_loss_rate: Option<f64>, packet_loss_rate: Option<f64>,
packet_delay: u32, packet_delay: u32,
use_ipv6: bool,
) -> Self { ) -> Self {
let channel_num = main_udp_socket.len(); let channel_num = main_udp_socket.len();
assert_ne!(channel_num, 0, "not channel"); assert_ne!(channel_num, 0, "not channel");
@@ -51,6 +52,7 @@ impl Context {
packet_loss_rate, packet_loss_rate,
packet_delay, packet_delay,
main_index: AtomicUsize::new(0), main_index: AtomicUsize::new(0),
use_ipv6,
}; };
Self { Self {
inner: Arc::new(inner), inner: Arc::new(inner),
@@ -70,7 +72,7 @@ impl Deref for Context {
} }
/// 对称网络增加的udp socket数目,有助于增加打洞成功率 /// 对称网络增加的udp socket数目,有助于增加打洞成功率
pub const SYMMETRIC_CHANNEL_NUM: usize = 64; pub const SYMMETRIC_CHANNEL_NUM: usize = 100;
const PACKET_LOSS_RATE_DENOMINATOR: u32 = 100_0000; const PACKET_LOSS_RATE_DENOMINATOR: u32 = 100_0000;
pub struct ContextInner { pub struct ContextInner {
// 核心udp socket // 核心udp socket
@@ -90,6 +92,7 @@ pub struct ContextInner {
//控制延迟 //控制延迟
packet_delay: u32, packet_delay: u32,
main_index: AtomicUsize, main_index: AtomicUsize,
use_ipv6: bool,
} }
impl ContextInner { impl ContextInner {
@@ -172,15 +175,16 @@ impl ContextInner {
} }
} }
pub fn send_main_udp(&self, index: usize, buf: &[u8], mut addr: SocketAddr) -> io::Result<()> { pub fn send_main_udp(&self, index: usize, buf: &[u8], mut addr: SocketAddr) -> io::Result<()> {
//核心udp socket都是ipv6模式,如果是v4地址则需要转换成v6 if self.use_ipv6 {
//只有服务器地址可能需要这样转换 //如果是v4地址则需要转换成v6
if let SocketAddr::V4(ipv4) = addr { if let SocketAddr::V4(ipv4) = addr {
addr = SocketAddr::V6(SocketAddrV6::new( addr = SocketAddr::V6(SocketAddrV6::new(
ipv4.ip().to_ipv6_mapped(), ipv4.ip().to_ipv6_mapped(),
ipv4.port(), ipv4.port(),
0, 0,
0, 0,
)); ));
}
} }
self.main_udp_socket[index].send_to(buf, addr)?; self.main_udp_socket[index].send_to(buf, addr)?;
Ok(()) Ok(())
@@ -208,17 +212,9 @@ impl ContextInner {
thread::sleep(Duration::from_millis(1)); thread::sleep(Duration::from_millis(1));
} }
} }
pub fn try_send_all_main(&self, buf: &[u8], mut addr: SocketAddr) { pub fn try_send_all_main(&self, buf: &[u8], addr: SocketAddr) {
if let SocketAddr::V4(ipv4) = addr { for index in 0..self.channel_num() {
addr = SocketAddr::V6(SocketAddrV6::new( if let Err(e) = self.send_main_udp(index, buf, addr) {
ipv4.ip().to_ipv6_mapped(),
ipv4.port(),
0,
0,
));
}
for udp in &self.main_udp_socket {
if let Err(e) = udp.send_to(buf, addr) {
log::warn!("{:?},add={:?}", e, addr); log::warn!("{:?},add={:?}", e, addr);
} }
} }
+42 -11
View File
@@ -154,13 +154,31 @@ pub fn init_context(
) -> io::Result<(Context, mio::net::TcpListener)> { ) -> io::Result<(Context, mio::net::TcpListener)> {
assert!(!ports.is_empty(), "not channel"); assert!(!ports.is_empty(), "not channel");
let mut udps = Vec::with_capacity(ports.len()); let mut udps = Vec::with_capacity(ports.len());
//检查系统是否支持ipv6
let use_ipv6 = match socket2::Socket::new(socket2::Domain::IPV6, socket2::Type::DGRAM, None) {
Ok(_) => true,
Err(e) => {
log::warn!("{:?}", e);
false
}
};
for port in &ports { for port in &ports {
//监听v6+v4双栈 //监听v6+v4双栈
let address: SocketAddr = format!("[::]:{}", port).parse().unwrap(); let (socket, address) = if use_ipv6 {
let socket = socket2::Socket::new(socket2::Domain::IPV6, socket2::Type::DGRAM, None)?; let address: SocketAddr = format!("[::]:{}", port).parse().unwrap();
io_convert(socket.set_only_v6(false), |_| { let socket = socket2::Socket::new(socket2::Domain::IPV6, socket2::Type::DGRAM, None)?;
format!("set_only_v6 failed: {}", &address) io_convert(socket.set_only_v6(false), |_| {
})?; format!("set_only_v6 failed: {}", &address)
})?;
(socket, address)
} else {
let address: SocketAddr = format!("0.0.0.0:{}", port).parse().unwrap();
(
socket2::Socket::new(socket2::Domain::IPV4, socket2::Type::DGRAM, None)?,
address,
)
};
io_convert(socket.set_reuse_address(true), |_| { io_convert(socket.set_reuse_address(true), |_| {
format!("set_reuse_address failed: {}", &address) format!("set_reuse_address failed: {}", &address)
})?; })?;
@@ -184,15 +202,24 @@ pub fn init_context(
is_tcp, is_tcp,
packet_loss_rate, packet_loss_rate,
packet_delay, packet_delay,
use_ipv6,
); );
let port = context.main_local_udp_port()?[0]; let port = context.main_local_udp_port()?[0];
//监听v6+v4双栈,tcp通道使用异步io //监听v6+v4双栈,tcp通道使用异步io
let address: SocketAddr = format!("[::]:{}", port).parse().unwrap(); let (socket, address) = if use_ipv6 {
let socket = socket2::Socket::new(socket2::Domain::IPV6, socket2::Type::STREAM, None)?; let address: SocketAddr = format!("[::]:{}", port).parse().unwrap();
io_convert(socket.set_only_v6(false), |_| { let socket = socket2::Socket::new(socket2::Domain::IPV6, socket2::Type::STREAM, None)?;
format!("set_only_v6 failed: {}", &address) io_convert(socket.set_only_v6(false), |_| {
})?; format!("set_only_v6 failed: {}", &address)
})?;
(socket, address)
} else {
let address: SocketAddr = format!("0.0.0.0:{}", port).parse().unwrap();
let socket = socket2::Socket::new(socket2::Domain::IPV4, socket2::Type::STREAM, None)?;
(socket, address)
};
io_convert(socket.set_reuse_address(true), |_| { io_convert(socket.set_reuse_address(true), |_| {
format!("set_reuse_address failed: {}", &address) format!("set_reuse_address failed: {}", &address)
})?; })?;
@@ -200,7 +227,11 @@ pub fn init_context(
if ports[0] == 0 { if ports[0] == 0 {
//端口可能冲突,则使用任意端口 //端口可能冲突,则使用任意端口
log::warn!("监听tcp端口失败 {:?},重试一次", address); log::warn!("监听tcp端口失败 {:?},重试一次", address);
let address: SocketAddr = format!("[::]:{}", 0).parse().unwrap(); let address: SocketAddr = if use_ipv6 {
format!("[::]:{}", 0).parse().unwrap()
} else {
format!("0.0.0.0:{}", port).parse().unwrap()
};
io_convert(socket.bind(&address.into()), |_| { io_convert(socket.bind(&address.into()), |_| {
format!("bind failed: {}", &address) format!("bind failed: {}", &address)
})?; })?;
+18 -17
View File
@@ -226,9 +226,9 @@ impl Punch {
} }
pub fn punch(&mut self, buf: &[u8], id: Ipv4Addr, nat_info: NatInfo) -> io::Result<()> { pub fn punch(&mut self, buf: &[u8], id: Ipv4Addr, nat_info: NatInfo) -> io::Result<()> {
if !self.context.route_table.need_punch(&id) { if !self.context.route_table.need_punch(&id) {
log::info!("已打洞成功,无需打洞:{:?}", id);
return Ok(()); return Ok(());
} }
if self.is_tcp && nat_info.tcp_port != 0 { if self.is_tcp && nat_info.tcp_port != 0 {
//向tcp发起连接 //向tcp发起连接
if let Some(ipv6_addr) = nat_info.local_tcp_ipv6addr() { if let Some(ipv6_addr) = nat_info.local_tcp_ipv6addr() {
@@ -301,25 +301,26 @@ impl Punch {
} }
let start = *self.port_index.entry(id.clone()).or_insert(0); let start = *self.port_index.entry(id.clone()).or_insert(0);
let mut end = start + max_k2; let mut end = start + max_k2;
let mut index = end; if end > self.port_vec.len() {
if end >= self.port_vec.len() {
end = self.port_vec.len(); end = self.port_vec.len();
}
let mut index = start
+ self.punch_symmetric(
&self.port_vec[start..end],
buf,
&nat_info.public_ips,
max_k2,
)?;
if index >= self.port_vec.len() {
index = 0 index = 0
} }
self.punch_symmetric(
&self.port_vec[start..end],
buf,
&nat_info.public_ips,
max_k2,
)?;
self.port_index.insert(id, index); self.port_index.insert(id, index);
} }
NatType::Cone => { NatType::Cone => {
let is_cone = self.context.is_cone(); let is_cone = self.context.is_cone();
for index in 0..channel_num { 'a: for index in 0..nat_info.public_ports.len().min(channel_num) {
let len = nat_info.public_ports.len();
for ip in &nat_info.public_ips { for ip in &nat_info.public_ips {
let port = nat_info.public_ports[index % len]; let port = nat_info.public_ports[index];
if port == 0 || ip.is_unspecified() { if port == 0 || ip.is_unspecified() {
continue; continue;
} }
@@ -334,7 +335,7 @@ impl Punch {
} }
if !is_cone { if !is_cone {
//对称网络数据只发一遍 //对称网络数据只发一遍
break; break 'a;
} }
} }
} }
@@ -348,19 +349,19 @@ impl Punch {
buf: &[u8], buf: &[u8],
ips: &Vec<Ipv4Addr>, ips: &Vec<Ipv4Addr>,
max: usize, max: usize,
) -> io::Result<()> { ) -> io::Result<usize> {
let mut count = 0; let mut count = 0;
for port in ports { for (index, port) in ports.iter().enumerate() {
for pub_ip in ips { for pub_ip in ips {
count += 1; count += 1;
if count == max { if count == max {
return Ok(()); return Ok(index);
} }
let addr = SocketAddr::V4(SocketAddrV4::new(*pub_ip, *port)); let addr = SocketAddr::V4(SocketAddrV4::new(*pub_ip, *port));
self.context.send_main_udp(0, buf, addr)?; self.context.send_main_udp(0, buf, addr)?;
thread::sleep(Duration::from_millis(2)); thread::sleep(Duration::from_millis(2));
} }
} }
Ok(()) Ok(ports.len())
} }
} }
+2 -2
View File
@@ -49,7 +49,7 @@ where
}; };
thread::Builder::new() thread::Builder::new()
.name("tcp读事件处理线程".into()) .name("tcpRead".into())
.spawn(move || { .spawn(move || {
if let Err(e) = tcp_listen0( if let Err(e) = tcp_listen0(
poll, poll,
@@ -173,7 +173,7 @@ fn init_writable_handler(
{ {
let writable_notify = writable_notify.clone(); let writable_notify = writable_notify.clone();
thread::Builder::new() thread::Builder::new()
.name("tcp-writeable-listen".into()) .name("tcpWriteableListen".into())
.spawn(move || { .spawn(move || {
if let Err(e) = tcp_writable_listen(receiver, poll, writable_notify, &context) { if let Err(e) = tcp_writable_listen(receiver, poll, writable_notify, &context) {
log::error!("{:?}", e); log::error!("{:?}", e);
+2 -2
View File
@@ -49,7 +49,7 @@ where
}; };
let accept = AcceptSocketSender::new(waker.clone(), udp_sender); let accept = AcceptSocketSender::new(waker.clone(), udp_sender);
thread::Builder::new() thread::Builder::new()
.name("sub_udp读事件处理线程".into()) .name("subUdp".into())
.spawn(move || { .spawn(move || {
if let Err(e) = sub_udp_listen0(poll, recv_handler, context, waker, udp_receiver) { if let Err(e) = sub_udp_listen0(poll, recv_handler, context, waker, udp_receiver) {
log::error!("{:?}", e); log::error!("{:?}", e);
@@ -153,7 +153,7 @@ where
} }
})?; })?;
thread::Builder::new() thread::Builder::new()
.name("main_udp".into()) .name("mainUdp".into())
.spawn(move || { .spawn(move || {
if let Err(e) = main_udp_listen0(poll, recv_handler, context) { if let Err(e) = main_udp_listen0(poll, recv_handler, context) {
log::error!("{:?}", e); log::error!("{:?}", e);
+2 -2
View File
@@ -19,8 +19,8 @@ pub fn addr_request(
config: BaseConfigInfo, config: BaseConfigInfo,
) { ) {
addr_request0(&context, &current_device_info, &server_cipher, &config); addr_request0(&context, &current_device_info, &server_cipher, &config);
// 9秒发送一次 // 17秒发送一次
let rs = scheduler.timeout(Duration::from_secs(9), |s| { let rs = scheduler.timeout(Duration::from_secs(17), |s| {
addr_request(s, context, current_device_info, server_cipher, config) addr_request(s, context, current_device_info, server_cipher, config)
}); });
if !rs { if !rs {
+2 -2
View File
@@ -12,7 +12,7 @@ use crate::cipher::Cipher;
use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo};
use crate::protocol::body::ENCRYPTION_RESERVED; use crate::protocol::body::ENCRYPTION_RESERVED;
use crate::protocol::control_packet::PingPacket; use crate::protocol::control_packet::PingPacket;
use crate::protocol::{control_packet, NetPacket, Protocol, Version, MAX_TTL}; use crate::protocol::{control_packet, NetPacket, Protocol, Version};
use crate::util::Scheduler; use crate::util::Scheduler;
/// 定时发送心跳包 /// 定时发送心跳包
@@ -216,7 +216,7 @@ fn heartbeat_packet(
net_packet.set_version(Version::V1); net_packet.set_version(Version::V1);
net_packet.set_protocol(Protocol::Control); net_packet.set_protocol(Protocol::Control);
net_packet.set_transport_protocol(control_packet::Protocol::Ping.into()); net_packet.set_transport_protocol(control_packet::Protocol::Ping.into());
net_packet.first_set_ttl(MAX_TTL); net_packet.first_set_ttl(5);
net_packet.set_source(src); net_packet.set_source(src);
net_packet.set_destination(dest); net_packet.set_destination(dest);
let mut ping = PingPacket::new(net_packet.payload_mut())?; let mut ping = PingPacket::new(net_packet.payload_mut())?;
-1
View File
@@ -86,7 +86,6 @@ fn idle_gateway0<Call: VntCallback>(
ErrorType::Disconnect, ErrorType::Disconnect,
format!("connect:{},error:{:?}", cur.connect_server, e), format!("connect:{},error:{:?}", cur.connect_server, e),
)); ));
log::warn!("{:?}", e);
} }
} }
fn idle_route0<Call: VntCallback>( fn idle_route0<Call: VntCallback>(
+51 -16
View File
@@ -11,7 +11,7 @@ use protobuf::Message;
use rand::prelude::SliceRandom; use rand::prelude::SliceRandom;
use crate::channel::context::Context; use crate::channel::context::Context;
use crate::channel::punch::{NatInfo, Punch}; use crate::channel::punch::{NatInfo, NatType, Punch};
use crate::cipher::Cipher; use crate::cipher::Cipher;
use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo}; use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo};
use crate::nat::NatTest; use crate::nat::NatTest;
@@ -24,31 +24,59 @@ use crate::util::Scheduler;
pub struct PunchSender { pub struct PunchSender {
sender_self: SyncSender<(Ipv4Addr, NatInfo)>, sender_self: SyncSender<(Ipv4Addr, NatInfo)>,
sender_peer: SyncSender<(Ipv4Addr, NatInfo)>, sender_peer: SyncSender<(Ipv4Addr, NatInfo)>,
sender_cone_self: SyncSender<(Ipv4Addr, NatInfo)>,
sender_cone_peer: SyncSender<(Ipv4Addr, NatInfo)>,
} }
impl PunchSender { impl PunchSender {
pub fn send(&self, src_peer: bool, ip: Ipv4Addr, info: NatInfo) -> bool { pub fn send(&self, src_peer: bool, ip: Ipv4Addr, info: NatInfo) -> bool {
if src_peer { log::info!(
self.sender_peer.send((ip, info)).is_ok() "发送打洞协商消息,是否对端发起:{},ip:{},info:{:?}",
} else { src_peer,
self.sender_self.send((ip, info)).is_ok() ip,
} info
);
let sender = match info.nat_type {
NatType::Symmetric => {
if src_peer {
&self.sender_peer
} else {
&self.sender_self
}
}
NatType::Cone => {
if src_peer {
&self.sender_cone_peer
} else {
&self.sender_cone_self
}
}
};
sender.try_send((ip, info)).is_ok()
} }
} }
pub struct PunchReceiver { pub struct PunchReceiver {
receiver_peer: Receiver<(Ipv4Addr, NatInfo)>, receiver_peer: Receiver<(Ipv4Addr, NatInfo)>,
receiver_self: Receiver<(Ipv4Addr, NatInfo)>, receiver_self: Receiver<(Ipv4Addr, NatInfo)>,
receiver_cone_peer: Receiver<(Ipv4Addr, NatInfo)>,
receiver_cone_self: Receiver<(Ipv4Addr, NatInfo)>,
} }
pub fn punch_channel() -> (PunchSender, PunchReceiver) { pub fn punch_channel() -> (PunchSender, PunchReceiver) {
let (sender_self, receiver_self) = sync_channel(1); let (sender_self, receiver_self) = sync_channel(1);
let (sender_peer, receiver_peer) = sync_channel(1); let (sender_peer, receiver_peer) = sync_channel(1);
let (sender_cone_peer, receiver_cone_peer) = sync_channel(1);
let (sender_cone_self, receiver_cone_self) = sync_channel(1);
( (
PunchSender { PunchSender {
sender_self, sender_self,
sender_peer, sender_peer,
sender_cone_peer,
sender_cone_self,
}, },
PunchReceiver { PunchReceiver {
receiver_peer, receiver_peer,
receiver_self, receiver_self,
receiver_cone_peer,
receiver_cone_self,
}, },
) )
} }
@@ -72,19 +100,21 @@ pub fn punch(
client_cipher.clone(), client_cipher.clone(),
0, 0,
); );
let receiver_peer = receiver.receiver_peer; let f = |receiver: Receiver<(Ipv4Addr, NatInfo)>| {
let receiver_self = receiver.receiver_self;
{
let punch = punch.clone(); let punch = punch.clone();
let current_device = current_device.clone(); let current_device = current_device.clone();
let client_cipher = client_cipher.clone(); let client_cipher = client_cipher.clone();
thread::spawn(move || { thread::Builder::new()
punch_start(receiver_peer, punch, current_device, client_cipher); .name("punch".into())
}); .spawn(move || {
} punch_start(receiver, punch, current_device, client_cipher);
thread::spawn(move || { })
punch_start(receiver_self, punch, current_device, client_cipher); .expect("punch");
}); };
f(receiver.receiver_peer);
f(receiver.receiver_self);
f(receiver.receiver_cone_peer);
f(receiver.receiver_cone_self);
} }
/// 接收打洞消息,配合对端打洞 /// 接收打洞消息,配合对端打洞
@@ -198,6 +228,11 @@ fn punch0(
&nat_info, &nat_info,
info.virtual_ip, info.virtual_ip,
)?; )?;
log::info!(
"发起打洞协商请求,目标:{:?},{:?}",
info.virtual_ip,
nat_info
);
context.send_default(packet.buffer(), current_device.connect_server)?; context.send_default(packet.buffer(), current_device.connect_server)?;
} }
Ok(()) Ok(())
+19 -16
View File
@@ -25,21 +25,24 @@ fn retrieve_nat_type0(
nat_test: NatTest, nat_test: NatTest,
udp_socket_sender: AcceptSocketSender<Option<Vec<mio::net::UdpSocket>>>, udp_socket_sender: AcceptSocketSender<Option<Vec<mio::net::UdpSocket>>>,
) { ) {
thread::spawn(move || { thread::Builder::new()
if nat_test.can_update() { .name("natTest".into())
let local_ipv4 = nat::local_ipv4(); .spawn(move || {
let local_ipv6 = nat::local_ipv6(); if nat_test.can_update() {
match nat_test.re_test(local_ipv4, local_ipv6) { let local_ipv4 = nat::local_ipv4();
Ok(nat_info) => { let local_ipv6 = nat::local_ipv6();
log::info!("当前nat信息:{:?}", nat_info); match nat_test.re_test(local_ipv4, local_ipv6) {
if let Err(e) = context.switch(nat_info.nat_type, &udp_socket_sender) { Ok(nat_info) => {
log::warn!("{:?}", e); log::info!("当前nat信息:{:?}", nat_info);
if let Err(e) = context.switch(nat_info.nat_type, &udp_socket_sender) {
log::warn!("{:?}", e);
}
} }
} Err(e) => {
Err(e) => { log::warn!("nat re_test {:?}", e);
log::warn!("nat re_test {:?}", e); }
} };
}; }
} })
}); .expect("natTest");
} }
+3 -5
View File
@@ -176,8 +176,8 @@ impl ClientPacketHandler {
net_packet.first_set_ttl(MAX_TTL); net_packet.first_set_ttl(MAX_TTL);
self.client_cipher.encrypt_ipv4(&mut net_packet)?; self.client_cipher.encrypt_ipv4(&mut net_packet)?;
context.send_by_key(net_packet.buffer(), route_key)?; context.send_by_key(net_packet.buffer(), route_key)?;
// let route = Route::from_default_rt(route_key, metric); let route = Route::from_default_rt(route_key, metric);
// context.route_table.add_route_if_absent(source, route); context.route_table.add_route_if_absent(source, route);
} }
ControlPacket::PongPacket(pong_packet) => { ControlPacket::PongPacket(pong_packet) => {
let current_time = crate::handle::now_time() as u16; let current_time = crate::handle::now_time() as u16;
@@ -320,13 +320,11 @@ impl ClientPacketHandler {
punch_packet.set_source(current_device.virtual_ip()); punch_packet.set_source(current_device.virtual_ip());
punch_packet.set_destination(source); punch_packet.set_destination(source);
punch_packet.set_payload(&bytes)?; punch_packet.set_payload(&bytes)?;
log::info!("接收打洞请求={:?}", peer_nat_info); self.client_cipher.encrypt_ipv4(&mut punch_packet)?;
if self.punch_sender.send(true, source, peer_nat_info) { if self.punch_sender.send(true, source, peer_nat_info) {
self.client_cipher.encrypt_ipv4(&mut punch_packet)?;
context.send_by_key(punch_packet.buffer(), route_key)?; context.send_by_key(punch_packet.buffer(), route_key)?;
} }
} else { } else {
log::info!("接收打洞请求回复={:?}", peer_nat_info);
self.punch_sender.send(false, source, peer_nat_info); self.punch_sender.send(false, source, peer_nat_info);
} }
} }
+1 -2
View File
@@ -268,6 +268,7 @@ impl<Call: VntCallback> ServerPacketHandler<Call> {
log::info!("ip发生变化,old:{:?},response={:?}", old, response); log::info!("ip发生变化,old:{:?},response={:?}", old, response);
} }
if let Err(e) = self.device.set_ip(virtual_ip, virtual_netmask) { if let Err(e) = self.device.set_ip(virtual_ip, virtual_netmask) {
log::error!("LocalIpExists {:?}", e);
self.callback.error(ErrorInfo::new_msg( self.callback.error(ErrorInfo::new_msg(
ErrorType::LocalIpExists, ErrorType::LocalIpExists,
format!("set_ip {:?}", e), format!("set_ip {:?}", e),
@@ -421,12 +422,10 @@ impl<Call: VntCallback> ServerPacketHandler<Call> {
self.callback.error(err); self.callback.error(err);
} }
InErrorPacket::IpAlreadyExists => { InErrorPacket::IpAlreadyExists => {
log::error!("IpAlreadyExists");
let err = ErrorInfo::new(ErrorType::IpAlreadyExists); let err = ErrorInfo::new(ErrorType::IpAlreadyExists);
self.callback.error(err); self.callback.error(err);
} }
InErrorPacket::InvalidIp => { InErrorPacket::InvalidIp => {
log::error!("InvalidIp");
let err = ErrorInfo::new(ErrorType::InvalidIp); let err = ErrorInfo::new(ErrorType::InvalidIp);
self.callback.error(err); self.callback.error(err);
} }
+11 -1
View File
@@ -18,7 +18,7 @@ impl PacketHandler for TurnPacketHandler {
fn handle( fn handle(
&self, &self,
mut net_packet: NetPacket<&mut [u8]>, mut net_packet: NetPacket<&mut [u8]>,
_route_key: RouteKey, route_key: RouteKey,
context: &Context, context: &Context,
_current_device: &CurrentDeviceInfo, _current_device: &CurrentDeviceInfo,
) -> std::io::Result<()> { ) -> std::io::Result<()> {
@@ -27,6 +27,16 @@ impl PacketHandler for TurnPacketHandler {
if ttl > 0 { if ttl > 0 {
let destination = net_packet.destination(); let destination = net_packet.destination();
if let Some(route) = context.route_table.route_one(&destination) { if let Some(route) = context.route_table.route_one(&destination) {
if route.addr == route_key.addr {
//防止环路
log::warn!(
"来源和目标相同 {:?},{},{}",
route_key,
net_packet.source(),
net_packet.destination()
);
return Ok(());
}
if route.metric <= ttl { if route.metric <= ttl {
context.send_by_key(net_packet.buffer(), route.route_key())?; context.send_by_key(net_packet.buffer(), route.route_key())?;
} }
+1 -1
View File
@@ -97,7 +97,7 @@ pub fn base_handle(
net_packet.set_version(Version::V1); net_packet.set_version(Version::V1);
net_packet.set_protocol(protocol::Protocol::IpTurn); net_packet.set_protocol(protocol::Protocol::IpTurn);
net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into()); net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into());
net_packet.first_set_ttl(3); net_packet.first_set_ttl(6);
net_packet.set_source(src_ip); net_packet.set_source(src_ip);
net_packet.set_destination(dest_ip); net_packet.set_destination(dest_ip);
if dest_ip == current_device.virtual_gateway { if dest_ip == current_device.virtual_gateway {
+3 -3
View File
@@ -111,7 +111,7 @@ pub fn start(
let client_cipher = client_cipher.clone(); let client_cipher = client_cipher.clone();
let server_cipher = server_cipher.clone(); let server_cipher = server_cipher.clone();
thread::Builder::new() thread::Builder::new()
.name(format!("tun_handler_{}", index)) .name(format!("tunHandler-{}", index))
.spawn(move || { .spawn(move || {
while let Ok((mut buf, len)) = receiver.recv() { while let Ok((mut buf, len)) = receiver.recv() {
#[cfg(not(target_os = "macos"))] #[cfg(not(target_os = "macos"))]
@@ -139,7 +139,7 @@ pub fn start(
})?; })?;
} }
thread::Builder::new() thread::Builder::new()
.name("tun_handler".into()) .name("tunHandlerM".into())
.spawn(move || { .spawn(move || {
if let Err(e) = start_multi(stop_manager, device, sender, &mut up_counter) { if let Err(e) = start_multi(stop_manager, device, sender, &mut up_counter) {
log::warn!("stop:{}", e); log::warn!("stop:{}", e);
@@ -148,7 +148,7 @@ pub fn start(
})?; })?;
} else { } else {
thread::Builder::new() thread::Builder::new()
.name("tun_handler".into()) .name("tunHandlerS".into())
.spawn(move || { .spawn(move || {
if let Err(e) = start_simple( if let Err(e) = start_simple(
stop_manager, stop_manager,
+15 -12
View File
@@ -47,18 +47,21 @@ impl IcmpProxy {
Arc::new(Mutex::new(HashMap::with_capacity(16))); Arc::new(Mutex::new(HashMap::with_capacity(16)));
{ {
let nat_map = nat_map.clone(); let nat_map = nat_map.clone();
thread::spawn(move || { thread::Builder::new()
if let Err(e) = icmp_proxy( .name("icmpProxy".into())
mio_icmp_socket, .spawn(move || {
nat_map, if let Err(e) = icmp_proxy(
context, mio_icmp_socket,
stop_manager, nat_map,
current_device, context,
client_cipher, stop_manager,
) { current_device,
log::warn!("icmp_proxy:{:?}", e); client_cipher,
} ) {
}); log::warn!("icmp_proxy:{:?}", e);
}
})
.expect("icmpProxy");
} }
Ok(Self { Ok(Self {
icmp_socket: Arc::new(std_socket), icmp_socket: Arc::new(std_socket),
+8 -5
View File
@@ -38,11 +38,14 @@ impl TcpProxy {
let port = tcp_listener.local_addr()?.port(); let port = tcp_listener.local_addr()?.port();
{ {
let nat_map = nat_map.clone(); let nat_map = nat_map.clone();
thread::spawn(move || { thread::Builder::new()
if let Err(e) = tcp_proxy(tcp_listener, nat_map, stop_manager) { .name("tcpProxy".into())
log::warn!("tcp_proxy:{:?}", e); .spawn(move || {
} if let Err(e) = tcp_proxy(tcp_listener, nat_map, stop_manager) {
}); log::warn!("tcp_proxy:{:?}", e);
}
})
.expect("tcpProxy");
} }
Ok(Self { port, nat_map }) Ok(Self { port, nat_map })
} }
+8 -5
View File
@@ -40,11 +40,14 @@ impl UdpProxy {
let port = udp.local_addr()?.port(); let port = udp.local_addr()?.port();
{ {
let nat_map = nat_map.clone(); let nat_map = nat_map.clone();
thread::spawn(move || { thread::Builder::new()
if let Err(e) = udp_proxy(udp, nat_map, scheduler, stop_manager) { .name("udpProxy".into())
log::warn!("udp_proxy:{:?}", e); .spawn(move || {
} if let Err(e) = udp_proxy(udp, nat_map, scheduler, stop_manager) {
}); log::warn!("udp_proxy:{:?}", e);
}
})
.expect("udpProxy");
} }
Ok(Self { port, nat_map }) Ok(Self { port, nat_map })
} }
+1 -1
View File
@@ -52,7 +52,7 @@ impl Scheduler {
run(receiver, s_inner); run(receiver, s_inner);
worker.stop_all(); worker.stop_all();
}) })
.unwrap(); .expect("Scheduler");
Ok(s) Ok(s)
} }
pub fn timeout<F>(&self, time: Duration, f: F) -> bool pub fn timeout<F>(&self, time: Duration, f: F) -> bool
+3 -1
View File
@@ -96,7 +96,9 @@ impl Device {
.map_err(|e| io::Error::new(e.kind(), format!("TAP_WIN_IOCTL_GET_MAC,err={:?}", e)))?; .map_err(|e| io::Error::new(e.kind(), format!("TAP_WIN_IOCTL_GET_MAC,err={:?}", e)))?;
let index = ffi::luid_to_index(&luid).map(|index| index as u32)?; let index = ffi::luid_to_index(&luid).map(|index| index as u32)?;
// 设置网卡跃点 // 设置网卡跃点
netsh::set_interface_metric(index, 0)?; if let Err(e) = netsh::set_interface_metric(index, 0) {
log::warn!("{:?}",e);
}
let device = Self { let device = Self {
handle, handle,
index, index,
+3 -1
View File
@@ -117,7 +117,9 @@ impl Device {
win_tun.WintunGetAdapterLUID(adapter, &mut luid as *mut wintun_raw::NET_LUID); win_tun.WintunGetAdapterLUID(adapter, &mut luid as *mut wintun_raw::NET_LUID);
let index = ffi::luid_to_index(&std::mem::transmute(luid)).map(|index| index as u32)?; let index = ffi::luid_to_index(&std::mem::transmute(luid)).map(|index| index as u32)?;
// 设置网卡跃点 // 设置网卡跃点
netsh::set_interface_metric(index, 0)?; if let Err(e) = netsh::set_interface_metric(index, 0) {
log::warn!("{:?}",e);
}
Ok(Self { Ok(Self {
luid: std::mem::transmute(luid), luid: std::mem::transmute(luid),
index, index,