Merge branch 'master' into next

This commit is contained in:
Johnathan Corgan
2026-07-11 05:05:22 +00:00
9 changed files with 100 additions and 53 deletions
+3 -1
View File
@@ -2,6 +2,8 @@
use portable_atomic::{AtomicU64, Ordering};
use serde::Serialize;
/// Statistics for an Ethernet transport instance.
///
/// Uses atomic counters for lock-free updates from the receive loop
@@ -92,7 +94,7 @@ impl Default for EthernetStats {
}
/// Point-in-time snapshot of Ethernet stats (non-atomic, copyable).
#[derive(Clone, Debug, Default)]
#[derive(Clone, Debug, Default, Serialize)]
pub struct EthernetStatsSnapshot {
pub frames_sent: u64,
pub frames_recv: u64,
@@ -3,8 +3,8 @@
//! Recovers FIPS packet boundaries from a TCP byte stream using the
//! existing 4-byte FMP common prefix `[ver+phase:1][flags:1][payload_len:2 LE]`.
//!
//! This module is deliberately separate from the TCP transport so it can
//! be reused by the future Tor transport.
//! This module is deliberately separate from any single transport so it can
//! be shared by the stream-oriented transports (TCP, Tor, Nym).
use tokio::io::{AsyncRead, AsyncReadExt};
+6 -13
View File
@@ -34,6 +34,11 @@ use tor::TorTransport;
use tor::control::TorMonitoringInfo;
use udp::UdpTransport;
pub(crate) mod framing;
mod stats_common;
pub(crate) use stats_common::PoolCounters;
mod types;
pub use types::*;
@@ -1022,19 +1027,7 @@ impl TransportHandle {
}
#[cfg(unix)]
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,
})
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
}
TransportHandle::Tcp(t) => {
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
+2 -2
View File
@@ -9,7 +9,7 @@
//!
//! Outbound-only: connects to remote TCP peers through the local
//! nym-socks5-client SOCKS5 proxy. Like the Tor transport, reuses FMP
//! stream framing from `tcp::stream` and follows the same connection
//! stream framing from `transport::framing` and follows the same connection
//! pool pattern. No inbound service is supported.
pub mod stats;
@@ -22,7 +22,7 @@ use super::{
TransportError, TransportId, TransportState, TransportType,
};
use crate::config::NymConfig;
use crate::transport::tcp::stream::read_fmp_packet;
use crate::transport::framing::read_fmp_packet;
use stats::NymStats;
use futures::FutureExt;
+56
View File
@@ -0,0 +1,56 @@
//! Statistics helpers shared across transports.
//!
//! Counter logic that is identical across multiple transports lives here
//! rather than being copied into each transport's `stats.rs`.
use portable_atomic::{AtomicU64, Ordering};
/// Inbound/outbound connection-pool occupancy counters.
///
/// The connection-oriented transports (TCP, Tor) each track how many inbound
/// and outbound connections are currently held in their pool; the inbound
/// count drives the `max_inbound_connections` admission gate. The counter
/// logic is identical, so it is defined once here and embedded in each
/// transport's stats struct.
#[derive(Default)]
pub struct PoolCounters {
inbound: AtomicU64,
outbound: AtomicU64,
}
impl PoolCounters {
/// Create counters at zero.
pub fn new() -> Self {
Self::default()
}
/// Increment the inbound pool count (called on accept).
pub fn record_inbound_added(&self) {
self.inbound.fetch_add(1, Ordering::Relaxed);
}
/// Decrement the inbound pool count (called on inbound receive-loop exit).
pub fn record_inbound_removed(&self) {
self.inbound.fetch_sub(1, Ordering::Relaxed);
}
/// Increment the outbound pool count (called on connect-on-send / promote).
pub fn record_outbound_added(&self) {
self.outbound.fetch_add(1, Ordering::Relaxed);
}
/// Decrement the outbound pool count (called on outbound receive-loop exit).
pub fn record_outbound_removed(&self) {
self.outbound.fetch_sub(1, Ordering::Relaxed);
}
/// Load the current inbound pool count for the admission gate.
pub fn inbound_count(&self) -> u64 {
self.inbound.load(Ordering::Relaxed)
}
/// Load the current outbound pool count.
pub fn outbound_count(&self) -> u64 {
self.outbound.load(Ordering::Relaxed)
}
}
+1 -2
View File
@@ -23,7 +23,6 @@
//! TCP stream and the receiver uses phase-dependent size computation.
pub mod stats;
pub mod stream;
use super::resolve_socket_addr;
use super::{
@@ -31,8 +30,8 @@ use super::{
TransportError, TransportId, TransportState, TransportType,
};
use crate::config::TcpConfig;
use crate::transport::framing::read_fmp_packet;
use stats::TcpStats;
use stream::read_fmp_packet;
use futures::FutureExt;
use socket2::TcpKeepalive;
+13 -15
View File
@@ -4,6 +4,8 @@ use portable_atomic::{AtomicU64, Ordering};
use serde::Serialize;
use crate::transport::PoolCounters;
/// Statistics for a TCP transport instance.
///
/// Uses atomic counters for lock-free updates from per-connection
@@ -21,12 +23,9 @@ pub struct TcpStats {
pub connections_rejected: AtomicU64,
pub connect_timeouts: AtomicU64,
pub connect_refused: AtomicU64,
/// Current number of inbound (accepted) connections held in the pool.
/// Drives the `max_inbound_connections` admission check.
pub pool_inbound: AtomicU64,
/// Current number of outbound (connect-on-send / promoted) connections
/// held in the pool. Independent of the inbound cap.
pub pool_outbound: AtomicU64,
/// Inbound/outbound connection-pool occupancy. Inbound drives the
/// `max_inbound_connections` admission check.
pub pool: PoolCounters,
}
impl TcpStats {
@@ -45,8 +44,7 @@ impl TcpStats {
connections_rejected: AtomicU64::new(0),
connect_timeouts: AtomicU64::new(0),
connect_refused: AtomicU64::new(0),
pool_inbound: AtomicU64::new(0),
pool_outbound: AtomicU64::new(0),
pool: PoolCounters::new(),
}
}
@@ -104,27 +102,27 @@ impl TcpStats {
/// Increment the inbound pool count (called on accept).
pub fn record_pool_inbound_added(&self) {
self.pool_inbound.fetch_add(1, Ordering::Relaxed);
self.pool.record_inbound_added();
}
/// Decrement the inbound pool count (called on inbound receive-loop exit).
pub fn record_pool_inbound_removed(&self) {
self.pool_inbound.fetch_sub(1, Ordering::Relaxed);
self.pool.record_inbound_removed();
}
/// Increment the outbound pool count (called on connect-on-send / promote).
pub fn record_pool_outbound_added(&self) {
self.pool_outbound.fetch_add(1, Ordering::Relaxed);
self.pool.record_outbound_added();
}
/// Decrement the outbound pool count (called on outbound receive-loop exit).
pub fn record_pool_outbound_removed(&self) {
self.pool_outbound.fetch_sub(1, Ordering::Relaxed);
self.pool.record_outbound_removed();
}
/// Load the current inbound pool count for the admission gate.
pub fn pool_inbound_count(&self) -> u64 {
self.pool_inbound.load(Ordering::Relaxed)
self.pool.inbound_count()
}
/// Take a snapshot of all counters.
@@ -142,8 +140,8 @@ impl TcpStats {
connections_rejected: self.connections_rejected.load(Ordering::Relaxed),
connect_timeouts: self.connect_timeouts.load(Ordering::Relaxed),
connect_refused: self.connect_refused.load(Ordering::Relaxed),
pool_inbound: self.pool_inbound.load(Ordering::Relaxed),
pool_outbound: self.pool_outbound.load(Ordering::Relaxed),
pool_inbound: self.pool.inbound_count(),
pool_outbound: self.pool.outbound_count(),
}
}
}
+3 -3
View File
@@ -14,8 +14,8 @@
//! ## Architecture
//!
//! Like TCP, each peer has its own connection. The transport reuses FMP
//! stream framing from `tcp::stream` and follows the same connection pool
//! pattern as the TCP transport. Inbound connections arrive via a local
//! stream framing from `transport::framing` and follows the same connection
//! pool pattern as the TCP transport. Inbound connections arrive via a local
//! TCP listener that the Tor daemon forwards onion service traffic to.
pub mod control;
@@ -31,7 +31,7 @@ use super::{
TransportError, TransportId, TransportState, TransportType,
};
use crate::config::TorConfig;
use crate::transport::tcp::stream::read_fmp_packet;
use crate::transport::framing::read_fmp_packet;
use control::{ControlAuth, TorControlClient, TorMonitoringInfo};
use stats::TorStats;
+14 -15
View File
@@ -4,6 +4,8 @@ use portable_atomic::{AtomicU64, Ordering};
use serde::Serialize;
use crate::transport::PoolCounters;
/// Statistics for a Tor transport instance.
///
/// Uses atomic counters for lock-free updates from per-connection
@@ -23,12 +25,10 @@ pub struct TorStats {
pub connections_accepted: AtomicU64,
pub connections_rejected: AtomicU64,
pub control_errors: AtomicU64,
/// Current number of inbound (accepted via onion service) connections
/// held in the pool. Drives the `max_inbound_connections` admission check.
pub pool_inbound: AtomicU64,
/// Current number of outbound (SOCKS5 connect) connections held in
/// the pool. Independent of the inbound cap.
pub pool_outbound: AtomicU64,
/// Inbound (accepted via onion service) / outbound (SOCKS5 connect)
/// connection-pool occupancy. Inbound drives the
/// `max_inbound_connections` admission check.
pub pool: PoolCounters,
}
impl TorStats {
@@ -49,8 +49,7 @@ impl TorStats {
connections_accepted: AtomicU64::new(0),
connections_rejected: AtomicU64::new(0),
control_errors: AtomicU64::new(0),
pool_inbound: AtomicU64::new(0),
pool_outbound: AtomicU64::new(0),
pool: PoolCounters::new(),
}
}
@@ -118,27 +117,27 @@ impl TorStats {
/// Increment the inbound pool count (called on accept).
pub fn record_pool_inbound_added(&self) {
self.pool_inbound.fetch_add(1, Ordering::Relaxed);
self.pool.record_inbound_added();
}
/// Decrement the inbound pool count (called on inbound receive-loop exit).
pub fn record_pool_inbound_removed(&self) {
self.pool_inbound.fetch_sub(1, Ordering::Relaxed);
self.pool.record_inbound_removed();
}
/// Increment the outbound pool count (called on SOCKS5-connect promote).
pub fn record_pool_outbound_added(&self) {
self.pool_outbound.fetch_add(1, Ordering::Relaxed);
self.pool.record_outbound_added();
}
/// Decrement the outbound pool count (called on outbound receive-loop exit).
pub fn record_pool_outbound_removed(&self) {
self.pool_outbound.fetch_sub(1, Ordering::Relaxed);
self.pool.record_outbound_removed();
}
/// Load the current inbound pool count for the admission gate.
pub fn pool_inbound_count(&self) -> u64 {
self.pool_inbound.load(Ordering::Relaxed)
self.pool.inbound_count()
}
/// Take a snapshot of all counters.
@@ -158,8 +157,8 @@ impl TorStats {
connections_accepted: self.connections_accepted.load(Ordering::Relaxed),
connections_rejected: self.connections_rejected.load(Ordering::Relaxed),
control_errors: self.control_errors.load(Ordering::Relaxed),
pool_inbound: self.pool_inbound.load(Ordering::Relaxed),
pool_outbound: self.pool_outbound.load(Ordering::Relaxed),
pool_inbound: self.pool.inbound_count(),
pool_outbound: self.pool.outbound_count(),
}
}
}