Yay, SeedService makes a remote 'connect' happy

This commit is contained in:
Deirdre Connolly 2019-11-11 18:08:59 -05:00 committed by Deirdre Connolly
parent 4d3ab201e6
commit d6ab549fd5
1 changed files with 61 additions and 36 deletions

View File

@ -8,22 +8,26 @@ use std::{
}; };
use abscissa_core::{config, Command, FrameworkError, Options, Runnable}; use abscissa_core::{config, Command, FrameworkError, Options, Runnable};
use futures::stream::StreamExt; use futures::channel::oneshot;
use tower::{buffer::Buffer, Service, ServiceExt}; use tower::{buffer::Buffer, Service, ServiceExt};
use zebra_network::{AddressBook, BoxedStdError, Request, Response}; use zebra_network::{AddressBook, BoxedStdError, Request, Response};
use crate::{config::ZebradConfig, prelude::*}; use crate::{config::ZebradConfig, prelude::*};
#[derive(Clone)] /// Whether our `SeedService` is poll_ready or not.
struct SeedService { #[derive(Debug)]
address_book: Option<Arc<Mutex<AddressBook>>>, enum SeederState {
///
TempState,
///
AwaitingAddressBook(oneshot::Receiver<Arc<Mutex<AddressBook>>>),
///
Ready(Arc<Mutex<AddressBook>>),
} }
impl SeedService { #[derive(Debug)]
fn set_address_book(&mut self, address_book: Arc<Mutex<AddressBook>>) { struct SeedService {
debug!("Settings SeedService.address_book: {:?}", address_book); state: SeederState,
self.address_book = Some(address_book);
}
} }
impl Service<Request> for SeedService { impl Service<Request> for SeedService {
@ -32,33 +36,51 @@ impl Service<Request> for SeedService {
type Future = type Future =
Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send + 'static>>; Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send + 'static>>;
fn poll_ready(&mut self, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> { fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
Ok(()).into() info!("State: {:?}", self.state);
let mut poll_result = Poll::Pending;
// We want to be able to consume the state, but it's behind a mutable
// reference, so we can't move it out of self without swapping in a
// placeholder, even if we immediately overwrite the placeholder.
let tmp_state = std::mem::replace(&mut self.state, SeederState::TempState);
self.state = match tmp_state {
SeederState::AwaitingAddressBook(mut rx) => match rx.try_recv() {
Ok(Some(address_book)) => {
info!("Message received! {:?}", address_book);
poll_result = Poll::Ready(Ok(()));
SeederState::Ready(address_book)
}
_ => SeederState::AwaitingAddressBook(rx),
},
SeederState::Ready(_) => {
poll_result = Poll::Ready(Ok(()));
tmp_state
}
SeederState::TempState => tmp_state,
};
return poll_result;
} }
fn call(&mut self, req: Request) -> Self::Future { fn call(&mut self, req: Request) -> Self::Future {
info!("SeedService handling a request: {:?}", req); info!("SeedService handling a request: {:?}", req);
match &self.address_book { let response = match (req, &self.state) {
Some(address_book) => trace!( (Request::GetPeers, SeederState::Ready(address_book)) => {
"SeedService address_book total: {:?}", info!("Responding to GetPeers");
address_book.lock().unwrap().len()
),
_ => (),
};
let response = match req { Ok::<Response, Self::Error>(Response::Peers(
Request::GetPeers => match &self.address_book { address_book.lock().unwrap().peers().collect(),
Some(address_book) => { ))
info!("Responding to GetPeers"); }
_ => {
trace!("Where is my address_book??? {:?}", &self.state);
Ok::<Response, Self::Error>(Response::Peers( Ok::<Response, Self::Error>(Response::Ok)
address_book.lock().unwrap().peers().collect(), }
))
}
_ => Ok::<Response, Self::Error>(Response::Ok),
},
_ => Ok::<Response, Self::Error>(Response::Ok),
}; };
info!("SeedService response: {:?}", response); info!("SeedService response: {:?}", response);
@ -71,7 +93,7 @@ impl Service<Request> for SeedService {
/// ///
/// A DNS seeder command to spider and collect as many valid peer /// A DNS seeder command to spider and collect as many valid peer
/// addresses as we can. /// addresses as we can.
#[derive(Command, Debug, Options)] #[derive(Command, Debug, Default, Options)]
pub struct SeedCmd { pub struct SeedCmd {
/// Filter strings /// Filter strings
#[options(free)] #[options(free)]
@ -130,15 +152,17 @@ impl SeedCmd {
info!("begin tower-based peer handling test stub"); info!("begin tower-based peer handling test stub");
let mut seed_service = SeedService { address_book: None }; let (addressbook_tx, addressbook_rx) = oneshot::channel();
// let node = Buffer::new(seed_service, 1); let seed_service = SeedService {
state: SeederState::AwaitingAddressBook(addressbook_rx),
};
let node = Buffer::new(seed_service, 1);
let config = app_config().network.clone(); let config = app_config().network.clone();
info!("{:?}", config);
let (mut peer_set, address_book) = zebra_network::init(config, seed_service.clone()).await; let (mut peer_set, address_book) = zebra_network::init(config, node).await;
seed_service.set_address_book(address_book.clone()); let _ = addressbook_tx.send(address_book);
// XXX Do not tell our DNS seed queries about gossiped addrs // XXX Do not tell our DNS seed queries about gossiped addrs
// that we have not connected to before? // that we have not connected to before?
@ -147,10 +171,11 @@ impl SeedCmd {
info!("peer_set became ready"); info!("peer_set became ready");
#[cfg(dos)]
use std::time::Duration; use std::time::Duration;
use tokio::timer::Interval; use tokio::timer::Interval;
//#[cfg(dos)] #[cfg(dos)]
// Fire GetPeers requests at ourselves, for testing. // Fire GetPeers requests at ourselves, for testing.
tokio::spawn(async move { tokio::spawn(async move {
let mut interval_stream = Interval::new_interval(Duration::from_secs(1)); let mut interval_stream = Interval::new_interval(Duration::from_secs(1));