diff --git a/src/control/queries.rs b/src/control/queries.rs index 46227c22..53f276eb 100644 --- a/src/control/queries.rs +++ b/src/control/queries.rs @@ -1303,6 +1303,7 @@ pub fn show_connections(node: &Node) -> Value { let now = now_ms(); let connections: Vec = node .connections() + .filter_map(|(_, machine)| machine.leg()) .map(|conn| { let mut conn_json = json!({ "link_id": conn.link_id().as_u64(), diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index 4bc9ad7b..3b78bf69 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -39,11 +39,14 @@ impl EstablishView for Node { rekey_in_progress: existing.map(|p| p.rekey_in_progress()).unwrap_or(false), existing_msg2: existing.and_then(|p| p.handshake_msg2().map(|m| m.to_vec())), at_max_peers: max_peers > 0 && self.peers.len() >= max_peers, - has_pending_outbound_to_peer: self.connections().any(|conn| { - conn.expected_identity() - .map(|id| id.node_addr() == peer_addr) - .unwrap_or(false) - }), + has_pending_outbound_to_peer: self + .connections() + .filter_map(|(_, machine)| machine.leg()) + .any(|conn| { + conn.expected_identity() + .map(|id| id.node_addr() == peer_addr) + .unwrap_or(false) + }), rekey_enabled: self.config().node.rekey.enabled, our_node_addr: *self.identity().node_addr(), } @@ -1552,6 +1555,7 @@ impl Node { // the 30s handshake timeout. let pending_to_same_peer: Vec = self .connections() + .filter_map(|(_, machine)| machine.leg()) .filter(|conn| { conn.expected_identity() .map(|id| *id.node_addr() == peer_node_addr) diff --git a/src/node/lifecycle/mod.rs b/src/node/lifecycle/mod.rs index 2b06c8ae..22fa2190 100644 --- a/src/node/lifecycle/mod.rs +++ b/src/node/lifecycle/mod.rs @@ -372,11 +372,13 @@ impl Node { } fn is_connecting_to_peer(&self, peer_node_addr: &NodeAddr) -> bool { - self.connections().any(|conn| { - conn.expected_identity() - .map(|id| id.node_addr() == peer_node_addr) - .unwrap_or(false) - }) + self.connections() + .filter_map(|(_, machine)| machine.leg()) + .any(|conn| { + conn.expected_identity() + .map(|id| id.node_addr() == peer_node_addr) + .unwrap_or(false) + }) } fn is_connecting_to_peer_on_path( @@ -928,6 +930,7 @@ impl Node { let now_ms = Self::now_ms(); let stale: Vec = self .connections() + .filter_map(|(_, machine)| machine.leg()) .filter(|conn| { conn.expected_identity() .map(|id| id.node_addr() == &peer_addr) @@ -2720,10 +2723,11 @@ impl Node { let connected: HashSet = self.peers.keys().copied().collect(); let connecting: HashSet = self .connections() + .filter_map(|(_, machine)| machine.leg()) .filter_map(|conn| conn.expected_identity().map(|id| *id.node_addr())) .collect(); let mut in_flight_by_peer: HashMap = HashMap::new(); - for conn in self.connections() { + for conn in self.connections().filter_map(|(_, machine)| machine.leg()) { if let Some(id) = conn.expected_identity() { *in_flight_by_peer.entry(*id.node_addr()).or_default() += 1; } @@ -2793,6 +2797,7 @@ impl Node { let in_flight_for_peer = self .connections() + .filter_map(|(_, machine)| machine.leg()) .filter(|conn| { conn.expected_identity() .map(|identity| identity.node_addr() == peer_node_addr) diff --git a/src/node/mod.rs b/src/node/mod.rs index 9799b02b..bb7d2dc8 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -1978,6 +1978,7 @@ impl Node { // --- connections (show_connections) --- let connection_rows: Vec = self .connections() + .filter_map(|(_, machine)| machine.leg()) .map(|conn| snap::ConnectionRow { link_id: conn.link_id().as_u64(), direction: format!("{}", conn.direction()), @@ -2557,11 +2558,20 @@ impl Node { connection } - /// Iterate over all connections. - pub fn connections(&self) -> impl Iterator { + /// Iterate over the control machines that carry a pending connection. + /// + /// Carrying a pending connection is what makes a machine handshake-phase, + /// so the filter below is the membership rule. It is the same predicate + /// that `connection_count` applies, and the one the stale-connection sweep + /// narrows further. + /// + /// Internal to the crate: this yields the control machine, which is not + /// part of the published surface. Callers outside the crate that need a + /// view of the pending handshakes go through the operator queries. + pub(crate) fn connections(&self) -> impl Iterator { self.peer_machines - .values() - .filter_map(|machine| machine.leg()) + .iter() + .filter(|(_, machine)| machine.leg().is_some()) } // === Peer Management (Active Phase) === diff --git a/src/node/tests/unit.rs b/src/node/tests/unit.rs index 138119ef..da4b34a3 100644 --- a/src/node/tests/unit.rs +++ b/src/node/tests/unit.rs @@ -138,6 +138,7 @@ async fn test_try_peer_addresses_races_all_concrete_udp_candidates() { let mut addrs = node .connections() + .filter_map(|(_, machine)| machine.leg()) .filter_map(|conn| conn.source_addr().and_then(|addr| addr.as_str())) .collect::>(); addrs.sort(); @@ -1342,6 +1343,7 @@ async fn update_peers_races_new_alternative_without_dropping_active_peer() { assert_eq!(node.connection_count(), 1); assert_eq!( node.connections() + .filter_map(|(_, machine)| machine.leg()) .next() .and_then(|conn| conn.source_addr()), Some(&new_addr)