diff --git a/CHANGELOG.md b/CHANGELOG.md index 474c1da..fe9f9cd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- `Node::update_peers` for runtime peer-list refresh, returning an + `UpdatePeersOutcome` summarizing added, removed, and retained peers. + Re-derives active peer connections from a new peer configuration + without dropping links to peers that remain in the set. + `PeerAddress` gains a `seen_at_ms` recency field (with + `with_seen_at_ms`) used to prefer more recently observed addresses. - `pool_inbound` and `pool_outbound` counters on the TCP and Tor transport stats (`TcpStats`, `TorStats`). Per-direction accounting is updated at every pool-insert and receive-loop-exit site, plus on @@ -25,6 +31,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed +- `complete_rekey_msg2` now returns the remote peer's startup epoch + alongside the new Noise session, so the rekey path can detect a peer + restart and clear stale session state. +- Active-peer path selection now sorts address candidates by recency + (`seen_at_ms`), preferring the most recently observed address when + racing concurrent path probes. +- Per-tick work budgets bound the connection churn done in a single + node tick: `MAX_DISCOVERY_CONNECTS_PER_TICK`, + `MAX_RETRY_CONNECTIONS_PER_TICK`, and + `MAX_PARALLEL_PATH_CANDIDATES_PER_PEER`. Work beyond a tick's budget + is deferred to the next tick rather than discarded. - Nostr discovery startup is now non-blocking. `Node::start` no longer waits for relay connect, subscribe, or initial advert publish before returning. A slow or unreachable relay no longer @@ -138,6 +155,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- A stale FSP (session-layer) session is now cleared when a peer + restart is detected during FMP rekey or cross-connection promotion. + Previously the old session could linger after the peer came back + with a new startup epoch, leaving the session-layer map out of sync + with the freshly promoted peer. - TCP and Tor `max_inbound_connections` admission cap is now compared against the per-direction inbound count (`pool_inbound`) rather than the combined pool size. Outbound connect-on-send connections share diff --git a/src/config/peer.rs b/src/config/peer.rs index 51a0903..3247737 100644 --- a/src/config/peer.rs +++ b/src/config/peer.rs @@ -45,8 +45,28 @@ pub struct PeerAddress { /// are tried first. #[serde(default = "default_priority")] pub priority: u8, + + /// Wall-clock observation timestamp (Unix ms) for ranking by recency. + /// + /// `None` means "no freshness signal", typically an operator-edited + /// static config. The dialer sorts candidates by this field descending + /// so freshly observed overlay or runtime hints can be tried before stale + /// static addresses. This field is runtime-only and is ignored when + /// comparing peer-address lists for config changes. + #[serde(default, skip_serializing_if = "Option::is_none", skip_deserializing)] + pub seen_at_ms: Option, } +impl PartialEq for PeerAddress { + fn eq(&self, other: &Self) -> bool { + self.transport == other.transport + && self.addr == other.addr + && self.priority == other.priority + } +} + +impl Eq for PeerAddress {} + fn default_priority() -> u8 { 100 } @@ -62,6 +82,7 @@ impl PeerAddress { transport: transport.into(), addr: addr.into(), priority: default_priority(), + seen_at_ms: None, } } @@ -75,8 +96,16 @@ impl PeerAddress { transport: transport.into(), addr: addr.into(), priority, + seen_at_ms: None, } } + + /// Tag this address with a freshness timestamp for source-agnostic + /// candidate ranking. + pub fn with_seen_at_ms(mut self, seen_at_ms: u64) -> Self { + self.seen_at_ms = Some(seen_at_ms); + self + } } /// Configuration for a known peer. diff --git a/src/lib.rs b/src/lib.rs index d3906e5..54252f3 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -67,4 +67,4 @@ pub use peer::{ }; // Re-export node types -pub use node::{Node, NodeError, NodeState}; +pub use node::{Node, NodeError, NodeState, UpdatePeersOutcome}; diff --git a/src/node/handlers/discovery.rs b/src/node/handlers/discovery.rs index d4ff004..933dda2 100644 --- a/src/node/handlers/discovery.rs +++ b/src/node/handlers/discovery.rs @@ -11,6 +11,8 @@ use crate::transport::{TransportAddr, TransportId}; use crate::{NodeAddr, PeerIdentity}; use tracing::{debug, info, trace, warn}; +const MAX_RECENT_DISCOVERY_REQUESTS: usize = 4096; + impl Node { /// Handle an incoming LookupRequest from a peer. /// @@ -34,6 +36,7 @@ impl Node { }; let now_ms = Self::now_ms(); + self.purge_expired_requests(now_ms); // Dedup: drop if we've already seen this request_id. // Also serves as loop protection — tree routing is loop-free, @@ -48,13 +51,21 @@ impl Node { return; } + if self.recent_requests.len() >= MAX_RECENT_DISCOVERY_REQUESTS { + debug!( + request_id = request.request_id, + from = %self.peer_display_name(from), + recent_requests = self.recent_requests.len(), + max_recent_requests = MAX_RECENT_DISCOVERY_REQUESTS, + "Discovery request dedup cache full, dropping LookupRequest" + ); + return; + } + // Record for reverse-path forwarding and dedup self.recent_requests .insert(request.request_id, RecentRequest::new(*from, now_ms)); - // Lazy purge expired entries - self.purge_expired_requests(now_ms); - // Are we the target? if request.target == *self.node_addr() { self.stats_mut().discovery.req_target_is_us += 1; diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 0743bea..e76ebe5 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -658,8 +658,15 @@ impl Node { // Complete the rekey handshake on the ActivePeer if let Some(peer) = self.peers.get_mut(&peer_node_addr) { match peer.complete_rekey_msg2(noise_msg2) { - Ok(session) => { + Ok((session, remote_epoch)) => { let our_index = peer.rekey_our_index().unwrap_or(header.receiver_idx); + let remote_epoch_changed = matches!( + (peer.remote_epoch(), remote_epoch), + (Some(old), Some(new)) if old != new + ); + if remote_epoch.is_some() { + peer.set_remote_epoch(remote_epoch); + } peer.set_pending_session(session, our_index, header.sender_idx); if let Some(transport_id) = peer.transport_id() { @@ -667,6 +674,19 @@ impl Node { .insert((transport_id, our_index.as_u32()), peer_node_addr); } + if remote_epoch_changed { + if self.sessions.remove(&peer_node_addr).is_some() { + debug!( + peer = %display_name, + "Cleared stale FSP session after peer restart during FMP rekey" + ); + } + info!( + peer = %display_name, + "Peer restart detected during FMP rekey, replacing stale endpoint session" + ); + } + debug!( peer = %display_name, new_our_index = %our_index, @@ -1024,9 +1044,15 @@ impl Node { if let Some(existing_peer) = self.peers.get(&peer_node_addr) { let existing_link_id = existing_peer.link_id(); - // Determine which connection wins - let this_wins = - cross_connection_winner(self.identity.node_addr(), &peer_node_addr, is_outbound); + let remote_epoch_changed = matches!((existing_peer.remote_epoch(), remote_epoch), (Some(old), Some(new)) if old != new); + + // Determine which connection wins. A peer restart (different + // startup epoch) is not a normal cross-connection: the old link + // and FSP sessions are cryptographically stale, so the freshly + // authenticated connection must replace them regardless of the + // tie-breaker direction. + let this_wins = remote_epoch_changed + || cross_connection_winner(self.identity.node_addr(), &peer_node_addr, is_outbound); if this_wins { // This connection wins, replace the existing peer @@ -1052,6 +1078,21 @@ impl Node { let _ = self.index_allocator.free(old_idx); } + if remote_epoch_changed { + if self.sessions.remove(&peer_node_addr).is_some() { + debug!( + peer = %self.peer_display_name(&peer_node_addr), + "Cleared stale FSP session after peer restart during promotion" + ); + } + info!( + peer = %self.peer_display_name(&peer_node_addr), + winner_link = %link_id, + loser_link = %loser_link_id, + "Peer restart detected during promotion, replacing stale active peer" + ); + } + self.seed_path_mtu_for_link_peer(&peer_node_addr, transport_id, ¤t_addr); let mut new_peer = ActivePeer::with_session( diff --git a/src/node/lifecycle.rs b/src/node/lifecycle.rs index 103b87e..4d3c326 100644 --- a/src/node/lifecycle.rs +++ b/src/node/lifecycle.rs @@ -14,14 +14,177 @@ use crate::protocol::{Disconnect, DisconnectReason}; use crate::transport::{Link, LinkDirection, LinkId, TransportAddr, TransportId, packet_channel}; use crate::upper::tun::{TunDevice, TunState, run_tun_reader, shutdown_tun_interface}; use crate::{NodeAddr, PeerIdentity}; -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::thread; use std::time::Duration; use tracing::{debug, info, warn}; const OPEN_DISCOVERY_RETRY_LIFETIME_MULTIPLIER: u64 = 2; +const MAX_PARALLEL_PATH_CANDIDATES_PER_PEER: usize = 4; +const MAX_DISCOVERY_CONNECTS_PER_TICK: usize = 16; impl Node { + /// Replace the runtime peer list. + /// + /// Newly added auto-connect peers are dialed immediately, removed peers + /// are dropped from retry bookkeeping, and existing peers get fresh + /// address hints without tearing down an active link. If an existing peer + /// is already connected and a new concrete candidate appears, FIPS starts + /// an alternate handshake in parallel; promotion switches only after that + /// handshake authenticates. + pub async fn update_peers( + &mut self, + new_peers: Vec, + ) -> Result { + let mut new_by_addr: HashMap = + HashMap::with_capacity(new_peers.len()); + for peer in new_peers { + let identity = + PeerIdentity::from_npub(&peer.npub).map_err(|e| NodeError::InvalidPeerNpub { + npub: peer.npub.clone(), + reason: e.to_string(), + })?; + new_by_addr.insert(*identity.node_addr(), peer); + } + + let current_by_addr: HashMap = self + .config + .peers() + .iter() + .filter_map(|peer| { + PeerIdentity::from_npub(&peer.npub) + .ok() + .map(|identity| (*identity.node_addr(), peer.clone())) + }) + .collect(); + + let new_addrs: HashSet<_> = new_by_addr.keys().copied().collect(); + let current_addrs: HashSet<_> = current_by_addr.keys().copied().collect(); + + let removed: Vec<_> = current_addrs.difference(&new_addrs).copied().collect(); + let added: Vec<_> = new_addrs.difference(¤t_addrs).copied().collect(); + let kept: Vec<_> = new_addrs.intersection(¤t_addrs).copied().collect(); + + let mut outcome = crate::node::UpdatePeersOutcome::default(); + let mut refresh_configs = Vec::new(); + + for node_addr in &removed { + if self.retry_pending.remove(node_addr).is_some() { + debug!( + peer = %self.peer_display_name(node_addr), + "Dropping retry entry for peer removed from runtime peer list" + ); + } + self.peer_aliases.remove(node_addr); + outcome.removed += 1; + } + + for node_addr in &kept { + let new_peer = &new_by_addr[node_addr]; + let current_peer = ¤t_by_addr[node_addr]; + let changed = new_peer.addresses != current_peer.addresses + || new_peer.alias != current_peer.alias + || new_peer.connect_policy != current_peer.connect_policy + || new_peer.auto_reconnect != current_peer.auto_reconnect + || new_peer.via_nostr != current_peer.via_nostr; + + if changed { + outcome.updated += 1; + if let Some(state) = self.retry_pending.get_mut(node_addr) { + state.peer_config = new_peer.clone(); + state.retry_after_ms = Self::now_ms(); + } + if let Some(alias) = new_peer.alias.clone() { + self.peer_aliases.insert(*node_addr, alias); + } + } else { + outcome.unchanged += 1; + } + + if new_peer.is_auto_connect() && (!new_peer.addresses.is_empty() || new_peer.via_nostr) + { + refresh_configs.push(new_peer.clone()); + } + } + + let added_configs: Vec<_> = added + .iter() + .map(|node_addr| new_by_addr[node_addr].clone()) + .collect(); + + self.config.peers = new_by_addr.into_values().collect(); + + for peer_config in added_configs { + outcome.added += 1; + let Ok(identity) = PeerIdentity::from_npub(&peer_config.npub) else { + continue; + }; + let name = peer_config + .alias + .clone() + .unwrap_or_else(|| identity.short_npub()); + self.peer_aliases.insert(*identity.node_addr(), name); + self.register_identity(*identity.node_addr(), identity.pubkey_full()); + + if peer_config.is_auto_connect() + && let Err(err) = self.initiate_peer_connection(&peer_config).await + { + debug!( + npub = %peer_config.npub, + error = %err, + "Failed to initiate connection for newly added runtime peer" + ); + self.schedule_retry(*identity.node_addr(), Self::now_ms()); + } + } + + for peer_config in refresh_configs { + let Ok(identity) = PeerIdentity::from_npub(&peer_config.npub) else { + continue; + }; + let node_addr = *identity.node_addr(); + + if self.peers.contains_key(&node_addr) { + match self + .try_active_peer_alternative_addresses(&peer_config, identity) + .await + { + Ok(true) => debug!( + peer = %self.peer_display_name(&node_addr), + "Started alternate-path handshake for active peer" + ), + Ok(false) => {} + Err(err) => debug!( + npub = %peer_config.npub, + error = %err, + "Active peer alternate-path refresh did not start" + ), + } + } else { + match self.initiate_peer_connection(&peer_config).await { + Ok(()) => { + if let Some(state) = self.retry_pending.get_mut(&node_addr) { + state.peer_config = peer_config; + state.retry_after_ms = Self::now_ms().saturating_add( + self.config.node.rate_limit.handshake_timeout_secs * 1000, + ); + } + } + Err(err) => { + debug!( + npub = %peer_config.npub, + error = %err, + "Refreshed peer addresses did not initiate a direct connection" + ); + self.schedule_retry(node_addr, Self::now_ms()); + } + } + } + } + + Ok(outcome) + } + /// Initiate connections to configured static peers. /// /// For each peer configured with AutoConnect policy, creates a link and @@ -142,6 +305,25 @@ impl Node { }) } + fn is_connecting_to_peer_on_path( + &self, + peer_node_addr: &NodeAddr, + transport_id: TransportId, + remote_addr: &TransportAddr, + ) -> bool { + self.connections.values().any(|conn| { + conn.expected_identity() + .map(|id| id.node_addr() == peer_node_addr) + .unwrap_or(false) + && conn.transport_id() == Some(transport_id) + && conn.source_addr() == Some(remote_addr) + }) || self.pending_connects.iter().any(|pending| { + pending.peer_identity.node_addr() == peer_node_addr + && pending.transport_id == transport_id + && pending.remote_addr == *remote_addr + }) + } + /// Initiate a connection to a peer on a specific transport and address. /// /// For connectionless transports (UDP, Ethernet): allocates a link, starts @@ -340,6 +522,9 @@ impl Node { pub(super) async fn poll_transport_discovery(&mut self) { // Collect discoveries first to avoid borrow conflict with self let mut to_connect = Vec::new(); + let mut queued_per_peer: HashMap = HashMap::new(); + let mut connect_budget = self.discovery_connect_budget(); + let mut skipped_budget = 0usize; for (transport_id, transport) in &self.transports { if !transport.is_operational() { @@ -366,29 +551,80 @@ impl Node { if node_addr == *self.identity.node_addr() { continue; } - // Skip if already connected + + let candidate_transport_id = *transport_id; + let remote_addr = peer.addr; + if self.peers.contains_key(&node_addr) { - continue; - } - // Skip if connection already in progress - let connecting = self.connections.values().any(|c| { - c.expected_identity() - .map(|id| id.node_addr() == &node_addr) - .unwrap_or(false) - }); - if connecting { + let transport_name = transport.transport_type().name; + let candidate = PeerAddress::new(transport_name, remote_addr.to_string()); + if self.active_peer_candidate_is_fresh_enough_to_skip( + &node_addr, + std::slice::from_ref(&candidate), + ) { + continue; + } + if self.is_connecting_to_peer_on_path( + &node_addr, + candidate_transport_id, + &remote_addr, + ) { + continue; + } + let queued_for_peer = queued_per_peer.get(&node_addr).copied().unwrap_or(0); + if connect_budget == 0 + || self + .path_candidate_attempt_budget(&node_addr) + .saturating_sub(queued_for_peer) + == 0 + { + skipped_budget = skipped_budget.saturating_add(1); + continue; + } + to_connect.push((candidate_transport_id, remote_addr, identity, true)); + *queued_per_peer.entry(node_addr).or_default() += 1; + connect_budget = connect_budget.saturating_sub(1); continue; } - to_connect.push((*transport_id, peer.addr, identity)); + if self.is_connecting_to_peer_on_path( + &node_addr, + candidate_transport_id, + &remote_addr, + ) { + continue; + } + let queued_for_peer = queued_per_peer.get(&node_addr).copied().unwrap_or(0); + if connect_budget == 0 + || self + .path_candidate_attempt_budget(&node_addr) + .saturating_sub(queued_for_peer) + == 0 + { + skipped_budget = skipped_budget.saturating_add(1); + continue; + } + + to_connect.push((candidate_transport_id, remote_addr, identity, false)); + *queued_per_peer.entry(node_addr).or_default() += 1; + connect_budget = connect_budget.saturating_sub(1); } } - for (transport_id, remote_addr, identity) in to_connect { + if skipped_budget > 0 { + debug!( + skipped = skipped_budget, + queued = to_connect.len(), + "Transport discovery connect budget exhausted" + ); + } + + for (transport_id, remote_addr, identity, active_refresh) in to_connect { info!( peer = %self.peer_display_name(identity.node_addr()), transport_id = %transport_id, remote_addr = %remote_addr, + active_refresh, "Auto-connecting to discovered peer" ); if let Err(e) = self @@ -1228,8 +1464,10 @@ impl Node { .max() .unwrap_or(100) .saturating_add(1); + let seen_at_ms = Self::now_ms(); for endpoint in endpoints { - let Some(candidate) = Self::overlay_endpoint_to_peer_address(&endpoint, next_priority) + let Some(candidate) = + Self::overlay_endpoint_to_peer_address(&endpoint, next_priority, seen_at_ms) else { continue; }; @@ -1251,17 +1489,27 @@ impl Node { fn overlay_endpoint_to_peer_address( endpoint: &OverlayEndpointAdvert, priority: u8, + seen_at_ms: u64, ) -> Option { let transport = match endpoint.transport { OverlayTransportKind::Udp => "udp", OverlayTransportKind::Tcp => "tcp", OverlayTransportKind::Tor => "tor", }; - Some(PeerAddress::with_priority( - transport, - endpoint.addr.clone(), - priority, - )) + Some( + PeerAddress::with_priority(transport, endpoint.addr.clone(), priority) + .with_seen_at_ms(seen_at_ms), + ) + } + + async fn request_nostr_bootstrap(&self, peer_config: &PeerConfig) -> bool { + let Some(bootstrap) = self.nostr_discovery.clone() else { + debug!(npub = %peer_config.npub, "No Nostr overlay runtime for udp:nat address"); + return false; + }; + bootstrap.request_connect(peer_config.clone()).await; + info!(npub = %peer_config.npub, "Started Nostr UDP NAT traversal attempt"); + true } async fn attempt_peer_address_list( @@ -1271,18 +1519,28 @@ impl Node { allow_bootstrap_nat: bool, addresses: &[PeerAddress], ) -> Result<(), NodeError> { + let peer_node_addr = *peer_identity.node_addr(); + let mut attempted = 0usize; + let max_attempts = self.path_candidate_attempt_budget(&peer_node_addr); + if max_attempts == 0 { + return Err(NodeError::NoTransportForType(format!( + "no outbound slots available for {}", + peer_config.npub + ))); + } + for addr in addresses { + if attempted >= max_attempts { + break; + } if addr.transport == "udp" && addr.addr.eq_ignore_ascii_case("nat") { if !allow_bootstrap_nat { continue; } - let Some(bootstrap) = self.nostr_discovery.clone() else { - debug!(npub = %peer_config.npub, "No Nostr overlay runtime for udp:nat address"); - continue; - }; - bootstrap.request_connect(peer_config.clone()).await; - info!(npub = %peer_config.npub, "Started Nostr UDP NAT traversal attempt"); - return Ok(()); + if self.request_nostr_bootstrap(peer_config).await { + attempted = attempted.saturating_add(1); + } + continue; } let (transport_id, remote_addr) = if addr.transport == "ethernet" { @@ -1334,11 +1592,15 @@ impl Node { (tid, TransportAddr::from_string(&addr.addr)) }; + if self.is_connecting_to_peer_on_path(&peer_node_addr, transport_id, &remote_addr) { + continue; + } + match self .initiate_connection(transport_id, remote_addr, peer_identity) .await { - Ok(()) => return Ok(()), + Ok(()) => attempted = attempted.saturating_add(1), Err(e @ NodeError::AccessDenied(_)) => return Err(e), Err(e) => { debug!( @@ -1351,6 +1613,10 @@ impl Node { } } + if attempted > 0 { + return Ok(()); + } + Err(NodeError::NoTransportForType(format!( "no operational transport for any of {}'s addresses", peer_config.npub @@ -1427,6 +1693,21 @@ impl Node { } if configured_npubs.contains(&npub) { + if let Ok(peer_identity) = PeerIdentity::from_npub(&npub) { + let node_addr = *peer_identity.node_addr(); + if !self.peers.contains_key(&node_addr) + && !self.is_connecting_to_peer(&node_addr) + && let Some(state) = self.retry_pending.get_mut(&node_addr) + && state.retry_after_ms > now_ms + { + state.retry_after_ms = now_ms; + debug!( + peer = %peer_identity.short_npub(), + caller = %caller, + "open-discovery sweep: fresh configured-peer advert expedited retry" + ); + } + } skipped_configured = skipped_configured.saturating_add(1); continue; } @@ -1467,8 +1748,10 @@ impl Node { let mut addresses = Vec::new(); let mut priority = 120u8; + let seen_at_ms = Self::now_ms(); for endpoint in endpoints { - let Some(candidate) = Self::overlay_endpoint_to_peer_address(&endpoint, priority) + let Some(candidate) = + Self::overlay_endpoint_to_peer_address(&endpoint, priority, seen_at_ms) else { continue; }; @@ -1607,6 +1890,61 @@ impl Node { connection_slots.min(peer_slots) } + fn outbound_handshake_slots(&self) -> usize { + let used = self + .connections + .len() + .saturating_add(self.pending_connects.len()); + if self.max_connections == 0 { + usize::MAX + } else { + self.max_connections.saturating_sub(used) + } + } + + fn outbound_link_slots(&self) -> usize { + if self.max_links == 0 { + usize::MAX + } else { + self.max_links.saturating_sub(self.links.len()) + } + } + + fn path_candidate_attempt_budget(&self, peer_node_addr: &NodeAddr) -> usize { + if !self.peers.contains_key(peer_node_addr) + && self.max_peers > 0 + && self.peers.len() >= self.max_peers + { + return 0; + } + + let in_flight_for_peer = self + .connections + .values() + .filter(|conn| { + conn.expected_identity() + .map(|identity| identity.node_addr() == peer_node_addr) + .unwrap_or(false) + }) + .count() + .saturating_add( + self.pending_connects + .iter() + .filter(|pending| pending.peer_identity.node_addr() == peer_node_addr) + .count(), + ); + + self.outbound_handshake_slots() + .min(self.outbound_link_slots()) + .min(MAX_PARALLEL_PATH_CANDIDATES_PER_PEER.saturating_sub(in_flight_for_peer)) + } + + fn discovery_connect_budget(&self) -> usize { + self.outbound_handshake_slots() + .min(self.outbound_link_slots()) + .min(MAX_DISCOVERY_CONNECTS_PER_TICK) + } + fn open_discovery_enqueue_budget(&self, configured_npubs: &HashSet) -> usize { let current_open_discovery_pending = self .retry_pending @@ -1841,47 +2179,159 @@ impl Node { return Ok(()); } - // Static-first dialing: avoid delaying configured address attempts on - // advert fetch/network latency. - let static_addresses = self.static_peer_addresses(peer_config); + let candidates = self.peer_address_candidates(peer_config).await; + + if candidates.is_empty() { + return Err(NodeError::NoTransportForType(format!( + "no addresses known for {}", + peer_config.npub + ))); + } + if self - .attempt_peer_address_list( - peer_config, - peer_identity, - allow_bootstrap_nat, - &static_addresses, - ) + .attempt_peer_address_list(peer_config, peer_identity, allow_bootstrap_nat, &candidates) .await .is_ok() { return Ok(()); } - { - let fallback = self - .nostr_peer_fallback_addresses(peer_config, &static_addresses) - .await; - if !fallback.is_empty() - && self - .attempt_peer_address_list( - peer_config, - peer_identity, - allow_bootstrap_nat, - &fallback, - ) - .await - .is_ok() - { - return Ok(()); - } - } - Err(NodeError::NoTransportForType(format!( "no operational transport for any of {}'s addresses", peer_config.npub ))) } + async fn try_active_peer_alternative_addresses( + &mut self, + peer_config: &PeerConfig, + peer_identity: PeerIdentity, + ) -> Result { + let peer_node_addr = *peer_identity.node_addr(); + let candidates = self.peer_address_candidates(peer_config).await; + + if candidates.is_empty() { + return Err(NodeError::NoTransportForType(format!( + "no addresses known for {}", + peer_config.npub + ))); + } + + let concrete: Vec<_> = candidates + .into_iter() + .filter(|addr| !(addr.transport == "udp" && addr.addr.eq_ignore_ascii_case("nat"))) + .collect(); + let has_alternative = concrete + .iter() + .any(|addr| !self.active_peer_matches_candidate(&peer_node_addr, addr)); + let attempt_candidates: Vec<_> = if has_alternative { + concrete + .into_iter() + .filter(|addr| !self.active_peer_matches_candidate(&peer_node_addr, addr)) + .collect() + } else if self.active_peer_needs_same_path_refresh(&peer_node_addr) { + concrete + } else { + Vec::new() + }; + + if attempt_candidates.is_empty() { + return Ok(false); + } + + self.attempt_peer_address_list(peer_config, peer_identity, false, &attempt_candidates) + .await?; + Ok(true) + } + + async fn peer_address_candidates(&self, peer_config: &PeerConfig) -> Vec { + let static_addresses = self.static_peer_addresses(peer_config); + let overlay_addresses = self + .nostr_peer_fallback_addresses(peer_config, &static_addresses) + .await; + + let mut candidates = Vec::with_capacity(overlay_addresses.len() + static_addresses.len()); + for addr in overlay_addresses.into_iter().chain(static_addresses) { + if !candidates.iter().any(|existing: &PeerAddress| { + existing.transport == addr.transport && existing.addr == addr.addr + }) { + candidates.push(addr); + } + } + + candidates.sort_by(|a, b| match (a.seen_at_ms, b.seen_at_ms) { + (Some(a_ts), Some(b_ts)) => b_ts.cmp(&a_ts), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => std::cmp::Ordering::Equal, + }); + candidates + } + + pub(in crate::node) fn active_peer_candidate_is_fresh_enough_to_skip( + &self, + peer_node_addr: &NodeAddr, + candidates: &[PeerAddress], + ) -> bool { + if !self.active_peer_matches_any_candidate(peer_node_addr, candidates) { + return false; + } + !self.active_peer_needs_same_path_refresh(peer_node_addr) + } + + fn active_peer_needs_same_path_refresh(&self, peer_node_addr: &NodeAddr) -> bool { + let Some(peer) = self.peers.get(peer_node_addr) else { + return false; + }; + let stale_after_ms = self + .config + .node + .heartbeat_interval_secs + .saturating_mul(1000) + .max(1000); + peer.idle_time(Self::now_ms()) > stale_after_ms + } + + fn active_peer_matches_any_candidate( + &self, + peer_node_addr: &NodeAddr, + candidates: &[PeerAddress], + ) -> bool { + candidates + .iter() + .any(|candidate| self.active_peer_matches_candidate(peer_node_addr, candidate)) + } + + fn active_peer_matches_candidate( + &self, + peer_node_addr: &NodeAddr, + candidate: &PeerAddress, + ) -> bool { + let Some(peer) = self.peers.get(peer_node_addr) else { + return false; + }; + let Some(current_addr) = peer.current_addr() else { + return false; + }; + if peer + .transport_id() + .map(|id| self.bootstrap_transports.contains(&id)) + .unwrap_or(false) + { + return false; + } + let current_addr = current_addr.to_string(); + let current_transport = peer + .transport_id() + .and_then(|id| self.transports.get(&id)) + .map(|transport| transport.transport_type().name); + + candidate.addr == current_addr + && current_transport + .map(|transport| transport == candidate.transport) + .unwrap_or(true) + } + // === Control API methods === /// Connect to a peer via the control API. diff --git a/src/node/mod.rs b/src/node/mod.rs index 84d3bef..b66a27f 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -199,6 +199,19 @@ impl fmt::Display for NodeState { } } +/// Reports what changed when replacing the runtime peer list. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct UpdatePeersOutcome { + /// Peers present in the new list but not the previous list. + pub added: usize, + /// Peers removed from the previous list. + pub removed: usize, + /// Existing peers whose configured behavior changed. + pub updated: usize, + /// Existing peers whose comparable config did not change. + pub unchanged: usize, +} + /// Recent request tracking for dedup and reverse-path forwarding. /// /// When a LookupRequest is forwarded through a node, the node stores the diff --git a/src/node/retry.rs b/src/node/retry.rs index 34f7e0e..91f2f08 100644 --- a/src/node/retry.rs +++ b/src/node/retry.rs @@ -11,6 +11,7 @@ use crate::identity::NodeAddr; use tracing::{debug, info, warn}; // MAX_BACKOFF_MS is now derived from config: node.retry.max_backoff_secs * 1000 +const MAX_RETRY_CONNECTIONS_PER_TICK: usize = 16; /// Tracks retry state for a peer across connection attempts. pub struct RetryState { @@ -248,8 +249,17 @@ impl Node { .filter(|(_, state)| now_ms >= state.retry_after_ms) .map(|(addr, _)| *addr) .collect(); + let deferred = due.len().saturating_sub(MAX_RETRY_CONNECTIONS_PER_TICK); + if deferred > 0 { + debug!( + due = due.len(), + processing = MAX_RETRY_CONNECTIONS_PER_TICK, + deferred, + "Retry processing budget exhausted; deferring remaining peers" + ); + } - for node_addr in due { + for node_addr in due.into_iter().take(MAX_RETRY_CONNECTIONS_PER_TICK) { // Peer may have connected inbound while we waited if self.peers.contains_key(&node_addr) { self.retry_pending.remove(&node_addr); diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index ed654a4..acbb403 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -97,6 +97,58 @@ async fn test_nat_bootstrap_failure_falls_back_to_direct_udp_address() { } } +#[tokio::test] +async fn test_try_peer_addresses_races_all_concrete_udp_candidates() { + let peer_identity = Identity::generate(); + let mut node = make_node(); + let (packet_tx, packet_rx) = packet_channel(64); + node.packet_tx = Some(packet_tx.clone()); + node.packet_rx = Some(packet_rx); + + let transport_id = TransportId::new(1); + let mut udp = UdpTransport::new( + transport_id, + Some("main".to_string()), + crate::config::UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + ..Default::default() + }, + packet_tx, + ); + udp.start_async().await.unwrap(); + node.transports + .insert(transport_id, TransportHandle::Udp(udp)); + + let peer_config = crate::config::PeerConfig { + npub: peer_identity.npub(), + alias: None, + addresses: vec![ + crate::config::PeerAddress::with_priority("udp", "127.0.0.1:9", 1), + crate::config::PeerAddress::with_priority("udp", "127.0.0.1:10", 2), + ], + connect_policy: crate::config::ConnectPolicy::AutoConnect, + auto_reconnect: true, + via_nostr: false, + }; + let peer_identity = PeerIdentity::from_npub(&peer_config.npub).unwrap(); + + node.try_peer_addresses(&peer_config, peer_identity, false) + .await + .unwrap(); + + let mut addrs = node + .connections + .values() + .filter_map(|conn| conn.source_addr().and_then(|addr| addr.as_str())) + .collect::>(); + addrs.sort(); + assert_eq!(addrs, vec!["127.0.0.1:10", "127.0.0.1:9"]); + + for transport in node.transports.values_mut() { + transport.stop().await.ok(); + } +} + #[tokio::test] async fn test_node_state_transitions() { let mut node = make_node(); @@ -712,6 +764,48 @@ fn test_schedule_retry_increments() { assert_eq!(state.retry_after_ms, 11_000 + 20_000); } +/// Retry processing is paced so a large due set cannot start every +/// handshake candidate in one maintenance tick. +#[tokio::test] +async fn test_process_pending_retries_is_budgeted_per_tick() { + let mut node = make_node(); + let mut addrs = Vec::new(); + + for _ in 0..20 { + let identity = Identity::generate(); + let npub = identity.npub(); + let peer_identity = PeerIdentity::from_npub(&npub).unwrap(); + let node_addr = *peer_identity.node_addr(); + node.retry_pending.insert( + node_addr, + crate::node::retry::RetryState { + peer_config: crate::config::PeerConfig::new(npub, "udp", "10.0.0.2:2121"), + retry_count: 0, + retry_after_ms: 0, + reconnect: true, + expires_at_ms: None, + }, + ); + addrs.push(node_addr); + } + + node.process_pending_retries(1).await; + + let processed = addrs + .iter() + .filter(|addr| { + node.retry_pending + .get(addr) + .is_some_and(|state| state.retry_count > 0) + }) + .count(); + let deferred = addrs.len().saturating_sub(processed); + + assert_eq!(processed, 16); + assert_eq!(deferred, 4); + assert_eq!(node.retry_pending.len(), 20); +} + /// Test that auto-connect peers retry indefinitely (never exhaust). #[test] fn test_schedule_retry_auto_connect_never_exhausts() { @@ -864,6 +958,130 @@ async fn test_try_peer_addresses_skips_connecting_peer() { ); } +#[test] +fn active_peer_same_path_discovery_skips_fresh_peer() { + let mut node = make_node(); + let peer_full = Identity::generate(); + let peer_identity = PeerIdentity::from_pubkey_full(peer_full.pubkey_full()); + let peer_node_addr = *peer_identity.node_addr(); + let transport_id = TransportId::new(1); + let current_addr = TransportAddr::from_string("127.0.0.1:9"); + let mut active_peer = ActivePeer::new(peer_identity, LinkId::new(7), Node::now_ms()); + active_peer.set_current_addr(transport_id, current_addr.clone()); + node.peers.insert(peer_node_addr, active_peer); + let candidate = crate::config::PeerAddress::new("udp", "127.0.0.1:9"); + + assert!(node.active_peer_candidate_is_fresh_enough_to_skip( + &peer_node_addr, + std::slice::from_ref(&candidate), + )); +} + +#[test] +fn active_peer_same_path_discovery_refreshes_stale_peer() { + let mut node = make_node(); + let peer_full = Identity::generate(); + let peer_identity = PeerIdentity::from_pubkey_full(peer_full.pubkey_full()); + let peer_node_addr = *peer_identity.node_addr(); + let transport_id = TransportId::new(1); + let current_addr = TransportAddr::from_string("127.0.0.1:9"); + let stale_at = Node::now_ms().saturating_sub( + node.config + .node + .heartbeat_interval_secs + .saturating_add(1) + .saturating_mul(1000), + ); + let mut active_peer = ActivePeer::new(peer_identity, LinkId::new(7), stale_at); + active_peer.set_current_addr(transport_id, current_addr.clone()); + node.peers.insert(peer_node_addr, active_peer); + let candidate = crate::config::PeerAddress::new("udp", "127.0.0.1:9"); + + assert!(!node.active_peer_candidate_is_fresh_enough_to_skip( + &peer_node_addr, + std::slice::from_ref(&candidate), + )); +} + +#[tokio::test] +async fn update_peers_races_new_alternative_without_dropping_active_peer() { + let mut node = make_node(); + let (packet_tx, packet_rx) = packet_channel(64); + node.packet_tx = Some(packet_tx.clone()); + node.packet_rx = Some(packet_rx); + + let transport_id = TransportId::new(1); + let mut udp = UdpTransport::new( + transport_id, + Some("main".to_string()), + crate::config::UdpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + ..Default::default() + }, + packet_tx, + ); + udp.start_async().await.unwrap(); + node.transports + .insert(transport_id, TransportHandle::Udp(udp)); + + let peer_full = Identity::generate(); + let peer_identity = PeerIdentity::from_pubkey_full(peer_full.pubkey_full()); + let peer_node_addr = *peer_identity.node_addr(); + let current_addr = TransportAddr::from_string("127.0.0.1:9"); + let new_addr = TransportAddr::from_string("127.0.0.1:10"); + let old_link_id = LinkId::new(7); + let mut active_peer = ActivePeer::new(peer_identity, old_link_id, Node::now_ms()); + active_peer.set_current_addr(transport_id, current_addr.clone()); + node.peers.insert(peer_node_addr, active_peer); + node.links.insert( + old_link_id, + Link::connectionless( + old_link_id, + transport_id, + current_addr.clone(), + LinkDirection::Outbound, + Duration::from_millis(100), + ), + ); + + let old_peer = crate::config::PeerConfig { + npub: peer_full.npub(), + alias: None, + addresses: vec![crate::config::PeerAddress::new("udp", "127.0.0.1:9")], + connect_policy: crate::config::ConnectPolicy::AutoConnect, + auto_reconnect: true, + via_nostr: false, + }; + let new_peer = crate::config::PeerConfig { + addresses: vec![ + crate::config::PeerAddress::new("udp", "127.0.0.1:9"), + crate::config::PeerAddress::new("udp", "127.0.0.1:10"), + ], + ..old_peer.clone() + }; + node.config.peers = vec![old_peer]; + + let outcome = node.update_peers(vec![new_peer]).await.unwrap(); + + assert_eq!(outcome.updated, 1); + assert_eq!(node.peer_count(), 1, "existing link must stay live"); + assert_eq!(node.connection_count(), 1); + assert_eq!( + node.connections + .values() + .next() + .and_then(|conn| conn.source_addr()), + Some(&new_addr) + ); + let active = node.get_peer(&peer_node_addr).unwrap(); + assert_eq!(active.link_id(), old_link_id); + assert_eq!(active.current_addr(), Some(¤t_addr)); + + for transport in node.transports.values_mut() { + transport.stop().await.ok(); + } +} + #[tokio::test] async fn test_nostr_traversal_failure_skips_connected_peer() { let mut node = make_node(); diff --git a/src/peer/active.rs b/src/peer/active.rs index 2af00f0..3fa1535 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -593,6 +593,13 @@ impl ActivePeer { self.remote_epoch } + /// Update the remote peer's startup epoch after a successful in-place + /// rekey. Initial handshakes set this through `with_session`, but recovery + /// rekeys also exchange epochs and must keep restart detection current. + pub(crate) fn set_remote_epoch(&mut self, remote_epoch: Option<[u8; 8]>) { + self.remote_epoch = remote_epoch; + } + // === Tree Accessors === /// Get the peer's tree coordinates, if known. @@ -1099,7 +1106,10 @@ impl ActivePeer { /// Takes the stored handshake state, reads msg2, and returns the /// completed NoiseSession. Clears the handshake-related fields but /// leaves rekey_our_index for set_pending_session to use. - pub fn complete_rekey_msg2(&mut self, msg2_bytes: &[u8]) -> Result { + pub fn complete_rekey_msg2( + &mut self, + msg2_bytes: &[u8], + ) -> Result<(NoiseSession, Option<[u8; 8]>), NoiseError> { let mut hs = self .rekey_handshake .take() @@ -1109,13 +1119,14 @@ impl ActivePeer { })?; hs.read_message_2(msg2_bytes)?; + let remote_epoch = hs.remote_epoch(); let session = hs.into_session()?; // Clear msg1 resend state self.rekey_msg1 = None; self.rekey_msg1_next_resend = 0; - Ok(session) + Ok((session, remote_epoch)) } /// Check if msg1 needs resending.