From 87d1af02698815b05fa73f662396089edabaf05a Mon Sep 17 00:00:00 2001 From: Martti Malmi Date: Fri, 15 May 2026 15:06:54 +0000 Subject: [PATCH] nostr: ignore stale traversal for active peers Skip BootstrapEvent::Established and BootstrapEvent::Failed dispatch in poll_nostr_discovery for peers that are already connected or actively handshaking. Without these guards, stale traversal events arriving after a peer connected through a different path would either attempt to adopt a redundant socket against the live connection (Established) or poison the per-peer failure-state cooldown and trigger redundant retraversal via schedule_retry / try_peer_addresses (Failed). The four guard sites use a new is_connecting_to_peer helper extracted from the existing closure inside initiate_peer_connection; the helper checks for an in-flight outbound handshake state. adopt_established_traversal gains a defense-in-depth check returning PeerAlreadyExists when called against an already-promoted peer, so the invariant holds if a future caller bypasses the outer dispatch guard. Side benefit: narrows a cooldown-poisoning vector previously available to an attacker injecting stale failure events for an active peer. Test coverage for the new behavior: - test_try_peer_addresses_skips_connected_peer - test_try_peer_addresses_skips_connecting_peer - test_nostr_traversal_failure_skips_connected_peer (Failed-arm event injection) - test_nostr_traversal_established_skips_connected_peer (Established-arm event injection, mirror of the Failed test) - test_adopted_traversal_skips_already_connected_peer (adopt_established_traversal defense-in-depth) CHANGELOG entry under [Unreleased] / Fixed. Closes #87 --- CHANGELOG.md | 14 ++++ src/discovery/nostr/runtime.rs | 6 ++ src/node/lifecycle.rs | 90 +++++++++++++++++++--- src/node/tests/bootstrap.rs | 45 +++++++++++ src/node/tests/unit.rs | 131 +++++++++++++++++++++++++++++++++ 5 files changed, 274 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5a6d592..dd53edd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -77,6 +77,20 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 after the test's previous one-shot grep gave up, producing a pre-existing flake on next-branch CI. Success-path cost is unchanged — the helper returns as soon as the pattern appears. +- Nostr-discovered NAT-traversal events (`BootstrapEvent::Established` + and `BootstrapEvent::Failed`) for peers that are already connected + or actively handshaking are now short-circuited at the + `poll_nostr_discovery` dispatch sites before any cooldown + bookkeeping or fallback retry scheduling runs. Stale `Failed` events + previously poisoned the per-peer failure-state cooldown of healthy + peers and could trigger redundant retraversal attempts via + `schedule_retry` / `try_peer_addresses`; stale `Established` + handoffs could attempt to adopt a second socket against a live + connection. A defense-in-depth guard was added to + `adopt_established_traversal` so the same invariant holds if a + future caller bypasses the outer dispatch check. As a side benefit, + narrows a cooldown-poisoning vector previously available to an + attacker injecting stale failure events for an active peer. ## [0.3.0] - 2026-05-11 diff --git a/src/discovery/nostr/runtime.rs b/src/discovery/nostr/runtime.rs index eea46f5..26f3b4b 100644 --- a/src/discovery/nostr/runtime.rs +++ b/src/discovery/nostr/runtime.rs @@ -1537,4 +1537,10 @@ impl NostrDiscovery { let mut cache = self.advert_cache.write().await; cache.insert(npub, advert); } + + /// Queue a bootstrap event directly for lifecycle tests without live relays + /// or a running traversal task. + pub(crate) fn push_event_for_test(&self, event: BootstrapEvent) { + let _ = self.event_tx.send(event); + } } diff --git a/src/node/lifecycle.rs b/src/node/lifecycle.rs index 4e92418..a8ee430 100644 --- a/src/node/lifecycle.rs +++ b/src/node/lifecycle.rs @@ -122,12 +122,7 @@ impl Node { } // Check if connection already in progress to this peer - let already_connecting = self.connections.values().any(|conn| { - conn.expected_identity() - .map(|id| id.node_addr() == &peer_node_addr) - .unwrap_or(false) - }); - if already_connecting { + if self.is_connecting_to_peer(&peer_node_addr) { debug!( npub = %peer_config.npub, "Connection already in progress, skipping" @@ -139,6 +134,14 @@ impl Node { .await } + fn is_connecting_to_peer(&self, peer_node_addr: &NodeAddr) -> bool { + self.connections.values().any(|conn| { + conn.expected_identity() + .map(|id| id.node_addr() == peer_node_addr) + .unwrap_or(false) + }) + } + /// Initiate a connection to a peer on a specific transport and address. /// /// For connectionless transports (UDP, Ethernet): allocates a link, starts @@ -410,6 +413,23 @@ impl Node { match event { BootstrapEvent::Established { traversal } => { let peer_npub = traversal.peer_npub.clone(); + if let Ok(peer_identity) = PeerIdentity::from_npub(&peer_npub) { + let peer_addr = *peer_identity.node_addr(); + if self.peers.contains_key(&peer_addr) { + debug!( + peer_npub = %peer_npub, + "Ignoring established NAT traversal for already-connected peer" + ); + continue; + } + if self.is_connecting_to_peer(&peer_addr) { + debug!( + peer_npub = %peer_npub, + "Ignoring established NAT traversal while peer handshake is already in progress" + ); + continue; + } + } match self.adopt_established_traversal(traversal).await { Ok(_) => { info!(peer_npub = %peer_npub, "Adopted NAT traversal socket"); @@ -426,6 +446,28 @@ impl Node { peer_config, reason, } => { + let peer_identity = match PeerIdentity::from_npub(&peer_config.npub) { + Ok(identity) => identity, + Err(_) => continue, + }; + let node_addr = *peer_identity.node_addr(); + if self.peers.contains_key(&node_addr) { + debug!( + npub = %peer_config.npub, + error = %reason, + "Ignoring failed NAT traversal for already-connected peer" + ); + continue; + } + if self.is_connecting_to_peer(&node_addr) { + debug!( + npub = %peer_config.npub, + error = %reason, + "Ignoring failed NAT traversal while peer handshake is already in progress" + ); + continue; + } + let now_ms = Self::now_ms(); let decision = bootstrap.record_traversal_failure(&peer_config.npub, now_ms); if decision.should_warn { @@ -476,11 +518,6 @@ impl Node { }); } - let peer_identity = match PeerIdentity::from_npub(&peer_config.npub) { - Ok(identity) => identity, - Err(_) => continue, - }; - if self .try_peer_addresses(&peer_config, peer_identity, false) .await @@ -489,7 +526,6 @@ impl Node { continue; } - let node_addr = *peer_identity.node_addr(); self.schedule_retry(node_addr, now_ms); if let Some(cooldown_until_ms) = decision.cooldown_until_ms && let Some(state) = self.retry_pending.get_mut(&node_addr) @@ -1694,6 +1730,22 @@ impl Node { peer_identity: PeerIdentity, allow_bootstrap_nat: bool, ) -> Result<(), NodeError> { + let peer_node_addr = *peer_identity.node_addr(); + if self.peers.contains_key(&peer_node_addr) { + debug!( + npub = %peer_config.npub, + "Peer already exists, skipping address attempts" + ); + return Ok(()); + } + if self.is_connecting_to_peer(&peer_node_addr) { + debug!( + npub = %peer_config.npub, + "Connection already in progress, skipping address attempts" + ); + return Ok(()); + } + // Static-first dialing: avoid delaying configured address attempts on // advert fetch/network latency. let static_addresses = self.static_peer_addresses(peer_config); @@ -1837,6 +1889,20 @@ impl Node { } })?; let peer_node_addr = *peer_identity.node_addr(); + if self.peers.contains_key(&peer_node_addr) { + debug!( + peer_npub = %traversal.peer_npub, + "Ignoring NAT traversal handoff for already-connected peer" + ); + return Err(NodeError::PeerAlreadyExists(peer_node_addr)); + } + if self.is_connecting_to_peer(&peer_node_addr) { + debug!( + peer_npub = %traversal.peer_npub, + "Ignoring NAT traversal handoff while peer handshake is already in progress" + ); + return Err(NodeError::PeerAlreadyExists(peer_node_addr)); + } self.peer_aliases .insert(peer_node_addr, peer_identity.short_npub()); diff --git a/src/node/tests/bootstrap.rs b/src/node/tests/bootstrap.rs index 04906fa..4e59cda 100644 --- a/src/node/tests/bootstrap.rs +++ b/src/node/tests/bootstrap.rs @@ -117,6 +117,51 @@ async fn test_failed_adopted_traversal_cleans_up_transport() { ); } +#[tokio::test] +async fn test_adopted_traversal_skips_already_connected_peer() { + let mut node = make_node(); + let (packet_tx, packet_rx) = packet_channel(64); + node.packet_tx = Some(packet_tx); + node.packet_rx = Some(packet_rx); + node.state = NodeState::Running; + + let transport_id = TransportId::new(1); + let link_id = LinkId::new(1); + let (conn, peer_identity) = make_completed_connection(&mut node, link_id, transport_id, 1_000); + let peer_node_addr = *peer_identity.node_addr(); + node.add_connection(conn).unwrap(); + node.promote_connection(link_id, peer_identity, 2_000) + .unwrap(); + + let link_count = node.link_count(); + let transport_count = node.transport_count(); + + let adopted_socket = std::net::UdpSocket::bind("127.0.0.1:0").unwrap(); + let handoff = EstablishedTraversal::new( + "sess-stale", + peer_identity.npub(), + "127.0.0.1:9".parse().unwrap(), + adopted_socket, + ) + .with_transport_name("nostr-stale"); + + let result = node.adopt_established_traversal(handoff).await; + assert!( + matches!(result, Err(NodeError::PeerAlreadyExists(addr)) if addr == peer_node_addr), + "stale traversal handoff should be ignored once the peer is already active" + ); + assert_eq!( + node.link_count(), + link_count, + "ignored traversal must not create a duplicate link" + ); + assert_eq!( + node.transport_count(), + transport_count, + "ignored traversal must not leak an adopted transport" + ); +} + #[tokio::test] async fn test_third_peer_can_handshake_via_adopted_transport_socket() { let mut node_a = make_node(); // Existing traversal peer (Alice) diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 03f9558..a4a7a2b 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -1,7 +1,9 @@ use super::*; +use crate::discovery::nostr::{BootstrapEvent, NostrDiscovery}; use crate::peer::PromotionResult; use crate::transport::udp::UdpTransport; use crate::transport::{TransportHandle, packet_channel}; +use std::sync::Arc; #[test] fn test_node_creation() { @@ -779,6 +781,135 @@ fn test_schedule_retry_skips_connected_peer() { ); } +#[tokio::test] +async fn test_try_peer_addresses_skips_connected_peer() { + let mut node = make_node(); + let transport_id = TransportId::new(1); + let link_id = LinkId::new(1); + let (conn, peer_identity) = make_completed_connection(&mut node, link_id, transport_id, 1000); + let peer_config = crate::config::PeerConfig::new(peer_identity.npub(), "udp", "127.0.0.1:9"); + + node.add_connection(conn).unwrap(); + node.promote_connection(link_id, peer_identity, 2000) + .unwrap(); + let link_count = node.link_count(); + let connection_count = node.connection_count(); + + node.try_peer_addresses(&peer_config, peer_identity, true) + .await + .unwrap(); + + assert_eq!( + node.link_count(), + link_count, + "stale retry/traversal fallback must not create a duplicate link" + ); + assert_eq!( + node.connection_count(), + connection_count, + "stale retry/traversal fallback must not create a duplicate handshake" + ); +} + +#[tokio::test] +async fn test_try_peer_addresses_skips_connecting_peer() { + let mut node = make_node(); + let peer_identity = make_peer_identity(); + let peer_config = crate::config::PeerConfig::new(peer_identity.npub(), "udp", "127.0.0.1:9"); + let pending = PeerConnection::outbound(LinkId::new(1), peer_identity, 1000); + node.add_connection(pending).unwrap(); + + node.try_peer_addresses(&peer_config, peer_identity, true) + .await + .unwrap(); + + assert_eq!( + node.connection_count(), + 1, + "stale retry/traversal fallback must not start a second handshake" + ); + assert_eq!( + node.link_count(), + 0, + "stale retry/traversal fallback must not allocate a link while a handshake is pending" + ); +} + +#[tokio::test] +async fn test_nostr_traversal_failure_skips_connected_peer() { + let mut node = make_node(); + let transport_id = TransportId::new(1); + let link_id = LinkId::new(1); + let (conn, peer_identity) = make_completed_connection(&mut node, link_id, transport_id, 1000); + node.add_connection(conn).unwrap(); + node.promote_connection(link_id, peer_identity, 2000) + .unwrap(); + + let bootstrap = Arc::new(NostrDiscovery::new_for_test()); + bootstrap.push_event_for_test(BootstrapEvent::Failed { + peer_config: crate::config::PeerConfig::new(peer_identity.npub(), "udp", "127.0.0.1:9"), + reason: "stale traversal failure".to_string(), + }); + node.nostr_discovery = Some(bootstrap.clone()); + + node.poll_nostr_discovery().await; + + assert!( + bootstrap.failure_state_snapshot().is_empty(), + "stale failures for connected peers must not affect traversal cooldown" + ); + assert!( + node.retry_pending.is_empty(), + "stale failures for connected peers must not enqueue reconnect attempts" + ); +} + +#[tokio::test] +async fn test_nostr_traversal_established_skips_connected_peer() { + use crate::discovery::EstablishedTraversal; + use std::net::UdpSocket; + + let mut node = make_node(); + let transport_id = TransportId::new(1); + let link_id = LinkId::new(1); + let (conn, peer_identity) = make_completed_connection(&mut node, link_id, transport_id, 1000); + node.add_connection(conn).unwrap(); + node.promote_connection(link_id, peer_identity, 2000) + .unwrap(); + let link_count = node.link_count(); + let connection_count = node.connection_count(); + + let bootstrap = Arc::new(NostrDiscovery::new_for_test()); + let socket = UdpSocket::bind("127.0.0.1:0").expect("bind local UDP socket"); + let remote_addr = "127.0.0.1:9999".parse().expect("parse remote addr"); + bootstrap.push_event_for_test(BootstrapEvent::Established { + traversal: EstablishedTraversal::new( + "test-session", + peer_identity.npub(), + remote_addr, + socket, + ), + }); + node.nostr_discovery = Some(bootstrap.clone()); + + node.poll_nostr_discovery().await; + + assert_eq!( + node.link_count(), + link_count, + "stale established handoff must not allocate a new link" + ); + assert_eq!( + node.connection_count(), + connection_count, + "stale established handoff must not start a new handshake" + ); + assert!( + node.retry_pending.is_empty(), + "stale established handoff must not enqueue a reconnect" + ); +} + #[tokio::test] async fn test_process_pending_retries_drops_expired_entries() { let mut node = make_node();