use std::{
cmp::min,
task::{Context, Poll},
};
use futures::{
future::BoxFuture,
stream::{FuturesUnordered, StreamExt},
FutureExt,
};
use tokio::time::timeout;
use tower::{Service, ServiceExt};
use zebra_chain::serialization::DateTime32;
use crate::{
address_book_updater::{AddressBookRequest, AddressBookService},
constants,
meta_addr::MetaAddrChange,
peer_set::set::MorePeers,
types::MetaAddr,
BoxError, Request, Response,
};
use super::RateLimitBySkipping;
pub(crate) type CrawlService<S> = RateLimitBySkipping<CrawlFanout<S>>;
pub(crate) async fn crawl_once<S>(
crawl_service: &mut CrawlService<S>,
fanout_limit: Option<usize>,
) -> Result<Option<MorePeers>, BoxError>
where
S: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
S::Future: Send + 'static,
{
crawl_service.ready().await?.call(fanout_limit).await
}
#[derive(Clone)]
pub(crate) struct CrawlFanout<S>
where
S: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
S::Future: Send + 'static,
{
peer_service: S,
address_book_service: AddressBookService,
}
impl<S> CrawlFanout<S>
where
S: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
S::Future: Send + 'static,
{
pub(super) fn service(
peer_service: S,
address_book_service: AddressBookService,
) -> CrawlService<S> {
RateLimitBySkipping::new(
Self {
peer_service,
address_book_service,
},
constants::MIN_PEER_GET_ADDR_INTERVAL,
)
}
}
impl<S> Service<Option<usize>> for CrawlFanout<S>
where
S: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
S::Future: Send + 'static,
{
type Response = Option<MorePeers>;
type Error = BoxError;
type Future = BoxFuture<'static, Result<Option<MorePeers>, BoxError>>;
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
Poll::Ready(Ok(()))
}
fn call(&mut self, fanout_limit: Option<usize>) -> Self::Future {
let peer_service = self.peer_service.clone();
let address_book_service = self.address_book_service.clone();
async move {
match timeout(
constants::PEER_GET_ADDR_TIMEOUT,
crawl_fanout(peer_service, address_book_service, fanout_limit),
)
.await
{
Ok(fanout_result) => fanout_result,
Err(_elapsed) => {
info!("timeout waiting for peer service readiness or peer responses");
Ok(None)
}
}
}
.boxed()
}
}
async fn crawl_fanout<S>(
mut peer_service: S,
address_book_service: AddressBookService,
fanout_limit: Option<usize>,
) -> Result<Option<MorePeers>, BoxError>
where
S: Service<Request, Response = Response, Error = BoxError> + Clone + Send + 'static,
S::Future: Send + 'static,
{
let fanout_limit = fanout_limit
.map(|fanout_limit| min(fanout_limit, constants::GET_ADDR_FANOUT))
.unwrap_or(constants::GET_ADDR_FANOUT);
debug!(?fanout_limit, "sending GetPeers requests");
let mut responses = FuturesUnordered::new();
let mut more_peers = None;
for attempt in 0..fanout_limit {
if attempt > 0 {
tokio::task::yield_now().await;
}
let peer_service = peer_service.ready().await?;
responses.push(peer_service.call(Request::Peers));
}
let mut address_book_updates = FuturesUnordered::new();
while let Some(rsp) = responses.next().await {
match rsp {
Ok(Response::Peers(addrs)) => {
trace!(
addr_count = ?addrs.len(),
?addrs,
"got response to GetPeers"
);
let addrs = validate_addrs(addrs, DateTime32::now());
address_book_updates.push(send_addrs(address_book_service.clone(), addrs));
more_peers = Some(MorePeers);
}
Err(e) => {
trace!(?e, "got error in GetPeers request");
}
Ok(_) => unreachable!("Peers requests always return Peers responses"),
}
}
while let Some(()) = address_book_updates.next().await {}
Ok(more_peers)
}
async fn send_addrs(
address_book_service: AddressBookService,
addrs: impl IntoIterator<Item = MetaAddr>,
) {
let addrs: Vec<MetaAddrChange> = addrs
.into_iter()
.map(MetaAddr::new_gossiped_change)
.map(|maybe_addr| maybe_addr.expect("Received gossiped peers always have services set"))
.collect();
debug!(count = ?addrs.len(), "sending gossiped addresses to the address book");
if addrs.is_empty() {
return;
}
let result = address_book_service
.oneshot(AddressBookRequest::ExtendGossiped(addrs))
.boxed()
.await;
if let Err(error) = result {
debug!(
?error,
"error sending gossiped addresses to the address book, is Zebra shutting down?"
);
}
}
pub(super) fn validate_addrs(
addrs: impl IntoIterator<Item = MetaAddr>,
last_seen_limit: DateTime32,
) -> impl Iterator<Item = MetaAddr> {
let mut addrs: Vec<_> = addrs.into_iter().collect();
limit_last_seen_times(&mut addrs, last_seen_limit);
addrs.into_iter()
}
fn limit_last_seen_times(addrs: &mut Vec<MetaAddr>, last_seen_limit: DateTime32) {
let last_seen_times = addrs.iter().map(|meta_addr| {
meta_addr
.untrusted_last_seen()
.expect("unexpected missing last seen: should be provided by deserialization")
});
let oldest_seen = last_seen_times.clone().min().unwrap_or(DateTime32::MIN);
let newest_seen = last_seen_times.max().unwrap_or(DateTime32::MAX);
if newest_seen > last_seen_limit {
let offset = newest_seen
.checked_duration_since(last_seen_limit)
.expect("unexpected underflow: just checked newest_seen is greater");
if oldest_seen.checked_sub(offset).is_some() {
for addr in addrs {
let last_seen = addr
.untrusted_last_seen()
.expect("unexpected missing last seen: should be provided by deserialization");
let last_seen = last_seen
.checked_sub(offset)
.expect("unexpected underflow: just checked oldest_seen");
addr.set_untrusted_last_seen(last_seen);
}
} else {
addrs.clear();
}
}
}