p2p.rs (6215B)
1 use std::{ 2 collections::BTreeMap, 3 net::SocketAddr, 4 sync::{Arc, Mutex as StdMutex}, 5 time::Duration, 6 }; 7 8 use tokio::{ 9 sync::{Mutex, OwnedSemaphorePermit, Semaphore, mpsc, watch}, 10 task::JoinHandle, 11 }; 12 13 use crate::app::{GossipEnvelope, SharedNode, SharedPeerBook}; 14 15 mod error; 16 mod fetch; 17 mod handshake; 18 mod identity; 19 mod inbound_limiter; 20 mod line_codec; 21 mod metrics; 22 mod network; 23 mod peer_addr; 24 mod peer_status; 25 mod process; 26 mod session; 27 mod sync; 28 #[cfg(test)] 29 mod test_support; 30 mod writer; 31 use error::SyncError; 32 pub use fetch::{ 33 fetch_peer_height, fetch_snapshot, fetch_snapshot_with_announcement, validate_chain_snapshot, 34 }; 35 use fetch::{ 36 network_adjusted_time_ms, validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf, 37 }; 38 #[cfg(test)] 39 use handshake::verify_advertised_peer_node_id; 40 use handshake::{ 41 forget_stale_self_peer, process_hello, process_hello_with_verification, record_peer_status, 42 }; 43 #[cfg(test)] 44 use identity::peer_verification_response_for_node_id; 45 use identity::{new_node_id, peer_verification_response}; 46 #[cfg(test)] 47 use identity::{new_verification_nonce, peer_verification_response_is_valid}; 48 use inbound_limiter::{InboundConnectionLimiter, InboundSessionPermit, InboundSessionRejection}; 49 use line_codec::{LimitedLineReader, parse_envelope, read_session_envelope}; 50 pub use metrics::P2pMetrics; 51 use metrics::P2pMetricsCounters; 52 use peer_addr::{ 53 inbound_error_counts_as_misbehavior, is_possible_fork_error, is_quiet_disconnect, 54 is_self_peer_address_for, next_reconnect_delay as next_reconnect_delay_with_max, 55 normalize_advertised_peer, 56 }; 57 use peer_status::PeerStatus; 58 use process::{process_envelope, respond_to_peer_verification_challenge}; 59 use session::{accept_loop, outbound_session, outbound_supervisor}; 60 use sync::{apply_peer_list, envelopes_for_peer, maybe_request_catchup, write_peer_exchange}; 61 use writer::{byte_bounded_block_page, write_envelope, write_payload}; 62 63 const MAX_BLOCK_BATCH: usize = 128; 64 const MAX_OBJECT_REQUESTS: usize = 128; 65 const MAX_INVENTORY_ITEMS: usize = 512; 66 const MAX_BLOCK_LOCATOR_HASHES: usize = 64; 67 const MAX_PEER_LIST: usize = 128; 68 const MAX_GOSSIP_LINE_BYTES: usize = 8 * 1024 * 1024; 69 const MAX_INBOUND_SESSIONS: usize = 64; 70 const MAX_INBOUND_SESSIONS_PER_IP: usize = 8; 71 const MAX_INBOUND_ACCEPTS_PER_IP_PER_WINDOW: usize = 24; 72 const INBOUND_ACCEPT_RATE_WINDOW_MS: u64 = 10_000; 73 const PEER_QUEUE_SIZE: usize = 256; 74 const INBOUND_PEER_QUEUE_SIZE: usize = 16; 75 const MAX_OUTBOUND_BATCH_BYTES: usize = MAX_GOSSIP_LINE_BYTES + 1; 76 const PEER_QUEUE_BYTES: usize = 4 * MAX_OUTBOUND_BATCH_BYTES; 77 const INBOUND_PEER_QUEUE_BYTES: usize = 2 * MAX_OUTBOUND_BATCH_BYTES; 78 const MAX_CONCURRENT_CHAIN_VALIDATIONS: usize = 2; 79 const STALE_INBOUND_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000; 80 const STALE_DISCOVERED_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000; 81 const MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE: usize = 32; 82 const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); 83 const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5); 84 const WRITE_TIMEOUT: Duration = Duration::from_secs(10); 85 const SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2); 86 const CATCHUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(10); 87 const PEER_EXCHANGE_INTERVAL: Duration = Duration::from_secs(30); 88 const JOIN_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5); 89 const MAX_JOIN_RESPONSE_ENVELOPES: usize = 16; 90 const MAX_PEER_VERIFICATION_ENVELOPES: usize = 8; 91 const INITIAL_RECONNECT_DELAY: Duration = Duration::from_secs(1); 92 const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30); 93 const INBOUND_SESSION_PREFIX: &str = "inbound://"; 94 struct OutboundBatch { 95 envelopes: Arc<[GossipEnvelope]>, 96 _queued_bytes: OwnedSemaphorePermit, 97 } 98 99 #[derive(Clone)] 100 struct GossipSession { 101 peer: String, 102 sender: mpsc::Sender<OutboundBatch>, 103 shutdown: watch::Sender<bool>, 104 queue_bytes: Arc<Semaphore>, 105 } 106 107 struct ChainValidationCoordinator { 108 permits: Arc<Semaphore>, 109 active: StdMutex<BTreeMap<String, watch::Sender<bool>>>, 110 } 111 112 impl Default for ChainValidationCoordinator { 113 fn default() -> Self { 114 Self { 115 permits: Arc::new(Semaphore::new(MAX_CONCURRENT_CHAIN_VALIDATIONS)), 116 active: StdMutex::new(BTreeMap::new()), 117 } 118 } 119 } 120 121 struct ChainValidationGuard { 122 coordinator: Arc<ChainValidationCoordinator>, 123 key: String, 124 _permit: OwnedSemaphorePermit, 125 } 126 127 impl Drop for ChainValidationGuard { 128 fn drop(&mut self) { 129 let sender = self 130 .coordinator 131 .active 132 .lock() 133 .expect("chain validation mutex poisoned") 134 .remove(&self.key); 135 if let Some(sender) = sender { 136 let _ = sender.send(true); 137 } 138 } 139 } 140 141 #[cfg(feature = "fuzzing")] 142 pub fn fuzz_parse_envelope(line: &str) -> anyhow::Result<GossipEnvelope> { 143 parse_envelope(line) 144 } 145 146 #[derive(Clone)] 147 pub struct GossipNetwork { 148 inner: Arc<GossipNetworkInner>, 149 } 150 151 pub(super) struct SyncProgressGuard { 152 network: GossipNetwork, 153 id: u64, 154 generation: u64, 155 } 156 157 impl SyncProgressGuard { 158 pub(super) fn id(&self) -> u64 { 159 self.id 160 } 161 162 pub(super) fn is_current(&self) -> bool { 163 self.network.sync_generation_is_current(self.generation) 164 } 165 } 166 167 impl Drop for SyncProgressGuard { 168 fn drop(&mut self) { 169 self.network.finish_sync_progress(self.id); 170 } 171 } 172 173 struct GossipNetworkInner { 174 node: SharedNode, 175 peers: SharedPeerBook, 176 listen_addr: SocketAddr, 177 p2p_announce_addr: Mutex<Option<SocketAddr>>, 178 node_id: String, 179 accept_task: Mutex<Option<JoinHandle<()>>>, 180 sessions: Mutex<BTreeMap<String, GossipSession>>, 181 inbound_limiter: Arc<StdMutex<InboundConnectionLimiter>>, 182 metrics: P2pMetricsCounters, 183 sync_progress: StdMutex<SyncProgressState>, 184 chain_validation: Arc<ChainValidationCoordinator>, 185 } 186 187 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 188 pub struct SyncProgress { 189 pub start_height: u64, 190 pub validated_height: u64, 191 pub target_height: u64, 192 } 193 194 #[derive(Default)] 195 struct SyncProgressState { 196 next_id: u64, 197 generation: u64, 198 active: BTreeMap<u64, SyncProgress>, 199 last_activity: Option<std::time::Instant>, 200 } 201 202 #[cfg(test)] 203 mod tests;