iuna

iuna

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

process.rs (18213B)


      1 use std::net::SocketAddr;
      2 
      3 use anyhow::Result;
      4 use sha2::{Digest, Sha256};
      5 use tokio::net::tcp::OwnedWriteHalf;
      6 
      7 use crate::{
      8     app::{GossipEnvelope, debug_logging_enabled},
      9     domain::{BurnBundle, Transaction},
     10 };
     11 
     12 use super::{
     13     GossipNetwork, MAX_BLOCK_BATCH, P2pMetricsCounters, apply_peer_list, forget_stale_self_peer,
     14     is_possible_fork_error, normalize_advertised_peer, peer_verification_response, process_hello,
     15     validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf, write_envelope,
     16 };
     17 
     18 pub(super) async fn respond_to_peer_verification_challenge(
     19     network: &GossipNetwork,
     20     writer: &mut OwnedWriteHalf,
     21     envelope: &GossipEnvelope,
     22 ) -> Result<bool> {
     23     let GossipEnvelope::PeerVerificationChallenge { address, nonce } = envelope else {
     24         return Ok(false);
     25     };
     26     if let Some(response) = peer_verification_response(network, address, nonce) {
     27         write_envelope(writer, &response).await?;
     28     }
     29     Ok(true)
     30 }
     31 
     32 pub(super) async fn process_envelope(
     33     network: &GossipNetwork,
     34     writer: &mut OwnedWriteHalf,
     35     remote_addr: SocketAddr,
     36     known_peer: &mut Option<String>,
     37     envelope: GossipEnvelope,
     38 ) -> Result<bool> {
     39     let mut requested_chain_data = false;
     40     match envelope {
     41         GossipEnvelope::Hello(hello) => {
     42             let _ = process_hello(network, remote_addr, known_peer, hello).await?;
     43         }
     44         GossipEnvelope::ChainBootstrapRequest => {
     45             let bootstrap = network.inner.node.lock().await.chain_bootstrap();
     46             write_envelope(writer, &GossipEnvelope::ChainBootstrap(bootstrap)).await?;
     47         }
     48         GossipEnvelope::BlockLocatorRequest { locator, limit } => {
     49             let blocks = network
     50                 .inner
     51                 .node
     52                 .lock()
     53                 .await
     54                 .blocks_after_locator(&locator, limit.min(MAX_BLOCK_BATCH));
     55             let blocks = super::byte_bounded_block_page(blocks);
     56             write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?;
     57         }
     58         GossipEnvelope::BlockRangeRequest { from_height, limit } => {
     59             let blocks = network
     60                 .inner
     61                 .node
     62                 .lock()
     63                 .await
     64                 .blocks_from(from_height, limit.min(MAX_BLOCK_BATCH));
     65             let blocks = super::byte_bounded_block_page(blocks);
     66             write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?;
     67         }
     68         GossipEnvelope::BlockRequest { hashes } => {
     69             let blocks = network.inner.node.lock().await.blocks_by_hash(&hashes);
     70             if !blocks.is_empty() {
     71                 let blocks = super::byte_bounded_block_page(blocks);
     72                 write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?;
     73             }
     74         }
     75         GossipEnvelope::PeerAnnouncement { address, node_id } => {
     76             let peer = normalize_advertised_peer(&address, remote_addr)?;
     77             if network.is_self_peer(&peer).await {
     78                 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections);
     79                 forget_stale_self_peer(network, known_peer).await;
     80             } else if node_id.is_some() && debug_logging_enabled() {
     81                 eprintln!("p2p peer announcement for {peer} ignored until hello verification");
     82             }
     83             let bootstrap = network.inner.node.lock().await.chain_bootstrap();
     84             write_envelope(writer, &GossipEnvelope::ChainBootstrap(bootstrap)).await?;
     85         }
     86         GossipEnvelope::PeerVerificationChallenge { address, nonce } => {
     87             if let Some(response) = peer_verification_response(network, &address, &nonce) {
     88                 write_envelope(writer, &response).await?;
     89             }
     90         }
     91         GossipEnvelope::PeerVerificationResponse { .. } => {}
     92         GossipEnvelope::PeerList { peers } => {
     93             apply_peer_list(network, remote_addr, peers).await?;
     94         }
     95         GossipEnvelope::Transaction(tx) => {
     96             process_transactions(network, remote_addr, known_peer, vec![tx]).await;
     97         }
     98         GossipEnvelope::Transactions { transactions } => {
     99             process_transactions(network, remote_addr, known_peer, transactions).await;
    100         }
    101         GossipEnvelope::TransactionV2 { envelope } => {
    102             process_transactions_v2(network, remote_addr, known_peer, vec![envelope]).await;
    103         }
    104         GossipEnvelope::TransactionsV2 { envelopes } => {
    105             process_transactions_v2(network, remote_addr, known_peer, envelopes).await;
    106         }
    107         GossipEnvelope::BurnBundle(bundle) => {
    108             process_burn_bundles(network, remote_addr, known_peer, vec![bundle]).await;
    109         }
    110         GossipEnvelope::BurnBundles { bundles } => {
    111             process_burn_bundles(network, remote_addr, known_peer, bundles).await;
    112         }
    113         GossipEnvelope::BurnBundleRequest {
    114             height,
    115             prev_hash,
    116             slots,
    117         } => {
    118             let bundles = network
    119                 .inner
    120                 .node
    121                 .lock()
    122                 .await
    123                 .burn_bundles_for_request(height, &prev_hash, &slots);
    124             if !bundles.is_empty() {
    125                 write_envelope(writer, &GossipEnvelope::BurnBundles { bundles }).await?;
    126             }
    127         }
    128         GossipEnvelope::Block(block) => {
    129             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    130             let (needs_vdf, base_tip, sync_generation) = {
    131                 let node = network.inner.node.lock().await;
    132                 (
    133                     node.block_requires_vdf_verification_at(&block, adjusted_time_ms),
    134                     node.ledger().tip_hash().to_string(),
    135                     network.sync_generation(),
    136                 )
    137             };
    138             let result = match needs_vdf {
    139                 Ok(false) => Ok(()),
    140                 Ok(true) => {
    141                     let validation_key =
    142                         chain_validation_key("block", &base_tip, std::slice::from_ref(&block));
    143                     let Some(_validation) = network.claim_chain_validation(validation_key).await
    144                     else {
    145                         return Ok(false);
    146                     };
    147                     let base_is_current = {
    148                         let node = network.inner.node.lock().await;
    149                         network.sync_generation_is_current(sync_generation)
    150                             && node.ledger().tip_hash() == base_tip
    151                     };
    152                     if !base_is_current {
    153                         return Ok(false);
    154                     }
    155                     match verify_block_vdf(block).await {
    156                         Ok(block) => {
    157                             let mut node = network.inner.node.lock().await;
    158                             if network.sync_generation_is_current(sync_generation)
    159                                 && node.ledger().tip_hash() == base_tip
    160                             {
    161                                 node.receive_preverified_block_at(block, adjusted_time_ms)
    162                             } else {
    163                                 Ok(())
    164                             }
    165                         }
    166                         Err(error) => Err(error),
    167                     }
    168                 }
    169                 Err(error) => Err(error),
    170             };
    171             let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error);
    172             record_rejected_chain_payload(
    173                 network,
    174                 &network.inner.metrics.rejected_blocks,
    175                 "block",
    176                 &result,
    177             );
    178             record_inbound_result(network, known_peer, remote_addr, result).await;
    179             if request_locator {
    180                 request_fork_blocks(network, writer).await?;
    181                 requested_chain_data = true;
    182             }
    183             network.forward_outbox().await;
    184         }
    185         GossipEnvelope::Blocks { blocks } => {
    186             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    187             let (base_tip, sync_generation) = {
    188                 let node = network.inner.node.lock().await;
    189                 if blocks.is_empty() || node.ledger().contains_block_sequence(&blocks) {
    190                     drop(node);
    191                     record_inbound_result(network, known_peer, remote_addr, Ok(())).await;
    192                     return Ok(false);
    193                 }
    194                 (
    195                     node.ledger().tip_hash().to_string(),
    196                     network.sync_generation(),
    197                 )
    198             };
    199             let validation_key = chain_validation_key("blocks", &base_tip, &blocks);
    200             let Some(_validation) = network.claim_chain_validation(validation_key).await else {
    201                 return Ok(false);
    202             };
    203             let (local_ledger, progress_guard) = {
    204                 let node = network.inner.node.lock().await;
    205                 if !network.sync_generation_is_current(sync_generation)
    206                     || node.ledger().tip_hash() != base_tip
    207                     || node.ledger().contains_block_sequence(&blocks)
    208                 {
    209                     drop(node);
    210                     record_inbound_result(network, known_peer, remote_addr, Ok(())).await;
    211                     return Ok(false);
    212                 }
    213                 let start_height = blocks
    214                     .first()
    215                     .map(|block| block.height.saturating_sub(1))
    216                     .unwrap_or_else(|| node.ledger().height());
    217                 let target_height = blocks
    218                     .last()
    219                     .map(|block| block.height)
    220                     .unwrap_or(start_height);
    221                 let progress_guard = network.begin_sync_progress(start_height, target_height);
    222                 (node.clone_ledger(), progress_guard)
    223             };
    224             let progress_id = progress_guard.id();
    225             let progress_network = network.clone();
    226             let result = match validate_blocks_extension(
    227                 local_ledger,
    228                 blocks,
    229                 adjusted_time_ms,
    230                 move |height| progress_network.update_sync_progress(progress_id, height),
    231             )
    232             .await
    233             {
    234                 Ok(ledger) => {
    235                     let mut node = network.inner.node.lock().await;
    236                     if progress_guard.is_current()
    237                         && network.sync_generation_is_current(sync_generation)
    238                         && node.ledger().tip_hash() == base_tip
    239                     {
    240                         node.import_verified_ledger(ledger).map(|_| ())
    241                     } else {
    242                         Ok(())
    243                     }
    244                 }
    245                 Err(error) => Err(error),
    246             };
    247             drop(progress_guard);
    248             let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error);
    249             record_rejected_chain_payload(
    250                 network,
    251                 &network.inner.metrics.rejected_block_batches,
    252                 "block batch",
    253                 &result,
    254             );
    255             record_inbound_result(network, known_peer, remote_addr, result).await;
    256             if request_locator {
    257                 request_fork_blocks(network, writer).await?;
    258                 requested_chain_data = true;
    259             }
    260             network.forward_outbox().await;
    261         }
    262         GossipEnvelope::ChainBootstrap(bootstrap) => {
    263             let sync_generation = network.sync_generation();
    264             let (base_tip, expected_profile_id) = {
    265                 let node = network.inner.node.lock().await;
    266                 (
    267                     node.ledger().tip_hash().to_string(),
    268                     node.ledger().launch_profile().profile_id.clone(),
    269                 )
    270             };
    271             let validation_key = chain_validation_key(
    272                 "bootstrap",
    273                 &base_tip,
    274                 std::slice::from_ref(&bootstrap.genesis_block),
    275             );
    276             let Some(_validation) = network.claim_chain_validation(validation_key).await else {
    277                 return Ok(false);
    278             };
    279             let base_is_current = {
    280                 let node = network.inner.node.lock().await;
    281                 network.sync_generation_is_current(sync_generation)
    282                     && node.ledger().tip_hash() == base_tip
    283             };
    284             if !base_is_current {
    285                 return Ok(false);
    286             }
    287             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    288             let result =
    289                 match validate_chain_bootstrap(&expected_profile_id, bootstrap, adjusted_time_ms)
    290                     .await
    291                 {
    292                     Ok(ledger) => {
    293                         let mut node = network.inner.node.lock().await;
    294                         if network.sync_generation_is_current(sync_generation)
    295                             && node.ledger().tip_hash() == base_tip
    296                         {
    297                             node.import_verified_ledger(ledger).map(|_| ())
    298                         } else {
    299                             Ok(())
    300                         }
    301                     }
    302                     Err(error) => Err(error),
    303                 };
    304             record_rejected_chain_payload(
    305                 network,
    306                 &network.inner.metrics.rejected_snapshots,
    307                 "chain bootstrap",
    308                 &result,
    309             );
    310             record_inbound_result(network, known_peer, remote_addr, result).await;
    311             network.forward_outbox().await;
    312         }
    313         other => {
    314             let result = network.inner.node.lock().await.receive(other);
    315             record_inbound_result(network, known_peer, remote_addr, result).await;
    316             network.forward_outbox().await;
    317         }
    318     }
    319     Ok(requested_chain_data)
    320 }
    321 
    322 fn chain_validation_key(kind: &str, base_tip: &str, blocks: &[crate::domain::Block]) -> String {
    323     let mut digest = Sha256::new();
    324     digest.update(kind.as_bytes());
    325     digest.update([0]);
    326     digest.update(base_tip.as_bytes());
    327     for block in blocks {
    328         digest.update(block.height.to_be_bytes());
    329         digest.update(block.prev_hash.as_bytes());
    330         digest.update([0]);
    331         digest.update(block.hash.as_bytes());
    332         digest.update([0]);
    333     }
    334     format!("{:x}", digest.finalize())
    335 }
    336 
    337 async fn request_fork_blocks(network: &GossipNetwork, writer: &mut OwnedWriteHalf) -> Result<()> {
    338     let locator = network.inner.node.lock().await.block_locator();
    339     write_envelope(
    340         writer,
    341         &GossipEnvelope::BlockLocatorRequest {
    342             locator,
    343             limit: MAX_BLOCK_BATCH,
    344         },
    345     )
    346     .await
    347 }
    348 
    349 async fn process_transactions(
    350     network: &GossipNetwork,
    351     remote_addr: SocketAddr,
    352     known_peer: &Option<String>,
    353     transactions: Vec<Transaction>,
    354 ) {
    355     let first_error = {
    356         let mut node = network.inner.node.lock().await;
    357         let mut first_error = None;
    358         for tx in transactions {
    359             if let Err(error) = node.receive_gossiped_transaction(tx) {
    360                 first_error.get_or_insert(error);
    361             }
    362         }
    363         first_error
    364     };
    365     record_inbound_result(
    366         network,
    367         known_peer,
    368         remote_addr,
    369         first_error.map(Err).unwrap_or(Ok(())),
    370     )
    371     .await;
    372     network.forward_outbox().await;
    373 }
    374 
    375 async fn process_transactions_v2(
    376     network: &GossipNetwork,
    377     remote_addr: SocketAddr,
    378     known_peer: &Option<String>,
    379     envelopes: Vec<String>,
    380 ) {
    381     let first_error = {
    382         let mut node = network.inner.node.lock().await;
    383         let mut first_error = None;
    384         for envelope in envelopes {
    385             if let Err(error) = node.receive_gossiped_transaction_v2(envelope) {
    386                 first_error.get_or_insert(error);
    387             }
    388         }
    389         first_error
    390     };
    391     record_inbound_result(
    392         network,
    393         known_peer,
    394         remote_addr,
    395         first_error.map(Err).unwrap_or(Ok(())),
    396     )
    397     .await;
    398     network.forward_outbox().await;
    399 }
    400 
    401 async fn process_burn_bundles(
    402     network: &GossipNetwork,
    403     remote_addr: SocketAddr,
    404     known_peer: &Option<String>,
    405     bundles: Vec<BurnBundle>,
    406 ) {
    407     let first_error = {
    408         let mut node = network.inner.node.lock().await;
    409         let mut first_error = None;
    410         for bundle in bundles {
    411             if let Err(error) = node.receive_burn_bundle(bundle) {
    412                 first_error.get_or_insert(error);
    413             }
    414         }
    415         first_error
    416     };
    417     record_inbound_result(
    418         network,
    419         known_peer,
    420         remote_addr,
    421         first_error.map(Err).unwrap_or(Ok(())),
    422     )
    423     .await;
    424     network.forward_outbox().await;
    425 }
    426 
    427 fn record_rejected_chain_payload(
    428     network: &GossipNetwork,
    429     counter: &std::sync::atomic::AtomicU64,
    430     kind: &str,
    431     result: &Result<()>,
    432 ) {
    433     let Err(error) = result else {
    434         return;
    435     };
    436     P2pMetricsCounters::inc(counter);
    437     P2pMetricsCounters::set_last(
    438         &network.inner.metrics.last_chain_payload_error,
    439         format!("{kind}: {error:#}"),
    440     );
    441 }
    442 
    443 async fn record_inbound_result(
    444     network: &GossipNetwork,
    445     known_peer: &Option<String>,
    446     remote_addr: SocketAddr,
    447     result: Result<()>,
    448 ) {
    449     let peer = known_peer
    450         .clone()
    451         .unwrap_or_else(|| remote_addr.to_string());
    452     match result {
    453         Ok(()) => {
    454             if known_peer.is_some() {
    455                 network.inner.peers.lock().await.record_received(&peer, 1);
    456             }
    457         }
    458         Err(error) => {
    459             let message = format!("{error:#}");
    460             let mut peers = network.inner.peers.lock().await;
    461             if super::inbound_error_counts_as_misbehavior(&error) {
    462                 if known_peer.is_some() {
    463                     peers.record_misbehavior(&peer, message.clone());
    464                 } else {
    465                     peers.record_inbound_misbehavior(&peer, message.clone());
    466                 }
    467             } else {
    468                 peers.record_inbound_error(&peer, message.clone());
    469             }
    470             if debug_logging_enabled() {
    471                 eprintln!("p2p envelope from {peer} ignored: {message}");
    472             }
    473         }
    474     }
    475 }