#![allow(unused_results, missing_docs)]
#[cfg(not(target_family = "wasm"))]
use crate::network::gossip::{
GossipHandle, IntraNodePayload, MyBehaviour, NetworkServiceWithoutSwarm, MAX_MESSAGE_SIZE,
};
use futures::StreamExt;
#[cfg(not(target_family = "wasm"))]
use libp2p::{
gossipsub, gossipsub::IdentTopic, kad::store::MemoryStore, mdns, request_response,
swarm::dial_opts::DialOpts, StreamProtocol,
};
use gadget_io::tokio::select;
use gadget_io::tokio::sync::{Mutex, RwLock};
use gadget_io::tokio::task::{spawn, JoinHandle};
use libp2p::Multiaddr;
use lru_mem::LruCache;
use sp_core::ecdsa;
use std::collections::BTreeMap;
use std::error::Error;
use std::io;
use std::net::IpAddr;
use std::str::FromStr;
use std::sync::atomic::AtomicU32;
use std::sync::Arc;
use std::time::Duration;
pub const AGENT_VERSION: &str = "tangle/gadget-sdk/1.0.0";
pub const CLIENT_VERSION: &str = "1.0.0";
pub struct NetworkConfig {
pub identity: libp2p::identity::Keypair,
pub ecdsa_key: ecdsa::Pair,
pub bootnodes: Vec<Multiaddr>,
pub bind_port: u16,
pub topics: Vec<String>,
}
impl std::fmt::Debug for NetworkConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("NetworkConfig")
.field("identity", &self.identity)
.field("bootnodes", &self.bootnodes)
.field("bind_port", &self.bind_port)
.field("topics", &self.topics)
.finish_non_exhaustive()
}
}
impl NetworkConfig {
#[must_use]
pub fn new(
identity: libp2p::identity::Keypair,
ecdsa_key: ecdsa::Pair,
bootnodes: Vec<Multiaddr>,
bind_port: u16,
topics: Vec<String>,
) -> Self {
Self {
identity,
ecdsa_key,
bootnodes,
bind_port,
topics,
}
}
pub fn new_service_network<T: Into<String>>(
identity: libp2p::identity::Keypair,
ecdsa_key: ecdsa::Pair,
bootnodes: Vec<Multiaddr>,
bind_port: u16,
service_name: T,
) -> Self {
Self::new(
identity,
ecdsa_key,
bootnodes,
bind_port,
vec![service_name.into()],
)
}
}
pub fn start_p2p_network(config: NetworkConfig) -> Result<GossipHandle, Box<dyn Error>> {
if config.topics.len() != 1 {
return Err("Only one network topic is allowed when running this function".into());
}
let (networks, _) = multiplexed_libp2p_network(config)?;
let network = networks.into_iter().next().ok_or("No network found")?.1;
Ok(network)
}
pub type NetworkResult = Result<(BTreeMap<String, GossipHandle>, JoinHandle<()>), Box<dyn Error>>;
#[allow(clippy::collapsible_else_if, clippy::too_many_lines)]
#[cfg(not(target_family = "wasm"))]
pub fn multiplexed_libp2p_network(config: NetworkConfig) -> NetworkResult {
use std::collections::BTreeMap;
let NetworkConfig {
identity,
bootnodes,
bind_port,
topics,
ecdsa_key,
} = config;
let topics_unique = topics
.iter()
.cloned()
.collect::<std::collections::BTreeSet<_>>()
.into_iter()
.collect::<Vec<_>>();
if topics_unique.len() != topics.len() {
return Err("All topics must be unique".into());
}
let networks = topics;
let my_id = identity.public().to_peer_id();
let mut swarm = libp2p::SwarmBuilder::with_existing_identity(identity)
.with_tokio()
.with_tcp(
libp2p::tcp::Config::default().nodelay(true), libp2p::noise::Config::new,
libp2p::yamux::Config::default,
)?
.with_quic_config(|mut config| {
config.handshake_timeout = Duration::from_secs(30);
config
})
.with_dns()?
.with_relay_client(libp2p::noise::Config::new, libp2p::yamux::Config::default)?
.with_behaviour(|key, relay_client| {
let gossipsub_config = gossipsub::ConfigBuilder::default()
.protocol_id_prefix("/tangle/gadget-binary-sdk/meshsub")
.max_transmit_size(MAX_MESSAGE_SIZE)
.validate_messages()
.validation_mode(gossipsub::ValidationMode::Strict) .build()
.map_err(|msg| io::Error::new(io::ErrorKind::Other, msg))?;
let gossipsub = gossipsub::Behaviour::new(
gossipsub::MessageAuthenticity::Signed(key.clone()),
gossipsub_config,
)?;
let mdns =
mdns::tokio::Behaviour::new(mdns::Config::default(), key.public().to_peer_id())?;
let p2p_config = request_response::Config::default();
let protocols = networks
.iter()
.map(|n| {
(
StreamProtocol::try_from_owned(n.clone()).expect("Invalid network name"),
request_response::ProtocolSupport::Full,
)
})
.collect::<Vec<_>>();
let p2p = request_response::Behaviour::new(protocols, p2p_config);
let identify = libp2p::identify::Behaviour::new(
libp2p::identify::Config::new(CLIENT_VERSION.into(), key.public())
.with_agent_version(AGENT_VERSION.into()),
);
let memory_db = MemoryStore::new(key.public().to_peer_id());
let kadmelia = libp2p::kad::Behaviour::new(key.public().to_peer_id(), memory_db);
let dcutr = libp2p::dcutr::Behaviour::new(key.public().to_peer_id());
let relay_config = libp2p::relay::Config::default();
let relay = libp2p::relay::Behaviour::new(key.public().to_peer_id(), relay_config);
let ping = libp2p::ping::Behaviour::new(libp2p::ping::Config::default());
Ok(MyBehaviour {
gossipsub,
mdns,
p2p,
identify,
kadmelia,
dcutr,
relay,
relay_client,
ping,
})
})?
.with_swarm_config(|c| c.with_idle_connection_timeout(Duration::from_secs(60)))
.build();
let mut inbound_mapping = Vec::new();
let (tx_to_outbound, mut rx_to_outbound) =
gadget_io::tokio::sync::mpsc::unbounded_channel::<IntraNodePayload>();
let ecdsa_peer_id_to_libp2p_id = Arc::new(RwLock::new(BTreeMap::new()));
let mut handles_ret = BTreeMap::new();
for network in networks {
let topic = IdentTopic::new(network.clone());
swarm.behaviour_mut().gossipsub.subscribe(&topic)?;
let (inbound_tx, inbound_rx) = gadget_io::tokio::sync::mpsc::unbounded_channel();
let connected_peers = Arc::new(AtomicU32::new(0));
inbound_mapping.push((topic.clone(), inbound_tx, connected_peers.clone()));
handles_ret.insert(
network,
GossipHandle {
connected_peers,
topic,
tx_to_outbound: tx_to_outbound.clone(),
rx_from_inbound: Arc::new(Mutex::new(inbound_rx)),
ecdsa_peer_id_to_libp2p_id: ecdsa_peer_id_to_libp2p_id.clone(),
recent_messages: LruCache::new(16 * 1024).into(),
my_id,
},
);
}
let ips_to_bind_to = [
IpAddr::from_str("::").unwrap(), IpAddr::from_str("0.0.0.0").unwrap(), ];
for addr in ips_to_bind_to {
let ip_label = if addr.is_ipv4() { "ip4" } else { "ip6" };
swarm.listen_on(format!("/{ip_label}/{addr}/udp/{bind_port}/quic-v1").parse()?)?;
swarm.listen_on(format!("/{ip_label}/{addr}/tcp/{bind_port}").parse()?)?;
}
for bootnode in &bootnodes {
swarm.dial(
DialOpts::unknown_peer_id()
.address(bootnode.clone())
.build(),
)?;
}
let worker = async move {
let span = tracing::debug_span!("network_worker");
let _enter = span.enter();
let service = NetworkServiceWithoutSwarm {
inbound_mapping: &inbound_mapping,
ecdsa_peer_id_to_libp2p_id,
ecdsa_key: &ecdsa_key,
span: tracing::debug_span!(parent: &span, "network_service"),
my_id,
};
loop {
select! {
Some(msg) = rx_to_outbound.recv() => {
service.with_swarm(&mut swarm).handle_intra_node_payload(msg);
}
event = swarm.select_next_some() => {
service.with_swarm(&mut swarm).handle_swarm_event(event).await;
}
}
}
};
let spawn_handle = spawn(worker);
Ok((handles_ret, spawn_handle))
}