From 4708ed2b1693a100c6ab145b1ae01742c2c540f7 Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Mon, 23 Feb 2026 20:15:27 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B3=A8=E5=86=8C=E5=A4=B1=E8=B4=A5=E6=97=B6?= =?UTF-8?q?=E9=87=8D=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main_cli.rs | 27 +++++++-- vnt-core/src/core/mod.rs | 60 +++++++++---------- .../tunnel_core/server/connection_manager.rs | 29 ++++----- vnt-web/src/service_http.rs | 36 +++++++---- 4 files changed, 90 insertions(+), 62 deletions(-) diff --git a/src/main_cli.rs b/src/main_cli.rs index 55abef0..81223ee 100644 --- a/src/main_cli.rs +++ b/src/main_cli.rs @@ -1,11 +1,13 @@ -use anyhow::Context; -use args_config::{build_config_from_args_and_file, Args, FileConfig}; +use anyhow::{Context, bail}; +use args_config::{Args, FileConfig, build_config_from_args_and_file}; use route_manager::Route; use std::path::Path; use vnt_ipc as vnt_core; use vnt_core::core::NetworkManager; use vnt_core::utils::task_control::TaskGroupManager; +use vnt_ipc::core::RegisterResponse; + pub mod args_config; #[cfg(windows)] @@ -83,8 +85,25 @@ async fn main0() -> anyhow::Result<()> { let mut network_manager = NetworkManager::create_network(Box::new(config), task_group) .await .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() { log::info!("启动网络:{}/{}", reg_msg.ip, reg_msg.prefix_len); network_manager.start_tun().await.context("start tun")?; diff --git a/vnt-core/src/core/mod.rs b/vnt-core/src/core/mod.rs index e2b562f..1c8d6b7 100644 --- a/vnt-core/src/core/mod.rs +++ b/vnt-core/src/core/mod.rs @@ -9,6 +9,7 @@ use crate::enhanced_tunnel::outbound::EnhancedOutbound; use crate::fec::{FecDecoder, FecEncoder}; use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager}; use crate::nat::{AllowSubnetExternalRoute, SubnetExternalRoute}; +use crate::protocol::control_message::ErrorResponseMsg; use crate::tun::enhanced_tun::EnhancedTunInbound; use crate::tun::{DeviceConfig, DeviceIOManager, TunDataInbound, TunReceiver, tun_channel}; use crate::tunnel_core::outbound::{BasicOutbound, HybridOutbound}; @@ -47,6 +48,10 @@ pub struct NetworkManager { tun_receiver: Option, registration_context: Option>, } +pub enum RegisterResponse { + Success(NetworkAddr), + Failed(ErrorResponseMsg), +} impl NetworkManager { pub async fn create_network( @@ -202,47 +207,41 @@ impl NetworkManager { /// Register with server(s) and start data handling tasks. /// This method can only be called once. /// Returns the registration response on success. - pub async fn register(&mut self) -> anyhow::Result { + pub async fn register(&mut self) -> anyhow::Result { let Some(mut ctx) = self.registration_context.take() else { bail!("register can only be called once"); }; 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 log::info!( "Multi-server mode: performing coordinated registration for {} servers", ctx.server_managers.len() ); - let reg_response = coordinated_registration(&mut ctx.server_managers).await?; - log::info!( - "Coordinated registration completed, IP: {}, prefix_len: {}", - reg_response.ip, - reg_response.prefix_len - ); - reg_response + coordinated_registration(&mut ctx.server_managers).await? } else { // Single-server: 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) - .await?; - match response { - crate::protocol::control_message::ResponseMessage::Reg(reg) => { - log::info!( - "Registration completed, IP: {}, prefix_len: {}", - reg.ip, - reg.prefix_len - ); - reg - } - crate::protocol::control_message::ResponseMessage::Error(e) => { - bail!("Registration failed: {}", e.message); - } - crate::protocol::control_message::ResponseMessage::ConfirmReg(_) => { - bail!("Unexpected ConfirmReg response"); - } + .await? + }; + let reg_response = match response { + crate::protocol::control_message::ResponseMessage::Reg(reg) => { + log::info!( + "Registration completed, IP: {}, prefix_len: {}", + reg.ip, + reg.prefix_len + ); + reg + } + crate::protocol::control_message::ResponseMessage::Error(e) => { + return Ok(RegisterResponse::Failed(e)); + } + crate::protocol::control_message::ResponseMessage::ConfirmReg(_) => { + bail!("Unexpected ConfirmReg response"); } }; let network_addr = NetworkAddr { @@ -256,10 +255,9 @@ impl NetworkManager { // 保存服务器版本信息 if !reg_response.server_version.is_empty() { for (index, _) in ctx.server_managers.iter().enumerate() { - self.app_state.server_info_collection.set_server_version( - index as u32, - reg_response.server_version.clone(), - ); + self.app_state + .server_info_collection + .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); } - Ok(network_addr) + Ok(RegisterResponse::Success(network_addr)) } pub fn is_no_tun(&self) -> bool { diff --git a/vnt-core/src/tunnel_core/server/connection_manager.rs b/vnt-core/src/tunnel_core/server/connection_manager.rs index 22e6298..35c9c51 100644 --- a/vnt-core/src/tunnel_core/server/connection_manager.rs +++ b/vnt-core/src/tunnel_core/server/connection_manager.rs @@ -6,7 +6,7 @@ use crate::crypto::PacketCrypto; use crate::enhanced_tunnel::inbound::EnhancedInbound; use crate::fec::FecDecoder; 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::server::inbound::ServerTurnInboundHandler; @@ -117,9 +117,7 @@ impl ServerTurnManager { let request_msg = RequestMessage::Reg(reg_msg); let encoded = request_msg.encode(); - self.transport_client - .send(encoded.freeze()) - .await?; + self.transport_client.send(encoded.freeze()).await?; let buf = self .transport_client .next_timeout(Duration::from_secs(10)) @@ -164,8 +162,7 @@ impl ServerTurnManager { config: Box, initial_response: NetworkAddr, ) { - let data_handler = - ServerTurnInboundHandler::new(self.server_id, initial_response, config); + let data_handler = ServerTurnInboundHandler::new(self.server_id, initial_response, config); let task_group_ = task_group.clone(); let Some(mut receiver) = self.receiver.take() else { unreachable!() @@ -212,10 +209,7 @@ impl ServerTurnManager { log::info!("已连接服务器:{}", self.config.server_addr); data_handler.handle_connected(); - if let Err(e) = self - .data_handle_loop(&mut receiver, &data_handler) - .await - { + if let Err(e) = self.data_handle_loop(&mut receiver, &data_handler).await { log::error!("Error on data_handle_loop: {:?}", e); } already_connected = false; @@ -271,7 +265,7 @@ impl ServerTurnManager { /// 4. Return the registration response pub async fn coordinated_registration( managers: &mut Vec, -) -> anyhow::Result { +) -> anyhow::Result { if managers.is_empty() { bail!("No servers to register"); } @@ -287,7 +281,10 @@ pub async fn coordinated_registration( let ip = match &first_response { 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"), }; 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); } 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), _ => bail!("Unexpected response from server {}", i + 1), @@ -339,8 +337,5 @@ pub async fn coordinated_registration( log::info!("Coordinated registration completed successfully"); // Return first server's response (contains IP info) - match first_response { - ResponseMessage::Reg(reg) => Ok(reg), - _ => unreachable!(), - } + Ok(first_response) } diff --git a/vnt-web/src/service_http.rs b/vnt-web/src/service_http.rs index c0bcd67..460bad8 100644 --- a/vnt-web/src/service_http.rs +++ b/vnt-web/src/service_http.rs @@ -1,5 +1,5 @@ use crate::defer; -use anyhow::{Context, anyhow}; +use anyhow::{Context, anyhow, bail}; use axum::body::Body; use axum::http::{HeaderMap, HeaderValue, StatusCode, Uri, header}; use axum::response::IntoResponse; @@ -26,7 +26,7 @@ use tokio::net::TcpListener; use tower_http::cors::{Any, CorsLayer}; use vnt_core::api::VntApi; 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::port_mapping::PortMapping; 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_guard: vnt_core::utils::task_control::TaskGroupGuard, ) -> anyhow::Result<()> { - let mut network_manager = NetworkManager::create_network(Box::new(core_config), task_group.clone()) - .await - .map_err(|e| anyhow!("Create network failed: {:?}", e))?; + let mut network_manager = + NetworkManager::create_network(Box::new(core_config), task_group.clone()) + .await + .map_err(|e| anyhow!("Create network failed: {:?}", e))?; let vnt_api = network_manager.vnt_api(); @@ -562,11 +563,26 @@ async fn start_vnt_network( state.record_log("连接服务器,执行注册"); log::info!("Registering with server"); - let reg_msg = network_manager - .register() - .await - .context("Registration failed")?; - + let reg_msg = loop { + let reg_msg = match network_manager.register().await { + Ok(rs) => rs, + 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)); log::info!("Network Started: {}/{}", reg_msg.ip, reg_msg.prefix_len); if !network_manager.is_no_tun() {