From 57a2428f21072de5556a30fbd62c620454be6d8a Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Mon, 23 Feb 2026 00:06:46 +0100 Subject: [PATCH] node: add server election peer_mgr server: add peer_mgr peer_mgr --- easytier/src/dns/config/mod.rs | 8 +++- easytier/src/dns/node.rs | 70 ++++++++++++++++++++++++++-------- easytier/src/dns/peer_mgr.rs | 17 ++++----- easytier/src/dns/server.rs | 8 ++-- 4 files changed, 73 insertions(+), 30 deletions(-) diff --git a/easytier/src/dns/config/mod.rs b/easytier/src/dns/config/mod.rs index 86331ae4..4f5b30f9 100644 --- a/easytier/src/dns/config/mod.rs +++ b/easytier/src/dns/config/mod.rs @@ -4,6 +4,7 @@ use hickory_proto::xfer::Protocol; use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4}; use std::str::FromStr; use std::sync::LazyLock; +use std::time::Duration; use url::Url; mod dns; @@ -11,14 +12,17 @@ pub use dns::*; mod policy; mod zone; +pub static DNS_DEFAULT_TLD: LazyLock = + LazyLock::new(|| LowerName::from_str("et.net.").unwrap()); pub const DNS_DEFAULT_ADDRESS: NameServerAddr = NameServerAddr { protocol: Protocol::Udp, addr: SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(100, 100, 100, 101), 53)), }; -pub static DNS_DEFAULT_TLD: LazyLock = - LazyLock::new(|| LowerName::from_str("et.net.").unwrap()); + pub static DNS_SERVER_RPC_ADDR: LazyLock = LazyLock::new(|| Url::parse("tcp://127.0.0.1:49813").unwrap()); +pub const DNS_SERVER_ELECTION_INTERVAL: Duration = Duration::from_secs(5); + pub(super) static DNS_SUPPORTED_PROTOCOLS: [Protocol; 2] = [ Protocol::Udp, Protocol::Tcp, diff --git a/easytier/src/dns/node.rs b/easytier/src/dns/node.rs index ae969c4a..f557953a 100644 --- a/easytier/src/dns/node.rs +++ b/easytier/src/dns/node.rs @@ -1,28 +1,30 @@ use crate::common::global_ctx::GlobalCtxEvent; use crate::common::PeerId; -use crate::dns::config::DNS_SERVER_RPC_ADDR; +use crate::dns::config::{DNS_SERVER_ELECTION_INTERVAL, DNS_SERVER_RPC_ADDR}; use crate::dns::peer_mgr::DnsPeerMgr; +use crate::dns::server::DnsServer; use crate::peers::peer_manager::PeerManager; +use crate::peers::NicPacketFilter; use crate::proto::dns::{DnsNodeMgrRpcClientFactory, DnsPeerMgrRpcServer, HeartbeatRequest}; -use crate::proto::rpc_impl::standalone::StandAloneClient; +use crate::proto::rpc_impl::standalone::{StandAloneClient, StandAloneServer}; use crate::proto::rpc_types::controller::BaseController; -use crate::tunnel::tcp::TcpTunnelConnector; +use crate::tunnel::tcp::{TcpTunnelConnector, TcpTunnelListener}; use std::sync::Arc; use std::time::Duration; use tokio::sync::{broadcast, Notify}; use tokio::task::JoinSet; -use tokio::time::{sleep_until, Instant}; +use tokio::time::{sleep, sleep_until, Instant}; use uuid::Uuid; #[derive(Debug)] pub struct DnsNode { mgr: Arc, - election: Arc, + peer_mgr: Arc, } impl DnsNode { - pub fn new(peer_mgr: Arc, election: Arc) -> Self { + pub fn new(peer_mgr: Arc) -> Self { let mgr = Arc::new(DnsPeerMgr::new(peer_mgr.clone())); peer_mgr .get_peer_rpc_mgr() @@ -33,17 +35,55 @@ impl DnsNode { &peer_mgr.get_global_ctx_ref().get_network_name(), ); - Self { - mgr, - election, - } + Self { mgr, peer_mgr } } pub fn id(&self) -> Uuid { - self.mgr.get_global_ctx_ref().get_id() + self.peer_mgr.get_global_ctx_ref().get_id() } pub async fn run(&self) { + let election = Notify::new(); + + tokio::join!(self.run_election(&election), self.run_node(&election)); + } + + async fn run_election(&self, election: &Notify) { + loop { + tokio::select! { + biased; + _ = election.notified() => {} + _ = sleep(DNS_SERVER_ELECTION_INTERVAL) => {} + } + + let mut rpc = + StandAloneServer::new(TcpTunnelListener::new(DNS_SERVER_RPC_ADDR.clone())); + + if rpc.serve().await.is_err() { + // Another instance already owns the address — that's fine. + continue; + } + + tracing::info!("won DNS server election, starting DnsServer"); + + let server = Arc::new(DnsServer::new(self.peer_mgr.clone(), rpc)); + + self.peer_mgr + .add_nic_packet_process_pipeline(Box::new(server.clone())) + .await; + + server.run().await; + + let _ = self + .peer_mgr + .remove_nic_packet_process_pipeline(server.id()) + .await; + + tracing::warn!("DnsServer exited, will retry election"); + } + } + + async fn run_node(&self, election: &Notify) { let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone())); let mut heartbeat = HeartbeatRequest { id: Some(self.id().into()), @@ -55,7 +95,7 @@ impl DnsNode { let sleep = sleep_until(last_heartbeat); tokio::pin!(sleep); - let mut subscriber = self.mgr.get_global_ctx_ref().subscribe(); + let mut subscriber = self.peer_mgr.get_global_ctx_ref().subscribe(); let mut tasks = JoinSet::new(); loop { @@ -73,7 +113,7 @@ impl DnsNode { _ = &mut sleep => { if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await { tracing::error!("heartbeat failed: {:?}", e); - self.election.notify_one(); + election.notify_one(); } last_heartbeat = Instant::now(); @@ -143,13 +183,13 @@ impl DnsNode { } fn refresh(&self, tasks: &mut JoinSet<()>, peer_ids: Vec) { - let my_peer_id = self.mgr.my_peer_id(); + let my_peer_id = self.peer_mgr.my_peer_id(); for peer_id in peer_ids { if peer_id == my_peer_id { continue; } let mgr = self.mgr.clone(); - let route = mgr.get_route(); + let route = self.peer_mgr.get_route(); tasks.spawn(async move { if let Some(peer_info) = route.get_peer_info(peer_id).await { mgr.refresh(peer_id, peer_info.dns).await; diff --git a/easytier/src/dns/peer_mgr.rs b/easytier/src/dns/peer_mgr.rs index b9e43edf..83237db5 100644 --- a/easytier/src/dns/peer_mgr.rs +++ b/easytier/src/dns/peer_mgr.rs @@ -13,7 +13,6 @@ use crate::proto::rpc_types; use crate::proto::rpc_types::controller::BaseController; use crate::utils::DeterministicDigest; use anyhow::Context; -use derive_more::Deref; use itertools::Itertools; use moka::future::Cache; use std::sync::Arc; @@ -39,26 +38,25 @@ impl TryFrom for DnsPeerInfo { const DNS_PEER_TTL: Duration = Duration::from_secs(3); -#[derive(Debug, Deref)] +#[derive(Debug)] pub struct DnsPeerMgr { peers: Cache, pub(super) dirty: DirtyFlag, - #[deref] - mgr: Arc, + peer_mgr: Arc, } impl DnsPeerMgr { pub fn new(peer_mgr: Arc) -> Self { Self { - mgr: peer_mgr.clone(), peers: Cache::builder().time_to_live(DNS_PEER_TTL).build(), dirty: Default::default(), + peer_mgr: peer_mgr.clone(), } } pub fn snapshot(&self) -> DnsSnapshot { - let global_ctx = self.get_global_ctx_ref(); + let global_ctx = self.peer_mgr.get_global_ctx_ref(); let config = global_ctx.config.get_dns(); let zones = config @@ -107,10 +105,11 @@ impl DnsPeerMgr { } async fn fetch(&self, peer_id: PeerId) -> anyhow::Result { - self.get_peer_rpc_mgr() + self.peer_mgr + .get_peer_rpc_mgr() .rpc_client() .scoped_client::>( - self.mgr.my_peer_id(), + self.peer_mgr.my_peer_id(), peer_id, "".to_string(), ) @@ -130,6 +129,6 @@ impl DnsPeerMgrRpc for DnsPeerMgr { _: Self::Controller, _: GetExportConfigRequest, ) -> rpc_types::error::Result { - Ok(self.get_global_ctx_ref().dns_export_config()) + Ok(self.peer_mgr.get_global_ctx_ref().dns_export_config()) } } diff --git a/easytier/src/dns/server.rs b/easytier/src/dns/server.rs index c7dbcb6b..ac44dfe6 100644 --- a/easytier/src/dns/server.rs +++ b/easytier/src/dns/server.rs @@ -1,4 +1,3 @@ -use crate::common::PeerId; use crate::dns::node_mgr::DnsNodeMgr; use crate::dns::utils::addr::NameServerAddr; use crate::peer_center::instance::PeerCenterPeerManagerTrait; @@ -154,10 +153,11 @@ impl ResponseHandler for Response { pub struct DnsServer { mgr: Arc, + peer_mgr: Arc, + #[derivative(Debug = "ignore")] catalog: DynamicCatalog, - my_peer_id: PeerId, addresses: Arc>>, } @@ -172,8 +172,8 @@ impl DnsServer { Self { mgr, + peer_mgr, catalog: DynamicCatalog::new(), - my_peer_id: peer_mgr.my_peer_id(), addresses: Arc::new(Default::default()), } } @@ -349,7 +349,7 @@ impl DnsServer { ip_packet.set_checksum(ipv4::checksum(&ip_packet.to_immutable())); // Route the response back to ourselves so it goes through the tun device. - zc_packet.mut_peer_manager_header().unwrap().to_peer_id = self.my_peer_id.into(); + zc_packet.mut_peer_manager_header().unwrap().to_peer_id = self.peer_mgr.my_peer_id().into(); Some(()) }