mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-06 11:38:24 +00:00
node: yield the control machine when iterating pending connections
`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.
This commit is contained in:
@@ -1303,6 +1303,7 @@ pub fn show_connections(node: &Node) -> Value {
|
||||
let now = now_ms();
|
||||
let connections: Vec<Value> = node
|
||||
.connections()
|
||||
.filter_map(|(_, machine)| machine.leg())
|
||||
.map(|conn| {
|
||||
let mut conn_json = json!({
|
||||
"link_id": conn.link_id().as_u64(),
|
||||
|
||||
@@ -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<LinkId> = self
|
||||
.connections()
|
||||
.filter_map(|(_, machine)| machine.leg())
|
||||
.filter(|conn| {
|
||||
conn.expected_identity()
|
||||
.map(|id| *id.node_addr() == peer_node_addr)
|
||||
|
||||
@@ -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<LinkId> = 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<NodeAddr> = self.peers.keys().copied().collect();
|
||||
let connecting: HashSet<NodeAddr> = 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<NodeAddr, usize> = 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)
|
||||
|
||||
+14
-4
@@ -1978,6 +1978,7 @@ impl Node {
|
||||
// --- connections (show_connections) ---
|
||||
let connection_rows: Vec<snap::ConnectionRow> = 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<Item = &PeerConnection> {
|
||||
/// 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<Item = (&LinkId, &PeerMachine)> {
|
||||
self.peer_machines
|
||||
.values()
|
||||
.filter_map(|machine| machine.leg())
|
||||
.iter()
|
||||
.filter(|(_, machine)| machine.leg().is_some())
|
||||
}
|
||||
|
||||
// === Peer Management (Active Phase) ===
|
||||
|
||||
@@ -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::<Vec<_>>();
|
||||
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)
|
||||
|
||||
Reference in New Issue
Block a user