mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-02 09:09:17 +00:00
client: rewrite heartbeat, add resync
This commit is contained in:
+32
-34
@@ -20,7 +20,6 @@ use moka::future::Cache;
|
|||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::sync::Mutex;
|
|
||||||
use tokio::task::JoinSet;
|
use tokio::task::JoinSet;
|
||||||
use url::Url;
|
use url::Url;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
@@ -145,17 +144,10 @@ pub struct DnsClient {
|
|||||||
mgr: Arc<DnsPeerManager>,
|
mgr: Arc<DnsPeerManager>,
|
||||||
|
|
||||||
tasks: JoinSet<()>,
|
tasks: JoinSet<()>,
|
||||||
|
|
||||||
#[derivative(Debug = "ignore")]
|
|
||||||
// Client to talk to local DnsServer
|
|
||||||
server_rpc_client: Arc<Mutex<StandAloneClient<TcpTunnelConnector>>>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DnsClient {
|
impl DnsClient {
|
||||||
pub fn new(peer_mgr: Arc<PeerManager>) -> Self {
|
pub fn new(peer_mgr: Arc<PeerManager>) -> 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()));
|
let mgr = Arc::new(DnsPeerManager::new(peer_mgr.clone()));
|
||||||
peer_mgr
|
peer_mgr
|
||||||
.get_peer_rpc_mgr()
|
.get_peer_rpc_mgr()
|
||||||
@@ -169,7 +161,6 @@ impl DnsClient {
|
|||||||
Self {
|
Self {
|
||||||
mgr,
|
mgr,
|
||||||
tasks: JoinSet::new(),
|
tasks: JoinSet::new(),
|
||||||
server_rpc_client,
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -178,41 +169,48 @@ impl DnsClient {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn run(&self) {
|
pub async fn run(&self) {
|
||||||
let mut snapshot = Default::default();
|
let mut rpc = StandAloneClient::new(TcpTunnelConnector::new(DNS_SERVER_RPC_ADDR.clone()));
|
||||||
let mut digest = Vec::new();
|
let mut heartbeat = HeartbeatRequest {
|
||||||
|
id: Some(self.id().into()),
|
||||||
|
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
loop {
|
loop {
|
||||||
let snapshot = self.mgr.dirty.swap(false, Ordering::Release).then(|| {
|
if let Err(e) = self.heartbeat(&mut rpc, &mut heartbeat).await {
|
||||||
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 {
|
|
||||||
tracing::error!("DnsClient heartbeat failed: {:?}", e);
|
tracing::error!("DnsClient heartbeat failed: {:?}", e);
|
||||||
}
|
}
|
||||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn heartbeat(&self, heartbeat: HeartbeatRequest) -> anyhow::Result<()> {
|
async fn heartbeat(
|
||||||
// scoped_client of StandAloneClient takes &mut self
|
&self,
|
||||||
let client = {
|
rpc: &mut StandAloneClient<TcpTunnelConnector>,
|
||||||
let mut client = self.server_rpc_client.lock().await; // Lock the mutex
|
heartbeat: &mut HeartbeatRequest,
|
||||||
client
|
) -> anyhow::Result<()> {
|
||||||
.scoped_client::<DnsServerRpcClientFactory<BaseController>>("".to_string())
|
let request = if self.mgr.dirty.swap(false, Ordering::Release) {
|
||||||
.await?
|
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
|
let client = rpc
|
||||||
.heartbeat(BaseController::default(), heartbeat)
|
.scoped_client::<DnsServerRpcClientFactory<BaseController>>("".to_string())
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
|
let response = client.heartbeat(BaseController::default(), request).await?;
|
||||||
|
if response.resync {
|
||||||
|
client
|
||||||
|
.heartbeat(BaseController::default(), heartbeat.clone())
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -12,12 +12,6 @@ message ZoneConfigPb {
|
|||||||
repeated string forwarders = 5;
|
repeated string forwarders = 5;
|
||||||
}
|
}
|
||||||
|
|
||||||
message DnsSnapshot {
|
|
||||||
repeated ZoneConfigPb zones = 1;
|
|
||||||
repeated common.SocketAddr addresses = 2;
|
|
||||||
repeated common.Url listeners = 3;
|
|
||||||
}
|
|
||||||
|
|
||||||
message GetExportConfigRequest {}
|
message GetExportConfigRequest {}
|
||||||
|
|
||||||
message GetExportConfigResponse {
|
message GetExportConfigResponse {
|
||||||
@@ -29,13 +23,21 @@ service DnsPeerManagerRpc {
|
|||||||
rpc GetExportConfig(GetExportConfigRequest) returns (GetExportConfigResponse) {}
|
rpc GetExportConfig(GetExportConfigRequest) returns (GetExportConfigResponse) {}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
message DnsSnapshot {
|
||||||
|
repeated ZoneConfigPb zones = 1;
|
||||||
|
repeated common.SocketAddr addresses = 2;
|
||||||
|
repeated common.Url listeners = 3;
|
||||||
|
}
|
||||||
|
|
||||||
message HeartbeatRequest {
|
message HeartbeatRequest {
|
||||||
common.UUID id = 1;
|
common.UUID id = 1;
|
||||||
bytes digest = 2;
|
bytes digest = 2;
|
||||||
optional DnsSnapshot snapshot = 3;
|
optional DnsSnapshot snapshot = 3;
|
||||||
}
|
}
|
||||||
|
|
||||||
message HeartbeatResponse {}
|
message HeartbeatResponse {
|
||||||
|
bool resync = 1;
|
||||||
|
}
|
||||||
|
|
||||||
service DnsServerRpc {
|
service DnsServerRpc {
|
||||||
rpc Heartbeat(HeartbeatRequest) returns (HeartbeatResponse) {}
|
rpc Heartbeat(HeartbeatRequest) returns (HeartbeatResponse) {}
|
||||||
|
|||||||
Reference in New Issue
Block a user