iuna

iuna

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

commit 2c76d938262bd4341cb08144ab29a9d898664e23
parent a7756116ffbcdef8500eb6b404403bc56182607f
Author: Joris Hartog <jorishartog@hotmail.com>
Date:   Sun, 16 Aug 2026 20:50:07 +0200

Harden peer discovery address handling

Diffstat:
MTHANKS.md | 2+-
Msrc/adapters/p2p.rs | 2++
Msrc/adapters/p2p/handshake.rs | 4++--
Msrc/adapters/p2p/network.rs | 22+++++++++++++---------
Msrc/adapters/p2p/sync.rs | 48+++++++++++++++++++++++++++++++++++++++++-------
Msrc/adapters/p2p/tests.rs | 2+-
Msrc/app/peer_book.rs | 221++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Mtests/iuna.rs | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
8 files changed, 351 insertions(+), 41 deletions(-)

diff --git a/THANKS.md b/THANKS.md @@ -5,6 +5,6 @@ iuna is better because people run it, break it, test it, and tell us what is con Special thanks to: - [IkkeMcwood](https://github.com/IkkeMcwood) - thorough UI testing, good feedback, and thoughtful protocol/product input. -- [radbnl](https://github.com/radbnl) - running a public testnet node, good feedback, and thoughtful protocol/product input. +- [radbnl](https://github.com/radbnl) - running a public testnet node, security testing, and thoughtful protocol/product input. - Robin - Valuable UX/UI and protocol testing - Pim - Windows testing diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs @@ -74,6 +74,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 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 SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2); diff --git a/src/adapters/p2p/handshake.rs b/src/adapters/p2p/handshake.rs @@ -485,7 +485,7 @@ mod tests { assert_eq!(listed.len(), 1); let peer = &listed[0]; assert_eq!(peer.address, "142.132.164.59:9444"); - assert_eq!(peer.direction, PeerDirection::Discovered); + assert_eq!(peer.direction, PeerDirection::Outbound); assert_eq!(peer.last_known_height, Some(7)); assert_eq!(peer.messages_received, 0); @@ -501,7 +501,7 @@ mod tests { assert!(repeated); assert_eq!( peers.lock().await.list()[0].direction, - PeerDirection::Discovered + PeerDirection::Outbound ); } diff --git a/src/adapters/p2p/network.rs b/src/adapters/p2p/network.rs @@ -13,7 +13,8 @@ use crate::app::{ use super::{ GossipNetwork, GossipNetworkInner, InboundConnectionLimiter, InboundSessionPermit, - InboundSessionRejection, OutboundBatch, P2pMetrics, P2pMetricsCounters, PEER_QUEUE_SIZE, + 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, }; @@ -220,17 +221,20 @@ impl GossipNetwork { } pub(super) async fn ensure_outbound_sessions(&self) { - self.inner - .peers - .lock() - .await - .prune_stale_inbound_peers_at(crate::app::now_ms(), STALE_INBOUND_PEER_RETENTION_MS); + self.inner.peers.lock().await.prune_stale_peers_at( + crate::app::now_ms(), + STALE_INBOUND_PEER_RETENTION_MS, + STALE_DISCOVERED_PEER_RETENTION_MS, + ); let addresses = self .inner .peers .lock() .await - .connectable_addresses_at(crate::app::now_ms()); + .outbound_session_candidates_at( + crate::app::now_ms(), + MAX_DISCOVERED_OUTBOUND_DIALS_PER_CYCLE, + ); let address_set = addresses.iter().cloned().collect::<BTreeSet<_>>(); let self_filter_addr = self.self_filter_addr().await; let mut sessions = self.inner.sessions.lock().await; @@ -316,7 +320,7 @@ mod tests { } #[tokio::test] - async fn peer_exchange_advertises_discovered_listening_peers() { + async fn peer_exchange_does_not_advertise_discovered_listening_peers() { let alice = Wallet::from_seed("px-discovered-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); @@ -329,7 +333,7 @@ mod tests { match network.peer_exchange().await { GossipEnvelope::PeerList { peers } => { - assert!(peers.contains(&"127.0.0.1:9546".to_string())); + assert!(!peers.contains(&"127.0.0.1:9546".to_string())); } other => panic!("expected peer list, got {other:?}"), } diff --git a/src/adapters/p2p/sync.rs b/src/adapters/p2p/sync.rs @@ -117,7 +117,7 @@ pub(super) async fn apply_peer_list( if is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr) { P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips); } else if peer_list_address_is_discoverable(&peer, remote_addr)? { - peerbook.add_peer(peer); + peerbook.add_discovered_peer(peer); } else { P2pMetricsCounters::inc(&network.inner.metrics.self_peer_skips); } @@ -339,7 +339,7 @@ mod tests { } #[tokio::test] - async fn peer_list_adds_stable_outbound_peers() { + async fn peer_list_adds_discovered_peers() { let alice = Wallet::from_seed("px-recv-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations))); @@ -358,9 +358,19 @@ mod tests { .await .unwrap(); - let addresses = peers.lock().await.addresses(); + let addresses = peers + .lock() + .await + .list() + .into_iter() + .map(|peer| peer.address) + .collect::<Vec<_>>(); assert!(!addresses.contains(&"127.0.0.1:9544".to_string())); assert!(addresses.contains(&"127.0.0.1:9546".to_string())); + assert_eq!( + peers.lock().await.direction_for_tests("127.0.0.1:9546"), + Some(crate::app::PeerDirection::Discovered) + ); } #[tokio::test] @@ -386,7 +396,13 @@ mod tests { .await .unwrap(); - let addresses = peers.lock().await.addresses(); + let addresses = peers + .lock() + .await + .list() + .into_iter() + .map(|peer| peer.address) + .collect::<Vec<_>>(); assert!(!addresses.contains(&"iuna.jhx.app:9444".to_string())); assert!(addresses.contains(&"127.0.0.1:9546".to_string())); } @@ -412,7 +428,13 @@ mod tests { .await .unwrap(); - let addresses = peers.lock().await.addresses(); + let addresses = peers + .lock() + .await + .list() + .into_iter() + .map(|peer| peer.address) + .collect::<Vec<_>>(); assert!(!addresses.contains(&"8.8.8.8:9444".to_string())); assert!(addresses.contains(&"8.8.4.4:9445".to_string())); network.set_accept_inbound(false).await.unwrap(); @@ -442,7 +464,13 @@ mod tests { .await .unwrap(); - let addresses = peers.lock().await.addresses(); + let addresses = peers + .lock() + .await + .list() + .into_iter() + .map(|peer| peer.address) + .collect::<Vec<_>>(); assert!(!addresses.contains(&"10.42.1.1:10091".to_string())); assert!(addresses.contains(&"142.132.164.59:9444".to_string())); } @@ -468,7 +496,13 @@ mod tests { .await .unwrap(); - let addresses = peers.lock().await.addresses(); + let addresses = peers + .lock() + .await + .list() + .into_iter() + .map(|peer| peer.address) + .collect::<Vec<_>>(); assert!(!addresses.contains(&"127.0.0.1:9545".to_string())); assert!(addresses.contains(&"127.0.0.1:9546".to_string())); assert_eq!(network.metrics().self_peer_skips, 1); diff --git a/src/adapters/p2p/tests.rs b/src/adapters/p2p/tests.rs @@ -414,7 +414,7 @@ async fn hello_remembers_advertised_address_after_signed_session_and_dialback() .iter() .find(|peer| peer.address == remote_addr.to_string()) .unwrap(); - assert_eq!(peer.direction, PeerDirection::Discovered); + assert_eq!(peer.direction, PeerDirection::Outbound); assert!( peers .lock() diff --git a/src/app/peer_book.rs b/src/app/peer_book.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeMap; +use std::{collections::BTreeMap, net::SocketAddr}; use serde::{Deserialize, Serialize}; @@ -7,6 +7,11 @@ use super::{ PEER_MISBEHAVIOR_BAN_SCORE, now_ms, }; +pub const MAX_DISCOVERED_PEERS: usize = 256; +pub const MAX_DISCOVERED_PEERS_PER_IP: usize = 4; +pub const MAX_DISCOVERED_PEERS_PER_IPV4_PREFIX: usize = 16; +pub const MAX_DISCOVERED_PEERS_PER_IPV6_PREFIX: usize = 16; + #[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)] pub struct PeerBook { peers: BTreeMap<String, PeerInfo>, @@ -32,15 +37,136 @@ impl PeerBook { } } - pub fn add_discovered_peer(&mut self, address: impl Into<String>) { + pub fn add_discovered_peer(&mut self, address: impl Into<String>) -> bool { let address = address.into(); - let peer = self + if let Some(direction) = self.peers.get(&address).map(|peer| peer.direction.clone()) { + if direction == PeerDirection::Inbound && !self.discovered_peer_has_room(&address) { + return false; + } + if let Some(peer) = self.peers.get_mut(&address) { + if peer.direction == PeerDirection::Inbound { + peer.direction = PeerDirection::Discovered; + } + peer.last_contact_ms = peer.last_contact_ms.or_else(|| Some(now_ms())); + } + return true; + } + if !self.discovered_peer_has_room(&address) { + return false; + } + let mut peer = PeerInfo::new(address.clone(), PeerDirection::Discovered); + peer.last_contact_ms = Some(now_ms()); + self.peers.insert(address, peer); + true + } + + fn discovered_peer_has_room(&self, address: &str) -> bool { + if self.discovered_peer_count() >= MAX_DISCOVERED_PEERS { + return false; + } + let Ok(candidate) = address.parse::<SocketAddr>() else { + return false; + }; + let candidate_ip = candidate.ip(); + let mut same_ip = 0usize; + let mut same_group = 0usize; + for peer in self .peers - .entry(address.clone()) - .or_insert_with(|| PeerInfo::new(address, PeerDirection::Discovered)); - if peer.direction == PeerDirection::Inbound { - peer.direction = PeerDirection::Discovered; + .values() + .filter(|peer| peer.direction == PeerDirection::Discovered) + { + let Ok(existing) = peer.address.parse::<SocketAddr>() else { + continue; + }; + let existing_ip = existing.ip(); + if existing_ip == candidate_ip { + same_ip += 1; + } + if same_discovery_group(existing, candidate) { + same_group += 1; + } + } + if same_ip >= MAX_DISCOVERED_PEERS_PER_IP { + return false; + } + match candidate { + SocketAddr::V4(_) => same_group < MAX_DISCOVERED_PEERS_PER_IPV4_PREFIX, + SocketAddr::V6(_) => same_group < MAX_DISCOVERED_PEERS_PER_IPV6_PREFIX, + } + } + + fn discovered_peer_count(&self) -> usize { + self.peers + .values() + .filter(|peer| peer.direction == PeerDirection::Discovered) + .count() + } + + pub fn promote_discovered_peer(&mut self, address: &str) { + if let Some(peer) = self.peers.get_mut(address) { + if peer.direction == PeerDirection::Discovered { + peer.direction = PeerDirection::Outbound; + } + } + } + + pub fn discovered_peer_count_for_tests(&self) -> usize { + self.discovered_peer_count() + } + + pub fn discovered_peer_capacity_for_tests(&self) -> usize { + MAX_DISCOVERED_PEERS + } + + pub fn direction_for_tests(&self, address: &str) -> Option<PeerDirection> { + self.peers.get(address).map(|peer| peer.direction.clone()) + } + + pub fn peer_count_for_tests(&self) -> usize { + self.peers.len() + } + + pub fn add_discovered_peer_at(&mut self, address: impl Into<String>, now_ms: u64) -> bool { + let address = address.into(); + let added = self.add_discovered_peer(address.clone()); + if added { + if let Some(peer) = self.peers.get_mut(&address) { + peer.last_contact_ms = Some(now_ms); + } } + added + } + + fn prune_stale_discovered_peer(peer: &PeerInfo, now_ms: u64, max_age_ms: u64) -> bool { + if peer.direction != PeerDirection::Discovered || peer.is_banned_at(now_ms) { + return true; + } + let Some(last_contact) = peer.last_success_ms.or(peer.last_contact_ms) else { + return false; + }; + now_ms.saturating_sub(last_contact) <= max_age_ms + } + + fn prune_stale_inbound_peer(peer: &PeerInfo, now_ms: u64, max_age_ms: u64) -> bool { + if peer.direction != PeerDirection::Inbound || peer.is_banned_at(now_ms) { + return true; + } + peer.last_contact_ms + .is_some_and(|last_contact| now_ms.saturating_sub(last_contact) <= max_age_ms) + } + + pub fn prune_stale_peers_at( + &mut self, + now_ms: u64, + inbound_max_age_ms: u64, + discovered_max_age_ms: u64, + ) -> usize { + let before = self.peers.len(); + self.peers.retain(|_, peer| { + Self::prune_stale_inbound_peer(peer, now_ms, inbound_max_age_ms) + && Self::prune_stale_discovered_peer(peer, now_ms, discovered_max_age_ms) + }); + before.saturating_sub(self.peers.len()) } pub fn observe_inbound_peer(&mut self, address: impl Into<String>) { @@ -126,24 +252,67 @@ impl PeerBook { } pub fn addresses(&self) -> Vec<String> { + self.outbound_addresses_at(now_ms()) + } + + pub fn connectable_addresses_at(&self, now_ms: u64) -> Vec<String> { self.peers .values() .filter(|peer| peer.direction != PeerDirection::Inbound) + .filter(|peer| !peer.is_banned_at(now_ms)) .map(|peer| peer.address.clone()) .collect() } - pub fn connectable_addresses_at(&self, now_ms: u64) -> Vec<String> { + pub fn outbound_addresses_at(&self, now_ms: u64) -> Vec<String> { self.peers .values() - .filter(|peer| peer.direction != PeerDirection::Inbound) + .filter(|peer| peer.direction == PeerDirection::Outbound) .filter(|peer| !peer.is_banned_at(now_ms)) .map(|peer| peer.address.clone()) .collect() } + pub fn outbound_session_candidates_at( + &self, + now_ms: u64, + max_discovered: usize, + ) -> Vec<String> { + let mut outbound = Vec::new(); + let mut discovered = self + .peers + .values() + .filter(|peer| peer.direction == PeerDirection::Discovered) + .filter(|peer| !peer.is_banned_at(now_ms)) + .cloned() + .collect::<Vec<_>>(); + discovered.sort_by(|left, right| { + right + .last_success_ms + .cmp(&left.last_success_ms) + .then_with(|| left.last_error_ms.cmp(&right.last_error_ms)) + .then_with(|| left.address.cmp(&right.address)) + }); + + for peer in self + .peers + .values() + .filter(|peer| peer.direction == PeerDirection::Outbound) + .filter(|peer| !peer.is_banned_at(now_ms)) + { + outbound.push(peer.address.clone()); + } + outbound.extend( + discovered + .into_iter() + .take(max_discovered) + .map(|peer| peer.address), + ); + outbound + } + pub fn addresses_except(&self, excluded: &str) -> Vec<String> { - self.connectable_addresses_at(now_ms()) + self.outbound_addresses_at(now_ms()) .into_iter() .filter(|address| address != excluded) .collect() @@ -154,20 +323,15 @@ impl PeerBook { } pub fn prune_stale_inbound_peers_at(&mut self, now_ms: u64, max_age_ms: u64) -> usize { - let before = self.peers.len(); - self.peers.retain(|_, peer| { - if peer.direction != PeerDirection::Inbound || peer.is_banned_at(now_ms) { - return true; - } - peer.last_contact_ms - .is_some_and(|last_contact| now_ms.saturating_sub(last_contact) <= max_age_ms) - }); - before.saturating_sub(self.peers.len()) + self.prune_stale_peers_at(now_ms, max_age_ms, u64::MAX) } pub fn record_sent(&mut self, address: &str, count: u64) { let now = now_ms(); let peer = self.ensure(address, PeerDirection::Outbound); + if peer.direction == PeerDirection::Discovered { + peer.direction = PeerDirection::Outbound; + } peer.messages_sent += count; peer.last_contact_ms = Some(now); peer.last_success_ms = Some(now); @@ -180,6 +344,9 @@ impl PeerBook { pub fn record_status(&mut self, address: &str, height: u64, tip_hash: String) { let now = now_ms(); let peer = self.ensure(address, PeerDirection::Outbound); + if peer.direction == PeerDirection::Discovered { + peer.direction = PeerDirection::Outbound; + } peer.last_known_height = Some(height); peer.last_known_tip_hash = Some(tip_hash); peer.last_contact_ms = Some(now); @@ -388,6 +555,22 @@ fn median_i64(mut values: Vec<i64>) -> Option<i64> { Some(values[values.len() / 2]) } +fn same_discovery_group(left: SocketAddr, right: SocketAddr) -> bool { + match (left, right) { + (SocketAddr::V4(left), SocketAddr::V4(right)) => { + let left = left.ip().octets(); + let right = right.ip().octets(); + left[0] == right[0] && left[1] == right[1] + } + (SocketAddr::V6(left), SocketAddr::V6(right)) => { + let left = left.ip().segments(); + let right = right.ip().segments(); + left[0] == right[0] && left[1] == right[1] + } + _ => false, + } +} + #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] #[serde(rename_all = "snake_case")] pub enum PeerDirection { diff --git a/tests/iuna.rs b/tests/iuna.rs @@ -2157,7 +2157,7 @@ fn peer_book_tracks_multiple_peers_without_networking() { } #[test] -fn peer_book_reports_only_connectable_peers_as_outbound() { +fn peer_book_reports_only_outbound_peers_as_addresses() { let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]); peers.record_received("127.0.0.1:56666", 1); peers.add_discovered_peer("127.0.0.1:9445"); @@ -2167,13 +2167,100 @@ fn peer_book_reports_only_connectable_peers_as_outbound() { assert!(peers.is_connectable_peer("127.0.0.1:9445")); assert!(!peers.is_connectable_peer("127.0.0.1:56666")); assert!(!peers.is_connectable_peer("127.0.0.1:57777")); - assert!(peers.addresses().contains(&"127.0.0.1:9445".to_string())); + assert_eq!(peers.addresses(), vec!["127.0.0.1:9444"]); assert!(peers.remove_peer("127.0.0.1:9444")); assert!(!peers.is_connectable_peer("127.0.0.1:9444")); } #[test] +fn peer_book_caps_discovered_peer_hints() { + let mut peers = PeerBook::default(); + for index in 0..600 { + let octet2 = (index / 16) % 256; + let octet3 = index % 16; + peers.add_discovered_peer(format!("8.{octet2}.{octet3}.1:9444")); + } + + assert_eq!( + peers.discovered_peer_count_for_tests(), + peers.discovered_peer_capacity_for_tests() + ); +} + +#[test] +fn peer_book_caps_inbound_to_discovered_conversions() { + let mut peers = PeerBook::default(); + for index in 0..600 { + let octet2 = (index / 16) % 256; + let octet3 = index % 16; + peers.add_discovered_peer(format!("8.{octet2}.{octet3}.1:9444")); + } + peers.observe_inbound_peer("9.9.9.9:9444"); + + assert!(!peers.add_discovered_peer("9.9.9.9:9444")); + assert_eq!( + peers.discovered_peer_count_for_tests(), + peers.discovered_peer_capacity_for_tests() + ); + assert_eq!( + peers.direction_for_tests("9.9.9.9:9444"), + Some(PeerDirection::Inbound) + ); +} + +#[test] +fn peer_book_limits_discovered_peers_per_ip() { + let mut peers = PeerBook::default(); + for port in 20_000..20_010 { + peers.add_discovered_peer(format!("8.8.8.8:{port}")); + } + + assert_eq!(peers.discovered_peer_count_for_tests(), 4); +} + +#[test] +fn peer_book_prunes_stale_discovered_peer_hints() { + let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]); + peers.add_discovered_peer_at("8.8.8.8:9444", 1_000); + peers.add_discovered_peer_at("8.8.4.4:9444", 10_000); + peers.record_status("8.8.4.4:9444", 1, "tip".to_string()); + + assert_eq!(peers.prune_stale_peers_at(3_700_000, 60_000, 3_600_000), 1); + assert!(peers.addresses().contains(&"127.0.0.1:9444".to_string())); + assert!(peers.addresses().contains(&"8.8.4.4:9444".to_string())); + assert!(!peers.addresses().contains(&"8.8.8.8:9444".to_string())); +} + +#[test] +fn peer_book_limits_discovered_outbound_session_candidates() { + let mut peers = PeerBook::from_addresses(vec![ + "127.0.0.1:9444".to_string(), + "127.0.0.1:9445".to_string(), + ]); + for index in 0..50 { + peers.add_discovered_peer(format!("8.{}.{}.1:9444", index / 16, index % 16)); + } + + let candidates = peers.outbound_session_candidates_at(iuna::app::now_ms(), 8); + + assert_eq!(candidates.len(), 10); + assert!(candidates.contains(&"127.0.0.1:9444".to_string())); + assert!(candidates.contains(&"127.0.0.1:9445".to_string())); +} + +#[test] +fn peer_book_advertises_only_outbound_peers() { + let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]); + peers.add_discovered_peer("8.8.8.8:9444"); + + assert_eq!( + peers.addresses_except("127.0.0.1:9445"), + vec!["127.0.0.1:9444"] + ); +} + +#[test] fn peer_book_prunes_stale_inbound_observations() { let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]); peers.record_status("127.0.0.1:9444", 1, "tip".to_string());