离线时才检测服务器地址

This commit is contained in:
lubeilin
2024-03-24 09:18:12 +08:00
parent ae983f014b
commit 69da6de1ed
2 changed files with 77 additions and 50 deletions
+8 -41
View File
@@ -1,4 +1,3 @@
use std::net::ToSocketAddrs;
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
@@ -16,10 +15,14 @@ pub fn addr_request(
context: Context, context: Context,
current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>, current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>,
server_cipher: Cipher, server_cipher: Cipher,
config: BaseConfigInfo, _config: BaseConfigInfo,
) { ) {
pub_address_request(scheduler, context, current_device_info.clone(), server_cipher); pub_address_request(
domain_request(scheduler,current_device_info,config); scheduler,
context,
current_device_info.clone(),
server_cipher,
);
} }
pub fn pub_address_request( pub fn pub_address_request(
scheduler: &Scheduler, scheduler: &Scheduler,
@@ -36,43 +39,7 @@ pub fn pub_address_request(
log::info!("定时任务停止"); log::info!("定时任务停止");
} }
} }
pub fn domain_request(
scheduler: &Scheduler,
current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>,
config: BaseConfigInfo,
) {
domain_request0(&current_device_info, &config);
// 120秒发送一次
let rs = scheduler.timeout(Duration::from_secs(120), |s| {
domain_request(s, current_device_info, config)
});
if !rs {
log::info!("定时任务停止");
}
}
pub fn domain_request0(
current_device: &AtomicCell<CurrentDeviceInfo>,
config: &BaseConfigInfo,
) {
let mut current_dev = current_device.load();
// 探测服务端地址变化
if let Ok(mut addr) = config.server_addr.to_socket_addrs() {
if let Some(addr) = addr.next() {
if addr != current_dev.connect_server {
let mut tmp = current_dev.clone();
tmp.connect_server = addr;
let rs = current_device.compare_exchange(current_dev, tmp);
current_dev.connect_server = addr;
log::info!(
"服务端地址变化,旧地址:{},新地址:{},替换结果:{}",
current_dev.connect_server,
addr,
rs.is_ok()
);
}
}
}
}
pub fn addr_request0( pub fn addr_request0(
context: &Context, context: &Context,
current_device: &AtomicCell<CurrentDeviceInfo>, current_device: &AtomicCell<CurrentDeviceInfo>,
+69 -9
View File
@@ -1,3 +1,11 @@
use std::io;
use std::net::{SocketAddr, ToSocketAddrs};
use std::sync::Arc;
use std::time::{Duration, Instant};
use crossbeam_utils::atomic::AtomicCell;
use mio::net::TcpStream;
use crate::channel::context::Context; use crate::channel::context::Context;
use crate::channel::idle::{Idle, IdleType}; use crate::channel::idle::{Idle, IdleType};
use crate::channel::sender::AcceptSocketSender; use crate::channel::sender::AcceptSocketSender;
@@ -6,12 +14,6 @@ use crate::handle::handshaker::Handshake;
use crate::handle::{handshaker, BaseConfigInfo, ConnectStatus, CurrentDeviceInfo}; use crate::handle::{handshaker, BaseConfigInfo, ConnectStatus, CurrentDeviceInfo};
use crate::util::Scheduler; use crate::util::Scheduler;
use crate::{ErrorInfo, VntCallback}; use crate::{ErrorInfo, VntCallback};
use crossbeam_utils::atomic::AtomicCell;
use mio::net::TcpStream;
use std::io;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
pub fn idle_route<Call: VntCallback>( pub fn idle_route<Call: VntCallback>(
scheduler: &Scheduler, scheduler: &Scheduler,
@@ -29,6 +31,29 @@ pub fn idle_route<Call: VntCallback>(
} }
} }
pub fn idle_gateway<Call: VntCallback>( pub fn idle_gateway<Call: VntCallback>(
scheduler: &Scheduler,
context: Context,
current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>,
config: BaseConfigInfo,
tcp_socket_sender: AcceptSocketSender<(TcpStream, SocketAddr, Option<Vec<u8>>)>,
call: Call,
connect_count: usize,
handshake: Handshake,
) {
let time = Instant::now();
idle_gateway_(
scheduler,
context,
current_device_info,
config,
tcp_socket_sender,
call,
connect_count,
handshake,
time,
);
}
pub fn idle_gateway_<Call: VntCallback>(
scheduler: &Scheduler, scheduler: &Scheduler,
context: Context, context: Context,
current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>, current_device_info: Arc<AtomicCell<CurrentDeviceInfo>>,
@@ -37,6 +62,7 @@ pub fn idle_gateway<Call: VntCallback>(
call: Call, call: Call,
mut connect_count: usize, mut connect_count: usize,
handshake: Handshake, handshake: Handshake,
mut time: Instant,
) { ) {
idle_gateway0( idle_gateway0(
&context, &context,
@@ -46,9 +72,10 @@ pub fn idle_gateway<Call: VntCallback>(
&call, &call,
&mut connect_count, &mut connect_count,
&handshake, &handshake,
&mut time,
); );
let rs = scheduler.timeout(Duration::from_secs(5), move |s| { let rs = scheduler.timeout(Duration::from_secs(5), move |s| {
idle_gateway( idle_gateway_(
s, s,
context, context,
current_device_info, current_device_info,
@@ -57,6 +84,7 @@ pub fn idle_gateway<Call: VntCallback>(
call, call,
connect_count, connect_count,
handshake, handshake,
time,
) )
}); });
if !rs { if !rs {
@@ -71,6 +99,7 @@ fn idle_gateway0<Call: VntCallback>(
call: &Call, call: &Call,
connect_count: &mut usize, connect_count: &mut usize,
handshake: &Handshake, handshake: &Handshake,
time: &mut Instant,
) { ) {
if let Err(e) = check_gateway_channel( if let Err(e) = check_gateway_channel(
context, context,
@@ -80,6 +109,7 @@ fn idle_gateway0<Call: VntCallback>(
call, call,
connect_count, connect_count,
handshake, handshake,
time,
) { ) {
let cur = current_device.load(); let cur = current_device.load();
call.error(ErrorInfo::new_msg( call.error(ErrorInfo::new_msg(
@@ -113,16 +143,22 @@ fn idle_route0<Call: VntCallback>(
fn check_gateway_channel<Call: VntCallback>( fn check_gateway_channel<Call: VntCallback>(
context: &Context, context: &Context,
current_device: &AtomicCell<CurrentDeviceInfo>, current_device_info: &AtomicCell<CurrentDeviceInfo>,
config: &BaseConfigInfo, config: &BaseConfigInfo,
tcp_socket_sender: &AcceptSocketSender<(TcpStream, SocketAddr, Option<Vec<u8>>)>, tcp_socket_sender: &AcceptSocketSender<(TcpStream, SocketAddr, Option<Vec<u8>>)>,
call: &Call, call: &Call,
count: &mut usize, count: &mut usize,
handshake: &Handshake, handshake: &Handshake,
time: &mut Instant,
) -> io::Result<()> { ) -> io::Result<()> {
let current_device = current_device.load(); let mut current_device = current_device_info.load();
if current_device.status.offline() { if current_device.status.offline() {
*count += 1; *count += 1;
if time.elapsed() < Duration::from_secs(6 * 60) {
// 探测服务器地址
current_device = domain_request0(current_device_info, config);
*time = Instant::now()
}
//需要重连 //需要重连
call.connect(ConnectInfo::new(*count, current_device.connect_server)); call.connect(ConnectInfo::new(*count, current_device.connect_server));
log::info!("发送握手请求,{:?}", config); log::info!("发送握手请求,{:?}", config);
@@ -149,3 +185,27 @@ fn check_gateway_channel<Call: VntCallback>(
} }
Ok(()) Ok(())
} }
pub fn domain_request0(
current_device: &AtomicCell<CurrentDeviceInfo>,
config: &BaseConfigInfo,
) -> CurrentDeviceInfo {
let mut current_dev = current_device.load();
// 探测服务端地址变化
if let Ok(mut addr) = config.server_addr.to_socket_addrs() {
if let Some(addr) = addr.next() {
if addr != current_dev.connect_server {
let mut tmp = current_dev.clone();
tmp.connect_server = addr;
let rs = current_device.compare_exchange(current_dev, tmp);
current_dev.connect_server = addr;
log::info!(
"服务端地址变化,旧地址:{},新地址:{},替换结果:{}",
current_dev.connect_server,
addr,
rs.is_ok()
);
}
}
}
current_dev
}