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 }