#[cfg(test)]
pub(crate) mod tests;
pub(crate) mod command;
pub(super) mod config;
mod dispatcher;
pub(super) mod event;
pub(super) mod event_stream;
use self::{
command::Command,
config::Config,
dispatcher::Dispatcher,
event::{Elders, Event, NodeElderChange},
event_stream::EventStream,
};
use crate::dbs::UsedSpace;
use crate::messaging::{system::SystemMsg, DstLocation, WireMsg};
use crate::routing::{
core::{join_network, Comm, ConnectionEvent, Core},
ed25519,
error::{Error, Result},
log_markers::LogMarker,
messages::WireMsgUtils,
network_knowledge::SectionAuthorityProvider,
node::Node,
Peer, MIN_ADULT_AGE,
};
use ed25519_dalek::{PublicKey, Signature, Signer, KEYPAIR_LENGTH};
use crate::types::PublicKey as TypesPublicKey;
use itertools::Itertools;
use secured_linked_list::SecuredLinkedList;
use std::path::PathBuf;
use std::{collections::BTreeSet, net::SocketAddr, sync::Arc};
use tokio::{sync::mpsc, task};
use xor_name::{Prefix, XorName};
#[allow(missing_debug_implementations)]
pub struct Routing {
dispatcher: Arc<Dispatcher>,
}
static EVENT_CHANNEL_SIZE: usize = 20;
impl Routing {
pub async fn new(
config: Config,
used_space: UsedSpace,
root_storage_dir: PathBuf,
) -> Result<(Self, EventStream)> {
let (event_tx, event_rx) = mpsc::channel(EVENT_CHANNEL_SIZE);
let (connection_event_tx, mut connection_event_rx) = mpsc::channel(1);
let core = if config.first {
let keypair = ed25519::gen_keypair(&Prefix::default().range_inclusive(), 255);
let node_name = ed25519::name(&keypair.public);
info!(
"{} Starting a new network as the genesis node (PID: {}).",
node_name,
std::process::id()
);
let comm = Comm::new(
config.local_addr,
config.network_config,
connection_event_tx,
)
.await?;
let node = Node::new(keypair, comm.our_connection_info());
let genesis_sk_set = bls::SecretKeySet::random(0, &mut rand::thread_rng());
let core = Core::first_node(
comm,
node,
event_tx,
used_space,
root_storage_dir.clone(),
genesis_sk_set,
)
.await?;
let network_knowledge = core.network_knowledge();
let elders = Elders {
prefix: network_knowledge.prefix().await,
key: network_knowledge.section_key().await,
remaining: BTreeSet::new(),
added: network_knowledge.authority_provider().await.names(),
removed: BTreeSet::new(),
};
trace!("{}", LogMarker::PromotedToElder);
core.send_event(Event::EldersChanged {
elders,
self_status_change: NodeElderChange::Promoted,
})
.await;
let genesis_key = network_knowledge.genesis_key();
info!(
"{} Genesis node started!. Genesis key {:?}, hex: {}",
node_name,
genesis_key,
hex::encode(genesis_key.to_bytes())
);
core
} else {
let genesis_key_str = config.genesis_key.ok_or_else(|| {
Error::Configuration("Network's genesis key was not provided.".to_string())
})?;
let genesis_key = TypesPublicKey::bls_from_hex(&genesis_key_str)?
.bls()
.ok_or_else(|| {
Error::Configuration(
"Unexpectedly failed to obtain genesis key from configuration.".to_string(),
)
})?;
let keypair = config.keypair.unwrap_or_else(|| {
ed25519::gen_keypair(&Prefix::default().range_inclusive(), MIN_ADULT_AGE)
});
let node_name = ed25519::name(&keypair.public);
info!("{} Bootstrapping a new node.", node_name);
let (comm, bootstrap_peer) = Comm::bootstrap(
config.local_addr,
config
.bootstrap_nodes
.iter()
.copied()
.collect_vec()
.as_slice(),
config.network_config,
connection_event_tx,
)
.await?;
info!(
"{} Joining as a new node (PID: {}) our socket: {}, bootstrapper was: {}, network's genesis key: {:?}",
node_name,
std::process::id(),
comm.our_connection_info(),
bootstrap_peer.addr(),
genesis_key
);
let joining_node = Node::new(keypair, comm.our_connection_info());
let (node, network_knowledge) = join_network(
joining_node,
&comm,
&mut connection_event_rx,
bootstrap_peer,
genesis_key,
)
.await?;
let core = Core::new(
comm,
node,
network_knowledge,
None,
event_tx,
used_space,
root_storage_dir.to_path_buf(),
false,
)
.await?;
info!("{} Joined the network!", core.node.read().await.name());
info!("Our AGE: {}", core.node.read().await.age());
core
};
let dispatcher = Arc::new(Dispatcher::new(core));
let event_stream = EventStream::new(event_rx);
let _handle = task::spawn(handle_connection_events(
dispatcher.clone(),
connection_event_rx,
));
dispatcher.clone().start_network_probing().await;
dispatcher.clone().write_prefixmap_to_disk().await;
let routing = Self { dispatcher };
Ok((routing, event_stream))
}
pub async fn set_joins_allowed(&self, joins_allowed: bool) -> Result<()> {
let command = Command::SetJoinsAllowed(joins_allowed);
self.dispatcher.clone().handle_commands(command, None).await
}
pub async fn start_connectivity_test(&self, name: XorName) -> Result<()> {
let command = Command::StartConnectivityTest(name);
self.dispatcher.clone().handle_commands(command, None).await
}
pub async fn age(&self) -> u8 {
self.dispatcher.core.node.read().await.age()
}
pub async fn public_key(&self) -> PublicKey {
self.dispatcher.core.node.read().await.keypair.public
}
pub async fn keypair_as_bytes(&self) -> [u8; KEYPAIR_LENGTH] {
self.dispatcher.core.node.read().await.keypair.to_bytes()
}
pub async fn sign_as_node(&self, data: &[u8]) -> Signature {
self.dispatcher.core.node.read().await.keypair.sign(data)
}
pub async fn sign_as_elder(
&self,
data: &[u8],
public_key: &bls::PublicKey,
) -> Result<(usize, bls::SignatureShare)> {
self.dispatcher
.core
.sign_with_section_key_share(data, public_key)
.await
}
pub async fn verify(&self, data: &[u8], signature: &Signature) -> bool {
self.dispatcher
.core
.node
.read()
.await
.keypair
.verify(data, signature)
.is_ok()
}
pub async fn name(&self) -> XorName {
self.dispatcher.core.node.read().await.name()
}
pub async fn our_connection_info(&self) -> SocketAddr {
self.dispatcher.core.our_connection_info()
}
pub async fn section_chain(&self) -> SecuredLinkedList {
self.dispatcher.core.section_chain().await
}
pub async fn genesis_key(&self) -> bls::PublicKey {
*self.dispatcher.core.network_knowledge().genesis_key()
}
pub async fn our_prefix(&self) -> Prefix {
self.dispatcher.core.network_knowledge().prefix().await
}
pub async fn matches_our_prefix(&self, name: &XorName) -> bool {
self.our_prefix().await.matches(name)
}
pub async fn is_elder(&self) -> bool {
self.dispatcher.core.is_elder().await
}
pub async fn our_elders(&self) -> Vec<Peer> {
self.dispatcher
.core
.network_knowledge()
.authority_provider()
.await
.elders_vec()
}
pub async fn our_elders_sorted_by_distance_to(&self, name: &XorName) -> Vec<Peer> {
self.our_elders()
.await
.into_iter()
.sorted_by(|lhs, rhs| name.cmp_distance(&lhs.name(), &rhs.name()))
.collect()
}
pub async fn our_adults(&self) -> Vec<Peer> {
self.dispatcher.core.network_knowledge().adults().await
}
pub async fn our_adults_sorted_by_distance_to(&self, name: &XorName) -> Vec<Peer> {
self.our_adults()
.await
.into_iter()
.sorted_by(|lhs, rhs| name.cmp_distance(&lhs.name(), &rhs.name()))
.collect()
}
pub async fn matching_section(&self, name: &XorName) -> Result<SectionAuthorityProvider> {
self.dispatcher.core.matching_section(name).await
}
pub async fn sign_single_src_msg(
&self,
node_msg: SystemMsg,
dst: DstLocation,
) -> Result<WireMsg> {
let src_section_pk = *self.section_chain().await.last_key();
WireMsg::single_src(
&self.dispatcher.core.node.read().await.clone(),
dst,
node_msg,
src_section_pk,
)
}
pub async fn sign_msg_for_dst_accumulation(
&self,
node_msg: SystemMsg,
dst: DstLocation,
) -> Result<WireMsg> {
let src = self.name().await;
let src_section_pk = *self.section_chain().await.last_key();
WireMsg::for_dst_accumulation(
&self.dispatcher.core.key_share().await.map_err(|err| err)?,
src,
dst,
node_msg,
src_section_pk,
)
}
pub async fn send_message(&self, wire_msg: WireMsg) -> Result<()> {
trace!(
"{:?} {:?}",
LogMarker::DispatchSendMsgCmd,
wire_msg.msg_id()
);
self.dispatcher
.clone()
.handle_commands(Command::ParseAndSendWireMsg(wire_msg), None)
.await
}
pub async fn public_key_set(&self) -> Result<bls::PublicKeySet> {
self.dispatcher.core.public_key_set().await
}
pub async fn our_index(&self) -> Result<usize> {
self.dispatcher.core.our_index().await
}
}
async fn handle_connection_events(
dispatcher: Arc<Dispatcher>,
mut incoming_conns: mpsc::Receiver<ConnectionEvent>,
) {
while let Some(event) = incoming_conns.recv().await {
match event {
ConnectionEvent::Received((sender, bytes)) => {
trace!(
"New message ({} bytes) received from: {:?}",
bytes.len(),
sender
);
let wire_msg = match WireMsg::from(bytes.clone()) {
Ok(wire_msg) => wire_msg,
Err(error) => {
error!("Failed to deserialize message header: {:?}", error);
continue;
}
};
let span = {
let core = &dispatcher.core;
trace_span!("handle_message", name = %core.node.read().await.name(), ?sender, msg_id = ?wire_msg.msg_id())
};
let _span_guard = span.enter();
trace!(
"{:?} from {:?} length {}",
LogMarker::DispatchHandleMsgCmd,
sender,
bytes.len(),
);
let command = Command::HandleMessage {
sender,
wire_msg,
original_bytes: Some(bytes),
};
let _handle = dispatcher.clone().handle_commands(command, None).await;
}
}
}
error!("Fatal error, the stream for incoming connections has been unexpectedly closed. No new connections or messages can be received from the network from here on.");
}