From 99358e1e945ebf831e5677dfe1f05a78ad0970fd Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Thu, 24 Sep 2026 22:55:55 +0800 Subject: [PATCH 1/9] refactor subscription listener and network lifecycle --- src/args_config.rs | 8 +- src/cli_subscription.rs | 102 ++ src/main_cli.rs | 637 +------ vnt-core/proto/control_message.proto | 31 +- vnt-core/src/api/mod.rs | 522 +----- vnt-core/src/context/config.rs | 155 +- vnt-core/src/context/mod.rs | 2 +- vnt-core/src/core/change_runtime.rs | 753 +++++++++ vnt-core/src/core/ip_update.rs | 76 + vnt-core/src/core/mod.rs | 683 +++++--- vnt-core/src/enhanced_tunnel/inbound.rs | 4 - vnt-core/src/enhanced_tunnel/mod.rs | 112 +- vnt-core/src/enhanced_tunnel/outbound.rs | 4 - .../src/enhanced_tunnel/quic_over/boot.rs | 43 +- .../enhanced_tunnel/quic_over/quic_client.rs | 4 - vnt-core/src/event_script.rs | 15 +- vnt-core/src/fec/encoder.rs | 5 - vnt-core/src/lib.rs | 2 + vnt-core/src/log_manager.rs | 173 ++ vnt-core/src/managed_config.rs | 688 +++++++- vnt-core/src/nat/mod.rs | 45 +- vnt-core/src/nat/subnet_mapping.rs | 5 - vnt-core/src/network_info.rs | 637 +++++++ vnt-core/src/protocol/client_message.rs | 14 +- vnt-core/src/protocol/control_message.rs | 106 +- vnt-core/src/runtime_config.rs | 442 +---- vnt-core/src/system_subnet_routes.rs | 63 +- vnt-core/src/tun/general.rs | 475 ++---- .../tunnel_core/server/connection_manager.rs | 225 ++- vnt-core/src/tunnel_core/server/inbound.rs | 896 +--------- vnt-core/src/tunnel_core/server/outbound.rs | 16 +- vnt-core/src/tunnel_core/server/rpc.rs | 17 +- .../tunnel_core/server/transport/config.rs | 42 +- vnt-jni/java_example/AndroidVpnExample.java | 268 +-- .../com/vnt/ChangeApplyResult.java | 66 + vnt-jni/java_example/com/vnt/LogEntry.java | 56 + .../java_example/com/vnt/NetworkResult.java | 128 ++ .../java_example/com/vnt/RegisterResult.java | 97 -- .../java_example/com/vnt/RuntimeChange.java | 117 ++ .../java_example/com/vnt/RuntimeEvent.java | 59 + .../com/vnt/TunRebuildListener.java | 6 - .../com/vnt/TunRebuildRequest.java | 45 - vnt-jni/java_example/com/vnt/VntApi.java | 44 +- vnt-jni/java_example/com/vnt/VntConfig.java | 8 +- vnt-jni/java_example/com/vnt/VntManager.java | 5 +- vnt-jni/java_example/com/vnt/VntNetwork.java | 151 +- vnt-jni/src/lib.rs | 1043 +++++------- vnt-web/src/service_http.rs | 1497 ++++++----------- vnt-web/ui/src/App.vue | 10 +- vnt-web/ui/src/api/index.js | 8 + vnt-web/ui/src/components/InstanceCard.vue | 4 +- vnt-web/ui/src/stores/app.js | 15 +- vnt-web/ui/src/stores/startLog.js | 112 +- vnt-web/ui/src/utils/configHelp.js | 4 +- vnt-web/ui/src/views/ConfigEditor.vue | 4 +- vnt-web/ui/src/views/ConfigView.vue | 61 +- 56 files changed, 5387 insertions(+), 5423 deletions(-) create mode 100644 src/cli_subscription.rs create mode 100644 vnt-core/src/core/change_runtime.rs create mode 100644 vnt-core/src/core/ip_update.rs create mode 100644 vnt-core/src/log_manager.rs create mode 100644 vnt-core/src/network_info.rs create mode 100644 vnt-jni/java_example/com/vnt/ChangeApplyResult.java create mode 100644 vnt-jni/java_example/com/vnt/LogEntry.java create mode 100644 vnt-jni/java_example/com/vnt/NetworkResult.java delete mode 100644 vnt-jni/java_example/com/vnt/RegisterResult.java create mode 100644 vnt-jni/java_example/com/vnt/RuntimeChange.java create mode 100644 vnt-jni/java_example/com/vnt/RuntimeEvent.java delete mode 100644 vnt-jni/java_example/com/vnt/TunRebuildListener.java delete mode 100644 vnt-jni/java_example/com/vnt/TunRebuildRequest.java diff --git a/src/args_config.rs b/src/args_config.rs index 36f7d76d..7c0f4aea 100644 --- a/src/args_config.rs +++ b/src/args_config.rs @@ -625,7 +625,7 @@ impl FileConfig { # ================================== # 可选的服务端配置源。本文件中明确填写的字段会覆盖订阅链接下发的同名字段。 -# subscription = "vnt2://join/1/..." +# subscription = "vnt2://join/2/..." # --- 网络配置 --- # 网络编号,相同网络编号的会组在同一个虚拟网 (必填) @@ -776,8 +776,8 @@ mod tests { #[test] fn subscription_uses_sub_cli_flag() { - let args = Args::try_parse_from(["vnt", "--sub", "vnt2://join/1/example"]).unwrap(); - assert_eq!(args.subscription.as_deref(), Some("vnt2://join/1/example")); + let args = Args::try_parse_from(["vnt", "--sub", "vnt2://join/2/example"]).unwrap(); + assert_eq!(args.subscription.as_deref(), Some("vnt2://join/2/example")); } #[test] @@ -787,7 +787,7 @@ mod tests { ) .unwrap(); let local: FileConfig = toml::from_str( - "subscription='vnt2://join/1/example'\nnetwork_code='local-net'\nmtu=1400\nno_punch=false", + "subscription='vnt2://join/2/example'\nnetwork_code='local-net'\nmtu=1400\nno_punch=false", ) .unwrap(); let merged = remote.overlay(local); diff --git a/src/cli_subscription.rs b/src/cli_subscription.rs new file mode 100644 index 00000000..ad8a2ace --- /dev/null +++ b/src/cli_subscription.rs @@ -0,0 +1,102 @@ +use super::{Args, FileConfig, resolve_configuration}; +use anyhow::Context; +use std::sync::Arc; +use vnt_core::api::VntApi; +use vnt_core::log_manager::InstanceLog; +use vnt_core::managed_config::Subscription; +use vnt_core::network_info::{ChangeOutcome, RuntimeChangeManager, RuntimeEvent}; +use vnt_ipc as vnt_core; + +pub(crate) async fn run( + args: Args, + local_file: Option, + subscription: Option, + log: Arc, + ctrl_port: Option, +) -> anyhow::Result<()> { + // 配置管理器:持有订阅连接与组网实例;托管模式在此等待服务端首份配置 + // (本地配置的身份字段已剥离,resolved 时 network_code 仅为占位符, + // 首份信封到达后由 merge_present_config 以信封身份覆盖) + let (local_config, _, _) = + resolve_configuration(&args, local_file.as_ref(), subscription.is_some())?; + let mut manager = RuntimeChangeManager::new(local_config, subscription, log.clone()).await?; + let network = manager.start_device().await?; + log::info!( + "启动网络:{}/{} (设备模式 {})", + network.ip, + network.prefix_len, + manager.device_mode() + ); + let mut ipc = IpcPublisher::new(ctrl_port); + ipc.publish(manager.api()); + loop { + tokio::select! { + result = tokio::signal::ctrl_c() => { + result.context("install Ctrl+C handler")?; + log::info!("Ctrl+c received!"); + break; + } + event = manager.next_event() => { + match event { + Ok(RuntimeEvent::InstanceStopped) => break, + Ok(RuntimeEvent::Changed(change)) => { + match manager.apply_change(&change).await { + Ok(ChangeOutcome::Applied) => { + log::info!("已应用运行期变化(完整入栈路由 {} 条)", change.routes.len()); + // 实例可能被内部重建,刷新 IPC 对外 API + ipc.publish(manager.api()); + } + // 桌面平台的 rebuild/need_fd 均由 apply_change 内部处理 + Ok(ChangeOutcome::Rebuild) | Ok(ChangeOutcome::NeedFd(_)) => {} + Err(error) => { + log::warn!("应用运行期变化失败: {error:#}"); + } + } + } + Err(error) => { + log::warn!("运行期变化监听结束: {error:#}"); + break; + } + } + } + } + } + manager.stop().await; + log::info!("stop network"); + Ok(()) +} + +struct IpcPublisher { + ctrl_port: Option, + sender: Option>, +} + +impl IpcPublisher { + fn new(ctrl_port: Option) -> Self { + Self { + ctrl_port, + sender: None, + } + } + + fn publish(&mut self, api: Option) { + let Some(api) = api else { + return; + }; + if let Some(sender) = &self.sender { + let _ = sender.send(api); + return; + } + if self.ctrl_port == Some(0) { + return; + } + let (sender, receiver) = tokio::sync::watch::channel(api); + self.sender = Some(sender); + let ctrl_port = self.ctrl_port; + tokio::spawn(async move { + if let Err(error) = vnt_ipc::server::run_server_dynamic(ctrl_port, receiver).await { + log::error!("ipc:{error:?}"); + } + }); + } +} \ No newline at end of file diff --git a/src/main_cli.rs b/src/main_cli.rs index 6d2d4674..da6c42d4 100644 --- a/src/main_cli.rs +++ b/src/main_cli.rs @@ -1,41 +1,18 @@ -use anyhow::{Context, bail}; +use anyhow::Context; use args_config::{Args, CtrlConfig, FileConfig, build_config_from_args_and_file}; -use sha2::{Digest, Sha256}; use std::path::Path; -use std::time::Duration; use vnt_ipc as vnt_core; -use vnt_core::api::{ApplyAction, ReconfigureFallback}; -use vnt_core::context::config::{Config, VirtualIp}; -use vnt_core::core::NetworkManager; +use vnt_core::context::config::Config; +use vnt_core::log_manager::InstanceLog; use vnt_core::managed_config::Subscription; -use vnt_core::protocol::control_message::{SubscriptionConfigAck, SubscriptionConfigApplyStatus}; -use vnt_core::utils::task_control::TaskGroupManager; -use vnt_ipc::core::RegisterResponse; pub mod args_config; +mod cli_subscription; #[cfg(windows)] mod extract_wintun_dll; -#[derive(Clone)] -struct SubscriptionContext { - link: Subscription, - remote_toml: String, - revision: u64, - managed_ip: VirtualIp, - managed_device_name: String, -} - -#[derive(Clone)] -struct RollbackState { - config: Config, - signature: String, - subscription_context: Option, - ctrl_port: Option, - failed_revision: u64, -} - #[tokio::main] pub async fn main() { if let Err(error) = main0().await { @@ -65,537 +42,93 @@ async fn main0() -> anyhow::Result<()> { log::info!("loaded config from {path:?}"); } - let subscription = args.subscription.as_deref().or_else(|| { + let subscription = args.subscription.clone().or_else(|| { local_file .as_ref() - .and_then(|file| file.subscription.as_deref()) + .and_then(|file| file.subscription.clone()) }); - let mut subscription_context = if let Some(value) = subscription { - let link = Subscription::parse(value)?; - let Some(envelope) = fetch_subscription_config_with_retry(&link).await? else { - log::info!("启动已取消"); - return Ok(()); - }; - let managed_ip = VirtualIp::new(envelope.managed_ip, envelope.managed_prefix_len) - .context("订阅配置中的 IP 无效")?; - Some(SubscriptionContext { - link, - remote_toml: envelope.toml, - revision: envelope.revision, - managed_ip, - managed_device_name: envelope.managed_device_name, - }) - } else { - None - }; - - let (mut current_config, initial_ctrl, mut current_signature) = - resolve_configuration(&args, local_file.as_ref(), subscription_context.as_ref())?; - validate_configuration(¤t_config)?; - log_configuration(¤t_config); - - let group_manager = TaskGroupManager::new(); - let mut ipc_sender: Option> = None; - let mut ctrl_port = initial_ctrl.ctrl_port; - let mut rollback: Option = None; - let mut failed_revision = None; - let mut pending_error_ack: Option<(u64, String)> = None; - - 'supervisor: loop { - let (task_group, task_group_guard) = group_manager.create_task()?; - let mut network_manager = match NetworkManager::create_network( - Box::new(current_config.clone()), - task_group, - ) - .await - { - Ok(manager) => manager, - Err(error) => { - drop(task_group_guard); - if let Some(previous) = rollback.take() { - let message = format!("create network: {error:#}"); - log::error!( - "应用订阅链接配置 revision {} 失败,恢复上一版本: {message}", - previous.failed_revision - ); - current_config = previous.config; - current_signature = previous.signature; - subscription_context = previous.subscription_context; - ctrl_port = previous.ctrl_port; - failed_revision = Some(previous.failed_revision); - pending_error_ack = Some((previous.failed_revision, message)); - continue 'supervisor; - } - return Err(error).context("create network"); - } - }; - - let reg_msg = loop { - match network_manager.register().await { - Ok(RegisterResponse::Success(response)) => break response, - Ok(RegisterResponse::Failed(error)) => { - let message = format!("注册失败:{}", error.message); - if let Some(previous) = rollback.take() { - log::error!( - "应用订阅链接配置 revision {} 失败,恢复上一版本: {message}", - previous.failed_revision - ); - group_manager.stop(); - network_manager.wait_all_stopped().await; - drop(network_manager); - drop(task_group_guard); - current_config = previous.config; - current_signature = previous.signature; - subscription_context = previous.subscription_context; - ctrl_port = previous.ctrl_port; - failed_revision = Some(previous.failed_revision); - pending_error_ack = Some((previous.failed_revision, message)); - continue 'supervisor; - } - bail!(message) - } - Err(error) => { - log::error!("Register failed: {error:?}"); - tokio::time::sleep(tokio::time::Duration::from_secs(5)).await; - } - } - }; - if network_manager.device_mode().has_device() { - log::info!( - "启动网络:{}/{} ({})", - reg_msg.ip, - reg_msg.prefix_len, - network_manager.device_mode() - ); - if let Err(error) = network_manager.start_device().await.context("start device") { - if let Some(previous) = rollback.take() { - let message = format!("{error:#}"); - log::error!( - "应用订阅链接配置 revision {} 失败,恢复上一版本: {message}", - previous.failed_revision - ); - group_manager.stop(); - network_manager.wait_all_stopped().await; - drop(network_manager); - drop(task_group_guard); - current_config = previous.config; - current_signature = previous.signature; - subscription_context = previous.subscription_context; - ctrl_port = previous.ctrl_port; - failed_revision = Some(previous.failed_revision); - pending_error_ack = Some((previous.failed_revision, message)); - continue 'supervisor; - } - return Err(error); - } - if let Err(error) = network_manager - .set_device_network_ip(reg_msg.ip, reg_msg.prefix_len) - .await - .context("set network ip") - { - if let Some(previous) = rollback.take() { - let message = format!("{error:#}"); - log::error!( - "应用订阅链接配置 revision {} 失败,恢复上一版本: {message}", - previous.failed_revision - ); - group_manager.stop(); - network_manager.wait_all_stopped().await; - drop(network_manager); - drop(task_group_guard); - current_config = previous.config; - current_signature = previous.signature; - subscription_context = previous.subscription_context; - ctrl_port = previous.ctrl_port; - failed_revision = Some(previous.failed_revision); - pending_error_ack = Some((previous.failed_revision, message)); - continue 'supervisor; - } - return Err(error); - } - } else { - log::info!( - "启动网络:{}/{} (无虚拟网卡)", - reg_msg.ip, - reg_msg.prefix_len - ); - } - - let api = network_manager.vnt_api(); - if subscription_context.is_some() && !api.has_verified_config_server() { - log::warn!("当前服务器不支持此订阅链接的实时同步;组网将继续运行"); - } - if let Some(sender) = &ipc_sender { - let _ = sender.send(api.clone()); - } else if ctrl_port.is_none_or(|port| port != 0) { - let (sender, receiver) = tokio::sync::watch::channel(api.clone()); - ipc_sender = Some(sender); - tokio::spawn(async move { - if let Err(error) = vnt_ipc::server::run_server_dynamic(ctrl_port, receiver).await { - log::error!("ipc:{error:?}"); - } - }); - } - - rollback.take(); - if let Some((revision, error)) = pending_error_ack.take() { - acknowledge( - &api, - revision, - SubscriptionConfigApplyStatus::SubscriptionConfigError, - Some(error), - ) - .await; - } - - tokio::select! { - _ = network_manager.wait_all_stopped() => { - break 'supervisor; - } - _ = tokio::signal::ctrl_c() => { - log::info!("Ctrl+c received!"); - group_manager.stop(); - break 'supervisor; - } - update = wait_subscription_update(api.clone()), if subscription_context.is_some() => { - let Some(update) = update else { continue; }; - let context = subscription_context.as_ref().expect("subscription branch"); - if update.revision <= context.revision { - continue; - } - if failed_revision == Some(update.revision) { - continue; - } - if failed_revision.is_some_and(|revision| update.revision > revision) { - failed_revision = None; - } - let candidate_subscription = SubscriptionContext { - link: context.link.clone(), - remote_toml: update.toml, - revision: update.revision, - managed_ip: match VirtualIp::new( - update.managed_ip, - update.managed_prefix_len, - ) { - Ok(value) => value, - Err(error) => { - log::error!( - "订阅链接配置 revision {} 中的 IP 无效: {error:#}", - update.revision - ); - acknowledge( - &api, - update.revision, - SubscriptionConfigApplyStatus::SubscriptionConfigError, - Some(error.to_string()), - ).await; - continue; - } - }, - managed_device_name: update.managed_device_name, - }; - let resolved = resolve_configuration( - &args, - local_file.as_ref(), - Some(&candidate_subscription), - ).and_then(|(config, ctrl, signature)| { - validate_configuration(&config)?; - Ok((config, ctrl, signature)) - }); - let (candidate, candidate_ctrl, candidate_signature) = match resolved { - Ok(value) => value, - Err(error) => { - log::error!("订阅链接配置 revision {} 校验失败: {error:#}", update.revision); - acknowledge( - &api, - update.revision, - SubscriptionConfigApplyStatus::SubscriptionConfigError, - Some(error.to_string()), - ).await; - continue; - } - }; - - if candidate_signature == current_signature { - subscription_context = Some(candidate_subscription); - acknowledge_applied( - &api, - update.revision, - ¤t_config, - ¤t_signature, - "NO_CHANGE", - Vec::new(), - ).await; - continue; - } - if candidate_ctrl.ctrl_port != ctrl_port { - log::warn!("订阅链接更新了 ctrl_port;现有控制监听端口会保持到本进程退出"); - } - match api.reconfigure(Box::new(candidate.clone())).await { - Ok(report) if matches!( - report.action, - ApplyAction::NoChange | ApplyAction::Live | ApplyAction::ComponentReload - ) => { - log::info!( - "订阅链接配置 revision {} 已在线应用({:?}): {}", - update.revision, - report.action, - report.changed_fields.join(", ") - ); - subscription_context = Some(candidate_subscription); - current_config = candidate; - current_signature = candidate_signature; - ctrl_port = candidate_ctrl.ctrl_port; - acknowledge_applied( - &api, - update.revision, - ¤t_config, - ¤t_signature, - match report.action { - ApplyAction::NoChange => "NO_CHANGE", - ApplyAction::Live => "LIVE", - ApplyAction::ComponentReload => "COMPONENT_RELOAD", - ApplyAction::InstanceRestart => unreachable!(), - }, - report.changed_fields, - ).await; - continue; - } - Err(error) if error.fallback == ReconfigureFallback::KeepCurrent => { - log::error!( - "订阅链接配置 revision {} 在线应用失败,旧配置保持运行: {error:#}", - update.revision - ); - acknowledge( - &api, - update.revision, - SubscriptionConfigApplyStatus::SubscriptionConfigError, - Some(error.to_string()), - ).await; - continue; - } - Ok(_) | Err(_) => {} - } - acknowledge( - &api, - update.revision, - SubscriptionConfigApplyStatus::SubscriptionConfigStaged, - None, - ).await; - log::info!("收到订阅链接配置 revision {},正在进程内重建网络实例", update.revision); - rollback = Some(RollbackState { - config: current_config.clone(), - signature: current_signature.clone(), - subscription_context: subscription_context.clone(), - ctrl_port, - failed_revision: update.revision, - }); - subscription_context = Some(candidate_subscription); - current_config = candidate; - current_signature = candidate_signature; - ctrl_port = candidate_ctrl.ctrl_port; - group_manager.stop(); - network_manager.wait_all_stopped().await; - drop(network_manager); - drop(task_group_guard); - continue 'supervisor; - } - } - } - - log::info!("stop network"); - Ok(()) -} - -/// Fetches the initial managed configuration until it succeeds or the user cancels startup. -async fn fetch_subscription_config_with_retry( - link: &Subscription, -) -> anyhow::Result> { - let mut attempts = 0_u64; - loop { - let result = tokio::select! { - _ = tokio::signal::ctrl_c() => return Ok(None), - result = link.fetch() => result, - }; - match result { - Ok(envelope) => return Ok(Some(envelope)), - Err(error) => { - attempts += 1; - log::error!("获取订阅配置失败(第 {attempts} 次):{error:#};5 秒后重试"); - tokio::select! { - _ = tokio::signal::ctrl_c() => return Ok(None), - _ = tokio::time::sleep(Duration::from_secs(5)) => {}, - } - } - } - } + // 实例日志:网络核心与订阅监听器的错误写入其中(CLI 无查看界面,仅记录) + let instance_log = std::sync::Arc::new(InstanceLog::new("cli")); + let subscription = subscription + .as_deref() + .map(Subscription::parse) + .transpose()?; + let ctrl_port = args + .ctrl_port + .or_else(|| local_file.as_ref().and_then(|file| file.ctrl_port)); + cli_subscription::run(args, local_file, subscription, instance_log, ctrl_port).await } +/// 解析本地配置。`managed` 为 true 时身份字段(network_code/device_id/ip/ +/// device_name)一律让位给订阅信封:本地与 CLI 的同类值全部剥除,缺失的 +/// network_code 以空占位符通过必填校验。实例只会在首份信封合并后才启动, +/// 合并规则始终以信封身份为准(见 `merge_present_config`)。 fn resolve_configuration( args: &Args, local_file: Option<&FileConfig>, - subscription_context: Option<&SubscriptionContext>, + managed: bool, ) -> anyhow::Result<(Config, CtrlConfig, String)> { - let mut remote = if let Some(context) = subscription_context { - parse_remote_config(&context.remote_toml)? - } else { - FileConfig::default() - }; - remote.subscription = None; - remote.event_script = None; - // Managed identity belongs to the subscription link, never to either TOML - // layer. Clear both inputs before merging, then inject it below. - remote.network_code = None; - remote.device_id = None; - remote.ip = None; - remote.device_name = None; + let mut effective_args = args.clone(); let mut merged = if let Some(local) = local_file { - let mut local = local.clone(); - local.network_code = None; - local.device_id = None; - local.ip = None; - local.device_name = None; - remote.overlay(local) + local.clone() } else { - remote + FileConfig::default() }; - let mut effective_args = args.clone(); - if let Some(context) = subscription_context { - effective_args.network_code = Some(context.link.network_code.clone()); - effective_args.device_id = Some(context.link.device_id.clone()); + if managed { + effective_args.network_code = None; + effective_args.device_id = None; effective_args.ip = None; effective_args.device_name = None; - // Also populate the file layer so all configuration builders observe - // the same effective identity if their precedence changes later. - merged.network_code = Some(context.link.network_code.clone()); - merged.device_id = Some(context.link.device_id.clone()); + merged.network_code = None; + merged.device_id = None; + merged.ip = None; + merged.device_name = None; + // 占位身份:首份订阅信封到达后由 merge_present_config 覆盖 + if effective_args.network_code.is_none() && merged.network_code.is_none() { + effective_args.network_code = Some(String::new()); + } } - let (mut config, ctrl) = build_config_from_args_and_file(Some(effective_args), Some(merged)) + let (config, ctrl) = build_config_from_args_and_file(Some(effective_args), Some(merged)) .context("invalid configuration")?; - if let Some(context) = subscription_context { - config.network_code.clone_from(&context.link.network_code); - config.device_id.clone_from(&context.link.device_id); - config.ip = Some(context.managed_ip); - config.device_name.clone_from(&context.managed_device_name); - config.managed = Some(context.link.registration(context.revision)); - } let mut comparable = config.clone(); comparable.managed = None; let signature = format!("{comparable:?}|ctrl_port={:?}", ctrl.ctrl_port); Ok((config, ctrl, signature)) } -fn parse_remote_config(content: &str) -> anyhow::Result { - let value: toml::Value = toml::from_str(content).context("服务端配置 TOML 无效")?; - let table = value.as_table().context("服务端配置 TOML 根节点必须是表")?; - for forbidden in ["subscription", "event_script", "config_name"] { - if table.contains_key(forbidden) { - bail!("服务端禁止下发字段 '{forbidden}'"); - } - } - toml::from_str(content).context("服务端配置字段或类型无效") -} - -fn validate_configuration(config: &Config) -> anyhow::Result<()> { - let mut config = config.clone(); - config.normalize()?; - config.check() -} +#[cfg(test)] +mod tests { + use super::{Args, FileConfig, Subscription, resolve_configuration}; + use clap::Parser; -async fn acknowledge( - api: &vnt_core::api::VntApi, - revision: u64, - status: SubscriptionConfigApplyStatus, - error: Option, -) { - let _ = api - .acknowledge_subscription_config(SubscriptionConfigAck::new( - revision, - status, - error.unwrap_or_default(), - Vec::new(), - )) - .await; -} + /// v:2 载荷:{v, server, cert_mode, join_id, credential_key} + const SUBSCRIPTION: &str = "vnt2://join/2/eyJ2IjoyLCJzZXJ2ZXIiOiJ0Y3A6Ly8xMjcuMC4wLjE6Mjk4NzIiLCJjZXJ0X21vZGUiOiJzdGFuZGFyZCIsImpvaW5faWQiOiIxMTExMTExMS0yMjIyLTMzMzMtNDQ0NC01NTU1NTU1NTU1NTUiLCJjcmVkZW50aWFsX2tleSI6IkFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUEifQ"; -async fn acknowledge_applied( - api: &vnt_core::api::VntApi, - revision: u64, - config: &Config, - effective_signature: &str, - apply_mode: &str, - changed_fields: Vec, -) { - let mut ack = SubscriptionConfigAck::new( - revision, - SubscriptionConfigApplyStatus::SubscriptionConfigApplied, - String::new(), - Vec::new(), - ); - ack.apply_mode = apply_mode.to_string(); - ack.changed_fields = changed_fields; - ack.effective_device_name = config.device_name.clone(); - if let Some(ip) = config.ip { - ack.effective_ip = ip.ip(); - ack.effective_prefix_len = ip.prefix_len().into(); - } else if let Some(network) = api.network() { - ack.effective_ip = network.ip; - ack.effective_prefix_len = network.prefix_len.into(); + #[test] + fn managed_mode_allows_missing_local_identity() { + let args = Args::try_parse_from(["vnt2_cli"]).unwrap(); + let (config, _, _) = resolve_configuration(&args, None, true).unwrap(); + // 占位身份:首份订阅信封到达后由 merge_present_config 覆盖 + assert!(config.network_code.is_empty()); } - ack.effective_output = config.output.clone(); - ack.allow_ikev2 = config.allow_ikev2; - ack.allow_wireguard = config.allow_wireguard; - ack.allow_mapping = config.allow_port_mapping; - ack.effective_config_sha256 = Sha256::digest(effective_signature.as_bytes()).to_vec(); - let _ = api.acknowledge_subscription_config(ack).await; -} -async fn wait_subscription_update( - api: vnt_core::api::VntApi, -) -> Option { - loop { - let updates = api.next_subscription_config_updates().await?; - if let Some(update) = updates.into_iter().max_by_key(|value| value.revision) { - return Some(update); - } + #[test] + fn ordinary_mode_still_requires_network_code() { + let args = Args::try_parse_from(["vnt2_cli"]).unwrap(); + let error = resolve_configuration(&args, None, false) + .err() + .expect("ordinary mode must fail without network_code"); + assert!(format!("{error:#}").contains("network_code")); } -} - -fn log_configuration(config: &Config) { - log::info!( - "server: {}", - config - .server_addr - .iter() - .map(ToString::to_string) - .collect::>() - .join(", ") - ); - log::info!("network code: {}", config.network_code); - log::info!("device id: {}", config.device_id); - log::info!("device name: {}", config.device_name); - log::info!("cert mode: {}", config.cert_mode); -} - -#[cfg(test)] -mod tests { - use super::{Args, FileConfig, Subscription, SubscriptionContext, resolve_configuration}; - use clap::Parser; - - const SUBSCRIPTION: &str = "vnt2://join/1/eyJ2IjoxLCJzZXJ2ZXIiOiJ0Y3A6Ly8xMjcuMC4wLjE6Mjk4NzIiLCJjZXJ0X21vZGUiOiJzdGFuZGFyZCIsIm5ldHdvcmtfY29kZSI6Im1hbmFnZWQtbmV0IiwiZGV2aWNlX2lkIjoibWFuYWdlZC1kZXYiLCJjcmVkZW50aWFsX2tleSI6IkFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUEifQ"; #[test] - fn subscription_identity_overrides_remote_local_and_cli_values() { + fn managed_mode_ignores_local_and_cli_identity_fields() { let args = Args::try_parse_from([ "vnt2_cli", "--network-code", "cli-net", "--device-id", "cli-dev", - "--ip", - "10.20.0.2/24", - "--device-name", - "cli-name", ]) .unwrap(); let local: FileConfig = toml::from_str( @@ -604,57 +137,23 @@ network_code = "local-net" device_id = "local-dev" ip = "10.20.0.3/24" device_name = "local-name" +mtu = 1400 "#, ) .unwrap(); - let context = SubscriptionContext { - link: Subscription::parse(SUBSCRIPTION).unwrap(), - remote_toml: r#" -server = ["tcp://127.0.0.1:29872"] -cert_mode = "standard" -network_code = "remote-net" -device_id = "remote-dev" -ip = "10.20.0.4/24" -device_name = "remote-name" -"# - .to_string(), - revision: 7, - managed_ip: "10.20.0.9/24".parse().unwrap(), - managed_device_name: "managed-name".to_string(), - }; - - let (config, _, _) = resolve_configuration(&args, Some(&local), Some(&context)).unwrap(); - assert_eq!(config.network_code, "managed-net"); - assert_eq!(config.device_id, "managed-dev"); - assert_eq!(config.ip, Some("10.20.0.9/24".parse().unwrap())); - assert_eq!(config.device_name, "managed-name"); - let managed = config.managed.unwrap(); - assert_eq!(managed.network_code, "managed-net"); - assert_eq!(managed.device_id, "managed-dev"); + let (config, _, _) = resolve_configuration(&args, Some(&local), true).unwrap(); + assert!(config.network_code.is_empty()); + // 非身份字段仍然来自本地配置 + assert_eq!(config.mtu, Some(1400)); } #[test] - fn subscription_envelope_metadata_changes_effective_signature() { - let args = Args::try_parse_from(["vnt2_cli"]).unwrap(); - let mut context = SubscriptionContext { - link: Subscription::parse(SUBSCRIPTION).unwrap(), - remote_toml: r#" -server = ["tcp://127.0.0.1:29872"] -cert_mode = "standard" -"# - .to_string(), - revision: 7, - managed_ip: "10.20.0.8/24".parse().unwrap(), - managed_device_name: "managed-name".to_string(), - }; - - let (_, _, old_signature) = resolve_configuration(&args, None, Some(&context)).unwrap(); - context.revision = 8; - context.managed_ip = "10.20.0.9/24".parse().unwrap(); - let (config, _, new_signature) = - resolve_configuration(&args, None, Some(&context)).unwrap(); - - assert_eq!(config.ip, Some("10.20.0.9/24".parse().unwrap())); - assert_ne!(old_signature, new_signature); + fn subscription_link_parses_v2_payload() { + let subscription = Subscription::parse(SUBSCRIPTION).unwrap(); + assert_eq!( + subscription.join_id, + "11111111-2222-3333-4444-555555555555" + ); + assert_eq!(subscription.server, "tcp://127.0.0.1:29872"); } } diff --git a/vnt-core/proto/control_message.proto b/vnt-core/proto/control_message.proto index 5e038082..7db7a24d 100644 --- a/vnt-core/proto/control_message.proto +++ b/vnt-core/proto/control_message.proto @@ -73,14 +73,35 @@ message SubscriptionConfigV1 { } message SubscriptionConfigFetchRequest { - string network_code = 1; - string device_id = 2; + string join_id = 1; + bytes client_nonce = 3; + bytes client_proof = 4; + bytes instance_id = 5; + uint64 applied_revision = 6; +} + +// Opens the long-lived subscription control connection. Unlike RegRequestMsg +// this request never creates a virtual-network/traffic session. +message SubscriptionRegisterRequest { + string join_id = 1; bytes client_nonce = 3; bytes client_proof = 4; bytes instance_id = 5; uint64 applied_revision = 6; } +message SubscriptionRegisterResponse { + SubscriptionConfigEnvelope config = 1; +} + +message SubscriptionPing { + uint64 nonce = 1; +} + +message SubscriptionPong { + uint64 nonce = 1; +} + message SubscriptionConfigEnvelope { uint64 revision = 1; SubscriptionConfigV1 config = 2; @@ -153,6 +174,9 @@ message RequestMessage{ ConfirmRegMsg confirm_reg = 2; FastRegRequestMsg fast_reg = 3; SubscriptionConfigFetchRequest subscription_config = 4; + SubscriptionRegisterRequest subscription_register = 5; + SubscriptionConfigAck subscription_ack = 6; + SubscriptionPing subscription_ping = 7; } } message ResponseMessage{ @@ -162,6 +186,9 @@ message ResponseMessage{ ConfirmRegResponseMsg confirm_reg = 3; FastRegResponseMsg fast_reg = 4; SubscriptionConfigEnvelope subscription_config = 5; + SubscriptionRegisterResponse subscription_register = 6; + SubscriptionConfigEnvelope subscription_push = 7; + SubscriptionPong subscription_pong = 8; } } diff --git a/vnt-core/src/api/mod.rs b/vnt-core/src/api/mod.rs index d2191f8f..d159a4a6 100644 --- a/vnt-core/src/api/mod.rs +++ b/vnt-core/src/api/mod.rs @@ -1,21 +1,12 @@ -use crate::context::config::{Config, VirtualIp}; +use crate::context::config::Config; use crate::context::{ AppState, NetworkAddr, PacketLossInfo, ServerNodeInfo, TrafficInfo, TunnelListenAddr, }; -use crate::core::DEFAULT_MTU; use crate::nat::NetInput; -use crate::nat::advertised_subnets; -use crate::protocol::client_message::{NodeIdentityTemplate, SharedNodeIdentity}; use crate::protocol::control_message::{ ClientSimpleInfo, SubscriptionConfigAck, SubscriptionConfigEnvelope, }; -use crate::runtime_config::{PortMappingReload, RuntimeConfigController}; -#[cfg(not(target_os = "android"))] -use crate::tun::DeviceNetworkUpdateError; use crate::tunnel_core::p2p::route_table::Route; -#[cfg(target_os = "android")] -use crate::tunnel_core::server::inbound::AndroidTunRebuildError; -use crate::tunnel_core::server::inbound::{IpUpdateContext, ManagedNetworkApply}; use crate::tunnel_core::server::rpc::ServerRPC; use anyhow::Context; use ipnet::Ipv4Net; @@ -34,453 +25,16 @@ pub struct ApiNodeInfo { pub struct VntApi { app_state: AppState, server_rpc: ServerRPC, - ip_update: IpUpdateContext, - node_identity: SharedNodeIdentity, - runtime_config: RuntimeConfigController, -} - -#[derive(Clone, Copy, Debug, Eq, PartialEq, serde::Serialize)] -#[serde(rename_all = "SCREAMING_SNAKE_CASE")] -pub enum ApplyAction { - NoChange, - Live, - ComponentReload, - InstanceRestart, -} - -#[derive(Clone, Debug, Eq, PartialEq, serde::Serialize)] -pub struct ReconfigureReport { - pub action: ApplyAction, - pub changed_fields: Vec, -} - -#[derive(Clone, Copy, Debug, Eq, PartialEq, serde::Serialize)] -#[serde(rename_all = "SCREAMING_SNAKE_CASE")] -pub enum ReconfigureStage { - Validate, - Prepare, - Commit, - Rollback, -} - -#[derive(Clone, Copy, Debug, Eq, PartialEq, serde::Serialize)] -#[serde(rename_all = "SCREAMING_SNAKE_CASE")] -pub enum ReconfigureFallback { - KeepCurrent, - InstanceRestart, -} - -#[derive(Debug, serde::Serialize)] -pub struct ReconfigureError { - pub stage: ReconfigureStage, - pub fallback: ReconfigureFallback, - pub message: String, -} - -impl std::fmt::Display for ReconfigureError { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(formatter, "{}", self.message) - } -} - -impl std::error::Error for ReconfigureError {} - -impl ReconfigureError { - fn new( - stage: ReconfigureStage, - fallback: ReconfigureFallback, - error: impl std::fmt::Display, - ) -> Self { - Self { - stage, - fallback, - message: error.to_string(), - } - } -} - -fn config_diff(current: &Config, candidate: &Config) -> Vec { - let mut fields = Vec::new(); - macro_rules! changed { - ($field:ident, $name:expr) => { - if current.$field != candidate.$field { - fields.push($name.to_string()); - } - }; - ($field:ident) => { - changed!($field, stringify!($field)); - }; - } - changed!(server_addr, "server"); - changed!(peer_address); - changed!(turn); - changed!(punch_model); - changed!(cert_mode); - changed!(device_name); - changed!(tun_name); - changed!(outbound_interface); - changed!(ip); - changed!(password); - changed!(no_punch); - changed!(no_broadcast); - changed!(allow_ikev2); - changed!(allow_wireguard); - changed!(compress); - changed!(rtx); - changed!(fec); - changed!(input); - changed!(subnet_mapping); - changed!(output); - changed!(auto_sync_subnet); - changed!(no_nat); - changed!(device_mode); - if current.mtu.unwrap_or(DEFAULT_MTU) != candidate.mtu.unwrap_or(DEFAULT_MTU) { - fields.push("mtu".to_string()); - } - changed!(port_mapping); - changed!(allow_port_mapping, "allow_mapping"); - changed!(udp_stun); - changed!(tcp_stun); - changed!(tunnel_addr); - changed!(tunnel_port); - fields -} - -fn port_mapping_update_requires_restart( - port_mapping_changed: bool, - ip_changed: bool, - server_changed: bool, - mtu_changed: bool, - tun_name_changed: bool, -) -> bool { - port_mapping_changed && (ip_changed || server_changed || mtu_changed || tun_name_changed) } impl VntApi { - pub(crate) fn new( - app_state: AppState, - server_rpc: ServerRPC, - ip_update: IpUpdateContext, - node_identity: SharedNodeIdentity, - runtime_config: RuntimeConfigController, - ) -> Self { + pub(crate) fn new(app_state: AppState, server_rpc: ServerRPC) -> Self { Self { app_state, server_rpc, - ip_update, - node_identity, - runtime_config, } } - /// Atomically applies the subset of configuration backed by live runtime - /// controls. If any changed field needs a restart, no live mutation is - /// performed and the caller receives `InstanceRestart`. - pub async fn reconfigure( - &self, - mut candidate: Box, - ) -> Result { - let _apply_guard = self.runtime_config.lock().await; - let current = self.app_state.get_config().ok_or_else(|| { - ReconfigureError::new( - ReconfigureStage::Prepare, - ReconfigureFallback::KeepCurrent, - "网络实例尚未运行", - ) - })?; - // Runtime identity belongs to the active subscription/instance. Do - // this before validation as well, so ignored remote identity cannot - // make an otherwise valid revision fail. - candidate.network_code = current.network_code.clone(); - candidate.device_id = current.device_id.clone(); - candidate.managed = current.managed.clone(); - // Event scripts are process-local and deliberately excluded from - // managed configuration. Never let a generic JNI/API candidate make - // AppState claim that a script changed while IpUpdateContext still - // owns the original executable. - candidate.event_script = current.event_script.clone(); - candidate.normalize().map_err(|error| { - ReconfigureError::new( - ReconfigureStage::Validate, - ReconfigureFallback::KeepCurrent, - error, - ) - })?; - candidate.check().map_err(|error| { - ReconfigureError::new( - ReconfigureStage::Validate, - ReconfigureFallback::KeepCurrent, - error, - ) - })?; - let changed_fields = config_diff(¤t, &candidate); - - if changed_fields.is_empty() { - return Ok(ReconfigureReport { - action: ApplyAction::NoChange, - changed_fields, - }); - } - - let has_restart = changed_fields - .iter() - .any(|field| matches!(field.as_str(), "password" | "device_mode")); - if has_restart { - return Ok(ReconfigureReport { - action: ApplyAction::InstanceRestart, - changed_fields, - }); - } - - let has_conditional_restart = changed_fields.iter().any(|field| { - matches!( - field.as_str(), - "tunnel_addr" | "tunnel_port" | "outbound_interface" - ) - }); - if has_conditional_restart { - return Ok(ReconfigureReport { - action: ApplyAction::InstanceRestart, - changed_fields, - }); - } - let port_mapping_changed = current.port_mapping != candidate.port_mapping; - let ip_changed = current.ip != candidate.ip; - let tun_name_changed = current.tun_name != candidate.tun_name; - let server_changed = current.server_addr != candidate.server_addr - || current.cert_mode != candidate.cert_mode; - let mtu_changed = - current.mtu.unwrap_or(DEFAULT_MTU) != candidate.mtu.unwrap_or(DEFAULT_MTU); - if port_mapping_update_requires_restart( - port_mapping_changed, - ip_changed, - server_changed, - mtu_changed, - tun_name_changed, - ) { - // These independently replace resources. Until their preparation - // can be committed through one dispatcher transaction, use the - // one-shot generation-safe restart path instead of risking a - // partially applied revision. - return Ok(ReconfigureReport { - action: ApplyAction::InstanceRestart, - changed_fields, - }); - } - let mut component_reload = changed_fields.iter().any(|field| { - matches!( - field.as_str(), - "server" - | "cert_mode" - | "peer_address" - | "no_punch" - | "udp_stun" - | "tcp_stun" - | "port_mapping" - | "rtx" - | "no_nat" - | "mtu" - ) - }); - - if port_mapping_changed { - match self - .runtime_config - .reload_port_mappings(&candidate.port_mapping) - .await - { - Ok(PortMappingReload::Reloaded) => {} - Ok(PortMappingReload::RestartRequired) => { - return Ok(ReconfigureReport { - action: ApplyAction::InstanceRestart, - changed_fields, - }); - } - Err(error) => { - return Err(ReconfigureError::new( - ReconfigureStage::Prepare, - ReconfigureFallback::KeepCurrent, - error, - )); - } - } - } - - let mut prepared_mtu = if mtu_changed { - Some( - self.runtime_config - .prepare_mtu_reload(&candidate) - .await - .map_err(|error| { - ReconfigureError::new( - ReconfigureStage::Prepare, - ReconfigureFallback::KeepCurrent, - error, - ) - })?, - ) - } else { - None - }; - - let previous_network = self - .app_state - .get_network() - .and_then(|network| VirtualIp::new(network.ip, network.prefix_len).ok()); - let mut network_applied = false; - if ip_changed || mtu_changed || tun_name_changed { - let target = match candidate.ip { - Some(target) => target, - None if !ip_changed => match previous_network { - Some(target) => target, - None => { - if let Some(prepared) = prepared_mtu.take() { - self.runtime_config.abort_mtu_reload(prepared).await; - } - return Err(ReconfigureError::new( - ReconfigureStage::Prepare, - ReconfigureFallback::KeepCurrent, - "客户端尚未完成网络注册", - )); - } - }, - None => { - if let Some(prepared) = prepared_mtu.take() { - self.runtime_config.abort_mtu_reload(prepared).await; - } - return Ok(ReconfigureReport { - action: ApplyAction::InstanceRestart, - changed_fields, - }); - } - }; - let network_apply = self - .ip_update - .apply_managed_network( - target, - candidate.mtu.unwrap_or(DEFAULT_MTU), - self.app_state.subnet_route.all_route(), - tun_name_changed.then(|| candidate.tun_name.clone()), - ) - .await; - let network_apply = match network_apply { - Ok(action) => action, - Err(error) => { - if let Some(prepared) = prepared_mtu.take() { - self.runtime_config.abort_mtu_reload(prepared).await; - } - #[cfg(not(target_os = "android"))] - let fallback = if error - .downcast_ref::() - .is_some_and(|error| error.rollback_failed) - { - ReconfigureFallback::InstanceRestart - } else { - ReconfigureFallback::KeepCurrent - }; - #[cfg(target_os = "android")] - let fallback = if error - .downcast_ref::() - .is_some_and(|error| error.restart_required) - { - ReconfigureFallback::InstanceRestart - } else { - ReconfigureFallback::KeepCurrent - }; - return Err(ReconfigureError::new( - ReconfigureStage::Commit, - fallback, - error, - )); - } - }; - if network_apply == ManagedNetworkApply::ComponentReload { - component_reload = true; - } - network_applied = true; - if ip_changed { - self.node_identity.notify_changed(); - } - } - if server_changed && let Err(error) = self.runtime_config.reload_servers(&candidate).await { - if network_applied { - let rollback = previous_network.ok_or_else(|| { - ReconfigureError::new( - ReconfigureStage::Rollback, - ReconfigureFallback::InstanceRestart, - "cannot roll back an unavailable virtual address", - ) - })?; - if let Err(rollback_error) = self - .ip_update - .apply_managed_network( - rollback, - current.mtu.unwrap_or(DEFAULT_MTU), - self.app_state.subnet_route.all_route(), - tun_name_changed.then(|| current.tun_name.clone()), - ) - .await - { - return Err(ReconfigureError::new( - ReconfigureStage::Rollback, - ReconfigureFallback::InstanceRestart, - format!( - "server reload failed ({error:#}) and IP rollback failed ({rollback_error:#})" - ), - )); - } - self.node_identity.notify_changed(); - } - if let Some(prepared) = prepared_mtu.take() { - self.runtime_config.abort_mtu_reload(prepared).await; - } - return Err(ReconfigureError::new( - ReconfigureStage::Prepare, - ReconfigureFallback::KeepCurrent, - error, - )); - } - if let Some(prepared) = prepared_mtu.take() - && let Err(error) = self.runtime_config.commit_mtu_reload(prepared).await - { - return Err(ReconfigureError::new( - ReconfigureStage::Commit, - ReconfigureFallback::InstanceRestart, - error, - )); - } - if current.input != candidate.input { - self.app_state - .subnet_route - .set_route_table(candidate.input.clone()); - } - if current.device_name != candidate.device_name - || current.output != candidate.output - || current.subnet_mapping != candidate.subnet_mapping - { - self.node_identity.set(NodeIdentityTemplate { - name: candidate.device_name.clone(), - version: env!("CARGO_PKG_VERSION").to_string(), - network_code: current.network_code.clone(), - advertised_subnets: advertised_subnets( - &candidate.output, - &candidate.subnet_mapping, - ), - }); - } - self.runtime_config.commit_policy(&candidate); - // Managed identity and revision state are runtime-owned and cannot be - // replaced by the parsed candidate. - self.app_state.set_config(candidate); - Ok(ReconfigureReport { - action: if component_reload { - ApplyAction::ComponentReload - } else { - ApplyAction::Live - }, - changed_fields, - }) - } pub fn server_rpc(&self) -> &ServerRPC { &self.server_rpc } @@ -641,75 +195,3 @@ impl VntApi { self.app_state.traffic_stats.reset_all() } } - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn managed_diff_is_stable_and_ignores_identity() { - let current = Config { - network_code: "network-a".into(), - device_id: "device-a".into(), - ..Default::default() - }; - let mut candidate = current.clone(); - candidate.network_code = "forged-network".into(); - candidate.device_id = "forged-device".into(); - candidate.device_name = "new-name".into(); - candidate.no_nat = true; - candidate.allow_port_mapping = true; - - assert_eq!( - config_diff(¤t, &candidate), - vec!["device_name", "no_nat", "allow_mapping"] - ); - } - - #[test] - fn implicit_and_explicit_default_mtu_are_semantically_equal() { - let current = Config::default(); - let mut explicit_default = current.clone(); - explicit_default.mtu = Some(DEFAULT_MTU); - assert!( - !config_diff(¤t, &explicit_default) - .iter() - .any(|field| field == "mtu") - ); - - explicit_default.mtu = Some(DEFAULT_MTU - 1); - assert!( - config_diff(¤t, &explicit_default) - .iter() - .any(|field| field == "mtu") - ); - } - - #[test] - fn port_mapping_and_tun_name_combination_requires_instance_restart() { - assert!(!port_mapping_update_requires_restart( - true, false, false, false, false - )); - assert!(!port_mapping_update_requires_restart( - false, false, false, false, true - )); - assert!(port_mapping_update_requires_restart( - true, false, false, false, true - )); - } - - #[test] - fn reconfigure_reports_have_protocol_stable_names() { - let report = ReconfigureReport { - action: ApplyAction::ComponentReload, - changed_fields: vec!["server".into()], - }; - assert_eq!( - serde_json::to_value(report).unwrap(), - serde_json::json!({ - "action": "COMPONENT_RELOAD", - "changed_fields": ["server"] - }) - ); - } -} diff --git a/vnt-core/src/context/config.rs b/vnt-core/src/context/config.rs index 6302536e..7cb6d47e 100644 --- a/vnt-core/src/context/config.rs +++ b/vnt-core/src/context/config.rs @@ -583,6 +583,16 @@ pub struct ManagedRegistration { revision: Arc, } +impl PartialEq for ManagedRegistration { + fn eq(&self, other: &Self) -> bool { + // revision 是运行期回执状态而非配置,比较时刻意排除 + self.credential_key == other.credential_key + && self.network_code == other.network_code + && self.device_id == other.device_id + && self.instance_id == other.instance_id + } +} + impl ManagedRegistration { pub fn new( credential_key: Vec, @@ -617,6 +627,11 @@ impl ManagedRegistration { } } + /// 已应用的服务端配置 revision(0 表示尚未应用)。 + pub(crate) fn applied_revision(&self) -> u64 { + self.revision.load(Ordering::Acquire) + } + pub(crate) fn verify_server_proof( &self, registration: &crate::protocol::control_message::SubscriptionRegistration, @@ -648,6 +663,96 @@ impl std::fmt::Debug for ManagedRegistration { } } impl Config { + /// 序列化为配置文件(TOML)文本,用于查看实例当前生效的配置。 + /// + /// 只输出与默认值不同的字段(布尔项仅在为 true 时输出),身份字段 + /// (network_code/device_id)始终输出;订阅托管时把已应用的 revision + /// 以注释形式标注在头部。键名与配置文件一致(如 server 对应 + /// server_addr、allow_mapping 对应 allow_port_mapping)。 + pub fn to_toml_string(&self) -> String { + let mut table = toml::Table::new(); + table.insert("network_code".into(), self.network_code.clone().into()); + table.insert("device_id".into(), self.device_id.clone().into()); + if !self.device_name.is_empty() { + table.insert("device_name".into(), self.device_name.clone().into()); + } + if !self.server_addr.is_empty() { + table.insert( + "server".into(), + toml_string_array(self.server_addr.iter().map(ToString::to_string)), + ); + } + if let Some(ip) = &self.ip { + table.insert("ip".into(), ip.to_string().into()); + } + insert_string_array(&mut table, "peer_address", &self.peer_address); + insert_string_array(&mut table, "turn", &self.turn); + insert_string_array(&mut table, "punch_model", &self.punch_model); + if self.cert_mode != CertValidationMode::default() { + table.insert("cert_mode".into(), self.cert_mode.to_string().into()); + } + if let Some(password) = &self.password { + table.insert("password".into(), password.clone().into()); + } + if let Some(outbound_interface) = &self.outbound_interface { + table.insert( + "outbound_interface".into(), + outbound_interface.clone().into(), + ); + } + if let Some(tun_name) = &self.tun_name { + table.insert("tun_name".into(), tun_name.clone().into()); + } + if self.device_mode != DeviceMode::default() { + table.insert("device_mode".into(), self.device_mode.to_string().into()); + } + if let Some(mtu) = self.mtu { + table.insert("mtu".into(), i64::from(mtu).into()); + } + if let Some(tunnel_port) = self.tunnel_port { + table.insert("tunnel_port".into(), i64::from(tunnel_port).into()); + } + insert_string_array(&mut table, "tunnel_addr", &self.tunnel_addr); + insert_string_array(&mut table, "input", &self.input); + insert_string_array(&mut table, "subnet_mapping", &self.subnet_mapping); + insert_string_array(&mut table, "output", &self.output); + insert_string_array(&mut table, "port_mapping", &self.port_mapping); + insert_string_array(&mut table, "udp_stun", &self.udp_stun); + insert_string_array(&mut table, "tcp_stun", &self.tcp_stun); + for (key, value) in [ + ("no_punch", self.no_punch), + ("no_broadcast", self.no_broadcast), + ("allow_ikev2", self.allow_ikev2), + ("allow_wireguard", self.allow_wireguard), + ("rtx", self.rtx), + ("compress", self.compress), + ("fec", self.fec), + ("auto_sync_subnet", self.auto_sync_subnet), + ("no_nat", self.no_nat), + // 配置文件里的键名是 allow_mapping + ("allow_mapping", self.allow_port_mapping), + ] { + if value { + table.insert(key.into(), value.into()); + } + } + if let Some(event_script) = &self.event_script { + table.insert("event_script".into(), event_script.clone().into()); + } + let mut text = toml::to_string(&table).unwrap_or_default(); + let mut header = String::from("# 当前生效配置(本地配置与服务端下发合并后的结果) +"); + if let Some(managed) = &self.managed { + header.push_str(&format!( + "# 服务端管理: 已应用 revision {} +", + managed.applied_revision() + )); + } + text.insert_str(0, &header); + text + } + pub fn normalize(&mut self) -> anyhow::Result<()> { self.check_turn_rules()?; self.check_tunnel_addr()?; @@ -769,7 +874,7 @@ impl Config { &self, index: usize, default_interface: Option, - registration_ip: crate::tunnel_core::server::transport::config::SharedRegistrationIp, + network: crate::context::SharedNetworkAddr, identity: crate::protocol::client_message::SharedNodeIdentity, client_instance_id: std::sync::Arc>, ) -> ConnectRegConfig { @@ -779,7 +884,7 @@ impl Config { network_code: self.network_code.clone(), device_id: self.device_id.clone(), identity, - ip: registration_ip, + ip: network, key_sign: self.key_sign(), ip_variable: self.ip.is_none(), allow_ikev2: self.allow_ikev2, @@ -791,6 +896,26 @@ impl Config { } } +/// 把实现了 Display 的集合序列化为 TOML 字符串数组。 +fn insert_string_array(table: &mut toml::Table, key: &str, values: &[T]) { + if values.is_empty() { + return; + } + table.insert( + key.into(), + toml::Value::Array( + values + .iter() + .map(|value| toml::Value::String(value.to_string())) + .collect(), + ), + ); +} + +fn toml_string_array>(values: I) -> toml::Value { + toml::Value::Array(values.into_iter().map(toml::Value::String).collect()) +} + #[cfg(test)] mod tests { use super::*; @@ -1153,4 +1278,30 @@ mod tests { assert!(allow_punch(&rules, &Ipv4Addr::new(10, 26, 0, 2))); assert!(allow_punch(&rules, &Ipv4Addr::new(10, 27, 0, 9))); } + + #[test] + fn to_toml_string_emits_non_default_fields_and_identity() { + let config = Config { + network_code: "net".to_string(), + device_id: "dev".to_string(), + device_name: "node".to_string(), + compress: true, + mtu: Some(1380), + server_addr: vec!["tcp://127.0.0.1:29872".parse().unwrap()], + ..Config::default() + }; + let text = config.to_toml_string(); + assert!(text.contains("network_code = \"net\""), "{text}"); + assert!(text.contains("device_id = \"dev\""), "{text}"); + assert!(text.contains("device_name = \"node\""), "{text}"); + assert!(text.contains("server = [\"tcp://127.0.0.1:29872\"]"), "{text}"); + assert!(text.contains("compress = true"), "{text}"); + assert!(text.contains("mtu = 1380"), "{text}"); + // 默认值字段不输出 + assert!(!text.contains("rtx"), "{text}"); + assert!(!text.contains("no_punch"), "{text}"); + assert!(!text.contains("password"), "{text}"); + // 输出本身是合法 TOML + assert!(toml::from_str::(&text).is_ok(), "{text}"); + } } diff --git a/vnt-core/src/context/mod.rs b/vnt-core/src/context/mod.rs index eedb0570..7b605655 100644 --- a/vnt-core/src/context/mod.rs +++ b/vnt-core/src/context/mod.rs @@ -320,7 +320,7 @@ pub struct TunnelListenAddr { pub protocol: &'static str, pub addr: SocketAddr, } -#[derive(Clone, Default)] +#[derive(Debug, Clone, Default)] pub(crate) struct SharedNetworkAddr { inner: Arc>>, } diff --git a/vnt-core/src/core/change_runtime.rs b/vnt-core/src/core/change_runtime.rs new file mode 100644 index 00000000..185237e4 --- /dev/null +++ b/vnt-core/src/core/change_runtime.rs @@ -0,0 +1,753 @@ +//! [`NetworkManager`] 的运行期变更应用:虚拟 IP/网段/入站路由的应用、 +//! 虚拟网卡任务重启与服务端地址增删、策略类配置热更新,以及 +//! "能否增量应用"的判断方法。 +//! +//! 每类配置一个独立的 apply 方法,方法内部各自完成"比较 → 提交对应 +//! 字段 → 执行变更动作":配置先行(动作失败时配置已是新值,不回滚, +//! 错误上抛由调用方记录),可单独调用;组合多项变更时按以下顺序 +//! (理由见各方法文档): +//! +//! 1. [`NetworkManager::needs_instance_rebuild`] 为 true:停止实例并重建, +//! 不要调用任何 apply 方法(等价于所有服务端连接全部重启)。 +//! 2. [`NetworkManager::apply_policy_change`]:先提交策略字段,让后续新建 +//! 的服务端连接携带最新的注册宣告字段。 +//! 3. [`NetworkManager::apply_server_change`]:依赖第 2 步已提交的配置。 +//! 4. [`NetworkManager::needs_network_change`] 为 true 时调用 +//! `apply_network_change`:与第 2、3 步无依赖;桌面平台的 MTU 也在 +//! 这一步热更新(tun-rs 支持动态修改,无需重建设备)。 +//! 5. [`NetworkManager::needs_device_restart`] 为 true 时最后调用 +//! `restart_device`:其内部会按快照重新应用地址与路由,覆盖第 4 步 +//! 的结果。 +//! +//! 平台划分:`apply_network_change` / `restart_device` 仅桌面平台 +//! (非 Android/iOS/tvOS)提供;移动平台改用 +//! [`NetworkManager::apply_mobile_change_fd`]——它同时承担变更应用与 +//! fd 型虚拟网卡重建(携带宿主新建的 TUN fd)。 + +use super::{DEFAULT_MTU, NetworkManager, ServerLinks}; +use crate::event_script::EventScriptType; +#[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] +use crate::context::config::DeviceMode; +use crate::context::{NetworkAddr, NetworkRoute}; +use crate::network_info::RuntimeChange; +use crate::runtime_config::RuntimePolicy; +#[cfg(not(target_os = "android"))] +use crate::tun::DeviceConfig; +use crate::tunnel_core::server::connection_manager::{ + InboundHandlerConfig, NewServerLink, create_server_manager, +}; +use crate::tunnel_core::server::transport::config::ProtocolAddress; +use anyhow::{Context, bail}; +use crate::context::config::Config; + +impl NetworkManager { + /// 判断一次运行期快照是否包含需要改动网络的字段:虚拟 IP/网段 + /// (`config.ip`)相对当前注册地址变化,配置的入站路由表(`input`) + /// 变化,或入站路由快照相对已应用路由变化。仅身份等非路由字段变化 + /// 时返回 false,调用方可以整体跳过应用。 + pub fn needs_network_change(&self, change: &RuntimeChange) -> bool { + let (_, address) = self.runtime_address(change.config.ip); + let structural_change = address.is_some() + || change.config.input != self.config.input + || normalized_routes(self.app_state.subnet_route.applied_routes()) + != normalized_routes(change.routes.clone()); + // 桌面平台 MTU 可热更新(tun-rs set_mtu),纳入网络变更判断; + // 移动平台 MTU 由宿主的 VPN 接口决定,归 needs_device_restart 覆盖 + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + let mtu_change = change.config.mtu.is_some() && change.config.mtu != self.config.mtu; + #[cfg(any(target_os = "android", target_os = "ios", target_os = "tvos"))] + let mtu_change = false; + // Windows/Linux/FreeBSD 支持网卡名热变更(tun-rs set_name); + // 恢复默认名(None)无法热应用,归 needs_device_restart 走重建 + #[cfg(any(windows, target_os = "linux", target_os = "freebsd"))] + let name_change = + change.config.tun_name.is_some() && change.config.tun_name != self.config.tun_name; + #[cfg(not(any(windows, target_os = "linux", target_os = "freebsd")))] + let name_change = false; + structural_change || mtu_change || name_change + } + + /// 判断快照是否包含只能通过重建实例生效的字段:password/cert_mode、 + /// 托管身份(network_code/device_id/device_name/managed)、绑定出口 + /// 网卡(outbound_interface)、本地隧道监听(tunnel_addr/tunnel_port)、 + /// device_mode、port_mapping、subnet_mapping、output、event_script。 + /// + /// 这些字段烘焙在长期存活对象中(加密与 QUIC 证书、连接参数、监听 + /// socket、设备分支、映射表),没有运行期应用路径。为 true 时调用方 + /// 应停止当前实例并重建(等价于所有服务端连接全部重启),不要调用 + /// 任何 apply 方法。 + pub fn needs_instance_rebuild(&self, change: &RuntimeChange) -> bool { + instance_rebuild_fields_differ(&self.config, &change.config) + } + + /// 判断快照中的服务端地址列表相对当前连接是否变化。 + /// + /// 应用顺序:无前置依赖;建议在 [`NetworkManager::apply_policy_change`] + /// 之后调用,使新增连接携带最新的注册宣告字段(allow_ikev2 等)。 + pub fn needs_server_change(&self, change: &RuntimeChange) -> bool { + server_fields_differ(&self.config, &change.config) + } + + /// 判断快照的策略类字段是否变化(可通过 apply_policy_change 热应用): + /// turn/punch_model/peer_address/compress/rtx/fec/no_punch/no_nat/ + /// no_broadcast/auto_sync_subnet/allow_port_mapping/allow_ikev2/ + /// allow_wireguard/udp_stun/tcp_stun。 + pub fn needs_policy_change(&self, change: &RuntimeChange) -> bool { + policy_fields_differ(&self.config, &change.config) + } + + /// 判断快照是否要求重建虚拟网卡任务:Windows/Linux/FreeBSD 的 MTU 与 + /// 网卡名均可热更新(仅恢复默认名需要重建);macOS 网卡名变化需要重建; + /// 移动平台的 MTU 与网卡名都要求重建。 + /// + /// 应用顺序:应在所有其他 apply 方法之后最后调用 + /// [`NetworkManager::restart_device`]——其内部会按快照重新应用虚拟 + /// 地址与路由,覆盖 apply_network_change 的结果。 + pub fn needs_device_restart(&self, change: &RuntimeChange) -> bool { + // 需要重建虚拟网卡任务的字段随平台热更新能力变化: + // - Windows/Linux/FreeBSD:MTU 与网卡名均可热变更,仅“恢复默认 + // 名”(tun_name 置空)无法热应用,需要重建; + // - macOS 等桌面平台:不支持热改名,网卡名变化需要重建; + // - 移动平台:MTU 与网卡名都由宿主 VPN 接口决定,都需重建 + #[cfg(any(windows, target_os = "linux", target_os = "freebsd"))] + { + change.config.tun_name.is_none() && self.config.tun_name.is_some() + } + #[cfg(all( + not(any(windows, target_os = "linux", target_os = "freebsd")), + not(any(target_os = "android", target_os = "ios", target_os = "tvos")) + ))] + { + self.config.tun_name != change.config.tun_name + } + #[cfg(any(target_os = "android", target_os = "ios", target_os = "tvos"))] + { + device_restart_fields_differ(&self.config, &change.config) + } + } + + /// 热应用策略类配置:turn/punch_model/peer_address/compress/rtx/fec/ + /// no_punch/no_nat/no_broadcast/auto_sync_subnet/allow_port_mapping/ + /// allow_ikev2/allow_wireguard/udp_stun/tcp_stun。 + /// + /// 通过 [`RuntimePolicyStore::store`] 整体换策略快照,各组件逐包读取 + /// 立即生效;peer/stun/no_punch 变化会唤醒对端探测、NAT 检测与打洞 + /// 循环。提交后把变更字段写入实例配置。返回是否发生了应用。 + /// + /// 边界:subnet_mapping/output 的句柄不可通过本方法变更(它们的差异 + /// 属于 [`NetworkManager::needs_instance_rebuild`]);allow_ikev2/ + /// allow_wireguard 的中继行为立即生效,但已建连服务器的注册宣告侧 + /// 要等下次重连或重建实例才更新。 + /// + /// 应用顺序:必须在 [`NetworkManager::apply_server_change`] **之前** + /// 调用——后者新建的服务端连接从已提交的实例配置读取注册宣告字段。 + pub fn apply_policy_change(&mut self, change: &RuntimeChange) -> bool { + if !policy_fields_differ(&self.config, &change.config) { + return false; + } + // subnet_mapping/output 的句柄从当前策略沿用:它们的差异属于 + // needs_instance_rebuild,不会走到这里 + let (subnet_mapping, relay_subnets) = { + let current = self.runtime_policy.load(); + ( + current.subnet_mapping.clone(), + current.relay_subnets.clone(), + ) + }; + let policy = RuntimePolicy::from_config_with_handles( + &change.config, + change.config.fec.then_some(self.fec_encoder.clone()), + subnet_mapping, + relay_subnets.clone(), + relay_subnets, + ); + self.runtime_policy.store(policy); + self.config.turn = change.config.turn.clone(); + self.config.punch_model = change.config.punch_model.clone(); + self.config.peer_address = change.config.peer_address.clone(); + self.config.compress = change.config.compress; + self.config.rtx = change.config.rtx; + self.config.fec = change.config.fec; + self.config.no_punch = change.config.no_punch; + self.config.no_nat = change.config.no_nat; + self.config.no_broadcast = change.config.no_broadcast; + self.config.auto_sync_subnet = change.config.auto_sync_subnet; + self.config.allow_port_mapping = change.config.allow_port_mapping; + self.config.allow_ikev2 = change.config.allow_ikev2; + self.config.allow_wireguard = change.config.allow_wireguard; + self.config.udp_stun = change.config.udp_stun.clone(); + self.config.tcp_stun = change.config.tcp_stun.clone(); + true + } + + /// 应用一次完整运行期快照中的网络状态:虚拟 IP/网段(设备虚拟地址)、 + /// 完整入站路由(系统路由 + 转发路由表),并把配置中影响路由的字段 + /// (`ip`、`input`)提交到实例。非安卓平台使用。 + /// + /// 相同地址的 set_network 是 no-op,系统路由按完整快照对账。执行顺序: + /// 先提交配置(路由快照、`ip`/`input`、MTU/网卡名与共享网络地址,配置 + /// 即新值),再执行变更动作(虚拟网卡、系统路由与服务端通告);动作 + /// 失败不回滚,错误上抛由调用方记录。快照未携带虚拟 IP(`config.ip` + /// 为空)时保持当前地址,只应用路由。 + /// + /// 应用顺序:与 apply_policy_change / apply_server_change 无依赖; + /// 若 [`NetworkManager::needs_device_restart`] 也为真,本方法的结果会 + /// 被 `restart_device` 覆盖,可直接只调用后者。 + /// 虚拟网卡被(重新)应用后触发事件脚本:热更新、重启网卡任务或换卡 + /// 重建之后调用,把当前生效的网卡状态通告宿主。无设备模式没有网卡可 + /// 应用,直接跳过。 + pub(crate) async fn notify_device_applied(&self) { + if !self.device_mode().has_device() { + return; + } + let Some(network) = self.app_state.network.get() else { + return; + }; + let mut params = vec![ + ("ip", network.ip.to_string()), + ("prefix-length", network.prefix_len.to_string()), + ( + "gateway", + network + .gateway + .map(|gateway| gateway.to_string()) + .unwrap_or_else(|| "-".to_string()), + ), + ("broadcast", network.broadcast.to_string()), + ("mtu", self.config.mtu.unwrap_or(DEFAULT_MTU).to_string()), + ]; + if let Some(tun_name) = &self.config.tun_name { + params.push(("tun-name", tun_name.clone())); + } + self.event_script + .notify(EventScriptType::DeviceApplied, ¶ms) + .await; + } + + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + pub async fn apply_network_change(&mut self, change: &RuntimeChange) -> anyhow::Result<()> { + let previous = self.app_state.network.get(); + let (desired, _) = self.runtime_address(change.config.ip); + // 先提交配置:MTU/网卡名、路由快照与 `ip`/`input`,随后同步共享 + // 网络地址(配置即新值) + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + if change.config.mtu != self.config.mtu { + self.config.mtu = change.config.mtu; + } + #[cfg(any(windows, target_os = "linux", target_os = "freebsd"))] + if change.config.tun_name != self.config.tun_name { + self.config.tun_name = change.config.tun_name.clone(); + } + self.commit_runtime_snapshot(change); + if let Some(address) = desired { + // 网关是注册时确定的虚拟网络属性,不随快照重算,沿用当前注册值 + let gateway = previous.and_then(|network| network.gateway); + self.app_state + .network + .set(NetworkAddr { gateway, ..address }); + } + // 变更动作:应用到虚拟网卡与系统路由(失败上抛,不回滚配置) + if self.device_mode().has_device() { + if let Some(address) = desired { + self.device_io_manager + .set_network(address.ip, address.prefix_len) + .await?; + } + self.device_io_manager + .apply_system_routes(change.routes.clone()) + .await?; + // MTU 热更新:桌面平台 tun-rs 支持动态修改,无需重建设备 + let mtu = change.config.mtu.or(self.config.mtu).unwrap_or(DEFAULT_MTU); + self.device_io_manager.set_mtu(mtu).await?; + // 网卡名热更新:Windows/Linux/FreeBSD 由 tun-rs 支持,无需重建 + #[cfg(any(windows, target_os = "linux", target_os = "freebsd"))] + if let Some(tun_name) = &change.config.tun_name { + self.device_io_manager.set_tun_name(tun_name).await?; + } + } + if let Some(address) = desired { + // 同步注册 IP,并快速注册到所有已连接服务端(连接不断、无 + // 周期性重注册,服务端映射靠这条通告收敛) + self.ip_update.announce_ip(address.ip).await; + } + Ok(()) + } + + /// 桌面平台专用:重启虚拟网卡任务。先停止当前任务,再按快照参数重建 + /// 虚拟网卡并启动新任务,最后把快照的虚拟 IP/网段与完整入站路由应用 + /// 到新设备。快照未携带虚拟 IP 时沿用当前注册地址。mtu/tun_name 在 + /// 重启动作前先提交到实例配置(配置即新值)。 + /// + /// 应用顺序:**最后调用**——当 [`NetworkManager::needs_device_restart`] + /// 为真时使用;其内部会按快照重新应用虚拟地址与路由,覆盖 + /// apply_network_change 已应用的结果。 + /// + /// Android/iOS/tvOS 不提供本方法:移动平台使用 + /// [`NetworkManager::apply_mobile_change_fd`] 应用变更(必要时重建 + /// fd 型虚拟网卡)。 + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + pub async fn restart_device(&mut self, change: &RuntimeChange) -> anyhow::Result<()> { + let config = self.restart_device_config(change)?; + self.restart_device_with(change, config).await + } + + /// 桌面 unix 专用变体:携带宿主新建的 TUN fd 时直接以该 fd 构建新 + /// 虚拟网卡,不经过系统设备创建流程。 + #[cfg(all(unix, not(any(target_os = "android", target_os = "ios", target_os = "tvos"))))] + pub async fn restart_device_fd( + &mut self, + change: &RuntimeChange, + tun_fd: Option, + ) -> anyhow::Result<()> { + let mut config = self.restart_device_config(change)?; + if let Some(tun_fd) = tun_fd { + config = config.set_tun_fd(tun_fd); + } + self.restart_device_with(change, config).await + } + + /// 依据快照与实例配置构建重启用的设备配置(模式、MTU、TAP MAC、名称)。 + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + fn restart_device_config(&self, change: &RuntimeChange) -> anyhow::Result { + let (desired, _) = self.runtime_address(change.config.ip); + let address = desired.or(self.app_state.network.get()); + let mut config = DeviceConfig::default() + .set_device_mode(self.config.device_mode) + .set_mtu(change.config.mtu.or(self.config.mtu).unwrap_or(DEFAULT_MTU)); + if self.config.device_mode == DeviceMode::Tap { + let ip = address + .map(|network| network.ip) + .context("网络尚未注册,无法确定 TAP 网卡 MAC")?; + config = config.set_mac_addr(crate::ethernet::mac_from_ip(ip).octets()); + } + if let Some(tun_name) = change.config.tun_name.clone() { + config = config.set_tun_name(tun_name); + } + Ok(config) + } + + #[cfg(not(any(target_os = "android", target_os = "ios", target_os = "tvos")))] + async fn restart_device_with( + &mut self, + change: &RuntimeChange, + config: DeviceConfig, + ) -> anyhow::Result<()> { + if !self.device_mode().has_device() { + bail!("当前为无虚拟网卡模式,没有可重启的虚拟网卡任务"); + } + if self.tun_receiver.is_some() || self.enhanced_outbound.is_some() { + bail!("虚拟网卡尚未启动"); + } + let (desired, _) = self.runtime_address(change.config.ip); + let address = desired.or(self.app_state.network.get()); + // 先把设备字段提交到实例配置(配置即新值),再关闭当前虚拟网卡 + // 任务并启动新任务,随后把地址与路由应用到新设备;动作失败不回滚 + self.config.mtu = change.config.mtu; + self.config.tun_name = change.config.tun_name.clone(); + self.device_io_manager.restart_task(config).await?; + if let Some(address) = &address { + self.device_io_manager + .set_network(address.ip, address.prefix_len) + .await?; + } + #[cfg(not(any(target_os = "ios", target_os = "tvos")))] + self.apply_network_change(change).await?; + Ok(()) + } + + /// 全平台通用:应用服务端地址变更。 + /// + /// 先把快照中的服务端地址列表提交到实例配置(配置即新值),再执行 + /// 连接增删动作:快照与当前连接对比分为删除、新增、未变三类——已 + /// 删除地址的连接任务先停止并从出站/RPC 登记中摘除;新增地址创建 + /// 新的连接任务(进入常规注册重连流程);未变地址保持原 server_id + /// 与连接完全不动。最后统一刷新服务端信息集合与出站表。动作失败 + /// 不回滚,错误上抛由调用方记录。 + /// + /// 应用顺序:建议在 [`NetworkManager::apply_policy_change`] **之后** + /// 调用(新建连接读取已提交实例配置中的注册宣告字段);与网络变更 + /// (apply_network_change)和设备重启相互独立。 + pub async fn apply_server_change(&mut self, change: &RuntimeChange) -> anyhow::Result<()> { + let new_addresses = change.config.server_addr.clone(); + if new_addresses == self.config.server_addr { + return Ok(()); + } + if self.server_task_group.is_stopped() { + bail!("服务端连接任务已停止,无法应用地址变更"); + } + self.config.server_addr = new_addresses.clone(); + // 阶段一(持锁、无 await):差分出删除与新增,摘除删除项的登记, + // 并为新增地址预先构建连接构件 + let (removals, created) = { + let mut links = self.server_links.lock(); + let removals = links.registry.plan_removals(&new_addresses); + let added: Vec = new_addresses + .iter() + .filter(|address| !links.registry.is_known(address)) + .cloned() + .collect(); + let created: Vec<(u32, ProtocolAddress, NewServerLink)> = added + .into_iter() + .map(|address| { + let server_id = links.registry.next_id(); + // 单地址配置,让 to_connect_config 以该地址构建连接参数 + let mut connect_config = (*self.config).clone(); + connect_config.server_addr = vec![address.clone()]; + let link = create_server_manager( + server_id, + &connect_config, + links.default_interface.clone(), + links.identity.clone(), + links.client_instance_id.clone(), + links.network.clone(), + ); + (server_id, address, link) + }) + .collect(); + (removals, created) + }; + // 阶段二(无锁):停止已删除地址的连接任务 + for (server_id, address, task) in removals { + if let Some(task) = task { + task.stop().await; + } + log::info!("已停止服务端连接任务: {server_id} {address}"); + } + // 阶段三(持锁、无 await):启动新增连接任务并登记,统一发布 + let mut links = self.server_links.lock(); + for (server_id, address, link) in created { + let NewServerLink { + manager, + sender, + notifier, + subscription_verified, + } = link; + let task = manager.data_handle_task( + &self.server_task_group, + self.server_handler_config(&links), + false, + ); + links + .registry + .add(server_id, address.clone(), sender, notifier, subscription_verified, task); + log::info!("已为新增服务端创建连接任务: {server_id} {address}"); + } + self.app_state + .server_info_collection + .update_server(links.registry.id_address_pairs()); + self.server_rpc.update_server_links(links.registry.publish()); + Ok(()) + } + + /// 组装服务端数据任务的处理配置,内容与初始注册时一致。 + fn server_handler_config(&self, links: &ServerLinks) -> Box { + Box::new(InboundHandlerConfig { + network_route: NetworkRoute::new( + self.app_state.network.clone(), + self.app_state.subnet_route.clone(), + ), + server_info: self.app_state.server_info_collection.clone(), + nat_info: self.app_state.nat_info.clone(), + peer_map: self.app_state.peer_map.clone(), + punch_backoff: self.app_state.punch_backoff.clone(), + puncher: links.puncher.clone(), + packet_crypto: links.packet_crypto.clone(), + enhanced_inbound: links.enhanced_inbound.clone(), + fec_decoder: links.fec_decoder.clone(), + policy: links.policy.clone(), + basic_outbound: links.basic_outbound.clone(), + app_state: self.app_state.clone(), + }) + } + + /// Android/iOS/tvOS 专用:应用一次完整运行期快照,额外携带宿主新建 + /// 的 TUN fd。移动平台没有 restart_device,本方法同时承担变更应用与 + /// fd 型虚拟网卡重建。 + /// + /// 设备模式:携带 fd 时以快照参数整体重建 fd 型虚拟网卡(地址/MTU 以 + /// 快照为准,快照未携带则沿用当前值):Android 走 + /// replace_android_virtual_device(同步注册地址,IP 变化时发送 + /// fast_reg);iOS/tvOS 经 restart_task 重建后设置地址。未携带 fd 时 + /// 网卡相关变更无法应用(返回错误,调用方据此让宿主重建接口后重试), + /// 仅提交与网卡无关的快照;无设备模式只做地址登记。执行顺序:先提交 + /// 路由快照、`ip`/`input`、共享地址与设备字段(`mtu`、`tun_name`, + /// 配置即新值),再执行 fd 重建动作;动作失败不回滚,错误上抛由 + /// 调用方记录。 + /// + /// 应用顺序:与 apply_policy_change / apply_server_change 无依赖; + /// 是否需要新 fd 由 [`NetworkManager::needs_vpn_rebuild`] 判断。 + #[cfg(any(target_os = "android", target_os = "ios", target_os = "tvos"))] + pub async fn apply_mobile_change_fd( + &mut self, + change: &RuntimeChange, + tun_fd: Option, + ) -> anyhow::Result<()> { + let (desired, address) = self.runtime_address(change.config.ip); + if !self.device_mode().has_device() { + // 无虚拟网卡:先提交配置(地址登记与快照),再执行通告动作 + // (策略/服务器变更已由调用方应用) + if let Some(address) = address { + self.app_state.network.set(address); + } + self.commit_runtime_snapshot(change); + if let Some(address) = address { + // 地址变化同样要通告:服务端映射不依赖本机网卡存在 + self.ip_update.announce_ip(address.ip).await; + } + return Ok(()); + } + let Some(tun_fd) = tun_fd else { + // 设备模式:VpnService 接口的地址/MTU/路由在 establish 时确定, + // 网卡相关变更必须由宿主整体重建;没有新 fd 时仅允许与网卡无关的 + // 变更(已由调用方先行应用),这里只提交快照 + if self.needs_vpn_rebuild(change) { + bail!("虚拟 IP/网段/MTU/路由/网卡名变更需要传入新的 TUN fd 以重建 VPN 接口"); + } + self.commit_runtime_snapshot(change); + return Ok(()); + }; + // 设备模式 + 新 fd:先提交配置(共享地址、路由快照与设备字段, + // 配置即新值),再执行 fd 重建动作;动作失败不回滚 + let mut network = desired + .or(self.app_state.network.get()) + .context("网络尚未注册,无法重建虚拟网卡")?; + // 网关是注册时确定的虚拟网络属性,不随快照重算,沿用当前注册值 + network.gateway = self.app_state.network.get().and_then(|value| value.gateway); + let mtu = change.config.mtu.or(self.config.mtu).unwrap_or(DEFAULT_MTU); + let (ip, prefix_len) = (network.ip, network.prefix_len); + let ip_changed = self + .app_state + .network + .get() + .map(|current| current.ip != network.ip) + .unwrap_or(true); + // 写入共享网络地址(重注册据此声称同一 IP) + self.app_state.network.set(network); + self.commit_runtime_snapshot(change); + self.config.mtu = change.config.mtu; + self.config.tun_name = change.config.tun_name.clone(); + #[cfg(target_os = "android")] + { + self.device_io_manager + .replace_task_fd(tun_fd, ip, prefix_len, mtu) + .await?; + } + #[cfg(any(target_os = "ios", target_os = "tvos"))] + { + let mut config = DeviceConfig::default() + .set_device_mode(self.config.device_mode) + .set_mtu(mtu) + .set_tun_fd(tun_fd); + if let Some(tun_name) = change.config.tun_name.clone() { + config = config.set_tun_name(tun_name); + } + self.device_io_manager.restart_task(config).await?; + self.device_io_manager.set_network(ip, prefix_len).await?; + } + // IP 实际变化时向所有已连接服务端发送快速注册 + if ip_changed { + self.ip_update.announce_ip(ip).await; + } + Ok(()) + } + + /// 移动平台:快照是否包含需要宿主重建 VPN 接口的变更。VpnService 的 + /// 地址/MTU/路由在 establish 时确定,虚拟 IP/网段、MTU、路由、网卡名 + /// 任何一项变化都需要新接口(即调用方需要传入新的 TUN fd)。 + /// + /// 无设备模式(device_mode=no)没有 VPN 接口可重建,网卡类变更随 + /// 快照直接提交,恒为 false。 + #[cfg(any(target_os = "android", target_os = "ios", target_os = "tvos"))] + pub fn needs_vpn_rebuild(&self, change: &RuntimeChange) -> bool { + if !self.device_mode().has_device() { + return false; + } + let (_, address) = self.runtime_address(change.config.ip); + address.is_some() + || (change.config.mtu.is_some() && change.config.mtu != self.config.mtu) + || change.config.tun_name != self.config.tun_name + || normalized_routes(self.app_state.subnet_route.applied_routes()) + != normalized_routes(change.routes.clone()) + } + + /// 依据快照中的虚拟 IP 推导目标网络地址,并返回“仅在实际变化时保留” + /// 的地址值。变化判断只看 ip 与前缀长度;虚拟网关是注册时确定的虚拟 + /// 网络属性,不参与网卡变更判断,也不随快照重算(沿用注册值)。 + fn runtime_address( + &self, + ip: Option, + ) -> (Option, Option) { + let previous = self.app_state.network.get(); + let desired = ip.map(|ip| { + let net = ip.network(); + NetworkAddr { + ip: ip.ip(), + gateway: None, + prefix_len: ip.prefix_len(), + broadcast: net.broadcast(), + } + }); + let changed = match (&desired, previous) { + (Some(desired), Some(previous)) => { + desired.ip != previous.ip || desired.prefix_len != previous.prefix_len + } + (Some(_), None) => true, + (None, _) => false, + }; + let address = desired.filter(|_| changed); + (desired, address) + } + + /// 提交路由快照,并把配置中影响路由的字段(`ip`、`input`)应用到实例。 + fn commit_runtime_snapshot(&mut self, change: &RuntimeChange) { + self.app_state.subnet_route.apply_routes(change.routes.clone()); + if change.config.ip != self.config.ip || change.config.input != self.config.input { + self.config.ip = change.config.ip; + self.config.input = change.config.input.clone(); + self.app_state + .subnet_route + .set_route_table(change.config.input.clone()); + self.app_state.set_config(self.config.clone()); + } + } +} + +/// 便于集合比较的路由规范化(排序去重;仅用于比较,不代表应用顺序)。 +fn normalized_routes(mut routes: Vec) -> Vec { + routes.sort_by_key(|route| { + ( + u32::from(route.net.network()), + route.net.prefix_len(), + u32::from(route.target_ip), + ) + }); + routes.dedup(); + routes +} + +/// 服务端地址差分:可通过 apply_server_change 热应用(增删连接)。 +fn server_fields_differ(current: &Config, target: &Config) -> bool { + current.server_addr != target.server_addr +} + +/// 策略类字段差分:可通过 apply_policy_change 热应用的字段集合。 +fn policy_fields_differ(current: &Config, target: &Config) -> bool { + current.turn != target.turn + || current.punch_model != target.punch_model + || current.peer_address != target.peer_address + || current.compress != target.compress + || current.rtx != target.rtx + || current.fec != target.fec + || current.no_punch != target.no_punch + || current.no_nat != target.no_nat + || current.no_broadcast != target.no_broadcast + || current.auto_sync_subnet != target.auto_sync_subnet + || current.allow_port_mapping != target.allow_port_mapping + || current.allow_ikev2 != target.allow_ikev2 + || current.allow_wireguard != target.allow_wireguard + || current.udp_stun != target.udp_stun + || current.tcp_stun != target.tcp_stun +} + +/// 重建实例字段差分:没有运行期应用路径、必须重建 NetworkManager 的字段集合。 +fn instance_rebuild_fields_differ(current: &Config, target: &Config) -> bool { + current.password != target.password + || current.cert_mode != target.cert_mode + || current.network_code != target.network_code + || current.device_id != target.device_id + || current.device_name != target.device_name + || current.managed != target.managed + || current.outbound_interface != target.outbound_interface + || current.tunnel_addr != target.tunnel_addr + || current.tunnel_port != target.tunnel_port + || current.device_mode != target.device_mode + || current.port_mapping != target.port_mapping + || current.subnet_mapping != target.subnet_mapping + || current.output != target.output + || current.event_script != target.event_script +} + +/// 设备重建字段差分:重启虚拟网卡任务(或宿主重建 fd 接口)才能生效的 +/// 字段集合。桌面平台的 MTU 已支持热更新,不参与本差分;该函数目前仅 +/// 移动平台调用。 +#[cfg_attr( + not(any(target_os = "android", target_os = "ios", target_os = "tvos")), + allow(dead_code) +)] +fn device_restart_fields_differ(current: &Config, target: &Config) -> bool { + current.mtu != target.mtu || current.tun_name != target.tun_name +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn policy_field_changes_are_detected() { + let current = Config::default(); + assert!(!policy_fields_differ(¤t, ¤t)); + + let mut target = current.clone(); + target.no_punch = true; + assert!(policy_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.compress = true; + assert!(policy_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.udp_stun = vec!["stun:127.0.0.1:3478".to_string()]; + assert!(policy_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.turn = vec!["10.26.0.0/24,10.26.0.1".parse().unwrap()]; + assert!(policy_fields_differ(¤t, &target)); + } + + #[test] + fn rebuild_field_changes_are_detected() { + let current = Config::default(); + assert!(!instance_rebuild_fields_differ(¤t, ¤t)); + + let mut target = current.clone(); + target.password = Some("changed".to_string()); + assert!(instance_rebuild_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.outbound_interface = Some("eth1".to_string()); + assert!(instance_rebuild_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.tunnel_addr = vec!["127.0.0.1:9999".parse().unwrap()]; + assert!(instance_rebuild_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.managed = Some(crate::context::config::ManagedRegistration::new( + vec![1; 32], + "network".to_string(), + "device".to_string(), + vec![2; 32], + 0, + )); + assert!(instance_rebuild_fields_differ(¤t, &target)); + } + + #[test] + fn device_restart_field_changes_are_detected() { + let current = Config::default(); + assert!(!device_restart_fields_differ(¤t, ¤t)); + + let mut target = current.clone(); + target.mtu = Some(1400); + assert!(device_restart_fields_differ(¤t, &target)); + + let mut target = current.clone(); + target.tun_name = Some("vnt1".to_string()); + assert!(device_restart_fields_differ(¤t, &target)); + + // policy 维度的变化不触发设备重启 + let mut target = current.clone(); + target.no_punch = true; + assert!(!device_restart_fields_differ(¤t, &target)); + } +} diff --git a/vnt-core/src/core/ip_update.rs b/vnt-core/src/core/ip_update.rs new file mode 100644 index 00000000..17ce7a4b --- /dev/null +++ b/vnt-core/src/core/ip_update.rs @@ -0,0 +1,76 @@ +//! 虚拟 IP 变更的传播:注册 IP 同步与向全部已连接服务端的快速注册通告。 +//! +//! 诞生时由服务端下发的 UpdateIp 消息在 inbound 里就地触发(服务端主动改 +//! IP);现在触发方是订阅统一流程(apply_network_change / +//! apply_mobile_change_fd),本结构只负责效果的落地。IP 变化后连接保持 +//! 不断、客户端也没有周期性重注册,服务端的 ip→会话映射全靠 +//! [`IpUpdateContext::announce_ip`] 收敛:不发则发往新 IP 的入站流量在 +//! 服务端无映射可投递,Ikev2/WireGuard 中继还会按 src 校验拒包。 + +use crate::protocol::control_message::{FastRegRequestMsg, RequestMessage}; +use crate::protocol::ip_packet_protocol::{HEAD_LENGTH, MsgType, NetPacket}; +use crate::protocol::transmission::TransmissionBytes; +use crate::tunnel_core::server::outbound::ServerOutbound; +use bytes::Bytes; +use std::net::Ipv4Addr; +use std::time::Duration; + +/// IP 变更的快速注册通告。重注册声称的 IP 直接读 +/// [`crate::context::SharedNetworkAddr`],本结构只负责“变化之后通知 +/// 所有已连接服务端”。 +#[derive(Clone)] +pub(crate) struct IpUpdateContext { + server_outbound: ServerOutbound, +} + +impl IpUpdateContext { + pub fn new(server_outbound: ServerOutbound) -> Self { + Self { server_outbound } + } + + /// 向所有已连接服务端发送快速注册(IP 已由调用方写入共享网络地址)。 + pub(crate) async fn announce_ip(&self, ip: Ipv4Addr) { + self.send_fast_reg(ip).await; + } + + fn fast_reg_packet(ip: Ipv4Addr) -> anyhow::Result { + let payload = RequestMessage::FastReg(FastRegRequestMsg { ip }).encode(); + let mut packet = NetPacket::new(TransmissionBytes::zeroed(HEAD_LENGTH + payload.len()))?; + packet.set_msg_type(MsgType::FastReg); + packet.set_ttl(1); + packet.set_gateway_flag(true); + packet.set_payload(&payload)?; + Ok(packet.into_buffer().into_bytes().freeze()) + } + + async fn send_fast_reg(&self, ip: Ipv4Addr) { + let result = match Self::fast_reg_packet(ip) { + Ok(packet) => { + self.server_outbound + .send_gateway_to_all(packet, Duration::from_secs(2)) + .await + } + Err(error) => Err(error), + }; + match result { + Ok(sent) => log::info!("快速注册已发送到 {sent} 台服务端,新 IP: {ip}"), + Err(error) => log::warn!("发送快速注册失败,保留新 IP {ip}: {error:#}"), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn fast_reg_packet_carries_ip_with_gateway_flag() { + let buffer = IpUpdateContext::fast_reg_packet(Ipv4Addr::new(10, 26, 0, 9)).unwrap(); + let packet = NetPacket::new(TransmissionBytes::from(buffer)).unwrap(); + assert_eq!(packet.msg_type().unwrap(), MsgType::FastReg); + assert!(packet.is_gateway()); + assert_eq!(packet.ttl(), 1); + // 网关包载荷只携带新 IP(proto FastRegRequestMsg 单字段) + assert!(!packet.payload().is_empty()); + } +} diff --git a/vnt-core/src/core/mod.rs b/vnt-core/src/core/mod.rs index 4341f11b..ec04e0b3 100644 --- a/vnt-core/src/core/mod.rs +++ b/vnt-core/src/core/mod.rs @@ -1,16 +1,13 @@ use crate::api::VntApi; use crate::context::config::{Config, DeviceMode}; -use crate::context::{AppState, NetworkAddr, NetworkRoute}; +use crate::context::{AppState, NetworkAddr, NetworkRoute, SharedNetworkAddr}; use crate::crypto::PacketCrypto; use crate::enhanced_tunnel::inbound::EnhancedInbound; use crate::enhanced_tunnel::outbound::EnhancedOutbound; -use crate::enhanced_tunnel::{ - MtuTunnelComponents, TunnelComponents, TunnelConfig, enhanced_ipv4_tunnel, -}; +use crate::enhanced_tunnel::{TunnelComponents, TunnelConfig, enhanced_ipv4_tunnel}; use crate::event_script::EventScript; -#[cfg(not(target_os = "android"))] -use crate::event_script::EventScriptType; use crate::fec::{FecDecoder, FecEncoder}; +use crate::log_manager::InstanceLog; use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager}; use crate::nat::subnet_packet::SubnetPacketMapper; use crate::nat::{ @@ -18,10 +15,7 @@ use crate::nat::{ }; use crate::protocol::client_message::{NodeIdentityTemplate, SharedNodeIdentity}; use crate::protocol::control_message::ErrorResponseMsg; -use crate::runtime_config::{ - MtuComponentController, RuntimeConfigController, RuntimePolicy, RuntimePolicyStore, - ServerComponentController, -}; +use crate::runtime_config::{RuntimePolicy, RuntimePolicyStore}; use crate::tun::enhanced_tun::EnhancedTunInbound; use crate::tun::{DeviceConfig, DeviceIOManager, TunDataInbound, TunReceiver, tun_channel}; use crate::tunnel_core::outbound::{BasicOutbound, HybridOutbound}; @@ -31,24 +25,30 @@ use crate::tunnel_core::p2p::transport::task::{ P2pInitConfig, init_tunnel, node_announcement_task, }; use crate::tunnel_core::server::connection_manager::{ - InboundHandlerConfig, ServerTurnManager, create_server_tunnel, register_with_first_available, - server_addresses, + InboundHandlerConfig, ServerLinkRegistry, ServerTurnManager, create_server_tunnel, + register_with_first_available, server_addresses, }; -use crate::tunnel_core::server::inbound::IpUpdateContext; use crate::tunnel_core::server::rpc::ServerRPC; use crate::utils::task_control::TaskGroup; use anyhow::{Context, bail}; use ipnet::Ipv4Net; +use parking_lot::Mutex; use rand::RngExt; -use std::net::Ipv4Addr; +use std::sync::Arc; +pub mod change_runtime; +mod ip_update; +use ip_update::IpUpdateContext; pub const DEFAULT_MTU: u16 = 1380; +/// P2P 栈(rustp2p-core IpStack)要求 IPv6 MTU >= 1280,低于它的配置 +/// 无法创建组网实例 +pub const MIN_MTU: u16 = 1280; /// Context for deferred registration struct RegistrationContext { server_task_group: TaskGroup, + server_links: Arc>, server_managers: Vec, - ip_update: IpUpdateContext, subnet_external_route: SubnetExternalRoute, puncher: NatPuncher, packet_crypto: PacketCrypto, @@ -58,29 +58,125 @@ struct RegistrationContext { basic_outbound: BasicOutbound, } +/// 运行期服务端连接登记与其按需重建所需的构件。初始服务端的连接任务 +/// 由注册任务挂入,之后可通过 [`NetworkManager::apply_server_change`] +/// 按快照增删。 +struct ServerLinks { + registry: ServerLinkRegistry, + packet_crypto: PacketCrypto, + puncher: NatPuncher, + enhanced_inbound: EnhancedInbound, + fec_decoder: FecDecoder, + policy: RuntimePolicyStore, + basic_outbound: BasicOutbound, + default_interface: Option, + identity: SharedNodeIdentity, + client_instance_id: Arc>, + network: SharedNetworkAddr, +} + +impl ServerLinks { + #[allow(clippy::too_many_arguments)] + fn new( + registry: ServerLinkRegistry, + packet_crypto: PacketCrypto, + puncher: NatPuncher, + enhanced_inbound: EnhancedInbound, + fec_decoder: FecDecoder, + policy: RuntimePolicyStore, + basic_outbound: BasicOutbound, + default_interface: Option, + identity: SharedNodeIdentity, + client_instance_id: Arc>, + network: SharedNetworkAddr, + ) -> Self { + Self { + registry, + packet_crypto, + puncher, + enhanced_inbound, + fec_decoder, + policy, + basic_outbound, + default_interface, + identity, + client_instance_id, + network, + } + } +} + pub struct NetworkManager { config: Box, - event_script: EventScript, app_state: AppState, + log: Arc, task_group: TaskGroup, device_io_manager: DeviceIOManager, ip_update: IpUpdateContext, - node_identity: SharedNodeIdentity, + /// 运行期事件脚本句柄:网卡被(重新)应用时触发 device_applied + event_script: EventScript, enhanced_outbound: Option, server_rpc: ServerRPC, - runtime_config: RuntimeConfigController, + server_task_group: TaskGroup, + server_links: Arc>, tun_receiver: Option, - registration_context: Option>, + registration_status: tokio::sync::watch::Receiver, + /// 运行期策略存储:apply_policy_change 通过 store() 热更新策略字段 + runtime_policy: RuntimePolicyStore, + /// FEC 编码句柄:保留它使 `fec` 开关可以在运行期双向切换 + /// (policy 中是否引用该句柄决定编码是否生效) + fec_encoder: FecEncoder, } -pub enum RegisterResponse { - Success(NetworkAddr), + +pub type NetworkInstance = NetworkManager; +enum RegisterResponse { + /// 注册成功,网段信息已写入 `app_state.network` + Registered, Failed(ErrorResponseMsg), } +/// `create_network` 启动的后台注册任务状态。注册动作对调用方透明, +/// 等待方只通过该状态拿到「注册成功」或「服务端拒绝」的结果。 +enum RegistrationStatus { + /// 注册进行中(含连接失败重试中) + Pending, + /// 注册成功,`app_state.network` 已写入网段信息 + Ready, + /// 服务端明确拒绝注册,携带错误消息,不再重试 + Failed(String), +} + impl NetworkManager { + pub async fn stop(mut self) { + self.task_group.stop(); + self.wait_all_stopped().await; + } + + /// 创建网络实例,并在实例任务组内后台连接服务器注册。 + /// 建网过程中的错误会写入实例日志。 + /// + /// `routes_tx` 是运行期变化通道的完整入栈路由发送端:核心把子网路由的 + /// 每一次生效快照推送给调用方持有的 [`crate::network_info::RuntimeChangeListener`]。 pub async fn create_network( + config: Box, + task_group: TaskGroup, + log: Arc, + routes_tx: tokio::sync::watch::Sender>, + ) -> anyhow::Result { + match Self::create_network_impl(config, task_group, log.clone(), routes_tx).await { + Ok(manager) => Ok(manager), + Err(error) => { + log.error(format!("创建网络失败: {error:#}")); + Err(error) + } + } + } + + async fn create_network_impl( mut config: Box, task_group: TaskGroup, + log: Arc, + routes_tx: tokio::sync::watch::Sender>, ) -> anyhow::Result { let app_state = AppState::default(); // 本机 NAT 身份变化(换网/NAT 重启等)时,把所有对端的 @@ -90,6 +186,16 @@ impl NetworkManager { app_state.nat_info.set_on_change(move || backoff.cap_all()); config.normalize()?; config.check()?; + // Install the route source before any server or gossip producer can + // publish automatic routes. Every source is therefore observed even + // when it changes during network task startup. + let subnet_external_route = app_state.subnet_route.clone(); + subnet_external_route.set_route_table(config.input.clone()); + let route_changes = subnet_external_route.subscribe(); + task_group.spawn(forward_network_route_changes( + route_changes, + routes_tx, + )); let outbound_interface_name = config .outbound_interface .as_deref() @@ -108,6 +214,9 @@ impl NetworkManager { log::info!("绑定出口网卡: {name}"); } let mtu = config.mtu.unwrap_or(DEFAULT_MTU); + if mtu < MIN_MTU { + bail!("MTU 必须 >= {MIN_MTU}(P2P 栈 IPv6 MTU 下限),当前配置: {mtu}"); + } let packet_crypto = PacketCrypto::new_from_str(config.password.as_deref())?; let allow_subnet = AllowSubnetExternalRoute::new(config.output.clone()); let relay_subnets = AllowSubnetExternalRoute::new(advertised_subnets( @@ -131,7 +240,7 @@ impl NetworkManager { let mut instance = vec![0_u8; 32]; rand::rng().fill(instance.as_mut_slice()); let client_instance_id = std::sync::Arc::new(instance); - let (server_manager_list, tunnel_to_server, server_rpc, registration_ip) = + let (server_manager_list, tunnel_to_server, server_rpc) = create_server_tunnel( app_state.clone(), &config, @@ -145,19 +254,7 @@ impl NetworkManager { .update_server(server_addresses(&config)); let server_task_group = task_group.child_scope(); let device_io_manager = DeviceIOManager::new(task_group.clone()); - let ip_update = IpUpdateContext::new( - app_state.network.clone(), - registration_ip.clone(), - tunnel_to_server.clone(), - device_io_manager.clone(), - config.device_mode, - EventScript::new(config.event_script.clone()), - config - .server_addr - .iter() - .map(|addr| addr.to_string()) - .collect(), - ); + let ip_update = IpUpdateContext::new(tunnel_to_server.clone()); // Keep the listener resident so no_punch/peer/turn can be changed // without replacing the instance or its virtual network device. let (puncher, p2p_socket_manager, p2p_task) = init_tunnel( @@ -185,8 +282,6 @@ impl NetworkManager { runtime_policy.clone(), node_identity.get(), ); - let subnet_external_route = app_state.subnet_route.clone(); - subnet_external_route.set_route_table(config.input.clone()); let subnet_packet_mapper = SubnetPacketMapper::default(); let fec_decoder = FecDecoder::new(packet_crypto.clone()); @@ -235,13 +330,8 @@ impl NetworkManager { allow_subnet.clone(), relay_subnets.clone(), )); - let runtime_config = RuntimeConfigController::new( - runtime_policy.clone(), - shared_fec_encoder, - subnet_mapping.clone(), - allow_subnet.clone(), - relay_subnets.clone(), - ); + // shared_fec_encoder 不丢弃:由 manager 保留句柄, + // apply_policy_change 才能双向热切换 fec 开关 let hybrid_outbound = HybridOutbound::new( app_state.network.clone(), @@ -299,19 +389,6 @@ impl NetworkManager { } }; - let mtu_components = MtuTunnelComponents { - hybrid_outbound: hybrid_outbound.clone(), - external_route: subnet_external_route.clone(), - subnet_mapping: subnet_mapping.clone(), - subnet_packet_mapper: subnet_packet_mapper.clone(), - allow_subnet: allow_subnet.clone(), - network: app_state.network.clone(), - no_tun: config.device_mode == DeviceMode::No, - default_interface: default_interface.clone(), - port_mapping_manager: port_mapping_manager.clone(), - policy: runtime_policy.clone(), - runtime_config: runtime_config.clone(), - }; let tunnel_components = TunnelComponents { hybrid_outbound: hybrid_outbound.clone(), external_route: subnet_external_route.clone(), @@ -320,9 +397,8 @@ impl NetworkManager { internal_nat_inbound, port_mapping_manager, policy: runtime_policy.clone(), - runtime_config: runtime_config.clone(), }; - let (enhanced_inbound, enhanced_outbound, quic_client) = enhanced_ipv4_tunnel( + let (enhanced_inbound, enhanced_outbound, _quic_client) = enhanced_ipv4_tunnel( app_state.clone(), mtu_task_group.clone(), task_group.clone(), @@ -337,17 +413,6 @@ impl NetworkManager { ) .await?; - runtime_config.attach_mtu_components(MtuComponentController { - app_state: app_state.clone(), - root_task_group: task_group.clone(), - active_scope: std::sync::Arc::new(tokio::sync::Mutex::new(mtu_task_group)), - tun_data_inbound: enhanced_tun_inbound, - components: mtu_components, - enhanced_inbound: enhanced_inbound.clone(), - enhanced_outbound: enhanced_outbound.clone(), - quic_client, - }); - if let Some(p2p_task) = p2p_task { let handler = P2pInboundHandler::new(P2pInboundConfig { network_route: NetworkRoute::new( @@ -374,75 +439,135 @@ impl NetworkManager { p2p_task.start(handler); } + // 运行期服务端连接登记:初始 server_id 与地址按配置顺序登记, + // 连接任务由注册任务挂入;之后可通过 apply_server_change 增删。 + let server_link_registry = { + let mut registry = ServerLinkRegistry::new(); + for (index, address) in config.server_addr.iter().enumerate() { + registry.insert_initial(index as u32, address.clone()); + } + registry + }; + let server_links = Arc::new(Mutex::new(ServerLinks::new( + server_link_registry, + packet_crypto.clone(), + puncher.clone(), + enhanced_inbound.clone(), + fec_decoder.clone(), + runtime_policy.clone(), + basic_outbound.clone(), + default_interface.clone(), + node_identity.clone(), + client_instance_id.clone(), + app_state.network.clone(), + ))); + let registration_context = Box::new(RegistrationContext { - server_task_group, + server_task_group: server_task_group.clone(), + server_links: server_links.clone(), server_managers: server_manager_list, - ip_update: ip_update.clone(), subnet_external_route, puncher, packet_crypto, enhanced_inbound, fec_decoder, - policy: runtime_policy, - basic_outbound: basic_outbound.clone(), - }); - - runtime_config.attach_server_components(ServerComponentController { - app_state: app_state.clone(), - root_task_group: task_group.clone(), - active_scope: std::sync::Arc::new(tokio::sync::Mutex::new( - registration_context.server_task_group.clone(), - )), - registration_ip, - default_interface, - identity: node_identity.clone(), - packet_crypto: registration_context.packet_crypto.clone(), - external_route: registration_context.subnet_external_route.clone(), - client_instance_id, - puncher: registration_context.puncher.clone(), - enhanced_inbound: registration_context.enhanced_inbound.clone(), - fec_decoder: registration_context.fec_decoder.clone(), + policy: runtime_policy.clone(), basic_outbound: basic_outbound.clone(), - server_outbound: tunnel_to_server, - server_rpc: server_rpc.clone(), }); app_state.set_config(config.clone()); + let fixed_ip = config.ip; + let (registration_status, registration_status_rx) = + tokio::sync::watch::channel(RegistrationStatus::Pending); + // 事件脚本路径来自构造配置(变更它需要重建实例,由 + // needs_instance_rebuild 覆盖) let event_script = EventScript::new(config.event_script.clone()); let manager = Self { config, - event_script, - app_state, - task_group, + app_state: app_state.clone(), + log: log.clone(), + task_group: task_group.clone(), device_io_manager, ip_update, - node_identity, + event_script, enhanced_outbound, server_rpc, - runtime_config, + server_task_group, + server_links, tun_receiver, - registration_context: Some(registration_context), + registration_status: registration_status_rx, + runtime_policy: runtime_policy.clone(), + fec_encoder: shared_fec_encoder, }; - #[cfg(target_os = "android")] - manager.start_android_tun_route_watch(); + // 连接服务器并注册:注册在实例任务组内后台进行,调用方无需关心注册 + // 动作,通过 current_network 获取网段信息即可 + manager.task_group.spawn(async move { + Self::registration_task( + app_state, + registration_context, + fixed_ip, + registration_status, + log, + ) + .await; + }); Ok(manager) } - /// Register with server(s) and start data handling tasks. - /// Returns the registration response on success. - /// On connection-level failure the internal state is kept, so the call can be retried. - pub async fn register(&mut self) -> anyhow::Result { - let Some(mut ctx) = self.registration_context.take() else { - bail!("register can only be called once"); - }; - match Self::register_impl(&self.app_state, &mut ctx, self.config.ip).await { - Ok(response) => Ok(response), - Err(e) => { - // 注册失败时归还上下文,允许调用方重试 - self.registration_context = Some(ctx); - Err(e) + /// 后台注册任务:连接服务器并注册,把结果写入共享状态与 watch 通道。 + /// 连接级错误按 5 秒间隔重试;服务端明确拒绝时记录错误并不再重试。 + async fn registration_task( + app_state: AppState, + mut ctx: Box, + fixed_ip: Option, + status: tokio::sync::watch::Sender, + log: Arc, + ) { + loop { + match Self::register_impl(&app_state, &mut ctx, fixed_ip).await { + Ok(RegisterResponse::Registered) => { + if let Some(addr) = app_state.network.get() { + log.info(format!("注册成功 {}/{}", addr.ip, addr.prefix_len)); + } + let _ = status.send(RegistrationStatus::Ready); + return; + } + Ok(RegisterResponse::Failed(e)) => { + log::error!("注册失败: {}", e.message); + log.error(format!("注册失败: {}", e.message)); + let _ = status.send(RegistrationStatus::Failed(e.message)); + return; + } + Err(e) => { + log::error!("Register failed: {e:?}, 5 秒后重试"); + log.warn(format!("连接服务器失败,5 秒后重试: {e:#}")); + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + } + } + } + + /// 获取当前网络(网段信息)。已配置时立即返回; + /// 否则等待 create_network 启动的后台注册结果。 + pub async fn current_network(&self) -> anyhow::Result { + if let Some(addr) = self.app_state.network.get() { + return Ok(addr); + } + let mut status = self.registration_status.clone(); + loop { + match &*status.borrow_and_update() { + RegistrationStatus::Ready => break, + RegistrationStatus::Failed(message) => bail!("{message}"), + RegistrationStatus::Pending => {} + } + if status.changed().await.is_err() { + bail!("网络注册在完成前被停止"); } } + self.app_state + .network + .get() + .context("network is not registered") } async fn register_impl( @@ -451,7 +576,7 @@ impl NetworkManager { fixed_ip: Option, ) -> anyhow::Result { let mut initially_connected_server = None; - let network_addr = if let Some(fixed_ip) = fixed_ip { + if let Some(fixed_ip) = fixed_ip { let network = fixed_ip.network(); let addr = NetworkAddr { gateway: None, @@ -460,9 +585,7 @@ impl NetworkManager { prefix_len: fixed_ip.prefix_len(), }; app_state.network.set(addr); - ctx.ip_update.set_registered_ip(addr.ip); log::info!("Local fixed network activated: {fixed_ip}"); - addr } else { let is_multi_server = ctx.server_managers.len() > 1; let (server_index, response) = if is_multi_server { @@ -496,6 +619,11 @@ impl NetworkManager { crate::protocol::control_message::ResponseMessage::SubscriptionConfig(_) => { bail!("Unexpected subscription configuration response during registration"); } + crate::protocol::control_message::ResponseMessage::SubscriptionRegister(_) + | crate::protocol::control_message::ResponseMessage::SubscriptionPush(_) + | crate::protocol::control_message::ResponseMessage::SubscriptionPong(_) => { + bail!("Unexpected subscription control response during traffic registration"); + } }; let addr = NetworkAddr { gateway: Some(reg_response.gateway), @@ -504,7 +632,6 @@ impl NetworkManager { prefix_len: reg_response.prefix_len, }; app_state.network.set(addr); - ctx.ip_update.set_registered_ip(addr.ip); initially_connected_server = Some(server_index); if !reg_response.server_version.is_empty() { app_state @@ -516,11 +643,10 @@ impl NetworkManager { reg_response.server_instance_id, reg_response.multi_link_supported, ); - addr - }; + } // Start data handling tasks for all servers - for (index, turn_manager) in ctx.server_managers.drain(..).enumerate() { + for (server_id, turn_manager) in ctx.server_managers.drain(..).enumerate() { let handler_config = Box::new(InboundHandlerConfig { network_route: NetworkRoute::new( app_state.network.clone(), @@ -538,14 +664,18 @@ impl NetworkManager { basic_outbound: ctx.basic_outbound.clone(), app_state: app_state.clone(), }); - turn_manager.data_handle_task( + let task = turn_manager.data_handle_task( &ctx.server_task_group, handler_config, - initially_connected_server == Some(index), + initially_connected_server == Some(server_id), ); + ctx.server_links + .lock() + .registry + .attach_task(server_id as u32, task); } - Ok(RegisterResponse::Success(network_addr)) + Ok(RegisterResponse::Registered) } pub fn device_mode(&self) -> DeviceMode { @@ -556,39 +686,44 @@ impl NetworkManager { if self.tun_receiver.is_none() || self.enhanced_outbound.is_none() { bail!("start_device requires tun/tap mode and can only be called once"); } + let network = self.current_network().await?; + let routes = self.app_state.subnet_route.all_route(); let mut config = DeviceConfig::default(); config = config .set_device_mode(self.config.device_mode) .set_mtu(self.config.mtu.unwrap_or(DEFAULT_MTU)); if self.config.device_mode == DeviceMode::Tap { - let net = self - .app_state - .get_network() - .context("network is not registered")?; - config = config.set_mac_addr(crate::ethernet::mac_from_ip(net.ip).octets()); + config = config.set_mac_addr(crate::ethernet::mac_from_ip(network.ip).octets()); } if let Some(tun_name) = self.config.tun_name.clone() { config = config.set_tun_name(tun_name); } // 失败时 tun_receiver/enhanced_outbound 不会被消耗,可以重试 - self.device_io_manager + if let Err(e) = self + .device_io_manager .start_task(config, &mut self.tun_receiver, &mut self.enhanced_outbound) .await + { + self.log.error(format!("启动虚拟网卡失败: {e:#}")); + return Err(e); + } + self.apply_network_state(&network, routes).await } #[cfg(unix)] - pub async fn start_device_fd(&mut self, tun_fd: Option) -> anyhow::Result<()> { + pub async fn start_device_fd( + &mut self, + tun_fd: Option, + ) -> anyhow::Result<()> { if self.tun_receiver.is_none() || self.enhanced_outbound.is_none() { bail!("start_device_fd requires tun/tap mode and can only be called once"); } + let network = self.current_network().await?; + let routes = self.app_state.subnet_route.all_route(); let mut config = DeviceConfig::default() .set_device_mode(self.config.device_mode) .set_mtu(self.config.mtu.unwrap_or(DEFAULT_MTU)); if self.config.device_mode == DeviceMode::Tap { - let net = self - .app_state - .get_network() - .context("network is not registered")?; - config = config.set_mac_addr(crate::ethernet::mac_from_ip(net.ip).octets()); + config = config.set_mac_addr(crate::ethernet::mac_from_ip(network.ip).octets()); } if let Some(tun_fd) = tun_fd { config = config.set_tun_fd(tun_fd); @@ -596,119 +731,83 @@ impl NetworkManager { if let Some(tun_name) = self.config.tun_name.clone() { config = config.set_tun_name(tun_name); } - self.device_io_manager + if let Err(e) = self + .device_io_manager .start_task(config, &mut self.tun_receiver, &mut self.enhanced_outbound) .await - } - #[cfg(not(target_os = "android"))] - pub async fn set_device_network_ip(&self, ip: Ipv4Addr, prefix_len: u8) -> anyhow::Result<()> { - // 服务端数据处理任务早于虚拟网卡初始化启动。启动期间配置协调器可能已 - // 应用了新地址,因此必须以共享状态中的最新地址为准,不能再用最初注册 - // 响应覆盖它。 - let (ip, prefix_len) = self - .app_state - .get_network() - .map(|network| (network.ip, network.prefix_len)) - .unwrap_or((ip, prefix_len)); - self.device_io_manager.set_network(ip, prefix_len).await?; - #[cfg(not(any(target_os = "ios", target_os = "tvos")))] - self.device_io_manager - .start_system_routes(self.app_state.subnet_route.subscribe()) - .await?; - // 网卡设置成功(应用 IP)后触发事件脚本 - if let Some(network) = self.app_state.get_network() { - let server = self - .config - .server_addr - .iter() - .map(|addr| addr.to_string()) - .collect::>() - .join(","); - self.event_script - .notify( - EventScriptType::NetCardCreated, - &[ - ("ip", network.ip.to_string()), - ("prefix-length", network.prefix_len.to_string()), - ( - "gateway", - network - .gateway - .map(|gateway| gateway.to_string()) - .unwrap_or_else(|| "-".to_string()), - ), - ("broadcast", network.broadcast.to_string()), - ("server", server), - ], - ) - .await; + { + self.log.error(format!("启动虚拟网卡失败: {e:#}")); + return Err(e); } - Ok(()) - } - - #[cfg(target_os = "android")] - fn start_android_tun_route_watch(&self) { - let mut routes = self.app_state.subnet_route.subscribe(); - let ip_update = self.ip_update.clone(); - self.task_group.spawn(async move { - loop { - if routes.changed().await.is_err() { - break; - } - let snapshot = routes.borrow_and_update().clone(); - if let Err(error) = ip_update.request_android_route_rebuild(snapshot).await { - log::warn!("请求 Android TUN 路由重建失败: {error:#}"); - } - } - }); + self.apply_network_state(&network, routes).await } - #[cfg(target_os = "android")] - pub fn tun_rebuild_coordinator( - &self, - ) -> crate::tunnel_core::server::inbound::TunRebuildCoordinator { - self.ip_update.tun_rebuild_coordinator() + /// Applies the registered network address together with the current + /// complete inbound route set, without creating a virtual device. + pub async fn apply_initial_network_info(&mut self) -> anyhow::Result<()> { + let network = self.current_network().await?; + let routes = self.app_state.subnet_route.all_route(); + self.apply_network_state(&network, routes).await } - #[cfg(target_os = "android")] - pub async fn replace_tun_task( - &self, - request_id: u64, - tun_fd: std::os::fd::OwnedFd, + /// Commits the registered address and complete route set to the platform. + async fn apply_network_state( + &mut self, + network: &NetworkAddr, + routes: Vec, ) -> anyhow::Result<()> { - self.ip_update - .replace_android_tun_task(request_id, tun_fd) - .await - } - - #[cfg(target_os = "android")] - pub async fn reject_tun_rebuild(&self, request_id: u64, reason: String) -> anyhow::Result<()> { - self.ip_update - .reject_android_tun_rebuild(request_id, reason) - .await + #[cfg(target_os = "android")] + let _ = network; + // 先提交转发路由表(配置即新值),再执行网卡与系统路由变更动作; + // 动作失败不回滚,错误上抛由调用方记录 + self.app_state.subnet_route.apply_routes(routes.clone()); + #[cfg(not(target_os = "android"))] + if self.device_mode().has_device() { + self.device_io_manager + .set_network(network.ip, network.prefix_len) + .await?; + #[cfg(not(any(target_os = "ios", target_os = "tvos")))] + self.device_io_manager + .apply_system_routes(routes) + .await?; + } + Ok(()) } fn stop_network(&mut self) { - #[cfg(target_os = "android")] - self.ip_update.close_android_tun_rebuild(); self.task_group.stop(); self.app_state.stop_network(); } - #[cfg(not(target_os = "android"))] - pub async fn device_if_index(&self) -> anyhow::Result { - self.device_io_manager.device_if_index().await - } pub async fn wait_all_stopped(&mut self) { self.task_group.wait_all_stopped().await; } pub fn vnt_api(&self) -> VntApi { - VntApi::new( - self.app_state.clone(), - self.server_rpc.clone(), - self.ip_update.clone(), - self.node_identity.clone(), - self.runtime_config.clone(), - ) + VntApi::new(self.app_state.clone(), self.server_rpc.clone()) + } +} + +/// 把子网路由源的最新完整生效快照转发给运行期变化消费者。 +/// 消费端被丢弃或路由源关闭时退出。 +async fn forward_network_route_changes( + mut route_changes: tokio::sync::watch::Receiver>, + routes_tx: tokio::sync::watch::Sender>, +) { + loop { + let routes = route_changes.borrow_and_update().clone(); + if routes_tx.is_closed() { + break; + } + routes_tx.send_if_modified(|current| { + if *current == routes { + false + } else { + *current = routes; + true + } + }); + if route_changes.changed().await.is_err() { + break; + } } } impl Drop for NetworkManager { @@ -717,11 +816,96 @@ impl Drop for NetworkManager { } } +#[cfg(test)] +mod network_route_change_tests { + use super::forward_network_route_changes; + use crate::context::config::Config; + use crate::nat::SubnetExternalRoute; + use crate::network_info::RuntimeChangeManager; + use std::time::Duration; + + #[tokio::test] + async fn subnet_sync_route_source_wakes_the_unified_listener() { + let routes = SubnetExternalRoute::default(); + let (routes_tx, mut routes_rx) = tokio::sync::watch::channel(Vec::new()); + let task = tokio::spawn(forward_network_route_changes( + routes.subscribe(), + routes_tx, + )); + + let route: crate::nat::NetInput = "192.168.50.0/24,10.26.0.3".parse().unwrap(); + routes.set_automatic_routes(vec![route.clone()]); + // 独立模式(无订阅)下 changed_with 等待路由源变更并返回完整快照 + let change = + tokio::time::timeout( + Duration::from_secs(1), + RuntimeChangeManager::changed_with(&mut None, &mut routes_rx, &Config::default()), + ) + .await + .unwrap() + .unwrap(); + assert_eq!(change.routes, vec![route]); + // 独立模式下配置保持构造时传入的本地配置 + assert_eq!(change.config.network_code, Config::default().network_code); + task.abort(); + } +} + +#[cfg(test)] +mod inspect_change_tests { + use super::InstanceLog; + use crate::context::config::{Config, DeviceMode}; + use crate::network_info::{RuntimeChange, RuntimeChangeManager}; + use std::sync::Arc; + + fn standalone_config() -> Config { + Config { + network_code: "inspect-change-test".to_string(), + device_id: "inspect-device".to_string(), + device_name: "inspect-node".to_string(), + ip: Some("10.26.0.2/24".parse().unwrap()), + device_mode: DeviceMode::No, + ..Config::default() + } + } + + #[tokio::test] + async fn identical_snapshot_needs_nothing_and_password_flags_rebuild() { + let config = standalone_config(); + let manager = RuntimeChangeManager::new( + config.clone(), + None, + Arc::new(InstanceLog::new("inspect-change-test")), + ) + .await + .unwrap(); + + // 相同快照:两个标志都为 false + let plan = manager.inspect_change(&RuntimeChange { + config: config.clone(), + routes: Vec::new(), + }); + assert!(!plan.instance_rebuild, "相同快照不应要求重建实例"); + // 桌面平台的网卡类变更由 apply 内部热应用,vpn_rebuild 恒为 false + assert!(!plan.vpn_rebuild); + + // 密码只能重建实例生效:instance_rebuild 置位,vpn_rebuild 不受影响 + let mut changed = config; + changed.password = Some("another-password".to_string()); + let plan = manager.inspect_change(&RuntimeChange { + config: changed, + routes: Vec::new(), + }); + assert!(plan.instance_rebuild); + assert!(!plan.vpn_rebuild); + } +} + #[cfg(test)] mod decentralized_loopback_tests { - use super::{NetworkManager, RegisterResponse}; + use super::InstanceLog; use crate::context::config::{Config, DeviceMode}; - use crate::utils::task_control::{TaskGroupGuard, TaskGroupManager}; + use crate::network_info::RuntimeChangeManager; use std::net::{Ipv4Addr, SocketAddr, TcpListener, UdpSocket}; use std::time::Duration; @@ -747,11 +931,8 @@ mod decentralized_loopback_tests { let ports = reserve_loopback_ports(4); let peers = [ports[1], ports[2], ports[3], ports[2]]; let mut managers = Vec::new(); - let mut guards: Vec = Vec::new(); for index in 0..4 { - let manager = TaskGroupManager::new(); - let (task_group, guard) = manager.create_task().unwrap(); let ip = Ipv4Addr::new(10, 26, 0, index as u8 + 2); let config = Config { network_code: "loopback-gossip-test".to_string(), @@ -765,23 +946,24 @@ mod decentralized_loopback_tests { peer_address: vec![format!("tcp://127.0.0.1:{}", peers[index]).parse().unwrap()], ..Config::default() }; - let mut network = NetworkManager::create_network(Box::new(config), task_group) - .await - .unwrap(); - assert!(matches!( - network.register().await.unwrap(), - RegisterResponse::Success(_) - )); - managers.push(network); - guards.push(guard); + // 配置管理器创建并持有组网实例(含路由通道与任务组) + let manager = RuntimeChangeManager::new( + config, + None, + std::sync::Arc::new(InstanceLog::new("loopback-test")), + ) + .await + .unwrap(); + manager.current_network().await.unwrap(); + managers.push(manager); } // Passive graph discovery waits for the first 25-35 second periodic // announcement; route creation no longer triggers an immediate one. tokio::time::timeout(Duration::from_secs(45), async { loop { - let a = managers[0].vnt_api(); - let d = managers[3].vnt_api(); + let a = managers[0].api().unwrap(); + let d = managers[3].api().unwrap(); let a_to_d = a.find_route(&Ipv4Addr::new(10, 26, 0, 5)); let d_to_a = d.find_route(&Ipv4Addr::new(10, 26, 0, 2)); let a_knows_d = a @@ -801,11 +983,12 @@ mod decentralized_loopback_tests { .unwrap_or_else(|_| { panic!( "loopback graph did not converge: A={:?}, D={:?}", - managers[0].vnt_api().route_table(), - managers[3].vnt_api().route_table() + managers[0].api().unwrap().route_table(), + managers[3].api().unwrap().route_table() ) }); - drop(guards); + // 管理器 drop 即停止组网实例 + drop(managers); } } diff --git a/vnt-core/src/enhanced_tunnel/inbound.rs b/vnt-core/src/enhanced_tunnel/inbound.rs index 408775fe..6752cf6a 100644 --- a/vnt-core/src/enhanced_tunnel/inbound.rs +++ b/vnt-core/src/enhanced_tunnel/inbound.rs @@ -65,10 +65,6 @@ impl EnhancedInbound { })), } } - - pub(crate) fn replace_from(&self, prepared: &Self) { - self.inner.store(prepared.inner.load_full()); - } } impl EnhancedInboundInner { diff --git a/vnt-core/src/enhanced_tunnel/mod.rs b/vnt-core/src/enhanced_tunnel/mod.rs index abce7523..58a4ec23 100644 --- a/vnt-core/src/enhanced_tunnel/mod.rs +++ b/vnt-core/src/enhanced_tunnel/mod.rs @@ -5,13 +5,12 @@ use crate::enhanced_tunnel::outbound::EnhancedOutbound; use crate::ethernet::MacTable; use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager}; use crate::nat::subnet_packet::SubnetPacketMapper; -use crate::nat::{AllowSubnetExternalRoute, SubnetExternalRoute, SubnetMappingTable}; +use crate::nat::{SubnetExternalRoute, SubnetMappingTable}; use crate::port_mapping::PortMapping; -use crate::runtime_config::{RuntimeConfigController, RuntimePolicyStore}; +use crate::runtime_config::RuntimePolicyStore; use crate::tun::enhanced_tun::EnhancedTunInbound; use crate::tunnel_core::outbound::HybridOutbound; use crate::utils::task_control::TaskGroup; -use rustp2p_core::socket::LocalInterface; pub(crate) mod quic_over; @@ -35,30 +34,6 @@ pub(crate) struct TunnelComponents { pub internal_nat_inbound: Option, pub port_mapping_manager: PortMappingManager, pub policy: RuntimePolicyStore, - pub runtime_config: RuntimeConfigController, -} - -/// Inputs retained by the runtime controller to prepare an MTU-dependent -/// replacement without touching the stable network instance. -#[derive(Clone)] -pub(crate) struct MtuTunnelComponents { - pub hybrid_outbound: HybridOutbound, - pub external_route: SubnetExternalRoute, - pub subnet_mapping: SubnetMappingTable, - pub subnet_packet_mapper: SubnetPacketMapper, - pub allow_subnet: AllowSubnetExternalRoute, - pub network: crate::context::SharedNetworkAddr, - pub no_tun: bool, - pub default_interface: Option, - pub port_mapping_manager: PortMappingManager, - pub policy: RuntimePolicyStore, - pub runtime_config: RuntimeConfigController, -} - -pub(crate) struct PreparedMtuTunnel { - pub inbound: EnhancedInbound, - pub outbound: Option, - pub quic_client: quic_over::quic_client::QuicTunnelClient, } pub(crate) async fn enhanced_ipv4_tunnel( @@ -72,31 +47,6 @@ pub(crate) async fn enhanced_ipv4_tunnel( EnhancedInbound, Option, quic_over::quic_client::QuicTunnelClient, -)> { - build_enhanced_ipv4_tunnel( - app_state, - task_group.clone(), - tun_data_sender, - config, - components, - true, - Some(port_mapping_root), - ) - .await -} - -async fn build_enhanced_ipv4_tunnel( - app_state: AppState, - task_group: TaskGroup, - tun_data_sender: EnhancedTunInbound, - config: TunnelConfig, - components: TunnelComponents, - initialize_port_mappings: bool, - port_mapping_root: Option, -) -> anyhow::Result<( - EnhancedInbound, - Option, - quic_over::quic_client::QuicTunnelClient, )> { let password = config.password.unwrap_or_else(|| "password".to_string()); let tun = match &tun_data_sender { @@ -119,9 +69,7 @@ async fn build_enhanced_ipv4_tunnel( internal_nat_manager: components.internal_nat_inbound.clone(), port_mapping_manager: components.port_mapping_manager, policy: components.policy.clone(), - runtime_config: components.runtime_config, }, - initialize_port_mappings, port_mapping_root, ) .await?; @@ -150,59 +98,3 @@ async fn build_enhanced_ipv4_tunnel( }); Ok((enhanced_inbound, enhanced_outbound, quic_client)) } - -pub(crate) async fn prepare_mtu_tunnel( - app_state: AppState, - task_group: TaskGroup, - tun_data_sender: EnhancedTunInbound, - config: TunnelConfig, - components: MtuTunnelComponents, -) -> anyhow::Result { - let internal_nat_inbound = Some( - InternalNatInbound::create( - &task_group, - config.mtu, - components.hybrid_outbound.clone(), - components.allow_subnet.clone(), - components.network.clone(), - components.no_tun, - components.default_interface.clone(), - ) - .await?, - ); - // In no-device mode the enhanced TUN input is the internal NAT stack - // itself, so it must follow the candidate MTU plane rather than retaining - // the old stack through the stable dispatcher. - let tun_data_sender = match tun_data_sender { - EnhancedTunInbound::Nat(_) => EnhancedTunInbound::Nat( - internal_nat_inbound - .clone() - .expect("internal NAT is constructed for every MTU plane"), - ), - other => other, - }; - let (inbound, outbound, quic_client) = build_enhanced_ipv4_tunnel( - app_state, - task_group, - tun_data_sender, - config, - TunnelComponents { - hybrid_outbound: components.hybrid_outbound, - external_route: components.external_route, - subnet_mapping: components.subnet_mapping, - subnet_packet_mapper: components.subnet_packet_mapper, - internal_nat_inbound, - port_mapping_manager: components.port_mapping_manager, - policy: components.policy, - runtime_config: components.runtime_config, - }, - false, - None, - ) - .await?; - Ok(PreparedMtuTunnel { - inbound, - outbound, - quic_client, - }) -} diff --git a/vnt-core/src/enhanced_tunnel/outbound.rs b/vnt-core/src/enhanced_tunnel/outbound.rs index 8fbb1643..e5ef1558 100644 --- a/vnt-core/src/enhanced_tunnel/outbound.rs +++ b/vnt-core/src/enhanced_tunnel/outbound.rs @@ -74,10 +74,6 @@ impl EnhancedOutbound { } } - pub(crate) fn replace_from(&self, prepared: &Self) { - self.inner.store(prepared.inner.load_full()); - } - pub async fn ipv4_outbound(&self, data: TransmissionBytes) { let inner = self.inner.load(); inner.ipv4_outbound(data).await; diff --git a/vnt-core/src/enhanced_tunnel/quic_over/boot.rs b/vnt-core/src/enhanced_tunnel/quic_over/boot.rs index dd185dc6..1d77fdf9 100644 --- a/vnt-core/src/enhanced_tunnel/quic_over/boot.rs +++ b/vnt-core/src/enhanced_tunnel/quic_over/boot.rs @@ -11,9 +11,7 @@ use crate::enhanced_tunnel::quic_over::{quic_client, quic_server}; use crate::nat::internal_nat::{InternalNatInbound, PortMappingManager}; use crate::nat::{SubnetExternalRoute, SubnetMappingTable}; use crate::port_mapping::PortMapping; -use crate::runtime_config::{ - PortMappingComponentController, RuntimeConfigController, RuntimePolicyStore, -}; +use crate::runtime_config::RuntimePolicyStore; use crate::tls; use crate::tun::TunDataInbound; use crate::tunnel_core::outbound::HybridOutbound; @@ -43,7 +41,6 @@ pub(crate) struct QuicTunnelComponents { pub internal_nat_manager: Option, pub port_mapping_manager: PortMappingManager, pub policy: RuntimePolicyStore, - pub runtime_config: RuntimeConfigController, } pub(crate) async fn quic_tunnel_start( @@ -52,8 +49,7 @@ pub(crate) async fn quic_tunnel_start( tun_data_sender: Option, config: QuicTunnelConfig, components: QuicTunnelComponents, - initialize_port_mappings: bool, - port_mapping_root: Option, + port_mapping_root: TaskGroup, ) -> anyhow::Result<( EnhancedQuicInbound, Option, @@ -106,35 +102,12 @@ pub(crate) async fn quic_tunnel_start( ) .await; } - if initialize_port_mappings { - let port_mapping_root = port_mapping_root - .context("port mapping root task group is required during initialization")?; - let mut active_mappings: Vec<(PortMapping, TaskGroup)> = Vec::new(); - for mapping in &config.port_mapping { - let scope = port_mapping_root.child_scope(); - if let Err(error) = crate::port_mapping::port_mapping_start( - &scope, - vec![mapping.clone()], - quic_client.clone(), - ) - .await - { - scope.stop(); - for (_, prepared_scope) in active_mappings { - prepared_scope.stop(); - } - return Err(error); - } - active_mappings.push((mapping.clone(), scope)); - } - components - .runtime_config - .attach_port_mapping_components(PortMappingComponentController { - root_task_group: port_mapping_root, - active: Arc::new(tokio::sync::Mutex::new(active_mappings)), - quic_client: quic_client.clone(), - }); - } + crate::port_mapping::port_mapping_start( + &port_mapping_root, + config.port_mapping, + quic_client.clone(), + ) + .await?; let quic_inbound = EnhancedQuicInbound::new(inbound); Ok((quic_inbound, quic_outbound, quic_client)) diff --git a/vnt-core/src/enhanced_tunnel/quic_over/quic_client.rs b/vnt-core/src/enhanced_tunnel/quic_over/quic_client.rs index b8ecfad4..a947b751 100644 --- a/vnt-core/src/enhanced_tunnel/quic_over/quic_client.rs +++ b/vnt-core/src/enhanced_tunnel/quic_over/quic_client.rs @@ -55,10 +55,6 @@ impl QuicTunnelClient { } } - pub(crate) fn replace_from(&self, prepared: &Self) { - self.inner.store(prepared.inner.load_full()); - } - pub async fn open_bi(&self, dest: Ipv4Addr) -> anyhow::Result<(SendStream, RecvStream)> { let inner = self.inner.load(); inner.open_bi(dest).await diff --git a/vnt-core/src/event_script.rs b/vnt-core/src/event_script.rs index dc5456f4..f7ba91e0 100644 --- a/vnt-core/src/event_script.rs +++ b/vnt-core/src/event_script.rs @@ -1,11 +1,10 @@ -/// 事件脚本:在指定事件(网卡创建成功、掉线、重连成功、IP 变化)发生时调用外部脚本。 +/// 事件脚本:在掉线或重连成功时调用外部脚本。 /// /// 脚本通过命令行参数接收事件名和事件数据,例如: /// ```text -/// @@ -187,6 +216,13 @@ onMounted(async () => { await app.fetchConfigList(); await refreshSubscriptionSt
+
+
+
+ VNT 订阅配置二维码 +
+
+

{{ qrSubscription }}

+
+
+ +
+

二维码包含接入凭据,请仅分享给可信设备。

+
From e686748b59c305fc5ba962fd505f5eeda7fa72bb Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Sat, 26 Sep 2026 11:11:42 +0800 Subject: [PATCH 6/9] =?UTF-8?q?feat(web):=20=E5=90=AF=E5=8A=A8=E7=BB=84?= =?UTF-8?q?=E7=BD=91=E9=9D=A2=E6=9D=BF=E6=96=B0=E5=A2=9E=E9=85=8D=E7=BD=AE?= =?UTF-8?q?=E9=A2=84=E8=A7=88=E6=8C=89=E9=92=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 选择配置下拉框旁新增预览按钮,点击后弹窗展示该配置文件的 TOML 原文, 支持一键复制;读取失败时 toast 提示并关闭弹窗。 --- vnt-web/ui/src/components/StartPanel.vue | 57 ++++++++++++++++++++++++ 1 file changed, 57 insertions(+) diff --git a/vnt-web/ui/src/components/StartPanel.vue b/vnt-web/ui/src/components/StartPanel.vue index d89b3460..1ef23238 100644 --- a/vnt-web/ui/src/components/StartPanel.vue +++ b/vnt-web/ui/src/components/StartPanel.vue @@ -2,13 +2,19 @@ import { ref, computed, watch } from "vue"; import { useAppStore } from "../stores/app"; import { useUiStore } from "../stores/ui"; +import { getConfig } from "../api"; import AppSelect from "./AppSelect.vue"; +import AppModal from "./AppModal.vue"; // 启动组网面板:选择配置 + 启动,总览页与实例页复用 const app = useAppStore(); const ui = useUiStore(); const localSelectedConfig = ref(""); +const showPreview = ref(false); +const previewText = ref(""); +const previewName = ref(""); +const previewLoading = ref(false); // 只列出没有对应实例的配置(同一配置最多一个实例) const availableConfigs = computed(() => @@ -41,6 +47,31 @@ const handleStart = () => { } app.startVnt(localSelectedConfig.value); }; + +// 预览配置文件的原始内容(尚未启动,没有合并后的生效配置) +const openPreview = async (fileName) => { + previewName.value = fileName; + previewText.value = ""; + previewLoading.value = true; + showPreview.value = true; + try { + previewText.value = await getConfig(fileName); + } catch (e) { + ui.toast.error(e.message); + showPreview.value = false; + } finally { + previewLoading.value = false; + } +}; + +const copyPreview = async () => { + try { + await navigator.clipboard.writeText(previewText.value); + ui.toast.success("已复制"); + } catch { + ui.toast.error("复制失败"); + } +}; From 15ca3a4565b6855b65b87e7eb2690ae8f67868d2 Mon Sep 17 00:00:00 2001 From: lbl <1791778603@qq.com> Date: Sat, 26 Sep 2026 11:11:46 +0800 Subject: [PATCH 7/9] =?UTF-8?q?feat(web):=20=E7=BB=84=E7=BD=91=E5=AE=9E?= =?UTF-8?q?=E4=BE=8B=E5=8D=A1=E7=89=87=E6=96=B0=E5=A2=9E=E7=BC=96=E8=BE=91?= =?UTF-8?q?=E6=8C=89=E9=92=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 实例卡片顶部新增编辑按钮,点击跳转到组网配置页,通过 ?edit=文件名 打开对应配置的编辑器,消费参数后立即清理 URL,避免刷新重复弹出。 ConfigEditor 加载配置的 watcher 补充 immediate:从总览页跳转时编辑器 带着 show=true 首次挂载,非 immediate 的 watcher 不会触发,导致配置 不加载、弹窗显示为空的新建表单。 --- vnt-web/ui/src/components/InstanceCard.vue | 8 ++++++++ vnt-web/ui/src/views/ConfigEditor.vue | 3 ++- vnt-web/ui/src/views/ConfigView.vue | 16 +++++++++++++++- 3 files changed, 25 insertions(+), 2 deletions(-) diff --git a/vnt-web/ui/src/components/InstanceCard.vue b/vnt-web/ui/src/components/InstanceCard.vue index 1e2f8832..ccd08d0c 100644 --- a/vnt-web/ui/src/components/InstanceCard.vue +++ b/vnt-web/ui/src/components/InstanceCard.vue @@ -1,5 +1,6 @@