iuna

iuna

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

commit 23fe372f63a604c4b78c1cca46789f3c5140f50a
parent 497172772da5b6e425070fb6ebda60389494b82f
Author: Joris Hartog <jorishartog@hotmail.com>
Date:   Sat, 29 Aug 2026 00:30:59 +0200

Improve P2P chain synchronization

Diffstat:
Msrc/adapters/p2p.rs | 68+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Msrc/adapters/p2p/handshake.rs | 29+++++++++--------------------
Msrc/adapters/p2p/network.rs | 264++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Msrc/adapters/p2p/peer_addr.rs | 16----------------
Msrc/adapters/p2p/peer_status.rs | 14++++++++------
Msrc/adapters/p2p/process.rs | 151++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Msrc/adapters/p2p/session.rs | 377+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Msrc/adapters/p2p/sync.rs | 85+++++++++++++++----------------------------------------------------------------
Msrc/adapters/p2p/test_support.rs | 1+
Msrc/adapters/p2p/tests.rs | 576++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Msrc/adapters/p2p/writer.rs | 25++++++++++++++++++++-----
Msrc/app/gossip.rs | 65++++++++++++++++++++++++++++++++++++++++++++++-------------------
Msrc/domain/ledger_queries.rs | 15++++++++++++++-
Msrc/main.rs | 95++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Msrc/main_tests.rs | 19++++++++++++++++++-
15 files changed, 1515 insertions(+), 285 deletions(-)

diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs @@ -6,7 +6,7 @@ use std::{ }; use tokio::{ - sync::{Mutex, mpsc}, + sync::{Mutex, OwnedSemaphorePermit, Semaphore, mpsc, watch}, task::JoinHandle, }; @@ -53,12 +53,7 @@ use peer_addr::{ use peer_status::PeerStatus; use process::{process_envelope, respond_to_peer_verification_challenge}; use session::{accept_loop, outbound_session, outbound_supervisor}; -#[cfg(test)] -use sync::catchup_payload_for_peer; -use sync::{ - apply_peer_list, envelopes_for_peer, maybe_request_catchup, push_catchup_to_peer, - write_peer_exchange, -}; +use sync::{apply_peer_list, envelopes_for_peer, maybe_request_catchup, write_peer_exchange}; use writer::{byte_bounded_block_page, write_envelope, write_payload}; const MAX_BLOCK_BATCH: usize = 128; @@ -72,19 +67,72 @@ const MAX_INBOUND_SESSIONS_PER_IP: usize = 8; const MAX_INBOUND_ACCEPTS_PER_IP_PER_WINDOW: usize = 24; const INBOUND_ACCEPT_RATE_WINDOW_MS: u64 = 10_000; const PEER_QUEUE_SIZE: usize = 256; +const INBOUND_PEER_QUEUE_SIZE: usize = 16; +const MAX_OUTBOUND_BATCH_BYTES: usize = MAX_GOSSIP_LINE_BYTES + 1; +const PEER_QUEUE_BYTES: usize = 4 * MAX_OUTBOUND_BATCH_BYTES; +const INBOUND_PEER_QUEUE_BYTES: usize = 2 * MAX_OUTBOUND_BATCH_BYTES; +const MAX_CONCURRENT_CHAIN_VALIDATIONS: usize = 2; const STALE_INBOUND_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000; const STALE_DISCOVERED_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000; const MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE: usize = 32; const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5); +const WRITE_TIMEOUT: Duration = Duration::from_secs(10); const SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2); +const CATCHUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(10); const PEER_EXCHANGE_INTERVAL: Duration = Duration::from_secs(30); const JOIN_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5); const MAX_JOIN_RESPONSE_ENVELOPES: usize = 16; const MAX_PEER_VERIFICATION_ENVELOPES: usize = 8; const INITIAL_RECONNECT_DELAY: Duration = Duration::from_secs(1); const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30); -type OutboundBatch = Vec<GossipEnvelope>; +const INBOUND_SESSION_PREFIX: &str = "inbound://"; +struct OutboundBatch { + envelopes: Arc<[GossipEnvelope]>, + _queued_bytes: OwnedSemaphorePermit, +} + +#[derive(Clone)] +struct GossipSession { + peer: String, + sender: mpsc::Sender<OutboundBatch>, + shutdown: watch::Sender<bool>, + queue_bytes: Arc<Semaphore>, +} + +struct ChainValidationCoordinator { + permits: Arc<Semaphore>, + active: StdMutex<BTreeMap<String, watch::Sender<bool>>>, +} + +impl Default for ChainValidationCoordinator { + fn default() -> Self { + Self { + permits: Arc::new(Semaphore::new(MAX_CONCURRENT_CHAIN_VALIDATIONS)), + active: StdMutex::new(BTreeMap::new()), + } + } +} + +struct ChainValidationGuard { + coordinator: Arc<ChainValidationCoordinator>, + key: String, + _permit: OwnedSemaphorePermit, +} + +impl Drop for ChainValidationGuard { + fn drop(&mut self) { + let sender = self + .coordinator + .active + .lock() + .expect("chain validation mutex poisoned") + .remove(&self.key); + if let Some(sender) = sender { + let _ = sender.send(true); + } + } +} #[cfg(feature = "fuzzing")] pub fn fuzz_parse_envelope(line: &str) -> anyhow::Result<GossipEnvelope> { @@ -125,10 +173,11 @@ struct GossipNetworkInner { p2p_announce_addr: Mutex<Option<SocketAddr>>, node_id: String, accept_task: Mutex<Option<JoinHandle<()>>>, - sessions: Mutex<BTreeMap<String, mpsc::Sender<OutboundBatch>>>, + sessions: Mutex<BTreeMap<String, GossipSession>>, inbound_limiter: Arc<StdMutex<InboundConnectionLimiter>>, metrics: P2pMetricsCounters, sync_progress: StdMutex<SyncProgressState>, + chain_validation: Arc<ChainValidationCoordinator>, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -143,6 +192,7 @@ struct SyncProgressState { next_id: u64, generation: u64, active: BTreeMap<u64, SyncProgress>, + last_activity: Option<std::time::Instant>, } #[cfg(test)] diff --git a/src/adapters/p2p/handshake.rs b/src/adapters/p2p/handshake.rs @@ -120,7 +120,7 @@ async fn process_hello_inner( { P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections); forget_stale_self_peer(network, known_peer).await; - return Ok(PeerStatus::with_time( + return Ok(PeerStatus::rejected( hello.height, hello.tip_hash, hello.time_ms, @@ -137,7 +137,6 @@ async fn process_hello_inner( let remote_is_setup_placeholder = hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash(); let request_bootstrap = genesis_mismatch && local_accepts_remote_genesis; - let push_bootstrap = genesis_mismatch && remote_is_setup_placeholder; if genesis_mismatch && !local_accepts_remote_genesis && !remote_is_setup_placeholder { anyhow::bail!( "wrong genesis {}; expected {local_genesis}", @@ -146,11 +145,13 @@ async fn process_hello_inner( } let remote_node_id = hello.node_id.clone(); + let mut reject_session = false; if let Some(listen_addr) = &hello.listen_addr { let peer = normalize_advertised_peer(listen_addr, remote_addr)?; if network.is_self_peer(&peer).await { P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections); forget_stale_self_peer(network, known_peer).await; + reject_session = true; } else { let verified = match verification_session.as_mut() { Some(session) => { @@ -180,25 +181,13 @@ async fn process_hello_inner( &PeerStatus::with_time(hello.height, hello.tip_hash.clone(), hello.time_ms), ) .await; - if request_bootstrap { - Ok(PeerStatus::with_bootstrap_request( - hello.height, - hello.tip_hash, - hello.time_ms, - )) - } else if push_bootstrap { - Ok(PeerStatus::with_bootstrap_push( - hello.height, - hello.tip_hash, - hello.time_ms, - )) + let mut status = if request_bootstrap { + PeerStatus::with_bootstrap_request(hello.height, hello.tip_hash, hello.time_ms) } else { - Ok(PeerStatus::with_time( - hello.height, - hello.tip_hash, - hello.time_ms, - )) - } + PeerStatus::with_time(hello.height, hello.tip_hash, hello.time_ms) + }; + status.reject_session = reject_session; + Ok(status) } fn setup_placeholder_genesis_hash() -> String { diff --git a/src/adapters/p2p/network.rs b/src/adapters/p2p/network.rs @@ -5,18 +5,22 @@ use std::{ }; use anyhow::{Context, Result}; -use tokio::{net::TcpListener, sync::mpsc}; +use tokio::{ + net::TcpListener, + sync::{mpsc, watch}, +}; use crate::app::{ BlockInventory, GossipEnvelope, SharedNode, SharedPeerBook, debug_logging_enabled, now_ms, }; use super::{ - GossipNetwork, GossipNetworkInner, InboundConnectionLimiter, InboundSessionPermit, - InboundSessionRejection, MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE, OutboundBatch, P2pMetrics, - P2pMetricsCounters, PEER_QUEUE_SIZE, STALE_DISCOVERED_PEER_RETENTION_MS, - STALE_INBOUND_PEER_RETENTION_MS, accept_loop, is_self_peer_address_for, new_node_id, - outbound_session, outbound_supervisor, + ChainValidationCoordinator, ChainValidationGuard, GossipNetwork, GossipNetworkInner, + GossipSession, INBOUND_SESSION_PREFIX, InboundConnectionLimiter, InboundSessionPermit, + InboundSessionRejection, MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE, MAX_GOSSIP_LINE_BYTES, + MAX_OUTBOUND_BATCH_BYTES, OutboundBatch, P2pMetrics, P2pMetricsCounters, PEER_QUEUE_BYTES, + PEER_QUEUE_SIZE, STALE_DISCOVERED_PEER_RETENTION_MS, STALE_INBOUND_PEER_RETENTION_MS, + accept_loop, is_self_peer_address_for, new_node_id, outbound_session, outbound_supervisor, }; impl GossipNetwork { @@ -35,12 +39,11 @@ impl GossipNetwork { p2p_announce_addr: tokio::sync::Mutex::new(p2p_announce_addr), node_id: new_node_id(), accept_task: tokio::sync::Mutex::new(None), - sessions: tokio::sync::Mutex::new( - BTreeMap::<String, mpsc::Sender<OutboundBatch>>::new(), - ), + sessions: tokio::sync::Mutex::new(BTreeMap::new()), inbound_limiter: Arc::new(StdMutex::new(InboundConnectionLimiter::default())), metrics: P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(ChainValidationCoordinator::default()), }), }; @@ -139,6 +142,7 @@ impl GossipNetwork { target_height: target_height.max(start_height), }, ); + state.last_activity = Some(std::time::Instant::now()); super::SyncProgressGuard { network: self.clone(), id, @@ -177,6 +181,7 @@ impl GossipNetwork { if let Some(progress) = state.active.get_mut(&id) { progress.validated_height = validated_height.clamp(progress.start_height, progress.target_height); + state.last_activity = Some(std::time::Instant::now()); } } @@ -187,6 +192,21 @@ impl GossipNetwork { .lock() .expect("sync progress mutex poisoned"); state.active.remove(&id); + state.last_activity = Some(std::time::Instant::now()); + } + + pub fn chain_sync_active_or_recent(&self, quiet_period: std::time::Duration) -> bool { + let state = self + .inner + .sync_progress + .lock() + .expect("sync progress mutex poisoned"); + if !state.active.is_empty() { + return true; + } + state + .last_activity + .is_some_and(|last_activity| last_activity.elapsed() <= quiet_period) } pub(super) fn try_acquire_inbound_session( @@ -204,24 +224,118 @@ impl GossipNetwork { }) } + pub(super) async fn claim_chain_validation(&self, key: String) -> Option<ChainValidationGuard> { + loop { + let active = self + .inner + .chain_validation + .active + .lock() + .expect("chain validation mutex poisoned") + .get(&key) + .map(watch::Sender::subscribe); + if let Some(mut active) = active { + while !*active.borrow() { + if active.changed().await.is_err() { + break; + } + } + continue; + } + + let permit = Arc::clone(&self.inner.chain_validation.permits) + .acquire_owned() + .await + .ok()?; + let duplicate = { + let mut active = self + .inner + .chain_validation + .active + .lock() + .expect("chain validation mutex poisoned"); + if let Some(sender) = active.get(&key) { + Some(sender.subscribe()) + } else { + let (sender, _) = watch::channel(false); + active.insert(key.clone(), sender); + None + } + }; + if let Some(mut receiver) = duplicate { + drop(permit); + while !*receiver.borrow() { + if receiver.changed().await.is_err() { + break; + } + } + continue; + } + return Some(ChainValidationGuard { + coordinator: Arc::clone(&self.inner.chain_validation), + key, + _permit: permit, + }); + } + } + pub async fn broadcast(&self, envelopes: Vec<GossipEnvelope>) -> Result<()> { let envelopes = self.prepare_gossip(envelopes).await; if envelopes.is_empty() { return Ok(()); } + let batches = byte_bounded_gossip_batches(envelopes)?; let sessions = self.inner.sessions.lock().await.clone(); - for (peer, sender) in sessions { - if self.inner.peers.lock().await.is_banned(&peer) { + let mut disconnect = Vec::new(); + for (session_id, session) in sessions { + if self.inner.peers.lock().await.is_banned(&session.peer) { continue; } - match sender.try_send(envelopes.clone()) { - Ok(()) => {} - Err(mpsc::error::TrySendError::Full(_)) => { - P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); + for (envelopes, encoded_bytes) in &batches { + let permit = + match Arc::clone(&session.queue_bytes).try_acquire_many_owned(*encoded_bytes) { + Ok(permit) => permit, + Err(_) => { + P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); + if session_id.starts_with(INBOUND_SESSION_PREFIX) { + let _ = session.shutdown.send(true); + disconnect.push((session_id.clone(), session.sender.clone())); + } + break; + } + }; + let batch = OutboundBatch { + envelopes: Arc::clone(envelopes), + _queued_bytes: permit, + }; + match session.sender.try_send(batch) { + Ok(()) => {} + Err(mpsc::error::TrySendError::Full(_)) => { + P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full); + if session_id.starts_with(INBOUND_SESSION_PREFIX) { + let _ = session.shutdown.send(true); + disconnect.push((session_id.clone(), session.sender.clone())); + } + break; + } + Err(mpsc::error::TrySendError::Closed(_)) => { + P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed); + let _ = session.shutdown.send(true); + disconnect.push((session_id.clone(), session.sender.clone())); + break; + } } - Err(mpsc::error::TrySendError::Closed(_)) => { - P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed); + } + } + if !disconnect.is_empty() { + let mut sessions = self.inner.sessions.lock().await; + for (session_id, sender) in disconnect { + if sessions + .get(&session_id) + .is_some_and(|session| session.sender.same_channel(&sender)) + { + sessions.remove(&session_id); } } } @@ -307,8 +421,9 @@ impl GossipNetwork { let self_filter_addr = self.self_filter_addr().await; let mut sessions = self.inner.sessions.lock().await; sessions.retain(|peer, _| { - let keep = address_set.contains(peer) - && !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr); + let keep = peer.starts_with(INBOUND_SESSION_PREFIX) + || (address_set.contains(peer) + && !is_self_peer_address_for(peer, self.inner.listen_addr, self_filter_addr)); if !keep { P2pMetricsCounters::inc(&self.inner.metrics.self_peer_skips); } @@ -324,8 +439,23 @@ impl GossipNetwork { } let (sender, receiver) = mpsc::channel(PEER_QUEUE_SIZE); - sessions.insert(peer.clone(), sender); - tokio::spawn(outbound_session(self.clone(), peer, receiver)); + let (shutdown, shutdown_receiver) = watch::channel(false); + let queue_bytes = Arc::new(tokio::sync::Semaphore::new(PEER_QUEUE_BYTES)); + sessions.insert( + peer.clone(), + GossipSession { + peer: peer.clone(), + sender, + shutdown, + queue_bytes, + }, + ); + tokio::spawn(outbound_session( + self.clone(), + peer, + receiver, + shutdown_receiver, + )); } } @@ -339,6 +469,35 @@ impl GossipNetwork { } } +fn byte_bounded_gossip_batches( + envelopes: Vec<GossipEnvelope>, +) -> Result<Vec<(Arc<[GossipEnvelope]>, u32)>> { + let mut batches = Vec::new(); + let mut batch = Vec::new(); + let mut batch_bytes = 0_usize; + for envelope in envelopes { + let encoded_bytes = serde_json::to_vec(&envelope)?.len().saturating_add(1); + if encoded_bytes > MAX_GOSSIP_LINE_BYTES.saturating_add(1) { + anyhow::bail!( + "p2p message is {} bytes, exceeding {} byte limit", + encoded_bytes.saturating_sub(1), + MAX_GOSSIP_LINE_BYTES + ); + } + if !batch.is_empty() && batch_bytes.saturating_add(encoded_bytes) > MAX_OUTBOUND_BATCH_BYTES + { + batches.push((Arc::from(std::mem::take(&mut batch)), batch_bytes as u32)); + batch_bytes = 0; + } + batch_bytes = batch_bytes.saturating_add(encoded_bytes); + batch.push(envelope); + } + if !batch.is_empty() { + batches.push((Arc::from(batch), batch_bytes as u32)); + } + Ok(batches) +} + #[cfg(test)] mod tests { use std::{collections::BTreeMap, sync::Arc}; @@ -350,6 +509,30 @@ mod tests { use super::super::test_support::{allocations, gossip_network, node}; + #[test] + fn outbound_batches_are_bounded_by_encoded_bytes() { + let large_tip = "a".repeat(super::super::MAX_GOSSIP_LINE_BYTES / 2); + let envelopes = vec![ + GossipEnvelope::PeerStatus { + height: 1, + tip_hash: large_tip.clone(), + time_ms: 1, + }, + GossipEnvelope::PeerStatus { + height: 2, + tip_hash: large_tip, + time_ms: 2, + }, + ]; + + let batches = super::byte_bounded_gossip_batches(envelopes).unwrap(); + + assert_eq!(batches.len(), 2); + assert!(batches.iter().all(|(_, bytes)| { + usize::try_from(*bytes).unwrap() <= super::super::MAX_OUTBOUND_BATCH_BYTES + })); + } + #[tokio::test] async fn sync_progress_is_incremental_and_scoped_to_the_active_validation() { let alice = Wallet::from_seed("sync-progress-alice"); @@ -407,6 +590,45 @@ mod tests { } #[tokio::test] + async fn identical_chain_validations_cannot_run_concurrently() { + let alice = Wallet::from_seed("validation-coordinator-alice"); + let allocations = allocations(std::slice::from_ref(&alice), 1_000); + let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); + let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); + let network = gossip_network(node, peers, "127.0.0.1:9544".parse().unwrap(), None); + + let first = network + .claim_chain_validation("same-page".to_string()) + .await + .unwrap(); + let waiting_network = network.clone(); + let waiting = tokio::spawn(async move { + waiting_network + .claim_chain_validation("same-page".to_string()) + .await + .unwrap() + }); + tokio::time::sleep(std::time::Duration::from_millis(25)).await; + assert!(!waiting.is_finished()); + + drop(first); + let second = tokio::time::timeout(std::time::Duration::from_secs(1), waiting) + .await + .unwrap() + .unwrap(); + drop(second); + assert!( + network + .inner + .chain_validation + .active + .lock() + .unwrap() + .is_empty() + ); + } + + #[tokio::test] async fn peer_exchange_does_not_advertise_self_when_outbound_only() { let alice = Wallet::from_seed("px-private-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); diff --git a/src/adapters/p2p/peer_addr.rs b/src/adapters/p2p/peer_addr.rs @@ -6,26 +6,10 @@ use std::{ use anyhow::{Context, Result}; -use crate::app::GossipEnvelope; - pub(super) fn next_reconnect_delay(current: Duration, max_delay: Duration) -> Duration { (current * 2).min(max_delay) } -pub(super) fn peer_has_block_gap(peer_height: u64, envelopes: &[GossipEnvelope]) -> bool { - envelopes - .iter() - .filter_map(|envelope| match envelope { - GossipEnvelope::Block(block) => Some(block.height), - GossipEnvelope::Inventory { blocks, .. } => { - blocks.iter().map(|block| block.height).min() - } - _ => None, - }) - .min() - .is_some_and(|first_block_height| peer_height + 1 < first_block_height) -} - pub(super) fn reachable_advertised_addr( advertised_addr: SocketAddr, remote_addr: SocketAddr, diff --git a/src/adapters/p2p/peer_status.rs b/src/adapters/p2p/peer_status.rs @@ -1,3 +1,4 @@ +#[cfg(test)] use crate::app::now_ms; #[derive(Clone, Debug, Eq, PartialEq)] @@ -6,10 +7,11 @@ pub(super) struct PeerStatus { pub(super) tip_hash: String, pub(super) time_ms: u64, pub(super) request_bootstrap: bool, - pub(super) push_bootstrap: bool, + pub(super) reject_session: bool, } impl PeerStatus { + #[cfg(test)] pub(super) fn new(height: u64, tip_hash: String) -> Self { Self::with_time(height, tip_hash, now_ms()) } @@ -20,7 +22,7 @@ impl PeerStatus { tip_hash, time_ms, request_bootstrap: false, - push_bootstrap: false, + reject_session: false, } } @@ -30,7 +32,7 @@ impl PeerStatus { tip_hash, time_ms, request_bootstrap: false, - push_bootstrap: false, + reject_session: false, } } @@ -40,17 +42,17 @@ impl PeerStatus { tip_hash, time_ms, request_bootstrap: true, - push_bootstrap: false, + reject_session: false, } } - pub(super) fn with_bootstrap_push(height: u64, tip_hash: String, time_ms: u64) -> Self { + pub(super) fn rejected(height: u64, tip_hash: String, time_ms: u64) -> Self { Self { height, tip_hash, time_ms, request_bootstrap: false, - push_bootstrap: true, + reject_session: true, } } } diff --git a/src/adapters/p2p/process.rs b/src/adapters/p2p/process.rs @@ -1,6 +1,7 @@ use std::net::SocketAddr; use anyhow::{Result, anyhow}; +use sha2::{Digest, Sha256}; use tokio::net::tcp::OwnedWriteHalf; use crate::{ @@ -12,7 +13,6 @@ use super::{ GossipNetwork, MAX_BLOCK_BATCH, P2pMetricsCounters, apply_peer_list, forget_stale_self_peer, is_possible_fork_error, normalize_advertised_peer, peer_verification_response, process_hello, validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf, write_envelope, - write_payload, }; pub(super) async fn respond_to_peer_verification_challenge( @@ -35,7 +35,8 @@ pub(super) async fn process_envelope( remote_addr: SocketAddr, known_peer: &mut Option<String>, envelope: GossipEnvelope, -) -> Result<()> { +) -> Result<bool> { + let mut requested_chain_data = false; match envelope { GossipEnvelope::Hello(hello) => { let _ = process_hello(network, remote_addr, known_peer, hello).await?; @@ -71,15 +72,6 @@ pub(super) async fn process_envelope( write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?; } } - GossipEnvelope::Inventory { blocks } => { - let requests = network - .inner - .node - .lock() - .await - .missing_inventory_requests(&blocks); - write_payload(writer, &requests).await?; - } GossipEnvelope::PeerAnnouncement { address, node_id } => { let peer = normalize_advertised_peer(&address, remote_addr)?; if network.is_self_peer(&peer).await { @@ -129,21 +121,45 @@ pub(super) async fn process_envelope( } GossipEnvelope::Block(block) => { let adjusted_time_ms = super::network_adjusted_time_ms(network).await; - let needs_vdf = { + let (needs_vdf, base_tip, sync_generation) = { let node = network.inner.node.lock().await; - node.block_requires_vdf_verification_at(&block, adjusted_time_ms) + ( + node.block_requires_vdf_verification_at(&block, adjusted_time_ms), + node.ledger().tip_hash().to_string(), + network.sync_generation(), + ) }; let result = match needs_vdf { Ok(false) => Ok(()), - Ok(true) => match verify_block_vdf(block).await { - Ok(block) => network - .inner - .node - .lock() - .await - .receive_preverified_block_at(block, adjusted_time_ms), - Err(error) => Err(error), - }, + Ok(true) => { + let validation_key = + chain_validation_key("block", &base_tip, std::slice::from_ref(&block)); + let Some(_validation) = network.claim_chain_validation(validation_key).await + else { + return Ok(false); + }; + let base_is_current = { + let node = network.inner.node.lock().await; + network.sync_generation_is_current(sync_generation) + && node.ledger().tip_hash() == base_tip + }; + if !base_is_current { + return Ok(false); + } + match verify_block_vdf(block).await { + Ok(block) => { + let mut node = network.inner.node.lock().await; + if network.sync_generation_is_current(sync_generation) + && node.ledger().tip_hash() == base_tip + { + node.receive_preverified_block_at(block, adjusted_time_ms) + } else { + Ok(()) + } + } + Err(error) => Err(error), + } + } Err(error) => Err(error), }; let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error); @@ -156,24 +172,48 @@ pub(super) async fn process_envelope( record_inbound_result(network, known_peer, remote_addr, result).await; if request_locator { request_fork_blocks(network, writer).await?; + requested_chain_data = true; } network.forward_outbox().await; } GossipEnvelope::Blocks { blocks } => { let adjusted_time_ms = super::network_adjusted_time_ms(network).await; + let (base_tip, sync_generation) = { + let node = network.inner.node.lock().await; + if blocks.is_empty() || node.ledger().contains_block_sequence(&blocks) { + drop(node); + record_inbound_result(network, known_peer, remote_addr, Ok(())).await; + return Ok(false); + } + ( + node.ledger().tip_hash().to_string(), + network.sync_generation(), + ) + }; + let validation_key = chain_validation_key("blocks", &base_tip, &blocks); + let Some(_validation) = network.claim_chain_validation(validation_key).await else { + return Ok(false); + }; let (local_ledger, progress_guard) = { let node = network.inner.node.lock().await; - let local_ledger = node.clone_ledger(); + if !network.sync_generation_is_current(sync_generation) + || node.ledger().tip_hash() != base_tip + || node.ledger().contains_block_sequence(&blocks) + { + drop(node); + record_inbound_result(network, known_peer, remote_addr, Ok(())).await; + return Ok(false); + } let start_height = blocks .first() .map(|block| block.height.saturating_sub(1)) - .unwrap_or_else(|| local_ledger.height()); + .unwrap_or_else(|| node.ledger().height()); let target_height = blocks .last() .map(|block| block.height) .unwrap_or(start_height); let progress_guard = network.begin_sync_progress(start_height, target_height); - (local_ledger, progress_guard) + (node.clone_ledger(), progress_guard) }; let progress_id = progress_guard.id(); let progress_network = network.clone(); @@ -187,7 +227,10 @@ pub(super) async fn process_envelope( { Ok(ledger) => { let mut node = network.inner.node.lock().await; - if progress_guard.is_current() { + if progress_guard.is_current() + && network.sync_generation_is_current(sync_generation) + && node.ledger().tip_hash() == base_tip + { node.import_verified_ledger(ledger).map(|_| ()) } else { Ok(()) @@ -206,16 +249,43 @@ pub(super) async fn process_envelope( record_inbound_result(network, known_peer, remote_addr, result).await; if request_locator { request_fork_blocks(network, writer).await?; + requested_chain_data = true; } network.forward_outbox().await; } GossipEnvelope::ChainBootstrap(bootstrap) => { let sync_generation = network.sync_generation(); + let base_tip = network + .inner + .node + .lock() + .await + .ledger() + .tip_hash() + .to_string(); + let validation_key = chain_validation_key( + "bootstrap", + &base_tip, + std::slice::from_ref(&bootstrap.genesis_block), + ); + let Some(_validation) = network.claim_chain_validation(validation_key).await else { + return Ok(false); + }; + let base_is_current = { + let node = network.inner.node.lock().await; + network.sync_generation_is_current(sync_generation) + && node.ledger().tip_hash() == base_tip + }; + if !base_is_current { + return Ok(false); + } let adjusted_time_ms = super::network_adjusted_time_ms(network).await; let result = match validate_chain_bootstrap(bootstrap, adjusted_time_ms).await { Ok(ledger) => { let mut node = network.inner.node.lock().await; - if network.sync_generation_is_current(sync_generation) { + if network.sync_generation_is_current(sync_generation) + && node.ledger().tip_hash() == base_tip + { node.import_verified_ledger(ledger).map(|_| ()) } else { Ok(()) @@ -238,7 +308,22 @@ pub(super) async fn process_envelope( network.forward_outbox().await; } } - Ok(()) + Ok(requested_chain_data) +} + +fn chain_validation_key(kind: &str, base_tip: &str, blocks: &[crate::domain::Block]) -> String { + let mut digest = Sha256::new(); + digest.update(kind.as_bytes()); + digest.update([0]); + digest.update(base_tip.as_bytes()); + for block in blocks { + digest.update(block.height.to_be_bytes()); + digest.update(block.prev_hash.as_bytes()); + digest.update([0]); + digest.update(block.hash.as_bytes()); + digest.update([0]); + } + format!("{:x}", digest.finalize()) } async fn request_fork_blocks(network: &GossipNetwork, writer: &mut OwnedWriteHalf) -> Result<()> { @@ -342,13 +427,15 @@ async fn record_inbound_result( } Err(error) => { let message = format!("{error:#}"); - if known_peer.is_some() { - let mut peers = network.inner.peers.lock().await; - if super::inbound_error_counts_as_misbehavior(&message) { + let mut peers = network.inner.peers.lock().await; + if super::inbound_error_counts_as_misbehavior(&message) { + if known_peer.is_some() { peers.record_misbehavior(&peer, message.clone()); } else { - peers.record_inbound_error(&peer, message.clone()); + peers.record_inbound_misbehavior(&peer, message.clone()); } + } else { + peers.record_inbound_error(&peer, message.clone()); } if debug_logging_enabled() { eprintln!("p2p envelope from {peer} ignored: {message}"); diff --git a/src/adapters/p2p/session.rs b/src/adapters/p2p/session.rs @@ -1,24 +1,30 @@ use std::{net::SocketAddr, time::Duration}; -use anyhow::Result; +use anyhow::{Context, Result, bail}; use tokio::{ net::{TcpListener, TcpStream}, - sync::mpsc, + sync::{mpsc, watch}, time::{Instant, interval, interval_at, sleep, timeout}, }; use crate::app::{GossipEnvelope, debug_logging_enabled}; use super::{ - CONNECT_TIMEOUT, GossipNetwork, HANDSHAKE_TIMEOUT, INITIAL_RECONNECT_DELAY, - MAX_RECONNECT_DELAY, PEER_EXCHANGE_INTERVAL, PEER_QUEUE_SIZE, PeerStatus, - SESSION_SYNC_INTERVAL, is_self_peer_address_for, next_reconnect_delay_with_max, - process_envelope, process_hello_with_verification, push_catchup_to_peer, read_session_envelope, - record_peer_status, respond_to_peer_verification_challenge, write_envelope, write_payload, - write_peer_exchange, + CATCHUP_REQUEST_TIMEOUT, CONNECT_TIMEOUT, GossipNetwork, GossipSession, HANDSHAKE_TIMEOUT, + INBOUND_PEER_QUEUE_BYTES, INBOUND_PEER_QUEUE_SIZE, INBOUND_SESSION_PREFIX, + INITIAL_RECONNECT_DELAY, MAX_RECONNECT_DELAY, OutboundBatch, PEER_EXCHANGE_INTERVAL, + PEER_QUEUE_BYTES, PEER_QUEUE_SIZE, PeerStatus, SESSION_SYNC_INTERVAL, is_self_peer_address_for, + next_reconnect_delay_with_max, process_envelope, process_hello_with_verification, + read_session_envelope, record_peer_status, respond_to_peer_verification_challenge, + write_envelope, write_payload, write_peer_exchange, }; -type OutboundBatch = Vec<GossipEnvelope>; +struct InboundRegistration { + key: String, + sender: mpsc::Sender<OutboundBatch>, + shutdown: watch::Sender<bool>, + queue_bytes: std::sync::Arc<tokio::sync::Semaphore>, +} pub(super) async fn accept_loop(network: GossipNetwork, listener: TcpListener) { loop { @@ -46,6 +52,17 @@ pub(super) async fn accept_loop(network: GossipNetwork, listener: TcpListener) { } }; super::P2pMetricsCounters::inc(&network.inner.metrics.inbound_sessions_started); + let session_key = format!("{INBOUND_SESSION_PREFIX}{remote_addr}"); + let (sender, receiver) = mpsc::channel(INBOUND_PEER_QUEUE_SIZE); + let (shutdown, mut shutdown_receiver) = watch::channel(false); + let queue_bytes = + std::sync::Arc::new(tokio::sync::Semaphore::new(INBOUND_PEER_QUEUE_BYTES)); + let registration = InboundRegistration { + key: session_key.clone(), + sender, + shutdown, + queue_bytes, + }; tokio::spawn(async move { let _permit = permit; let result = session_loop( @@ -53,9 +70,12 @@ pub(super) async fn accept_loop(network: GossipNetwork, listener: TcpListener) { stream, remote_addr, None, - mpsc::channel(1).1, + receiver, + &mut shutdown_receiver, + Some(registration), ) .await; + network.inner.sessions.lock().await.remove(&session_key); match result { Ok(()) => { super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); @@ -98,6 +118,7 @@ pub(super) async fn outbound_session( network: GossipNetwork, peer: String, mut receiver: mpsc::Receiver<OutboundBatch>, + mut shutdown: watch::Receiver<bool>, ) { let mut reconnect_delay = INITIAL_RECONNECT_DELAY; loop { @@ -156,8 +177,14 @@ pub(super) async fn outbound_session( remote_addr, Some(peer.clone()), receiver, + &mut shutdown, + None, ) .await; + if *shutdown.borrow() { + network.inner.sessions.lock().await.remove(&peer); + return; + } match result { Ok(()) => { super::P2pMetricsCounters::inc(&network.inner.metrics.sessions_closed); @@ -185,17 +212,23 @@ pub(super) async fn outbound_session( } let (sender, next_receiver) = mpsc::channel(PEER_QUEUE_SIZE); + let (next_shutdown, next_shutdown_receiver) = watch::channel(false); + let queue_bytes = std::sync::Arc::new(tokio::sync::Semaphore::new(PEER_QUEUE_BYTES)); receiver = next_receiver; + shutdown = next_shutdown_receiver; if !peer_is_connectable(&network, &peer).await { network.inner.sessions.lock().await.remove(&peer); return; } - network - .inner - .sessions - .lock() - .await - .insert(peer.clone(), sender); + network.inner.sessions.lock().await.insert( + peer.clone(), + GossipSession { + peer: peer.clone(), + sender, + shutdown: next_shutdown, + queue_bytes, + }, + ); sleep(reconnect_delay).await; reconnect_delay = next_reconnect_delay(reconnect_delay); } @@ -207,6 +240,8 @@ async fn session_loop( remote_addr: SocketAddr, stable_peer: Option<String>, mut outbound: mpsc::Receiver<OutboundBatch>, + shutdown: &mut watch::Receiver<bool>, + inbound_registration: Option<InboundRegistration>, ) -> Result<()> { let (reader, mut writer) = stream.into_split(); let connection_label = stable_peer @@ -230,10 +265,68 @@ async fn session_loop( ); let mut outbound_closed = false; let mut peer_status: Option<PeerStatus> = None; + let mut catchup_requested_at: Option<Instant> = None; let is_outbound_session = stable_peer.is_some(); let mut known_peer = stable_peer; + let mut handshake_complete = false; + let mut shutdown_closed = false; + + if !is_outbound_session { + let envelope = timeout( + HANDSHAKE_TIMEOUT, + read_session_envelope(&network, &connection_label, &mut reader), + ) + .await + .context("inbound p2p handshake timed out")?? + .context("inbound peer closed before sending Hello")?; + let hello = match envelope { + GossipEnvelope::Hello(hello) => hello, + challenge @ GossipEnvelope::PeerVerificationChallenge { .. } => { + respond_to_peer_verification_challenge(&network, &mut writer, &challenge).await?; + return Ok(()); + } + _ => bail!("inbound peer sent data before Hello"), + }; + let status = process_hello_with_verification( + &network, + &mut writer, + &mut reader, + &connection_label, + remote_addr, + &mut known_peer, + hello, + ) + .await?; + if status.reject_session { + return Ok(()); + } + peer_status = Some(status); + handshake_complete = true; + if let Some(registration) = inbound_registration { + let peer = known_peer + .clone() + .unwrap_or_else(|| remote_addr.to_string()); + network.inner.sessions.lock().await.insert( + registration.key, + GossipSession { + peer, + sender: registration.sender, + shutdown: registration.shutdown, + queue_bytes: registration.queue_bytes, + }, + ); + } + maybe_start_catchup( + &network, + &mut writer, + peer_status.as_ref().unwrap(), + &mut catchup_requested_at, + ) + .await?; + write_peer_exchange(&network, &mut writer, &known_peer).await?; + } - if known_peer.is_some() { + if is_outbound_session { if let Ok(Ok(Some(envelope))) = timeout( HANDSHAKE_TIMEOUT, read_session_envelope(&network, &connection_label, &mut reader), @@ -241,23 +334,31 @@ async fn session_loop( .await { if let GossipEnvelope::Hello(hello) = envelope { - peer_status = Some( - process_hello_with_verification( - &network, - &mut writer, - &mut reader, - &connection_label, - remote_addr, - &mut known_peer, - hello, - ) - .await?, - ); + let status = process_hello_with_verification( + &network, + &mut writer, + &mut reader, + &connection_label, + remote_addr, + &mut known_peer, + hello, + ) + .await?; + if status.reject_session { + return Ok(()); + } + peer_status = Some(status); + handshake_complete = true; if is_outbound_session && known_peer.is_none() { return Ok(()); } - super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()) - .await?; + maybe_start_catchup( + &network, + &mut writer, + peer_status.as_ref().unwrap(), + &mut catchup_requested_at, + ) + .await?; write_peer_exchange(&network, &mut writer, &known_peer).await?; } else if let GossipEnvelope::PeerStatus { height, @@ -268,8 +369,14 @@ async fn session_loop( let status = PeerStatus::from_envelope(height, tip_hash, time_ms); record_peer_status(&network, &known_peer, remote_addr, &status).await; peer_status = Some(status); - super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()) - .await?; + handshake_complete = true; + maybe_start_catchup( + &network, + &mut writer, + peer_status.as_ref().unwrap(), + &mut catchup_requested_at, + ) + .await?; write_peer_exchange(&network, &mut writer, &known_peer).await?; } else if respond_to_peer_verification_challenge(&network, &mut writer, &envelope) .await? @@ -277,8 +384,15 @@ async fn session_loop( if known_peer.is_none() { return Ok(()); } - } else { - process_envelope( + } else if !maybe_request_inventory( + &network, + &mut writer, + &envelope, + &mut catchup_requested_at, + ) + .await? + { + let requested_chain_data = process_envelope( &network, &mut writer, remote_addr, @@ -286,6 +400,9 @@ async fn session_loop( envelope, ) .await?; + if requested_chain_data { + catchup_requested_at = Some(Instant::now()); + } if is_outbound_session && known_peer.is_none() { return Ok(()); } @@ -295,13 +412,21 @@ async fn session_loop( loop { tokio::select! { + result = shutdown.changed(), if !shutdown_closed => { + match result { + Ok(()) if *shutdown.borrow() => return Ok(()), + Ok(()) => {} + Err(_) if !is_outbound_session => return Ok(()), + Err(_) => shutdown_closed = true, + } + } maybe_batch = outbound.recv(), if !outbound_closed => { match maybe_batch { Some(batch) => { let payload = super::envelopes_for_peer( Some(&network.inner.node), peer_status.clone(), - &batch, + &batch.envelopes, ).await; write_payload(&mut writer, &payload).await?; if let Some(peer) = &known_peer { @@ -314,10 +439,11 @@ async fn session_loop( _ = sync_tick.tick() => { let status = network.inner.node.lock().await.peer_status(); write_envelope(&mut writer, &status).await?; - if let Some(status) = peer_status.as_mut() { - if let Some(updated_status) = push_catchup_to_peer(&network, &mut writer, status).await? { - *status = updated_status; - } + if catchup_requested_at.is_some_and(|started| started.elapsed() >= CATCHUP_REQUEST_TIMEOUT) { + catchup_requested_at = None; + } + if let Some(status) = peer_status.as_ref() { + maybe_start_catchup(&network, &mut writer, status, &mut catchup_requested_at).await?; } } _ = peer_exchange_tick.tick() => { @@ -328,8 +454,10 @@ async fn session_loop( return Ok(()); }; if let GossipEnvelope::Hello(hello) = envelope { - peer_status = Some( - process_hello_with_verification( + if handshake_complete { + bail!("peer sent duplicate Hello"); + } + let status = process_hello_with_verification( &network, &mut writer, &mut reader, @@ -338,12 +466,21 @@ async fn session_loop( &mut known_peer, hello, ) - .await?, - ); + .await?; + if status.reject_session { + return Ok(()); + } + peer_status = Some(status); + handshake_complete = true; if is_outbound_session && known_peer.is_none() { return Ok(()); } - super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?; + maybe_start_catchup( + &network, + &mut writer, + peer_status.as_ref().unwrap(), + &mut catchup_requested_at, + ).await?; write_peer_exchange(&network, &mut writer, &known_peer).await?; continue; } @@ -356,7 +493,12 @@ async fn session_loop( let status = PeerStatus::from_envelope(*height, tip_hash.clone(), *time_ms); record_peer_status(&network, &known_peer, remote_addr, &status).await; peer_status = Some(status); - super::maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?; + maybe_start_catchup( + &network, + &mut writer, + peer_status.as_ref().unwrap(), + &mut catchup_requested_at, + ).await?; write_peer_exchange(&network, &mut writer, &known_peer).await?; continue; } @@ -367,13 +509,70 @@ async fn session_loop( } continue; } - process_envelope( + if maybe_request_inventory( + &network, + &mut writer, + &envelope, + &mut catchup_requested_at, + ) + .await? + { + continue; + } + let continue_catchup = matches!( + &envelope, + GossipEnvelope::Block(_) | GossipEnvelope::Blocks { .. } | GossipEnvelope::ChainBootstrap(_) + ); + let completes_catchup_request = matches!( + &envelope, + GossipEnvelope::Blocks { .. } | GossipEnvelope::ChainBootstrap(_) + ); + let (height_before, response_already_applied) = if continue_catchup { + let node = network.inner.node.lock().await; + let status = node.ledger().status(); + let already_applied = match &envelope { + GossipEnvelope::Blocks { blocks } => { + !blocks.is_empty() && node.ledger().contains_block_sequence(blocks) + } + _ => false, + }; + (Some((status.height, status.tip_hash)), already_applied) + } else { + (None, false) + }; + let requested_chain_data = process_envelope( &network, &mut writer, remote_addr, &mut known_peer, envelope, ).await?; + let chain_changed = if let Some(height_before) = height_before { + let status = network.inner.node.lock().await.ledger().status(); + status.height != height_before.0 || status.tip_hash != height_before.1 + } else { + false + }; + let response_satisfied = chain_changed || response_already_applied; + update_catchup_request_state( + &mut catchup_requested_at, + completes_catchup_request, + response_satisfied, + requested_chain_data, + ); + if continue_catchup && response_satisfied && !requested_chain_data { + if let Some(status) = peer_status.as_ref() { + maybe_start_catchup( + &network, + &mut writer, + status, + &mut catchup_requested_at, + ).await?; + } + } + if session_peer_is_banned(&network, &known_peer, remote_addr).await { + return Ok(()); + } if is_outbound_session && known_peer.is_none() { return Ok(()); } @@ -382,10 +581,72 @@ async fn session_loop( } } +async fn maybe_start_catchup( + network: &GossipNetwork, + writer: &mut tokio::net::tcp::OwnedWriteHalf, + peer_status: &PeerStatus, + requested_at: &mut Option<Instant>, +) -> Result<()> { + if requested_at.is_none() && super::maybe_request_catchup(network, writer, peer_status).await? { + *requested_at = Some(Instant::now()); + } + Ok(()) +} + +fn update_catchup_request_state( + requested_at: &mut Option<Instant>, + completes_request: bool, + response_satisfied: bool, + requested_chain_data: bool, +) { + if requested_chain_data || (completes_request && !response_satisfied) { + *requested_at = Some(Instant::now()); + } else if completes_request { + *requested_at = None; + } +} + +async fn maybe_request_inventory( + network: &GossipNetwork, + writer: &mut tokio::net::tcp::OwnedWriteHalf, + envelope: &GossipEnvelope, + requested_at: &mut Option<Instant>, +) -> Result<bool> { + let GossipEnvelope::Inventory { blocks } = envelope else { + return Ok(false); + }; + if requested_at.is_some() { + return Ok(true); + } + let request = network + .inner + .node + .lock() + .await + .missing_inventory_request(blocks); + if let Some(request) = request { + write_envelope(writer, &request).await?; + *requested_at = Some(Instant::now()); + } + Ok(true) +} + async fn peer_is_connectable(network: &GossipNetwork, peer: &str) -> bool { network.inner.peers.lock().await.is_connectable_peer(peer) } +async fn session_peer_is_banned( + network: &GossipNetwork, + known_peer: &Option<String>, + remote_addr: SocketAddr, +) -> bool { + let peer = known_peer + .as_deref() + .map(str::to_owned) + .unwrap_or_else(|| remote_addr.to_string()); + network.inner.peers.lock().await.is_banned(&peer) +} + pub(super) fn next_reconnect_delay(current: Duration) -> Duration { next_reconnect_delay_with_max(current, MAX_RECONNECT_DELAY) } @@ -394,8 +655,10 @@ pub(super) fn next_reconnect_delay(current: Duration) -> Duration { mod tests { use std::time::Duration; + use tokio::time::Instant; + use super::super::{INITIAL_RECONNECT_DELAY, MAX_RECONNECT_DELAY}; - use super::next_reconnect_delay; + use super::{next_reconnect_delay, update_catchup_request_state}; #[test] fn reconnect_backoff_is_capped() { @@ -408,4 +671,22 @@ mod tests { MAX_RECONNECT_DELAY ); } + + #[test] + fn already_applied_response_clears_in_flight_request() { + let mut requested_at = Some(Instant::now()); + + update_catchup_request_state(&mut requested_at, true, true, false); + + assert!(requested_at.is_none()); + } + + #[test] + fn fork_recovery_request_is_marked_in_flight() { + let mut requested_at = None; + + update_catchup_request_state(&mut requested_at, false, false, true); + + assert!(requested_at.is_some()); + } } diff --git a/src/adapters/p2p/sync.rs b/src/adapters/p2p/sync.rs @@ -5,20 +5,16 @@ use tokio::net::tcp::OwnedWriteHalf; use super::metrics::P2pMetricsCounters; use super::peer_addr::{ - is_self_peer_address_for, normalize_advertised_peer, peer_has_block_gap, - peer_list_address_is_discoverable, -}; -use super::{ - GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, byte_bounded_block_page, write_envelope, - write_payload, + is_self_peer_address_for, normalize_advertised_peer, peer_list_address_is_discoverable, }; +use super::{GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, write_envelope}; use crate::app::{GossipEnvelope, SharedNode, debug_logging_enabled}; pub(super) async fn maybe_request_catchup( network: &GossipNetwork, writer: &mut OwnedWriteHalf, peer_status: &PeerStatus, -) -> Result<()> { +) -> Result<bool> { let (local_height, local_tip_hash) = { let node = network.inner.node.lock().await; let status = node.ledger().status(); @@ -26,6 +22,7 @@ pub(super) async fn maybe_request_catchup( }; if peer_status.request_bootstrap { write_envelope(writer, &GossipEnvelope::ChainBootstrapRequest).await?; + return Ok(true); } else if peer_status.height > local_height { write_envelope( writer, @@ -35,6 +32,7 @@ pub(super) async fn maybe_request_catchup( }, ) .await?; + return Ok(true); } else if peer_status.height == local_height && peer_status.tip_hash != local_tip_hash { let locator = network.inner.node.lock().await.block_locator(); write_envelope( @@ -45,58 +43,9 @@ pub(super) async fn maybe_request_catchup( }, ) .await?; + return Ok(true); } - Ok(()) -} - -pub(super) async fn push_catchup_to_peer( - network: &GossipNetwork, - writer: &mut OwnedWriteHalf, - peer_status: &PeerStatus, -) -> Result<Option<PeerStatus>> { - let payload = catchup_payload_for_peer(&network.inner.node, peer_status).await; - if payload.is_empty() { - return Ok(None); - } - - let updated_status = payload.iter().find_map(|envelope| match envelope { - GossipEnvelope::Blocks { blocks } => blocks - .last() - .map(|block| PeerStatus::new(block.height, block.hash.clone())), - GossipEnvelope::ChainBootstrap(bootstrap) => { - Some(PeerStatus::new(0, bootstrap.genesis_block.hash.clone())) - } - _ => None, - }); - write_payload(writer, &payload).await?; - Ok(updated_status) -} - -pub(super) async fn catchup_payload_for_peer( - node: &SharedNode, - peer_status: &PeerStatus, -) -> Vec<GossipEnvelope> { - let mut node = node.lock().await; - let local_status = node.ledger().status(); - if node.ledger().is_setup_placeholder() { - return Vec::new(); - } - if peer_status.push_bootstrap { - return vec![GossipEnvelope::ChainBootstrap(node.chain_bootstrap())]; - } - if peer_status.height < local_status.height { - let blocks = - byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); - return (!blocks.is_empty()) - .then_some(GossipEnvelope::Blocks { blocks }) - .into_iter() - .collect(); - } else if peer_status.height == local_status.height - && peer_status.tip_hash != local_status.tip_hash - { - return Vec::new(); - } - node.mempool_gossip() + Ok(false) } pub(super) async fn apply_peer_list( @@ -168,11 +117,15 @@ pub(super) async fn envelopes_for_peer( .collect(); } if peer_status.height < local_status.height { - let blocks = - byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); - return (!blocks.is_empty()) - .then_some(GossipEnvelope::Blocks { blocks }) - .into_iter() + return envelopes + .iter() + .filter(|envelope| { + matches!( + envelope, + GossipEnvelope::PeerStatus { .. } | GossipEnvelope::Inventory { .. } + ) + }) + .cloned() .collect(); } @@ -180,12 +133,6 @@ pub(super) async fn envelopes_for_peer( return Vec::new(); } - if peer_has_block_gap(peer_status.height, envelopes) { - let blocks = - byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); - return vec![GossipEnvelope::Blocks { blocks }]; - } - envelopes .iter() .filter(|envelope| match envelope { diff --git a/src/adapters/p2p/test_support.rs b/src/adapters/p2p/test_support.rs @@ -56,6 +56,7 @@ pub(super) fn gossip_network( inbound_limiter: Arc::new(std::sync::Mutex::new(InboundConnectionLimiter::default())), metrics: P2pMetricsCounters::default(), sync_progress: std::sync::Mutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), } } diff --git a/src/adapters/p2p/tests.rs b/src/adapters/p2p/tests.rs @@ -11,7 +11,7 @@ use crate::{ }, domain::{Ledger, Wallet, run_vdf}, }; -use tokio::io::AsyncWriteExt; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use super::test_support::{allocations, gossip_network, node, queue_plaintext_burn}; @@ -45,6 +45,340 @@ async fn block_batch_validation_reports_each_validated_height() { } #[tokio::test] +async fn session_serializes_inventory_and_paginated_catchup_requests() { + let alice = Wallet::from_seed("immediate-page-sync-alice"); + let allocations = allocations(std::slice::from_ref(&alice), 1_000); + let local = Arc::new(tokio::sync::Mutex::new(node( + "immediate-page-local", + alice.clone(), + allocations.clone(), + ))); + let mut remote = node("immediate-page-remote", alice.clone(), allocations); + for timestamp_ms in [1, 2] { + queue_plaintext_burn(&mut remote, &alice, 1); + remote.drain_outbox(); + remote.mine_one_at(timestamp_ms).unwrap(); + remote.drain_outbox(); + } + let hello = remote.hello(None, None); + let blocks = remote.ledger().blocks_from(1, 2); + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let peer_addr = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let (reader, mut writer) = stream.into_split(); + let mut lines = BufReader::new(reader).lines(); + let first = lines.next_line().await.unwrap().unwrap(); + assert!(matches!( + super::parse_envelope(&first).unwrap(), + GossipEnvelope::Hello(_) + )); + super::write_envelope(&mut writer, &hello).await.unwrap(); + + let mut sent_first_page = false; + loop { + let line = tokio::time::timeout(std::time::Duration::from_secs(30), lines.next_line()) + .await + .expect("next block request did not follow the imported page") + .unwrap() + .unwrap(); + match super::parse_envelope(&line).unwrap() { + GossipEnvelope::BlockRangeRequest { from_height: 1, .. } => { + assert!(!sent_first_page); + sent_first_page = true; + super::write_envelope( + &mut writer, + &GossipEnvelope::Inventory { + blocks: vec![crate::app::BlockInventory { + height: blocks[1].height, + hash: blocks[1].hash.clone(), + }], + }, + ) + .await + .unwrap(); + let duplicate_request = + tokio::time::timeout(std::time::Duration::from_millis(200), async { + loop { + let line = lines.next_line().await.unwrap().unwrap(); + let envelope = super::parse_envelope(&line).unwrap(); + if matches!( + envelope, + GossipEnvelope::BlockRequest { .. } + | GossipEnvelope::BlockRangeRequest { .. } + | GossipEnvelope::BlockLocatorRequest { .. } + ) { + break envelope; + } + } + }) + .await; + assert!( + duplicate_request.is_err(), + "Inventory bypassed the outstanding paginated catchup request: {duplicate_request:?}" + ); + super::write_envelope( + &mut writer, + &GossipEnvelope::Blocks { + blocks: vec![blocks[0].clone()], + }, + ) + .await + .unwrap(); + } + GossipEnvelope::BlockRangeRequest { from_height: 2, .. } => { + assert!(sent_first_page); + super::write_envelope( + &mut writer, + &GossipEnvelope::Blocks { + blocks: vec![blocks[1].clone()], + }, + ) + .await + .unwrap(); + break; + } + _ => {} + } + } + }); + + let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ + peer_addr.to_string(), + ]))); + let _network = super::GossipNetwork::start( + Arc::clone(&local), + peers, + "127.0.0.1:0".parse().unwrap(), + None, + false, + ) + .await + .unwrap(); + + tokio::time::timeout(std::time::Duration::from_secs(60), server) + .await + .unwrap() + .unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(30), async { + loop { + if local.lock().await.ledger().height() == 2 { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); +} + +#[tokio::test] +async fn invalid_bootstrap_response_is_not_retried_without_backoff() { + let wallet = Wallet::from_seed("invalid-bootstrap-backoff"); + let remote = node( + "invalid-bootstrap-remote", + wallet.clone(), + allocations(std::slice::from_ref(&wallet), 1_000), + ); + let hello = remote.hello(None, None); + let mut invalid_bootstrap = remote.chain_bootstrap(); + invalid_bootstrap.genesis_block.hash = "0".repeat(64); + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let peer_addr = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let (reader, mut writer) = stream.into_split(); + let mut lines = BufReader::new(reader).lines(); + let first = lines.next_line().await.unwrap().unwrap(); + assert!(matches!( + super::parse_envelope(&first).unwrap(), + GossipEnvelope::Hello(_) + )); + super::write_envelope(&mut writer, &hello).await.unwrap(); + + loop { + let line = lines.next_line().await.unwrap().unwrap(); + if matches!( + super::parse_envelope(&line).unwrap(), + GossipEnvelope::ChainBootstrapRequest + ) { + break; + } + } + super::write_envelope( + &mut writer, + &GossipEnvelope::ChainBootstrap(invalid_bootstrap), + ) + .await + .unwrap(); + + tokio::time::timeout(std::time::Duration::from_secs(3), async { + loop { + let Some(line) = lines.next_line().await.unwrap() else { + break; + }; + assert!( + !matches!( + super::parse_envelope(&line).unwrap(), + GossipEnvelope::ChainBootstrapRequest + ), + "invalid bootstrap was retried immediately" + ); + } + }) + .await + .ok(); + }); + + let local = Arc::new(tokio::sync::Mutex::new(NodeCore::from_ledger( + wallet, + Ledger::new(BTreeMap::new(), 1), + 0, + ))); + let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![ + peer_addr.to_string(), + ]))); + let _network = + super::GossipNetwork::start(local, peers, "127.0.0.1:0".parse().unwrap(), None, false) + .await + .unwrap(); + + tokio::time::timeout(std::time::Duration::from_secs(5), server) + .await + .unwrap() + .unwrap(); +} + +#[tokio::test] +async fn inbound_session_receives_new_block_inventory_without_waiting_for_status_tick() { + let alice = Wallet::from_seed("inbound-fast-block-relay-alice"); + let allocations = allocations(std::slice::from_ref(&alice), 1_000); + let mut source = node( + "inbound-fast-block-relay-source", + alice.clone(), + allocations.clone(), + ); + let remote = node( + "inbound-fast-block-relay-remote", + alice.clone(), + allocations, + ); + queue_plaintext_burn(&mut source, &alice, 1); + source.drain_outbox(); + let block = source.mine_one_at(1).unwrap(); + source.drain_outbox(); + + let reserved = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let listen_addr = reserved.local_addr().unwrap(); + drop(reserved); + let network = super::GossipNetwork::start( + Arc::new(tokio::sync::Mutex::new(source)), + Arc::new(tokio::sync::Mutex::new(PeerBook::default())), + listen_addr, + None, + true, + ) + .await + .unwrap(); + + let stream = tokio::net::TcpStream::connect(listen_addr).await.unwrap(); + let (reader, mut writer) = stream.into_split(); + let mut lines = BufReader::new(reader).lines(); + let hello = lines.next_line().await.unwrap().unwrap(); + assert!(matches!( + super::parse_envelope(&hello).unwrap(), + GossipEnvelope::Hello(_) + )); + + network + .broadcast(vec![GossipEnvelope::Block(block.clone())]) + .await + .unwrap(); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(150), lines.next_line()) + .await + .is_err(), + "inbound peer received gossip before completing Hello" + ); + + super::write_envelope(&mut writer, &remote.hello(None, None)) + .await + .unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + if !network.inner.sessions.lock().await.is_empty() { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("inbound session was not registered after Hello"); + + network + .broadcast(vec![GossipEnvelope::Block(block.clone())]) + .await + .unwrap(); + let relayed = tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + let line = lines.next_line().await.unwrap().unwrap(); + if let GossipEnvelope::Inventory { blocks } = super::parse_envelope(&line).unwrap() { + break blocks; + } + } + }) + .await + .expect("inbound block relay waited for periodic anti-entropy"); + + assert!(relayed.iter().any(|item| item.hash == block.hash)); + network.set_accept_inbound(false).await.unwrap(); +} + +#[tokio::test] +async fn inbound_self_connection_is_closed_before_relay_registration() { + let wallet = Wallet::from_seed("inbound-self-connection"); + let reserved = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let listen_addr = reserved.local_addr().unwrap(); + drop(reserved); + let network = super::GossipNetwork::start( + Arc::new(tokio::sync::Mutex::new(node( + "inbound-self-connection", + wallet.clone(), + allocations(&[wallet], 1_000), + ))), + Arc::new(tokio::sync::Mutex::new(PeerBook::default())), + listen_addr, + None, + true, + ) + .await + .unwrap(); + + let stream = tokio::net::TcpStream::connect(listen_addr).await.unwrap(); + let (reader, mut writer) = stream.into_split(); + let mut lines = BufReader::new(reader).lines(); + let local_hello = + match super::parse_envelope(&lines.next_line().await.unwrap().unwrap()).unwrap() { + GossipEnvelope::Hello(hello) => hello, + other => panic!("expected Hello, got {other:?}"), + }; + super::write_envelope(&mut writer, &GossipEnvelope::Hello(local_hello)) + .await + .unwrap(); + + let closed = tokio::time::timeout(std::time::Duration::from_secs(1), lines.next_line()) + .await + .expect("self connection remained open") + .unwrap(); + assert!(closed.is_none()); + assert!(network.inner.sessions.lock().await.is_empty()); + assert_eq!(network.metrics().self_peer_rejections, 1); + network.set_accept_inbound(false).await.unwrap(); +} + +#[tokio::test] async fn full_outbound_queue_is_metric_not_peer_error() { let wallet = Wallet::from_seed("full-outbound-queue"); let node = Arc::new(tokio::sync::Mutex::new(node( @@ -67,22 +401,31 @@ async fn full_outbound_queue_is_metric_not_peer_error() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let (sender, _receiver) = tokio::sync::mpsc::channel(1); + let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::PEER_QUEUE_BYTES)); + let queued_bytes = Arc::clone(&queue_bytes).try_acquire_many_owned(1).unwrap(); sender - .try_send(vec![GossipEnvelope::PeerStatus { - height: 1, - tip_hash: "queued".to_string(), - time_ms: 1_000, - }]) + .try_send(super::OutboundBatch { + envelopes: Arc::from(vec![GossipEnvelope::PeerStatus { + height: 1, + tip_hash: "queued".to_string(), + time_ms: 1_000, + }]), + _queued_bytes: queued_bytes, + }) .unwrap(); - network - .inner - .sessions - .lock() - .await - .insert("127.0.0.1:9444".to_string(), sender); + network.inner.sessions.lock().await.insert( + "127.0.0.1:9444".to_string(), + super::GossipSession { + peer: "127.0.0.1:9444".to_string(), + sender, + shutdown: tokio::sync::watch::channel(false).0, + queue_bytes, + }, + ); network .broadcast(vec![GossipEnvelope::PeerStatus { @@ -100,6 +443,167 @@ async fn full_outbound_queue_is_metric_not_peer_error() { } #[tokio::test] +async fn full_inbound_queue_disconnects_the_session() { + let wallet = Wallet::from_seed("full-inbound-queue"); + let network = gossip_network( + Arc::new(tokio::sync::Mutex::new(node( + "full-inbound-queue", + wallet.clone(), + allocations(&[wallet], 1_000), + ))), + Arc::new(tokio::sync::Mutex::new(PeerBook::default())), + "127.0.0.1:9544".parse().unwrap(), + None, + ); + let session_id = format!("{}127.0.0.1:51234", super::INBOUND_SESSION_PREFIX); + let (sender, _receiver) = tokio::sync::mpsc::channel(1); + let (shutdown, mut shutdown_receiver) = tokio::sync::watch::channel(false); + let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::INBOUND_PEER_QUEUE_BYTES)); + let queued_bytes = Arc::clone(&queue_bytes).try_acquire_many_owned(1).unwrap(); + sender + .try_send(super::OutboundBatch { + envelopes: Arc::from(vec![GossipEnvelope::PeerStatus { + height: 1, + tip_hash: "queued".to_string(), + time_ms: 1_000, + }]), + _queued_bytes: queued_bytes, + }) + .unwrap(); + network.inner.sessions.lock().await.insert( + session_id.clone(), + super::GossipSession { + peer: "127.0.0.1:51234".to_string(), + sender, + shutdown, + queue_bytes, + }, + ); + + network + .broadcast(vec![GossipEnvelope::PeerStatus { + height: 2, + tip_hash: "new".to_string(), + time_ms: 2_000, + }]) + .await + .unwrap(); + + assert_eq!(network.metrics().outbound_queue_full, 1); + assert!( + !network + .inner + .sessions + .lock() + .await + .contains_key(&session_id) + ); + shutdown_receiver.changed().await.unwrap(); + assert!(*shutdown_receiver.borrow()); +} + +#[tokio::test] +async fn exhausted_inbound_byte_budget_disconnects_the_session() { + let wallet = Wallet::from_seed("inbound-byte-budget"); + let network = gossip_network( + Arc::new(tokio::sync::Mutex::new(node( + "inbound-byte-budget", + wallet.clone(), + allocations(&[wallet], 1_000), + ))), + Arc::new(tokio::sync::Mutex::new(PeerBook::default())), + "127.0.0.1:9544".parse().unwrap(), + None, + ); + let session_id = format!("{}127.0.0.1:51235", super::INBOUND_SESSION_PREFIX); + let (sender, mut receiver) = tokio::sync::mpsc::channel(4); + let (shutdown, mut shutdown_receiver) = tokio::sync::watch::channel(false); + network.inner.sessions.lock().await.insert( + session_id.clone(), + super::GossipSession { + peer: "127.0.0.1:51235".to_string(), + sender, + shutdown, + queue_bytes: Arc::new(tokio::sync::Semaphore::new(1)), + }, + ); + + network + .broadcast(vec![GossipEnvelope::PeerStatus { + height: 2, + tip_hash: "larger-than-one-byte".to_string(), + time_ms: 2_000, + }]) + .await + .unwrap(); + + assert_eq!(network.metrics().outbound_queue_full, 1); + assert!( + !network + .inner + .sessions + .lock() + .await + .contains_key(&session_id) + ); + shutdown_receiver.changed().await.unwrap(); + assert!(*shutdown_receiver.borrow()); + assert!(matches!( + receiver.try_recv(), + Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) + )); +} + +#[tokio::test] +async fn broadcast_skips_banned_inbound_session_identity() { + let wallet = Wallet::from_seed("banned-inbound-relay"); + let node = Arc::new(tokio::sync::Mutex::new(node( + "banned-inbound-relay", + wallet.clone(), + allocations(&[wallet], 1_000), + ))); + let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default())); + let network = gossip_network( + node, + Arc::clone(&peers), + "127.0.0.1:9544".parse().unwrap(), + None, + ); + let peer = "127.0.0.1:51234"; + for _ in 0..crate::app::PEER_MISBEHAVIOR_BAN_SCORE { + peers + .lock() + .await + .record_inbound_misbehavior(peer, "invalid block"); + } + let (sender, mut receiver) = tokio::sync::mpsc::channel(1); + let queue_bytes = Arc::new(tokio::sync::Semaphore::new(super::INBOUND_PEER_QUEUE_BYTES)); + network.inner.sessions.lock().await.insert( + format!("{}{peer}", super::INBOUND_SESSION_PREFIX), + super::GossipSession { + peer: peer.to_string(), + sender, + shutdown: tokio::sync::watch::channel(false).0, + queue_bytes, + }, + ); + + network + .broadcast(vec![GossipEnvelope::PeerStatus { + height: 1, + tip_hash: "tip".to_string(), + time_ms: 1_000, + }]) + .await + .unwrap(); + + assert!(matches!( + receiver.try_recv(), + Err(tokio::sync::mpsc::error::TryRecvError::Empty) + )); +} + +#[tokio::test] async fn single_block_fork_error_requests_blocks_by_locator() { let alice = Wallet::from_seed("single-block-fork-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); @@ -144,7 +648,7 @@ async fn single_block_fork_error_requests_blocks_by_locator() { let mut client_reader = super::LimitedLineReader::new(client_reader); let mut known_peer = None; - super::process_envelope( + let requested_chain_data = super::process_envelope( &network, &mut server_writer, remote_addr, @@ -153,6 +657,7 @@ async fn single_block_fork_error_requests_blocks_by_locator() { ) .await .unwrap(); + assert!(requested_chain_data); let line = tokio::time::timeout(std::time::Duration::from_secs(1), client_reader.read_line()) .await @@ -215,7 +720,7 @@ async fn block_page_without_local_ancestor_requests_blocks_by_locator() { let mut client_reader = super::LimitedLineReader::new(client_reader); let mut known_peer = None; - super::process_envelope( + let requested_chain_data = super::process_envelope( &network, &mut server_writer, remote_addr, @@ -226,6 +731,7 @@ async fn block_page_without_local_ancestor_requests_blocks_by_locator() { ) .await .unwrap(); + assert!(requested_chain_data); let line = tokio::time::timeout(std::time::Duration::from_secs(1), client_reader.read_line()) .await @@ -476,6 +982,7 @@ async fn hello_rejects_wrong_network_or_genesis_without_banning() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; @@ -582,6 +1089,7 @@ async fn hello_records_remote_clock_observation() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let remote_time_ms = crate::app::now_ms().saturating_add(60_000); @@ -640,6 +1148,7 @@ async fn hello_remembers_advertised_address_after_signed_session_and_dialback() inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let remote_node_id = super::new_node_id(); @@ -717,6 +1226,7 @@ async fn hello_ignores_advertised_address_when_connected_peer_cannot_sign_claime inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let victim_node_id = super::new_node_id(); @@ -788,6 +1298,7 @@ async fn dialback_rejects_address_that_signs_with_different_node_id() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let honest_node_id = super::new_node_id(); @@ -900,6 +1411,7 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_bootstrap() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; @@ -932,7 +1444,6 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_bootstrap() { .unwrap(); assert!(peer_status.request_bootstrap); - assert!(!peer_status.push_bootstrap); assert_eq!(known_peer.as_deref(), Some("iuna.jhx.app:9444")); let listed = network.inner.peers.lock().await.list(); assert_eq!(listed.len(), 1); @@ -963,7 +1474,7 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_bootstrap() { } #[tokio::test] -async fn real_node_accepts_setup_placeholder_peer_and_pushes_bootstrap() { +async fn real_node_accepts_setup_placeholder_peer_without_requesting_its_chain() { let wallet = Wallet::from_seed("setup-placeholder-peer-real-node"); let node = Arc::new(tokio::sync::Mutex::new(node( "real", @@ -982,6 +1493,7 @@ async fn real_node_accepts_setup_placeholder_peer_and_pushes_bootstrap() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let setup_ledger = Ledger::new(BTreeMap::new(), 1); @@ -1006,19 +1518,13 @@ async fn real_node_accepts_setup_placeholder_peer_and_pushes_bootstrap() { .unwrap(); assert!(!peer_status.request_bootstrap); - assert!(peer_status.push_bootstrap); - let payload = super::catchup_payload_for_peer(&node, &peer_status).await; - assert!(matches!( - payload.as_slice(), - [GossipEnvelope::ChainBootstrap(_)] - )); } #[tokio::test] -async fn lagging_peer_receives_block_pages_before_mempool() { - let wallet = Wallet::from_seed("lagging-peer-blocks-before-mempool"); +async fn periodic_broadcast_does_not_push_duplicate_block_pages_to_lagging_peer() { + let wallet = Wallet::from_seed("lagging-peer-pull-only"); let mut local = node( - "lagging-peer-source", + "lagging-peer-pull-only-source", wallet.clone(), allocations(std::slice::from_ref(&wallet), 1_000), ); @@ -1027,16 +1533,22 @@ async fn lagging_peer_receives_block_pages_before_mempool() { local.drain_outbox(); local.mine_one_at(1).unwrap(); local.drain_outbox(); - queue_plaintext_burn(&mut local, &wallet, 1); - assert!(!local.ledger().pending().is_empty()); let local = Arc::new(tokio::sync::Mutex::new(local)); - let peer_status = super::PeerStatus::new(0, genesis_hash); - let payload = super::catchup_payload_for_peer(&local, &peer_status).await; + let payload = super::envelopes_for_peer( + Some(&local), + Some(super::PeerStatus::new(0, genesis_hash)), + &[GossipEnvelope::PeerStatus { + height: 1, + tip_hash: "tip".to_string(), + time_ms: 1, + }], + ) + .await; assert!(matches!( payload.as_slice(), - [GossipEnvelope::Blocks { .. }] + [GossipEnvelope::PeerStatus { .. }] )); } @@ -1058,6 +1570,7 @@ async fn hello_ignores_private_advertised_listen_address() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let status = node.lock().await.ledger().status(); @@ -1104,6 +1617,7 @@ async fn hello_ignores_loopback_alias_for_unspecified_self() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let hello = ProtocolHello { @@ -1160,6 +1674,7 @@ async fn hello_removes_outbound_peer_that_announces_self_address() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let hello = ProtocolHello { @@ -1215,6 +1730,7 @@ async fn hello_removes_outbound_peer_with_same_node_id() { inbound_limiter: Arc::new(StdMutex::new(super::InboundConnectionLimiter::default())), metrics: super::P2pMetricsCounters::default(), sync_progress: StdMutex::new(super::SyncProgressState::default()), + chain_validation: Arc::new(super::ChainValidationCoordinator::default()), }), }; let hello = ProtocolHello { diff --git a/src/adapters/p2p/writer.rs b/src/adapters/p2p/writer.rs @@ -1,17 +1,22 @@ use anyhow::Result; -use tokio::{io::AsyncWriteExt, net::tcp::OwnedWriteHalf}; +use tokio::{io::AsyncWriteExt, net::tcp::OwnedWriteHalf, time::timeout}; use crate::app::GossipEnvelope; -use super::MAX_GOSSIP_LINE_BYTES; +use super::{MAX_GOSSIP_LINE_BYTES, WRITE_TIMEOUT}; pub(super) async fn write_payload( writer: &mut OwnedWriteHalf, payload: &[GossipEnvelope], ) -> Result<()> { - for envelope in payload { - write_envelope(writer, envelope).await?; - } + timeout(WRITE_TIMEOUT, async { + for envelope in payload { + write_envelope_inner(writer, envelope).await?; + } + Ok::<(), anyhow::Error>(()) + }) + .await + .map_err(|_| anyhow::anyhow!("p2p batch write timed out after {WRITE_TIMEOUT:?}"))??; Ok(()) } @@ -19,6 +24,16 @@ pub(super) async fn write_envelope( writer: &mut OwnedWriteHalf, envelope: &GossipEnvelope, ) -> Result<()> { + timeout(WRITE_TIMEOUT, write_envelope_inner(writer, envelope)) + .await + .map_err(|_| anyhow::anyhow!("p2p write timed out after {WRITE_TIMEOUT:?}"))??; + Ok(()) +} + +async fn write_envelope_inner( + writer: &mut OwnedWriteHalf, + envelope: &GossipEnvelope, +) -> Result<()> { let line = serde_json::to_string(envelope)?; if line.len() > MAX_GOSSIP_LINE_BYTES { anyhow::bail!( diff --git a/src/app/gossip.rs b/src/app/gossip.rs @@ -82,7 +82,7 @@ impl NodeCore { .collect() } - pub fn missing_inventory_requests(&self, blocks: &[BlockInventory]) -> Vec<GossipEnvelope> { + pub fn missing_inventory_request(&self, blocks: &[BlockInventory]) -> Option<GossipEnvelope> { let local_height = self.ledger.height(); let first_height_gap = blocks .iter() @@ -97,19 +97,16 @@ impl NodeCore { .map(|block| block.hash.clone()) .collect::<Vec<_>>(); - let mut requests = Vec::new(); - if !missing_blocks.is_empty() { - requests.push(GossipEnvelope::BlockRequest { - hashes: missing_blocks, - }); - } if first_height_gap.is_some() { - requests.push(GossipEnvelope::BlockRangeRequest { + return Some(GossipEnvelope::BlockRangeRequest { from_height: local_height + 1, limit: BLOCK_REQUEST_LIMIT, }); } - requests + + (!missing_blocks.is_empty()).then_some(GossipEnvelope::BlockRequest { + hashes: missing_blocks, + }) } } @@ -203,12 +200,11 @@ mod tests { hash: block.hash.clone(), }]; - let requests = remote.missing_inventory_requests(&inventory); - assert_eq!(requests.len(), 1); - assert!(matches!(requests[0], GossipEnvelope::BlockRequest { .. })); + let request = remote.missing_inventory_request(&inventory); + assert!(matches!(request, Some(GossipEnvelope::BlockRequest { .. }))); remote.receive(GossipEnvelope::Block(block)).unwrap(); - assert!(remote.missing_inventory_requests(&inventory).is_empty()); + assert!(remote.missing_inventory_request(&inventory).is_none()); } #[test] @@ -226,21 +222,52 @@ mod tests { } let latest = latest.unwrap(); - let requests = remote.missing_inventory_requests(&[BlockInventory { + let request = remote.missing_inventory_request(&[BlockInventory { height: latest.height, hash: latest.hash, }]); - assert_eq!(requests.len(), 1); - match &requests[0] { - GossipEnvelope::BlockRangeRequest { from_height, limit } => { - assert_eq!(*from_height, 1); - assert_eq!(*limit, BLOCK_REQUEST_LIMIT); + match request { + Some(GossipEnvelope::BlockRangeRequest { from_height, limit }) => { + assert_eq!(from_height, 1); + assert_eq!(limit, BLOCK_REQUEST_LIMIT); } other => panic!("expected block range request, got {other:?}"), } } + #[test] + fn multi_block_inventory_starts_exactly_one_range_request() { + let alice = Wallet::from_seed("multi-inventory-alice"); + let bob = Wallet::from_seed("multi-inventory-bob"); + let allocations = allocations(&[alice.clone(), bob.clone()], 1_000); + let mut source = node(alice.clone(), allocations.clone()); + let receiver = node(bob, allocations); + for timestamp_ms in [1, 2] { + queue_plaintext_burn(&mut source, &alice, 1); + source.mine_one_at(timestamp_ms).unwrap(); + } + let inventory = source + .ledger() + .blocks_from(1, 2) + .into_iter() + .map(|block| BlockInventory { + height: block.height, + hash: block.hash, + }) + .collect::<Vec<_>>(); + + let request = receiver.missing_inventory_request(&inventory); + + assert!(matches!( + request, + Some(GossipEnvelope::BlockRangeRequest { + from_height: 1, + limit: BLOCK_REQUEST_LIMIT + }) + )); + } + fn node(wallet: Wallet, allocations: BTreeMap<String, Amount>) -> NodeCore { let ledger = Ledger::new_with_genesis_burns( allocations, diff --git a/src/domain/ledger_queries.rs b/src/domain/ledger_queries.rs @@ -431,14 +431,27 @@ impl Ledger { if limit == 0 { return Vec::new(); } + let Ok(start) = usize::try_from(from_height) else { + return Vec::new(); + }; self.chain + .get(start..) + .unwrap_or_default() .iter() - .filter(|block| block.height >= from_height) .take(limit) .cloned() .collect() } + pub(crate) fn contains_block_sequence(&self, blocks: &[Block]) -> bool { + blocks.iter().all(|block| { + usize::try_from(block.height) + .ok() + .and_then(|height| self.chain.get(height)) + .is_some_and(|known| known.hash == block.hash) + }) + } + pub(crate) fn block_locator(&self) -> Vec<String> { let mut locator = Vec::new(); let mut index = self.chain.len().saturating_sub(1); diff --git a/src/main.rs b/src/main.rs @@ -44,6 +44,7 @@ const VDF_MEASUREMENT_INITIAL_ROUNDS: u64 = 1_000; const VDF_MEASUREMENT_MAX_ROUNDS: u64 = 10_000_000; const VDF_MEASUREMENT_MIN_ELAPSED: Duration = Duration::from_millis(150); const VDF_PROGRESS_LOG_INTERVAL: Duration = Duration::from_secs(10); +const SYNC_CHAIN_CHECKPOINT_INTERVAL: Duration = Duration::from_secs(30); const AUTOMATIC_BURN_ENABLED_ENV: &str = "IUNA_AUTOMATIC_BURN_ENABLED"; const POW_MINING_ENABLED_ENV: &str = "IUNA_POW_MINING_ENABLED"; const POW_MINING_WORKERS_ENV: &str = "IUNA_POW_MINING_WORKERS"; @@ -227,6 +228,7 @@ async fn main() -> Result<()> { let persistence_store = chain_store.clone(); let persistence_ui_data_store = ui_data_store.clone(); let persistence_config = Arc::clone(&ui_config); + let persistence_gossip = gossip.clone(); let persistence_initial_tip = { let node = node.lock().await; if node.has_real_chain() { @@ -242,6 +244,7 @@ async fn main() -> Result<()> { persistence_store, persistence_ui_data_store, persistence_config, + persistence_gossip, persistence_initial_tip, persistence_initial_keep_metrics, ) @@ -971,21 +974,29 @@ async fn run_chain_persistence( store: SqliteChainStore, ui_data_store: SqliteUiDataStore, ui_config: Arc<Mutex<config_store::UiConfig>>, + gossip: p2p::GossipNetwork, initial_saved_tip: Option<String>, initial_projected_keep_metrics: bool, ) { - run_chain_persistence_with_interval( + let initial_state = ChainPersistenceState { + projected_tip: initial_saved_tip.clone(), + saved_tip: initial_saved_tip, + projected_keep_metrics: initial_projected_keep_metrics, + sync_checkpoint_interval: SYNC_CHAIN_CHECKPOINT_INTERVAL, + }; + run_chain_persistence_loop( node, store, ui_data_store, ui_config, Duration::from_secs(2), - initial_saved_tip, - initial_projected_keep_metrics, + Some(gossip), + initial_state, ) .await; } +#[cfg(test)] async fn run_chain_persistence_with_interval( node: SharedNode, store: SqliteChainStore, @@ -995,15 +1006,66 @@ async fn run_chain_persistence_with_interval( initial_saved_tip: Option<String>, initial_projected_keep_metrics: bool, ) { - let mut last_saved_tip = initial_saved_tip; - let mut last_projected_keep_metrics = initial_projected_keep_metrics; + let initial_state = ChainPersistenceState { + projected_tip: initial_saved_tip.clone(), + saved_tip: initial_saved_tip, + projected_keep_metrics: initial_projected_keep_metrics, + sync_checkpoint_interval: SYNC_CHAIN_CHECKPOINT_INTERVAL, + }; + run_chain_persistence_loop( + node, + store, + ui_data_store, + ui_config, + interval, + None, + initial_state, + ) + .await; +} + +struct ChainPersistenceState { + saved_tip: Option<String>, + projected_tip: Option<String>, + projected_keep_metrics: bool, + sync_checkpoint_interval: Duration, +} + +async fn run_chain_persistence_loop( + node: SharedNode, + store: SqliteChainStore, + ui_data_store: SqliteUiDataStore, + ui_config: Arc<Mutex<config_store::UiConfig>>, + interval: Duration, + gossip: Option<p2p::GossipNetwork>, + initial_state: ChainPersistenceState, +) { + let mut last_saved_tip = initial_state.saved_tip; + let mut last_projected_tip = initial_state.projected_tip; + let mut last_projected_keep_metrics = initial_state.projected_keep_metrics; + let mut last_chain_checkpoint = Instant::now(); loop { tokio::time::sleep(interval).await; + let syncing = gossip + .as_ref() + .is_some_and(|network| network.chain_sync_active_or_recent(Duration::from_secs(5))); + let defer_sync_checkpoint = should_defer_sync_checkpoint( + syncing, + last_chain_checkpoint.elapsed(), + initial_state.sync_checkpoint_interval, + ); let snapshot = { let node = node.lock().await; if !node.has_real_chain() { continue; } + let tip_hash = node.ledger().tip_hash(); + if syncing { + let tip_changed = last_saved_tip.as_deref() != Some(tip_hash); + if !tip_changed || defer_sync_checkpoint { + continue; + } + } node.chain_snapshot() }; let Some(tip_hash) = snapshot.blocks.last().map(|block| block.hash.clone()) else { @@ -1011,20 +1073,29 @@ async fn run_chain_persistence_with_interval( }; let keep_metrics = ui_config.lock().await.keep_track_of_metrics; let tip_changed = last_saved_tip.as_deref() != Some(tip_hash.as_str()); + let projected_tip_changed = last_projected_tip.as_deref() != Some(tip_hash.as_str()); let metrics_mode_changed = last_projected_keep_metrics != keep_metrics; - if !tip_changed && !metrics_mode_changed { + if !tip_changed && !projected_tip_changed && !metrics_mode_changed { continue; } - let result = if tip_changed { + let result = if syncing && tip_changed { + persist_chain_snapshot(&store, snapshot).await + } else if tip_changed { persist_chain_and_project_ui_data(&store, &ui_data_store, snapshot, keep_metrics).await } else { project_ui_data_store(&ui_data_store, snapshot, keep_metrics).await }; match result { Ok(()) => { + if !syncing { + last_projected_tip = Some(tip_hash.clone()); + last_projected_keep_metrics = keep_metrics; + } last_saved_tip = Some(tip_hash); - last_projected_keep_metrics = keep_metrics; + if tip_changed { + last_chain_checkpoint = Instant::now(); + } } Err(error) if debug_logging_enabled() => { eprintln!("chain persistence failed: {error:#}") @@ -1034,6 +1105,14 @@ async fn run_chain_persistence_with_interval( } } +fn should_defer_sync_checkpoint( + syncing: bool, + since_last_checkpoint: Duration, + checkpoint_interval: Duration, +) -> bool { + syncing && since_last_checkpoint < checkpoint_interval +} + async fn persist_chain_and_project_ui_data( store: &SqliteChainStore, ui_data_store: &SqliteUiDataStore, diff --git a/src/main_tests.rs b/src/main_tests.rs @@ -24,7 +24,7 @@ use super::{ initial_burn_per_block, initialize_ledger, load_startup_wallet, measure_vdf_rounds, parse_startup_bool_env_value, parse_startup_pow_mining_workers_env_value, persist_chain_snapshot, project_ui_data_store, run_chain_persistence_with_interval, - should_log_automatic_finalization_skip, validate_wallet_for_mode, + should_defer_sync_checkpoint, should_log_automatic_finalization_skip, validate_wallet_for_mode, }; fn parse(args: &[&str]) -> anyhow::Result<Option<CliOptions>> { @@ -45,6 +45,23 @@ fn help_mentions_dev_seed_verify_bypass_env() { } #[test] +fn active_sync_cannot_defer_chain_checkpoint_past_the_interval() { + let interval = Duration::from_secs(30); + + assert!(should_defer_sync_checkpoint( + true, + Duration::from_secs(29), + interval + )); + assert!(!should_defer_sync_checkpoint(true, interval, interval)); + assert!(!should_defer_sync_checkpoint( + false, + Duration::ZERO, + interval + )); +} + +#[test] fn local_testnet_compose_uses_one_obvious_test_password() { let compose = include_str!("../docker-compose.yml"); assert_eq!(