Implement comprehensive node and transport statistics

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.
This commit is contained in:
Johnathan Corgan
2026-02-28 18:23:04 +00:00
parent 5c1cbb4c30
commit 71a5c68fa9
15 changed files with 864 additions and 49 deletions
+28 -5
View File
@@ -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<String> = 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<Value> = 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(),
})
}
+16 -1
View File
@@ -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)
+20
View File
@@ -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 {
+10
View File
@@ -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());
}
}
+9 -1
View File
@@ -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.
+19
View File
@@ -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.
+366
View File
@@ -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,
}
+44 -18
View File
@@ -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<NodeAddr, f64> = 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<NodeAddr, f64> = 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<NodeAddr, f64> = self.peers.iter()
.map(|(addr, peer)| (*addr, peer.link_cost()))
.collect();
+30
View File
@@ -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()
}
}
}
}
// ============================================================================
+41 -5
View File
@@ -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<JoinHandle<()>>,
/// Local listener address (after start, if bind_addr configured).
local_addr: Option<SocketAddr>,
/// Transport statistics.
stats: Arc<TcpStats>,
}
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<TcpStats> {
&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<TcpStats>,
) {
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<TcpStats>,
) {
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,
+137
View File
@@ -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,
}
@@ -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<JoinHandle<()>>,
/// Local bound address (after start).
local_addr: Option<SocketAddr>,
/// Transport statistics.
stats: Arc<UdpStats>,
}
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<UdpStats> {
&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<UdpStats>,
) {
// 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,
+105
View File
@@ -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,
}
+2 -2
View File
@@ -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(
+4 -3
View File
@@ -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