iuna

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

sync.rs (16527B)


      1 use std::net::SocketAddr;
      2 
      3 use anyhow::Result;
      4 use tokio::net::tcp::OwnedWriteHalf;
      5 
      6 use super::metrics::P2pMetricsCounters;
      7 use super::peer_addr::{
      8     is_self_peer_address_for, normalize_advertised_peer, peer_list_address_is_discoverable,
      9     peer_needs_snapshot,
     10 };
     11 use super::{GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, write_envelope, write_payload};
     12 use crate::app::{GossipEnvelope, SharedNode, debug_logging_enabled};
     13 
     14 pub(super) async fn maybe_request_catchup(
     15     network: &GossipNetwork,
     16     writer: &mut OwnedWriteHalf,
     17     peer_status: &PeerStatus,
     18 ) -> Result<()> {
     19     let (local_height, local_tip_hash) = {
     20         let node = network.inner.node.lock().await;
     21         let status = node.ledger().status();
     22         (status.height, status.tip_hash)
     23     };
     24     if peer_status.request_snapshot {
     25         write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?;
     26     } else if peer_status.height > local_height {
     27         write_envelope(
     28             writer,
     29             &GossipEnvelope::BlockRangeRequest {
     30                 from_height: local_height + 1,
     31                 limit: MAX_BLOCK_BATCH,
     32             },
     33         )
     34         .await?;
     35     } else if peer_status.height == local_height && peer_status.tip_hash != local_tip_hash {
     36         write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?;
     37     }
     38     Ok(())
     39 }
     40 
     41 pub(super) async fn push_catchup_to_peer(
     42     network: &GossipNetwork,
     43     writer: &mut OwnedWriteHalf,
     44     peer_status: &PeerStatus,
     45 ) -> Result<Option<PeerStatus>> {
     46     let payload = catchup_payload_for_peer(&network.inner.node, peer_status).await;
     47     if payload.is_empty() {
     48         return Ok(None);
     49     }
     50 
     51     let updated_status = payload.iter().find_map(|envelope| match envelope {
     52         GossipEnvelope::Blocks { blocks } => blocks
     53             .last()
     54             .map(|block| PeerStatus::new(block.height, block.hash.clone())),
     55         GossipEnvelope::ChainSnapshot(snapshot) => snapshot
     56             .blocks
     57             .last()
     58             .map(|block| PeerStatus::new(block.height, block.hash.clone())),
     59         _ => None,
     60     });
     61     write_payload(writer, &payload).await?;
     62     Ok(updated_status)
     63 }
     64 
     65 pub(super) async fn catchup_payload_for_peer(
     66     node: &SharedNode,
     67     peer_status: &PeerStatus,
     68 ) -> Vec<GossipEnvelope> {
     69     let mut node = node.lock().await;
     70     let local_status = node.ledger().status();
     71     if node.ledger().is_setup_placeholder() {
     72         return Vec::new();
     73     }
     74     let mempool = node.mempool_gossip();
     75     if peer_status.push_snapshot {
     76         let mut payload = vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())];
     77         payload.extend(mempool);
     78         return payload;
     79     }
     80     if peer_status.height < local_status.height {
     81         let blocks = node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH);
     82         if blocks.is_empty() {
     83             mempool
     84         } else {
     85             let mut payload = vec![GossipEnvelope::Blocks { blocks }];
     86             payload.extend(mempool);
     87             payload
     88         }
     89     } else if peer_status.height == local_status.height
     90         && peer_status.tip_hash != local_status.tip_hash
     91     {
     92         let mut payload = vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())];
     93         payload.extend(mempool);
     94         payload
     95     } else {
     96         mempool
     97     }
     98 }
     99 
    100 pub(super) async fn apply_peer_list(
    101     network: &GossipNetwork,
    102     remote_addr: SocketAddr,
    103     peers: Vec<String>,
    104 ) -> Result<()> {
    105     let self_filter_addr = network.self_filter_addr().await;
    106     let mut peerbook = network.inner.peers.lock().await;
    107     for address in peers {
    108         let peer = match normalize_advertised_peer(&address, remote_addr) {
    109             Ok(peer) => peer,
    110             Err(error) => {
    111                 if debug_logging_enabled() {
    112                     eprintln!("p2p peer-list address {address} ignored: {error:#}");
    113                 }
    114                 continue;
    115             }
    116         };
    117         if is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr) {
    118             P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips);
    119         } else if peer_list_address_is_discoverable(&peer, remote_addr)? {
    120             peerbook.add_peer(peer);
    121         } else {
    122             P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips);
    123         }
    124     }
    125     Ok(())
    126 }
    127 
    128 pub(super) async fn write_peer_exchange(
    129     network: &GossipNetwork,
    130     writer: &mut OwnedWriteHalf,
    131     known_peer: &Option<String>,
    132 ) -> Result<()> {
    133     let envelope = network.peer_exchange().await;
    134     let GossipEnvelope::PeerList { peers } = &envelope else {
    135         return Ok(());
    136     };
    137     if peers.is_empty() {
    138         return Ok(());
    139     }
    140     write_envelope(writer, &envelope).await?;
    141     if let Some(peer) = known_peer {
    142         network.inner.peers.lock().await.record_sent(peer, 1);
    143     }
    144     Ok(())
    145 }
    146 
    147 pub(super) async fn envelopes_for_peer(
    148     node: Option<&SharedNode>,
    149     peer_status: Option<PeerStatus>,
    150     envelopes: &[GossipEnvelope],
    151 ) -> Vec<GossipEnvelope> {
    152     let Some(node) = node else {
    153         return envelopes.to_vec();
    154     };
    155     let Some(peer_status) = peer_status else {
    156         return envelopes.to_vec();
    157     };
    158 
    159     let node = node.lock().await;
    160     let local_status = node.ledger().status();
    161     if node.ledger().is_setup_placeholder() {
    162         return envelopes
    163             .iter()
    164             .filter(|envelope| !matches!(envelope, GossipEnvelope::Block(_)))
    165             .cloned()
    166             .collect();
    167     }
    168     if peer_status.height < local_status.height {
    169         let mut payload = vec![GossipEnvelope::Blocks {
    170             blocks: node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH),
    171         }];
    172         payload.extend(
    173             envelopes
    174                 .iter()
    175                 .filter(|envelope| !matches!(envelope, GossipEnvelope::Block(_)))
    176                 .cloned(),
    177         );
    178         return payload;
    179     }
    180 
    181     if peer_status.height == local_status.height && peer_status.tip_hash != local_status.tip_hash {
    182         return vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())];
    183     }
    184 
    185     if peer_needs_snapshot(peer_status.height, envelopes) {
    186         return vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())];
    187     }
    188 
    189     envelopes
    190         .iter()
    191         .filter(|envelope| match envelope {
    192             GossipEnvelope::Block(block) => block.height > peer_status.height,
    193             GossipEnvelope::Inventory { blocks, .. } => {
    194                 blocks.iter().any(|block| block.height > peer_status.height)
    195             }
    196             _ => true,
    197         })
    198         .map(|envelope| match envelope {
    199             GossipEnvelope::Inventory { blocks } => GossipEnvelope::Inventory {
    200                 blocks: blocks
    201                     .iter()
    202                     .filter(|block| block.height > peer_status.height)
    203                     .cloned()
    204                     .collect(),
    205             },
    206             other => other.clone(),
    207         })
    208         .filter(|envelope| match envelope {
    209             GossipEnvelope::Inventory { blocks } => !blocks.is_empty(),
    210             _ => true,
    211         })
    212         .collect()
    213 }
    214 
    215 #[cfg(test)]
    216 mod tests {
    217     use std::sync::Arc;
    218 
    219     use crate::{
    220         app::{GossipEnvelope, PeerBook},
    221         domain::Wallet,
    222     };
    223 
    224     use super::super::{
    225         PeerStatus,
    226         test_support::{allocations, gossip_network, node, queue_plaintext_burn},
    227     };
    228     use super::{apply_peer_list, catchup_payload_for_peer, envelopes_for_peer};
    229 
    230     #[tokio::test]
    231     async fn peer_payload_repairs_lagging_peer_without_networking() {
    232         let alice = Wallet::from_seed("p2p-alice");
    233         let bob = Wallet::from_seed("p2p-bob");
    234         let allocations = allocations(&[alice.clone(), bob.clone()], 1_000);
    235         let node = Arc::new(tokio::sync::Mutex::new(node(
    236             "alice",
    237             alice.clone(),
    238             allocations,
    239         )));
    240         let block = {
    241             let mut node = node.lock().await;
    242             queue_plaintext_burn(&mut node, &alice, 1);
    243             node.drain_outbox();
    244             let block = node.mine_one_at(1).unwrap();
    245             node.drain_outbox();
    246             block
    247         };
    248 
    249         let payload = envelopes_for_peer(
    250             Some(&node),
    251             Some(PeerStatus::new(0, "genesis".to_string())),
    252             &[GossipEnvelope::Block(block)],
    253         )
    254         .await;
    255 
    256         assert!(matches!(payload[0], GossipEnvelope::Blocks { .. }));
    257         match &payload[0] {
    258             GossipEnvelope::Blocks { blocks } => {
    259                 assert_eq!(blocks.len(), 1);
    260                 assert_eq!(blocks[0].height, 1);
    261             }
    262             _ => unreachable!(),
    263         }
    264     }
    265 
    266     #[tokio::test]
    267     async fn session_catchup_payload_pushes_missing_blocks_to_lagging_peer() {
    268         let alice = Wallet::from_seed("catchup-alice");
    269         let bob = Wallet::from_seed("catchup-bob");
    270         let allocations = allocations(&[alice.clone(), bob], 1_000);
    271         let node = Arc::new(tokio::sync::Mutex::new(node(
    272             "alice",
    273             alice.clone(),
    274             allocations,
    275         )));
    276         {
    277             let mut node = node.lock().await;
    278             for height in 1..=3 {
    279                 queue_plaintext_burn(&mut node, &alice, 1);
    280                 node.drain_outbox();
    281                 node.mine_one_at(height).unwrap();
    282                 node.drain_outbox();
    283             }
    284         }
    285 
    286         let payload =
    287             catchup_payload_for_peer(&node, &PeerStatus::new(1, "old-tip".to_string())).await;
    288 
    289         assert_eq!(payload.len(), 1);
    290         match &payload[0] {
    291             GossipEnvelope::Blocks { blocks } => {
    292                 assert_eq!(
    293                     blocks.iter().map(|block| block.height).collect::<Vec<_>>(),
    294                     vec![2, 3]
    295                 );
    296             }
    297             other => panic!("expected missing block payload, got {other:?}"),
    298         }
    299     }
    300 
    301     #[tokio::test]
    302     async fn session_catchup_payload_pushes_blinded_mempool_to_synced_peer() {
    303         let alice = Wallet::from_seed("catchup-mempool-alice");
    304         let bob = Wallet::from_seed("catchup-mempool-bob");
    305         let allocations = allocations(&[alice.clone(), bob], 1_000);
    306         let node = Arc::new(tokio::sync::Mutex::new(node(
    307             "alice",
    308             alice.clone(),
    309             allocations,
    310         )));
    311         let expected_commitment = {
    312             let mut node = node.lock().await;
    313             let tx = node.ledger().build_burn(&alice, 1, 0).unwrap();
    314             let built = node
    315                 .ledger()
    316                 .build_blinded_transaction(&alice, tx, 20)
    317                 .unwrap();
    318             let commitment = built.transaction.commitment.clone();
    319             node.receive_blinded_transaction(built.transaction).unwrap();
    320             node.drain_outbox();
    321             commitment
    322         };
    323         let peer_status = {
    324             let node = node.lock().await;
    325             let status = node.ledger().status();
    326             PeerStatus::new(status.height, status.tip_hash)
    327         };
    328 
    329         let payload = catchup_payload_for_peer(&node, &peer_status).await;
    330 
    331         assert_eq!(payload.len(), 1);
    332         match &payload[0] {
    333             GossipEnvelope::BlindedTransactions { transactions } => {
    334                 assert_eq!(transactions.len(), 1);
    335                 assert_eq!(transactions[0].commitment, expected_commitment);
    336             }
    337             other => panic!("expected blinded mempool payload, got {other:?}"),
    338         }
    339     }
    340 
    341     #[tokio::test]
    342     async fn peer_list_adds_stable_outbound_peers() {
    343         let alice = Wallet::from_seed("px-recv-alice");
    344         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    345         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    346         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    347         let network = gossip_network(
    348             node,
    349             Arc::clone(&peers),
    350             "127.0.0.1:9544".parse().unwrap(),
    351             None,
    352         );
    353         apply_peer_list(
    354             &network,
    355             "127.0.0.1:9545".parse().unwrap(),
    356             vec!["127.0.0.1:9544".to_string(), "127.0.0.1:9546".to_string()],
    357         )
    358         .await
    359         .unwrap();
    360 
    361         let addresses = peers.lock().await.addresses();
    362         assert!(!addresses.contains(&"127.0.0.1:9544".to_string()));
    363         assert!(addresses.contains(&"127.0.0.1:9546".to_string()));
    364     }
    365 
    366     #[tokio::test]
    367     async fn peer_list_ignores_invalid_peer_addresses() {
    368         let alice = Wallet::from_seed("px-list-invalid-alice");
    369         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    370         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    371         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    372         let network = gossip_network(
    373             node,
    374             Arc::clone(&peers),
    375             "127.0.0.1:9544".parse().unwrap(),
    376             None,
    377         );
    378         apply_peer_list(
    379             &network,
    380             "127.0.0.1:9545".parse().unwrap(),
    381             vec![
    382                 "iuna.jhx.app:9444".to_string(),
    383                 "127.0.0.1:9546".to_string(),
    384             ],
    385         )
    386         .await
    387         .unwrap();
    388 
    389         let addresses = peers.lock().await.addresses();
    390         assert!(!addresses.contains(&"iuna.jhx.app:9444".to_string()));
    391         assert!(addresses.contains(&"127.0.0.1:9546".to_string()));
    392     }
    393 
    394     #[tokio::test]
    395     async fn peer_list_ignores_announced_self_address() {
    396         let alice = Wallet::from_seed("px-list-announced-self-alice");
    397         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    398         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    399         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    400         let network = gossip_network(
    401             node,
    402             Arc::clone(&peers),
    403             "0.0.0.0:9444".parse().unwrap(),
    404             Some("8.8.8.8:9444".parse().unwrap()),
    405         );
    406 
    407         apply_peer_list(
    408             &network,
    409             "8.8.4.4:9444".parse().unwrap(),
    410             vec!["8.8.8.8:9444".to_string(), "8.8.4.4:9445".to_string()],
    411         )
    412         .await
    413         .unwrap();
    414 
    415         let addresses = peers.lock().await.addresses();
    416         assert!(!addresses.contains(&"8.8.8.8:9444".to_string()));
    417         assert!(addresses.contains(&"8.8.4.4:9445".to_string()));
    418         network.set_accept_inbound(false).await.unwrap();
    419     }
    420 
    421     #[tokio::test]
    422     async fn peer_list_ignores_private_ephemeral_addresses() {
    423         let alice = Wallet::from_seed("px-private-ephemeral-alice");
    424         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    425         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    426         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    427         let network = gossip_network(
    428             node,
    429             Arc::clone(&peers),
    430             "0.0.0.0:9444".parse().unwrap(),
    431             None,
    432         );
    433 
    434         apply_peer_list(
    435             &network,
    436             "142.132.164.59:9444".parse().unwrap(),
    437             vec![
    438                 "10.42.1.1:10091".to_string(),
    439                 "142.132.164.59:9444".to_string(),
    440             ],
    441         )
    442         .await
    443         .unwrap();
    444 
    445         let addresses = peers.lock().await.addresses();
    446         assert!(!addresses.contains(&"10.42.1.1:10091".to_string()));
    447         assert!(addresses.contains(&"142.132.164.59:9444".to_string()));
    448     }
    449 
    450     #[tokio::test]
    451     async fn peer_list_ignores_loopback_alias_for_unspecified_self() {
    452         let alice = Wallet::from_seed("px-self-alias-alice");
    453         let allocations = allocations(std::slice::from_ref(&alice), 1_000);
    454         let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
    455         let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
    456         let network = gossip_network(
    457             node,
    458             Arc::clone(&peers),
    459             "0.0.0.0:9545".parse().unwrap(),
    460             None,
    461         );
    462 
    463         apply_peer_list(
    464             &network,
    465             "127.0.0.1:9544".parse().unwrap(),
    466             vec!["127.0.0.1:9545".to_string(), "127.0.0.1:9546".to_string()],
    467         )
    468         .await
    469         .unwrap();
    470 
    471         let addresses = peers.lock().await.addresses();
    472         assert!(!addresses.contains(&"127.0.0.1:9545".to_string()));
    473         assert!(addresses.contains(&"127.0.0.1:9546".to_string()));
    474         assert_eq!(network.metrics().self_peer_skips, 1);
    475     }
    476 }