From 2e68a8db914d9d3dc14c239be5e3d4954b2e44fc Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Thu, 19 Feb 2026 14:11:13 +0100 Subject: [PATCH] client: rewrite heartbeat, add resync --- easytier/src/dns/client.rs | 66 +++++++++++++++++------------------- easytier/src/proto/dns.proto | 16 +++++---- 2 files changed, 41 insertions(+), 41 deletions(-) diff --git a/easytier/src/dns/client.rs b/easytier/src/dns/client.rs index 59fc20c2..e8a8b4d5 100644 --- a/easytier/src/dns/client.rs +++ b/easytier/src/dns/client.rs @@ -20,7 +20,6 @@ use moka::future::Cache; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::Duration; -use tokio::sync::Mutex; use tokio::task::JoinSet; use url::Url; use uuid::Uuid; @@ -145,17 +144,10 @@ pub struct DnsClient { mgr: Arc, tasks: JoinSet<()>, - - #[derivative(Debug = "ignore")] - // Client to talk to local DnsServer - server_rpc_client: Arc>>, } impl DnsClient { pub fn new(peer_mgr: Arc) -> Self { - let connector = TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone()); - let server_rpc_client = Arc::new(Mutex::new(StandAloneClient::new(connector))); - let mgr = Arc::new(DnsPeerManager::new(peer_mgr.clone())); peer_mgr .get_peer_rpc_mgr() @@ -169,7 +161,6 @@ impl DnsClient { Self { mgr, tasks: JoinSet::new(), - server_rpc_client, } } @@ -178,41 +169,48 @@ impl DnsClient { } pub async fn run(&self) { - let mut snapshot = Default::default(); - let mut digest = Vec::new(); + let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone())); + let mut heartbeat = HeartbeatRequest { + id: Some(self.id().into()), + + ..Default::default() + }; loop { - let snapshot = self.mgr.dirty.swap(false, Ordering::Release).then(|| { - snapshot = self.mgr.snapshot(); - digest = snapshot.digest(); - - snapshot.clone() - }); - - let heartbeat = HeartbeatRequest { - id: Some(self.id().into()), - digest: digest.clone(), - snapshot, - }; - - if let Err(e) = self.heartbeat(heartbeat).await { + if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await { tracing::error!("DnsClient heartbeat failed: {:?}", e); } tokio::time::sleep(Duration::from_secs(1)).await; } } - async fn heartbeat(&self, heartbeat: HeartbeatRequest) -> anyhow::Result<()> { - // scoped_client of StandAloneClient takes &mut self - let client = { - let mut client = self.server_rpc_client.lock().await; // Lock the mutex - client - .scoped_client::>("".to_string()) - .await? + async fn heartbeat( + &self, + rpc: &mut StandAloneClient, + heartbeat: &mut HeartbeatRequest, + ) -> anyhow::Result<()> { + let request = if self.mgr.dirty.swap(false, Ordering::Release) { + let snapshot = self.mgr.snapshot(); + heartbeat.digest = snapshot.digest(); + heartbeat.snapshot = Some(snapshot); + heartbeat.clone() + } else { + let snapshot = heartbeat.snapshot.take(); + let request = heartbeat.clone(); + heartbeat.snapshot = snapshot; + request }; - client - .heartbeat(BaseController::default(), heartbeat) + let client = rpc + .scoped_client::>("".to_string()) .await?; + + let response = client.heartbeat(BaseController::default(), request).await?; + if response.resync { + client + .heartbeat(BaseController::default(), heartbeat.clone()) + .await?; + } + Ok(()) } diff --git a/easytier/src/proto/dns.proto b/easytier/src/proto/dns.proto index 2e1bda06..91b867e7 100644 --- a/easytier/src/proto/dns.proto +++ b/easytier/src/proto/dns.proto @@ -12,12 +12,6 @@ message ZoneConfigPb { repeated string forwarders = 5; } -message DnsSnapshot { - repeated ZoneConfigPb zones = 1; - repeated common.SocketAddr addresses = 2; - repeated common.Url listeners = 3; -} - message GetExportConfigRequest {} message GetExportConfigResponse { @@ -29,13 +23,21 @@ service DnsPeerManagerRpc { rpc GetExportConfig(GetExportConfigRequest) returns (GetExportConfigResponse) {} } +message DnsSnapshot { + repeated ZoneConfigPb zones = 1; + repeated common.SocketAddr addresses = 2; + repeated common.Url listeners = 3; +} + message HeartbeatRequest { common.UUID id = 1; bytes digest = 2; optional DnsSnapshot snapshot = 3; } -message HeartbeatResponse {} +message HeartbeatResponse { + bool resync = 1; +} service DnsServerRpc { rpc Heartbeat(HeartbeatRequest) returns (HeartbeatResponse) {}