Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 6 additions & 6 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "vnt2"
version = "2.0.9"
version = "2.0.10"
edition = "2024"
license = "Apache-2.0"

Expand Down
8 changes: 4 additions & 4 deletions src/args_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -625,7 +625,7 @@ impl FileConfig {
# ==================================

# 可选的服务端配置源。本文件中明确填写的字段会覆盖订阅链接下发的同名字段。
# subscription = "vnt2://join/1/..."
# subscription = "vnt2://join/2/..."

# --- 网络配置 ---
# 网络编号,相同网络编号的会组在同一个虚拟网 (必填)
Expand Down Expand Up @@ -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]
Expand All @@ -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);
Expand Down
102 changes: 102 additions & 0 deletions src/cli_subscription.rs
Original file line number Diff line number Diff line change
@@ -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<FileConfig>,
subscription: Option<Subscription>,
log: Arc<InstanceLog>,
ctrl_port: Option<u16>,
) -> 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<u16>,
sender: Option<tokio::sync::watch::Sender<VntApi>>,
}

impl IpcPublisher {
fn new(ctrl_port: Option<u16>) -> Self {
Self {
ctrl_port,
sender: None,
}
}

fn publish(&mut self, api: Option<VntApi>) {
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:?}");
}
});
}
}
Loading
Loading