From 2e1f9bc9ccae007a21335bfbb6427e942e93588b Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Sun, 22 Feb 2026 03:34:45 +0100 Subject: [PATCH] node & server: separate dirty flag, remove DirtyState --- easytier/src/dns/node.rs | 2 +- easytier/src/dns/node_mgr.rs | 5 +-- easytier/src/dns/peer_mgr.rs | 6 +-- easytier/src/dns/server.rs | 78 +++++++++++++++++++-------------- easytier/src/dns/utils/dirty.rs | 21 ++++++--- 5 files changed, 66 insertions(+), 46 deletions(-) diff --git a/easytier/src/dns/node.rs b/easytier/src/dns/node.rs index 2f9aff1b..7152b78a 100644 --- a/easytier/src/dns/node.rs +++ b/easytier/src/dns/node.rs @@ -74,7 +74,7 @@ impl DnsNode { last_heartbeat = Instant::now(); } - _ = self.mgr.dirty.notify.notified() => {} + _ = self.mgr.dirty.notified() => {} event = subscriber.recv() => { match event { diff --git a/easytier/src/dns/node_mgr.rs b/easytier/src/dns/node_mgr.rs index 67c2bb4a..baa6b541 100644 --- a/easytier/src/dns/node_mgr.rs +++ b/easytier/src/dns/node_mgr.rs @@ -1,5 +1,5 @@ use crate::dns::utils::addr::NameServerAddr; -use crate::dns::utils::dirty::{DirtyFlag, DirtyState}; +use crate::dns::utils::dirty::DirtyFlag; use crate::dns::zone::{Zone, ZoneGroup}; use crate::proto::dns::DnsNodeMgrRpc; use crate::proto::dns::{DnsSnapshot, HeartbeatRequest, HeartbeatResponse}; @@ -47,7 +47,7 @@ pub struct DnsNodeMgrDirtyFlags { #[derive(Debug)] pub struct DnsNodeMgr { nodes: Cache, - pub(super) dirty: DirtyState, + pub(super) dirty: DnsNodeMgrDirtyFlags, } impl DnsNodeMgr { @@ -145,7 +145,6 @@ impl DnsNodeMgrRpc for DnsNodeMgr { } self.nodes.insert(id, new).await; - self.dirty.notify.notify_one(); } false } else { diff --git a/easytier/src/dns/peer_mgr.rs b/easytier/src/dns/peer_mgr.rs index b04cfe77..b9e43edf 100644 --- a/easytier/src/dns/peer_mgr.rs +++ b/easytier/src/dns/peer_mgr.rs @@ -1,7 +1,7 @@ use crate::common::config::ConfigLoader; use crate::common::PeerId; use crate::dns::config::{DnsExportConfig, DnsGlobalCtxExt}; -use crate::dns::utils::dirty::{DirtyFlag, DirtyState}; +use crate::dns::utils::dirty::DirtyFlag; use crate::dns::zone::ZoneGroup; use crate::peer_center::instance::PeerCenterPeerManagerTrait; use crate::peers::peer_manager::PeerManager; @@ -42,7 +42,7 @@ const DNS_PEER_TTL: Duration = Duration::from_secs(3); #[derive(Debug, Deref)] pub struct DnsPeerMgr { peers: Cache, - pub(super) dirty: DirtyState, + pub(super) dirty: DirtyFlag, #[deref] mgr: Arc, @@ -103,7 +103,7 @@ impl DnsPeerMgr { } } - self.dirty.notify.notify_one(); + self.dirty.notify_one(); } async fn fetch(&self, peer_id: PeerId) -> anyhow::Result { diff --git a/easytier/src/dns/server.rs b/easytier/src/dns/server.rs index 6cbea3c6..45562bc0 100644 --- a/easytier/src/dns/server.rs +++ b/easytier/src/dns/server.rs @@ -131,9 +131,6 @@ pub struct DnsServer { #[derivative(Debug = "ignore")] catalog: DynamicCatalog, - - /// Current set of hijacked addresses (only UDP protocol addresses). - addresses: RwLock>, } const DNS_SERVER_LISTENER_TCP_TIMEOUT: Duration = Duration::from_secs(5); @@ -153,23 +150,24 @@ impl DnsServer { Self { mgr, catalog: DynamicCatalog::new(), - addresses: RwLock::new(HashSet::new()), } } - async fn reload_addresses(&self, addresses: impl IntoIterator) { + async fn reload_addresses( + &self, + addresses: impl IntoIterator, + current: &mut HashSet, + ) { let addresses = addresses.into_iter().collect::>(); - let mut active = self.addresses.write().await; - - let added = addresses.difference(&active).cloned().collect_vec(); - let removed = active.difference(&addresses).cloned().collect_vec(); + let added = addresses.difference(¤t).cloned().collect_vec(); + let removed = current.difference(&addresses).cloned().collect_vec(); if added.is_empty() && removed.is_empty() { return; } - *active = addresses; + *current = addresses; // TODO } @@ -209,30 +207,44 @@ impl DnsServer { pub async fn run(&self) { let dirty = &self.mgr.dirty; - let mut runtime = None; - loop { - dirty.notify.notified().await; - if dirty.catalog.reset() { - self.catalog.replace(self.mgr.catalog()).await; - } - - if dirty.addresses.reset() { - self.reload_addresses(self.mgr.iter_addresses()).await; - } - - if dirty.listeners.reset() { - if let Err(e) = self - .reload_listeners(self.mgr.iter_listeners(), &mut runtime) - .await - { - tracing::error!("failed to reload listeners: {:?}", e); - dirty.listeners.mark(); - dirty.notify.notify_one(); + tokio::join!( + async { + loop { + dirty.catalog.notified().await; + if dirty.catalog.reset() { + self.catalog.replace(self.mgr.catalog()).await; + } + tokio::time::sleep(Duration::from_secs(1)).await; } - } - - tokio::time::sleep(Duration::from_secs(1)).await; - } + }, + async { + let mut addresses = HashSet::new(); + loop { + dirty.addresses.notified().await; + if dirty.addresses.reset() { + self.reload_addresses(self.mgr.iter_addresses(), &mut addresses) + .await; + } + tokio::time::sleep(Duration::from_secs(1)).await; + } + }, + async { + let mut runtime = None; + loop { + dirty.listeners.notified().await; + if dirty.listeners.reset() { + if let Err(e) = self + .reload_listeners(self.mgr.iter_listeners(), &mut runtime) + .await + { + tracing::error!("failed to reload listeners: {:?}", e); + dirty.listeners.mark(); + } + } + tokio::time::sleep(Duration::from_secs(1)).await; + } + }, + ); } } diff --git a/easytier/src/dns/utils/dirty.rs b/easytier/src/dns/utils/dirty.rs index 2b369512..97fd30e8 100644 --- a/easytier/src/dns/utils/dirty.rs +++ b/easytier/src/dns/utils/dirty.rs @@ -11,24 +11,33 @@ pub struct DirtyState { pub notify: Notify, } -#[derive(Derivative, Debug)] +#[derive(Derivative, Debug, Deref)] #[derivative(Default)] -pub struct DirtyFlag(#[derivative(Default(value = "AtomicBool::new(true)"))] AtomicBool); +pub struct DirtyFlag { + #[derivative(Default(value = "AtomicBool::new(true)"))] + dirty: AtomicBool, + #[deref] + notify: Notify, +} impl DirtyFlag { pub fn new(value: bool) -> Self { - Self(AtomicBool::new(value)) + Self { + dirty: AtomicBool::new(value), + notify: Notify::new(), + } } pub fn mark(&self) { - self.0.store(true, Ordering::Release); + self.dirty.store(true, Ordering::Release); + self.notify.notify_one(); } pub fn peek(&self) -> bool { - self.0.load(Ordering::Acquire) + self.dirty.load(Ordering::Acquire) } pub fn reset(&self) -> bool { - self.0.swap(false, Ordering::Acquire) + self.dirty.swap(false, Ordering::Acquire) } }