From 704121432179b031888b3ccae869af4cc5d4a7a9 Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Fri, 20 Feb 2026 16:21:53 +0100 Subject: [PATCH] client: add Heartbeat --- easytier/src/dns/client.rs | 74 +++++++++++++++++++++++++++++--------- 1 file changed, 57 insertions(+), 17 deletions(-) diff --git a/easytier/src/dns/client.rs b/easytier/src/dns/client.rs index 9838f06b..095d9983 100644 --- a/easytier/src/dns/client.rs +++ b/easytier/src/dns/client.rs @@ -1,13 +1,12 @@ use super::config::DNS_SERVER_RPC_ADDR; -use crate::dns::peer_mgr::DnsPeerMgr; +use crate::dns::peer_mgr::{DnsPeerMgr, DnsSnapshot}; use crate::peers::peer_manager::PeerManager; -use crate::proto::dns::{ - DeterministicDigest, DnsPeerManagerRpcServer, DnsServerRpcClientFactory, HeartbeatRequest, -}; +use crate::proto::dns::{DnsPeerManagerRpcServer, DnsServerRpcClientFactory, HeartbeatRequest}; use crate::proto::peer_rpc::RoutePeerInfo; use crate::proto::rpc_impl::standalone::StandAloneClient; use crate::proto::rpc_types::controller::BaseController; use crate::tunnel::tcp::TcpTunnelConnector; +use crate::utils::DeterministicDigest; use derivative::Derivative; use std::sync::atomic::Ordering; use std::sync::Arc; @@ -15,6 +14,53 @@ use std::time::Duration; use tokio::task::JoinSet; use uuid::Uuid; +#[derive(Debug, Clone, Default)] +pub struct Heartbeat { + pub(super) id: Uuid, + pub(super) digest: Vec, + pub(super) snapshot: Option, +} + +impl Heartbeat { + pub fn new(id: Uuid) -> Self { + Self { + id, + + ..Default::default() + } + } + + pub fn update(&mut self, snapshot: DnsSnapshot) { + self.digest = snapshot.digest(); + self.snapshot = Some(snapshot); + } +} + +impl From for HeartbeatRequest { + fn from(value: Heartbeat) -> Self { + Self { + id: Some(value.id.into()), + digest: value.digest, + snapshot: value.snapshot.map(Into::into), + } + } +} + +impl TryFrom for Heartbeat { + type Error = anyhow::Error; + + fn try_from(value: HeartbeatRequest) -> Result { + Ok(Self { + id: value + .id + .ok_or(anyhow::anyhow!("missing id in heartbeat"))? + .into(), + digest: value.digest, + snapshot: value.snapshot.map(TryInto::try_into).transpose()?, + }) + } +} + #[derive(Derivative)] #[derivative(Debug)] pub struct DnsClient { @@ -42,16 +88,12 @@ impl DnsClient { } pub fn id(&self) -> Uuid { - self.mgr.get_global_ctx_ref().get_id().into() + self.mgr.get_global_ctx_ref().get_id() } pub async fn run(&self) { let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone())); - let mut heartbeat = HeartbeatRequest { - id: Some(self.id().into()), - - ..Default::default() - }; + let mut heartbeat = Heartbeat::new(self.id()); loop { if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await { tracing::error!("DnsClient heartbeat failed: {:?}", e); @@ -63,17 +105,15 @@ impl DnsClient { async fn heartbeat( &self, rpc: &mut StandAloneClient, - heartbeat: &mut HeartbeatRequest, + heartbeat: &mut Heartbeat, ) -> anyhow::Result<()> { let request = if heartbeat.snapshot.is_none() || self.mgr.dirty.swap(false, Ordering::Release) { - let snapshot = self.mgr.snapshot(); - heartbeat.digest = snapshot.digest(); - heartbeat.snapshot = Some(snapshot); - heartbeat.clone() + heartbeat.update(self.mgr.snapshot()); + heartbeat.clone().into() } else { let snapshot = heartbeat.snapshot.take(); - let request = heartbeat.clone(); + let request = heartbeat.clone().into(); heartbeat.snapshot = snapshot; request }; @@ -85,7 +125,7 @@ impl DnsClient { let response = client.heartbeat(BaseController::default(), request).await?; if response.resync { client - .heartbeat(BaseController::default(), heartbeat.clone()) + .heartbeat(BaseController::default(), heartbeat.clone().into()) .await?; }