iuna

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

handshake.rs (18563B)


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