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 }