Add shared virtual NIC core with per-member routing

Share a named TUN across instances while keeping unnamed devices
dedicated. Track member addresses, routes, and MTU, and reject
overlapping destinations except the common Magic DNS host route.

Dispatch packets between the physical TUN and member tunnels. Keep
IPv4 and ordinary IPv6 source translation for mobile wrong-source
traffic. Update platform configuration and runtime host integration,
with unit and network namespace coverage for routing and lifecycle.

Add the packet dependency and keep core CI matrix jobs independent.
This commit is contained in:
KKRainbow committed 2026-09-29 22:19:38 +08:00
1 parent ff3921ce68
commit cc2d58bb00
25 files changed
+6357 -287

No files matched your search

+1 -1
View File
@@ -82,7 +82,7 @@ jobs:
easytier-web/frontend/dist/*
build:
strategy:
fail-fast: true
fail-fast: false
matrix:
include:
- TARGET: x86_64-unknown-linux-musl
Generated
+34
View File
@@ -2594,6 +2594,7 @@ dependencies = [
"percent-encoding",
"pin-project-lite",
"pnet_datalink",
"pnet_packet",
"prost 0.14.4",
"quanta",
"quinn",
@@ -6911,6 +6912,39 @@ dependencies = [
"winapi",
]
[[package]]
name = "pnet_macros"
version = "0.35.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "13325ac86ee1a80a480b0bc8e3d30c25d133616112bb16e86f712dcf8a71c863"
dependencies = [
"proc-macro2",
"quote",
"regex",
"syn 2.0.119",
]
[[package]]
name = "pnet_macros_support"
version = "0.35.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "eed67a952585d509dd0003049b1fc56b982ac665c8299b124b90ea2bdb3134ab"
dependencies = [
"pnet_base",
]
[[package]]
name = "pnet_packet"
version = "0.35.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4c96ebadfab635fcc23036ba30a7d33a80c39e8461b8bd7dc7bb186acb96560f"
dependencies = [
"glob",
"pnet_base",
"pnet_macros",
"pnet_macros_support",
]
[[package]]
name = "pnet_sys"
version = "0.35.0"
+2 -1
View File
@@ -153,6 +153,7 @@ rand.workspace = true
serde = { workspace = true, features = ["derive"] }
pnet_datalink = { version = "0.35.0", optional = true }
pnet_packet = "0.35.0"
smoltcp = { workspace = true, optional = true, features = [
"std",
"medium-ethernet",
@@ -362,7 +363,7 @@ mimalloc = ["dep:mimalloc"]
aes-gcm = ["easytier-core/aes-gcm"]
openssl-crypto = ["easytier-core/openssl-crypto"]
ring-crypto = ["easytier-core/ring-crypto"]
tun = ["dep:tun", "linux-netlink"]
tun = ["dep:tun", "linux-netlink", "tokio/rt-multi-thread"]
linux-netlink = ["dep:netlink-sys"]
proxy-cidr-monitor = ["easytier-core/proxy-cidr-monitor"]
websocket = [
+63 -13
View File
@@ -1,10 +1,38 @@
use std::net::Ipv4Addr;
use std::{collections::BTreeSet, net::Ipv4Addr};
use super::{Error, IfConfiguerTrait, cidr_to_subnet_mask, run_shell_cmd};
use async_trait::async_trait;
use cidr::{Ipv4Inet, Ipv6Inet};
use tokio::sync::Mutex;
#[derive(Default)]
pub struct MacIfConfiger {
configured_ipv4: Mutex<BTreeSet<Ipv4Inet>>,
}
impl MacIfConfiger {
fn build_add_ipv4_cmd(name: &str, addr: Ipv4Inet, has_configured_ipv4: bool) -> String {
let address = addr.address();
if has_configured_ipv4 {
format!(
"ifconfig {} alias {:?} {:?} netmask {}",
name,
address,
address,
cidr_to_subnet_mask(addr.network_length())
)
} else {
format!(
"ifconfig {} {:?}/{:?} {:?} up",
name,
address,
addr.network_length(),
address,
)
}
}
}
pub struct MacIfConfiger {}
#[async_trait]
impl IfConfiguerTrait for MacIfConfiger {
async fn add_ipv4_route(
@@ -14,12 +42,28 @@ impl IfConfiguerTrait for MacIfConfiger {
cidr_prefix: u8,
cost: Option<i32>,
) -> Result<(), Error> {
self.add_ipv4_route_with_source_hint(name, address, cidr_prefix, cost, None)
.await
}
async fn add_ipv4_route_with_source_hint(
&self,
name: &str,
address: Ipv4Addr,
cidr_prefix: u8,
cost: Option<i32>,
source_hint: Option<Ipv4Addr>,
) -> Result<(), Error> {
let source_hint = source_hint
.map(|source| format!(" -ifa {}", source))
.unwrap_or_default();
run_shell_cmd(
format!(
"route -n add {} -netmask {} -interface {} -hopcount {}",
"route -n add {} -netmask {} -interface {}{} -hopcount {}",
address,
cidr_to_subnet_mask(cidr_prefix),
name,
source_hint,
cost.unwrap_or(7)
)
.as_str(),
@@ -51,14 +95,15 @@ impl IfConfiguerTrait for MacIfConfiger {
address: Ipv4Addr,
cidr_prefix: u8,
) -> Result<(), Error> {
run_shell_cmd(
format!(
"ifconfig {} {:?}/{:?} {:?} up",
name, address, cidr_prefix, address,
)
.as_str(),
)
.await
let addr = Ipv4Inet::new(address, cidr_prefix).map_err(|err| {
anyhow::anyhow!("invalid IPv4 address {address}/{cidr_prefix}: {err:?}")
})?;
let mut configured_ipv4 = self.configured_ipv4.lock().await;
let cmd = Self::build_add_ipv4_cmd(name, addr, !configured_ipv4.is_empty());
run_shell_cmd(cmd.as_str()).await?;
configured_ipv4.insert(addr);
Ok(())
}
async fn set_link_status(&self, name: &str, up: bool) -> Result<(), Error> {
@@ -67,11 +112,16 @@ impl IfConfiguerTrait for MacIfConfiger {
}
async fn remove_ip(&self, name: &str, ip: Option<Ipv4Inet>) -> Result<(), Error> {
let mut configured_ipv4 = self.configured_ipv4.lock().await;
if let Some(ip) = ip {
run_shell_cmd(format!("ifconfig {} inet {} delete", name, ip.address()).as_str()).await
run_shell_cmd(format!("ifconfig {} inet {} delete", name, ip.address()).as_str())
.await?;
configured_ipv4.remove(&ip);
} else {
run_shell_cmd(format!("ifconfig {} inet delete", name).as_str()).await
run_shell_cmd(format!("ifconfig {} inet delete", name).as_str()).await?;
configured_ipv4.clear();
}
Ok(())
}
async fn set_mtu(&self, name: &str, mtu: u32) -> Result<(), Error> {
+21
View File
@@ -40,6 +40,16 @@ pub trait IfConfiguerTrait: Send + Sync {
) -> Result<(), Error> {
Ok(())
}
async fn add_ipv4_route_with_source_hint(
&self,
name: &str,
address: Ipv4Addr,
cidr_prefix: u8,
cost: Option<i32>,
_source_hint: Option<Ipv4Addr>,
) -> Result<(), Error> {
self.add_ipv4_route(name, address, cidr_prefix, cost).await
}
async fn remove_ipv4_route(
&self,
_name: &str,
@@ -48,6 +58,16 @@ pub trait IfConfiguerTrait: Send + Sync {
) -> Result<(), Error> {
Ok(())
}
async fn remove_ipv4_route_with_cost_and_source_hint(
&self,
name: &str,
address: Ipv4Addr,
cidr_prefix: u8,
_cost: Option<i32>,
_source_hint: Option<Ipv4Addr>,
) -> Result<(), Error> {
self.remove_ipv4_route(name, address, cidr_prefix).await
}
async fn add_ipv4_ip(
&self,
_name: &str,
@@ -157,6 +177,7 @@ async fn run_shell_cmd(cmd: &str) -> Result<(), Error> {
Ok(())
}
#[derive(Default)]
pub struct DummyIfConfiger {}
#[async_trait]
impl IfConfiguerTrait for DummyIfConfiger {}
+106 -9
View File
@@ -133,6 +133,7 @@ fn dump_netlink_messages<T: NetlinkDecode>(
receive_netlink_dump(builder)
}
#[derive(Default)]
pub struct NetlinkIfConfiger {}
impl NetlinkIfConfiger {
@@ -335,6 +336,27 @@ impl NetlinkIfConfiger {
})
.collect())
}
fn ipv4_route_message(
ifindex: u32,
address: Ipv4Addr,
cidr_prefix: u8,
cost: Option<i32>,
source_hint: Option<Ipv4Addr>,
) -> RouteMessage {
let mut builder = RouteMessageBuilder::new(libc::AF_INET as u8)
.destination(IpAddr::V4(address), cidr_prefix)
.oif(ifindex)
.priority(cost.unwrap_or(65535) as u32)
.table(libc::RT_TABLE_MAIN.into())
.static_protocol()
.universe_scope()
.route_type(RouteType::Unicast);
if let Some(source_hint) = source_hint {
builder = builder.preferred_source(IpAddr::V4(source_hint));
}
builder.build()
}
}
#[async_trait]
@@ -346,15 +368,25 @@ impl IfConfiguerTrait for NetlinkIfConfiger {
cidr_prefix: u8,
cost: Option<i32>,
) -> Result<(), Error> {
let message = RouteMessageBuilder::new(libc::AF_INET as u8)
.destination(IpAddr::V4(address), cidr_prefix)
.oif(Self::get_interface_index(name)?)
.priority(cost.unwrap_or(65535) as u32)
.table(libc::RT_TABLE_MAIN.into())
.static_protocol()
.universe_scope()
.route_type(RouteType::Unicast)
.build();
self.add_ipv4_route_with_source_hint(name, address, cidr_prefix, cost, None)
.await
}
async fn add_ipv4_route_with_source_hint(
&self,
name: &str,
address: Ipv4Addr,
cidr_prefix: u8,
cost: Option<i32>,
source_hint: Option<Ipv4Addr>,
) -> Result<(), Error> {
let message = NetlinkIfConfiger::ipv4_route_message(
NetlinkIfConfiger::get_interface_index(name)?,
address,
cidr_prefix,
cost,
source_hint,
);
let request = message_request(
RTM_NEWROUTE,
NLM_F_ACK | NLM_F_CREATE | NLM_F_EXCL | NLM_F_REQUEST,
@@ -390,6 +422,28 @@ impl IfConfiguerTrait for NetlinkIfConfiger {
Ok(())
}
async fn remove_ipv4_route_with_cost_and_source_hint(
&self,
name: &str,
address: Ipv4Addr,
cidr_prefix: u8,
cost: Option<i32>,
source_hint: Option<Ipv4Addr>,
) -> Result<(), Error> {
let message = Self::ipv4_route_message(
Self::get_interface_index(name)?,
address,
cidr_prefix,
cost,
source_hint,
);
let request = message_request(RTM_DELROUTE, NLM_F_ACK | NLM_F_REQUEST, &message)?;
match send_netlink_req_and_wait_ack(request) {
Err(Error::IOError(err)) if err.raw_os_error() == Some(libc::ESRCH) => Ok(()),
result => result,
}
}
async fn add_ipv4_ip(
&self,
name: &str,
@@ -696,6 +750,49 @@ mod tests {
assert!(!routes.contains(&IpAddr::V4("10.5.5.0".parse().unwrap())));
}
#[serial_test::serial]
#[tokio::test]
async fn remove_ipv4_route_with_source_hint_keeps_other_metric() {
let iface = test_iface_name("rm");
let _link = ScopedDummyLink::new(&iface);
let ifcfg = NetlinkIfConfiger {};
let address = "10.231.1.1".parse().unwrap();
let destination = "10.99.0.0".parse().unwrap();
ifcfg.add_ipv4_ip(&iface, address, 24).await.unwrap();
for cost in [123, 124] {
ifcfg
.add_ipv4_route_with_source_hint(&iface, destination, 24, Some(cost), Some(address))
.await
.unwrap();
}
ifcfg
.remove_ipv4_route_with_cost_and_source_hint(
&iface,
destination,
24,
Some(123),
Some(address),
)
.await
.unwrap();
let routes = run_cmd(&format!("ip -4 route show 10.99.0.0/24 dev {iface}"));
assert!(!routes.contains("metric 123"));
assert!(routes.contains("metric 124"));
ifcfg
.remove_ipv4_route_with_cost_and_source_hint(
&iface,
destination,
24,
Some(123),
Some(address),
)
.await
.unwrap();
}
#[serial_test::serial]
#[tokio::test]
async fn ipv6_addr_readback_test() {
+24
View File
@@ -37,6 +37,7 @@ const RTA_DST: u16 = 1;
const RTA_SRC: u16 = 2;
const RTA_OIF: u16 = 4;
const RTA_PRIORITY: u16 = 6;
const RTA_PREFSRC: u16 = 7;
const RTA_TABLE: u16 = 15;
const NDA_DST: u16 = 1;
@@ -488,6 +489,13 @@ impl RouteMessageBuilder {
self
}
pub(crate) fn preferred_source(mut self, address: IpAddr) -> Self {
self.message
.attributes
.push(Attribute::new(RTA_PREFSRC, ip_bytes(address)));
self
}
pub(crate) fn table(mut self, table: u32) -> Self {
if let Ok(table) = u8::try_from(table) {
self.message.table = table;
@@ -632,6 +640,22 @@ mod tests {
assert_eq!(encode(&decoded), bytes);
}
#[test]
fn route_builder_encodes_preferred_source() {
let source = "10.231.1.1".parse().unwrap();
let message = RouteMessageBuilder::new(libc::AF_INET as u8)
.destination("10.99.0.0".parse().unwrap(), 24)
.preferred_source(source)
.oif(7)
.table(libc::RT_TABLE_MAIN.into())
.build();
let bytes = encode(&message);
let attributes = parse_attributes(&bytes[12..]).unwrap();
assert!(attributes.iter().any(|attribute| {
attribute.kind == RTA_PREFSRC && attribute.value == ip_bytes(source)
}));
}
#[test]
fn route_parser_reads_ipv6_source_prefix() {
let mut bytes = vec![
+1
View File
@@ -22,6 +22,7 @@ use winreg::{
};
use super::{Error, IfConfiguerTrait};
#[derive(Default)]
pub struct WindowsIfConfiger {}
fn format_win_error(error: u32) -> String {
+8 -1
View File
@@ -32,6 +32,8 @@ use crate::{
};
use super::host::{NativeInstanceHost, native_instance_host};
#[cfg(feature = "tun")]
use super::shared_virtual_nic::ArcSharedVirtualNicRegistry;
#[cfg(feature = "kcp")]
use crate::gateway::kcp_proxy::KcpProxyService;
#[cfg(feature = "quic")]
@@ -45,6 +47,7 @@ pub(crate) type NativeCoreInstance = CoreInstance<NativeInstanceHost>;
pub(crate) fn compose_native_core_instance(
toml_config: TomlConfig,
process_runtime: Arc<CoreProcessRuntime>,
#[cfg(feature = "tun")] shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
compact_runtime: bool,
) -> anyhow::Result<Arc<NativeCoreInstance>> {
let host_config = if compact_runtime {
@@ -58,7 +61,11 @@ pub(crate) fn compose_native_core_instance(
&normalized,
&host_config,
));
let runtime_host = NativeInstanceRuntimeHost::new(global_ctx.clone());
let runtime_host = NativeInstanceRuntimeHost::new(
global_ctx.clone(),
#[cfg(feature = "tun")]
shared_virtual_nic_registry,
);
let mut adapters = runtime_core_host_adapters_with_packet_egress_and_config(
global_ctx.clone(),
process_runtime,
+111 -10
View File
@@ -5,7 +5,15 @@ use std::{net::Ipv4Addr, sync::Arc, time::Duration};
use easytier_core::instance::CorePacketPlane;
use crate::common::global_ctx::ArcGlobalCtx;
use crate::{
common::{
error::Error as EtError,
global_ctx::ArcGlobalCtx,
ifcfg::{IfConfiger, IfConfiguerTrait},
netns::NetNS,
},
instance::virtual_nic::NicBackend,
};
use super::{client_instance::MagicDnsClientInstance, server_instance::MagicDnsServerInstance};
@@ -17,6 +25,52 @@ pub struct DnsRunner {
tun_dev: Option<String>,
tun_inet: Ipv4Inet,
fake_ip: Ipv4Addr,
shared_route_backend: Option<NicBackend>,
}
#[derive(Clone)]
struct MagicDnsFakeIpRouteClaim {
tun_dev: Option<String>,
net_ns: NetNS,
fake_ip: Ipv4Addr,
route_backend: NicBackend,
}
impl MagicDnsFakeIpRouteClaim {
async fn add(&self) -> anyhow::Result<()> {
let cost = if cfg!(target_os = "windows") {
Some(4)
} else {
None
};
match self
.route_backend
.add_route_with_cost(self.fake_ip, 32, cost)
.await
{
Err(EtError::IOError(err))
if err.kind() == std::io::ErrorKind::AlreadyExists && self.tun_dev.is_some() =>
{
let ifcfg = IfConfiger::default();
let _guard = self.net_ns.guard();
ifcfg
.remove_ipv4_route(self.tun_dev.as_deref().unwrap(), self.fake_ip, 32)
.await?;
self.route_backend
.add_route_with_cost(self.fake_ip, 32, cost)
.await?;
Ok(())
}
result => result.map_err(Into::into),
}
}
async fn remove(&self) {
if let Err(err) = self.route_backend.remove_route(self.fake_ip, 32).await {
tracing::warn!(?err, fake_ip = ?self.fake_ip, "remove magic dns route failed");
}
}
}
impl DnsRunner {
@@ -35,9 +89,15 @@ impl DnsRunner {
tun_dev,
tun_inet,
fake_ip,
shared_route_backend: None,
}
}
pub(crate) fn with_shared_route_backend(mut self, route_backend: Option<NicBackend>) -> Self {
self.shared_route_backend = route_backend;
self
}
async fn clean_env(&mut self) {
if let Some(server) = self.server.take() {
server.clean_env().await;
@@ -45,17 +105,53 @@ impl DnsRunner {
self.client.take();
}
fn should_manage_fake_ip_route_externally(&self) -> bool {
self.shared_route_backend.is_some() && !self.tun_inet.contains(&self.fake_ip)
}
fn fake_ip_route_claim(&self) -> Option<MagicDnsFakeIpRouteClaim> {
if !self.should_manage_fake_ip_route_externally() {
return None;
}
Some(MagicDnsFakeIpRouteClaim {
tun_dev: self.tun_dev.clone(),
net_ns: self.global_ctx.net_ns.clone(),
fake_ip: self.fake_ip,
route_backend: self.shared_route_backend.clone()?,
})
}
async fn run_once(&mut self) -> anyhow::Result<()> {
if let Some(claim) = self.fake_ip_route_claim() {
claim
.add()
.await
.map_err(|err| anyhow::anyhow!("failed to add magic dns fake-ip route: {err}"))?;
}
// try server first
match MagicDnsServerInstance::new(
self.packet_plane.clone(),
self.global_ctx.clone(),
self.tun_dev.clone(),
self.tun_inet,
self.fake_ip,
)
.await
{
let server_result = if self.should_manage_fake_ip_route_externally() {
MagicDnsServerInstance::new_with_external_fake_ip_route(
self.packet_plane.clone(),
self.global_ctx.clone(),
self.tun_dev.clone(),
self.tun_inet,
self.fake_ip,
)
.await
} else {
MagicDnsServerInstance::new(
self.packet_plane.clone(),
self.global_ctx.clone(),
self.tun_dev.clone(),
self.tun_inet,
self.fake_ip,
)
.await
};
match server_result {
Ok(server) => {
self.server = Some(server);
tracing::info!("DnsRunner::run_once: server started");
@@ -74,11 +170,16 @@ impl DnsRunner {
}
pub async fn run(&mut self, canel_token: CancellationToken) {
let fake_ip_route_claim = self.fake_ip_route_claim();
loop {
tracing::info!("DnsRunner::run: start");
tokio::select! {
_ = canel_token.cancelled() => {
self.clean_env().await;
if let Some(claim) = &fake_ip_route_claim {
claim.remove().await;
}
tracing::info!("DnsRunner::run: cancelled");
return;
}
@@ -14,8 +14,10 @@ use super::{
};
use crate::{
common::{
error::Error as EtError,
global_ctx::ArcGlobalCtx,
ifcfg::{IfConfiger, IfConfiguerTrait},
netns::NetNS,
},
instance::dns_server::{
config::{Record, RecordBuilder, RecordType},
@@ -51,7 +53,9 @@ use std::{collections::BTreeMap, io, net::Ipv4Addr, str::FromStr, sync::Arc, tim
pub(super) struct MagicDnsServerInstanceData {
dns_server: Server,
tun_dev: Option<String>,
net_ns: NetNS,
fake_ip: Ipv4Addr,
manage_fake_ip_route: bool,
route_store: MagicDnsRecordStore,
record_apply: tokio::sync::Mutex<()>,
@@ -356,12 +360,66 @@ fn get_system_config(
}
impl MagicDnsServerInstance {
async fn add_fake_ip_route(
tun_dev_name: &str,
fake_ip: Ipv4Addr,
net_ns: &NetNS,
cost: Option<i32>,
) -> Result<(), anyhow::Error> {
let ifcfg = IfConfiger::default();
let _guard = net_ns.guard();
match ifcfg.add_ipv4_route(tun_dev_name, fake_ip, 32, cost).await {
Err(EtError::IOError(err)) if err.kind() == io::ErrorKind::AlreadyExists => {
ifcfg.remove_ipv4_route(tun_dev_name, fake_ip, 32).await?;
ifcfg
.add_ipv4_route(tun_dev_name, fake_ip, 32, cost)
.await?;
Ok(())
}
ret => ret.map_err(Into::into),
}
}
async fn remove_fake_ip_route(tun_dev_name: &str, fake_ip: Ipv4Addr, net_ns: &NetNS) {
let ifcfg = IfConfiger::default();
let _guard = net_ns.guard();
if let Err(err) = ifcfg.remove_ipv4_route(tun_dev_name, fake_ip, 32).await {
tracing::warn!(
?err,
?tun_dev_name,
?fake_ip,
"remove magic dns route failed"
);
}
}
pub(crate) async fn new(
packet_plane: Arc<CorePacketPlane>,
global_ctx: ArcGlobalCtx,
tun_dev: Option<String>,
tun_inet: Ipv4Inet,
fake_ip: Ipv4Addr,
) -> Result<Self, anyhow::Error> {
Self::new_inner(packet_plane, global_ctx, tun_dev, tun_inet, fake_ip, true).await
}
pub(crate) async fn new_with_external_fake_ip_route(
packet_plane: Arc<CorePacketPlane>,
global_ctx: ArcGlobalCtx,
tun_dev: Option<String>,
tun_inet: Ipv4Inet,
fake_ip: Ipv4Addr,
) -> Result<Self, anyhow::Error> {
Self::new_inner(packet_plane, global_ctx, tun_dev, tun_inet, fake_ip, false).await
}
async fn new_inner(
packet_plane: Arc<CorePacketPlane>,
global_ctx: ArcGlobalCtx,
tun_dev: Option<String>,
tun_inet: Ipv4Inet,
fake_ip: Ipv4Addr,
manage_fake_ip_route: bool,
) -> Result<Self, anyhow::Error> {
let tcp_listener = runtime_rpc_listener(MAGIC_DNS_INSTANCE_SOCKET_ADDR.parse()?);
let mut rpc_server = StandAloneServer::new(tcp_listener);
@@ -374,7 +432,8 @@ impl MagicDnsServerInstance {
let mut dns_server = Server::new(dns_config);
dns_server.run().await?;
if !tun_inet.contains(&fake_ip)
if manage_fake_ip_route
&& !tun_inet.contains(&fake_ip)
&& let Some(tun_dev_name) = &tun_dev
{
let cost = if cfg!(target_os = "windows") {
@@ -382,16 +441,15 @@ impl MagicDnsServerInstance {
} else {
None
};
let ifcfg = IfConfiger {};
ifcfg
.add_ipv4_route(tun_dev_name, fake_ip, 32, cost)
.await?;
Self::add_fake_ip_route(tun_dev_name, fake_ip, &global_ctx.net_ns, cost).await?;
}
let data = Arc::new(MagicDnsServerInstanceData {
dns_server,
tun_dev: tun_dev.clone(),
net_ns: global_ctx.net_ns.clone(),
fake_ip,
manage_fake_ip_route,
route_store: MagicDnsRecordStore::default(),
record_apply: tokio::sync::Mutex::new(()),
system_config: get_system_config(tun_dev.as_deref())?,
@@ -436,14 +494,13 @@ impl MagicDnsServerInstance {
if let Err(e) = ret {
tracing::error!("Failed to close system config: {:?}", e);
}
if !self.tun_inet.contains(&self.data.fake_ip)
&& let Some(tun_dev_name) = &self.data.tun_dev
{
let ifcfg = IfConfiger {};
let _ = ifcfg
.remove_ipv4_route(tun_dev_name, self.data.fake_ip, 32)
.await;
}
}
if self.data.manage_fake_ip_route
&& !self.tun_inet.contains(&self.data.fake_ip)
&& let Some(tun_dev_name) = &self.data.tun_dev
{
Self::remove_fake_ip_route(tun_dev_name, self.data.fake_ip, &self.data.net_ns).await;
}
self.packet_filter.close().await;
+66
View File
@@ -1,5 +1,10 @@
use std::sync::Arc;
#[cfg(feature = "tun")]
use tokio::sync::Mutex;
#[cfg(all(feature = "management-rpc", feature = "tun", mobile))]
use easytier_core::instance::CoreInstanceState;
#[cfg(any(feature = "management-rpc", test))]
use easytier_core::instance::manager::InstanceManager;
#[cfg(feature = "management-rpc")]
@@ -12,6 +17,8 @@ use easytier_core::{
use crate::common::global_ctx::EventBusSubscriber;
#[cfg(feature = "tun")]
use super::shared_virtual_nic::{ArcSharedVirtualNicRegistry, SharedVirtualNicRegistry};
use super::{
composition::compose_native_core_instance, host::NativeInstanceHost,
runtime_host::NativeInstanceRuntimeHost,
@@ -51,6 +58,58 @@ pub fn subscribe_native_instance_event(
.map(NativeInstanceRuntimeHost::subscribe_event)
}
#[cfg(all(feature = "management-rpc", feature = "tun", mobile))]
pub async fn attach_mobile_tun_fd(manager: &NativeInstanceManager, fd: i32) -> anyhow::Result<()> {
let instances = manager
.instances()
.into_iter()
.filter(|instance| instance.state() == CoreInstanceState::Running)
.filter(|instance| {
instance
.runtime_host::<NativeInstanceRuntimeHost>()
.is_some_and(NativeInstanceRuntimeHost::tun_enabled)
})
.collect::<Vec<_>>();
if fd > 0 && instances.is_empty() {
anyhow::bail!("no running TUN-enabled instance is available for fd attachment");
}
let mut errors = Vec::new();
// The first member opens the TUN; the remaining members join its dispatcher.
for (index, instance) in instances.iter().enumerate() {
if let Err(error) = attach_mobile_tun_fd_to_instance(instance, fd, index == 0).await {
errors.push(format!("{}: {error}", instance.instance_id()));
if fd > 0 {
break;
}
}
}
if errors.is_empty() {
return Ok(());
}
if fd > 0 {
for instance in &instances {
if let Err(error) = attach_mobile_tun_fd_to_instance(instance, 0, false).await {
errors.push(format!("cleanup {}: {error}", instance.instance_id()));
}
}
}
anyhow::bail!("failed to attach mobile TUN fd: {}", errors.join("; "))
}
#[cfg(all(feature = "management-rpc", feature = "tun", mobile))]
async fn attach_mobile_tun_fd_to_instance(
instance: &NativeCoreInstance,
fd: i32,
replace_tun_fd: bool,
) -> anyhow::Result<()> {
let runtime = instance
.runtime_host::<NativeInstanceRuntimeHost>()
.ok_or_else(|| anyhow::anyhow!("native runtime host is unavailable"))?;
runtime.attach_mobile_tun_fd(fd, replace_tun_fd).await
}
#[cfg(feature = "management-rpc")]
pub fn native_instance_manager_with_runtime(
runtime_handle: tokio::runtime::Handle,
@@ -97,6 +156,8 @@ fn native_instance_manager_with_optional_runtime(
/// Native construction Adapter for the canonical core InstanceManager.
pub struct NativeInstanceFactory {
process_runtime: Arc<CoreProcessRuntime>,
#[cfg(feature = "tun")]
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
runtime_handle: Option<tokio::runtime::Handle>,
compact_runtime: bool,
#[cfg(feature = "logging")]
@@ -107,6 +168,8 @@ impl NativeInstanceFactory {
pub fn new(process_runtime: Arc<CoreProcessRuntime>) -> Self {
Self {
process_runtime,
#[cfg(feature = "tun")]
shared_virtual_nic_registry: Arc::new(Mutex::new(SharedVirtualNicRegistry::new())),
runtime_handle: None,
compact_runtime: false,
#[cfg(feature = "logging")]
@@ -126,6 +189,7 @@ impl NativeInstanceFactory {
self
}
#[cfg(feature = "management-rpc")]
fn with_compact_runtime(mut self) -> Self {
self.compact_runtime = true;
self
@@ -149,6 +213,8 @@ impl InstanceFactory for NativeInstanceFactory {
let instance = compose_native_core_instance(
config,
self.process_runtime.clone(),
#[cfg(feature = "tun")]
self.shared_virtual_nic_registry.clone(),
self.compact_runtime,
)?;
#[cfg(feature = "logging")]
+3
View File
@@ -18,6 +18,9 @@ pub(crate) mod listeners;
#[cfg(feature = "public-ipv6-provider")]
pub(crate) mod public_ipv6_provider;
#[cfg(feature = "tun")]
pub mod shared_virtual_nic;
#[cfg(feature = "tun")]
pub mod virtual_nic;
+42 -7
View File
@@ -1,13 +1,16 @@
use std::sync::Arc;
#[cfg(feature = "web-client")]
use easytier_core::config::runtime::CoreInstanceRuntimeConfig;
use easytier_core::{
config::runtime::CoreInstanceRuntimeConfig, gateway::dhcp::DhcpIpv4Host,
host::packet::HostPacketReceiver, instance::CorePacketPlane,
gateway::dhcp::DhcpIpv4Host, host::packet::HostPacketReceiver, instance::CorePacketPlane,
};
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
use crate::common::global_ctx::ArcGlobalCtx;
#[cfg(feature = "tun")]
use crate::instance::shared_virtual_nic::ArcSharedVirtualNicRegistry;
mod event_journal;
mod implementation;
@@ -39,9 +42,17 @@ pub(crate) struct NativeInstanceRuntimeHost {
}
impl NativeInstanceRuntimeHost {
pub(crate) fn new(global_ctx: ArcGlobalCtx) -> Arc<Self> {
pub(crate) fn new(
global_ctx: ArcGlobalCtx,
#[cfg(feature = "tun")] shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
) -> Arc<Self> {
let cancel = CancellationToken::new();
let tun = NativeTunRuntime::new(global_ctx.clone(), cancel.clone());
let tun = NativeTunRuntime::new(
global_ctx.clone(),
cancel.clone(),
#[cfg(feature = "tun")]
shared_virtual_nic_registry,
);
let event_journal = EventJournal::new(&global_ctx);
Arc::new(Self {
global_ctx,
@@ -117,10 +128,24 @@ impl NativeInstanceRuntimeHost {
self.global_ctx.subscribe()
}
#[cfg(all(feature = "tun", mobile))]
pub(crate) fn tun_enabled(&self) -> bool {
!self.global_ctx.get_flags().no_tun
}
fn attach_runtime_tun_fd(&self, fd: i32) -> anyhow::Result<()> {
self.tun.attach_fd(fd)
}
#[cfg(all(feature = "tun", mobile))]
pub(crate) async fn attach_mobile_tun_fd(
&self,
fd: i32,
replace_tun_fd: bool,
) -> anyhow::Result<()> {
self.tun.attach_mobile_fd(fd, replace_tun_fd).await
}
fn install_packet_receiver(&self, receiver: HostPacketReceiver) -> anyhow::Result<()> {
self.tun.install_packet_receiver(receiver)
}
@@ -134,6 +159,16 @@ mod tests {
global_ctx::{GlobalCtx, GlobalCtxEvent},
};
fn runtime_host(global_ctx: ArcGlobalCtx) -> Arc<NativeInstanceRuntimeHost> {
NativeInstanceRuntimeHost::new(
global_ctx,
#[cfg(feature = "tun")]
Arc::new(tokio::sync::Mutex::new(
crate::instance::shared_virtual_nic::SharedVirtualNicRegistry::new(),
)),
)
}
#[cfg(feature = "web-client")]
fn runtime_config(config: &TomlConfig) -> CoreInstanceRuntimeConfig {
let normalized = easytier_core::instance::CoreInstanceConfig::from_toml(config).unwrap();
@@ -146,7 +181,7 @@ mod tests {
#[test]
fn runtime_host_owns_event_subscription_context() {
let global_ctx = Arc::new(GlobalCtx::new(TomlConfig::default()));
let runtime_host = NativeInstanceRuntimeHost::new(global_ctx.clone());
let runtime_host = runtime_host(global_ctx.clone());
let mut events = runtime_host.subscribe_event();
global_ctx.issue_event(GlobalCtxEvent::CredentialChanged);
@@ -167,7 +202,7 @@ mod tests {
config.set_ipv4(Some("10.20.0.1/24".parse().unwrap()));
config.set_ipv6(Some("fd00::1/64".parse().unwrap()));
let global_ctx = Arc::new(GlobalCtx::new(config.clone()));
let runtime_host = NativeInstanceRuntimeHost::new(global_ctx.clone());
let runtime_host = runtime_host(global_ctx.clone());
assert_eq!(global_ctx.get_hostname(), "before");
assert_eq!(global_ctx.get_ipv4(), Some("10.20.0.1/24".parse().unwrap()));
@@ -209,7 +244,7 @@ mod tests {
let config = TomlConfig::default();
config.set_dhcp(true);
let global_ctx = Arc::new(GlobalCtx::new(config.clone()));
let runtime_host = NativeInstanceRuntimeHost::new(global_ctx.clone());
let runtime_host = runtime_host(global_ctx.clone());
let lease = "10.20.0.7/24".parse().unwrap();
global_ctx.set_ipv4(Some(lease));
@@ -3,9 +3,9 @@ use easytier_core::instance::CorePacketPlane;
#[cfg(feature = "magic-dns")]
use tokio_util::{sync::CancellationToken, task::AbortOnDropHandle};
use crate::common::global_ctx::ArcGlobalCtx;
#[cfg(feature = "magic-dns")]
use crate::instance::dns_server::{MAGIC_DNS_FAKE_IP, runner::DnsRunner};
use crate::{common::global_ctx::ArcGlobalCtx, instance::virtual_nic::NicBackend};
#[derive(Default)]
pub(super) struct MagicDnsRuntime {
@@ -26,6 +26,7 @@ impl MagicDnsRuntime {
packet_plane: std::sync::Arc<CorePacketPlane>,
tun_dev: Option<String>,
tun_ip: Ipv4Inet,
shared_route_backend: Option<NicBackend>,
) -> Self {
let active = global_ctx.get_flags().accept_dns.then(|| {
let mut runner = DnsRunner::new(
@@ -34,7 +35,8 @@ impl MagicDnsRuntime {
tun_dev,
tun_ip,
MAGIC_DNS_FAKE_IP.parse().unwrap(),
);
)
.with_shared_route_backend(shared_route_backend);
let cancel = CancellationToken::new();
let task_cancel = cancel.clone();
let task = tokio::spawn(async move {
@@ -54,6 +56,7 @@ impl MagicDnsRuntime {
_packet_plane: std::sync::Arc<CorePacketPlane>,
_tun_dev: Option<String>,
_tun_ip: Ipv4Inet,
_shared_route_backend: Option<NicBackend>,
) -> Self {
Self::default()
}
@@ -3,11 +3,43 @@ use std::{
sync::{Arc, OnceLock},
};
use easytier_core::host::packet::HostPacketReceiver;
use easytier_core::{host::packet::HostPacketReceiver, instance::CorePacketPlane};
use tokio::{sync::Mutex, task::JoinSet};
use super::MagicDnsRuntime;
use crate::instance::virtual_nic::NicCtx;
use crate::{
common::{error::Error, global_ctx::ArcGlobalCtx},
instance::{shared_virtual_nic::ArcSharedVirtualNicRegistry, virtual_nic::NicCtx},
};
pub(super) async fn create_nic_ctx(
global_ctx: ArcGlobalCtx,
packet_plane: Arc<CorePacketPlane>,
receiver: Arc<Mutex<HostPacketReceiver>>,
close_notifier: Arc<tokio::sync::Notify>,
registry: ArcSharedVirtualNicRegistry,
) -> Result<NicCtx, Error> {
#[cfg(not(mobile))]
if global_ctx.get_flags().dev_name.is_empty() {
return Ok(NicCtx::new(
global_ctx,
packet_plane,
receiver,
close_notifier,
));
}
let member_id = global_ctx.get_id();
NicCtx::new_shared(
global_ctx,
packet_plane,
receiver,
close_notifier,
registry,
member_id,
)
.await
}
struct NicCtxContainer {
_nic_ctx: Option<Box<dyn Any + Send>>,
@@ -14,28 +14,37 @@ use tokio::{
};
use tokio_util::sync::CancellationToken;
use super::{MagicDnsRuntime, tun_common::TunNicState};
use super::{
MagicDnsRuntime,
tun_common::{TunNicState, create_nic_ctx},
};
use crate::{
common::{
error::Error,
global_ctx::{ArcGlobalCtx, GlobalCtxEvent},
},
instance::virtual_nic::NicCtx,
instance::shared_virtual_nic::ArcSharedVirtualNicRegistry,
};
pub(super) struct NativeTunRuntime {
global_ctx: ArcGlobalCtx,
cancel: CancellationToken,
nic: TunNicState,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
static_ip_task: Mutex<Option<JoinHandle<()>>>,
}
impl NativeTunRuntime {
pub(super) fn new(global_ctx: ArcGlobalCtx, cancel: CancellationToken) -> Self {
pub(super) fn new(
global_ctx: ArcGlobalCtx,
cancel: CancellationToken,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
) -> Self {
Self {
global_ctx,
cancel,
nic: TunNicState::empty(),
shared_virtual_nic_registry,
static_ip_task: Mutex::new(None),
}
}
@@ -67,6 +76,7 @@ impl NativeTunRuntime {
let cancel = self.cancel.clone();
let global_ctx = self.global_ctx.clone();
let receiver = self.nic.receiver();
let shared_virtual_nic_registry = self.shared_virtual_nic_registry.clone();
let (output, first_round) = oneshot::channel();
let task = tokio::spawn(async move {
let mut output = Some(output);
@@ -76,12 +86,29 @@ impl NativeTunRuntime {
return;
}
let closed = Arc::new(Notify::new());
let mut nic = NicCtx::new(
let mut nic = match create_nic_ctx(
global_ctx.clone(),
packet_plane.clone(),
receiver.clone(),
closed.clone(),
);
shared_virtual_nic_registry.clone(),
)
.await
{
Ok(nic) => nic,
Err(error) => {
if let Some(output) = output.take() {
let _ = output.send(Err(error));
return;
}
tracing::error!(?error, "failed to create native interface context");
tokio::select! {
_ = cancel.cancelled() => return,
_ = tokio::time::sleep(Duration::from_secs(1)) => {}
}
continue;
}
};
let result = tokio::select! {
biased;
_ = cancel.cancelled() => {
@@ -104,11 +131,13 @@ impl NativeTunRuntime {
}
let magic_dns = if let Some(ip) = ipv4 {
let shared_route_backend = nic.shared_route_backend_for_dns();
MagicDnsRuntime::start(
global_ctx.clone(),
packet_plane.clone(),
nic.ifname().await,
ip,
shared_route_backend,
)
} else {
MagicDnsRuntime::default()
@@ -161,6 +190,7 @@ impl NativeTunRuntime {
nic: self.nic.clone(),
closed: Arc::new(Notify::new()),
packet_plane,
shared_virtual_nic_registry: self.shared_virtual_nic_registry.clone(),
})
}
}
@@ -172,6 +202,7 @@ struct NativeDhcpIpv4Host {
nic: TunNicState,
closed: Arc<Notify>,
packet_plane: Arc<CorePacketPlane>,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
}
impl NativeDhcpIpv4Host {
@@ -199,21 +230,25 @@ impl NativeDhcpIpv4Host {
return Ok(Some(ip));
}
let mut nic = NicCtx::new(
let mut nic = create_nic_ctx(
self.global_ctx.clone(),
self.packet_plane.clone(),
self.nic.receiver(),
self.closed.clone(),
);
self.shared_virtual_nic_registry.clone(),
)
.await?;
tokio::select! {
_ = self.cancel.cancelled() => anyhow::bail!("instance is closing; DHCP update cancelled"),
result = nic.run(Some(ip), self.global_ctx.get_ipv6()) => result?,
}
let shared_route_backend = nic.shared_route_backend_for_dns();
let magic_dns = MagicDnsRuntime::start(
self.global_ctx.clone(),
self.packet_plane.clone(),
nic.ifname().await,
ip,
shared_route_backend,
);
self.nic.install(nic, magic_dns).await;
self.global_ctx.set_ipv4(Some(ip));
@@ -8,26 +8,40 @@ use easytier_core::{
instance::CorePacketPlane,
};
use futures::FutureExt as _;
use tokio::sync::{Mutex, Notify, mpsc};
use tokio::sync::{Mutex, Notify, mpsc, oneshot};
use tokio_util::sync::CancellationToken;
use super::{MagicDnsRuntime, tun_common::TunNicState};
use super::{
MagicDnsRuntime,
tun_common::{TunNicState, create_nic_ctx},
};
use crate::{
common::global_ctx::{ArcGlobalCtx, GlobalCtxEvent},
instance::virtual_nic::NicCtx,
instance::shared_virtual_nic::ArcSharedVirtualNicRegistry,
};
struct MobileTunAttachment {
fd: i32,
replace_tun_fd: bool,
completion: Option<oneshot::Sender<anyhow::Result<()>>>,
}
pub(super) struct NativeTunRuntime {
global_ctx: ArcGlobalCtx,
cancel: CancellationToken,
nic: TunNicState,
tun_fd: mpsc::Sender<i32>,
tun_fd_receiver: Mutex<Option<mpsc::Receiver<i32>>>,
tun_fd: mpsc::Sender<MobileTunAttachment>,
tun_fd_receiver: Mutex<Option<mpsc::Receiver<MobileTunAttachment>>>,
task: Mutex<Option<tokio::task::JoinHandle<()>>>,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
}
impl NativeTunRuntime {
pub(super) fn new(global_ctx: ArcGlobalCtx, cancel: CancellationToken) -> Self {
pub(super) fn new(
global_ctx: ArcGlobalCtx,
cancel: CancellationToken,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
) -> Self {
let (tun_fd, tun_fd_receiver) = mpsc::channel(16);
Self {
global_ctx,
@@ -36,6 +50,7 @@ impl NativeTunRuntime {
tun_fd,
tun_fd_receiver: Mutex::new(Some(tun_fd_receiver)),
task: Mutex::new(None),
shared_virtual_nic_registry,
}
}
@@ -51,23 +66,31 @@ impl NativeTunRuntime {
global_ctx: ArcGlobalCtx,
packet_plane: Arc<CorePacketPlane>,
fd: i32,
replace_tun_fd: bool,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
) -> anyhow::Result<()> {
nic_state.drain().await;
if fd <= 0 {
return Ok(());
}
let closed = Arc::new(Notify::new());
let mut nic = NicCtx::new(
let mut nic = create_nic_ctx(
global_ctx.clone(),
packet_plane.clone(),
nic_state.receiver(),
closed,
);
nic.run_for_mobile(fd).await.context("add ip failed")?;
let magic_dns = global_ctx
.get_ipv4()
.map(|ip| MagicDnsRuntime::start(global_ctx, packet_plane, None, ip))
.unwrap_or_default();
shared_virtual_nic_registry,
)
.await?;
nic.run_for_mobile(fd, replace_tun_fd)
.await
.context("attach mobile TUN failed")?;
let magic_dns = if let Some(ip) = global_ctx.get_ipv4() {
let shared_route_backend = nic.shared_route_backend_for_dns();
MagicDnsRuntime::start(global_ctx, packet_plane, None, ip, shared_route_backend)
} else {
MagicDnsRuntime::default()
};
nic_state.install(nic, magic_dns).await;
Ok(())
}
@@ -80,28 +103,43 @@ impl NativeTunRuntime {
let nic_state = self.nic.clone();
let global_ctx = self.global_ctx.clone();
let cancel = self.cancel.clone();
let shared_virtual_nic_registry = self.shared_virtual_nic_registry.clone();
self.task.lock().await.replace(tokio::spawn(async move {
loop {
let fd = tokio::select! {
let attachment = tokio::select! {
_ = cancel.cancelled() => return,
fd = tun_fds.recv() => match fd { Some(fd) => fd, None => return },
attachment = tun_fds.recv() => match attachment {
Some(attachment) => attachment,
None => return,
},
};
if let Err(error) = Self::install_mobile_tun(
nic_state.clone(),
global_ctx.clone(),
packet_plane.clone(),
fd,
)
.await
{
let result = if attachment.fd <= 0 {
nic_state.drain().await;
Ok(())
} else {
Self::install_mobile_tun(
nic_state.clone(),
global_ctx.clone(),
packet_plane.clone(),
attachment.fd,
attachment.replace_tun_fd,
shared_virtual_nic_registry.clone(),
)
.await
};
if let Err(error) = &result {
tracing::error!(?error, "failed to attach mobile TUN fd");
}
if let Some(completion) = attachment.completion {
let _ = completion.send(result);
}
}
}));
Ok(())
}
pub(super) async fn shutdown(&self) {
self.cancel.cancel();
if let Some(task) = self.task.lock().await.take() {
let _ = task.await;
}
@@ -110,10 +148,40 @@ impl NativeTunRuntime {
pub(super) fn attach_fd(&self, fd: i32) -> anyhow::Result<()> {
self.tun_fd
.try_send(fd)
.try_send(MobileTunAttachment {
fd,
replace_tun_fd: true,
completion: None,
})
.map_err(|error| anyhow::anyhow!("failed to send TUN fd: {error}"))
}
pub(super) async fn attach_mobile_fd(
&self,
fd: i32,
replace_tun_fd: bool,
) -> anyhow::Result<()> {
if self.task.lock().await.is_none() {
anyhow::bail!("mobile TUN runtime is not running");
}
let (completion, result) = oneshot::channel();
tokio::select! {
_ = self.cancel.cancelled() => anyhow::bail!("instance is closing; TUN attachment cancelled"),
send_result = self.tun_fd.send(MobileTunAttachment {
fd,
replace_tun_fd,
completion: Some(completion),
}) => send_result.map_err(|error| anyhow::anyhow!("failed to send TUN fd: {error}"))?,
}
tokio::select! {
_ = self.cancel.cancelled() => anyhow::bail!("instance is closing; TUN attachment cancelled"),
result = result => result
.map_err(|_| anyhow::anyhow!("mobile TUN runtime stopped before attachment completed"))?,
}
}
pub(super) fn dhcp_host(
&self,
operation: Arc<Mutex<()>>,
File diff suppressed because it is too large. Load diff
File diff suppressed because it is too large. Load diff
+40 -5
View File
@@ -7,6 +7,8 @@ use easytier_core::{
process_runtime::CoreProcessRuntime,
};
#[cfg(feature = "tun")]
use crate::instance::shared_virtual_nic::{ArcSharedVirtualNicRegistry, SharedVirtualNicRegistry};
use crate::{
common::global_ctx::{ArcGlobalCtx, GlobalCtx},
instance::{
@@ -15,6 +17,8 @@ use crate::{
},
socket::udp::RuntimeUdpSocket,
};
#[cfg(feature = "tun")]
use tokio::sync::Mutex;
pub(crate) struct TestInstance {
core: Arc<NativeCoreInstance>,
@@ -26,7 +30,27 @@ impl TestInstance {
config: TomlConfig,
process_runtime: Arc<CoreProcessRuntime>,
) -> Self {
Self::compose(config, process_runtime, |_| {})
Self::compose(
config,
process_runtime,
#[cfg(feature = "tun")]
Arc::new(Mutex::new(SharedVirtualNicRegistry::new())),
|_| {},
)
}
#[cfg(feature = "tun")]
pub fn new_with_process_runtime_and_shared_virtual_nic_registry(
config: TomlConfig,
process_runtime: Arc<CoreProcessRuntime>,
shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
) -> Self {
Self::compose(config, process_runtime, shared_virtual_nic_registry, |_| {})
}
#[cfg(feature = "tun")]
pub fn new_shared_virtual_nic_registry() -> ArcSharedVirtualNicRegistry {
Arc::new(Mutex::new(SharedVirtualNicRegistry::new()))
}
pub fn new_with_process_runtime_and_stun_provider(
@@ -35,14 +59,21 @@ impl TestInstance {
provider: Box<dyn StunSocketMapper<RuntimeUdpSocket>>,
) -> Self {
let provider: Arc<dyn StunSocketMapper<RuntimeUdpSocket>> = Arc::from(provider);
Self::compose(config, process_runtime, move |adapters| {
adapters.replace_stun_provider(provider);
})
Self::compose(
config,
process_runtime,
#[cfg(feature = "tun")]
Arc::new(Mutex::new(SharedVirtualNicRegistry::new())),
move |adapters| {
adapters.replace_stun_provider(provider);
},
)
}
fn compose(
config: TomlConfig,
process_runtime: Arc<CoreProcessRuntime>,
#[cfg(feature = "tun")] shared_virtual_nic_registry: ArcSharedVirtualNicRegistry,
customize: impl FnOnce(
&mut easytier_core::instance::CoreHostAdapters<
crate::instance::host::NativeInstanceHost,
@@ -50,7 +81,11 @@ impl TestInstance {
),
) -> Self {
let global_ctx = Arc::new(GlobalCtx::new(config.clone()));
let runtime_host = NativeInstanceRuntimeHost::new(global_ctx.clone());
let runtime_host = NativeInstanceRuntimeHost::new(
global_ctx.clone(),
#[cfg(feature = "tun")]
shared_virtual_nic_registry,
);
let mut adapters = runtime_core_host_adapters_with_packet_egress(
global_ctx.clone(),
process_runtime,
File diff suppressed because it is too large. Load diff
+3
View File
@@ -1,6 +1,9 @@
#[cfg(target_os = "linux")]
mod three_node;
#[cfg(all(target_os = "linux", feature = "tun"))]
mod shared_virtual_nic;
mod ipv6_test;
#[cfg(target_os = "linux")]
+476
View File
@@ -0,0 +1,476 @@
use std::{net::Ipv4Addr, process::Command, sync::Arc, time::Duration};
use easytier_core::{config::PeerId, process_runtime::CoreProcessRuntime};
use super::{
InstanceTestExt as _, add_ns_to_bridge, create_netns, del_netns, drop_insts, ping_test,
prepare_bridge,
};
use crate::{
common::{
config::{ConfigLoader, NetworkIdentity, TomlConfigLoader},
netns::{NetNS, ROOT_NETNS_NAME},
},
instance::{
shared_virtual_nic::{ArcSharedVirtualNicRegistry, SharedIpv4Route},
test_instance::TestInstance as Instance,
},
tunnel::common::tests::wait_for_condition,
};
const PROXY_CIDR: &str = "10.1.2.0/24";
const WAIT: Duration = Duration::from_secs(10);
#[derive(Clone)]
struct SharedTestRuntime {
process: Arc<CoreProcessRuntime>,
registry: ArcSharedVirtualNicRegistry,
}
impl SharedTestRuntime {
fn new() -> Self {
Self {
process: CoreProcessRuntime::new(),
registry: Instance::new_shared_virtual_nic_registry(),
}
}
fn instance(&self, config: TomlConfigLoader) -> Instance {
Instance::new_with_process_runtime_and_shared_virtual_nic_registry(
config,
self.process.clone(),
self.registry.clone(),
)
}
}
fn test_config(
instance_name: &str,
network_name: &str,
network_secret: &str,
netns: Option<&str>,
dev_name: Option<&str>,
ipv4: &str,
) -> TomlConfigLoader {
let config = TomlConfigLoader::default();
config.set_inst_name(instance_name.to_owned());
config.set_network_identity(NetworkIdentity::new(
network_name.to_owned(),
network_secret.to_owned(),
));
config.set_netns(netns.map(str::to_owned));
config.set_ipv4(Some(ipv4.parse().unwrap()));
config.set_ipv6(None);
config.set_dhcp(false);
config.set_listeners(vec![]);
config.set_socks5_portal(None);
let mut flags = config.get_flags();
flags.dev_name = dev_name.unwrap_or_default().to_owned();
flags.enable_ipv6 = false;
config.set_flags(flags);
config
}
fn test_dev_name() -> String {
format!("st{:08x}", rand::random::<u32>())
}
fn short_name(prefix: &str) -> String {
format!("{prefix}{:04x}", rand::random::<u16>())
}
struct TestNetnsGuard {
name: String,
}
impl TestNetnsGuard {
fn new(name: String, ipv4: &str, ipv6: &str) -> Self {
let guard = Self { name };
del_netns(&guard.name);
create_netns(&guard.name, ipv4, ipv6);
guard
}
}
impl Drop for TestNetnsGuard {
fn drop(&mut self) {
del_netns(&self.name);
}
}
struct ProxyLab {
source_ns: String,
owner_ns: String,
target_ns: String,
bridge: String,
}
impl ProxyLab {
fn new() -> Self {
let suffix = format!("{:04x}", rand::random::<u16>());
let lab = Self {
source_ns: format!("svs{suffix}"),
owner_ns: format!("svo{suffix}"),
target_ns: format!("svt{suffix}"),
bridge: format!("svb{suffix}"),
};
lab.cleanup();
create_netns(&lab.source_ns, "10.1.1.1/24", "fd11::1/64");
create_netns(&lab.owner_ns, "10.1.2.3/24", "fd12::3/64");
create_netns(&lab.target_ns, "10.1.2.4/24", "fd12::4/64");
prepare_bridge(&lab.bridge);
add_ns_to_bridge(&lab.bridge, &lab.owner_ns);
add_ns_to_bridge(&lab.bridge, &lab.target_ns);
lab
}
fn cleanup(&self) {
del_netns(&self.source_ns);
del_netns(&self.owner_ns);
del_netns(&self.target_ns);
let _ = Command::new("ip")
.args(["link", "del", &self.bridge])
.output();
}
}
impl Drop for ProxyLab {
fn drop(&mut self) {
self.cleanup();
}
}
async fn wait_tun_ready(instance: &Instance, expected: &str) {
wait_for_condition(
|| async { instance.get_global_ctx().get_tun_device_name().as_deref() == Some(expected) },
WAIT,
)
.await;
}
fn proxy_route_exists(
routes: &[easytier_proto::core_peer::peer::Route],
peer_id: PeerId,
proxy_cidr: &str,
) -> bool {
routes
.iter()
.any(|route| route.peer_id == peer_id && route.proxy_cidrs.iter().any(|c| c == proxy_cidr))
}
async fn wait_proxy_route(instance: &Instance, peer_id: PeerId, proxy_cidr: &str) {
wait_for_condition(
|| async {
proxy_route_exists(
&instance.get_core_instance().route_snapshots().await,
peer_id,
proxy_cidr,
)
},
WAIT,
)
.await;
}
async fn wait_proxy_route_absent(instance: &Instance, peer_id: PeerId, proxy_cidr: &str) {
wait_for_condition(
|| async {
!proxy_route_exists(
&instance.get_core_instance().route_snapshots().await,
peer_id,
proxy_cidr,
)
},
WAIT,
)
.await;
}
fn ipv4_route_exists_in_ns(ns: &str, needle: &str) -> bool {
let _root = NetNS::new(Some(ROOT_NETNS_NAME.to_owned())).guard();
let output = Command::new("ip")
.args(["netns", "exec", ns, "ip", "route", "show"])
.output()
.unwrap();
assert!(
output.status.success(),
"failed to list IPv4 routes in {ns}: {}",
String::from_utf8_lossy(&output.stderr)
);
String::from_utf8_lossy(&output.stdout)
.lines()
.any(|line| line.contains(needle))
}
#[cfg(feature = "proxy-cidr-monitor")]
async fn patch_proxy_cidr(
instance: &Instance,
action: crate::proto::api::config::ConfigPatchAction,
) {
use crate::proto::api::config::{InstanceConfigPatch, ProxyNetworkPatch};
instance
.get_config_patcher()
.apply_patch(InstanceConfigPatch {
proxy_networks: vec![ProxyNetworkPatch {
action: action as i32,
cidr: Some(PROXY_CIDR.parse().unwrap()),
mapped_cidr: None,
}],
..Default::default()
})
.await
.unwrap();
}
async fn shared_route_owner_count(
registry: &ArcSharedVirtualNicRegistry,
dev_name: &str,
route: &SharedIpv4Route,
) -> usize {
let nic = {
let registry = registry.lock().await;
registry.get_by_dev_name_for_test(dev_name)
};
let Some(nic) = nic else {
return 0;
};
nic.lock().await.ifcfg().owners_of_ipv4_route(route).len()
}
#[tokio::test]
#[serial_test::serial]
async fn same_namespace_members_share_tun_across_independent_networks() {
let dev_name = test_dev_name();
let first_peer_ns = TestNetnsGuard::new(short_name("sva"), "10.231.1.2/24", "fd31::2/64");
let second_peer_ns = TestNetnsGuard::new(short_name("svb"), "10.231.2.2/24", "fd32::2/64");
let runtime = SharedTestRuntime::new();
let mut first = runtime.instance(test_config(
"shared_tun_first",
"shared_tun_network_a",
"shared_tun_secret_a",
None,
Some(&dev_name),
"10.144.250.1/24",
));
let mut second = runtime.instance(test_config(
"shared_tun_second",
"shared_tun_network_b",
"shared_tun_secret_b",
None,
Some(&dev_name),
"10.144.251.1/24",
));
let mut first_peer = runtime.instance(test_config(
"shared_tun_first_peer",
"shared_tun_network_a",
"shared_tun_secret_a",
Some(&first_peer_ns.name),
None,
"10.144.250.2/24",
));
let mut second_peer = runtime.instance(test_config(
"shared_tun_second_peer",
"shared_tun_network_b",
"shared_tun_secret_b",
Some(&second_peer_ns.name),
None,
"10.144.251.2/24",
));
first.run().await.unwrap();
second.run().await.unwrap();
first_peer.run().await.unwrap();
second_peer.run().await.unwrap();
wait_tun_ready(&first, &dev_name).await;
wait_tun_ready(&second, &dev_name).await;
assert_eq!(
first.get_global_ctx().get_tun_device_name(),
second.get_global_ctx().get_tun_device_name()
);
first_peer.add_connector_url(first.ring_listener_url());
second_peer.add_connector_url(second.ring_listener_url());
wait_for_condition(
|| async {
first
.get_core_instance()
.route_snapshots()
.await
.iter()
.any(|route| route.peer_id == first_peer.peer_id())
&& second
.get_core_instance()
.route_snapshots()
.await
.iter()
.any(|route| route.peer_id == second_peer.peer_id())
},
WAIT,
)
.await;
wait_for_condition(
|| async { ping_test(&first_peer_ns.name, "10.144.250.1", None).await },
WAIT,
)
.await;
wait_for_condition(
|| async { ping_test(&second_peer_ns.name, "10.144.251.1", None).await },
WAIT,
)
.await;
drop_insts(vec![first, second, first_peer, second_peer]).await;
}
#[cfg(feature = "proxy-cidr-monitor")]
#[tokio::test]
#[serial_test::serial]
async fn runtime_proxy_patch_adds_and_removes_os_route() {
use crate::proto::api::config::ConfigPatchAction;
let lab = ProxyLab::new();
let source_dev = test_dev_name();
let destination_dev = test_dev_name();
let runtime = SharedTestRuntime::new();
let mut source = runtime.instance(test_config(
"shared_patch_source",
"shared_patch_network",
"shared_patch_secret",
Some(&lab.source_ns),
Some(&source_dev),
"10.144.244.1/24",
));
let mut destination = runtime.instance(test_config(
"shared_patch_destination",
"shared_patch_network",
"shared_patch_secret",
Some(&lab.owner_ns),
Some(&destination_dev),
"10.144.244.2/24",
));
source.run().await.unwrap();
destination.run().await.unwrap();
wait_tun_ready(&source, &source_dev).await;
wait_tun_ready(&destination, &destination_dev).await;
destination.add_connector_url(source.ring_listener_url());
wait_for_condition(
|| async {
source
.get_core_instance()
.route_snapshots()
.await
.iter()
.any(|route| route.peer_id == destination.peer_id())
},
WAIT,
)
.await;
assert!(!ipv4_route_exists_in_ns(
&lab.source_ns,
&format!("{PROXY_CIDR} dev {source_dev}")
));
patch_proxy_cidr(&destination, ConfigPatchAction::Add).await;
wait_proxy_route(&source, destination.peer_id(), PROXY_CIDR).await;
wait_for_condition(
|| async {
ipv4_route_exists_in_ns(&lab.source_ns, &format!("{PROXY_CIDR} dev {source_dev}"))
},
WAIT,
)
.await;
patch_proxy_cidr(&destination, ConfigPatchAction::Remove).await;
wait_proxy_route_absent(&source, destination.peer_id(), PROXY_CIDR).await;
wait_for_condition(
|| async {
!ipv4_route_exists_in_ns(&lab.source_ns, &format!("{PROXY_CIDR} dev {source_dev}"))
},
WAIT,
)
.await;
drop_insts(vec![source, destination]).await;
}
#[cfg(feature = "magic-dns")]
#[tokio::test]
#[serial_test::serial]
async fn magic_dns_route_lives_until_last_shared_owner_leaves() {
use crate::instance::dns_server::MAGIC_DNS_FAKE_IP;
let netns = TestNetnsGuard::new(short_name("svd"), "10.232.1.2/24", "fd42::2/64");
let dev_name = test_dev_name();
let runtime = SharedTestRuntime::new();
let first_config = test_config(
"shared_dns_first",
"shared_dns_network",
"shared_dns_secret",
Some(&netns.name),
Some(&dev_name),
"10.144.243.1/24",
);
let mut flags = first_config.get_flags();
flags.accept_dns = true;
first_config.set_flags(flags.clone());
let second_config = test_config(
"shared_dns_second",
"shared_dns_network",
"shared_dns_secret",
Some(&netns.name),
Some(&dev_name),
"10.144.242.2/24",
);
second_config.set_flags(flags);
let mut first = runtime.instance(first_config);
let mut second = runtime.instance(second_config);
first.run().await.unwrap();
second.run().await.unwrap();
wait_tun_ready(&first, &dev_name).await;
wait_tun_ready(&second, &dev_name).await;
let route = SharedIpv4Route::new(MAGIC_DNS_FAKE_IP.parse::<Ipv4Addr>().unwrap(), 32, None);
wait_for_condition(
|| async { shared_route_owner_count(&runtime.registry, &dev_name, &route).await == 2 },
WAIT,
)
.await;
assert!(ipv4_route_exists_in_ns(
&netns.name,
&format!("{MAGIC_DNS_FAKE_IP} dev {dev_name}")
));
drop_insts(vec![first]).await;
wait_for_condition(
|| async {
shared_route_owner_count(&runtime.registry, &dev_name, &route).await == 1
&& ipv4_route_exists_in_ns(
&netns.name,
&format!("{MAGIC_DNS_FAKE_IP} dev {dev_name}"),
)
},
WAIT,
)
.await;
drop_insts(vec![second]).await;
wait_for_condition(
|| async {
shared_route_owner_count(&runtime.registry, &dev_name, &route).await == 0
&& !ipv4_route_exists_in_ns(
&netns.name,
&format!("{MAGIC_DNS_FAKE_IP} dev {dev_name}"),
)
},
WAIT,
)
.await;
}
+151
View File
@@ -1017,6 +1017,157 @@ pub async fn public_ipv6_auto_addr_reconnect_reuses_same_address() {
drop_insts(vec![provider, client]).await;
}
#[cfg(feature = "tun")]
#[tokio::test]
#[serial_test::serial]
pub async fn shared_tun_public_ipv6_auto_addr_end_to_end() {
let lab = PublicIpv6Lab::setup_with_topology(PublicIpv6LabTopology::DelegatedPrefix);
let provider_dev = format!("st{:08x}", rand::random::<u32>());
let client_dev = format!("st{:08x}", rand::random::<u32>());
let process_runtime = CoreProcessRuntime::new();
let shared_virtual_nic_registry = Instance::new_shared_virtual_nic_registry();
let provider_cfg = get_public_ipv6_config(
"provider_shared_public_ipv6",
PublicIpv6Lab::PROVIDER_NS,
"10.144.144.1",
&provider_dev,
uuid::Uuid::parse_str("44444444-4444-4444-4444-444444444444").unwrap(),
);
provider_cfg.set_ipv6_public_addr_provider(true);
let client_cfg = get_public_ipv6_config(
"client_shared_public_ipv6",
PublicIpv6Lab::CLIENT_NS,
"10.144.144.2",
&client_dev,
uuid::Uuid::parse_str("55555555-5555-5555-5555-555555555555").unwrap(),
);
client_cfg.set_ipv6_public_addr_auto(true);
let client_peer_cfg = get_public_ipv6_config(
"client_shared_public_ipv6_peer",
PublicIpv6Lab::CLIENT_NS,
"10.144.145.3",
&client_dev,
uuid::Uuid::parse_str("66666666-6666-6666-6666-666666666666").unwrap(),
);
client_peer_cfg.set_listeners(vec![]);
let mut provider = Instance::new_with_process_runtime_and_shared_virtual_nic_registry(
provider_cfg,
process_runtime.clone(),
shared_virtual_nic_registry.clone(),
);
let mut client = Instance::new_with_process_runtime_and_shared_virtual_nic_registry(
client_cfg,
process_runtime.clone(),
shared_virtual_nic_registry.clone(),
);
let mut client_peer = Instance::new_with_process_runtime_and_shared_virtual_nic_registry(
client_peer_cfg,
process_runtime,
shared_virtual_nic_registry,
);
let mut client_events = client.get_global_ctx().subscribe();
let mut client_peer_events = client_peer.get_global_ctx().subscribe();
provider.run().await.unwrap();
client.run().await.unwrap();
client_peer.run().await.unwrap();
let shared_ifname = wait_for_tun_ready(&mut client_events).await;
assert_eq!(
shared_ifname,
wait_for_tun_ready(&mut client_peer_events).await
);
assert_eq!(shared_ifname, client_dev);
provider.add_connector_url("tcp://10.1.1.2:11010".parse().unwrap());
wait_for_condition(
|| async {
provider.get_core_instance().route_snapshots().await.len() == 1
&& client.get_core_instance().route_snapshots().await.len() == 1
},
Duration::from_secs(8),
)
.await;
wait_for_condition(
|| async {
provider
.get_core_instance()
.node_snapshot()
.await
.ipv6_public_addr_prefix
== Some(PublicIpv6Lab::PROVIDER_PREFIX.parse().unwrap())
},
Duration::from_secs(10),
)
.await;
let leased = wait_for_public_ipv6_addr(&client).await;
wait_for_public_ipv6_route(&provider, leased).await;
wait_for_condition(
|| async {
addr_exists_in_ns(PublicIpv6Lab::CLIENT_NS, &client_dev, &leased.to_string())
&& route_exists_in_ns(
PublicIpv6Lab::CLIENT_NS,
&format!("default dev {client_dev}"),
)
&& route_exists_in_ns(
PublicIpv6Lab::PROVIDER_NS,
&format!("{} dev {provider_dev}", leased.address()),
)
},
Duration::from_secs(10),
)
.await;
wait_for_condition(
|| async { ping6_test(PublicIpv6Lab::CLIENT_NS, PublicIpv6Lab::SERVER_IP, None).await },
Duration::from_secs(10),
)
.await;
wait_for_condition(
|| async {
ping6_test(
PublicIpv6Lab::SERVER_NS,
leased.address().to_string().as_str(),
None,
)
.await
},
Duration::from_secs(10),
)
.await;
drop_insts(vec![provider, client, client_peer]).await;
drop(lab);
}
#[cfg(feature = "tun")]
async fn wait_for_tun_ready(
receiver: &mut tokio::sync::broadcast::Receiver<crate::common::global_ctx::GlobalCtxEvent>,
) -> String {
tokio::time::timeout(Duration::from_secs(5), async {
loop {
match receiver.recv().await.unwrap() {
crate::common::global_ctx::GlobalCtxEvent::TunDeviceReady(ifname) => return ifname,
crate::common::global_ctx::GlobalCtxEvent::TunDeviceError(error) => {
panic!("tun device error: {error}")
}
_ => {}
}
}
})
.await
.expect("timed out waiting for tun ready")
}
#[rstest::rstest]
#[tokio::test]
#[serial_test::serial]