Merge upstream main into shared virtual NIC work

Port shared TUN ownership, routing, mobile source dispatch, and
Magic DNS route claims onto the native runtime-host architecture.

Preserve per-member lifecycle and add focused desktop, mobile, netlink,
GUI, FFI, and root integration validation for the merged code.
This commit is contained in:
KKRainbow
2026-08-13 00:56:52 +08:00
609 changed files with 130194 additions and 69014 deletions
+2 -2
View File
@@ -5,8 +5,8 @@
"private": true,
"packageManager": "pnpm@9.12.1+sha512.e5a7e52a4183a02d5931057f7a0dbff9d5e9ce3161e33fa68ae392125b79282a8a8a470a51dfc8a0ed86221442eb2fb57019b0990ed24fab519bf0e1bc5ccfc4",
"scripts": {
"dev": "vite",
"build": "vue-tsc --noEmit && vite build",
"dev": "pnpm --dir ../easytier-web/frontend-lib build && vite",
"build": "pnpm --dir ../easytier-web/frontend-lib build && vue-tsc --noEmit && vite build",
"preview": "vite preview",
"tauri": "tauri",
"lint": "eslint . --ignore-pattern src-tauri",
+1
View File
@@ -25,6 +25,7 @@ serde = { version = "1", features = ["derive"] }
serde_json = "1"
easytier = { path = "../../easytier" }
easytier-core = { path = "../../easytier-core" }
tokio = { version = "1", features = ["full"] }
anyhow = "1.0"
chrono = { version = "0.4.37", features = ["serde"] }
+169 -101
View File
@@ -4,30 +4,37 @@
mod elevate;
use anyhow::Context;
#[cfg(any(target_os = "android", target_os = "ios"))]
use easytier::instance::factory::attach_mobile_tun_fd;
#[cfg(target_os = "android")]
use easytier::instance::factory::subscribe_native_instance_event;
use easytier::proto::api::manage::{
CollectNetworkInfoResponse, ValidateConfigResponse, WebClientService,
WebClientServiceClientFactory,
};
use easytier::rpc_service::remote_client::{
GetNetworkMetasResponse, ListNetworkInstanceIdsJsonResp, ListNetworkProps, RemoteClientManager,
Storage,
};
use easytier::web_client::{self, WebClient};
use easytier::{
common::config::{NetworkConfig, NetworkConfigExt},
common::{
config::{
ConfigLoader, ConfigSource, FileLoggerConfig, LoggingConfigBuilder, TomlConfigLoader,
},
log,
},
instance_manager::NetworkInstanceManager,
launcher::NetworkConfig,
instance::factory::{NativeInstanceManager, native_instance_manager},
proto::rpc::standalone::{runtime_rpc_dialer, runtime_rpc_listener},
rpc_service::ApiRpcServer,
tunnel::TunnelListener,
tunnel::ring::RingTunnelListener,
tunnel::tcp::TcpTunnelListener,
utils::panic::setup_panic_handler,
};
use easytier_core::management::config_source_to_rpc;
use easytier_core::management::remote_client::{
GetNetworkMetasResponse, ListNetworkInstanceIdsJsonResp, ListNetworkProps, RemoteClientManager,
Storage,
};
use easytier_core::{
connectivity::protocol::raw::TunnelDialer as _, process_runtime::CoreProcessRuntime,
socket::SocketListener, tunnel::Tunnel,
};
use std::ops::Deref;
use std::sync::Arc;
use tokio::sync::{Mutex, RwLock, RwLockReadGuard};
@@ -38,7 +45,7 @@ use tauri::{AppHandle, Emitter, Manager as _};
#[cfg(not(target_os = "android"))]
use tauri::tray::{MouseButton, MouseButtonState, TrayIconBuilder, TrayIconEvent};
static INSTANCE_MANAGER: once_cell::sync::Lazy<RwLock<Option<Arc<NetworkInstanceManager>>>> =
static INSTANCE_MANAGER: once_cell::sync::Lazy<RwLock<Option<Arc<NativeInstanceManager>>>> =
once_cell::sync::Lazy::new(|| RwLock::new(None));
static RPC_RING_UUID: once_cell::sync::Lazy<uuid::Uuid> =
@@ -47,7 +54,7 @@ static RPC_RING_UUID: once_cell::sync::Lazy<uuid::Uuid> =
static CLIENT_MANAGER: once_cell::sync::Lazy<RwLock<Option<manager::GUIClientManager>>> =
once_cell::sync::Lazy::new(|| RwLock::new(None));
type BoxedTunnelListener = Box<dyn TunnelListener>;
type BoxedTunnelListener = Box<dyn SocketListener<Accepted = Box<dyn Tunnel>>>;
#[derive(Clone, Copy, PartialEq, Eq)]
enum RpcServerKind {
@@ -187,30 +194,31 @@ struct TunFdInstanceSources {
ipv6_routes: Vec<String>,
}
#[cfg(mobile)]
#[cfg(any(target_os = "android", target_os = "ios"))]
fn parse_tun_fd_instance_sources(
instance_sources: Vec<TunFdInstanceSources>,
) -> Result<
std::collections::HashMap<uuid::Uuid, easytier::instance::virtual_nic::MobileTunSources>,
String,
> {
instance_sources
.into_iter()
.map(|source| {
let instance_id = source
.instance_id
.parse::<uuid::Uuid>()
.map_err(|err| format!("invalid instance id {}: {err}", source.instance_id))?;
let sources = easytier::instance::virtual_nic::MobileTunSources::parse(
source.ipv4_addrs,
source.ipv6_addrs,
source.ipv4_routes,
source.ipv6_routes,
)
.map_err(|err| err.to_string())?;
Ok((instance_id, sources))
})
.collect()
let mut parsed = std::collections::HashMap::with_capacity(instance_sources.len());
for source in instance_sources {
let instance_id = source
.instance_id
.parse::<uuid::Uuid>()
.map_err(|err| format!("invalid instance id {}: {err}", source.instance_id))?;
let sources = easytier::instance::virtual_nic::MobileTunSources::parse(
source.ipv4_addrs,
source.ipv6_addrs,
source.ipv4_routes,
source.ipv6_routes,
)
.map_err(|err| err.to_string())?;
if parsed.insert(instance_id, sources).is_some() {
return Err(format!("duplicate tun sources for instance {instance_id}"));
}
}
Ok(parsed)
}
#[tauri::command]
@@ -233,45 +241,105 @@ async fn set_tun_fd(
.collect::<Result<Vec<_>, _>>()?,
_ => get_client_manager!()?.get_enabled_instances_for_tun_fd(),
};
if target_ids.is_empty() {
return Err("no TUN-enabled instance is available for fd attachment".to_string());
}
let target_id_set = target_ids
.iter()
.copied()
.collect::<std::collections::HashSet<_>>();
if target_id_set.len() != target_ids.len() {
return Err("duplicate instance id in TUN fd attachment group".to_string());
}
#[cfg(mobile)]
#[cfg(any(target_os = "android", target_os = "ios"))]
let mut source_map = parse_tun_fd_instance_sources(instance_sources.unwrap_or_default())?;
#[cfg(not(mobile))]
#[cfg(not(any(target_os = "android", target_os = "ios")))]
let _ = instance_sources;
let mut success_count = 0;
let mut errors = Vec::new();
for uuid in target_ids {
#[cfg(mobile)]
let set_result = match source_map.remove(&uuid) {
Some(sources) => instance_manager.set_tun_fd(&uuid, fd, sources),
None => Err(anyhow::anyhow!("missing tun sources for instance {uuid}")),
};
#[cfg(not(mobile))]
let set_result = instance_manager.set_tun_fd(&uuid, fd);
match set_result {
Ok(()) => {
success_count += 1;
#[cfg(any(target_os = "android", target_os = "ios"))]
{
for instance_id in &target_ids {
if instance_manager.instance(*instance_id).is_none() {
return Err(format!("instance {instance_id} not found"));
}
Err(err) => {
errors.push(format!("{}: {}", uuid, err));
if !source_map.contains_key(instance_id) {
return Err(format!("missing tun sources for instance {instance_id}"));
}
}
if let Some(unexpected_id) = source_map
.keys()
.find(|instance_id| !target_id_set.contains(instance_id))
{
return Err(format!(
"received tun sources for unexpected instance {unexpected_id}"
));
}
}
if success_count == 0 && !errors.is_empty() {
return Err(format!(
"failed to set tun fd for all instances: {}",
errors.join("; ")
));
#[cfg(any(target_os = "android", target_os = "ios"))]
{
let mut attach_error = None;
for instance_id in target_ids.iter().copied() {
let result = attach_mobile_tun_fd(
instance_manager.as_ref(),
instance_id,
fd,
source_map
.remove(&instance_id)
.expect("tun sources were validated before attachment"),
)
.await;
if let Err(error) = result {
attach_error = Some(format!("{instance_id}: {error}"));
break;
}
}
if let Some(attach_error) = attach_error {
let mut rollback_errors = Vec::new();
for instance_id in target_ids {
if let Err(error) = attach_mobile_tun_fd(
instance_manager.as_ref(),
instance_id,
0,
easytier::instance::virtual_nic::MobileTunSources::default(),
)
.await
{
rollback_errors.push(format!("{instance_id}: {error}"));
}
}
let rollback = if rollback_errors.is_empty() {
String::new()
} else {
format!("; rollback failed for {}", rollback_errors.join("; "))
};
return Err(format!(
"failed to set tun fd for the instance group: {attach_error}{rollback}"
));
}
return Ok(());
}
for err in errors {
eprintln!("set_tun_fd skipped instance: {err}");
#[cfg(not(any(target_os = "android", target_os = "ios")))]
{
let mut errors = Vec::new();
for instance_id in target_ids {
if let Err(error) = instance_manager.attach_tun_fd(instance_id, fd) {
errors.push(format!("{instance_id}: {error}"));
}
}
if errors.is_empty() {
Ok(())
} else {
Err(format!(
"failed to set tun fd for the instance group: {}",
errors.join("; ")
))
}
}
Ok(())
}
#[tauri::command]
@@ -475,6 +543,21 @@ fn normalize_normal_mode_rpc_portal(portal: &str) -> Result<(url::Url, url::Url)
Ok((bind_url, connect_url))
}
async fn resolve_rpc_bind_url(url: &url::Url) -> Result<std::net::SocketAddr, String> {
if url.scheme() != "tcp" {
return Err(format!("RPC portal requires tcp URL: {url}"));
}
let host = url
.host_str()
.ok_or_else(|| format!("RPC portal has no host: {url}"))?;
let port = url.port().unwrap_or(11010);
tokio::net::lookup_host((host, port))
.await
.map_err(|error| format!("failed to resolve RPC portal {url}: {error}"))?
.next()
.ok_or_else(|| format!("RPC portal has no resolved address: {url}"))
}
#[tauri::command]
async fn init_rpc_connection(
_app: AppHandle,
@@ -493,11 +576,12 @@ async fn init_rpc_connection(
.map_err(|_| "Failed to acquire lock for rpc server")?;
let mut client_url = url.clone();
let mut local_process_runtime = None;
if is_normal_mode {
let instance_manager = if let Some(im) = instance_manager_guard.take() {
im
} else {
Arc::new(NetworkInstanceManager::new())
Arc::new(native_instance_manager())
};
let portal = url.and_then(|s| {
@@ -525,12 +609,14 @@ async fn init_rpc_connection(
*rpc_server_guard = None;
let tunnel: BoxedTunnelListener = match desired_kind {
RpcServerKind::Ring => Box::new(RingTunnelListener::new(
format!("ring://{}", RPC_RING_UUID.deref()).parse().unwrap(),
)),
RpcServerKind::Tcp => Box::new(TcpTunnelListener::new(
bind_url.clone().expect("tcp rpc must have bind url"),
)),
RpcServerKind::Ring => instance_manager
.process_runtime()
.bind_ring_tunnel(*RPC_RING_UUID.deref())
.map_err(|error| error.to_string())?,
RpcServerKind::Tcp => {
let bind_url = bind_url.as_ref().expect("tcp rpc must have bind url");
Box::new(runtime_rpc_listener(resolve_rpc_bind_url(bind_url).await?))
}
};
let rpc_server = ApiRpcServer::from_tunnel(tunnel, instance_manager.clone())
@@ -545,6 +631,7 @@ async fn init_rpc_connection(
});
}
local_process_runtime = Some(instance_manager.process_runtime());
*instance_manager_guard = Some(instance_manager);
client_url = connect_url.map(|u| u.to_string());
} else {
@@ -553,7 +640,7 @@ async fn init_rpc_connection(
let client_manager = tokio::time::timeout(
std::time::Duration::from_millis(1000),
manager::GUIClientManager::new(client_url),
manager::GUIClientManager::new(client_url, local_process_runtime),
)
.await
.map_err(|_| "connect remote rpc timed out".to_string())?
@@ -565,7 +652,8 @@ async fn init_rpc_connection(
drop(WEB_CLIENT.write().await.take());
if let Some(instance_manager) = instance_manager_guard.take() {
instance_manager
.retain_network_instance(vec![])
.retain_network_instances(&[])
.await
.map_err(|e| e.to_string())?;
drop(instance_manager);
}
@@ -703,19 +791,13 @@ mod manager {
use super::*;
use async_trait::async_trait;
use dashmap::{DashMap, DashSet};
use easytier::common::global_ctx::GlobalCtx;
use easytier::common::stun::MockStunInfoCollector;
use easytier::launcher::NetworkConfig;
use easytier::common::config::{NetworkConfig, NetworkConfigExt};
use easytier::proto::api::logger::{LoggerRpc, LoggerRpcClientFactory, SetLoggerConfigRequest};
use easytier::proto::api::manage::RunNetworkInstanceRequest;
use easytier::proto::common::NatType;
use easytier::proto::rpc_impl::bidirect::BidirectRpcManager;
use easytier::proto::rpc::bidirect::BidirectRpcManager;
use easytier::proto::rpc_types::controller::BaseController;
use easytier::rpc_service::logger::LoggerRpcService;
use easytier::rpc_service::remote_client::PersistentConfig;
use easytier::tunnel::TunnelConnector;
use easytier::tunnel::ring::RingTunnelConnector;
use easytier::web_client::WebClientHooks;
use easytier_core::management::remote_client::PersistentConfig;
pub(super) struct GuiHooks {
pub(super) app: AppHandle,
@@ -980,27 +1062,16 @@ mod manager {
pub(super) rpc_manager: BidirectRpcManager,
}
impl GUIClientManager {
pub async fn new(rpc_url: Option<String>) -> Result<Self, anyhow::Error> {
let global_ctx = Arc::new(GlobalCtx::new(TomlConfigLoader::default()));
global_ctx.replace_stun_info_collector(Box::new(MockStunInfoCollector {
udp_nat_type: NatType::Unknown,
}));
let mut flags = global_ctx.get_flags();
flags.bind_device = false;
global_ctx.set_flags(flags);
pub async fn new(
rpc_url: Option<String>,
local_process_runtime: Option<Arc<CoreProcessRuntime>>,
) -> Result<Self, anyhow::Error> {
let tunnel = if let Some(url) = rpc_url {
let mut connector = easytier::connector::create_connector_by_url(
&url,
&global_ctx,
easytier::tunnel::IpVersion::Both,
)
.await?;
connector.connect().await?
runtime_rpc_dialer(url.parse()?).connect().await?
} else {
let mut connector = RingTunnelConnector::new(
format!("ring://{}", RPC_RING_UUID.deref()).parse().unwrap(),
);
connector.connect().await?
local_process_runtime
.context("local RPC requires a core process runtime")?
.connect_ring_tunnel(*RPC_RING_UUID.deref())?
};
let rpc_manager = BidirectRpcManager::new();
@@ -1111,7 +1182,7 @@ mod manager {
app: &AppHandle,
web_only: bool,
shared_dev_name: Option<&str>,
) -> Result<(), easytier::rpc_service::remote_client::RemoteClientError<anyhow::Error>>
) -> Result<(), easytier_core::management::remote_client::RemoteClientError<anyhow::Error>>
{
for inst_id in self.enabled_incompatible_tun_ids(web_only, shared_dev_name) {
self.handle_update_network_state(app.clone(), inst_id, true)
@@ -1202,11 +1273,8 @@ mod manager {
#[cfg(target_os = "android")]
if let Some(instance_manager) = super::INSTANCE_MANAGER.read().await.as_ref() {
let instance_uuid = *instance_id;
if let Some(instance_ref) = instance_manager
.iter()
.find(|item| *item.key() == instance_uuid)
{
if let Some(mut event_receiver) = instance_ref.value().subscribe_event() {
if let Some(instance) = instance_manager.instance(instance_uuid) {
if let Some(mut event_receiver) = subscribe_native_instance_event(&instance) {
let app_clone = app.clone();
let instance_id_clone = *instance_id;
tokio::spawn(async move {
@@ -1281,7 +1349,7 @@ mod manager {
.set_logger_config(
BaseController::default(),
SetLoggerConfigRequest {
level: LoggerRpcService::string_to_log_level(&level).into(),
level: easytier_core::management::parse_log_level(&level).into(),
},
)
.await?;
@@ -1332,7 +1400,7 @@ mod manager {
inst_id: None,
config: Some(config),
overwrite: false,
source: source.to_runtime_source().to_rpc(),
source: config_source_to_rpc(source.to_runtime_source()),
},
)
.await?;
+23 -17
View File
@@ -1,8 +1,8 @@
import type { NetworkTypes } from 'easytier-frontend-lib'
import type { TunFdInstanceSources } from './backend'
import { addPluginListener } from '@tauri-apps/api/core'
import { Utils } from 'easytier-frontend-lib'
import { get_vpn_status, prepare_vpn, start_vpn, stop_vpn } from 'tauri-plugin-vpnservice-api'
import type { TunFdInstanceSources } from './backend'
type Route = NetworkTypes.Route
@@ -182,9 +182,13 @@ async function onVpnServiceStart(payload: any) {
console.log('vpn service start', JSON.stringify(payload))
curVpnStatus.running = true
if (payload.fd) {
await setTunFd(payload.fd, curVpnStatus.instanceIds, curVpnStatus.instanceSources).catch((e) => {
try {
await setTunFd(payload.fd, curVpnStatus.instanceIds, curVpnStatus.instanceSources)
}
catch (e) {
console.error('set tun fd failed', e)
})
await doStopVpn(true).catch(stopError => console.error('stop vpn after tun attach failure', stopError))
}
}
}
@@ -209,13 +213,9 @@ async function registerVpnServiceListener() {
)
}
function getRoutesForVpn(routes: Route[], node_config: NetworkTypes.NetworkConfig): string[] {
if (!routes) {
return []
}
function getRoutesForVpn(routes: Route[] | undefined, node_config: NetworkTypes.NetworkConfig): string[] {
const ret = []
for (const r of routes) {
for (const r of routes ?? []) {
for (let cidr of r.proxy_cidrs ?? []) {
if (!cidr.includes('/')) {
cidr += '/32'
@@ -224,9 +224,9 @@ function getRoutesForVpn(routes: Route[], node_config: NetworkTypes.NetworkConfi
}
}
node_config.routes?.forEach(r => {
ret.push(r)
})
for (const route of node_config.routes ?? []) {
ret.push(route)
}
if (node_config.enable_magic_dns) {
ret.push('100.100.100.101/32')
@@ -340,7 +340,6 @@ async function applyNetworkInstanceChange(instanceId: string) {
}
return
}
const group = await findRunningTunInstanceGroup(instanceId)
if (!group.length) {
console.warn('vpn service skipped because no running tun instance is available', instanceId)
@@ -371,7 +370,8 @@ async function applyNetworkInstanceChange(instanceId: string) {
continue
}
const virtual_ip = Utils.ipv4ToString(curNetworkInfo?.my_node_info?.virtual_ipv4.address)
const virtualIpv4 = curNetworkInfo.my_node_info?.virtual_ipv4
const virtual_ip = virtualIpv4?.address?.addr ? Utils.ipv4ToString(virtualIpv4.address) : undefined
if (config.dhcp && (!virtual_ip || !virtual_ip.length)) {
console.log('DHCP enabled but no IP yet, will retry in', DHCP_POLLING_INTERVAL, 'ms')
retryInstanceIds.push(instanceId)
@@ -383,7 +383,7 @@ async function applyNetworkInstanceChange(instanceId: string) {
continue
}
let network_length = curNetworkInfo?.my_node_info?.virtual_ipv4.network_length
let network_length = virtualIpv4?.network_length
if (!network_length) {
network_length = 24
}
@@ -415,6 +415,12 @@ async function applyNetworkInstanceChange(instanceId: string) {
if (retryInstanceIds.length) {
scheduleVpnConfigRetry(retryInstanceIds[0])
console.warn('vpn service waiting for complete shared tun instance sources', retryInstanceIds)
if (hasQueuedVpnConfigChange()) {
return
}
await doStopVpn()
return
}
if (!ipv4Addrs.length) {
@@ -484,7 +490,7 @@ async function isNoTunEnabled(instanceId: string | undefined) {
async function findRunningTunInstanceId() {
const instanceIds = await listNetworkInstanceIds()
const runningIds = instanceIds.running_inst_ids.map(Utils.UuidToStr)
const runningIds = (instanceIds.running_inst_ids ?? []).map(Utils.UuidToStr)
console.log('vpn service sync running instances', JSON.stringify(runningIds))
for (const instanceId of runningIds) {
@@ -500,7 +506,7 @@ async function findRunningTunInstanceId() {
async function findRunningTunInstanceGroup(preferredInstanceId?: string) {
const instanceIds = await listNetworkInstanceIds()
const runningIds = instanceIds.running_inst_ids.map(Utils.UuidToStr)
const runningIds = (instanceIds.running_inst_ids ?? []).map(Utils.UuidToStr)
const runningTunInstances = []
for (const instanceId of runningIds) {
+2 -2
View File
@@ -9,7 +9,7 @@ export class GUIRemoteClient implements Api.RemoteClient {
await backend.runNetworkInstance(config, save);
}
async get_network_info(inst_id: string): Promise<NetworkTypes.NetworkInstanceRunningInfo | undefined> {
return backend.collectNetworkInfo(inst_id).then(infos => infos.info.map[inst_id]);
return backend.collectNetworkInfo(inst_id).then(infos => infos.info?.map?.[inst_id]);
}
async list_network_instance_ids(): Promise<Api.ListNetworkInstanceIdResponse> {
return backend.listNetworkInstanceIds();
@@ -44,4 +44,4 @@ export class GUIRemoteClient implements Api.RemoteClient {
return await backend.getNetworkMetas(instance_ids);
}
}
}