From 11ec16777c6e35b104d09878faf79e0230d66c26 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sat, 18 Jul 2026 22:06:46 +0000 Subject: [PATCH] node: yield the control machine when iterating pending connections MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `Node::connections()` yielded the pending connection itself, so its consumers reached the handshake-phase fields through that value. The pending connection is being folded into the control machine, and the fields will move off it, so yield the machine (keyed by its link) and let each consumer reach what it needs from there. Membership is unchanged: the iterator still selects machines that carry a pending connection, which is the predicate `connection_count` and the stale-connection sweep already use. Every consumer takes the same value from the same place, one hop further out, and no field changes carrier here. The method is now internal to the crate. It yields the control machine, which is not part of the published surface, and nothing outside the crate iterates pending connections — operator views of the handshake phase go through the query surface instead. --- src/control/queries.rs | 1 + src/node/handlers/handshake.rs | 14 +++++++++----- src/node/lifecycle/mod.rs | 17 +++++++++++------ src/node/mod.rs | 18 ++++++++++++++---- src/node/tests/unit.rs | 2 ++ 5 files changed, 37 insertions(+), 15 deletions(-) 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)