Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
707b07b8d3 | ||
|
|
de5a6971f0 | ||
|
|
056036c4d2 | ||
|
|
1cfb188845 | ||
|
|
73a2c31854 | ||
|
|
92eea536f8 | ||
|
|
00936a923e | ||
|
|
ba69ba78af | ||
|
|
17b206bace | ||
|
|
58d5a4f5da | ||
|
|
2438d14175 | ||
|
|
99f8526799 | ||
|
|
56fcbd64ed | ||
|
|
3766b2b7c1 | ||
|
|
baf0698fe4 | ||
|
|
301938b9fc | ||
|
|
16a37c713a | ||
|
|
d412a769dd | ||
|
|
c8eecc87fd | ||
|
|
6a11db70c8 |
@@ -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,
|
||||
|
||||
@@ -17,7 +17,7 @@ pub enum CommandEnum {
|
||||
|
||||
pub fn command(cmd: CommandEnum) {
|
||||
if let Err(e) = command_(cmd) {
|
||||
println!("cmd: {}", e);
|
||||
println!("cmd: {:?}", e);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -362,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);
|
||||
}
|
||||
});
|
||||
|
||||
+51
-47
@@ -522,10 +522,6 @@ impl Context {
|
||||
pub fn update_read_time(&self, id: &Ipv4Addr, route_key: &RouteKey) {
|
||||
if let Some(mut time) = self.inner.route_table_time.get_mut(&(*route_key, *id)) {
|
||||
*time.value_mut() = Instant::now();
|
||||
} else {
|
||||
self.inner
|
||||
.route_table_time
|
||||
.insert((*route_key, *id), Instant::now());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -595,8 +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);
|
||||
handler.handle(&mut buf, head_reserve, head_reserve + len, key, &context);
|
||||
}
|
||||
}
|
||||
async fn start_tcp(
|
||||
@@ -679,16 +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 || {
|
||||
while let Ok((mut buf, start, end, route_key)) = buf_receiver.recv() {
|
||||
handler
|
||||
.handle(&mut buf, start, end, route_key, &context);
|
||||
}
|
||||
log::warn!("异步处理停止");
|
||||
});
|
||||
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);
|
||||
}
|
||||
log::warn!("异步处理停止");
|
||||
})
|
||||
.unwrap();
|
||||
num += 1;
|
||||
}
|
||||
Some(buf_sender)
|
||||
} else {
|
||||
@@ -711,18 +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 || {
|
||||
log::info!("启动udp v6");
|
||||
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");
|
||||
@@ -730,18 +732,21 @@ impl Channel {
|
||||
let main_channel = main_channel.clone();
|
||||
let handler = handler.clone();
|
||||
let buf_sender = buf_sender.clone();
|
||||
std::thread::spawn(move || {
|
||||
log::info!("启动udp v4");
|
||||
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;
|
||||
@@ -815,14 +820,13 @@ impl Channel {
|
||||
break;
|
||||
}
|
||||
}
|
||||
handler
|
||||
.handle(
|
||||
&mut buf,
|
||||
head_reserve,
|
||||
end,
|
||||
RouteKey::new(id, addr),
|
||||
&context,
|
||||
);
|
||||
handler.handle(
|
||||
&mut buf,
|
||||
head_reserve,
|
||||
end,
|
||||
RouteKey::new(id, addr),
|
||||
&context,
|
||||
);
|
||||
}
|
||||
Err(e) => {
|
||||
log::error!("udp :{:?}", e);
|
||||
@@ -864,11 +868,11 @@ impl Channel {
|
||||
#[cfg(target_os = "windows")]
|
||||
use std::os::windows::io::AsRawSocket;
|
||||
#[cfg(target_os = "windows")]
|
||||
let id = 3 + udp.as_raw_socket() as usize;
|
||||
let id = 3 + udp.as_raw_socket() as usize;
|
||||
#[cfg(any(unix))]
|
||||
use std::os::fd::AsRawFd;
|
||||
#[cfg(any(unix))]
|
||||
let id = 3 + udp.as_raw_fd() as usize;
|
||||
let id = 3 + udp.as_raw_fd() as usize;
|
||||
|
||||
context.insert_udp(id, udp.clone());
|
||||
match buf_sender {
|
||||
|
||||
+9
-1
@@ -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(),
|
||||
);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,6 +95,7 @@ impl ChannelDataHandler {
|
||||
rsa_cipher,
|
||||
relay,
|
||||
token,
|
||||
time: Arc::new(AtomicCell::new(Instant::now())),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
@@ -273,8 +285,10 @@ impl ChannelDataHandler {
|
||||
}
|
||||
_ => {
|
||||
log::warn!(
|
||||
"不支持的ip代理Icmp协议:{}",
|
||||
destination
|
||||
"不支持的ip代理Icmp协议:{}->{}->{}",
|
||||
source,
|
||||
destination,
|
||||
dest_ip
|
||||
);
|
||||
return Err(Error::Warn(
|
||||
"不支持的ip代理Icmp协议".to_string(),
|
||||
@@ -283,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()));
|
||||
}
|
||||
}
|
||||
@@ -632,12 +664,14 @@ 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();
|
||||
std::thread::spawn(move ||{
|
||||
std::thread::spawn(move || {
|
||||
tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all().build().unwrap()
|
||||
.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);
|
||||
|
||||
@@ -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()
|
||||
{
|
||||
//短时间不重复注册
|
||||
|
||||
@@ -139,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,8 +1,9 @@
|
||||
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};
|
||||
@@ -84,15 +85,21 @@ impl TcpProxy {
|
||||
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).await {
|
||||
if let Err(e) = copy(client_read, server_write, &time1).await {
|
||||
log::warn!("{:?}", e);
|
||||
}
|
||||
});
|
||||
copy(server_read, client_write).await
|
||||
copy(server_read, client_write, &time).await
|
||||
}
|
||||
|
||||
async fn copy(mut read: OwnedReadHalf, mut write: OwnedWriteHalf) -> io::Result<()> {
|
||||
async fn copy(
|
||||
mut read: OwnedReadHalf,
|
||||
mut write: OwnedWriteHalf,
|
||||
time: &AtomicCell<Instant>,
|
||||
) -> io::Result<()> {
|
||||
let mut buf = [0; 10240];
|
||||
loop {
|
||||
tokio::select! {
|
||||
@@ -102,9 +109,13 @@ async fn copy(mut read: OwnedReadHalf, mut write: OwnedWriteHalf) -> io::Result<
|
||||
break;
|
||||
}
|
||||
write.write_all(&buf[..len]).await?;
|
||||
time.store(Instant::now());
|
||||
}
|
||||
_ = tokio::time::sleep(Duration::from_secs(300)) =>{
|
||||
break;
|
||||
_ = tokio::time::sleep(Duration::from_secs(600)) =>{
|
||||
if time.load().elapsed()>=Duration::from_secs(580){
|
||||
//读写均超时再退出
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
@@ -1,5 +1,5 @@
|
||||
use crate::error::Error;
|
||||
pub const VNT_VERSION: &'static str = "1.2.4";
|
||||
pub const VNT_VERSION: &'static str = "1.2.4.4";
|
||||
pub type Result<T> = std::result::Result<T, Error>;
|
||||
|
||||
pub mod channel;
|
||||
|
||||
+15
-2
@@ -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,7 +96,16 @@ 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()
|
||||
@@ -128,6 +140,7 @@ impl NatTest {
|
||||
ipv6_addr,
|
||||
)
|
||||
.await;
|
||||
log::info!("探测nat类型={:?}", info);
|
||||
*self.info.lock() = info.clone();
|
||||
info
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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]),
|
||||
|
||||
@@ -29,5 +29,6 @@ features = [
|
||||
"winerror",
|
||||
"ipexport",
|
||||
"iphlpapi",
|
||||
"handleapi"
|
||||
"handleapi",
|
||||
"ifdef"
|
||||
]
|
||||
@@ -119,7 +119,8 @@ impl TapDevice {
|
||||
}
|
||||
|
||||
pub fn delete(self) -> io::Result<()> {
|
||||
iface::delete_interface(&self.luid)
|
||||
// iface::delete_interface(&self.luid)
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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(())
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user