From 9b372d3087863b8ddde92074c9e1c4bbc71fba70 Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Sun, 22 Feb 2026 02:33:10 +0100 Subject: [PATCH] node: watch config/ip change --- easytier/src/dns/node.rs | 60 ++++++++++++++++++++++++++++++--- easytier/src/dns/utils/dirty.rs | 4 +++ 2 files changed, 60 insertions(+), 4 deletions(-) diff --git a/easytier/src/dns/node.rs b/easytier/src/dns/node.rs index b230e17d..b3b3e9de 100644 --- a/easytier/src/dns/node.rs +++ b/easytier/src/dns/node.rs @@ -1,3 +1,4 @@ +use crate::common::global_ctx::GlobalCtxEvent; use crate::dns::config::DNS_SERVER_RPC_ADDR; use crate::dns::peer_mgr::DnsPeerMgr; use crate::peers::peer_manager::PeerManager; @@ -8,7 +9,9 @@ use crate::proto::rpc_types::controller::BaseController; use crate::tunnel::tcp::TcpTunnelConnector; use std::sync::Arc; use std::time::Duration; +use tokio::sync::broadcast; use tokio::task::JoinSet; +use tokio::time::{sleep_until, Instant}; use uuid::Uuid; #[derive(Debug)] @@ -47,12 +50,61 @@ impl DnsNode { ..Default::default() }; + let rr_interval = Duration::from_secs(1); + let mut last_heartbeat = Instant::now(); + let sleep = sleep_until(last_heartbeat); + tokio::pin!(sleep); + + let mut subscriber = self.mgr.get_global_ctx_ref().subscribe(); + loop { - self.mgr.dirty.notify.notified().await; - if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await { - tracing::error!("heartbeat failed: {:?}", e); + let next_heartbeat = last_heartbeat + + if self.mgr.dirty.peek() { + rr_interval + } else { + rr_interval / 8 + }; + sleep.as_mut().reset(next_heartbeat); + + tokio::select! { + biased; + + _ = &mut sleep => { + if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await { + // TODO: try to start server + tracing::error!("heartbeat failed: {:?}", e); + } + + last_heartbeat = Instant::now(); + } + + _ = self.mgr.dirty.notify.notified() => {} + + event = subscriber.recv() => { + match event { + Ok( + GlobalCtxEvent::DhcpIpv4Changed(..) + | GlobalCtxEvent::DhcpIpv4Conflicted(..), + ) => { + tracing::info!(?event, "ip change detected, rebuilding snapshot"); + } + Ok(GlobalCtxEvent::ConfigPatched(patch)) => { + // TODO: inspect patch + tracing::info!(?patch, "config change detected, rebuilding snapshot"); + } + Err(broadcast::error::RecvError::Lagged(n)) => { + tracing::warn!("event listener lagged, skipped {n} events, rebuilding snapshot"); + } + Err(broadcast::error::RecvError::Closed) => { + tracing::info!("event bus closed"); + break; + } + _ => continue, + } + + self.mgr.dirty.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 3faa3247..2b369512 100644 --- a/easytier/src/dns/utils/dirty.rs +++ b/easytier/src/dns/utils/dirty.rs @@ -24,6 +24,10 @@ impl DirtyFlag { self.0.store(true, Ordering::Release); } + pub fn peek(&self) -> bool { + self.0.load(Ordering::Acquire) + } + pub fn reset(&self) -> bool { self.0.swap(false, Ordering::Acquire) }