From e1ae261eb2cd64a10d1629960635832100ad81c5 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 28 May 2026 21:21:45 +0000 Subject: [PATCH] node: wire the next-side XX handshake and rekey rejection sites Adapt the typed RejectReason coverage to the Noise XX handshake: wire the msg1/msg2/msg3 state-machine rejection sites in handlers/handshake.rs and the rekey-initiator outbound sites in handlers/rekey.rs, and add the XX-specific HandshakeStats counters. The shared RejectReason scaffold and the branch-agnostic clusters come from the merged refactor-hotpath work; this commit carries only the next-side delta. --- CHANGELOG.md | 12 ++-- src/control/snapshots/show_bloom.json | 1 + src/node/handlers/handshake.rs | 92 ++++++++++++++++++++++++++- src/node/handlers/rekey.rs | 11 ++++ src/node/stats.rs | 14 +++- src/node/tests/bloom_poison.rs | 8 +-- 6 files changed, 124 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 03ee4e6..55fbcda 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -110,13 +110,15 @@ with v0.3.x peers. existing `TreeStats`, `BloomStats`, `DiscoveryStats`, and `ForwardingStats`, and `TreeStats::ancestry_invalid` is now incremented from the `TreeAnnounce::validate_semantics` rejection - site that was previously silent. Several handshake, MMP, tree, and + site that was previously silent. The Noise XX handshake cluster + (msg1/msg2/msg3) and the rekey-initiator outbound sites are wired in + addition to the shared clusters. Several handshake, MMP, tree, and discovery rejection paths that had no counter at all are now counted, including the `send_lookup_response` no-route drop - (`DiscoveryStats::resp_no_route`). Existing - direct counters at the bloom / discovery / forwarding sites are - retained alongside the new dispatch while the rollout is in progress; - a later change collapses the duplicate increment. + (`DiscoveryStats::resp_no_route`). Existing direct counters at the + bloom / discovery / forwarding sites are retained alongside the new + dispatch while the rollout is in progress; a later change collapses + the duplicate increment. - `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 diff --git a/src/control/snapshots/show_bloom.json b/src/control/snapshots/show_bloom.json index 8c5f5ff..3e85739 100644 --- a/src/control/snapshots/show_bloom.json +++ b/src/control/snapshots/show_bloom.json @@ -16,6 +16,7 @@ "invalid": 0, "nacks_received": 0, "nacks_sent": 0, + "non_v1": 0, "received": 0, "send_failed": 0, "sent": 0, diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index d195873..3a2d697 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -7,6 +7,7 @@ use crate::PeerIdentity; use crate::node::acl::PeerAclContext; +use crate::node::reject::{HandshakeReject, RejectReason}; use crate::node::wire::{Msg1Header, Msg2Header, Msg3Header, build_msg2, build_msg3}; use crate::node::{Node, NodeError}; use crate::peer::{ActivePeer, PeerConnection, PromotionResult, cross_connection_winner}; @@ -86,6 +87,8 @@ impl Node { // deadlocks when the larger-NodeAddr side has accept_connections=false. if !self.should_admit_msg1(packet.transport_id, &packet.remote_addr) { self.msg1_rate_limiter.complete_handshake(); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } @@ -95,6 +98,8 @@ impl Node { None => { self.msg1_rate_limiter.complete_handshake(); debug!("Invalid msg1 header"); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -137,6 +142,9 @@ impl Node { remote_addr = %packet.remote_addr, "Duplicate msg1 but no stored msg2 to resend" ); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::UnknownConnection, + )); } self.msg1_rate_limiter.complete_handshake(); return; @@ -189,6 +197,8 @@ impl Node { error = %e, "Failed to process msg1" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -209,6 +219,8 @@ impl Node { Err(e) => { self.msg1_rate_limiter.complete_handshake(); warn!(error = %e, "Failed to allocate session index for inbound"); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -259,6 +271,8 @@ impl Node { .remove(&(packet.transport_id, packet.remote_addr)); let _ = self.index_allocator.free(our_index); self.msg1_rate_limiter.complete_handshake(); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -306,6 +320,8 @@ impl Node { Some(h) => h, None => { debug!("Invalid msg2 header"); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -319,6 +335,8 @@ impl Node { receiver_idx = %header.receiver_idx, "No pending outbound handshake for index" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection)); return; } }; @@ -407,6 +425,9 @@ impl Node { } let _ = self.index_allocator.free(idx); } + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::BadState, + )); } } Err(e) => { @@ -421,6 +442,8 @@ impl Node { } let _ = self.index_allocator.free(idx); } + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); } } } @@ -429,8 +452,13 @@ impl Node { return; } - // Not a rekey — stale pending_outbound entry + // Not a rekey — stale pending_outbound entry pointing at a + // removed connection and no rekey-in-progress peer claims the + // receiver_idx. State-machine inconsistency, not a fresh + // lookup miss. self.pending_outbound.remove(&key); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } @@ -438,6 +466,8 @@ impl Node { let Some(conn) = self.connections.get_mut(&link_id) else { warn!(link_id = %link_id, "Connection removed during msg2 processing"); self.pending_outbound.remove(&key); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection)); return; }; @@ -459,6 +489,8 @@ impl Node { "Handshake completion failed" ); conn.mark_failed(); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -470,6 +502,8 @@ impl Node { Err(e) => { warn!(link_id = %link_id, our_profile = %self.node_profile, error = %e, "FMP negotiation failed"); conn.mark_failed(); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -484,6 +518,8 @@ impl Node { Some(id) => *id, None => { warn!(link_id = %link_id, "No identity after handshake"); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -509,12 +545,16 @@ impl Node { self.pending_outbound.remove(&key); self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } if peer_node_addr == *self.identity.node_addr() { debug!(link_id = %link_id, "Discovered self via shared-media beacon, dropping"); self.connections.remove(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } @@ -542,6 +582,8 @@ impl Node { if let Some(conn) = self.connections.get_mut(&link_id) { conn.mark_failed(); } + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -576,6 +618,8 @@ impl Node { Some(c) => c, None => { self.pending_outbound.remove(&key); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection)); return; } }; @@ -594,6 +638,8 @@ impl Node { _ => { warn!(peer = %self.peer_display_name(&peer_node_addr), "Incomplete outbound connection"); self.pending_outbound.remove(&key); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -764,6 +810,8 @@ impl Node { error = %e, "Failed to promote connection" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); } } } @@ -780,6 +828,8 @@ impl Node { Some(h) => h, None => { debug!("Invalid msg3 header"); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -789,7 +839,10 @@ impl Node { let link_id = match self.pending_inbound.remove(&key) { Some(id) => id, None => { - // Check if this is a rekey msg3 for an active peer + // Check if this is a rekey msg3 for an active peer. + // handle_rekey_msg3 records its own UnknownConnection or + // BadState classification depending on whether a matching + // rekey-responder slot is found. self.handle_rekey_msg3(&packet, &header).await; return; } @@ -804,6 +857,8 @@ impl Node { link_id = %link_id, "No pending connection for msg3" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection)); return; } }; @@ -829,6 +884,8 @@ impl Node { if let Some(idx) = our_idx_to_free { let _ = self.index_allocator.free(idx); } + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -841,6 +898,8 @@ impl Node { warn!(link_id = %link_id, our_profile = %self.node_profile, error = %e, "FMP negotiation failed"); self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -853,6 +912,8 @@ impl Node { warn!("Identity not learned from msg3"); self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -904,6 +965,8 @@ impl Node { } self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } @@ -911,6 +974,8 @@ impl Node { debug!(link_id = %link_id, "Received msg3 from self, dropping"); self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } @@ -954,6 +1019,8 @@ impl Node { if let Some(idx) = our_idx_to_free { let _ = self.index_allocator.free(idx); } + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -1020,6 +1087,9 @@ impl Node { None => { self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::BadState, + )); return; } }; @@ -1032,6 +1102,9 @@ impl Node { let Some(transport_id) = peer.transport_id() else { self.connections.remove(&link_id); self.remove_link(&link_id); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::BadState, + )); return; }; if let Some(old_idx) = old_our_index { @@ -1098,6 +1171,9 @@ impl Node { ); self.connections.remove(&link_id); self.links.remove(&link_id); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::BadState, + )); return; } // We lose — abandon our rekey/pending, fall through as responder. @@ -1128,6 +1204,9 @@ impl Node { let Some(conn) = self.connections.get_mut(&link_id) else { warn!(link_id = %link_id, "Connection removed during rekey msg3 processing"); self.links.remove(&link_id); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::UnknownConnection, + )); return; }; conn.take_session() @@ -1140,6 +1219,9 @@ impl Node { warn!("Rekey msg3: no session from handshake"); self.connections.remove(&link_id); self.links.remove(&link_id); + self.stats_mut().record_reject(RejectReason::Handshake( + HandshakeReject::BadState, + )); return; } }; @@ -1296,6 +1378,8 @@ impl Node { // Clean up on promotion failure self.remove_link(&link_id); let _ = self.index_allocator.free(our_index); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); } } } @@ -1326,6 +1410,8 @@ impl Node { receiver_idx = %header.receiver_idx, "No pending inbound or rekey state for msg3" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::UnknownConnection)); return; } }; @@ -1362,6 +1448,8 @@ impl Node { "Rekey msg3 processing failed" ); peer.clear_rekey_responder(); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); } } } diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 9e92a88..2977ec4 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -7,6 +7,7 @@ use crate::NodeAddr; use crate::node::Node; +use crate::node::reject::{HandshakeReject, RejectReason}; use crate::node::wire::build_msg1; use crate::noise::HandshakeState; use crate::protocol::{SessionDatagram, SessionSetup}; @@ -194,6 +195,8 @@ impl Node { error = %e, "Failed to allocate index for rekey" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -212,6 +215,8 @@ impl Node { "Failed to generate rekey msg1" ); let _ = self.index_allocator.free(our_index); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -235,6 +240,8 @@ impl Node { "Failed to send rekey msg1" ); let _ = self.index_allocator.free(our_index); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } } @@ -546,6 +553,8 @@ impl Node { error = %e, "Failed to generate FSP rekey XX msg1" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } }; @@ -567,6 +576,8 @@ impl Node { error = %e, "Failed to send FSP rekey SessionSetup" ); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); return; } diff --git a/src/node/stats.rs b/src/node/stats.rs index 1bcebb9..74445e1 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -86,9 +86,9 @@ impl ForwardingStats { /// The byte-counted side of each outcome is recorded by the /// existing `record_*` methods at the call site (which know the /// payload size); `record_reject` only bumps the packet count and - /// is paired with the byte-aware call at the call site while the - /// typed-rejection rollout is in progress. A later change may - /// collapse the two calls into a single typed entry point. + /// is paired with the byte-aware call at the call site for the + /// duration of the typed-rejection rollout. A later change may collapse the + /// two calls into a single typed entry point. pub(super) fn record_reject(&mut self, reason: ForwardingReject) { match reason { ForwardingReject::DecodeError => self.decode_error_packets += 1, @@ -256,6 +256,12 @@ pub struct BloomStats { pub received: u64, pub decode_error: u64, pub invalid: u64, + /// Non-v1-compliant size class. Reserved for symmetry with the + /// master-side BloomStats shape; the next-side handler does not + /// presently check `is_v1_compliant()` so the counter stays at + /// zero. Kept as a field so `BloomReject::NonV1` dispatches to a + /// real target and the cross-line enum shape is identical. + pub non_v1: u64, pub unknown_peer: u64, pub stale: u64, pub fill_exceeded: u64, @@ -293,6 +299,7 @@ impl BloomStats { received: self.received, decode_error: self.decode_error, invalid: self.invalid, + non_v1: self.non_v1, unknown_peer: self.unknown_peer, stale: self.stale, fill_exceeded: self.fill_exceeded, @@ -647,6 +654,7 @@ pub struct BloomStatsSnapshot { pub received: u64, pub decode_error: u64, pub invalid: u64, + pub non_v1: u64, pub unknown_peer: u64, pub stale: u64, pub fill_exceeded: u64, diff --git a/src/node/tests/bloom_poison.rs b/src/node/tests/bloom_poison.rs index 07d9075..3a7c56d 100644 --- a/src/node/tests/bloom_poison.rs +++ b/src/node/tests/bloom_poison.rs @@ -49,10 +49,10 @@ async fn test_m1_rejects_all_ones_filter_announce() { node.handle_filter_announce(&peer_addr, &payload).await; let after = &node.stats().bloom; - // While the typed-rejection rollout is in progress the call site - // bumps the counter directly AND dispatches through record_reject, - // which hits the same counter. A later change will collapse this to - // a single increment by removing the legacy direct bump; for now + // While the typed-rejection rollout is in progress the call site bumps + // counter directly AND dispatches through record_reject, which + // hits the same counter. A later change will collapse this to a + // single increment by removing the legacy direct bump; for now // the rejection-path event yields a +2 delta. assert_eq!( after.fill_exceeded,