use crate::dht::errors::DhtError;
use crate::dht::DhtConfig;
use libp2p::{
allow_block_list,
core::muxing::StreamMuxerBox,
core::transport::Boxed,
core::upgrade,
dns,
identify::{Behaviour as Identify, Config as IdentifyConfig, Event as IdentifyEvent},
identity::Keypair,
kad::{self, Kademlia, KademliaConfig, KademliaEvent, KademliaStoreInserts},
noise,
swarm::{self, NetworkBehaviour, SwarmBuilder, SwarmEvent, THandlerErr},
tcp, yamux, PeerId, Swarm, Transport,
};
use std::time::Duration;
use std::{io, result::Result};
use void::Void;
const CONNECTION_TIMEOUT_SECONDS: u64 = 20;
#[derive(Debug)]
#[allow(clippy::large_enum_variant)]
pub enum DHTEvent {
Kademlia(KademliaEvent),
Identify(IdentifyEvent),
Void,
}
impl From<KademliaEvent> for DHTEvent {
fn from(event: KademliaEvent) -> Self {
DHTEvent::Kademlia(event)
}
}
impl From<IdentifyEvent> for DHTEvent {
fn from(event: IdentifyEvent) -> Self {
DHTEvent::Identify(event)
}
}
impl From<Void> for DHTEvent {
fn from(_: Void) -> Self {
DHTEvent::Void
}
}
#[derive(NetworkBehaviour)]
#[behaviour(out_event = "DHTEvent", event_process = false)]
pub struct DhtBehavior {
pub identify: Identify,
pub kad: Kademlia<kad::record::store::MemoryStore>,
blocked_peers: allow_block_list::Behaviour<allow_block_list::BlockedPeers>,
}
pub type DHTSwarmEvent =
SwarmEvent<<DhtBehavior as swarm::NetworkBehaviour>::OutEvent, THandlerErr<DhtBehavior>>;
impl DhtBehavior {
pub fn new(keypair: &Keypair, local_peer_id: &PeerId, config: &DhtConfig) -> Self {
let kad = {
let mut cfg = KademliaConfig::default();
cfg.set_query_timeout(Duration::from_secs(config.query_timeout.into()));
cfg.set_record_filtering(KademliaStoreInserts::FilterBoth);
cfg.set_record_ttl(Some(Duration::from_secs(config.record_ttl.into())));
cfg.set_publication_interval(Some(Duration::from_secs(
config.publication_interval.into(),
)));
cfg.set_replication_interval(Some(Duration::from_secs(
config.replication_interval.into(),
)));
cfg.set_provider_record_ttl(Some(Duration::from_secs(config.record_ttl.into())));
cfg.set_provider_publication_interval(Some(Duration::from_secs(
config.publication_interval.into(),
)));
let store = kad::record::store::MemoryStore::new(local_peer_id.to_owned());
Kademlia::with_config(local_peer_id.to_owned(), store, cfg)
};
let identify = {
let config = IdentifyConfig::new("ipfs/1.0.0".into(), keypair.public())
.with_agent_version(format!("noosphere-ns/{}", env!("CARGO_PKG_VERSION")));
Identify::new(config)
};
DhtBehavior {
kad,
identify,
blocked_peers: allow_block_list::Behaviour::default(),
}
}
}
fn build_transport(keypair: &Keypair) -> Result<Boxed<(PeerId, StreamMuxerBox)>, io::Error> {
let transport =
dns::TokioDnsConfig::system(tcp::tokio::Transport::new(tcp::Config::new().nodelay(true)))?;
let noise_keys = noise::Config::new(keypair).expect("Noise key generation failed.");
Ok(transport
.upgrade(upgrade::Version::V1)
.authenticate(noise_keys)
.multiplex(yamux::Config::default())
.timeout(std::time::Duration::from_secs(CONNECTION_TIMEOUT_SECONDS))
.boxed())
}
pub fn build_swarm(
keypair: &Keypair,
local_peer_id: &PeerId,
config: &DhtConfig,
) -> Result<Swarm<DhtBehavior>, DhtError> {
let transport = build_transport(keypair).map_err(DhtError::from)?;
let behaviour = DhtBehavior::new(keypair, local_peer_id, config);
let swarm =
SwarmBuilder::with_tokio_executor(transport, behaviour, local_peer_id.to_owned()).build();
Ok(swarm)
}