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 }