peer_mgr(refresh): add backoff retry

fmt

refresh
This commit is contained in:
Luna Yao
2026-04-28 15:03:53 +02:00
parent 40d3e48cc4
commit 80041086de
4 changed files with 98 additions and 43 deletions
+2
View File
@@ -24,3 +24,5 @@ pub const DNS_NODE_TTI: Duration = Duration::from_secs(5);
pub const DNS_NODE_RR_INTERVAL: Duration = Duration::from_secs(4); pub const DNS_NODE_RR_INTERVAL: Duration = Duration::from_secs(4);
pub const DNS_SERVER_ELECTION_INTERVAL: Duration = Duration::from_secs(5); pub const DNS_SERVER_ELECTION_INTERVAL: Duration = Duration::from_secs(5);
pub const DNS_PEER_TTI: Duration = Duration::from_secs(3); pub const DNS_PEER_TTI: Duration = Duration::from_secs(3);
pub const DNS_PEER_REFRESH_ATTEMPTS: usize = 3;
pub const DNS_PEER_REFRESH_BACKOFF: Duration = Duration::from_secs(1);
+17 -12
View File
@@ -1,5 +1,8 @@
use crate::common::global_ctx::{ArcGlobalCtx, GlobalCtxEvent}; use crate::common::global_ctx::{ArcGlobalCtx, GlobalCtxEvent};
use crate::dns::config::{DNS_NODE_RR_INTERVAL, DNS_SERVER_ELECTION_INTERVAL, DNS_SERVER_RPC_ADDR}; use crate::dns::config::{
DNS_NODE_RR_INTERVAL, DNS_PEER_REFRESH_ATTEMPTS, DNS_PEER_REFRESH_BACKOFF,
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;
#[cfg(feature = "tun")] #[cfg(feature = "tun")]
@@ -12,16 +15,16 @@ use crate::tunnel::tcp::{TcpTunnelConnector, TcpTunnelListener};
use crate::utils::task::CancellableTask; use crate::utils::task::CancellableTask;
use std::io; use std::io;
use std::sync::Arc; use std::sync::Arc;
use tokio::sync::{Notify, broadcast}; use tokio::sync::{broadcast, Notify};
use tokio::task::JoinSet; use tokio::task::JoinSet;
use tokio::time::{Instant, sleep, sleep_until}; use tokio::time::{sleep, sleep_until, Instant};
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use tracing::instrument; use tracing::instrument;
use uuid::Uuid; use uuid::Uuid;
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
struct DnsNodeRuntime { struct DnsNodeRuntime {
mgr: Arc<DnsPeerMgr>, mgr: DnsPeerMgr,
#[cfg(feature = "tun")] #[cfg(feature = "tun")]
nic_ctx: ArcNicCtx, // TODO: REMOVE THIS nic_ctx: ArcNicCtx, // TODO: REMOVE THIS
@@ -87,8 +90,8 @@ impl DnsNodeRuntime {
..Default::default() ..Default::default()
}; };
let mut last_heartbeat = Instant::now(); let mut last_heartbeat = Instant::now();
let sleep = sleep_until(last_heartbeat); let timer = sleep_until(last_heartbeat);
tokio::pin!(sleep); tokio::pin!(timer);
let mut subscriber = self.global_ctx.subscribe(); let mut subscriber = self.global_ctx.subscribe();
let mut tasks = JoinSet::new(); let mut tasks = JoinSet::new();
@@ -101,17 +104,17 @@ impl DnsNodeRuntime {
} else { } else {
DNS_NODE_RR_INTERVAL / 4 DNS_NODE_RR_INTERVAL / 4
}; };
sleep.as_mut().reset(next_heartbeat); timer.as_mut().reset(next_heartbeat);
tokio::select! { tokio::select! {
biased; biased;
_ = token.cancelled() => { _ = token.cancelled() => {
tracing::info!("DnsNode received shutdown signal, exiting node loop"); tracing::info!("DnsNode received shutdown signal, exiting main loop");
break; break;
} }
_ = &mut sleep => { _ = &mut timer => {
if let Err(error) = self.heartbeat(&mut rpc, &mut heartbeat).await { if let Err(error) = self.heartbeat(&mut rpc, &mut heartbeat).await {
tracing::error!(?error, "heartbeat failed"); tracing::error!(?error, "heartbeat failed");
self.elect.notify_one(); self.elect.notify_one();
@@ -128,7 +131,9 @@ impl DnsNodeRuntime {
for peer_id in peer_ids { for peer_id in peer_ids {
let mgr = self.mgr.clone(); let mgr = self.mgr.clone();
tasks.spawn(async move { tasks.spawn(async move {
mgr.refresh(peer_id).await; if let Err(error) = mgr.refresh(peer_id, DNS_PEER_REFRESH_ATTEMPTS, DNS_PEER_REFRESH_BACKOFF).await {
tracing::error!(?error, ?peer_id, "failed to refresh peer");
}
}); });
} }
continue; continue;
@@ -209,7 +214,7 @@ impl DnsNode {
#[cfg(feature = "tun")] nic_ctx: ArcNicCtx, // TODO: REMOVE THIS #[cfg(feature = "tun")] nic_ctx: ArcNicCtx, // TODO: REMOVE THIS
) -> Self { ) -> Self {
let runtime = DnsNodeRuntime { let runtime = DnsNodeRuntime {
mgr: Arc::new(DnsPeerMgr::new(peer_mgr.clone(), global_ctx.clone())), mgr: DnsPeerMgr::new(peer_mgr.clone(), global_ctx.clone()),
#[cfg(feature = "tun")] #[cfg(feature = "tun")]
nic_ctx, nic_ctx,
peer_mgr, peer_mgr,
@@ -291,7 +296,7 @@ mod tests {
let global_ctx = peer_mgr.get_global_ctx(); let global_ctx = peer_mgr.get_global_ctx();
let nic_ctx: ArcNicCtx = Arc::new(Mutex::new(None)); let nic_ctx: ArcNicCtx = Arc::new(Mutex::new(None));
DnsNodeRuntime { DnsNodeRuntime {
mgr: Arc::new(DnsPeerMgr::new(peer_mgr.clone(), global_ctx.clone())), mgr: DnsPeerMgr::new(peer_mgr.clone(), global_ctx.clone()),
nic_ctx, nic_ctx,
peer_mgr, peer_mgr,
global_ctx, global_ctx,
+56 -25
View File
@@ -1,6 +1,9 @@
use crate::common::PeerId; use crate::common::PeerId;
use crate::common::global_ctx::ArcGlobalCtx; use crate::common::global_ctx::ArcGlobalCtx;
use crate::dns::config::{DNS_PEER_TTI, DnsExportConfig, DnsGlobalCtxExt}; use crate::dns::config::{
DNS_PEER_REFRESH_ATTEMPTS, DNS_PEER_REFRESH_BACKOFF, DNS_PEER_TTI, DnsExportConfig,
DnsGlobalCtxExt,
};
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;
@@ -14,9 +17,13 @@ use crate::proto::rpc_types::controller::BaseController;
use crate::proto::utils::TransientDigest; use crate::proto::utils::TransientDigest;
use crate::utils::dirty::DirtyFlag; use crate::utils::dirty::DirtyFlag;
use anyhow::Context; use anyhow::Context;
use futures::StreamExt;
use futures::stream;
use moka::future::Cache; use moka::future::Cache;
use std::ops::Deref; use std::ops::Deref;
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration;
use tokio::time::sleep;
use tracing::instrument; use tracing::instrument;
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
@@ -68,18 +75,44 @@ impl DnsPeerMgrInner {
} }
} }
pub async fn refresh(&self, peer_id: PeerId) { pub async fn refresh(
&self,
peer_id: PeerId,
mut attempts: usize,
mut backoff: Duration,
) -> anyhow::Result<bool> {
loop {
attempts = attempts.saturating_sub(1);
let result = self.try_refresh(peer_id).await;
match &result {
Ok(_) => return result,
Err(_) if attempts == 0 => return result,
Err(error) => {
tracing::error!(
?error,
?peer_id,
"failed to refresh peer info, retrying in {:?}",
backoff
);
sleep(backoff).await;
backoff *= 2;
}
}
}
}
async fn try_refresh(&self, peer_id: PeerId) -> anyhow::Result<bool> {
if peer_id == self.peer_mgr.my_peer_id() { if peer_id == self.peer_mgr.my_peer_id() {
self.dirty.mark(); self.dirty.mark();
return; return Ok(true);
} }
let Some(route) = self.peer_mgr.get_route().get_peer_info(peer_id).await else { let Some(route) = self.peer_mgr.get_route().get_peer_info(peer_id).await else {
if self.peers.remove(&peer_id).await.is_some() { if self.peers.remove(&peer_id).await.is_some() {
tracing::debug!(%peer_id, "peer route disappeared, removing from cache"); tracing::debug!(?peer_id, "peer route disappeared, removing from cache");
self.dirty.mark(); self.dirty.mark();
} }
return; return Ok(true);
}; };
if self if self
@@ -88,25 +121,20 @@ impl DnsPeerMgrInner {
.await .await
.is_some_and(|info| route.dns == info.digest) .is_some_and(|info| route.dns == info.digest)
{ {
return; return Ok(false);
} }
let info = if !route.dns.is_empty() {
if !route.dns.is_empty() { let info = self.fetch(peer_id).await.with_context(|| {
self.fetch(peer_id).await.inspect_err(|error| { format!("failed to fetch dns export config from peer {}", peer_id)
tracing::warn!(%peer_id, ?error, "failed to fetch dns export config from peer"); })?;
}).ok()
} else {
None
};
if let Some(info) = info {
self.peers.insert(peer_id, info).await; self.peers.insert(peer_id, info).await;
} else { } else {
self.peers.invalidate(&peer_id).await; self.peers.invalidate(&peer_id).await;
} }
self.dirty.mark(); self.dirty.mark();
Ok(true)
} }
#[instrument(skip(self), level = "trace", ret)] #[instrument(skip(self), level = "trace", ret)]
@@ -139,7 +167,7 @@ impl DnsPeerMgrRpc for DnsPeerMgrInner {
} }
} }
#[derive(Debug)] #[derive(Debug, Clone)]
pub struct DnsPeerMgr(Arc<DnsPeerMgrInner>); pub struct DnsPeerMgr(Arc<DnsPeerMgrInner>);
impl DnsPeerMgr { impl DnsPeerMgr {
@@ -433,7 +461,7 @@ mod tests {
let mgr = DnsPeerMgr::new(peer_mgr.clone(), peer_mgr.get_global_ctx()); let mgr = DnsPeerMgr::new(peer_mgr.clone(), peer_mgr.get_global_ctx());
mgr.dirty.reset(); mgr.dirty.reset();
mgr.refresh(peer_mgr.my_peer_id()).await; mgr.try_refresh(peer_mgr.my_peer_id()).await.unwrap();
assert!(mgr.dirty.peek()); assert!(mgr.dirty.peek());
} }
@@ -449,7 +477,7 @@ mod tests {
let mgr = DnsPeerMgr::new(peer_mgr, get_mock_global_ctx()); let mgr = DnsPeerMgr::new(peer_mgr, get_mock_global_ctx());
mgr.dirty.reset(); mgr.dirty.reset();
mgr.refresh(987_654).await; mgr.try_refresh(987_654).await.unwrap();
assert!(!mgr.dirty.peek()); assert!(!mgr.dirty.peek());
} }
@@ -496,7 +524,7 @@ mod tests {
.await; .await;
mgr.dirty.reset(); mgr.dirty.reset();
mgr.refresh(remote_id).await; mgr.try_refresh(remote_id).await.unwrap();
sleep(Duration::from_millis(50)).await; sleep(Duration::from_millis(50)).await;
assert!(!mgr.dirty.peek()); assert!(!mgr.dirty.peek());
@@ -527,7 +555,7 @@ mod tests {
.expect("route should appear"); .expect("route should appear");
local_dns.dirty.reset(); local_dns.dirty.reset();
local_dns.refresh(remote.my_peer_id()).await; local_dns.try_refresh(remote.my_peer_id()).await.unwrap();
assert!(local_dns.dirty.peek()); assert!(local_dns.dirty.peek());
let snapshot = local_dns.snapshot(); let snapshot = local_dns.snapshot();
@@ -567,7 +595,7 @@ mod tests {
.await .await
.expect("route to peer_b should appear"); .expect("route to peer_b should appear");
local_dns.refresh(peer_a.my_peer_id()).await; local_dns.try_refresh(peer_a.my_peer_id()).await.unwrap();
let snapshot = local_dns.snapshot(); let snapshot = local_dns.snapshot();
assert!( assert!(
@@ -643,7 +671,7 @@ mod tests {
.expect("route to keep_peer should appear"); .expect("route to keep_peer should appear");
local_dns.dirty.reset(); local_dns.dirty.reset();
local_dns.refresh(fail_id).await; local_dns.try_refresh(fail_id).await.unwrap();
assert!(local_dns.dirty.peek()); assert!(local_dns.dirty.peek());
assert!(local_dns.peers.get(&fail_id).await.is_none()); assert!(local_dns.peers.get(&fail_id).await.is_none());
@@ -719,11 +747,14 @@ mod tests {
.await; .await;
local_dns.dirty.reset(); local_dns.dirty.reset();
local_dns.refresh(changed_peer.my_peer_id()).await; local_dns
.try_refresh(changed_peer.my_peer_id())
.await
.unwrap();
assert!(local_dns.dirty.peek()); assert!(local_dns.dirty.peek());
local_dns.dirty.reset(); local_dns.dirty.reset();
local_dns.refresh(unchanged_id).await; local_dns.try_refresh(unchanged_id).await.unwrap();
assert!(!local_dns.dirty.peek()); assert!(!local_dns.dirty.peek());
let unchanged_cache = local_dns let unchanged_cache = local_dns
+23 -6
View File
@@ -339,7 +339,9 @@ async fn wait_peer_zone_visibility(
let deadline = Instant::now() + Duration::from_secs(20); let deadline = Instant::now() + Duration::from_secs(20);
loop { loop {
dns.refresh(target_peer_id).await; dns.refresh(target_peer_id, Default::default(), Default::default())
.await
.unwrap();
let snapshot = dns.snapshot(); let snapshot = dns.snapshot();
@@ -486,7 +488,10 @@ records = ["secret IN A 10.99.0.9"]
// Verify from peer-sync view to avoid host-wide DNS-server election side effects. // Verify from peer-sync view to avoid host-wide DNS-server election side effects.
let _dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx()); let _dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx());
let dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx()); let dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx());
dns_b.refresh(peer_a.my_peer_id()).await; dns_b
.refresh(peer_a.my_peer_id(), Default::default(), Default::default())
.await
.unwrap();
let snapshot = dns_b.snapshot(); let snapshot = dns_b.snapshot();
assert!( assert!(
@@ -529,7 +534,10 @@ disabled = true
let _dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx()); let _dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx());
let dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx()); let dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx());
dns_b.refresh(peer_a.my_peer_id()).await; dns_b
.refresh(peer_a.my_peer_id(), Default::default(), Default::default())
.await
.unwrap();
let snapshot = dns_b.snapshot(); let snapshot = dns_b.snapshot();
assert!( assert!(
@@ -866,7 +874,10 @@ records = ["secret IN A 10.66.3.8"]
let _dns_c = DnsPeerMgr::new(peer_c.clone(), peer_c.get_global_ctx()); let _dns_c = DnsPeerMgr::new(peer_c.clone(), peer_c.get_global_ctx());
let dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx()); let dns_a = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx());
dns_a.refresh(peer_c.my_peer_id()).await; dns_a
.refresh(peer_c.my_peer_id(), Default::default(), Default::default())
.await
.unwrap();
let snapshot = dns_a.snapshot(); let snapshot = dns_a.snapshot();
assert!( assert!(
@@ -897,7 +908,10 @@ async fn config_string_two_nodes_peer_dns_offline_then_rejoin() {
let _dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx()); let _dns_b = DnsPeerMgr::new(peer_b.clone(), peer_b.get_global_ctx());
let dns_a_online = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx()); let dns_a_online = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx());
dns_a_online.refresh(peer_b.my_peer_id()).await; dns_a_online
.refresh(peer_b.my_peer_id(), Default::default(), Default::default())
.await
.unwrap();
assert!( assert!(
dns_a_online dns_a_online
.snapshot() .snapshot()
@@ -950,7 +964,10 @@ async fn config_string_two_nodes_peer_dns_offline_then_rejoin() {
.expect("route should re-appear"); .expect("route should re-appear");
let dns_a_rejoin = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx()); let dns_a_rejoin = DnsPeerMgr::new(peer_a.clone(), peer_a.get_global_ctx());
dns_a_rejoin.refresh(peer_b.my_peer_id()).await; dns_a_rejoin
.refresh(peer_b.my_peer_id(), Default::default(), Default::default())
.await
.unwrap();
assert!( assert!(
dns_a_rejoin dns_a_rejoin
.snapshot() .snapshot()