mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-02 17:15:43 +00:00
client: add Heartbeat
This commit is contained in:
+57
-17
@@ -1,13 +1,12 @@
|
|||||||
use super::config::DNS_SERVER_RPC_ADDR;
|
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::peers::peer_manager::PeerManager;
|
||||||
use crate::proto::dns::{
|
use crate::proto::dns::{DnsPeerManagerRpcServer, DnsServerRpcClientFactory, HeartbeatRequest};
|
||||||
DeterministicDigest, DnsPeerManagerRpcServer, DnsServerRpcClientFactory, HeartbeatRequest,
|
|
||||||
};
|
|
||||||
use crate::proto::peer_rpc::RoutePeerInfo;
|
use crate::proto::peer_rpc::RoutePeerInfo;
|
||||||
use crate::proto::rpc_impl::standalone::StandAloneClient;
|
use crate::proto::rpc_impl::standalone::StandAloneClient;
|
||||||
use crate::proto::rpc_types::controller::BaseController;
|
use crate::proto::rpc_types::controller::BaseController;
|
||||||
use crate::tunnel::tcp::TcpTunnelConnector;
|
use crate::tunnel::tcp::TcpTunnelConnector;
|
||||||
|
use crate::utils::DeterministicDigest;
|
||||||
use derivative::Derivative;
|
use derivative::Derivative;
|
||||||
use std::sync::atomic::Ordering;
|
use std::sync::atomic::Ordering;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
@@ -15,6 +14,53 @@ use std::time::Duration;
|
|||||||
use tokio::task::JoinSet;
|
use tokio::task::JoinSet;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Default)]
|
||||||
|
pub struct Heartbeat {
|
||||||
|
pub(super) id: Uuid,
|
||||||
|
pub(super) digest: Vec<u8>,
|
||||||
|
pub(super) snapshot: Option<DnsSnapshot>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Heartbeat> for HeartbeatRequest {
|
||||||
|
fn from(value: Heartbeat) -> Self {
|
||||||
|
Self {
|
||||||
|
id: Some(value.id.into()),
|
||||||
|
digest: value.digest,
|
||||||
|
snapshot: value.snapshot.map(Into::into),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TryFrom<HeartbeatRequest> for Heartbeat {
|
||||||
|
type Error = anyhow::Error;
|
||||||
|
|
||||||
|
fn try_from(value: HeartbeatRequest) -> Result<Self, Self::Error> {
|
||||||
|
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)]
|
#[derive(Derivative)]
|
||||||
#[derivative(Debug)]
|
#[derivative(Debug)]
|
||||||
pub struct DnsClient {
|
pub struct DnsClient {
|
||||||
@@ -42,16 +88,12 @@ impl DnsClient {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn id(&self) -> Uuid {
|
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) {
|
pub async fn run(&self) {
|
||||||
let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone()));
|
let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone()));
|
||||||
let mut heartbeat = HeartbeatRequest {
|
let mut heartbeat = Heartbeat::new(self.id());
|
||||||
id: Some(self.id().into()),
|
|
||||||
|
|
||||||
..Default::default()
|
|
||||||
};
|
|
||||||
loop {
|
loop {
|
||||||
if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await {
|
if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await {
|
||||||
tracing::error!("DnsClient heartbeat failed: {:?}", e);
|
tracing::error!("DnsClient heartbeat failed: {:?}", e);
|
||||||
@@ -63,17 +105,15 @@ impl DnsClient {
|
|||||||
async fn heartbeat(
|
async fn heartbeat(
|
||||||
&self,
|
&self,
|
||||||
rpc: &mut StandAloneClient<TcpTunnelConnector>,
|
rpc: &mut StandAloneClient<TcpTunnelConnector>,
|
||||||
heartbeat: &mut HeartbeatRequest,
|
heartbeat: &mut Heartbeat,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
let request =
|
let request =
|
||||||
if heartbeat.snapshot.is_none() || self.mgr.dirty.swap(false, Ordering::Release) {
|
if heartbeat.snapshot.is_none() || self.mgr.dirty.swap(false, Ordering::Release) {
|
||||||
let snapshot = self.mgr.snapshot();
|
heartbeat.update(self.mgr.snapshot());
|
||||||
heartbeat.digest = snapshot.digest();
|
heartbeat.clone().into()
|
||||||
heartbeat.snapshot = Some(snapshot);
|
|
||||||
heartbeat.clone()
|
|
||||||
} else {
|
} else {
|
||||||
let snapshot = heartbeat.snapshot.take();
|
let snapshot = heartbeat.snapshot.take();
|
||||||
let request = heartbeat.clone();
|
let request = heartbeat.clone().into();
|
||||||
heartbeat.snapshot = snapshot;
|
heartbeat.snapshot = snapshot;
|
||||||
request
|
request
|
||||||
};
|
};
|
||||||
@@ -85,7 +125,7 @@ impl DnsClient {
|
|||||||
let response = client.heartbeat(BaseController::default(), request).await?;
|
let response = client.heartbeat(BaseController::default(), request).await?;
|
||||||
if response.resync {
|
if response.resync {
|
||||||
client
|
client
|
||||||
.heartbeat(BaseController::default(), heartbeat.clone())
|
.heartbeat(BaseController::default(), heartbeat.clone().into())
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user