2021-11-22 16:28:30 +00:00
|
|
|
mod bucket;
|
|
|
|
mod bucket_entry;
|
2021-12-14 14:48:33 +00:00
|
|
|
mod debug;
|
2021-11-22 16:28:30 +00:00
|
|
|
mod find_nodes;
|
|
|
|
mod node_ref;
|
2021-11-26 14:54:38 +00:00
|
|
|
mod stats_accounting;
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
use crate::dht::*;
|
|
|
|
use crate::intf::*;
|
|
|
|
use crate::network_manager::*;
|
|
|
|
use crate::rpc_processor::*;
|
|
|
|
use crate::xx::*;
|
|
|
|
use crate::*;
|
|
|
|
use alloc::str::FromStr;
|
2021-11-26 14:54:38 +00:00
|
|
|
use bucket::*;
|
|
|
|
pub use bucket_entry::*;
|
2021-12-14 14:48:33 +00:00
|
|
|
pub use debug::*;
|
2021-11-26 14:54:38 +00:00
|
|
|
pub use find_nodes::*;
|
2021-11-22 16:28:30 +00:00
|
|
|
use futures_util::stream::{FuturesUnordered, StreamExt};
|
2021-11-26 14:54:38 +00:00
|
|
|
pub use node_ref::*;
|
|
|
|
pub use stats_accounting::*;
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
//////////////////////////////////////////////////////////////////////////
|
|
|
|
|
2022-05-28 14:07:57 +00:00
|
|
|
pub const BOOTSTRAP_TXT_VERSION: u8 = 0;
|
|
|
|
|
|
|
|
#[derive(Clone, Debug)]
|
|
|
|
pub struct BootstrapRecord {
|
|
|
|
min_version: u8,
|
|
|
|
max_version: u8,
|
|
|
|
dial_info_details: Vec<DialInfoDetail>,
|
|
|
|
}
|
|
|
|
pub type BootstrapRecordMap = BTreeMap<DHTKey, BootstrapRecord>;
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
#[derive(Debug, Copy, Clone, PartialEq, PartialOrd, Ord, Eq)]
|
|
|
|
pub enum RoutingDomain {
|
|
|
|
PublicInternet,
|
|
|
|
LocalNetwork,
|
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Debug, Default)]
|
|
|
|
pub struct RoutingDomainDetail {
|
|
|
|
dial_info_details: Vec<DialInfoDetail>,
|
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
struct RoutingTableInner {
|
|
|
|
network_manager: NetworkManager,
|
|
|
|
node_id: DHTKey,
|
|
|
|
node_id_secret: DHTKeySecret,
|
|
|
|
buckets: Vec<Bucket>,
|
2022-04-23 01:30:09 +00:00
|
|
|
public_internet_routing_domain: RoutingDomainDetail,
|
|
|
|
local_network_routing_domain: RoutingDomainDetail,
|
2021-11-22 16:28:30 +00:00
|
|
|
bucket_entry_count: usize,
|
2022-03-19 22:19:40 +00:00
|
|
|
|
2021-11-26 14:54:38 +00:00
|
|
|
// Transfer stats for this node
|
2022-03-19 22:19:40 +00:00
|
|
|
self_latency_stats_accounting: LatencyStatsAccounting,
|
|
|
|
self_transfer_stats_accounting: TransferStatsAccounting,
|
|
|
|
self_transfer_stats: TransferStatsDownUp,
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
2022-03-24 14:14:50 +00:00
|
|
|
#[derive(Clone, Debug, Default)]
|
|
|
|
pub struct RoutingTableHealth {
|
|
|
|
pub reliable_entry_count: usize,
|
|
|
|
pub unreliable_entry_count: usize,
|
|
|
|
pub dead_entry_count: usize,
|
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
struct RoutingTableUnlockedInner {
|
|
|
|
// Background processes
|
|
|
|
rolling_transfers_task: TickTask,
|
|
|
|
bootstrap_task: TickTask,
|
|
|
|
peer_minimum_refresh_task: TickTask,
|
|
|
|
ping_validator_task: TickTask,
|
2022-06-13 00:58:02 +00:00
|
|
|
node_info_update_single_future: MustJoinSingleFuture<()>,
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Clone)]
|
|
|
|
pub struct RoutingTable {
|
|
|
|
config: VeilidConfig,
|
|
|
|
inner: Arc<Mutex<RoutingTableInner>>,
|
|
|
|
unlocked_inner: Arc<RoutingTableUnlockedInner>,
|
|
|
|
}
|
|
|
|
|
|
|
|
impl RoutingTable {
|
|
|
|
fn new_inner(network_manager: NetworkManager) -> RoutingTableInner {
|
|
|
|
RoutingTableInner {
|
2021-11-26 14:54:38 +00:00
|
|
|
network_manager,
|
2021-11-22 16:28:30 +00:00
|
|
|
node_id: DHTKey::default(),
|
|
|
|
node_id_secret: DHTKeySecret::default(),
|
|
|
|
buckets: Vec::new(),
|
2022-04-23 01:30:09 +00:00
|
|
|
public_internet_routing_domain: RoutingDomainDetail::default(),
|
|
|
|
local_network_routing_domain: RoutingDomainDetail::default(),
|
2021-11-22 16:28:30 +00:00
|
|
|
bucket_entry_count: 0,
|
2022-03-19 22:19:40 +00:00
|
|
|
self_latency_stats_accounting: LatencyStatsAccounting::new(),
|
|
|
|
self_transfer_stats_accounting: TransferStatsAccounting::new(),
|
|
|
|
self_transfer_stats: TransferStatsDownUp::default(),
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
fn new_unlocked_inner(config: VeilidConfig) -> RoutingTableUnlockedInner {
|
|
|
|
let c = config.get();
|
|
|
|
RoutingTableUnlockedInner {
|
2021-11-26 14:54:38 +00:00
|
|
|
rolling_transfers_task: TickTask::new(ROLLING_TRANSFERS_INTERVAL_SECS),
|
2021-11-22 16:28:30 +00:00
|
|
|
bootstrap_task: TickTask::new(1),
|
2022-01-27 14:53:01 +00:00
|
|
|
peer_minimum_refresh_task: TickTask::new_ms(c.network.dht.min_peer_refresh_time_ms),
|
2021-11-22 16:28:30 +00:00
|
|
|
ping_validator_task: TickTask::new(1),
|
2022-06-13 00:58:02 +00:00
|
|
|
node_info_update_single_future: MustJoinSingleFuture::new(),
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
pub fn new(network_manager: NetworkManager) -> Self {
|
|
|
|
let config = network_manager.config();
|
|
|
|
let this = Self {
|
|
|
|
config: config.clone(),
|
|
|
|
inner: Arc::new(Mutex::new(Self::new_inner(network_manager))),
|
|
|
|
unlocked_inner: Arc::new(Self::new_unlocked_inner(config)),
|
|
|
|
};
|
|
|
|
// Set rolling transfers tick task
|
|
|
|
{
|
|
|
|
let this2 = this.clone();
|
|
|
|
this.unlocked_inner
|
|
|
|
.rolling_transfers_task
|
2022-06-13 00:58:02 +00:00
|
|
|
.set_routine(move |s, l, t| {
|
|
|
|
Box::pin(this2.clone().rolling_transfers_task_routine(s, l, t))
|
2021-11-22 16:28:30 +00:00
|
|
|
});
|
|
|
|
}
|
|
|
|
// Set bootstrap tick task
|
|
|
|
{
|
|
|
|
let this2 = this.clone();
|
|
|
|
this.unlocked_inner
|
|
|
|
.bootstrap_task
|
2022-06-13 00:58:02 +00:00
|
|
|
.set_routine(move |s, _l, _t| Box::pin(this2.clone().bootstrap_task_routine(s)));
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
// Set peer minimum refresh tick task
|
|
|
|
{
|
|
|
|
let this2 = this.clone();
|
|
|
|
this.unlocked_inner
|
|
|
|
.peer_minimum_refresh_task
|
2022-06-13 00:58:02 +00:00
|
|
|
.set_routine(move |s, _l, _t| {
|
|
|
|
Box::pin(this2.clone().peer_minimum_refresh_task_routine(s))
|
2021-11-22 16:28:30 +00:00
|
|
|
});
|
|
|
|
}
|
|
|
|
// Set ping validator tick task
|
|
|
|
{
|
|
|
|
let this2 = this.clone();
|
|
|
|
this.unlocked_inner
|
|
|
|
.ping_validator_task
|
2022-06-13 00:58:02 +00:00
|
|
|
.set_routine(move |s, l, t| {
|
|
|
|
Box::pin(this2.clone().ping_validator_task_routine(s, l, t))
|
|
|
|
});
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
this
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn network_manager(&self) -> NetworkManager {
|
|
|
|
self.inner.lock().network_manager.clone()
|
|
|
|
}
|
|
|
|
pub fn rpc_processor(&self) -> RPCProcessor {
|
|
|
|
self.network_manager().rpc_processor()
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn node_id(&self) -> DHTKey {
|
|
|
|
self.inner.lock().node_id
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn node_id_secret(&self) -> DHTKeySecret {
|
|
|
|
self.inner.lock().node_id_secret
|
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
fn with_routing_domain<F, R>(inner: &RoutingTableInner, domain: RoutingDomain, f: F) -> R
|
|
|
|
where
|
|
|
|
F: FnOnce(&RoutingDomainDetail) -> R,
|
|
|
|
{
|
|
|
|
match domain {
|
|
|
|
RoutingDomain::PublicInternet => f(&inner.public_internet_routing_domain),
|
|
|
|
RoutingDomain::LocalNetwork => f(&inner.local_network_routing_domain),
|
|
|
|
}
|
2021-12-24 23:02:53 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
fn with_routing_domain_mut<F, R>(
|
|
|
|
inner: &mut RoutingTableInner,
|
|
|
|
domain: RoutingDomain,
|
|
|
|
f: F,
|
|
|
|
) -> R
|
|
|
|
where
|
|
|
|
F: FnOnce(&mut RoutingDomainDetail) -> R,
|
|
|
|
{
|
|
|
|
match domain {
|
|
|
|
RoutingDomain::PublicInternet => f(&mut inner.public_internet_routing_domain),
|
|
|
|
RoutingDomain::LocalNetwork => f(&mut inner.local_network_routing_domain),
|
|
|
|
}
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn has_dial_info(&self, domain: RoutingDomain) -> bool {
|
2021-11-22 16:28:30 +00:00
|
|
|
let inner = self.inner.lock();
|
2022-04-23 01:30:09 +00:00
|
|
|
Self::with_routing_domain(&*inner, domain, |rd| !rd.dial_info_details.is_empty())
|
2021-12-24 01:34:52 +00:00
|
|
|
}
|
2022-04-16 15:18:54 +00:00
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn dial_info_details(&self, domain: RoutingDomain) -> Vec<DialInfoDetail> {
|
2021-11-22 16:28:30 +00:00
|
|
|
let inner = self.inner.lock();
|
2022-04-23 01:30:09 +00:00
|
|
|
Self::with_routing_domain(&*inner, domain, |rd| rd.dial_info_details.clone())
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn first_filtered_dial_info_detail(
|
2022-04-16 15:18:54 +00:00
|
|
|
&self,
|
2022-04-25 15:29:02 +00:00
|
|
|
domain: Option<RoutingDomain>,
|
2022-04-16 15:18:54 +00:00
|
|
|
filter: &DialInfoFilter,
|
|
|
|
) -> Option<DialInfoDetail> {
|
|
|
|
let inner = self.inner.lock();
|
2022-04-25 15:29:02 +00:00
|
|
|
// Prefer local network first if it isn't filtered out
|
|
|
|
if domain == None || domain == Some(RoutingDomain::LocalNetwork) {
|
|
|
|
Self::with_routing_domain(&*inner, RoutingDomain::LocalNetwork, |rd| {
|
|
|
|
for did in &rd.dial_info_details {
|
|
|
|
if did.matches_filter(filter) {
|
|
|
|
return Some(did.clone());
|
|
|
|
}
|
2022-04-23 01:30:09 +00:00
|
|
|
}
|
2022-04-25 15:29:02 +00:00
|
|
|
None
|
|
|
|
})
|
|
|
|
} else {
|
2022-04-23 01:30:09 +00:00
|
|
|
None
|
2022-04-25 15:29:02 +00:00
|
|
|
}
|
|
|
|
.or_else(|| {
|
|
|
|
if domain == None || domain == Some(RoutingDomain::PublicInternet) {
|
|
|
|
Self::with_routing_domain(&*inner, RoutingDomain::PublicInternet, |rd| {
|
|
|
|
for did in &rd.dial_info_details {
|
|
|
|
if did.matches_filter(filter) {
|
|
|
|
return Some(did.clone());
|
|
|
|
}
|
|
|
|
}
|
|
|
|
None
|
|
|
|
})
|
|
|
|
} else {
|
|
|
|
None
|
|
|
|
}
|
2022-04-23 01:30:09 +00:00
|
|
|
})
|
2022-04-16 15:18:54 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn all_filtered_dial_info_details(
|
2022-04-16 15:18:54 +00:00
|
|
|
&self,
|
2022-04-24 02:08:02 +00:00
|
|
|
domain: Option<RoutingDomain>,
|
2022-04-16 15:18:54 +00:00
|
|
|
filter: &DialInfoFilter,
|
|
|
|
) -> Vec<DialInfoDetail> {
|
|
|
|
let inner = self.inner.lock();
|
2022-04-24 02:08:02 +00:00
|
|
|
let mut ret = Vec::new();
|
|
|
|
|
2022-04-25 00:16:13 +00:00
|
|
|
if domain == None || domain == Some(RoutingDomain::LocalNetwork) {
|
|
|
|
Self::with_routing_domain(&*inner, RoutingDomain::LocalNetwork, |rd| {
|
2022-04-25 15:29:02 +00:00
|
|
|
for did in &rd.dial_info_details {
|
2022-04-24 02:08:02 +00:00
|
|
|
if did.matches_filter(filter) {
|
|
|
|
ret.push(did.clone());
|
|
|
|
}
|
2022-04-23 01:30:09 +00:00
|
|
|
}
|
2022-04-24 02:08:02 +00:00
|
|
|
});
|
|
|
|
}
|
|
|
|
if domain == None || domain == Some(RoutingDomain::PublicInternet) {
|
|
|
|
Self::with_routing_domain(&*inner, RoutingDomain::PublicInternet, |rd| {
|
2022-04-25 15:29:02 +00:00
|
|
|
for did in &rd.dial_info_details {
|
2022-04-24 02:08:02 +00:00
|
|
|
if did.matches_filter(filter) {
|
|
|
|
ret.push(did.clone());
|
|
|
|
}
|
|
|
|
}
|
|
|
|
});
|
|
|
|
}
|
|
|
|
ret.remove_duplicates();
|
|
|
|
ret
|
2022-04-16 15:18:54 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn register_dial_info(
|
2021-11-22 16:28:30 +00:00
|
|
|
&self,
|
2022-04-23 01:30:09 +00:00
|
|
|
domain: RoutingDomain,
|
2021-11-22 16:28:30 +00:00
|
|
|
dial_info: DialInfo,
|
2022-04-24 02:08:02 +00:00
|
|
|
class: DialInfoClass,
|
2022-04-26 13:16:48 +00:00
|
|
|
) -> Result<(), String> {
|
2022-05-26 00:56:13 +00:00
|
|
|
log_rtab!(debug
|
|
|
|
"Registering dial_info with:\n domain: {:?}\n dial_info: {:?}\n class: {:?}",
|
|
|
|
domain, dial_info, class
|
2022-04-26 13:16:48 +00:00
|
|
|
);
|
2022-04-16 15:18:54 +00:00
|
|
|
let enable_local_peer_scope = {
|
2022-04-17 17:28:39 +00:00
|
|
|
let config = self.network_manager().config();
|
|
|
|
let c = config.get();
|
2022-04-16 15:18:54 +00:00
|
|
|
c.network.enable_local_peer_scope
|
|
|
|
};
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
if !enable_local_peer_scope
|
|
|
|
&& matches!(domain, RoutingDomain::PublicInternet)
|
|
|
|
&& dial_info.is_local()
|
|
|
|
{
|
2022-04-26 13:16:48 +00:00
|
|
|
return Err("shouldn't be registering local addresses as public".to_owned())
|
|
|
|
.map_err(logthru_rtab!(error));
|
2022-04-17 17:28:39 +00:00
|
|
|
}
|
|
|
|
if !dial_info.is_valid() {
|
2022-04-26 13:16:48 +00:00
|
|
|
return Err(format!(
|
|
|
|
"shouldn't be registering invalid addresses: {:?}",
|
|
|
|
dial_info
|
|
|
|
))
|
|
|
|
.map_err(logthru_rtab!(error));
|
2022-04-17 17:28:39 +00:00
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
let mut inner = self.inner.lock();
|
2022-04-23 01:30:09 +00:00
|
|
|
Self::with_routing_domain_mut(&mut *inner, domain, |rd| {
|
|
|
|
rd.dial_info_details.push(DialInfoDetail {
|
|
|
|
dial_info: dial_info.clone(),
|
2022-04-24 02:08:02 +00:00
|
|
|
class,
|
2022-04-23 01:30:09 +00:00
|
|
|
});
|
2022-04-25 00:16:13 +00:00
|
|
|
rd.dial_info_details.sort();
|
2022-04-16 15:18:54 +00:00
|
|
|
});
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
let domain_str = match domain {
|
|
|
|
RoutingDomain::PublicInternet => "Public",
|
|
|
|
RoutingDomain::LocalNetwork => "Local",
|
|
|
|
};
|
2022-04-16 15:18:54 +00:00
|
|
|
info!(
|
2022-04-23 01:30:09 +00:00
|
|
|
"{} Dial Info: {}",
|
|
|
|
domain_str,
|
2022-04-16 15:18:54 +00:00
|
|
|
NodeDialInfo {
|
|
|
|
node_id: NodeId::new(inner.node_id),
|
|
|
|
dial_info
|
|
|
|
}
|
|
|
|
.to_string(),
|
|
|
|
);
|
2022-04-24 02:08:02 +00:00
|
|
|
debug!(" Class: {:?}", class);
|
2022-05-11 16:20:33 +00:00
|
|
|
|
|
|
|
// Public dial info changed, go through all nodes and reset their 'seen our node info' bit
|
|
|
|
if matches!(domain, RoutingDomain::PublicInternet) {
|
2022-05-24 21:13:52 +00:00
|
|
|
let cur_ts = intf::get_timestamp();
|
|
|
|
Self::with_entries(&mut *inner, cur_ts, BucketEntryState::Dead, |_, e| {
|
|
|
|
e.set_seen_our_node_info(false);
|
|
|
|
Option::<()>::None
|
|
|
|
});
|
2022-05-11 16:20:33 +00:00
|
|
|
}
|
|
|
|
|
2022-04-26 13:16:48 +00:00
|
|
|
Ok(())
|
2022-04-16 15:18:54 +00:00
|
|
|
}
|
|
|
|
|
2022-04-23 01:30:09 +00:00
|
|
|
pub fn clear_dial_info_details(&self, domain: RoutingDomain) {
|
2022-04-26 13:16:48 +00:00
|
|
|
trace!("clearing dial info domain: {:?}", domain);
|
|
|
|
|
2021-12-11 01:14:33 +00:00
|
|
|
let mut inner = self.inner.lock();
|
2022-04-23 01:30:09 +00:00
|
|
|
Self::with_routing_domain_mut(&mut *inner, domain, |rd| {
|
|
|
|
rd.dial_info_details.clear();
|
|
|
|
})
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
fn bucket_depth(index: usize) -> usize {
|
|
|
|
match index {
|
|
|
|
0 => 256,
|
|
|
|
1 => 128,
|
|
|
|
2 => 64,
|
|
|
|
3 => 32,
|
|
|
|
4 => 16,
|
|
|
|
5 => 8,
|
|
|
|
6 => 4,
|
|
|
|
7 => 4,
|
|
|
|
8 => 4,
|
|
|
|
9 => 4,
|
|
|
|
_ => 4,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
pub async fn init(&self) -> Result<(), String> {
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
// Size the buckets (one per bit)
|
|
|
|
inner.buckets.reserve(DHT_KEY_LENGTH * 8);
|
|
|
|
for _ in 0..DHT_KEY_LENGTH * 8 {
|
|
|
|
let bucket = Bucket::new(self.clone());
|
|
|
|
inner.buckets.push(bucket);
|
|
|
|
}
|
|
|
|
|
|
|
|
// make local copy of node id for easy access
|
|
|
|
let c = self.config.get();
|
|
|
|
inner.node_id = c.network.node_id;
|
|
|
|
inner.node_id_secret = c.network.node_id_secret;
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
pub async fn terminate(&self) {
|
2022-05-25 15:12:19 +00:00
|
|
|
// Cancel all tasks being ticked
|
2022-06-13 00:58:02 +00:00
|
|
|
if let Err(e) = self.unlocked_inner.rolling_transfers_task.stop().await {
|
|
|
|
error!("rolling_transfers_task not stopped: {}", e);
|
2022-05-25 15:12:19 +00:00
|
|
|
}
|
2022-06-13 00:58:02 +00:00
|
|
|
if let Err(e) = self.unlocked_inner.bootstrap_task.stop().await {
|
|
|
|
error!("bootstrap_task not stopped: {}", e);
|
2022-05-25 15:12:19 +00:00
|
|
|
}
|
2022-06-13 00:58:02 +00:00
|
|
|
if let Err(e) = self.unlocked_inner.peer_minimum_refresh_task.stop().await {
|
|
|
|
error!("peer_minimum_refresh_task not stopped: {}", e);
|
2022-05-25 15:12:19 +00:00
|
|
|
}
|
2022-06-13 00:58:02 +00:00
|
|
|
if let Err(e) = self.unlocked_inner.ping_validator_task.stop().await {
|
|
|
|
error!("ping_validator_task not stopped: {}", e);
|
2022-05-25 15:12:19 +00:00
|
|
|
}
|
|
|
|
if self
|
|
|
|
.unlocked_inner
|
|
|
|
.node_info_update_single_future
|
2022-06-13 00:58:02 +00:00
|
|
|
.join()
|
2022-05-25 15:12:19 +00:00
|
|
|
.await
|
|
|
|
.is_err()
|
|
|
|
{
|
2022-06-13 00:58:02 +00:00
|
|
|
error!("node_info_update_single_future not stopped");
|
2022-05-25 15:12:19 +00:00
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
*self.inner.lock() = Self::new_inner(self.network_manager());
|
|
|
|
}
|
|
|
|
|
2022-05-11 16:20:33 +00:00
|
|
|
// Inform routing table entries that our dial info has changed
|
2022-06-11 22:47:58 +00:00
|
|
|
pub async fn send_node_info_updates(&self) {
|
2022-05-11 16:20:33 +00:00
|
|
|
let this = self.clone();
|
2022-06-11 22:47:58 +00:00
|
|
|
// Run in background only once
|
|
|
|
let _ = self
|
|
|
|
.clone()
|
|
|
|
.unlocked_inner
|
|
|
|
.node_info_update_single_future
|
|
|
|
.single_spawn(async move {
|
|
|
|
// Only update if we actually have a valid network class
|
|
|
|
let netman = this.network_manager();
|
|
|
|
if matches!(
|
|
|
|
netman.get_network_class().unwrap_or(NetworkClass::Invalid),
|
|
|
|
NetworkClass::Invalid
|
|
|
|
) {
|
|
|
|
trace!(
|
|
|
|
"not sending node info update because our network class is not yet valid"
|
|
|
|
);
|
|
|
|
return;
|
|
|
|
}
|
2022-05-11 16:20:33 +00:00
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
// Get the list of refs to all nodes to update
|
|
|
|
let node_refs = {
|
|
|
|
let mut inner = this.inner.lock();
|
|
|
|
let mut node_refs = Vec::<NodeRef>::with_capacity(inner.bucket_entry_count);
|
|
|
|
let cur_ts = intf::get_timestamp();
|
|
|
|
Self::with_entries(
|
|
|
|
&mut *inner,
|
|
|
|
cur_ts,
|
|
|
|
BucketEntryState::Unreliable,
|
|
|
|
|k, e| {
|
2022-05-24 21:13:52 +00:00
|
|
|
// Only update nodes that haven't seen our node info yet
|
|
|
|
if !e.has_seen_our_node_info() {
|
2022-06-11 22:47:58 +00:00
|
|
|
node_refs.push(NodeRef::new(this.clone(), *k, e, None));
|
2022-05-11 16:20:33 +00:00
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
Option::<()>::None
|
2022-06-11 22:47:58 +00:00
|
|
|
},
|
|
|
|
);
|
|
|
|
node_refs
|
|
|
|
};
|
2022-05-11 16:20:33 +00:00
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
// Send the updates
|
|
|
|
log_rtab!("Sending node info updates to {} nodes", node_refs.len());
|
|
|
|
let mut unord = FuturesUnordered::new();
|
|
|
|
for nr in node_refs {
|
|
|
|
let rpc = this.rpc_processor();
|
|
|
|
unord.push(async move {
|
|
|
|
// Update the node
|
|
|
|
if let Err(e) = rpc
|
|
|
|
.rpc_call_node_info_update(Destination::Direct(nr.clone()), None)
|
|
|
|
.await
|
|
|
|
{
|
|
|
|
// Not fatal, but we should be able to see if this is happening
|
|
|
|
trace!("failed to send node info update to {:?}: {}", nr, e);
|
|
|
|
return;
|
|
|
|
}
|
2022-05-11 16:20:33 +00:00
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
// Mark the node as updated
|
|
|
|
nr.set_seen_our_node_info();
|
|
|
|
});
|
|
|
|
}
|
2022-05-11 16:20:33 +00:00
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
// Wait for futures to complete
|
|
|
|
while unord.next().await.is_some() {}
|
2022-05-11 16:20:33 +00:00
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
log_rtab!("Finished sending node updates");
|
|
|
|
})
|
|
|
|
.await;
|
2022-05-11 16:20:33 +00:00
|
|
|
}
|
|
|
|
|
2022-03-09 03:32:12 +00:00
|
|
|
// Attempt to empty the routing table
|
|
|
|
// should only be performed when there are no node_refs (detached)
|
|
|
|
pub fn purge(&self) {
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
log_rtab!(
|
|
|
|
"Starting routing table purge. Table currently has {} nodes",
|
|
|
|
inner.bucket_entry_count
|
|
|
|
);
|
|
|
|
for bucket in &mut inner.buckets {
|
|
|
|
bucket.kick(0);
|
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
log_rtab!(debug
|
|
|
|
"Routing table purge complete. Routing table now has {} nodes",
|
2022-03-09 03:32:12 +00:00
|
|
|
inner.bucket_entry_count
|
|
|
|
);
|
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
// Attempt to settle buckets and remove entries down to the desired number
|
|
|
|
// which may not be possible due extant NodeRefs
|
|
|
|
fn kick_bucket(inner: &mut RoutingTableInner, idx: usize) {
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let bucket_depth = Self::bucket_depth(idx);
|
|
|
|
|
|
|
|
if let Some(dead_node_ids) = bucket.kick(bucket_depth) {
|
|
|
|
// Remove counts
|
|
|
|
inner.bucket_entry_count -= dead_node_ids.len();
|
2022-05-24 21:13:52 +00:00
|
|
|
log_rtab!(debug "Routing table now has {} nodes", inner.bucket_entry_count);
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
// Now purge the routing table inner vectors
|
|
|
|
//let filter = |k: &DHTKey| dead_node_ids.contains(k);
|
|
|
|
//inner.closest_reliable_nodes.retain(filter);
|
|
|
|
//inner.fastest_reliable_nodes.retain(filter);
|
|
|
|
//inner.closest_nodes.retain(filter);
|
|
|
|
//inner.fastest_nodes.retain(filter);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
fn find_bucket_index(inner: &RoutingTableInner, node_id: DHTKey) -> usize {
|
|
|
|
distance(&node_id, &inner.node_id)
|
|
|
|
.first_nonzero_bit()
|
|
|
|
.unwrap()
|
|
|
|
}
|
|
|
|
|
2022-05-24 21:13:52 +00:00
|
|
|
fn get_entry_count(inner: &mut RoutingTableInner, min_state: BucketEntryState) -> usize {
|
|
|
|
let mut count = 0usize;
|
|
|
|
let cur_ts = intf::get_timestamp();
|
|
|
|
Self::with_entries(inner, cur_ts, min_state, |_, _| {
|
|
|
|
count += 1;
|
|
|
|
Option::<()>::None
|
|
|
|
});
|
|
|
|
count
|
|
|
|
}
|
|
|
|
|
|
|
|
fn with_entries<T, F: FnMut(&DHTKey, &mut BucketEntry) -> Option<T>>(
|
|
|
|
inner: &mut RoutingTableInner,
|
|
|
|
cur_ts: u64,
|
|
|
|
min_state: BucketEntryState,
|
|
|
|
mut f: F,
|
|
|
|
) -> Option<T> {
|
|
|
|
for bucket in &mut inner.buckets {
|
|
|
|
for entry in bucket.entries_mut() {
|
|
|
|
if entry.1.state(cur_ts) >= min_state {
|
|
|
|
if let Some(out) = f(entry.0, entry.1) {
|
|
|
|
return Some(out);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
None
|
|
|
|
}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
fn drop_node_ref(&self, node_id: DHTKey) {
|
|
|
|
// Reduce ref count on entry
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
let idx = Self::find_bucket_index(&*inner, node_id);
|
|
|
|
let new_ref_count = {
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let entry = bucket.entry_mut(&node_id).unwrap();
|
|
|
|
entry.ref_count -= 1;
|
|
|
|
entry.ref_count
|
|
|
|
};
|
|
|
|
|
|
|
|
// If this entry could possibly go away, kick the bucket
|
|
|
|
if new_ref_count == 0 {
|
|
|
|
// it important to do this in the same inner lock as the ref count decrease
|
|
|
|
Self::kick_bucket(&mut *inner, idx);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-05-03 20:43:15 +00:00
|
|
|
// Create a node reference, possibly creating a bucket entry
|
|
|
|
// the 'update_func' closure is called on the node, and, if created,
|
|
|
|
// in a locked fashion as to ensure the bucket entry state is always valid
|
|
|
|
pub fn create_node_ref<F>(&self, node_id: DHTKey, update_func: F) -> Result<NodeRef, String>
|
|
|
|
where
|
|
|
|
F: FnOnce(&mut BucketEntry),
|
|
|
|
{
|
2021-11-22 16:28:30 +00:00
|
|
|
// Ensure someone isn't trying register this node itself
|
|
|
|
if node_id == self.node_id() {
|
2021-12-18 00:18:25 +00:00
|
|
|
return Err("can't register own node".to_owned()).map_err(logthru_rtab!(error));
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
2022-05-03 20:43:15 +00:00
|
|
|
// Lock this entire operation
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
|
|
|
|
// Look up existing entry
|
|
|
|
let idx = Self::find_bucket_index(&*inner, node_id);
|
|
|
|
let noderef = {
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let entry = bucket.entry_mut(&node_id);
|
|
|
|
entry.map(|e| NodeRef::new(self.clone(), node_id, e, None))
|
|
|
|
};
|
|
|
|
|
|
|
|
// If one doesn't exist, insert into bucket, possibly evicting a bucket member
|
|
|
|
let noderef = match noderef {
|
2021-11-22 16:28:30 +00:00
|
|
|
None => {
|
|
|
|
// Make new entry
|
2022-05-03 20:43:15 +00:00
|
|
|
inner.bucket_entry_count += 1;
|
2022-05-24 21:13:52 +00:00
|
|
|
let cnt = inner.bucket_entry_count;
|
|
|
|
log_rtab!(debug "Routing table now has {} nodes, {} live", cnt, Self::get_entry_count(&mut *inner, BucketEntryState::Unreliable));
|
2022-05-03 20:43:15 +00:00
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let nr = bucket.add_entry(node_id);
|
|
|
|
|
|
|
|
// Update the entry
|
|
|
|
let entry = bucket.entry_mut(&node_id);
|
|
|
|
update_func(entry.unwrap());
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
// Kick the bucket
|
|
|
|
// It is important to do this in the same inner lock as the add_entry
|
|
|
|
Self::kick_bucket(&mut *inner, idx);
|
|
|
|
|
|
|
|
nr
|
|
|
|
}
|
2022-05-03 20:43:15 +00:00
|
|
|
Some(nr) => {
|
|
|
|
// Update the entry
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let entry = bucket.entry_mut(&node_id);
|
|
|
|
update_func(entry.unwrap());
|
|
|
|
|
|
|
|
nr
|
|
|
|
}
|
2021-11-22 16:28:30 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
Ok(noderef)
|
|
|
|
}
|
|
|
|
|
|
|
|
pub fn lookup_node_ref(&self, node_id: DHTKey) -> Option<NodeRef> {
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
let idx = Self::find_bucket_index(&*inner, node_id);
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
2021-11-26 14:54:38 +00:00
|
|
|
bucket
|
|
|
|
.entry_mut(&node_id)
|
2022-04-19 15:23:44 +00:00
|
|
|
.map(|e| NodeRef::new(self.clone(), node_id, e, None))
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Shortcut function to add a node to our routing table if it doesn't exist
|
|
|
|
// and add the dial info we have for it, since that's pretty common
|
2022-05-11 01:49:42 +00:00
|
|
|
pub fn register_node_with_signed_node_info(
|
2021-11-22 16:28:30 +00:00
|
|
|
&self,
|
|
|
|
node_id: DHTKey,
|
2022-05-11 01:49:42 +00:00
|
|
|
signed_node_info: SignedNodeInfo,
|
2021-11-22 16:28:30 +00:00
|
|
|
) -> Result<NodeRef, String> {
|
2022-05-11 13:37:54 +00:00
|
|
|
// validate signed node info is not something malicious
|
|
|
|
if node_id == self.node_id() {
|
|
|
|
return Err("can't register own node id in routing table".to_owned());
|
|
|
|
}
|
|
|
|
if let Some(rpi) = &signed_node_info.node_info.relay_peer_info {
|
|
|
|
if rpi.node_id.key == node_id {
|
|
|
|
return Err("node can not be its own relay".to_owned());
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-05-03 20:43:15 +00:00
|
|
|
let nr = self.create_node_ref(node_id, |e| {
|
2022-05-11 01:49:42 +00:00
|
|
|
e.update_node_info(signed_node_info);
|
2021-11-22 16:28:30 +00:00
|
|
|
})?;
|
|
|
|
|
|
|
|
Ok(nr)
|
|
|
|
}
|
|
|
|
|
|
|
|
// Shortcut function to add a node to our routing table if it doesn't exist
|
|
|
|
// and add the last peer address we have for it, since that's pretty common
|
|
|
|
pub fn register_node_with_existing_connection(
|
|
|
|
&self,
|
|
|
|
node_id: DHTKey,
|
|
|
|
descriptor: ConnectionDescriptor,
|
|
|
|
timestamp: u64,
|
|
|
|
) -> Result<NodeRef, String> {
|
2022-05-03 20:43:15 +00:00
|
|
|
let nr = self.create_node_ref(node_id, |e| {
|
2021-11-22 16:28:30 +00:00
|
|
|
// set the most recent node address for connection finding and udp replies
|
|
|
|
e.set_last_connection(descriptor, timestamp);
|
2022-05-03 20:43:15 +00:00
|
|
|
})?;
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
Ok(nr)
|
|
|
|
}
|
|
|
|
|
2022-05-03 20:43:15 +00:00
|
|
|
fn operate_on_bucket_entry_locked<T, F>(
|
|
|
|
inner: &mut RoutingTableInner,
|
|
|
|
node_id: DHTKey,
|
|
|
|
f: F,
|
|
|
|
) -> T
|
2021-11-22 16:28:30 +00:00
|
|
|
where
|
|
|
|
F: FnOnce(&mut BucketEntry) -> T,
|
|
|
|
{
|
|
|
|
let idx = Self::find_bucket_index(&*inner, node_id);
|
|
|
|
let bucket = &mut inner.buckets[idx];
|
|
|
|
let entry = bucket.entry_mut(&node_id).unwrap();
|
|
|
|
f(entry)
|
|
|
|
}
|
|
|
|
|
2022-05-03 20:43:15 +00:00
|
|
|
fn operate_on_bucket_entry<T, F>(&self, node_id: DHTKey, f: F) -> T
|
|
|
|
where
|
|
|
|
F: FnOnce(&mut BucketEntry) -> T,
|
|
|
|
{
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
Self::operate_on_bucket_entry_locked(&mut *inner, node_id, f)
|
|
|
|
}
|
|
|
|
|
2022-04-07 13:55:09 +00:00
|
|
|
pub fn find_inbound_relay(&self, cur_ts: u64) -> Option<NodeRef> {
|
|
|
|
let mut inner = self.inner.lock();
|
2022-05-26 00:56:13 +00:00
|
|
|
let inner = &mut *inner;
|
|
|
|
let mut best_inbound_relay: Option<(&DHTKey, &mut BucketEntry)> = None;
|
2022-04-07 13:55:09 +00:00
|
|
|
|
|
|
|
// Iterate all known nodes for candidates
|
2022-05-26 00:56:13 +00:00
|
|
|
for bucket in &mut inner.buckets {
|
|
|
|
for (k, e) in bucket.entries_mut() {
|
|
|
|
if e.state(cur_ts) >= BucketEntryState::Unreliable {
|
|
|
|
// Ensure this node is not on our local network
|
|
|
|
if !e
|
|
|
|
.local_node_info()
|
|
|
|
.map(|l| l.has_dial_info())
|
|
|
|
.unwrap_or(false)
|
|
|
|
{
|
|
|
|
// Ensure we have the node's status
|
|
|
|
if let Some(node_status) = &e.peer_stats().status {
|
|
|
|
// Ensure the node will relay
|
|
|
|
if node_status.will_relay {
|
|
|
|
// Compare against previous candidate
|
|
|
|
if let Some(best_inbound_relay) = best_inbound_relay.as_mut() {
|
|
|
|
// Less is faster
|
|
|
|
if BucketEntry::cmp_fastest_reliable(
|
|
|
|
cur_ts,
|
|
|
|
e,
|
|
|
|
best_inbound_relay.1,
|
|
|
|
) == std::cmp::Ordering::Less
|
|
|
|
{
|
|
|
|
*best_inbound_relay = (k, e);
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
// Always store the first candidate
|
|
|
|
best_inbound_relay = Some((k, e));
|
|
|
|
}
|
2022-04-07 13:55:09 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2022-05-26 00:56:13 +00:00
|
|
|
}
|
|
|
|
// Return the best inbound relay noderef
|
|
|
|
best_inbound_relay.map(|(k, e)| NodeRef::new(self.clone(), *k, e, None))
|
2022-04-07 13:55:09 +00:00
|
|
|
}
|
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), ret, err)]
|
2022-05-11 01:49:42 +00:00
|
|
|
pub fn register_find_node_answer(&self, fna: FindNodeAnswer) -> Result<Vec<NodeRef>, String> {
|
2021-11-22 16:28:30 +00:00
|
|
|
let node_id = self.node_id();
|
2022-05-11 01:49:42 +00:00
|
|
|
|
|
|
|
// register nodes we'd found
|
|
|
|
let mut out = Vec::<NodeRef>::with_capacity(fna.peers.len());
|
|
|
|
for p in fna.peers {
|
|
|
|
// if our own node if is in the list then ignore it, as we don't add ourselves to our own routing table
|
|
|
|
if p.node_id.key == node_id {
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
// register the node if it's new
|
|
|
|
let nr = self
|
|
|
|
.register_node_with_signed_node_info(p.node_id.key, p.signed_node_info.clone())
|
|
|
|
.map_err(map_to_string)
|
|
|
|
.map_err(logthru_rtab!(
|
|
|
|
"couldn't register node {} at {:?}",
|
|
|
|
p.node_id.key,
|
|
|
|
&p.signed_node_info
|
|
|
|
))?;
|
|
|
|
out.push(nr);
|
|
|
|
}
|
|
|
|
Ok(out)
|
|
|
|
}
|
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), ret, err)]
|
2022-05-11 01:49:42 +00:00
|
|
|
pub async fn find_node(
|
|
|
|
&self,
|
|
|
|
node_ref: NodeRef,
|
|
|
|
node_id: DHTKey,
|
|
|
|
) -> Result<Vec<NodeRef>, String> {
|
2021-11-22 16:28:30 +00:00
|
|
|
let rpc_processor = self.rpc_processor();
|
|
|
|
|
2021-12-18 00:18:25 +00:00
|
|
|
let res = rpc_processor
|
2022-03-25 02:07:55 +00:00
|
|
|
.clone()
|
2021-11-22 16:28:30 +00:00
|
|
|
.rpc_call_find_node(
|
|
|
|
Destination::Direct(node_ref.clone()),
|
|
|
|
node_id,
|
|
|
|
None,
|
2022-04-25 00:16:13 +00:00
|
|
|
rpc_processor.make_respond_to_sender(node_ref.clone()),
|
2021-11-22 16:28:30 +00:00
|
|
|
)
|
|
|
|
.await
|
2021-12-18 00:18:25 +00:00
|
|
|
.map_err(map_to_string)
|
|
|
|
.map_err(logthru_rtab!())?;
|
2021-11-22 16:28:30 +00:00
|
|
|
|
|
|
|
// register nodes we'd found
|
2022-03-27 01:25:24 +00:00
|
|
|
self.register_find_node_answer(res)
|
|
|
|
}
|
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), ret, err)]
|
2022-05-11 01:49:42 +00:00
|
|
|
pub async fn find_self(&self, node_ref: NodeRef) -> Result<Vec<NodeRef>, String> {
|
2022-03-27 01:25:24 +00:00
|
|
|
let node_id = self.node_id();
|
2022-05-11 01:49:42 +00:00
|
|
|
self.find_node(node_ref, node_id).await
|
|
|
|
}
|
2022-03-27 01:25:24 +00:00
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), ret, err)]
|
2022-05-11 01:49:42 +00:00
|
|
|
pub async fn find_target(&self, node_ref: NodeRef) -> Result<Vec<NodeRef>, String> {
|
|
|
|
let node_id = node_ref.node_id();
|
|
|
|
self.find_node(node_ref, node_id).await
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self))]
|
2021-11-22 16:28:30 +00:00
|
|
|
pub async fn reverse_find_node(&self, node_ref: NodeRef, wide: bool) {
|
|
|
|
// Ask bootstrap node to 'find' our own node so we can get some more nodes near ourselves
|
|
|
|
// and then contact those nodes to inform -them- that we exist
|
|
|
|
|
|
|
|
// Ask bootstrap server for nodes closest to our own node
|
|
|
|
let closest_nodes = match self.find_self(node_ref.clone()).await {
|
|
|
|
Err(e) => {
|
2021-12-18 00:18:25 +00:00
|
|
|
log_rtab!(error
|
2021-11-22 16:28:30 +00:00
|
|
|
"reverse_find_node: find_self failed for {:?}: {}",
|
|
|
|
&node_ref, e
|
|
|
|
);
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
Ok(v) => v,
|
|
|
|
};
|
|
|
|
|
|
|
|
// Ask each node near us to find us as well
|
|
|
|
if wide {
|
|
|
|
for closest_nr in closest_nodes {
|
|
|
|
match self.find_self(closest_nr.clone()).await {
|
|
|
|
Err(e) => {
|
2021-12-18 00:18:25 +00:00
|
|
|
log_rtab!(error
|
2021-11-22 16:28:30 +00:00
|
|
|
"reverse_find_node: closest node find_self failed for {:?}: {}",
|
|
|
|
&closest_nr, e
|
|
|
|
);
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
Ok(v) => v,
|
|
|
|
};
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-05-24 21:13:52 +00:00
|
|
|
// Bootstrap lookup process
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), ret, err)]
|
2022-05-28 14:07:57 +00:00
|
|
|
async fn resolve_bootstrap(
|
|
|
|
&self,
|
|
|
|
bootstrap: Vec<String>,
|
|
|
|
) -> Result<BootstrapRecordMap, String> {
|
2022-05-24 21:13:52 +00:00
|
|
|
// Resolve from bootstrap root to bootstrap hostnames
|
|
|
|
let mut bsnames = Vec::<String>::new();
|
2022-05-16 15:52:48 +00:00
|
|
|
for bh in bootstrap {
|
2022-05-24 21:13:52 +00:00
|
|
|
// Get TXT record for bootstrap (bootstrap.veilid.net, or similar)
|
|
|
|
let records = intf::txt_lookup(&bh).await?;
|
|
|
|
for record in records {
|
|
|
|
// Split the bootstrap name record by commas
|
|
|
|
for rec in record.split(',') {
|
|
|
|
let rec = rec.trim();
|
|
|
|
// If the name specified is fully qualified, go with it
|
|
|
|
let bsname = if rec.ends_with('.') {
|
|
|
|
rec.to_string()
|
|
|
|
}
|
|
|
|
// If the name is not fully qualified, prepend it to the bootstrap name
|
|
|
|
else {
|
|
|
|
format!("{}.{}", rec, bh)
|
|
|
|
};
|
|
|
|
|
|
|
|
// Add to the list of bootstrap name to look up
|
|
|
|
bsnames.push(bsname);
|
|
|
|
}
|
|
|
|
}
|
2022-05-16 15:52:48 +00:00
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
|
|
|
|
// Get bootstrap nodes from hostnames concurrently
|
|
|
|
let mut unord = FuturesUnordered::new();
|
|
|
|
for bsname in bsnames {
|
|
|
|
unord.push(async move {
|
|
|
|
// look up boostrap node txt records
|
|
|
|
let bsnirecords = match intf::txt_lookup(&bsname).await {
|
|
|
|
Err(e) => {
|
|
|
|
warn!("bootstrap node txt lookup failed for {}: {}", bsname, e);
|
|
|
|
return None;
|
|
|
|
}
|
|
|
|
Ok(v) => v,
|
|
|
|
};
|
2022-05-28 14:07:57 +00:00
|
|
|
// for each record resolve into key/bootstraprecord pairs
|
|
|
|
let mut bootstrap_records: Vec<(DHTKey, BootstrapRecord)> = Vec::new();
|
2022-05-24 21:13:52 +00:00
|
|
|
for bsnirecord in bsnirecords {
|
2022-05-28 14:07:57 +00:00
|
|
|
// Bootstrap TXT Record Format Version 0:
|
|
|
|
// txt_version,min_version,max_version,nodeid,hostname,dialinfoshort*
|
|
|
|
//
|
|
|
|
// Split bootstrap node record by commas. Example:
|
|
|
|
// 0,0,0,7lxDEabK_qgjbe38RtBa3IZLrud84P6NhGP-pRTZzdQ,bootstrap-dev-alpha.veilid.net,T5150,U5150,W5150/ws
|
|
|
|
let records: Vec<String> = bsnirecord
|
|
|
|
.trim()
|
|
|
|
.split(',')
|
|
|
|
.map(|x| x.trim().to_owned())
|
|
|
|
.collect();
|
|
|
|
if records.len() < 6 {
|
|
|
|
warn!("invalid number of fields in bootstrap txt record");
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Bootstrap TXT record version
|
|
|
|
let txt_version: u8 = match records[0].parse::<u8>() {
|
|
|
|
Ok(v) => v,
|
|
|
|
Err(e) => {
|
|
|
|
warn!(
|
|
|
|
"invalid txt_version specified in bootstrap node txt record: {}",
|
|
|
|
e
|
|
|
|
);
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
};
|
|
|
|
if txt_version != BOOTSTRAP_TXT_VERSION {
|
|
|
|
warn!("unsupported bootstrap txt record version");
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Min/Max wire protocol version
|
|
|
|
let min_version: u8 = match records[1].parse::<u8>() {
|
|
|
|
Ok(v) => v,
|
|
|
|
Err(e) => {
|
|
|
|
warn!(
|
|
|
|
"invalid min_version specified in bootstrap node txt record: {}",
|
|
|
|
e
|
|
|
|
);
|
2022-05-24 21:13:52 +00:00
|
|
|
continue;
|
|
|
|
}
|
|
|
|
};
|
2022-05-28 14:07:57 +00:00
|
|
|
let max_version: u8 = match records[2].parse::<u8>() {
|
|
|
|
Ok(v) => v,
|
|
|
|
Err(e) => {
|
|
|
|
warn!(
|
|
|
|
"invalid max_version specified in bootstrap node txt record: {}",
|
|
|
|
e
|
|
|
|
);
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
};
|
|
|
|
|
|
|
|
// Node Id
|
|
|
|
let node_id_str = &records[3];
|
2022-05-24 21:13:52 +00:00
|
|
|
let node_id_key = match DHTKey::try_decode(node_id_str) {
|
|
|
|
Ok(v) => v,
|
|
|
|
Err(e) => {
|
|
|
|
warn!(
|
|
|
|
"Invalid node id in bootstrap node record {}: {}",
|
|
|
|
node_id_str, e
|
|
|
|
);
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
};
|
|
|
|
|
2022-05-28 14:07:57 +00:00
|
|
|
// Hostname
|
|
|
|
let hostname_str = &records[4];
|
|
|
|
|
2022-05-24 21:13:52 +00:00
|
|
|
// If this is our own node id, then we skip it for bootstrap, in case we are a bootstrap node
|
|
|
|
if self.node_id() == node_id_key {
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
|
|
|
|
// Resolve each record and store in node dial infos list
|
2022-05-28 14:07:57 +00:00
|
|
|
let mut bootstrap_record = BootstrapRecord {
|
|
|
|
min_version,
|
|
|
|
max_version,
|
|
|
|
dial_info_details: Vec::new(),
|
|
|
|
};
|
|
|
|
for rec in &records[5..] {
|
2022-05-24 21:13:52 +00:00
|
|
|
let rec = rec.trim();
|
2022-05-28 14:07:57 +00:00
|
|
|
let dial_infos = match DialInfo::try_vec_from_short(rec, hostname_str) {
|
2022-05-24 21:13:52 +00:00
|
|
|
Ok(dis) => dis,
|
|
|
|
Err(e) => {
|
|
|
|
warn!("Couldn't resolve bootstrap node dial info {}: {}", rec, e);
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
};
|
|
|
|
|
2022-05-28 14:07:57 +00:00
|
|
|
for di in dial_infos {
|
|
|
|
bootstrap_record.dial_info_details.push(DialInfoDetail {
|
|
|
|
dial_info: di,
|
|
|
|
class: DialInfoClass::Direct,
|
|
|
|
});
|
2022-05-24 21:13:52 +00:00
|
|
|
}
|
|
|
|
}
|
2022-05-28 14:07:57 +00:00
|
|
|
bootstrap_records.push((node_id_key, bootstrap_record));
|
2022-05-24 21:13:52 +00:00
|
|
|
}
|
2022-05-28 14:07:57 +00:00
|
|
|
Some(bootstrap_records)
|
2022-05-24 21:13:52 +00:00
|
|
|
});
|
|
|
|
}
|
2022-05-28 14:07:57 +00:00
|
|
|
|
|
|
|
let mut bsmap = BootstrapRecordMap::new();
|
|
|
|
while let Some(bootstrap_records) = unord.next().await {
|
|
|
|
if let Some(bootstrap_records) = bootstrap_records {
|
|
|
|
for (bskey, mut bsrec) in bootstrap_records {
|
|
|
|
let rec = bsmap.entry(bskey).or_insert_with(|| BootstrapRecord {
|
|
|
|
min_version: bsrec.min_version,
|
|
|
|
max_version: bsrec.max_version,
|
|
|
|
dial_info_details: Vec::new(),
|
|
|
|
});
|
|
|
|
rec.dial_info_details.append(&mut bsrec.dial_info_details);
|
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-05-28 14:07:57 +00:00
|
|
|
Ok(bsmap)
|
2022-05-16 15:52:48 +00:00
|
|
|
}
|
|
|
|
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), err)]
|
2022-06-13 00:58:02 +00:00
|
|
|
async fn bootstrap_task_routine(self, stop_token: StopToken) -> Result<(), String> {
|
2022-05-16 15:52:48 +00:00
|
|
|
let (bootstrap, bootstrap_nodes) = {
|
2021-11-22 16:28:30 +00:00
|
|
|
let c = self.config.get();
|
2022-05-16 15:52:48 +00:00
|
|
|
(
|
|
|
|
c.network.bootstrap.clone(),
|
|
|
|
c.network.bootstrap_nodes.clone(),
|
|
|
|
)
|
2021-11-22 16:28:30 +00:00
|
|
|
};
|
|
|
|
|
2022-05-26 00:56:13 +00:00
|
|
|
log_rtab!(debug "--- bootstrap_task");
|
2021-12-09 21:00:47 +00:00
|
|
|
|
2022-05-16 15:52:48 +00:00
|
|
|
// If we aren't specifying a bootstrap node list explicitly, then pull from the bootstrap server(s)
|
2022-05-28 14:07:57 +00:00
|
|
|
|
|
|
|
let bsmap: BootstrapRecordMap = if !bootstrap_nodes.is_empty() {
|
|
|
|
let mut bsmap = BootstrapRecordMap::new();
|
|
|
|
let mut bootstrap_node_dial_infos = Vec::new();
|
2022-05-24 21:13:52 +00:00
|
|
|
for b in bootstrap_nodes {
|
|
|
|
let ndis = NodeDialInfo::from_str(b.as_str())
|
|
|
|
.map_err(map_to_string)
|
|
|
|
.map_err(logthru_rtab!(
|
|
|
|
"Invalid node dial info in bootstrap entry: {}",
|
|
|
|
b
|
|
|
|
))?;
|
2022-05-28 14:07:57 +00:00
|
|
|
bootstrap_node_dial_infos.push(ndis);
|
2022-05-24 21:13:52 +00:00
|
|
|
}
|
2022-05-28 14:07:57 +00:00
|
|
|
for ndi in bootstrap_node_dial_infos {
|
|
|
|
let node_id = ndi.node_id.key;
|
|
|
|
bsmap
|
|
|
|
.entry(node_id)
|
|
|
|
.or_insert_with(|| BootstrapRecord {
|
|
|
|
min_version: MIN_VERSION,
|
|
|
|
max_version: MAX_VERSION,
|
|
|
|
dial_info_details: Vec::new(),
|
|
|
|
})
|
|
|
|
.dial_info_details
|
|
|
|
.push(DialInfoDetail {
|
|
|
|
dial_info: ndi.dial_info,
|
|
|
|
class: DialInfoClass::Direct, // Bootstraps are always directly reachable
|
|
|
|
});
|
|
|
|
}
|
|
|
|
bsmap
|
2022-05-16 15:52:48 +00:00
|
|
|
} else {
|
|
|
|
// Resolve bootstrap servers and recurse their TXT entries
|
|
|
|
self.resolve_bootstrap(bootstrap).await?
|
|
|
|
};
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
// Map all bootstrap entries to a single key with multiple dialinfo
|
|
|
|
|
|
|
|
// Run all bootstrap operations concurrently
|
|
|
|
let mut unord = FuturesUnordered::new();
|
2022-05-31 23:54:52 +00:00
|
|
|
for (k, mut v) in bsmap {
|
|
|
|
// Sort dial info so we get the preferred order correct
|
|
|
|
v.dial_info_details.sort();
|
|
|
|
|
2022-05-26 00:56:13 +00:00
|
|
|
log_rtab!("--- bootstrapping {} with {:?}", k.encode(), &v);
|
2022-05-11 01:49:42 +00:00
|
|
|
|
|
|
|
// Make invalid signed node info (no signature)
|
2021-12-18 00:18:25 +00:00
|
|
|
let nr = self
|
2022-05-11 01:49:42 +00:00
|
|
|
.register_node_with_signed_node_info(
|
2022-04-08 14:17:09 +00:00
|
|
|
k,
|
2022-05-11 01:49:42 +00:00
|
|
|
SignedNodeInfo::with_no_signature(NodeInfo {
|
2022-04-25 00:16:13 +00:00
|
|
|
network_class: NetworkClass::InboundCapable, // Bootstraps are always inbound capable
|
2022-04-23 01:30:09 +00:00
|
|
|
outbound_protocols: ProtocolSet::empty(), // Bootstraps do not participate in relaying and will not make outbound requests
|
2022-05-28 14:07:57 +00:00
|
|
|
min_version: v.min_version, // Minimum protocol version specified in txt record
|
|
|
|
max_version: v.max_version, // Maximum protocol version specified in txt record
|
|
|
|
dial_info_detail_list: v.dial_info_details, // Dial info is as specified in the bootstrap list
|
|
|
|
relay_peer_info: None, // Bootstraps never require a relay themselves
|
2022-05-11 01:49:42 +00:00
|
|
|
}),
|
2022-04-08 14:17:09 +00:00
|
|
|
)
|
2022-05-26 00:56:13 +00:00
|
|
|
.map_err(logthru_rtab!(error "Couldn't add bootstrap node: {}", k))?;
|
2022-05-11 01:49:42 +00:00
|
|
|
|
|
|
|
// Add this our futures to process in parallel
|
|
|
|
let this = self.clone();
|
|
|
|
unord.push(async move {
|
|
|
|
// Need VALID signed peer info, so ask bootstrap to find_node of itself
|
|
|
|
// which will ensure it has the bootstrap's signed peer info as part of the response
|
|
|
|
let _ = this.find_target(nr.clone()).await;
|
|
|
|
|
|
|
|
// Ensure we got the signed peer info
|
|
|
|
if !nr.operate(|e| e.has_valid_signed_node_info()) {
|
2022-05-26 00:56:13 +00:00
|
|
|
log_rtab!(warn
|
2022-05-11 01:49:42 +00:00
|
|
|
"bootstrap at {:?} did not return valid signed node info",
|
|
|
|
nr
|
|
|
|
);
|
2022-05-24 21:13:52 +00:00
|
|
|
// If this node info is invalid, it will time out after being unpingable
|
2022-05-11 01:49:42 +00:00
|
|
|
} else {
|
|
|
|
// otherwise this bootstrap is valid, lets ask it to find ourselves now
|
|
|
|
this.reverse_find_node(nr, true).await
|
|
|
|
}
|
|
|
|
});
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
|
|
|
|
// Wait for all bootstrap operations to complete before we complete the singlefuture
|
2021-11-22 16:28:30 +00:00
|
|
|
while unord.next().await.is_some() {}
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
///////////////////////////////////////////////////////////
|
|
|
|
/// Peer ping validation
|
|
|
|
|
|
|
|
// Ask our remaining peers to give us more peers before we go
|
|
|
|
// back to the bootstrap servers to keep us from bothering them too much
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), err)]
|
2022-06-13 00:58:02 +00:00
|
|
|
async fn peer_minimum_refresh_task_routine(self, stop_token: StopToken) -> Result<(), String> {
|
2022-05-24 21:13:52 +00:00
|
|
|
// get list of all peers we know about, even the unreliable ones, and ask them to find nodes close to our node too
|
2021-11-22 16:28:30 +00:00
|
|
|
let noderefs = {
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
let mut noderefs = Vec::<NodeRef>::with_capacity(inner.bucket_entry_count);
|
2022-05-24 21:13:52 +00:00
|
|
|
let cur_ts = intf::get_timestamp();
|
|
|
|
Self::with_entries(
|
|
|
|
&mut *inner,
|
|
|
|
cur_ts,
|
|
|
|
BucketEntryState::Unreliable,
|
|
|
|
|k, entry| {
|
|
|
|
noderefs.push(NodeRef::new(self.clone(), *k, entry, None));
|
|
|
|
Option::<()>::None
|
|
|
|
},
|
|
|
|
);
|
2021-11-22 16:28:30 +00:00
|
|
|
noderefs
|
|
|
|
};
|
|
|
|
|
|
|
|
// do peer minimum search concurrently
|
|
|
|
let mut unord = FuturesUnordered::new();
|
|
|
|
for nr in noderefs {
|
2022-05-26 00:56:13 +00:00
|
|
|
log_rtab!("--- peer minimum search with {:?}", nr);
|
2021-11-22 16:28:30 +00:00
|
|
|
unord.push(self.reverse_find_node(nr, false));
|
|
|
|
}
|
|
|
|
while unord.next().await.is_some() {}
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
// Ping each node in the routing table if they need to be pinged
|
|
|
|
// to determine their reliability
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), err)]
|
2022-06-13 00:58:02 +00:00
|
|
|
async fn ping_validator_task_routine(
|
|
|
|
self,
|
|
|
|
stop_token: StopToken,
|
|
|
|
_last_ts: u64,
|
|
|
|
cur_ts: u64,
|
|
|
|
) -> Result<(), String> {
|
2022-05-24 21:13:52 +00:00
|
|
|
// log_rtab!("--- ping_validator task");
|
2022-05-01 19:33:14 +00:00
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
let rpc = self.rpc_processor();
|
2022-05-01 19:33:14 +00:00
|
|
|
let netman = self.network_manager();
|
|
|
|
let relay_node_id = netman.relay_node().map(|nr| nr.node_id());
|
|
|
|
|
2022-06-11 22:47:58 +00:00
|
|
|
let mut unord = FuturesUnordered::new();
|
|
|
|
{
|
|
|
|
let mut inner = self.inner.lock();
|
|
|
|
|
|
|
|
Self::with_entries(&mut *inner, cur_ts, BucketEntryState::Unreliable, |k, e| {
|
|
|
|
if e.needs_ping(k, cur_ts, relay_node_id) {
|
|
|
|
let nr = NodeRef::new(self.clone(), *k, e, None);
|
|
|
|
log_rtab!(
|
|
|
|
" --- ping validating: {:?} ({})",
|
|
|
|
nr,
|
|
|
|
e.state_debug_info(cur_ts)
|
|
|
|
);
|
2022-06-13 00:58:02 +00:00
|
|
|
unord.push(MustJoinHandle::new(intf::spawn_local(
|
|
|
|
rpc.clone().rpc_call_status(nr),
|
|
|
|
)));
|
2022-06-11 22:47:58 +00:00
|
|
|
}
|
|
|
|
Option::<()>::None
|
|
|
|
});
|
|
|
|
}
|
|
|
|
|
|
|
|
// Wait for futures to complete
|
|
|
|
while unord.next().await.is_some() {}
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
// Compute transfer statistics to determine how 'fast' a node is
|
2022-06-10 21:07:10 +00:00
|
|
|
#[instrument(level = "trace", skip(self), err)]
|
2022-06-13 00:58:02 +00:00
|
|
|
async fn rolling_transfers_task_routine(
|
|
|
|
self,
|
|
|
|
stop_token: StopToken,
|
|
|
|
last_ts: u64,
|
|
|
|
cur_ts: u64,
|
|
|
|
) -> Result<(), String> {
|
2022-05-24 21:13:52 +00:00
|
|
|
// log_rtab!("--- rolling_transfers task");
|
2021-11-26 14:54:38 +00:00
|
|
|
let inner = &mut *self.inner.lock();
|
|
|
|
|
|
|
|
// Roll our own node's transfers
|
2022-03-19 22:19:40 +00:00
|
|
|
inner.self_transfer_stats_accounting.roll_transfers(
|
|
|
|
last_ts,
|
|
|
|
cur_ts,
|
|
|
|
&mut inner.self_transfer_stats,
|
|
|
|
);
|
2021-11-26 14:54:38 +00:00
|
|
|
|
|
|
|
// Roll all bucket entry transfers
|
2021-11-22 16:28:30 +00:00
|
|
|
for b in &mut inner.buckets {
|
|
|
|
b.roll_transfers(last_ts, cur_ts);
|
|
|
|
}
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
|
|
|
// Ticks about once per second
|
|
|
|
// to run tick tasks which may run at slower tick rates as configured
|
|
|
|
pub async fn tick(&self) -> Result<(), String> {
|
|
|
|
// Do rolling transfers every ROLLING_TRANSFERS_INTERVAL_SECS secs
|
|
|
|
self.unlocked_inner.rolling_transfers_task.tick().await?;
|
|
|
|
|
2022-05-24 21:13:52 +00:00
|
|
|
// If routing table has no live entries, then add the bootstrap nodes to it
|
|
|
|
if Self::get_entry_count(&mut *self.inner.lock(), BucketEntryState::Unreliable) == 0 {
|
2021-11-22 16:28:30 +00:00
|
|
|
self.unlocked_inner.bootstrap_task.tick().await?;
|
|
|
|
}
|
|
|
|
|
|
|
|
// If we still don't have enough peers, find nodes until we do
|
|
|
|
let min_peer_count = {
|
|
|
|
let c = self.config.get();
|
|
|
|
c.network.dht.min_peer_count as usize
|
|
|
|
};
|
2022-05-24 21:13:52 +00:00
|
|
|
if Self::get_entry_count(&mut *self.inner.lock(), BucketEntryState::Unreliable)
|
|
|
|
< min_peer_count
|
|
|
|
{
|
2021-11-22 16:28:30 +00:00
|
|
|
self.unlocked_inner.peer_minimum_refresh_task.tick().await?;
|
|
|
|
}
|
|
|
|
// Ping validate some nodes to groom the table
|
|
|
|
self.unlocked_inner.ping_validator_task.tick().await?;
|
|
|
|
|
2022-04-07 13:55:09 +00:00
|
|
|
// Keepalive
|
|
|
|
|
2021-11-22 16:28:30 +00:00
|
|
|
Ok(())
|
|
|
|
}
|
2021-11-26 14:54:38 +00:00
|
|
|
|
|
|
|
//////////////////////////////////////////////////////////////////////
|
|
|
|
// Stats Accounting
|
2022-04-18 22:49:33 +00:00
|
|
|
pub fn stats_question_sent(
|
|
|
|
&self,
|
|
|
|
node_ref: NodeRef,
|
|
|
|
ts: u64,
|
|
|
|
bytes: u64,
|
|
|
|
expects_answer: bool,
|
|
|
|
) {
|
2022-03-19 22:19:40 +00:00
|
|
|
self.inner
|
|
|
|
.lock()
|
|
|
|
.self_transfer_stats_accounting
|
|
|
|
.add_up(bytes);
|
2021-11-26 14:54:38 +00:00
|
|
|
node_ref.operate(|e| {
|
2022-04-18 22:49:33 +00:00
|
|
|
e.question_sent(ts, bytes, expects_answer);
|
2021-11-26 14:54:38 +00:00
|
|
|
})
|
|
|
|
}
|
2022-03-20 14:52:03 +00:00
|
|
|
pub fn stats_question_rcvd(&self, node_ref: NodeRef, ts: u64, bytes: u64) {
|
2022-03-19 22:19:40 +00:00
|
|
|
self.inner
|
|
|
|
.lock()
|
|
|
|
.self_transfer_stats_accounting
|
|
|
|
.add_down(bytes);
|
2021-11-26 14:54:38 +00:00
|
|
|
node_ref.operate(|e| {
|
|
|
|
e.question_rcvd(ts, bytes);
|
|
|
|
})
|
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
pub fn stats_answer_sent(&self, node_ref: NodeRef, bytes: u64) {
|
2022-03-19 22:19:40 +00:00
|
|
|
self.inner
|
|
|
|
.lock()
|
|
|
|
.self_transfer_stats_accounting
|
|
|
|
.add_up(bytes);
|
2021-11-26 14:54:38 +00:00
|
|
|
node_ref.operate(|e| {
|
2022-05-24 21:13:52 +00:00
|
|
|
e.answer_sent(bytes);
|
2021-11-26 14:54:38 +00:00
|
|
|
})
|
|
|
|
}
|
2022-03-20 14:52:03 +00:00
|
|
|
pub fn stats_answer_rcvd(&self, node_ref: NodeRef, send_ts: u64, recv_ts: u64, bytes: u64) {
|
2022-03-19 22:19:40 +00:00
|
|
|
self.inner
|
|
|
|
.lock()
|
|
|
|
.self_transfer_stats_accounting
|
|
|
|
.add_down(bytes);
|
2022-03-20 14:52:03 +00:00
|
|
|
self.inner
|
|
|
|
.lock()
|
|
|
|
.self_latency_stats_accounting
|
|
|
|
.record_latency(recv_ts - send_ts);
|
2021-11-26 14:54:38 +00:00
|
|
|
node_ref.operate(|e| {
|
|
|
|
e.answer_rcvd(send_ts, recv_ts, bytes);
|
|
|
|
})
|
|
|
|
}
|
2022-05-24 21:13:52 +00:00
|
|
|
pub fn stats_question_lost(&self, node_ref: NodeRef) {
|
|
|
|
node_ref.operate(|e| {
|
|
|
|
e.question_lost();
|
|
|
|
})
|
|
|
|
}
|
|
|
|
pub fn stats_failed_to_send(&self, node_ref: NodeRef, ts: u64, expects_answer: bool) {
|
2021-11-26 14:54:38 +00:00
|
|
|
node_ref.operate(|e| {
|
2022-05-24 21:13:52 +00:00
|
|
|
e.failed_to_send(ts, expects_answer);
|
2021-11-26 14:54:38 +00:00
|
|
|
})
|
|
|
|
}
|
2022-03-24 14:14:50 +00:00
|
|
|
|
|
|
|
//////////////////////////////////////////////////////////////////////
|
|
|
|
// Routing Table Health Metrics
|
|
|
|
|
|
|
|
pub fn get_routing_table_health(&self) -> RoutingTableHealth {
|
|
|
|
let mut health = RoutingTableHealth::default();
|
|
|
|
let cur_ts = intf::get_timestamp();
|
|
|
|
let inner = self.inner.lock();
|
|
|
|
for bucket in &inner.buckets {
|
|
|
|
for entry in bucket.entries() {
|
|
|
|
match entry.1.state(cur_ts) {
|
|
|
|
BucketEntryState::Reliable => {
|
|
|
|
health.reliable_entry_count += 1;
|
|
|
|
}
|
|
|
|
BucketEntryState::Unreliable => {
|
|
|
|
health.unreliable_entry_count += 1;
|
|
|
|
}
|
|
|
|
BucketEntryState::Dead => {
|
|
|
|
health.dead_entry_count += 1;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
health
|
|
|
|
}
|
2021-11-22 16:28:30 +00:00
|
|
|
}
|