注册失败时重试

This commit is contained in:
lbl
2026-02-24 16:08:18 +08:00
parent b4301a8106
commit 4708ed2b16
4 changed files with 90 additions and 62 deletions
+23 -4
View File
@@ -1,11 +1,13 @@
use anyhow::Context; use anyhow::{Context, bail};
use args_config::{build_config_from_args_and_file, Args, FileConfig}; use args_config::{Args, FileConfig, build_config_from_args_and_file};
use route_manager::Route; use route_manager::Route;
use std::path::Path; use std::path::Path;
use vnt_ipc as vnt_core; use vnt_ipc as vnt_core;
use vnt_core::core::NetworkManager; use vnt_core::core::NetworkManager;
use vnt_core::utils::task_control::TaskGroupManager; use vnt_core::utils::task_control::TaskGroupManager;
use vnt_ipc::core::RegisterResponse;
pub mod args_config; pub mod args_config;
#[cfg(windows)] #[cfg(windows)]
@@ -83,8 +85,25 @@ async fn main0() -> anyhow::Result<()> {
let mut network_manager = NetworkManager::create_network(Box::new(config), task_group) let mut network_manager = NetworkManager::create_network(Box::new(config), task_group)
.await .await
.context("create network")?; .context("create network")?;
let reg_msg = network_manager.register().await.context("register")?; let reg_msg = loop {
let reg_msg = match network_manager.register().await {
Ok(rs) => rs,
Err(e) => {
log::error!("Register failed: {:?}", e);
tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
continue;
}
};
match reg_msg {
RegisterResponse::Success(reg_msg) => {
break reg_msg;
}
RegisterResponse::Failed(e) => {
log::error!("Register failed: {:?}", e);
bail!("注册失败:{}", e.message)
}
}
};
if !network_manager.is_no_tun() { if !network_manager.is_no_tun() {
log::info!("启动网络:{}/{}", reg_msg.ip, reg_msg.prefix_len); log::info!("启动网络:{}/{}", reg_msg.ip, reg_msg.prefix_len);
network_manager.start_tun().await.context("start tun")?; network_manager.start_tun().await.context("start tun")?;
+29 -31
View File
@@ -9,6 +9,7 @@ use crate::enhanced_tunnel::outbound::EnhancedOutbound;
use crate::fec::{FecDecoder, FecEncoder}; use crate::fec::{FecDecoder, FecEncoder};
use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager}; use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager};
use crate::nat::{AllowSubnetExternalRoute, SubnetExternalRoute}; use crate::nat::{AllowSubnetExternalRoute, SubnetExternalRoute};
use crate::protocol::control_message::ErrorResponseMsg;
use crate::tun::enhanced_tun::EnhancedTunInbound; use crate::tun::enhanced_tun::EnhancedTunInbound;
use crate::tun::{DeviceConfig, DeviceIOManager, TunDataInbound, TunReceiver, tun_channel}; use crate::tun::{DeviceConfig, DeviceIOManager, TunDataInbound, TunReceiver, tun_channel};
use crate::tunnel_core::outbound::{BasicOutbound, HybridOutbound}; use crate::tunnel_core::outbound::{BasicOutbound, HybridOutbound};
@@ -47,6 +48,10 @@ pub struct NetworkManager {
tun_receiver: Option<TunReceiver>, tun_receiver: Option<TunReceiver>,
registration_context: Option<Box<RegistrationContext>>, registration_context: Option<Box<RegistrationContext>>,
} }
pub enum RegisterResponse {
Success(NetworkAddr),
Failed(ErrorResponseMsg),
}
impl NetworkManager { impl NetworkManager {
pub async fn create_network( pub async fn create_network(
@@ -202,47 +207,41 @@ impl NetworkManager {
/// Register with server(s) and start data handling tasks. /// Register with server(s) and start data handling tasks.
/// This method can only be called once. /// This method can only be called once.
/// Returns the registration response on success. /// Returns the registration response on success.
pub async fn register(&mut self) -> anyhow::Result<NetworkAddr> { pub async fn register(&mut self) -> anyhow::Result<RegisterResponse> {
let Some(mut ctx) = self.registration_context.take() else { let Some(mut ctx) = self.registration_context.take() else {
bail!("register can only be called once"); bail!("register can only be called once");
}; };
let is_multi_server = ctx.server_managers.len() > 1; let is_multi_server = ctx.server_managers.len() > 1;
let reg_response = if is_multi_server { let response = if is_multi_server {
// Multi-server: coordinated pre-registration // Multi-server: coordinated pre-registration
log::info!( log::info!(
"Multi-server mode: performing coordinated registration for {} servers", "Multi-server mode: performing coordinated registration for {} servers",
ctx.server_managers.len() ctx.server_managers.len()
); );
let reg_response = coordinated_registration(&mut ctx.server_managers).await?; coordinated_registration(&mut ctx.server_managers).await?
log::info!(
"Coordinated registration completed, IP: {}, prefix_len: {}",
reg_response.ip,
reg_response.prefix_len
);
reg_response
} else { } else {
// Single-server: normal registration // Single-server: normal registration
log::info!("Single-server mode: performing normal registration"); log::info!("Single-server mode: performing normal registration");
let response = ctx.server_managers[0] ctx.server_managers[0]
.connect_and_reg(crate::protocol::control_message::RegistrationMode::Normal) .connect_and_reg(crate::protocol::control_message::RegistrationMode::Normal)
.await?; .await?
match response { };
crate::protocol::control_message::ResponseMessage::Reg(reg) => { let reg_response = match response {
log::info!( crate::protocol::control_message::ResponseMessage::Reg(reg) => {
"Registration completed, IP: {}, prefix_len: {}", log::info!(
reg.ip, "Registration completed, IP: {}, prefix_len: {}",
reg.prefix_len reg.ip,
); reg.prefix_len
reg );
} reg
crate::protocol::control_message::ResponseMessage::Error(e) => { }
bail!("Registration failed: {}", e.message); crate::protocol::control_message::ResponseMessage::Error(e) => {
} return Ok(RegisterResponse::Failed(e));
crate::protocol::control_message::ResponseMessage::ConfirmReg(_) => { }
bail!("Unexpected ConfirmReg response"); crate::protocol::control_message::ResponseMessage::ConfirmReg(_) => {
} bail!("Unexpected ConfirmReg response");
} }
}; };
let network_addr = NetworkAddr { let network_addr = NetworkAddr {
@@ -256,10 +255,9 @@ impl NetworkManager {
// 保存服务器版本信息 // 保存服务器版本信息
if !reg_response.server_version.is_empty() { if !reg_response.server_version.is_empty() {
for (index, _) in ctx.server_managers.iter().enumerate() { for (index, _) in ctx.server_managers.iter().enumerate() {
self.app_state.server_info_collection.set_server_version( self.app_state
index as u32, .server_info_collection
reg_response.server_version.clone(), .set_server_version(index as u32, reg_response.server_version.clone());
);
} }
} }
@@ -283,7 +281,7 @@ impl NetworkManager {
turn_manager.data_handle_task_connected(&self.task_group, handler_config, network_addr); turn_manager.data_handle_task_connected(&self.task_group, handler_config, network_addr);
} }
Ok(network_addr) Ok(RegisterResponse::Success(network_addr))
} }
pub fn is_no_tun(&self) -> bool { pub fn is_no_tun(&self) -> bool {
@@ -6,7 +6,7 @@ use crate::crypto::PacketCrypto;
use crate::enhanced_tunnel::inbound::EnhancedInbound; use crate::enhanced_tunnel::inbound::EnhancedInbound;
use crate::fec::FecDecoder; use crate::fec::FecDecoder;
use crate::protocol::control_message::{ use crate::protocol::control_message::{
ConfirmRegResponseMsg, RegResponseMsg, RegistrationMode, RequestMessage, ResponseMessage, ConfirmRegResponseMsg, RegistrationMode, RequestMessage, ResponseMessage,
}; };
use crate::tunnel_core::p2p::transport::punch::NatPuncher; use crate::tunnel_core::p2p::transport::punch::NatPuncher;
use crate::tunnel_core::server::inbound::ServerTurnInboundHandler; use crate::tunnel_core::server::inbound::ServerTurnInboundHandler;
@@ -117,9 +117,7 @@ impl ServerTurnManager {
let request_msg = RequestMessage::Reg(reg_msg); let request_msg = RequestMessage::Reg(reg_msg);
let encoded = request_msg.encode(); let encoded = request_msg.encode();
self.transport_client self.transport_client.send(encoded.freeze()).await?;
.send(encoded.freeze())
.await?;
let buf = self let buf = self
.transport_client .transport_client
.next_timeout(Duration::from_secs(10)) .next_timeout(Duration::from_secs(10))
@@ -164,8 +162,7 @@ impl ServerTurnManager {
config: Box<InboundHandlerConfig>, config: Box<InboundHandlerConfig>,
initial_response: NetworkAddr, initial_response: NetworkAddr,
) { ) {
let data_handler = let data_handler = ServerTurnInboundHandler::new(self.server_id, initial_response, config);
ServerTurnInboundHandler::new(self.server_id, initial_response, config);
let task_group_ = task_group.clone(); let task_group_ = task_group.clone();
let Some(mut receiver) = self.receiver.take() else { let Some(mut receiver) = self.receiver.take() else {
unreachable!() unreachable!()
@@ -212,10 +209,7 @@ impl ServerTurnManager {
log::info!("已连接服务器:{}", self.config.server_addr); log::info!("已连接服务器:{}", self.config.server_addr);
data_handler.handle_connected(); data_handler.handle_connected();
if let Err(e) = self if let Err(e) = self.data_handle_loop(&mut receiver, &data_handler).await {
.data_handle_loop(&mut receiver, &data_handler)
.await
{
log::error!("Error on data_handle_loop: {:?}", e); log::error!("Error on data_handle_loop: {:?}", e);
} }
already_connected = false; already_connected = false;
@@ -271,7 +265,7 @@ impl ServerTurnManager {
/// 4. Return the registration response /// 4. Return the registration response
pub async fn coordinated_registration( pub async fn coordinated_registration(
managers: &mut Vec<ServerTurnManager>, managers: &mut Vec<ServerTurnManager>,
) -> anyhow::Result<RegResponseMsg> { ) -> anyhow::Result<ResponseMessage> {
if managers.is_empty() { if managers.is_empty() {
bail!("No servers to register"); bail!("No servers to register");
} }
@@ -287,7 +281,10 @@ pub async fn coordinated_registration(
let ip = match &first_response { let ip = match &first_response {
ResponseMessage::Reg(reg) => reg.ip, ResponseMessage::Reg(reg) => reg.ip,
ResponseMessage::Error(e) => bail!("First server registration failed: {}", e.message), ResponseMessage::Error(e) => {
log::info!("First server registration failed: {}", e.message);
return Ok(first_response);
}
_ => bail!("Unexpected response from first server"), _ => bail!("Unexpected response from first server"),
}; };
log::info!("Got IP {} from first server", ip); log::info!("Got IP {} from first server", ip);
@@ -313,7 +310,8 @@ pub async fn coordinated_registration(
log::info!("Server {} pre-registered successfully", i + 1); log::info!("Server {} pre-registered successfully", i + 1);
} }
Ok(ResponseMessage::Error(e)) => { Ok(ResponseMessage::Error(e)) => {
bail!("Server {} registration failed: {}", i + 1, e.message) log::info!("Server {} registration failed: {}", i + 1, e.message);
return Ok(ResponseMessage::Error(e.clone()));
} }
Err(e) => bail!("Server {} registration failed: {}", i + 1, e), Err(e) => bail!("Server {} registration failed: {}", i + 1, e),
_ => bail!("Unexpected response from server {}", i + 1), _ => bail!("Unexpected response from server {}", i + 1),
@@ -339,8 +337,5 @@ pub async fn coordinated_registration(
log::info!("Coordinated registration completed successfully"); log::info!("Coordinated registration completed successfully");
// Return first server's response (contains IP info) // Return first server's response (contains IP info)
match first_response { Ok(first_response)
ResponseMessage::Reg(reg) => Ok(reg),
_ => unreachable!(),
}
} }
+26 -10
View File
@@ -1,5 +1,5 @@
use crate::defer; use crate::defer;
use anyhow::{Context, anyhow}; use anyhow::{Context, anyhow, bail};
use axum::body::Body; use axum::body::Body;
use axum::http::{HeaderMap, HeaderValue, StatusCode, Uri, header}; use axum::http::{HeaderMap, HeaderValue, StatusCode, Uri, header};
use axum::response::IntoResponse; use axum::response::IntoResponse;
@@ -26,7 +26,7 @@ use tokio::net::TcpListener;
use tower_http::cors::{Any, CorsLayer}; use tower_http::cors::{Any, CorsLayer};
use vnt_core::api::VntApi; use vnt_core::api::VntApi;
use vnt_core::context::config::Config as CoreConfig; use vnt_core::context::config::Config as CoreConfig;
use vnt_core::core::{DEFAULT_MTU, NetworkManager}; use vnt_core::core::{DEFAULT_MTU, NetworkManager, RegisterResponse};
use vnt_core::nat::NetInput; use vnt_core::nat::NetInput;
use vnt_core::port_mapping::PortMapping; use vnt_core::port_mapping::PortMapping;
use vnt_core::tls::verifier::CertValidationMode; use vnt_core::tls::verifier::CertValidationMode;
@@ -536,9 +536,10 @@ async fn start_vnt_network(
task_group: vnt_core::utils::task_control::TaskGroup, task_group: vnt_core::utils::task_control::TaskGroup,
task_group_guard: vnt_core::utils::task_control::TaskGroupGuard, task_group_guard: vnt_core::utils::task_control::TaskGroupGuard,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let mut network_manager = NetworkManager::create_network(Box::new(core_config), task_group.clone()) let mut network_manager =
.await NetworkManager::create_network(Box::new(core_config), task_group.clone())
.map_err(|e| anyhow!("Create network failed: {:?}", e))?; .await
.map_err(|e| anyhow!("Create network failed: {:?}", e))?;
let vnt_api = network_manager.vnt_api(); let vnt_api = network_manager.vnt_api();
@@ -562,11 +563,26 @@ async fn start_vnt_network(
state.record_log("连接服务器,执行注册"); state.record_log("连接服务器,执行注册");
log::info!("Registering with server"); log::info!("Registering with server");
let reg_msg = network_manager let reg_msg = loop {
.register() let reg_msg = match network_manager.register().await {
.await Ok(rs) => rs,
.context("Registration failed")?; Err(e) => {
log::error!("Register failed: {:?}", e);
state.record_log(format!("注册失败:{},5秒后重试", e));
tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
continue;
}
};
match reg_msg {
RegisterResponse::Success(reg_msg) => {
break reg_msg;
}
RegisterResponse::Failed(e) => {
log::error!("Register failed: {:?}", e);
bail!("注册失败:{}", e.message)
}
}
};
state.record_log(format!("注册成功 {}/{}", reg_msg.ip, reg_msg.prefix_len)); state.record_log(format!("注册成功 {}/{}", reg_msg.ip, reg_msg.prefix_len));
log::info!("Network Started: {}/{}", reg_msg.ip, reg_msg.prefix_len); log::info!("Network Started: {}/{}", reg_msg.ip, reg_msg.prefix_len);
if !network_manager.is_no_tun() { if !network_manager.is_no_tun() {