iuna

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

process.rs (12115B)


      1 use std::net::SocketAddr;
      2 
      3 use anyhow::{Result, anyhow};
      4 use tokio::net::tcp::OwnedWriteHalf;
      5 
      6 use crate::{
      7     app::{GossipEnvelope, debug_logging_enabled},
      8     domain::{BlindedReveal, BlindedTransaction, RevealBundle, Transaction},
      9 };
     10 
     11 use super::{
     12     GossipNetwork, MAX_BLOCK_BATCH, P2pMetricsCounters, apply_peer_list, forget_stale_self_peer,
     13     is_possible_fork_error, normalize_advertised_peer, peer_verification_response, process_hello,
     14     validate_blocks_extension, validate_snapshot_extension, verify_block_vdf, write_envelope,
     15     write_payload,
     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<()> {
     39     match envelope {
     40         GossipEnvelope::Hello(hello) => {
     41             let _ = process_hello(network, remote_addr, known_peer, hello).await?;
     42         }
     43         GossipEnvelope::ChainSnapshotRequest => {
     44             let snapshot = network.inner.node.lock().await.chain_snapshot();
     45             write_envelope(writer, &GossipEnvelope::ChainSnapshot(snapshot)).await?;
     46         }
     47         GossipEnvelope::BlockRangeRequest { from_height, limit } => {
     48             let blocks = network
     49                 .inner
     50                 .node
     51                 .lock()
     52                 .await
     53                 .blocks_from(from_height, limit.min(MAX_BLOCK_BATCH));
     54             write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?;
     55         }
     56         GossipEnvelope::BlockRequest { hashes } => {
     57             let blocks = network.inner.node.lock().await.blocks_by_hash(&hashes);
     58             if !blocks.is_empty() {
     59                 write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?;
     60             }
     61         }
     62         GossipEnvelope::Inventory { blocks } => {
     63             let requests = network
     64                 .inner
     65                 .node
     66                 .lock()
     67                 .await
     68                 .missing_inventory_requests(&blocks);
     69             write_payload(writer, &requests).await?;
     70         }
     71         GossipEnvelope::PeerAnnouncement { address, node_id } => {
     72             let peer = normalize_advertised_peer(&address, remote_addr)?;
     73             if network.is_self_peer(&peer).await {
     74                 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections);
     75                 forget_stale_self_peer(network, known_peer).await;
     76             } else if node_id.is_some() && debug_logging_enabled() {
     77                 eprintln!("p2p peer announcement for {peer} ignored until hello verification");
     78             }
     79             let snapshot = network.inner.node.lock().await.chain_snapshot();
     80             write_envelope(writer, &GossipEnvelope::ChainSnapshot(snapshot)).await?;
     81         }
     82         GossipEnvelope::PeerVerificationChallenge { address, nonce } => {
     83             if let Some(response) = peer_verification_response(network, &address, &nonce) {
     84                 write_envelope(writer, &response).await?;
     85             }
     86         }
     87         GossipEnvelope::PeerVerificationResponse { .. } => {}
     88         GossipEnvelope::PeerList { peers } => {
     89             apply_peer_list(network, remote_addr, peers).await?;
     90         }
     91         GossipEnvelope::BlindedTransaction(tx) => {
     92             process_blinded_transactions(network, remote_addr, known_peer, vec![tx]).await;
     93         }
     94         GossipEnvelope::BlindedTransactions { transactions } => {
     95             process_blinded_transactions(network, remote_addr, known_peer, transactions).await;
     96         }
     97         GossipEnvelope::MineAction(tx) => {
     98             process_mine_actions(network, remote_addr, known_peer, vec![tx]).await;
     99         }
    100         GossipEnvelope::MineActions { transactions } => {
    101             process_mine_actions(network, remote_addr, known_peer, transactions).await;
    102         }
    103         GossipEnvelope::BlindedReveal(reveal) => {
    104             process_blinded_reveals(network, remote_addr, known_peer, vec![reveal]).await;
    105         }
    106         GossipEnvelope::BlindedReveals { reveals } => {
    107             process_blinded_reveals(network, remote_addr, known_peer, reveals).await;
    108         }
    109         GossipEnvelope::RevealBundle(bundle) => {
    110             process_reveal_bundles(network, remote_addr, known_peer, vec![bundle]).await;
    111         }
    112         GossipEnvelope::RevealBundles { bundles } => {
    113             process_reveal_bundles(network, remote_addr, known_peer, bundles).await;
    114         }
    115         GossipEnvelope::Block(block) => {
    116             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    117             let needs_vdf = {
    118                 let node = network.inner.node.lock().await;
    119                 node.block_requires_vdf_verification_at(&block, adjusted_time_ms)
    120             };
    121             let result = match needs_vdf {
    122                 Ok(false) => Ok(()),
    123                 Ok(true) => match verify_block_vdf(block).await {
    124                     Ok(block) => network
    125                         .inner
    126                         .node
    127                         .lock()
    128                         .await
    129                         .receive_preverified_block_at(block, adjusted_time_ms),
    130                     Err(error) => Err(error),
    131                 },
    132                 Err(error) => Err(error),
    133             };
    134             let request_snapshot = result.as_ref().err().is_some_and(is_possible_fork_error);
    135             record_inbound_result(network, known_peer, remote_addr, result).await;
    136             if request_snapshot {
    137                 write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?;
    138             }
    139             network.forward_outbox().await;
    140         }
    141         GossipEnvelope::Blocks { blocks } => {
    142             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    143             let local_ledger = network.inner.node.lock().await.clone_ledger();
    144             let result =
    145                 match validate_blocks_extension(local_ledger, blocks, adjusted_time_ms).await {
    146                     Ok(ledger) => network
    147                         .inner
    148                         .node
    149                         .lock()
    150                         .await
    151                         .import_verified_ledger(ledger)
    152                         .map(|_| ()),
    153                     Err(error) => Err(error),
    154                 };
    155             let request_snapshot = result.as_ref().err().is_some_and(is_possible_fork_error);
    156             record_inbound_result(network, known_peer, remote_addr, result).await;
    157             if request_snapshot {
    158                 write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?;
    159             }
    160             network.forward_outbox().await;
    161         }
    162         GossipEnvelope::ChainSnapshot(snapshot) => {
    163             let adjusted_time_ms = super::network_adjusted_time_ms(network).await;
    164             let local_ledger = network.inner.node.lock().await.clone_ledger();
    165             let result =
    166                 match validate_snapshot_extension(local_ledger, snapshot, adjusted_time_ms).await {
    167                     Ok(ledger) => network
    168                         .inner
    169                         .node
    170                         .lock()
    171                         .await
    172                         .import_verified_ledger(ledger)
    173                         .map(|_| ()),
    174                     Err(error) => Err(error),
    175                 };
    176             record_inbound_result(network, known_peer, remote_addr, result).await;
    177             network.forward_outbox().await;
    178         }
    179         other => {
    180             let result = network.inner.node.lock().await.receive(other);
    181             record_inbound_result(network, known_peer, remote_addr, result).await;
    182             network.forward_outbox().await;
    183         }
    184     }
    185     Ok(())
    186 }
    187 
    188 async fn process_blinded_transactions(
    189     network: &GossipNetwork,
    190     remote_addr: SocketAddr,
    191     known_peer: &Option<String>,
    192     transactions: Vec<BlindedTransaction>,
    193 ) {
    194     let first_error = {
    195         let mut node = network.inner.node.lock().await;
    196         let mut first_error = None;
    197         for tx in transactions {
    198             if let Err(error) = node.receive_blinded_transaction(tx) {
    199                 first_error.get_or_insert(error);
    200             }
    201         }
    202         first_error
    203     };
    204     record_inbound_result(
    205         network,
    206         known_peer,
    207         remote_addr,
    208         first_error
    209             .map(|error| Err(anyhow!(format!("{error:#}"))))
    210             .unwrap_or(Ok(())),
    211     )
    212     .await;
    213     network.forward_outbox().await;
    214 }
    215 
    216 async fn process_mine_actions(
    217     network: &GossipNetwork,
    218     remote_addr: SocketAddr,
    219     known_peer: &Option<String>,
    220     transactions: Vec<Transaction>,
    221 ) {
    222     let first_error = {
    223         let mut node = network.inner.node.lock().await;
    224         let mut first_error = None;
    225         for tx in transactions {
    226             if let Err(error) = node.receive_mine_action(tx) {
    227                 first_error.get_or_insert(error);
    228             }
    229         }
    230         first_error
    231     };
    232     record_inbound_result(
    233         network,
    234         known_peer,
    235         remote_addr,
    236         first_error
    237             .map(|error| Err(anyhow!(format!("{error:#}"))))
    238             .unwrap_or(Ok(())),
    239     )
    240     .await;
    241     network.forward_outbox().await;
    242 }
    243 
    244 async fn process_blinded_reveals(
    245     network: &GossipNetwork,
    246     remote_addr: SocketAddr,
    247     known_peer: &Option<String>,
    248     reveals: Vec<BlindedReveal>,
    249 ) {
    250     let first_error = {
    251         let mut node = network.inner.node.lock().await;
    252         let mut first_error = None;
    253         for reveal in reveals {
    254             if let Err(error) = node.receive_blinded_reveal(reveal) {
    255                 first_error.get_or_insert(error);
    256             }
    257         }
    258         first_error
    259     };
    260     record_inbound_result(
    261         network,
    262         known_peer,
    263         remote_addr,
    264         first_error
    265             .map(|error| Err(anyhow!(format!("{error:#}"))))
    266             .unwrap_or(Ok(())),
    267     )
    268     .await;
    269     network.forward_outbox().await;
    270 }
    271 
    272 async fn process_reveal_bundles(
    273     network: &GossipNetwork,
    274     remote_addr: SocketAddr,
    275     known_peer: &Option<String>,
    276     bundles: Vec<RevealBundle>,
    277 ) {
    278     let first_error = {
    279         let mut node = network.inner.node.lock().await;
    280         let mut first_error = None;
    281         for bundle in bundles {
    282             if let Err(error) = node.receive_reveal_bundle(bundle) {
    283                 first_error.get_or_insert(error);
    284             }
    285         }
    286         first_error
    287     };
    288     record_inbound_result(
    289         network,
    290         known_peer,
    291         remote_addr,
    292         first_error
    293             .map(|error| Err(anyhow!(format!("{error:#}"))))
    294             .unwrap_or(Ok(())),
    295     )
    296     .await;
    297     network.forward_outbox().await;
    298 }
    299 
    300 async fn record_inbound_result(
    301     network: &GossipNetwork,
    302     known_peer: &Option<String>,
    303     remote_addr: SocketAddr,
    304     result: Result<()>,
    305 ) {
    306     let peer = known_peer
    307         .clone()
    308         .unwrap_or_else(|| remote_addr.to_string());
    309     match result {
    310         Ok(()) => {
    311             if known_peer.is_some() {
    312                 network.inner.peers.lock().await.record_received(&peer, 1);
    313             }
    314         }
    315         Err(error) => {
    316             let message = format!("{error:#}");
    317             if known_peer.is_some() {
    318                 let mut peers = network.inner.peers.lock().await;
    319                 if super::inbound_error_counts_as_misbehavior(&message) {
    320                     peers.record_misbehavior(&peer, message.clone());
    321                 } else {
    322                     peers.record_inbound_error(&peer, message.clone());
    323                 }
    324             }
    325             if debug_logging_enabled() {
    326                 eprintln!("p2p envelope from {peer} ignored: {message}");
    327             }
    328         }
    329     }
    330 }