mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-02 17:15:43 +00:00
node & peer_mgr: refactor refresh logic, add my_peer_id check
node
This commit is contained in:
@@ -1,6 +1,5 @@
|
|||||||
use crate::common::global_ctx::{ArcGlobalCtx, GlobalCtxEvent};
|
use crate::common::global_ctx::{ArcGlobalCtx, GlobalCtxEvent};
|
||||||
use crate::common::scoped_task::ScopedTask;
|
use crate::common::scoped_task::ScopedTask;
|
||||||
use crate::common::PeerId;
|
|
||||||
use crate::dns::config::{DNS_SERVER_ELECTION_INTERVAL, DNS_SERVER_RPC_ADDR};
|
use crate::dns::config::{DNS_SERVER_ELECTION_INTERVAL, DNS_SERVER_RPC_ADDR};
|
||||||
use crate::dns::peer_mgr::DnsPeerMgr;
|
use crate::dns::peer_mgr::DnsPeerMgr;
|
||||||
use crate::dns::server::DnsServer;
|
use crate::dns::server::DnsServer;
|
||||||
@@ -185,7 +184,12 @@ impl DnsNode {
|
|||||||
event = subscriber.recv() => {
|
event = subscriber.recv() => {
|
||||||
match event {
|
match event {
|
||||||
Ok(GlobalCtxEvent::PeerInfoUpdated(peer_ids)) => {
|
Ok(GlobalCtxEvent::PeerInfoUpdated(peer_ids)) => {
|
||||||
self.refresh(&mut tasks, peer_ids);
|
for peer_id in peer_ids {
|
||||||
|
let mgr = self.mgr.clone();
|
||||||
|
tasks.spawn(async move {
|
||||||
|
mgr.refresh(peer_id).await;
|
||||||
|
});
|
||||||
|
}
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
Ok(
|
Ok(
|
||||||
@@ -242,20 +246,4 @@ impl DnsNode {
|
|||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn refresh(&self, tasks: &mut JoinSet<()>, peer_ids: Vec<PeerId>) {
|
|
||||||
let my_peer_id = self.peer_mgr.my_peer_id();
|
|
||||||
for peer_id in peer_ids {
|
|
||||||
if peer_id == my_peer_id {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
let mgr = self.mgr.clone();
|
|
||||||
let route = self.peer_mgr.get_route();
|
|
||||||
tasks.spawn(async move {
|
|
||||||
if let Some(peer_info) = route.get_peer_info(peer_id).await {
|
|
||||||
mgr.refresh(peer_id, peer_info.dns).await;
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ use crate::dns::utils::dirty::DirtyFlag;
|
|||||||
use crate::dns::zone::ZoneGroup;
|
use crate::dns::zone::ZoneGroup;
|
||||||
use crate::peer_center::instance::PeerCenterPeerManagerTrait;
|
use crate::peer_center::instance::PeerCenterPeerManagerTrait;
|
||||||
use crate::peers::peer_manager::PeerManager;
|
use crate::peers::peer_manager::PeerManager;
|
||||||
|
use crate::peers::route_trait::Route;
|
||||||
use crate::proto::dns::{
|
use crate::proto::dns::{
|
||||||
DnsPeerMgrRpc, DnsPeerMgrRpcClientFactory, DnsSnapshot, GetExportConfigRequest,
|
DnsPeerMgrRpc, DnsPeerMgrRpcClientFactory, DnsSnapshot, GetExportConfigRequest,
|
||||||
GetExportConfigResponse, ZoneData,
|
GetExportConfigResponse, ZoneData,
|
||||||
@@ -78,12 +79,24 @@ impl DnsPeerMgr {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) async fn refresh(&self, peer_id: PeerId, digest: Vec<u8>) {
|
pub async fn refresh(&self, peer_id: PeerId) {
|
||||||
if let Some(info) = self.peers.get(&peer_id).await {
|
if peer_id == self.peer_mgr.my_peer_id() {
|
||||||
if info.digest == *digest {
|
self.dirty.mark();
|
||||||
return;
|
self.dirty.notify_one();
|
||||||
}
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some(route) = self.peer_mgr.get_route().get_peer_info(peer_id).await else {
|
||||||
|
return;
|
||||||
};
|
};
|
||||||
|
if self
|
||||||
|
.peers
|
||||||
|
.get(&peer_id)
|
||||||
|
.await
|
||||||
|
.is_some_and(|info| info.digest == *route.dns)
|
||||||
|
{
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
self.dirty.mark();
|
self.dirty.mark();
|
||||||
|
|
||||||
@@ -91,12 +104,8 @@ impl DnsPeerMgr {
|
|||||||
Ok(info) => {
|
Ok(info) => {
|
||||||
self.peers.insert(peer_id, info).await;
|
self.peers.insert(peer_id, info).await;
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(error) => {
|
||||||
tracing::warn!(
|
tracing::warn!(%peer_id, ?error, "failed to fetch dns export config from peer");
|
||||||
"failed to fetch dns export config from peer {}: {:?}",
|
|
||||||
peer_id,
|
|
||||||
e
|
|
||||||
);
|
|
||||||
self.peers.invalidate(&peer_id).await;
|
self.peers.invalidate(&peer_id).await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user