sync.rs (5658B)
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 }; 10 use super::{GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, write_envelope}; 11 use crate::app::{ 12 CAPABILITY_TRANSACTION_V2_MEMPOOL, GossipEnvelope, SharedNode, debug_logging_enabled, 13 }; 14 15 pub(super) async fn maybe_request_catchup( 16 network: &GossipNetwork, 17 writer: &mut OwnedWriteHalf, 18 peer_status: &PeerStatus, 19 ) -> Result<bool> { 20 let (local_height, local_tip_hash) = { 21 let node = network.inner.node.lock().await; 22 let status = node.ledger().status(); 23 (status.height, status.tip_hash) 24 }; 25 if peer_status.request_bootstrap { 26 write_envelope(writer, &GossipEnvelope::ChainBootstrapRequest).await?; 27 return Ok(true); 28 } else if peer_status.height > local_height { 29 write_envelope( 30 writer, 31 &GossipEnvelope::BlockRangeRequest { 32 from_height: local_height + 1, 33 limit: MAX_BLOCK_BATCH, 34 }, 35 ) 36 .await?; 37 return Ok(true); 38 } else if peer_status.height == local_height && peer_status.tip_hash != local_tip_hash { 39 let locator = network.inner.node.lock().await.block_locator(); 40 write_envelope( 41 writer, 42 &GossipEnvelope::BlockLocatorRequest { 43 locator, 44 limit: MAX_BLOCK_BATCH, 45 }, 46 ) 47 .await?; 48 return Ok(true); 49 } 50 Ok(false) 51 } 52 53 pub(super) async fn apply_peer_list( 54 network: &GossipNetwork, 55 remote_addr: SocketAddr, 56 peers: Vec<String>, 57 ) -> Result<()> { 58 let self_filter_addr = network.self_filter_addr().await; 59 let mut peerbook = network.inner.peers.lock().await; 60 for address in peers { 61 let peer = match normalize_advertised_peer(&address, remote_addr) { 62 Ok(peer) => peer, 63 Err(error) => { 64 if debug_logging_enabled() { 65 eprintln!("p2p peer-list address {address} ignored: {error:#}"); 66 } 67 continue; 68 } 69 }; 70 if is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr) { 71 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips); 72 } else if peer_list_address_is_discoverable(&peer, remote_addr)? { 73 peerbook.add_discovered_peer(peer); 74 } else { 75 P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips); 76 } 77 } 78 Ok(()) 79 } 80 81 pub(super) async fn write_peer_exchange( 82 network: &GossipNetwork, 83 writer: &mut OwnedWriteHalf, 84 known_peer: &Option<String>, 85 ) -> Result<()> { 86 let envelope = network.peer_exchange().await; 87 let GossipEnvelope::PeerList { peers } = &envelope else { 88 return Ok(()); 89 }; 90 if peers.is_empty() { 91 return Ok(()); 92 } 93 write_envelope(writer, &envelope).await?; 94 if let Some(peer) = known_peer { 95 network.inner.peers.lock().await.record_sent(peer, 1); 96 } 97 Ok(()) 98 } 99 100 pub(super) async fn envelopes_for_peer( 101 node: Option<&SharedNode>, 102 peer_status: Option<PeerStatus>, 103 envelopes: &[GossipEnvelope], 104 ) -> Vec<GossipEnvelope> { 105 let Some(node) = node else { 106 return envelopes.to_vec(); 107 }; 108 let supports_v2_mempool = peer_status 109 .as_ref() 110 .is_some_and(|status| status.supports(CAPABILITY_TRANSACTION_V2_MEMPOOL)); 111 let envelopes = envelopes 112 .iter() 113 .filter(|envelope| { 114 supports_v2_mempool 115 || !matches!( 116 envelope, 117 GossipEnvelope::TransactionV2 { .. } | GossipEnvelope::TransactionsV2 { .. } 118 ) 119 }) 120 .cloned() 121 .collect::<Vec<_>>(); 122 let Some(peer_status) = peer_status else { 123 return envelopes; 124 }; 125 126 let node = node.lock().await; 127 let local_status = node.ledger().status(); 128 if node.ledger().is_setup_placeholder() { 129 return envelopes 130 .iter() 131 .filter(|envelope| !matches!(envelope, GossipEnvelope::Block(_))) 132 .cloned() 133 .collect(); 134 } 135 if peer_status.height < local_status.height { 136 return envelopes 137 .iter() 138 .filter(|envelope| { 139 matches!( 140 envelope, 141 GossipEnvelope::PeerStatus { .. } | GossipEnvelope::Inventory { .. } 142 ) 143 }) 144 .cloned() 145 .collect(); 146 } 147 148 if peer_status.height == local_status.height && peer_status.tip_hash != local_status.tip_hash { 149 return Vec::new(); 150 } 151 152 envelopes 153 .iter() 154 .filter(|envelope| match envelope { 155 GossipEnvelope::Block(block) => block.height > peer_status.height, 156 GossipEnvelope::Inventory { blocks, .. } => { 157 blocks.iter().any(|block| block.height > peer_status.height) 158 } 159 _ => true, 160 }) 161 .map(|envelope| match envelope { 162 GossipEnvelope::Inventory { blocks } => GossipEnvelope::Inventory { 163 blocks: blocks 164 .iter() 165 .filter(|block| block.height > peer_status.height) 166 .cloned() 167 .collect(), 168 }, 169 other => other.clone(), 170 }) 171 .filter(|envelope| match envelope { 172 GossipEnvelope::Inventory { blocks } => !blocks.is_empty(), 173 _ => true, 174 }) 175 .collect() 176 }