iuna

iuna - experimental devnet protocol
git clone https://iuna.jhx.app/git/iuna.git
Log | Files | Refs | README | LICENSE

network.rs (14595B)


      1 use std::{
      2     collections::{BTreeMap, BTreeSet},
      3     net::{IpAddr, SocketAddr},
      4     sync::{Arc, Mutex as StdMutex},
      5 };
      6 
      7 use anyhow::{Context, Result};
      8 use tokio::{net::TcpListener, sync::mpsc};
      9 
     10 use crate::app::{
     11     BlockInventory, GossipEnvelope, SharedNode, SharedPeerBook, debug_logging_enabled, now_ms,
     12 };
     13 
     14 use super::{
     15     GossipNetwork, GossipNetworkInner, InboundConnectionLimiter, InboundSessionPermit,
     16     InboundSessionRejection, OutboundBatch, P2pMetrics, P2pMetricsCounters, PEER_QUEUE_SIZE,
     17     STALE_INBOUND_PEER_RETENTION_MS, accept_loop, is_self_peer_address_for, new_node_id,
     18     outbound_session, outbound_supervisor,
     19 };
     20 
     21 impl GossipNetwork {
     22     #[cfg(test)]
     23     pub(crate) fn new_for_tests(node: SharedNode, peers: SharedPeerBook) -> Self {
     24         Self {
     25             inner: Arc::new(GossipNetworkInner {
     26                 node,
     27                 peers,
     28                 listen_addr: "127.0.0.1:0".parse().unwrap(),
     29                 p2p_announce_addr: tokio::sync::Mutex::new(None),
     30                 node_id: new_node_id(),
     31                 accept_task: tokio::sync::Mutex::new(None),
     32                 sessions: tokio::sync::Mutex::new(
     33                     BTreeMap::<String, mpsc::Sender<OutboundBatch>>::new(),
     34                 ),
     35                 inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())),
     36                 metrics: P2pMetricsCounters::default(),
     37             }),
     38         }
     39     }
     40 
     41     pub async fn start(
     42         node: SharedNode,
     43         peers: SharedPeerBook,
     44         addr: SocketAddr,
     45         p2p_announce_addr: Option<SocketAddr>,
     46         accept_inbound: bool,
     47     ) -> Result<Self> {
     48         let network = Self {
     49             inner: Arc::new(GossipNetworkInner {
     50                 node,
     51                 peers,
     52                 listen_addr: addr,
     53                 p2p_announce_addr: tokio::sync::Mutex::new(p2p_announce_addr),
     54                 node_id: new_node_id(),
     55                 accept_task: tokio::sync::Mutex::new(None),
     56                 sessions: tokio::sync::Mutex::new(
     57                     BTreeMap::<String, mpsc::Sender<OutboundBatch>>::new(),
     58                 ),
     59                 inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())),
     60                 metrics: P2pMetricsCounters::default(),
     61             }),
     62         };
     63 
     64         if accept_inbound {
     65             network.set_accept_inbound(true).await?;
     66         }
     67         tokio::spawn(outbound_supervisor(network.clone()));
     68         network.ensure_outbound_sessions().await;
     69         Ok(network)
     70     }
     71 
     72     pub async fn set_accept_inbound(&self, enabled: bool) -> Result<()> {
     73         let mut accept_task = self.inner.accept_task.lock().await;
     74         if enabled {
     75             if accept_task.is_some() {
     76                 return Ok(());
     77             }
     78             let listener = TcpListener::bind(self.inner.listen_addr)
     79                 .await
     80                 .with_context(|| format!("binding p2p listener on {}", self.inner.listen_addr))?;
     81             *accept_task = Some(tokio::spawn(accept_loop(self.clone(), listener)));
     82         } else if let Some(task) = accept_task.take() {
     83             task.abort();
     84         }
     85         Ok(())
     86     }
     87 
     88     pub async fn accepts_inbound(&self) -> bool {
     89         self.inner.accept_task.lock().await.is_some()
     90     }
     91 
     92     pub fn listen_addr(&self) -> SocketAddr {
     93         self.inner.listen_addr
     94     }
     95 
     96     pub async fn set_p2p_announce_addr(&self, addr: Option<SocketAddr>) {
     97         *self.inner.p2p_announce_addr.lock().await = addr;
     98     }
     99 
    100     pub(super) async fn advertised_addr(&self) -> Option<SocketAddr> {
    101         if !self.accepts_inbound().await {
    102             return None;
    103         }
    104         Some((*self.inner.p2p_announce_addr.lock().await).unwrap_or(self.inner.listen_addr))
    105     }
    106 
    107     pub(super) async fn self_filter_addr(&self) -> Option<SocketAddr> {
    108         if let Some(addr) = *self.inner.p2p_announce_addr.lock().await {
    109             return Some(addr);
    110         }
    111         self.accepts_inbound()
    112             .await
    113             .then_some(self.inner.listen_addr)
    114     }
    115 
    116     pub(super) async fn is_self_peer(&self, address: &str) -> bool {
    117         is_self_peer_address_for(
    118             address,
    119             self.inner.listen_addr,
    120             self.self_filter_addr().await,
    121         )
    122     }
    123 
    124     pub fn metrics(&self) -> P2pMetrics {
    125         self.inner.metrics.snapshot()
    126     }
    127 
    128     pub(super) fn try_acquire_inbound_session(
    129         &self,
    130         ip: IpAddr,
    131     ) -> std::result::Result<InboundSessionPermit, InboundSessionRejection> {
    132         self.inner
    133             .inbound_limiter
    134             .lock()
    135             .expect("inbound limiter mutex poisoned")
    136             .try_acquire(ip, now_ms())?;
    137         Ok(InboundSessionPermit {
    138             limiter: Arc::clone(&self.inner.inbound_limiter),
    139             ip,
    140         })
    141     }
    142 
    143     pub async fn broadcast(&self, envelopes: Vec<GossipEnvelope>) -> Result<()> {
    144         let envelopes = self.prepare_gossip(envelopes).await;
    145         if envelopes.is_empty() {
    146             return Ok(());
    147         }
    148 
    149         let sessions = self.inner.sessions.lock().await.clone();
    150         for (peer, sender) in sessions {
    151             if self.inner.peers.lock().await.is_banned(&peer) {
    152                 continue;
    153             }
    154             match sender.try_send(envelopes.clone()) {
    155                 Ok(()) => {}
    156                 Err(mpsc::error::TrySendError::Full(_)) => {
    157                     P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full);
    158                 }
    159                 Err(mpsc::error::TrySendError::Closed(_)) => {
    160                     P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed);
    161                 }
    162             }
    163         }
    164         Ok(())
    165     }
    166 
    167     async fn prepare_gossip(&self, envelopes: Vec<GossipEnvelope>) -> Vec<GossipEnvelope> {
    168         let mut blocks = Vec::new();
    169         let mut passthrough = Vec::new();
    170 
    171         for envelope in envelopes {
    172             match envelope {
    173                 GossipEnvelope::Block(block) => blocks.push(BlockInventory {
    174                     height: block.height,
    175                     hash: block.hash,
    176                 }),
    177                 GossipEnvelope::Blocks { blocks: batch } => {
    178                     blocks.extend(batch.into_iter().map(|block| BlockInventory {
    179                         height: block.height,
    180                         hash: block.hash,
    181                     }));
    182                 }
    183                 GossipEnvelope::Inventory { blocks: inv_blocks } => blocks.extend(inv_blocks),
    184                 other => passthrough.push(other),
    185             }
    186         }
    187 
    188         blocks.sort_by(|left, right| {
    189             left.height
    190                 .cmp(&right.height)
    191                 .then_with(|| left.hash.cmp(&right.hash))
    192         });
    193         blocks.dedup_by(|left, right| left.hash == right.hash);
    194 
    195         if !blocks.is_empty() {
    196             passthrough.push(GossipEnvelope::Inventory { blocks });
    197         }
    198         passthrough
    199     }
    200 
    201     pub async fn peer_exchange(&self) -> GossipEnvelope {
    202         let advertised_addr = self.advertised_addr().await;
    203         let self_filter_addr = self.self_filter_addr().await;
    204         let self_addr = advertised_addr.map(|addr| addr.to_string());
    205         let peers = self
    206             .inner
    207             .peers
    208             .lock()
    209             .await
    210             .addresses_except(self_addr.as_deref().unwrap_or(""))
    211             .into_iter()
    212             .filter(|peer| peer.parse::<SocketAddr>().is_ok())
    213             .filter(|peer| {
    214                 !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr)
    215             })
    216             .collect::<Vec<_>>();
    217         GossipEnvelope::PeerList {
    218             peers: self_addr.into_iter().chain(peers.into_iter()).collect(),
    219         }
    220     }
    221 
    222     pub(super) async fn ensure_outbound_sessions(&self) {
    223         self.inner
    224             .peers
    225             .lock()
    226             .await
    227             .prune_stale_inbound_peers_at(crate::app::now_ms(), STALE_INBOUND_PEER_RETENTION_MS);
    228         let addresses = self
    229             .inner
    230             .peers
    231             .lock()
    232             .await
    233             .connectable_addresses_at(crate::app::now_ms());
    234         let address_set = addresses.iter().cloned().collect::<BTreeSet<_>>();
    235         let self_filter_addr = self.self_filter_addr().await;
    236         let mut sessions = self.inner.sessions.lock().await;
    237         sessions.retain(|peer, _| {
    238             let keep = address_set.contains(peer)
    239                 && !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr);
    240             if !keep {
    241                 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips);
    242             }
    243             keep
    244         });
    245         for peer in addresses {
    246             if is_self_peer_address_for(&peer, self.inner.listen_addr, self_filter_addr) {
    247                 P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips);
    248                 continue;
    249             }
    250             if sessions.contains_key(&peer) {
    251                 continue;
    252             }
    253 
    254             let (sender, receiver) = mpsc::channel(PEER_QUEUE_SIZE);
    255             sessions.insert(peer.clone(), sender);
    256             tokio::spawn(outbound_session(self.clone(), peer, receiver));
    257         }
    258     }
    259 
    260     pub(super) async fn forward_outbox(&self) {
    261         let outbox = self.inner.node.lock().await.drain_outbox();
    262         if let Err(error) = self.broadcast(outbox).await {
    263             if debug_logging_enabled() {
    264                 eprintln!("p2p rebroadcast failed: {error:#}");
    265             }
    266         }
    267     }
    268 }
    269 
    270 #[cfg(test)]
    271 mod tests {
    272     use std::sync::Arc;
    273 
    274     use crate::{
    275         app::{GossipEnvelope, PeerBook},
    276         domain::Wallet,
    277     };
    278 
    279     use super::super::test_support::{allocations, gossip_network, node};
    280 
    281     #[tokio::test]
    282     async fn peer_exchange_does_not_advertise_self_when_outbound_only() {
    283         let alice = Wallet::from_seed("px-private-alice");
    284         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    285         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    286         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
    287             "127.0.0.1:9545".to_string(),
    288         ])));
    289         let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None);
    290         match network.peer_exchange().await {
    291             GossipEnvelope::PeerList { peers } => {
    292                 assert!(!peers.contains(&"127.0.0.1:9544".to_string()));
    293                 assert!(peers.contains(&"127.0.0.1:9545".to_string()));
    294             }
    295             other => panic!("expected peer list, got {other:?}"),
    296         }
    297     }
    298 
    299     #[tokio::test]
    300     async fn peer_exchange_omits_hostname_bootstrap_peers() {
    301         let alice = Wallet::from_seed("px-hostname-alice");
    302         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    303         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    304         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
    305             "iuna.jhx.app:9444".to_string(),
    306             "127.0.0.1:9545".to_string(),
    307         ])));
    308         let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None);
    309         match network.peer_exchange().await {
    310             GossipEnvelope::PeerList { peers } => {
    311                 assert!(!peers.contains(&"iuna.jhx.app:9444".to_string()));
    312                 assert!(peers.contains(&"127.0.0.1:9545".to_string()));
    313             }
    314             other => panic!("expected peer list, got {other:?}"),
    315         }
    316     }
    317 
    318     #[tokio::test]
    319     async fn peer_exchange_advertises_discovered_listening_peers() {
    320         let alice = Wallet::from_seed("px-discovered-alice");
    321         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    322         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    323         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    324         peers
    325             .lock()
    326             .await
    327             .add_discovered_peer("127.0.0.1:9546".to_string());
    328         let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None);
    329 
    330         match network.peer_exchange().await {
    331             GossipEnvelope::PeerList { peers } => {
    332                 assert!(peers.contains(&"127.0.0.1:9546".to_string()));
    333             }
    334             other => panic!("expected peer list, got {other:?}"),
    335         }
    336     }
    337 
    338     #[tokio::test]
    339     async fn peer_exchange_advertises_stable_listen_and_known_peers() {
    340         let alice = Wallet::from_seed("px-alice");
    341         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    342         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    343         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
    344             "127.0.0.1:9545".to_string(),
    345         ])));
    346         let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None);
    347         network.set_accept_inbound(true).await.unwrap();
    348 
    349         match network.peer_exchange().await {
    350             GossipEnvelope::PeerList { peers } => {
    351                 assert!(peers.contains(&"127.0.0.1:9544".to_string()));
    352                 assert!(peers.contains(&"127.0.0.1:9545".to_string()));
    353             }
    354             other => panic!("expected peer list, got {other:?}"),
    355         }
    356         network.set_accept_inbound(false).await.unwrap();
    357     }
    358 
    359     #[tokio::test]
    360     async fn peer_exchange_filters_announced_self_from_known_peers() {
    361         let alice = Wallet::from_seed("px-announced-self-alice");
    362         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    363         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    364         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
    365             "8.8.8.8:9444".to_string(),
    366             "8.8.4.4:9444".to_string(),
    367         ])));
    368         let network = gossip_network(
    369             node,
    370             peers,
    371             "127.0.0.1:0".parse().unwrap(),
    372             Some("8.8.8.8:9444".parse().unwrap()),
    373         );
    374         network.set_accept_inbound(true).await.unwrap();
    375 
    376         match network.peer_exchange().await {
    377             GossipEnvelope::PeerList { peers } => {
    378                 assert_eq!(
    379                     peers.iter().filter(|peer| *peer == "8.8.8.8:9444").count(),
    380                     1
    381                 );
    382                 assert!(peers.contains(&"8.8.4.4:9444".to_string()));
    383             }
    384             other => panic!("expected peer list, got {other:?}"),
    385         }
    386         network.set_accept_inbound(false).await.unwrap();
    387     }
    388 }