From 71a5c68fa94dfc9247e29ed8c1d9c99bd15794bf Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Sat, 28 Feb 2026 18:23:04 +0000 Subject: [PATCH] Implement comprehensive node and transport statistics MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add 71 new counters (84 values) across three categories: Node statistics (NodeStats, plain u64 — single handler context): - Forwarding: 9 counters x (packets + bytes) = 18 values. Covers received, decode_error, ttl_exhausted, delivered, forwarded, drop_no_route, drop_mtu_exceeded, drop_send_error, originated. - Discovery: 17 counters (packets only). Request path: received, decode_error, duplicate, already_visited, target_is_us, forwarded, ttl_exhausted, initiated, deduplicated. Response path: received, decode_error, forwarded, identity_miss, proof_failed, accepted, timed_out. - Error signals: 3 counters — coords_required, path_broken, mtu_exceeded. - Spanning tree: 16 counters. Inbound announce handling (received through accepted, parent switch, loop detection, ancestry change), outbound (sent, rate limited, send failed), cumulative events (parent switches/losses, flap dampening). - Bloom filter: 10 counters. Inbound (received through accepted), outbound (sent, debounce suppressed, send failed). Transport statistics (AtomicU64 + Arc — shared with spawned tasks): - UDP (6 counters, 8 values): packets/bytes sent/recv, send_errors, recv_errors, mtu_exceeded, kernel_drops (stub for SO_MEMINFO). - TCP (10 counters, 12 values): packets/bytes sent/recv, send_errors, recv_errors, mtu_exceeded, plus connection lifecycle counters (established, accepted, rejected, timeouts, refused). Control socket integration: - show_routing: forwarding, discovery, error signal stats - show_tree: spanning tree stats + per-peer bloom metrics (estimated_count, set_bits, fill_ratio) and coordinate paths - show_bloom: bloom filter stats + per-peer snapshots - show_transports: per-transport stats snapshots Also refactor UDP transport from flat files (udp.rs + udp_stats.rs) into directory module (udp/mod.rs + udp/stats.rs) matching TCP structure, and fix pre-existing clippy warnings in tree/tests.rs. --- src/control/queries.rs | 33 ++- src/node/bloom.rs | 17 +- src/node/handlers/discovery.rs | 20 ++ src/node/handlers/forwarding.rs | 10 + src/node/handlers/session.rs | 10 +- src/node/mod.rs | 19 ++ src/node/stats.rs | 366 +++++++++++++++++++++++++++ src/node/tree.rs | 62 +++-- src/transport/mod.rs | 30 +++ src/transport/tcp/mod.rs | 46 +++- src/transport/tcp/stats.rs | 137 ++++++++++ src/transport/{udp.rs => udp/mod.rs} | 47 +++- src/transport/udp/stats.rs | 105 ++++++++ src/tree/tests.rs | 4 +- testing/sidecar/.env | 7 +- 15 files changed, 864 insertions(+), 49 deletions(-) create mode 100644 src/node/stats.rs create mode 100644 src/transport/tcp/stats.rs rename src/transport/{udp.rs => udp/mod.rs} (94%) create mode 100644 src/transport/udp/stats.rs diff --git a/src/control/queries.rs b/src/control/queries.rs index 08e33bc..b83583c 100644 --- a/src/control/queries.rs +++ b/src/control/queries.rs @@ -151,8 +151,13 @@ pub fn show_tree(node: &Node) -> Value { "display_name": node.peer_display_name(peer_id), }); if let Some(coords) = tree.peer_coords(peer_id) { + let coord_path: Vec = coords.entries() + .iter() + .map(|e| hex::encode(e.node_addr.as_bytes())) + .collect(); peer_json["depth"] = json!(coords.depth()); peer_json["root"] = json!(hex::encode(coords.root_id().as_bytes())); + peer_json["coords"] = json!(coord_path); peer_json["distance_to_us"] = json!(my_coords.distance_to(coords)); } peer_json @@ -163,6 +168,8 @@ pub fn show_tree(node: &Node) -> Value { let parent_hex = hex::encode(parent_addr.as_bytes()); let parent_display = node.peer_display_name(parent_addr); + let tree_stats = node.stats().snapshot().tree; + json!({ "my_node_addr": hex::encode(tree.my_node_addr().as_bytes()), "root": hex::encode(tree.root().as_bytes()), @@ -175,6 +182,7 @@ pub fn show_tree(node: &Node) -> Value { "declaration_signed": decl.is_signed(), "peer_tree_count": tree.peer_count(), "peers": peers, + "stats": serde_json::to_value(&tree_stats).unwrap_or_default(), }) } @@ -229,14 +237,22 @@ pub fn show_bloom(node: &Node) -> Value { // Build per-peer filter info let peer_filters: Vec = node.peers().map(|peer| { let addr = *peer.node_addr(); - json!({ + let mut pf = json!({ "peer": hex::encode(addr.as_bytes()), "display_name": node.peer_display_name(&addr), "has_filter": peer.filter_sequence() > 0, "filter_sequence": peer.filter_sequence(), - }) + }); + if let Some(filter) = peer.inbound_filter() { + pf["estimated_count"] = json!(filter.estimated_count()); + pf["set_bits"] = json!(filter.count_ones()); + pf["fill_ratio"] = json!(filter.fill_ratio()); + } + pf }).collect(); + let bloom_stats = node.stats().snapshot().bloom; + json!({ "own_node_addr": hex::encode(node.node_addr().as_bytes()), "is_leaf_only": node.is_leaf_only(), @@ -244,6 +260,7 @@ pub fn show_bloom(node: &Node) -> Value { "leaf_dependent_count": bloom.leaf_dependents().len(), "leaf_dependents": leaf_deps, "peer_filters": peer_filters, + "stats": serde_json::to_value(&bloom_stats).unwrap_or_default(), }) } @@ -376,22 +393,28 @@ pub fn show_transports(node: &Node) -> Value { t_json["local_addr"] = json!(format!("{}", addr)); } + t_json["stats"] = handle.transport_stats(); + t_json }).collect(); json!({ "transports": transports }) } -/// `show_routing` — Routing table summary. +/// `show_routing` — Routing table summary and node statistics. pub fn show_routing(node: &Node) -> Value { let cache = node.coord_cache(); - let stats = cache.stats(now_ms()); + let cache_stats = cache.stats(now_ms()); + let node_stats = node.stats().snapshot(); json!({ - "coord_cache_entries": stats.entries, + "coord_cache_entries": cache_stats.entries, "identity_cache_entries": node.identity_cache_len(), "pending_lookups": node.pending_lookup_count(), "recent_requests": node.recent_request_count(), + "forwarding": serde_json::to_value(&node_stats.forwarding).unwrap_or_default(), + "discovery": serde_json::to_value(&node_stats.discovery).unwrap_or_default(), + "error_signals": serde_json::to_value(&node_stats.errors).unwrap_or_default(), }) } diff --git a/src/node/bloom.rs b/src/node/bloom.rs index 56a48f9..3b1ad6c 100644 --- a/src/node/bloom.rs +++ b/src/node/bloom.rs @@ -57,6 +57,7 @@ impl Node { // Check debounce if !self.bloom_state.should_send_update(peer_addr, now_ms) { + self.stats_mut().bloom.debounce_suppressed += 1; // Either not pending or rate-limited; will retry on tick return Ok(()); } @@ -70,7 +71,12 @@ impl Node { })?; // Send - self.send_encrypted_link_message(peer_addr, &encoded).await?; + if let Err(e) = self.send_encrypted_link_message(peer_addr, &encoded).await { + self.stats_mut().bloom.send_failed += 1; + return Err(e); + } + + self.stats_mut().bloom.sent += 1; // Record send and store the filter for change detection debug!( @@ -123,9 +129,12 @@ impl Node { /// 3. Store the filter on the peer /// 4. Mark other peers for outgoing filter update pub(super) async fn handle_filter_announce(&mut self, from: &NodeAddr, payload: &[u8]) { + self.stats_mut().bloom.received += 1; + let announce = match FilterAnnounce::decode(payload) { Ok(a) => a, Err(e) => { + self.stats_mut().bloom.decode_error += 1; debug!(from = %self.peer_display_name(from), error = %e, "Malformed FilterAnnounce"); return; } @@ -133,10 +142,12 @@ impl Node { // Validate if !announce.is_valid() { + self.stats_mut().bloom.invalid += 1; debug!(from = %self.peer_display_name(from), "FilterAnnounce filter/size_class mismatch"); return; } if !announce.is_v1_compliant() { + self.stats_mut().bloom.non_v1 += 1; debug!(from = %self.peer_display_name(from), size_class = announce.size_class, "Non-v1 FilterAnnounce rejected"); return; } @@ -145,6 +156,7 @@ impl Node { let current_seq = match self.peers.get(from) { Some(peer) => peer.filter_sequence(), None => { + self.stats_mut().bloom.unknown_peer += 1; debug!(from = %self.peer_display_name(from), "FilterAnnounce from unknown peer"); return; } @@ -152,6 +164,7 @@ impl Node { // Reject stale/replay if announce.sequence <= current_seq { + self.stats_mut().bloom.stale += 1; debug!( from = %self.peer_display_name(from), received_seq = announce.sequence, @@ -161,6 +174,8 @@ impl Node { return; } + self.stats_mut().bloom.accepted += 1; + let now_ms = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_millis() as u64) diff --git a/src/node/handlers/discovery.rs b/src/node/handlers/discovery.rs index cf343f3..35c83ef 100644 --- a/src/node/handlers/discovery.rs +++ b/src/node/handlers/discovery.rs @@ -25,9 +25,12 @@ impl Node { from: &NodeAddr, payload: &[u8], ) { + self.stats_mut().discovery.req_received += 1; + let request = match LookupRequest::decode(payload) { Ok(req) => req, Err(e) => { + self.stats_mut().discovery.req_decode_error += 1; debug!(from = %self.peer_display_name(from), error = %e, "Malformed LookupRequest"); return; } @@ -37,6 +40,7 @@ impl Node { // Dedup: drop if we've already seen this request_id if self.recent_requests.contains_key(&request.request_id) { + self.stats_mut().discovery.req_duplicate += 1; trace!( request_id = request.request_id, from = %self.peer_display_name(from), @@ -56,6 +60,7 @@ impl Node { // Loop prevention: drop if we've already been visited if request.was_visited(self.node_addr()) { + self.stats_mut().discovery.req_already_visited += 1; trace!( request_id = request.request_id, target = %self.peer_display_name(&request.target), @@ -66,6 +71,7 @@ impl Node { // Are we the target? if request.target == *self.node_addr() { + self.stats_mut().discovery.req_target_is_us += 1; debug!( request_id = request.request_id, origin = %self.peer_display_name(&request.origin), @@ -77,8 +83,10 @@ impl Node { // Forward if TTL permits if request.can_forward() { + self.stats_mut().discovery.req_forwarded += 1; self.forward_lookup_request(request).await; } else { + self.stats_mut().discovery.req_ttl_exhausted += 1; trace!( request_id = request.request_id, target = %self.peer_display_name(&request.target), @@ -99,9 +107,12 @@ impl Node { from: &NodeAddr, payload: &[u8], ) { + self.stats_mut().discovery.resp_received += 1; + let mut response = match LookupResponse::decode(payload) { Ok(resp) => resp, Err(e) => { + self.stats_mut().discovery.resp_decode_error += 1; debug!(from = %self.peer_display_name(from), error = %e, "Malformed LookupResponse"); return; } @@ -113,6 +124,7 @@ impl Node { if let Some(recent) = self.recent_requests.get(&response.request_id) { // Transit node: reverse-path forward let from_peer = recent.from_peer; + self.stats_mut().discovery.resp_forwarded += 1; // Apply path_mtu min() from the outgoing link's transport MTU if let Some(peer) = self.peers.get(&from_peer) @@ -153,6 +165,7 @@ impl Node { let target_pubkey = match self.lookup_by_fips_prefix(&prefix) { Some((_addr, pubkey)) => pubkey, None => { + self.stats_mut().discovery.resp_identity_miss += 1; warn!( request_id = response.request_id, target = %self.peer_display_name(&target), @@ -171,6 +184,7 @@ impl Node { &response.target_coords, ); if !peer_id.verify(&proof_data, &response.proof) { + self.stats_mut().discovery.resp_proof_failed += 1; warn!( request_id = response.request_id, target = %self.peer_display_name(&target), @@ -179,6 +193,8 @@ impl Node { return; } + self.stats_mut().discovery.resp_accepted += 1; + debug!( request_id = response.request_id, target = %self.peer_display_name(&target), @@ -335,6 +351,8 @@ impl Node { /// response arrives, it's recognized as "our request" and the /// target's coordinates are cached in coord_cache. pub(in crate::node) async fn initiate_lookup(&mut self, target: &NodeAddr, ttl: u8) { + self.stats_mut().discovery.req_initiated += 1; + let origin = *self.node_addr(); let origin_coords = self.tree_state().my_coords().clone(); let mut request = LookupRequest::generate(*target, origin, origin_coords, ttl, 0); @@ -375,6 +393,7 @@ impl Node { if let Some(&initiated_at) = self.pending_lookups.get(dest) && now_ms.saturating_sub(initiated_at) < lookup_timeout_ms { + self.stats_mut().discovery.req_deduplicated += 1; return; } self.pending_lookups.insert(*dest, now_ms); @@ -396,6 +415,7 @@ impl Node { .collect(); for addr in timed_out { + self.stats_mut().discovery.resp_timed_out += 1; self.pending_lookups.remove(&addr); if let Some(packets) = self.pending_tun_packets.remove(&addr) { for pkt in &packets { diff --git a/src/node/handlers/forwarding.rs b/src/node/handlers/forwarding.rs index e182380..fc78307 100644 --- a/src/node/handlers/forwarding.rs +++ b/src/node/handlers/forwarding.rs @@ -22,9 +22,12 @@ impl Node { /// Called by `dispatch_link_message` for msg_type 0x00. The payload /// has already had its msg_type byte stripped by dispatch. pub(in crate::node) async fn handle_session_datagram(&mut self, _from: &NodeAddr, payload: &[u8]) { + self.stats_mut().forwarding.record_received(payload.len()); + let mut datagram = match SessionDatagram::decode(payload) { Ok(dg) => dg, Err(e) => { + self.stats_mut().forwarding.record_decode_error(payload.len()); debug!(error = %e, "Malformed SessionDatagram"); return; } @@ -32,6 +35,7 @@ impl Node { // TTL enforcement: decrement and drop if exhausted if !datagram.decrement_ttl() { + self.stats_mut().forwarding.record_ttl_exhausted(payload.len()); debug!( src = %datagram.src_addr, dest = %datagram.dest_addr, @@ -45,6 +49,7 @@ impl Node { // Local delivery: dispatch to session layer handlers if datagram.dest_addr == *self.node_addr() { + self.stats_mut().forwarding.record_delivered(payload.len()); self.handle_session_payload(&datagram.src_addr, &datagram.payload, datagram.path_mtu) .await; return; @@ -54,6 +59,7 @@ impl Node { let next_hop_addr = match self.find_next_hop(&datagram.dest_addr) { Some(peer) => *peer.node_addr(), None => { + self.stats_mut().forwarding.record_drop_no_route(payload.len()); self.send_routing_error(&datagram).await; return; } @@ -79,9 +85,11 @@ impl Node { { match e { NodeError::MtuExceeded { mtu, .. } => { + self.stats_mut().forwarding.record_drop_mtu_exceeded(payload.len()); self.send_mtu_exceeded_error(&datagram, mtu).await; } _ => { + self.stats_mut().forwarding.record_drop_send_error(payload.len()); debug!( next_hop = %next_hop_addr, dest = %datagram.dest_addr, @@ -90,6 +98,8 @@ impl Node { ); } } + } else { + self.stats_mut().forwarding.record_forwarded(encoded.len()); } } diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 9825940..00ba75d 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -703,6 +703,8 @@ impl Node { /// immediately (rate-limited), trigger discovery, and reset the /// warmup counter for subsequent data packets. async fn handle_coords_required(&mut self, inner: &[u8]) { + self.stats_mut().errors.coords_required += 1; + let msg = match CoordsRequired::decode(inner) { Ok(m) => m, Err(e) => { @@ -759,6 +761,8 @@ impl Node { /// Send a standalone CoordsWarmup immediately (rate-limited), invalidate /// cached coordinates, trigger re-discovery, and reset the warmup counter. async fn handle_path_broken(&mut self, inner: &[u8]) { + self.stats_mut().errors.path_broken += 1; + let msg = match PathBroken::decode(inner) { Ok(m) => m, Err(e) => { @@ -820,6 +824,8 @@ impl Node { /// next-hop transport MTU. Apply the reported bottleneck MTU to our /// PathMtuState for the affected session, causing an immediate decrease. async fn handle_mtu_exceeded(&mut self, inner: &[u8]) { + self.stats_mut().errors.mtu_exceeded += 1; + let msg = match MtuExceeded::decode(inner) { Ok(m) => m, Err(e) => { @@ -1231,7 +1237,9 @@ impl Node { } let encoded = datagram.encode(); - self.send_encrypted_link_message(&next_hop_addr, &encoded).await + self.send_encrypted_link_message(&next_hop_addr, &encoded).await?; + self.stats_mut().forwarding.record_originated(encoded.len()); + Ok(()) } /// Look up destination coordinates from available caches. diff --git a/src/node/mod.rs b/src/node/mod.rs index 1b18a1a..8c69011 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -13,6 +13,7 @@ mod routing_error_rate_limit; pub(crate) mod session; pub(crate) mod session_wire; pub(crate) mod wire; +pub(crate) mod stats; mod tree; #[cfg(test)] mod tests; @@ -299,6 +300,10 @@ pub struct Node { /// Next transport ID to allocate. next_transport_id: u32, + // === Node Statistics === + /// Routing, forwarding, discovery, and error signal counters. + stats: stats::NodeStats, + // === TUN Interface === /// TUN device state. tun_state: TunState, @@ -433,6 +438,7 @@ impl Node { max_links, next_link_id: 1, next_transport_id: 1, + stats: stats::NodeStats::new(), tun_state, tun_name: None, tun_tx: None, @@ -526,6 +532,7 @@ impl Node { max_links, next_link_id: 1, next_transport_id: 1, + stats: stats::NodeStats::new(), tun_state, tun_name: None, tun_tx: None, @@ -803,6 +810,18 @@ impl Node { &mut self.coord_cache } + // === Node Statistics === + + /// Get the node statistics. + pub fn stats(&self) -> &stats::NodeStats { + &self.stats + } + + /// Get mutable node statistics. + pub(crate) fn stats_mut(&mut self) -> &mut stats::NodeStats { + &mut self.stats + } + // === TUN Interface === /// Get the TUN state. diff --git a/src/node/stats.rs b/src/node/stats.rs new file mode 100644 index 0000000..8f1f657 --- /dev/null +++ b/src/node/stats.rs @@ -0,0 +1,366 @@ +//! Node-level statistics for routing, forwarding, and discovery operations. +//! +//! Unlike `EthernetStats` (which uses `AtomicU64` + `Arc` for cross-task +//! sharing), these counters use plain `u64` because `Node` handlers run +//! on a single `&mut self` context. A `snapshot()` method produces a +//! copyable struct for control socket queries. + +use serde::Serialize; + +/// Forwarding statistics — packets and bytes for each outcome. +#[derive(Default)] +pub struct ForwardingStats { + pub received_packets: u64, + pub received_bytes: u64, + pub decode_error_packets: u64, + pub decode_error_bytes: u64, + pub ttl_exhausted_packets: u64, + pub ttl_exhausted_bytes: u64, + pub delivered_packets: u64, + pub delivered_bytes: u64, + pub forwarded_packets: u64, + pub forwarded_bytes: u64, + pub drop_no_route_packets: u64, + pub drop_no_route_bytes: u64, + pub drop_mtu_exceeded_packets: u64, + pub drop_mtu_exceeded_bytes: u64, + pub drop_send_error_packets: u64, + pub drop_send_error_bytes: u64, + pub originated_packets: u64, + pub originated_bytes: u64, +} + +impl ForwardingStats { + pub fn record_received(&mut self, bytes: usize) { + self.received_packets += 1; + self.received_bytes += bytes as u64; + } + + pub fn record_decode_error(&mut self, bytes: usize) { + self.decode_error_packets += 1; + self.decode_error_bytes += bytes as u64; + } + + pub fn record_ttl_exhausted(&mut self, bytes: usize) { + self.ttl_exhausted_packets += 1; + self.ttl_exhausted_bytes += bytes as u64; + } + + pub fn record_delivered(&mut self, bytes: usize) { + self.delivered_packets += 1; + self.delivered_bytes += bytes as u64; + } + + pub fn record_forwarded(&mut self, bytes: usize) { + self.forwarded_packets += 1; + self.forwarded_bytes += bytes as u64; + } + + pub fn record_drop_no_route(&mut self, bytes: usize) { + self.drop_no_route_packets += 1; + self.drop_no_route_bytes += bytes as u64; + } + + pub fn record_drop_mtu_exceeded(&mut self, bytes: usize) { + self.drop_mtu_exceeded_packets += 1; + self.drop_mtu_exceeded_bytes += bytes as u64; + } + + pub fn record_drop_send_error(&mut self, bytes: usize) { + self.drop_send_error_packets += 1; + self.drop_send_error_bytes += bytes as u64; + } + + pub fn record_originated(&mut self, bytes: usize) { + self.originated_packets += 1; + self.originated_bytes += bytes as u64; + } + + pub fn snapshot(&self) -> ForwardingStatsSnapshot { + ForwardingStatsSnapshot { + received_packets: self.received_packets, + received_bytes: self.received_bytes, + decode_error_packets: self.decode_error_packets, + decode_error_bytes: self.decode_error_bytes, + ttl_exhausted_packets: self.ttl_exhausted_packets, + ttl_exhausted_bytes: self.ttl_exhausted_bytes, + delivered_packets: self.delivered_packets, + delivered_bytes: self.delivered_bytes, + forwarded_packets: self.forwarded_packets, + forwarded_bytes: self.forwarded_bytes, + drop_no_route_packets: self.drop_no_route_packets, + drop_no_route_bytes: self.drop_no_route_bytes, + drop_mtu_exceeded_packets: self.drop_mtu_exceeded_packets, + drop_mtu_exceeded_bytes: self.drop_mtu_exceeded_bytes, + drop_send_error_packets: self.drop_send_error_packets, + drop_send_error_bytes: self.drop_send_error_bytes, + originated_packets: self.originated_packets, + originated_bytes: self.originated_bytes, + } + } +} + +/// Discovery statistics — packet counts for request and response handling. +#[derive(Default)] +pub struct DiscoveryStats { + // Request counters + pub req_received: u64, + pub req_decode_error: u64, + pub req_duplicate: u64, + pub req_already_visited: u64, + pub req_target_is_us: u64, + pub req_forwarded: u64, + pub req_ttl_exhausted: u64, + pub req_initiated: u64, + pub req_deduplicated: u64, + // Response counters + pub resp_received: u64, + pub resp_decode_error: u64, + pub resp_forwarded: u64, + pub resp_identity_miss: u64, + pub resp_proof_failed: u64, + pub resp_accepted: u64, + pub resp_timed_out: u64, +} + +impl DiscoveryStats { + pub fn snapshot(&self) -> DiscoveryStatsSnapshot { + DiscoveryStatsSnapshot { + req_received: self.req_received, + req_decode_error: self.req_decode_error, + req_duplicate: self.req_duplicate, + req_already_visited: self.req_already_visited, + req_target_is_us: self.req_target_is_us, + req_forwarded: self.req_forwarded, + req_ttl_exhausted: self.req_ttl_exhausted, + req_initiated: self.req_initiated, + req_deduplicated: self.req_deduplicated, + resp_received: self.resp_received, + resp_decode_error: self.resp_decode_error, + resp_forwarded: self.resp_forwarded, + resp_identity_miss: self.resp_identity_miss, + resp_proof_failed: self.resp_proof_failed, + resp_accepted: self.resp_accepted, + resp_timed_out: self.resp_timed_out, + } + } +} + +/// Spanning tree statistics — announce handling and parent tracking. +#[derive(Default)] +pub struct TreeStats { + // Inbound announce handling + pub received: u64, + pub decode_error: u64, + pub unknown_peer: u64, + pub addr_mismatch: u64, + pub sig_failed: u64, + pub stale: u64, + pub accepted: u64, + pub parent_switched: u64, + pub loop_detected: u64, + pub ancestry_changed: u64, + // Outbound announce sending + pub sent: u64, + pub rate_limited: u64, + pub send_failed: u64, + // Cumulative events + pub parent_switches: u64, + pub parent_losses: u64, + pub flap_dampened: u64, +} + +impl TreeStats { + pub fn snapshot(&self) -> TreeStatsSnapshot { + TreeStatsSnapshot { + received: self.received, + decode_error: self.decode_error, + unknown_peer: self.unknown_peer, + addr_mismatch: self.addr_mismatch, + sig_failed: self.sig_failed, + stale: self.stale, + accepted: self.accepted, + parent_switched: self.parent_switched, + loop_detected: self.loop_detected, + ancestry_changed: self.ancestry_changed, + sent: self.sent, + rate_limited: self.rate_limited, + send_failed: self.send_failed, + parent_switches: self.parent_switches, + parent_losses: self.parent_losses, + flap_dampened: self.flap_dampened, + } + } +} + +/// Bloom filter statistics — filter announce handling. +#[derive(Default)] +pub struct BloomStats { + // Inbound announce handling + pub received: u64, + pub decode_error: u64, + pub invalid: u64, + pub non_v1: u64, + pub unknown_peer: u64, + pub stale: u64, + pub accepted: u64, + // Outbound announce sending + pub sent: u64, + pub debounce_suppressed: u64, + pub send_failed: u64, +} + +impl BloomStats { + pub fn snapshot(&self) -> BloomStatsSnapshot { + BloomStatsSnapshot { + received: self.received, + decode_error: self.decode_error, + invalid: self.invalid, + non_v1: self.non_v1, + unknown_peer: self.unknown_peer, + stale: self.stale, + accepted: self.accepted, + sent: self.sent, + debounce_suppressed: self.debounce_suppressed, + send_failed: self.send_failed, + } + } +} + +/// Error signal statistics — counts of each error signal type received. +#[derive(Default)] +pub struct ErrorSignalStats { + pub coords_required: u64, + pub path_broken: u64, + pub mtu_exceeded: u64, +} + +impl ErrorSignalStats { + pub fn snapshot(&self) -> ErrorSignalStatsSnapshot { + ErrorSignalStatsSnapshot { + coords_required: self.coords_required, + path_broken: self.path_broken, + mtu_exceeded: self.mtu_exceeded, + } + } +} + +/// Aggregate node statistics. +#[derive(Default)] +pub struct NodeStats { + pub forwarding: ForwardingStats, + pub discovery: DiscoveryStats, + pub tree: TreeStats, + pub bloom: BloomStats, + pub errors: ErrorSignalStats, +} + +impl NodeStats { + pub fn new() -> Self { + Self::default() + } + + pub fn snapshot(&self) -> NodeStatsSnapshot { + NodeStatsSnapshot { + forwarding: self.forwarding.snapshot(), + discovery: self.discovery.snapshot(), + tree: self.tree.snapshot(), + bloom: self.bloom.snapshot(), + errors: self.errors.snapshot(), + } + } +} + +// --- Snapshot types (copyable, serializable) --- + +#[derive(Clone, Debug, Default, Serialize)] +pub struct ForwardingStatsSnapshot { + pub received_packets: u64, + pub received_bytes: u64, + pub decode_error_packets: u64, + pub decode_error_bytes: u64, + pub ttl_exhausted_packets: u64, + pub ttl_exhausted_bytes: u64, + pub delivered_packets: u64, + pub delivered_bytes: u64, + pub forwarded_packets: u64, + pub forwarded_bytes: u64, + pub drop_no_route_packets: u64, + pub drop_no_route_bytes: u64, + pub drop_mtu_exceeded_packets: u64, + pub drop_mtu_exceeded_bytes: u64, + pub drop_send_error_packets: u64, + pub drop_send_error_bytes: u64, + pub originated_packets: u64, + pub originated_bytes: u64, +} + +#[derive(Clone, Debug, Default, Serialize)] +pub struct DiscoveryStatsSnapshot { + pub req_received: u64, + pub req_decode_error: u64, + pub req_duplicate: u64, + pub req_already_visited: u64, + pub req_target_is_us: u64, + pub req_forwarded: u64, + pub req_ttl_exhausted: u64, + pub req_initiated: u64, + pub req_deduplicated: u64, + pub resp_received: u64, + pub resp_decode_error: u64, + pub resp_forwarded: u64, + pub resp_identity_miss: u64, + pub resp_proof_failed: u64, + pub resp_accepted: u64, + pub resp_timed_out: u64, +} + +#[derive(Clone, Debug, Default, Serialize)] +pub struct TreeStatsSnapshot { + pub received: u64, + pub decode_error: u64, + pub unknown_peer: u64, + pub addr_mismatch: u64, + pub sig_failed: u64, + pub stale: u64, + pub accepted: u64, + pub parent_switched: u64, + pub loop_detected: u64, + pub ancestry_changed: u64, + pub sent: u64, + pub rate_limited: u64, + pub send_failed: u64, + pub parent_switches: u64, + pub parent_losses: u64, + pub flap_dampened: u64, +} + +#[derive(Clone, Debug, Default, Serialize)] +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 accepted: u64, + pub sent: u64, + pub debounce_suppressed: u64, + pub send_failed: u64, +} + +#[derive(Clone, Debug, Default, Serialize)] +pub struct ErrorSignalStatsSnapshot { + pub coords_required: u64, + pub path_broken: u64, + pub mtu_exceeded: u64, +} + +#[derive(Clone, Debug, Default, Serialize)] +pub struct NodeStatsSnapshot { + pub forwarding: ForwardingStatsSnapshot, + pub discovery: DiscoveryStatsSnapshot, + pub tree: TreeStatsSnapshot, + pub bloom: BloomStatsSnapshot, + pub errors: ErrorSignalStatsSnapshot, +} diff --git a/src/node/tree.rs b/src/node/tree.rs index 73df275..7605449 100644 --- a/src/node/tree.rs +++ b/src/node/tree.rs @@ -48,6 +48,7 @@ impl Node { if !peer.can_send_tree_announce(now_ms) { peer.mark_tree_announce_pending(); + self.stats_mut().tree.rate_limited += 1; debug!( peer = %self.peer_display_name(peer_addr), "TreeAnnounce rate-limited, marking pending" @@ -63,7 +64,12 @@ impl Node { })?; // Send - self.send_encrypted_link_message(peer_addr, &encoded).await?; + if let Err(e) = self.send_encrypted_link_message(peer_addr, &encoded).await { + self.stats_mut().tree.send_failed += 1; + return Err(e); + } + + self.stats_mut().tree.sent += 1; // Record send time if let Some(peer) = self.peers.get_mut(peer_addr) { @@ -122,9 +128,12 @@ impl Node { /// 4. Re-evaluate parent selection /// 5. If parent changed: increment seq, sign, recompute coords, announce to all pub(super) async fn handle_tree_announce(&mut self, from: &NodeAddr, payload: &[u8]) { + self.stats_mut().tree.received += 1; + let announce = match TreeAnnounce::decode(payload) { Ok(a) => a, Err(e) => { + self.stats_mut().tree.decode_error += 1; debug!(from = %self.peer_display_name(from), error = %e, "Malformed TreeAnnounce"); return; } @@ -134,6 +143,7 @@ impl Node { let pubkey = match self.peers.get(from) { Some(peer) => peer.pubkey(), None => { + self.stats_mut().tree.unknown_peer += 1; debug!(from = %self.peer_display_name(from), "TreeAnnounce from unknown peer"); return; } @@ -141,6 +151,7 @@ impl Node { // The declaring node_addr in the announce should match the sender if announce.declaration.node_addr() != from { + self.stats_mut().tree.addr_mismatch += 1; debug!( from = %self.peer_display_name(from), declared = %announce.declaration.node_addr(), @@ -150,6 +161,7 @@ impl Node { } if let Err(e) = announce.declaration.verify(&pubkey) { + self.stats_mut().tree.sig_failed += 1; warn!( from = %self.peer_display_name(from), error = %e, @@ -179,10 +191,13 @@ impl Node { ); if !updated { + self.stats_mut().tree.stale += 1; debug!(from = %self.peer_display_name(from), "TreeAnnounce not fresher than existing, ignored"); return; } + self.stats_mut().tree.accepted += 1; + info!( from = %self.peer_display_name(from), seq = announce.declaration.sequence(), @@ -215,6 +230,9 @@ impl Node { self.tree_state.recompute_coords(); self.coord_cache.clear(); + self.stats_mut().tree.parent_switched += 1; + self.stats_mut().tree.parent_switches += 1; + info!( new_parent = %self.peer_display_name(&new_parent), new_seq = new_seq, @@ -223,6 +241,7 @@ impl Node { "Parent switched, flushed coord cache, announcing to all peers" ); if flap_dampened { + self.stats_mut().tree.flap_dampened += 1; warn!("Flap dampening engaged: excessive parent switches detected"); } @@ -235,25 +254,26 @@ impl Node { && *self.tree_state.my_declaration().parent_id() == *from { // Check for loop: if parent's ancestry now contains us, drop parent - if let Some(parent_coords) = self.tree_state.peer_coords(from) { - if parent_coords.contains(self.identity.node_addr()) { - warn!( - parent = %self.peer_display_name(from), - "Parent ancestry contains us — loop detected, dropping parent" - ); - let peer_costs: HashMap = self.peers.iter() - .map(|(addr, peer)| (*addr, peer.link_cost())) - .collect(); - if self.tree_state.handle_parent_lost(&peer_costs) { - if let Err(e) = self.tree_state.sign_declaration(&self.identity) { - warn!(error = %e, "Failed to sign declaration after loop detection"); - return; - } - self.coord_cache.clear(); - self.send_tree_announce_to_all().await; + if let Some(parent_coords) = self.tree_state.peer_coords(from) + && parent_coords.contains(self.identity.node_addr()) + { + self.stats_mut().tree.loop_detected += 1; + warn!( + parent = %self.peer_display_name(from), + "Parent ancestry contains us — loop detected, dropping parent" + ); + let peer_costs: HashMap = self.peers.iter() + .map(|(addr, peer)| (*addr, peer.link_cost())) + .collect(); + if self.tree_state.handle_parent_lost(&peer_costs) { + if let Err(e) = self.tree_state.sign_declaration(&self.identity) { + warn!(error = %e, "Failed to sign declaration after loop detection"); + return; } - return; + self.coord_cache.clear(); + self.send_tree_announce_to_all().await; } + return; } // Our parent's ancestry changed but we're keeping the same parent. @@ -280,6 +300,7 @@ impl Node { let new_depth = self.tree_state.my_coords().depth(); if new_root != old_root || new_depth != old_depth { + self.stats_mut().tree.ancestry_changed += 1; info!( parent = %self.peer_display_name(from), old_root = %old_root, @@ -351,6 +372,9 @@ impl Node { self.tree_state.recompute_coords(); self.coord_cache.clear(); + self.stats_mut().tree.parent_switched += 1; + self.stats_mut().tree.parent_switches += 1; + info!( new_parent = %self.peer_display_name(&new_parent), new_seq = new_seq, @@ -360,6 +384,7 @@ impl Node { "Parent switched via periodic cost re-evaluation" ); if flap_dampened { + self.stats_mut().tree.flap_dampened += 1; warn!("Flap dampening engaged: excessive parent switches detected"); } @@ -383,6 +408,7 @@ impl Node { self.tree_state.remove_peer(node_addr); if was_parent { + self.stats_mut().tree.parent_losses += 1; let peer_costs: HashMap = self.peers.iter() .map(|(addr, peer)| (*addr, peer.link_cost())) .collect(); diff --git a/src/transport/mod.rs b/src/transport/mod.rs index 8804a79..b5a98ac 100644 --- a/src/transport/mod.rs +++ b/src/transport/mod.rs @@ -970,6 +970,36 @@ impl TransportHandle { pub fn is_operational(&self) -> bool { self.state().is_operational() } + + /// Get transport-specific stats as a JSON value. + /// + /// Returns a snapshot of counters for the specific transport type. + pub fn transport_stats(&self) -> serde_json::Value { + match self { + TransportHandle::Udp(t) => { + serde_json::to_value(t.stats().snapshot()).unwrap_or_default() + } + #[cfg(target_os = "linux")] + TransportHandle::Ethernet(t) => { + let snap = t.stats().snapshot(); + serde_json::json!({ + "frames_sent": snap.frames_sent, + "frames_recv": snap.frames_recv, + "bytes_sent": snap.bytes_sent, + "bytes_recv": snap.bytes_recv, + "send_errors": snap.send_errors, + "recv_errors": snap.recv_errors, + "beacons_sent": snap.beacons_sent, + "beacons_recv": snap.beacons_recv, + "frames_too_short": snap.frames_too_short, + "frames_too_long": snap.frames_too_long, + }) + } + TransportHandle::Tcp(t) => { + serde_json::to_value(t.stats().snapshot()).unwrap_or_default() + } + } + } } // ============================================================================ diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index 77ddeb7..5bb9057 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -22,6 +22,7 @@ //! No additional framing overhead — packets are written directly to the //! TCP stream and the receiver uses phase-dependent size computation. +pub mod stats; pub mod stream; use super::{ @@ -29,6 +30,7 @@ use super::{ TransportId, TransportState, TransportType, }; use crate::config::TcpConfig; +use stats::TcpStats; use stream::read_fmp_packet; use socket2::TcpKeepalive; @@ -91,6 +93,8 @@ pub struct TcpTransport { accept_task: Option>, /// Local listener address (after start, if bind_addr configured). local_addr: Option, + /// Transport statistics. + stats: Arc, } impl TcpTransport { @@ -110,6 +114,7 @@ impl TcpTransport { packet_tx, accept_task: None, local_addr: None, + stats: Arc::new(TcpStats::new()), } } @@ -123,6 +128,11 @@ impl TcpTransport { self.local_addr } + /// Get the transport statistics. + pub fn stats(&self) -> &Arc { + &self.stats + } + /// Start the transport asynchronously. /// /// If `bind_addr` is configured, binds a TCP listener and spawns @@ -154,6 +164,7 @@ impl TcpTransport { let transport_id = self.transport_id; let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); + let stats = self.stats.clone(); let cfg = AcceptConfig { mtu: self.config.mtu(), max_inbound: self.config.max_inbound_connections(), @@ -164,7 +175,7 @@ impl TcpTransport { }; let accept_task = tokio::spawn(async move { - accept_loop(listener, transport_id, packet_tx, pool, cfg).await; + accept_loop(listener, transport_id, packet_tx, pool, cfg, stats).await; }); self.accept_task = Some(accept_task); } @@ -245,6 +256,7 @@ impl TcpTransport { // disruptive reset-reconnect cycle. let mtu = self.config.mtu() as usize; if data.len() > mtu { + self.stats.record_mtu_exceeded(); return Err(TransportError::MtuExceeded { packet_size: data.len(), mtu: self.config.mtu(), @@ -269,6 +281,7 @@ impl TcpTransport { let mut w = writer.lock().await; match w.write_all(data).await { Ok(()) => { + self.stats.record_send(data.len()); trace!( transport_id = %self.transport_id, remote_addr = %addr, @@ -278,6 +291,7 @@ impl TcpTransport { Ok(data.len()) } Err(e) => { + self.stats.record_send_error(); drop(w); // Remove failed connection from pool let mut pool = self.pool.lock().await; @@ -301,13 +315,22 @@ impl TcpTransport { let timeout_ms = self.config.connect_timeout_ms(); // Connect with timeout - let stream = tokio::time::timeout( + let stream = match tokio::time::timeout( Duration::from_millis(timeout_ms), TcpStream::connect(socket_addr), ) .await - .map_err(|_| TransportError::Timeout)? - .map_err(|_| TransportError::ConnectionRefused)?; + { + Ok(Ok(stream)) => stream, + Ok(Err(_)) => { + self.stats.record_connect_refused(); + return Err(TransportError::ConnectionRefused); + } + Err(_) => { + self.stats.record_connect_timeout(); + return Err(TransportError::Timeout); + } + }; // Configure socket options via socket2 let std_stream = stream.into_std() @@ -328,11 +351,12 @@ impl TcpTransport { let transport_id = self.transport_id; let packet_tx = self.packet_tx.clone(); let pool = self.pool.clone(); + let recv_stats = self.stats.clone(); let remote_addr = addr.clone(); let mtu = mss_mtu; let recv_task = tokio::spawn(async move { - tcp_receive_loop(read_half, transport_id, remote_addr.clone(), packet_tx, pool, mtu).await; + tcp_receive_loop(read_half, transport_id, remote_addr.clone(), packet_tx, pool, mtu, recv_stats).await; }); let conn = TcpConnection { @@ -345,6 +369,8 @@ impl TcpTransport { let mut pool = self.pool.lock().await; pool.insert(addr.clone(), conn); + self.stats.record_connection_established(); + debug!( transport_id = %self.transport_id, remote_addr = %addr, @@ -447,6 +473,7 @@ async fn accept_loop( packet_tx: PacketTx, pool: ConnectionPool, cfg: AcceptConfig, + stats: Arc, ) { let AcceptConfig { mtu, max_inbound, nodelay, keepalive_secs, recv_buf, send_buf } = cfg; debug!(transport_id = %transport_id, "TCP accept loop starting"); @@ -458,6 +485,7 @@ async fn accept_loop( { let pool_guard = pool.lock().await; if pool_guard.len() >= max_inbound { + stats.record_connection_rejected(); warn!( transport_id = %transport_id, peer_addr = %peer_addr, @@ -514,6 +542,7 @@ async fn accept_loop( let recv_pool = pool.clone(); let recv_packet_tx = packet_tx.clone(); + let recv_stats = stats.clone(); let recv_addr = remote_addr.clone(); let recv_task = tokio::spawn(async move { @@ -524,6 +553,7 @@ async fn accept_loop( recv_packet_tx, recv_pool, conn_mtu, + recv_stats, ) .await; }); @@ -538,6 +568,8 @@ async fn accept_loop( let mut pool_guard = pool.lock().await; pool_guard.insert(remote_addr.clone(), conn); + stats.record_connection_accepted(); + debug!( transport_id = %transport_id, remote_addr = %remote_addr, @@ -572,6 +604,7 @@ async fn tcp_receive_loop( packet_tx: PacketTx, pool: ConnectionPool, mtu: u16, + stats: Arc, ) { debug!( transport_id = %transport_id, @@ -582,6 +615,8 @@ async fn tcp_receive_loop( loop { match read_fmp_packet(&mut reader, mtu).await { Ok(data) => { + stats.record_recv(data.len()); + trace!( transport_id = %transport_id, remote_addr = %remote_addr, @@ -604,6 +639,7 @@ async fn tcp_receive_loop( } } Err(e) => { + stats.record_recv_error(); // EOF or protocol error — remove connection from pool debug!( transport_id = %transport_id, diff --git a/src/transport/tcp/stats.rs b/src/transport/tcp/stats.rs new file mode 100644 index 0000000..badcaa9 --- /dev/null +++ b/src/transport/tcp/stats.rs @@ -0,0 +1,137 @@ +//! TCP transport statistics. + +use std::sync::atomic::{AtomicU64, Ordering}; + +use serde::Serialize; + +/// Statistics for a TCP transport instance. +/// +/// Uses atomic counters for lock-free updates from per-connection +/// receive loops and the send path concurrently. +pub struct TcpStats { + pub packets_sent: AtomicU64, + pub bytes_sent: AtomicU64, + pub packets_recv: AtomicU64, + pub bytes_recv: AtomicU64, + pub send_errors: AtomicU64, + pub recv_errors: AtomicU64, + pub mtu_exceeded: AtomicU64, + pub connections_established: AtomicU64, + pub connections_accepted: AtomicU64, + pub connections_rejected: AtomicU64, + pub connect_timeouts: AtomicU64, + pub connect_refused: AtomicU64, +} + +impl TcpStats { + /// Create a new stats instance with all counters at zero. + pub fn new() -> Self { + Self { + packets_sent: AtomicU64::new(0), + bytes_sent: AtomicU64::new(0), + packets_recv: AtomicU64::new(0), + bytes_recv: AtomicU64::new(0), + send_errors: AtomicU64::new(0), + recv_errors: AtomicU64::new(0), + mtu_exceeded: AtomicU64::new(0), + connections_established: AtomicU64::new(0), + connections_accepted: AtomicU64::new(0), + connections_rejected: AtomicU64::new(0), + connect_timeouts: AtomicU64::new(0), + connect_refused: AtomicU64::new(0), + } + } + + /// Record a successful send. + pub fn record_send(&self, bytes: usize) { + self.packets_sent.fetch_add(1, Ordering::Relaxed); + self.bytes_sent.fetch_add(bytes as u64, Ordering::Relaxed); + } + + /// Record a successful receive. + pub fn record_recv(&self, bytes: usize) { + self.packets_recv.fetch_add(1, Ordering::Relaxed); + self.bytes_recv.fetch_add(bytes as u64, Ordering::Relaxed); + } + + /// Record a send error. + pub fn record_send_error(&self) { + self.send_errors.fetch_add(1, Ordering::Relaxed); + } + + /// Record a receive error. + pub fn record_recv_error(&self) { + self.recv_errors.fetch_add(1, Ordering::Relaxed); + } + + /// Record an MTU exceeded rejection. + pub fn record_mtu_exceeded(&self) { + self.mtu_exceeded.fetch_add(1, Ordering::Relaxed); + } + + /// Record a successful outbound connection. + pub fn record_connection_established(&self) { + self.connections_established.fetch_add(1, Ordering::Relaxed); + } + + /// Record a successful inbound connection. + pub fn record_connection_accepted(&self) { + self.connections_accepted.fetch_add(1, Ordering::Relaxed); + } + + /// Record a rejected inbound connection. + pub fn record_connection_rejected(&self) { + self.connections_rejected.fetch_add(1, Ordering::Relaxed); + } + + /// Record a connect timeout. + pub fn record_connect_timeout(&self) { + self.connect_timeouts.fetch_add(1, Ordering::Relaxed); + } + + /// Record a connection refused. + pub fn record_connect_refused(&self) { + self.connect_refused.fetch_add(1, Ordering::Relaxed); + } + + /// Take a snapshot of all counters. + pub fn snapshot(&self) -> TcpStatsSnapshot { + TcpStatsSnapshot { + packets_sent: self.packets_sent.load(Ordering::Relaxed), + bytes_sent: self.bytes_sent.load(Ordering::Relaxed), + packets_recv: self.packets_recv.load(Ordering::Relaxed), + bytes_recv: self.bytes_recv.load(Ordering::Relaxed), + send_errors: self.send_errors.load(Ordering::Relaxed), + recv_errors: self.recv_errors.load(Ordering::Relaxed), + mtu_exceeded: self.mtu_exceeded.load(Ordering::Relaxed), + connections_established: self.connections_established.load(Ordering::Relaxed), + connections_accepted: self.connections_accepted.load(Ordering::Relaxed), + connections_rejected: self.connections_rejected.load(Ordering::Relaxed), + connect_timeouts: self.connect_timeouts.load(Ordering::Relaxed), + connect_refused: self.connect_refused.load(Ordering::Relaxed), + } + } +} + +impl Default for TcpStats { + fn default() -> Self { + Self::new() + } +} + +/// Point-in-time snapshot of TCP stats (non-atomic, copyable). +#[derive(Clone, Debug, Default, Serialize)] +pub struct TcpStatsSnapshot { + pub packets_sent: u64, + pub bytes_sent: u64, + pub packets_recv: u64, + pub bytes_recv: u64, + pub send_errors: u64, + pub recv_errors: u64, + pub mtu_exceeded: u64, + pub connections_established: u64, + pub connections_accepted: u64, + pub connections_rejected: u64, + pub connect_timeouts: u64, + pub connect_refused: u64, +} diff --git a/src/transport/udp.rs b/src/transport/udp/mod.rs similarity index 94% rename from src/transport/udp.rs rename to src/transport/udp/mod.rs index faf85de..61252dd 100644 --- a/src/transport/udp.rs +++ b/src/transport/udp/mod.rs @@ -6,6 +6,8 @@ use super::{ DiscoveredPeer, PacketTx, ReceivedPacket, Transport, TransportAddr, TransportError, TransportId, TransportState, TransportType, }; +mod stats; +use stats::UdpStats; use crate::config::UdpConfig; use socket2::{Domain, Protocol, Socket, Type}; use std::net::SocketAddr; @@ -36,6 +38,8 @@ pub struct UdpTransport { recv_task: Option>, /// Local bound address (after start). local_addr: Option, + /// Transport statistics. + stats: Arc, } impl UdpTransport { @@ -55,6 +59,7 @@ impl UdpTransport { packet_tx, recv_task: None, local_addr: None, + stats: Arc::new(UdpStats::new()), } } @@ -73,6 +78,11 @@ impl UdpTransport { self.socket.as_ref() } + /// Get the transport statistics. + pub fn stats(&self) -> &Arc { + &self.stats + } + /// Start the transport asynchronously. /// /// Binds the UDP socket and spawns the receive loop. @@ -145,9 +155,10 @@ impl UdpTransport { let transport_id = self.transport_id; let packet_tx = self.packet_tx.clone(); let mtu = self.config.mtu(); + let stats = self.stats.clone(); let recv_task = tokio::spawn(async move { - udp_receive_loop(socket, transport_id, packet_tx, mtu).await; + udp_receive_loop(socket, transport_id, packet_tx, mtu, stats).await; }); self.recv_task = Some(recv_task); @@ -210,6 +221,7 @@ impl UdpTransport { } if data.len() > self.config.mtu() as usize { + self.stats.record_mtu_exceeded(); return Err(TransportError::MtuExceeded { packet_size: data.len(), mtu: self.config.mtu(), @@ -219,19 +231,22 @@ impl UdpTransport { let socket_addr = parse_socket_addr(addr)?; let socket = self.socket.as_ref().ok_or(TransportError::NotStarted)?; - let bytes_sent = socket - .send_to(data, socket_addr) - .await - .map_err(|e| TransportError::SendFailed(format!("{}", e)))?; - - trace!( - transport_id = %self.transport_id, - remote_addr = %socket_addr, - bytes = bytes_sent, - "UDP packet sent" - ); - - Ok(bytes_sent) + match socket.send_to(data, socket_addr).await { + Ok(bytes_sent) => { + self.stats.record_send(bytes_sent); + trace!( + transport_id = %self.transport_id, + remote_addr = %socket_addr, + bytes = bytes_sent, + "UDP packet sent" + ); + Ok(bytes_sent) + } + Err(e) => { + self.stats.record_send_error(); + Err(TransportError::SendFailed(format!("{}", e))) + } + } } } @@ -294,6 +309,7 @@ async fn udp_receive_loop( transport_id: TransportId, packet_tx: PacketTx, mtu: u16, + stats: Arc, ) { // Buffer with headroom for slightly oversized packets let mut buf = vec![0u8; mtu as usize + 100]; @@ -303,6 +319,8 @@ async fn udp_receive_loop( loop { match socket.recv_from(&mut buf).await { Ok((len, remote_addr)) => { + stats.record_recv(len); + let data = buf[..len].to_vec(); let addr = TransportAddr::from_string(&remote_addr.to_string()); let packet = ReceivedPacket::new(transport_id, addr, data); @@ -324,6 +342,7 @@ async fn udp_receive_loop( } } Err(e) => { + stats.record_recv_error(); // Log error but continue - transient errors are expected warn!( transport_id = %transport_id, diff --git a/src/transport/udp/stats.rs b/src/transport/udp/stats.rs new file mode 100644 index 0000000..a85870c --- /dev/null +++ b/src/transport/udp/stats.rs @@ -0,0 +1,105 @@ +//! UDP transport statistics. + +use std::sync::atomic::{AtomicU64, Ordering}; + +use serde::Serialize; + +/// Statistics for a UDP transport instance. +/// +/// Uses atomic counters for lock-free updates from the receive loop +/// and send path concurrently. +pub struct UdpStats { + pub packets_sent: AtomicU64, + pub bytes_sent: AtomicU64, + pub packets_recv: AtomicU64, + pub bytes_recv: AtomicU64, + pub send_errors: AtomicU64, + pub recv_errors: AtomicU64, + pub mtu_exceeded: AtomicU64, + pub kernel_drops: AtomicU64, +} + +impl UdpStats { + /// Create a new stats instance with all counters at zero. + pub fn new() -> Self { + Self { + packets_sent: AtomicU64::new(0), + bytes_sent: AtomicU64::new(0), + packets_recv: AtomicU64::new(0), + bytes_recv: AtomicU64::new(0), + send_errors: AtomicU64::new(0), + recv_errors: AtomicU64::new(0), + mtu_exceeded: AtomicU64::new(0), + kernel_drops: AtomicU64::new(0), + } + } + + /// Record a successful send. + pub fn record_send(&self, bytes: usize) { + self.packets_sent.fetch_add(1, Ordering::Relaxed); + self.bytes_sent.fetch_add(bytes as u64, Ordering::Relaxed); + } + + /// Record a successful receive. + pub fn record_recv(&self, bytes: usize) { + self.packets_recv.fetch_add(1, Ordering::Relaxed); + self.bytes_recv.fetch_add(bytes as u64, Ordering::Relaxed); + } + + /// Record a send error. + pub fn record_send_error(&self) { + self.send_errors.fetch_add(1, Ordering::Relaxed); + } + + /// Record a receive error. + pub fn record_recv_error(&self) { + self.recv_errors.fetch_add(1, Ordering::Relaxed); + } + + /// Record an MTU exceeded rejection. + pub fn record_mtu_exceeded(&self) { + self.mtu_exceeded.fetch_add(1, Ordering::Relaxed); + } + + /// Update kernel drop count from SO_MEMINFO. + /// + /// Not yet wired up — requires `getsockopt(SO_MEMINFO)` on the raw fd + /// (via socket2 or libc) to read `SK_MEMINFO_DROPS`. Linux-only. + /// Until implemented, this counter will always be zero. + pub fn set_kernel_drops(&self, drops: u64) { + self.kernel_drops.store(drops, Ordering::Relaxed); + } + + /// Take a snapshot of all counters. + pub fn snapshot(&self) -> UdpStatsSnapshot { + UdpStatsSnapshot { + packets_sent: self.packets_sent.load(Ordering::Relaxed), + bytes_sent: self.bytes_sent.load(Ordering::Relaxed), + packets_recv: self.packets_recv.load(Ordering::Relaxed), + bytes_recv: self.bytes_recv.load(Ordering::Relaxed), + send_errors: self.send_errors.load(Ordering::Relaxed), + recv_errors: self.recv_errors.load(Ordering::Relaxed), + mtu_exceeded: self.mtu_exceeded.load(Ordering::Relaxed), + kernel_drops: self.kernel_drops.load(Ordering::Relaxed), + } + } +} + +impl Default for UdpStats { + fn default() -> Self { + Self::new() + } +} + +/// Point-in-time snapshot of UDP stats (non-atomic, copyable). +#[derive(Clone, Debug, Default, Serialize)] +pub struct UdpStatsSnapshot { + pub packets_sent: u64, + pub bytes_sent: u64, + pub packets_recv: u64, + pub bytes_recv: u64, + pub send_errors: u64, + pub recv_errors: u64, + pub mtu_exceeded: u64, + pub kernel_drops: u64, +} diff --git a/src/tree/tests.rs b/src/tree/tests.rs index 26b4635..2ebb4dc 100644 --- a/src/tree/tests.rs +++ b/src/tree/tests.rs @@ -498,7 +498,7 @@ fn test_evaluate_parent_rejects_loop_candidate() { let mut state = TreeState::new(my_node); let peer1 = make_node_addr(1); - let root = make_node_addr(0); + let _root = make_node_addr(0); // Peer 1's ancestry: [1, 5, 0] — contains us (node 5) state.update_peer( @@ -519,7 +519,7 @@ fn test_evaluate_parent_picks_loop_free_over_loopy() { let peer1 = make_node_addr(1); let peer2 = make_node_addr(2); - let root = make_node_addr(0); + let _root = make_node_addr(0); // Peer 1: depth 2, but ancestry contains us — loop state.update_peer( diff --git a/testing/sidecar/.env b/testing/sidecar/.env index dc8e0f0..40eb89d 100644 --- a/testing/sidecar/.env +++ b/testing/sidecar/.env @@ -7,6 +7,7 @@ FIPS_NSEC=e752b92aed3ac1595807f5d0eb5125589fbec0a2cfd3a2948d87ea076557deeb # Peer configuration (leave empty for standalone operation) -FIPS_PEER_NPUB= -FIPS_PEER_ADDR= -FIPS_PEER_ALIAS=peer +FIPS_PEER_NPUB=npub16xhnhwaxzu3w6dlf88eqnea52cqx9crdwhx4s807zd9nxmng3seqfc587p +FIPS_PEER_ADDR=217.77.8.91:2121 +FIPS_PEER_ALIAS=vps +FIPS_PEER_TRANSPORT=udp