iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

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;