iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

handshake.rs (19626B)


      1 use std::{collections::BTreeMap, net::SocketAddr};
      2 
      3 use anyhow::Result;
      4 use tokio::{
      5     net::{TcpStream, tcp::OwnedWriteHalf},
      6     time::timeout,
      7 };
      8 
      9 use super::identity::{
     10     new_verification_nonce, peer_verification_response, peer_verification_response_is_valid,
     11 };
     12 use super::line_codec::{LimitedLineReader, parse_envelope, read_session_envelope};
     13 use super::metrics::P2pMetricsCounters;
     14 use super::peer_addr::{advertised_peer_is_discoverable, normalize_advertised_peer};
     15 use super::{
     16     CONNECT_TIMEOUT, GossipNetwork, HANDSHAKE_TIMEOUT, MAX_PEER_VERIFICATION_ENVELOPES, PeerStatus,
     17     write_envelope,
     18 };
     19 use crate::{
     20     app::{
     21         GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, PeerDirection, ProtocolHello,
     22         debug_logging_enabled, now_ms, validate_protocol_capabilities,
     23         validate_transaction_v2_peer_capability,
     24     },
     25     domain::Ledger,
     26 };
     27 
     28 pub(super) struct PeerVerificationSession<'a> {
     29     pub(super) writer: &'a mut OwnedWriteHalf,
     30     pub(super) reader: &'a mut LimitedLineReader<tokio::net::tcp::OwnedReadHalf>,
     31     pub(super) connection_label: &'a str,
     32 }
     33 
     34 pub(super) async fn record_peer_status(
     35     network: &GossipNetwork,
     36     known_peer: &Option<String>,
     37     remote_addr: SocketAddr,
     38     peer_status: &PeerStatus,
     39 ) {
     40     let local_receive_time_ms = now_ms();
     41     if let Some(peer) = known_peer {
     42         let mut peers = network.inner.peers.lock().await;
     43         peers.record_status(peer, peer_status.height, peer_status.tip_hash.clone());
     44         peers.record_clock_observation(
     45             peer,
     46             PeerDirection::Outbound,
     47             peer_status.time_ms,
     48             local_receive_time_ms,
     49         );
     50     } else {
     51         let peer = remote_addr.to_string();
     52         let mut peers = network.inner.peers.lock().await;
     53         peers.record_clock_observation(
     54             &peer,
     55             PeerDirection::Inbound,
     56             peer_status.time_ms,
     57             local_receive_time_ms,
     58         );
     59         peers.record_received(&peer, 1);
     60     }
     61 }
     62 
     63 async fn record_peer_hello(
     64     network: &GossipNetwork,
     65     known_peer: &Option<String>,
     66     remote_addr: SocketAddr,
     67     hello: ProtocolHello,
     68 ) {
     69     let (peer, direction) = match known_peer {
     70         Some(peer) => (peer.clone(), PeerDirection::Outbound),
     71         None => (remote_addr.to_string(), PeerDirection::Inbound),
     72     };
     73     network
     74         .inner
     75         .peers
     76         .lock()
     77         .await
     78         .record_hello(&peer, direction, hello);
     79 }
     80 
     81 pub(super) async fn process_hello(
     82     network: &GossipNetwork,
     83     remote_addr: SocketAddr,
     84     known_peer: &mut Option<String>,
     85     hello: ProtocolHello,
     86 ) -> Result<PeerStatus> {
     87     process_hello_inner(network, None, remote_addr, known_peer, hello).await
     88 }
     89 
     90 pub(super) async fn process_hello_with_verification(
     91     network: &GossipNetwork,
     92     writer: &mut OwnedWriteHalf,
     93     reader: &mut LimitedLineReader<tokio::net::tcp::OwnedReadHalf>,
     94     connection_label: &str,
     95     remote_addr: SocketAddr,
     96     known_peer: &mut Option<String>,
     97     hello: ProtocolHello,
     98 ) -> Result<PeerStatus> {
     99     let mut verification_session = PeerVerificationSession {
    100         writer,
    101         reader,
    102         connection_label,
    103     };
    104     process_hello_inner(
    105         network,
    106         Some(&mut verification_session),
    107         remote_addr,
    108         known_peer,
    109         hello,
    110     )
    111     .await
    112 }
    113 
    114 async fn process_hello_inner(
    115     network: &GossipNetwork,
    116     mut verification_session: Option<&mut PeerVerificationSession<'_>>,
    117     remote_addr: SocketAddr,
    118     known_peer: &mut Option<String>,
    119     hello: ProtocolHello,
    120 ) -> Result<PeerStatus> {
    121     if hello.protocol_version != PROTOCOL_VERSION {
    122         anyhow::bail!(
    123             "unsupported protocol version {}; expected {}",
    124             hello.protocol_version,
    125             PROTOCOL_VERSION
    126         );
    127     }
    128     validate_protocol_capabilities(&hello.capabilities)?;
    129     if hello.network_id != NETWORK_ID {
    130         anyhow::bail!(
    131             "wrong network {}; expected {}",
    132             hello.network_id,
    133             NETWORK_ID
    134         );
    135     }
    136     if hello
    137         .node_id
    138         .as_deref()
    139         .is_some_and(|node_id| node_id == network.inner.node_id)
    140     {
    141         P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections);
    142         forget_stale_self_peer(network, known_peer).await;
    143         return Ok(PeerStatus::rejected(
    144             hello.height,
    145             hello.tip_hash,
    146             hello.time_ms,
    147         ));
    148     }
    149     let (local_genesis, local_accepts_remote_genesis, local_height) = {
    150         let node = network.inner.node.lock().await;
    151         (
    152             node.ledger().genesis_hash().to_string(),
    153             node.ledger().is_setup_placeholder(),
    154             node.ledger().height(),
    155         )
    156     };
    157     validate_transaction_v2_peer_capability(&hello.capabilities, local_height, hello.height)?;
    158     let peer_capabilities = hello.capabilities.clone();
    159     let genesis_mismatch = hello.genesis_hash != local_genesis;
    160     let remote_is_setup_placeholder =
    161         hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash();
    162     let request_bootstrap = genesis_mismatch && local_accepts_remote_genesis;
    163     if genesis_mismatch && !local_accepts_remote_genesis && !remote_is_setup_placeholder {
    164         anyhow::bail!(
    165             "wrong genesis {}; expected {local_genesis}",
    166             hello.genesis_hash
    167         );
    168     }
    169 
    170     let remote_node_id = hello.node_id.clone();
    171     let mut reject_session = false;
    172     if let Some(listen_addr) = &hello.listen_addr {
    173         let peer = normalize_advertised_peer(listen_addr, remote_addr)?;
    174         if network.is_self_peer(&peer).await {
    175             P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections);
    176             forget_stale_self_peer(network, known_peer).await;
    177             reject_session = true;
    178         } else {
    179             let verified = match verification_session.as_mut() {
    180                 Some(session) => {
    181                     remember_verified_advertised_peer(
    182                         network,
    183                         session,
    184                         remote_addr,
    185                         known_peer,
    186                         peer.clone(),
    187                         remote_node_id.as_deref(),
    188                     )
    189                     .await?
    190                 }
    191                 None => false,
    192             };
    193             if !verified && debug_logging_enabled() {
    194                 eprintln!(
    195                     "p2p advertised address {peer} ignored because ownership was not verified"
    196                 );
    197             }
    198         }
    199     }
    200     record_peer_status(
    201         network,
    202         known_peer,
    203         remote_addr,
    204         &PeerStatus::with_time(hello.height, hello.tip_hash.clone(), hello.time_ms),
    205     )
    206     .await;
    207     record_peer_hello(network, known_peer, remote_addr, hello.clone()).await;
    208     let mut status = if request_bootstrap {
    209         PeerStatus::with_bootstrap_request(hello.height, hello.tip_hash, hello.time_ms)
    210             .with_capabilities(peer_capabilities)
    211     } else {
    212         PeerStatus::with_time(hello.height, hello.tip_hash, hello.time_ms)
    213             .with_capabilities(peer_capabilities)
    214     };
    215     status.reject_session = reject_session;
    216     Ok(status)
    217 }
    218 
    219 fn setup_placeholder_genesis_hash() -> String {
    220     Ledger::new(BTreeMap::new(), 1).genesis_hash().to_string()
    221 }
    222 
    223 async fn remember_verified_advertised_peer(
    224     network: &GossipNetwork,
    225     session: &mut PeerVerificationSession<'_>,
    226     remote_addr: SocketAddr,
    227     known_peer: &mut Option<String>,
    228     peer: String,
    229     expected_node_id: Option<&str>,
    230 ) -> Result<bool> {
    231     if !advertised_peer_is_discoverable(&peer, remote_addr)? {
    232         return Ok(false);
    233     }
    234     if known_peer.as_deref() != Some(peer.as_str()) {
    235         let Some(expected_node_id) = expected_node_id else {
    236             return Ok(false);
    237         };
    238         if !verify_connected_peer_node_id(network, session, &peer, expected_node_id).await? {
    239             return Ok(false);
    240         }
    241         if !verify_advertised_peer_node_id(network, &peer, expected_node_id).await {
    242             return Ok(false);
    243         }
    244     }
    245     remember_discoverable_advertised_peer(network, remote_addr, known_peer, peer).await
    246 }
    247 
    248 async fn verify_connected_peer_node_id(
    249     network: &GossipNetwork,
    250     session: &mut PeerVerificationSession<'_>,
    251     peer: &str,
    252     expected_node_id: &str,
    253 ) -> Result<bool> {
    254     let nonce = new_verification_nonce();
    255     write_envelope(
    256         session.writer,
    257         &GossipEnvelope::PeerVerificationChallenge {
    258             address: peer.to_string(),
    259             nonce: nonce.clone(),
    260         },
    261     )
    262     .await?;
    263 
    264     for _ in 0..MAX_PEER_VERIFICATION_ENVELOPES {
    265         let envelope = match timeout(
    266             HANDSHAKE_TIMEOUT,
    267             read_session_envelope(network, session.connection_label, session.reader),
    268         )
    269         .await
    270         {
    271             Ok(Ok(Some(envelope))) => envelope,
    272             Ok(Ok(None)) | Err(_) => return Ok(false),
    273             Ok(Err(error)) => return Err(error),
    274         };
    275         match envelope {
    276             GossipEnvelope::PeerVerificationResponse {
    277                 address,
    278                 nonce: response_nonce,
    279                 node_id,
    280                 signature,
    281             } => {
    282                 return Ok(peer_verification_response_is_valid(
    283                     &address,
    284                     &response_nonce,
    285                     &node_id,
    286                     &signature,
    287                     peer,
    288                     &nonce,
    289                     expected_node_id,
    290                 ));
    291             }
    292             GossipEnvelope::PeerVerificationChallenge { address, nonce } => {
    293                 if let Some(response) = peer_verification_response(network, &address, &nonce) {
    294                     write_envelope(session.writer, &response).await?;
    295                 }
    296             }
    297             _ => {}
    298         }
    299     }
    300     Ok(false)
    301 }
    302 
    303 pub(super) async fn verify_advertised_peer_node_id(
    304     network: &GossipNetwork,
    305     peer: &str,
    306     expected_node_id: &str,
    307 ) -> bool {
    308     let stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(peer)).await {
    309         Ok(Ok(stream)) => stream,
    310         Ok(Err(error)) => {
    311             if debug_logging_enabled() {
    312                 eprintln!("p2p announced address {peer} failed verification: {error}");
    313             }
    314             return false;
    315         }
    316         Err(_) => {
    317             if debug_logging_enabled() {
    318                 eprintln!("p2p announced address {peer} failed verification: timeout");
    319             }
    320             return false;
    321         }
    322     };
    323     let (reader, mut writer) = stream.into_split();
    324     let mut reader = LimitedLineReader::new(reader);
    325     let line = match timeout(HANDSHAKE_TIMEOUT, reader.read_line()).await {
    326         Ok(Ok(Some(line))) => line,
    327         Ok(Ok(None)) => return false,
    328         Ok(Err(error)) => {
    329             if debug_logging_enabled() {
    330                 eprintln!(
    331                     "p2p announced address {peer} sent invalid verification hello: {error:#}"
    332                 );
    333             }
    334             return false;
    335         }
    336         Err(_) => return false,
    337     };
    338     let hello = match parse_envelope(&line) {
    339         Ok(GossipEnvelope::Hello(hello)) => hello,
    340         Ok(_) | Err(_) => return false,
    341     };
    342 
    343     if !advertised_peer_hello_is_compatible(network, &hello).await
    344         || hello.node_id.as_deref() != Some(expected_node_id)
    345     {
    346         return false;
    347     }
    348 
    349     let nonce = new_verification_nonce();
    350     if write_envelope(
    351         &mut writer,
    352         &GossipEnvelope::PeerVerificationChallenge {
    353             address: peer.to_string(),
    354             nonce: nonce.clone(),
    355         },
    356     )
    357     .await
    358     .is_err()
    359     {
    360         return false;
    361     }
    362     for _ in 0..MAX_PEER_VERIFICATION_ENVELOPES {
    363         let line = match timeout(HANDSHAKE_TIMEOUT, reader.read_line()).await {
    364             Ok(Ok(Some(line))) => line,
    365             Ok(Ok(None)) | Ok(Err(_)) | Err(_) => return false,
    366         };
    367         let envelope = match parse_envelope(&line) {
    368             Ok(envelope) => envelope,
    369             Err(_) => return false,
    370         };
    371         if let GossipEnvelope::PeerVerificationResponse {
    372             address,
    373             nonce: response_nonce,
    374             node_id,
    375             signature,
    376         } = envelope
    377         {
    378             return peer_verification_response_is_valid(
    379                 &address,
    380                 &response_nonce,
    381                 &node_id,
    382                 &signature,
    383                 peer,
    384                 &nonce,
    385                 expected_node_id,
    386             );
    387         }
    388     }
    389     false
    390 }
    391 
    392 async fn advertised_peer_hello_is_compatible(
    393     network: &GossipNetwork,
    394     hello: &ProtocolHello,
    395 ) -> bool {
    396     if hello.protocol_version != PROTOCOL_VERSION
    397         || hello.network_id != NETWORK_ID
    398         || validate_protocol_capabilities(&hello.capabilities).is_err()
    399     {
    400         return false;
    401     }
    402     let (local_genesis, local_accepts_remote_genesis, local_height) = {
    403         let node = network.inner.node.lock().await;
    404         (
    405             node.ledger().genesis_hash().to_string(),
    406             node.ledger().is_setup_placeholder(),
    407             node.ledger().height(),
    408         )
    409     };
    410     if validate_transaction_v2_peer_capability(&hello.capabilities, local_height, hello.height)
    411         .is_err()
    412     {
    413         return false;
    414     }
    415     let remote_is_setup_placeholder =
    416         hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash();
    417     hello.genesis_hash == local_genesis
    418         || local_accepts_remote_genesis
    419         || remote_is_setup_placeholder
    420 }
    421 
    422 pub(super) async fn remember_discoverable_advertised_peer(
    423     network: &GossipNetwork,
    424     remote_addr: SocketAddr,
    425     known_peer: &mut Option<String>,
    426     peer: String,
    427 ) -> Result<bool> {
    428     if !advertised_peer_is_discoverable(&peer, remote_addr)? {
    429         return Ok(false);
    430     }
    431     if let Some(previous_peer) = known_peer.as_deref() {
    432         network
    433             .inner
    434             .peers
    435             .lock()
    436             .await
    437             .replace_peer_address(previous_peer, peer.clone());
    438     } else {
    439         network
    440             .inner
    441             .peers
    442             .lock()
    443             .await
    444             .add_discovered_peer(peer.clone());
    445     }
    446     *known_peer = Some(peer);
    447     Ok(true)
    448 }
    449 
    450 pub(super) async fn forget_stale_self_peer(
    451     network: &GossipNetwork,
    452     known_peer: &mut Option<String>,
    453 ) {
    454     if let Some(previous_peer) = known_peer.take() {
    455         network.inner.peers.lock().await.remove_peer(&previous_peer);
    456     }
    457 }
    458 
    459 #[cfg(test)]
    460 mod tests {
    461     use std::sync::Arc;
    462 
    463     use crate::{
    464         app::{PeerBook, PeerDirection},
    465         domain::Wallet,
    466     };
    467 
    468     use super::super::{
    469         PeerStatus,
    470         test_support::{allocations, gossip_network, node},
    471     };
    472     use super::{
    473         forget_stale_self_peer, record_peer_status, remember_discoverable_advertised_peer,
    474     };
    475 
    476     #[tokio::test]
    477     async fn inbound_announced_address_replaces_gateway_address_for_ui() {
    478         let alice = Wallet::from_seed("hello-public-announced-inbound-alice");
    479         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    480         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    481         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    482         let network = gossip_network(
    483             node,
    484             Arc::clone(&peers),
    485             "0.0.0.0:9444".parse().unwrap(),
    486             None,
    487         );
    488         let mut known_peer = None;
    489 
    490         let remembered = remember_discoverable_advertised_peer(
    491             &network,
    492             "10.42.0.1:51234".parse().unwrap(),
    493             &mut known_peer,
    494             "142.132.164.59:9444".to_string(),
    495         )
    496         .await
    497         .unwrap();
    498         record_peer_status(
    499             &network,
    500             &known_peer,
    501             "10.42.0.1:51234".parse().unwrap(),
    502             &PeerStatus::with_time(7, "tip".to_string(), 1_000),
    503         )
    504         .await;
    505 
    506         assert!(remembered);
    507         assert_eq!(known_peer.as_deref(), Some("142.132.164.59:9444"));
    508         let listed = peers.lock().await.list();
    509         assert_eq!(listed.len(), 1);
    510         let peer = &listed[0];
    511         assert_eq!(peer.address, "142.132.164.59:9444");
    512         assert_eq!(peer.direction, PeerDirection::Outbound);
    513         assert_eq!(peer.last_known_height, Some(7));
    514         assert_eq!(peer.messages_received, 0);
    515 
    516         let repeated = remember_discoverable_advertised_peer(
    517             &network,
    518             "10.42.0.1:51234".parse().unwrap(),
    519             &mut known_peer,
    520             "142.132.164.59:9444".to_string(),
    521         )
    522         .await
    523         .unwrap();
    524 
    525         assert!(repeated);
    526         assert_eq!(
    527             peers.lock().await.list()[0].direction,
    528             PeerDirection::Outbound
    529         );
    530     }
    531 
    532     #[tokio::test]
    533     async fn inbound_status_does_not_create_outbound_ephemeral_peer() {
    534         let alice = Wallet::from_seed("inbound-status-alice");
    535         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    536         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    537         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    538         let network = gossip_network(
    539             node,
    540             Arc::clone(&peers),
    541             "127.0.0.1:9544".parse().unwrap(),
    542             None,
    543         );
    544 
    545         record_peer_status(
    546             &network,
    547             &None,
    548             "127.0.0.1:51729".parse().unwrap(),
    549             &PeerStatus::new(4, "tip".to_string()),
    550         )
    551         .await;
    552 
    553         let peers = peers.lock().await;
    554         assert!(peers.addresses().is_empty());
    555         let listed = peers.list();
    556         assert_eq!(listed.len(), 1);
    557         assert_eq!(listed[0].direction, PeerDirection::Inbound);
    558     }
    559 
    560     #[tokio::test]
    561     async fn peer_announcement_ignores_private_ephemeral_address() {
    562         let alice = Wallet::from_seed("px-private-announcement-alice");
    563         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    564         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    565         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    566         let network = gossip_network(
    567             node,
    568             Arc::clone(&peers),
    569             "0.0.0.0:9444".parse().unwrap(),
    570             None,
    571         );
    572         let mut known_peer = None;
    573 
    574         let remembered = remember_discoverable_advertised_peer(
    575             &network,
    576             "142.132.164.59:51234".parse().unwrap(),
    577             &mut known_peer,
    578             "10.42.1.1:10091".to_string(),
    579         )
    580         .await
    581         .unwrap();
    582 
    583         assert!(!remembered);
    584         assert!(known_peer.is_none());
    585         assert!(peers.lock().await.addresses().is_empty());
    586     }
    587 
    588     #[tokio::test]
    589     async fn peer_announcement_removes_outbound_peer_that_announces_self_address() {
    590         let alice = Wallet::from_seed("px-self-announcement-alice");
    591         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    592         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    593         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
    594             "10.42.1.1:30508".to_string(),
    595         ])));
    596         let network = gossip_network(
    597             node,
    598             Arc::clone(&peers),
    599             "0.0.0.0:9444".parse().unwrap(),
    600             None,
    601         );
    602         let mut known_peer = Some("10.42.1.1:30508".to_string());
    603 
    604         forget_stale_self_peer(&network, &mut known_peer).await;
    605 
    606         assert_eq!(known_peer, None);
    607         assert!(peers.lock().await.addresses().is_empty());
    608     }
    609 }