commit 43357e43202dd0bbbc94999b256a9bf7d1c7f292
parent 4867cad30a1b95f5c4daa21ca6e133a5b61d7520
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Sun, 9 Aug 2026 20:04:56 +0200
Prune stale inbound peers
Diffstat:
3 files changed, 92 insertions(+), 5 deletions(-)
diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs
@@ -42,6 +42,7 @@ 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 STALE_INBOUND_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000;
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);
@@ -492,11 +493,6 @@ impl GossipNetwork {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_full);
- self.inner
- .peers
- .lock()
- .await
- .record_error(&peer, "outbound gossip queue is full");
}
Err(mpsc::error::TrySendError::Closed(_)) => {
P2pMetricsCounters::inc(&self.inner.metrics.outbound_queue_closed);
@@ -562,6 +558,11 @@ impl GossipNetwork {
}
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);
let addresses = self
.inner
.peers
@@ -2757,6 +2758,62 @@ mod tests {
}
#[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(
+ "full-outbound-queue",
+ wallet.clone(),
+ allocations(&[wallet], 1_000),
+ )));
+ let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::from_addresses(vec![
+ "127.0.0.1:9444".to_string(),
+ ])));
+ let network = super::GossipNetwork {
+ inner: Arc::new(super::GossipNetworkInner {
+ node,
+ peers: Arc::clone(&peers),
+ listen_addr: "127.0.0.1:9544".parse().unwrap(),
+ p2p_announce_addr: tokio::sync::Mutex::new(None),
+ node_id: super::new_node_id(),
+ accept_task: tokio::sync::Mutex::new(None),
+ sessions: tokio::sync::Mutex::new(BTreeMap::new()),
+ inbound_limiter: Arc::new(
+ StdMutex::new(super::InboundConnectionLimiter::default()),
+ ),
+ metrics: super::P2pMetricsCounters::default(),
+ }),
+ };
+ let (sender, _receiver) = tokio::sync::mpsc::channel(1);
+ sender
+ .try_send(vec![GossipEnvelope::PeerStatus {
+ height: 1,
+ tip_hash: "queued".to_string(),
+ time_ms: 1_000,
+ }])
+ .unwrap();
+ network
+ .inner
+ .sessions
+ .lock()
+ .await
+ .insert("127.0.0.1:9444".to_string(), sender);
+
+ 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);
+ let peer = peers.lock().await.list().pop().unwrap();
+ assert_eq!(peer.last_error, None);
+ assert_eq!(peer.last_error_ms, None);
+ }
+
+ #[tokio::test]
async fn limited_line_reader_keeps_partial_line_after_cancelled_read() {
let (mut writer, reader) = tokio::io::duplex(1024);
let mut reader = super::LimitedLineReader::new(reader);
diff --git a/src/app.rs b/src/app.rs
@@ -2169,6 +2169,18 @@ impl PeerBook {
self.peers.values().cloned().collect()
}
+ 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())
+ }
+
pub fn record_sent(&mut self, address: &str, count: u64) {
let now = now_ms();
let peer = self.ensure(address, PeerDirection::Outbound);
diff --git a/tests/iuna.rs b/tests/iuna.rs
@@ -2171,6 +2171,24 @@ fn peer_book_reports_only_configured_outbound_peers_as_outbound() {
}
#[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());
+ peers.observe_inbound_peer("127.0.0.1:56666");
+ peers.record_received("127.0.0.1:57777", 1);
+
+ assert_eq!(
+ peers.prune_stale_inbound_peers_at(iuna::app::now_ms(), 60_000),
+ 1
+ );
+ let listed = peers.list();
+
+ assert!(listed.iter().any(|peer| peer.address == "127.0.0.1:9444"));
+ assert!(listed.iter().any(|peer| peer.address == "127.0.0.1:57777"));
+ assert!(!listed.iter().any(|peer| peer.address == "127.0.0.1:56666"));
+}
+
+#[test]
fn peer_book_bans_misbehaving_peer_temporarily_and_recovers_on_success() {
let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]);