mirror of
https://github.com/jmcorgan/fips.git
synced 2026-08-09 00:04:54 +00:00
Add Nym mixnet transport and single-container demo
Add an outbound-only Nym mixnet transport that tunnels FMP peer links through a local nym-socks5-client SOCKS5 proxy into the Nym mixnet. It structurally mirrors the Tor SOCKS5 transport (connection pool, connect-on-send background promotion, FMP-v0 framing reused from TCP) with the onion, inbound-listener, and control-port machinery removed. Wires the transport through the full TransportHandle dispatch, NymConfig (standard transport-instance pattern), and node instantiation, and surfaces its counters in fipstop. Includes a mock SOCKS5 harness and unit coverage for the address-parsing paths. Also adds an isolated single-container example (examples/sidecar-nostr-mixnet-relay/) demonstrating FIPS peering across the mixnet end to end. No new crate dependencies: tokio_socks, socket2, and futures are already pulled in by the Tor transport.
This commit is contained in:
committed by
Johnathan Corgan
parent
fb8bb4fb97
commit
4e43cb81e9
@@ -481,6 +481,26 @@ fn draw_transport_detail(frame: &mut Frame, app: &App, area: Rect, t: &serde_jso
|
||||
&helpers::nested_u64(t, "stats", "connect_refused"),
|
||||
));
|
||||
}
|
||||
"nym" => {
|
||||
lines.push(helpers::kv_line(
|
||||
"MTU Exceeded",
|
||||
&helpers::nested_u64(t, "stats", "mtu_exceeded"),
|
||||
));
|
||||
lines.push(helpers::kv_line(
|
||||
"SOCKS5 Errors",
|
||||
&helpers::nested_u64(t, "stats", "socks5_errors"),
|
||||
));
|
||||
lines.push(Line::from(""));
|
||||
lines.push(helpers::section_header("Connections"));
|
||||
lines.push(helpers::kv_line(
|
||||
"Established",
|
||||
&helpers::nested_u64(t, "stats", "connections_established"),
|
||||
));
|
||||
lines.push(helpers::kv_line(
|
||||
"Timeouts",
|
||||
&helpers::nested_u64(t, "stats", "connect_timeouts"),
|
||||
));
|
||||
}
|
||||
"ethernet" => {
|
||||
lines.push(Line::from(""));
|
||||
lines.push(helpers::section_header("Beacons"));
|
||||
|
||||
+2
-2
@@ -39,8 +39,8 @@ pub use node::{
|
||||
};
|
||||
pub use peer::{ConnectPolicy, PeerAddress, PeerConfig};
|
||||
pub use transport::{
|
||||
BleConfig, DirectoryServiceConfig, EthernetConfig, TcpConfig, TorConfig, TransportInstances,
|
||||
TransportsConfig, UdpConfig,
|
||||
BleConfig, DirectoryServiceConfig, EthernetConfig, NymConfig, TcpConfig, TorConfig,
|
||||
TransportInstances, TransportsConfig, UdpConfig,
|
||||
};
|
||||
|
||||
/// Default config filename.
|
||||
|
||||
@@ -801,6 +801,77 @@ impl BleConfig {
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Nym Transport Configuration
|
||||
// ============================================================================
|
||||
|
||||
/// Default Nym SOCKS5 proxy address (nym-socks5-client).
|
||||
const DEFAULT_NYM_SOCKS5_ADDR: &str = "127.0.0.1:1080";
|
||||
|
||||
/// Default Nym connect timeout in milliseconds (300s — Nym mixnet
|
||||
/// SOCKS5 connections require multiple round-trips through 3 mix nodes
|
||||
/// with timing obfuscation, which can take several minutes).
|
||||
const DEFAULT_NYM_CONNECT_TIMEOUT_MS: u64 = 300_000;
|
||||
|
||||
/// Default Nym MTU (same as TCP).
|
||||
const DEFAULT_NYM_MTU: u16 = 1400;
|
||||
|
||||
/// Default Nym startup timeout in seconds (time to wait for
|
||||
/// nym-socks5-client to become ready before giving up).
|
||||
const DEFAULT_NYM_STARTUP_TIMEOUT_SECS: u64 = 120;
|
||||
|
||||
/// Nym transport instance configuration.
|
||||
///
|
||||
/// Outbound-only connections through a nym-socks5-client SOCKS5 proxy.
|
||||
/// The nym-socks5-client must be running separately (e.g., as a sidecar
|
||||
/// process or container).
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct NymConfig {
|
||||
/// SOCKS5 proxy address (host:port). Defaults to "127.0.0.1:1080".
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub socks5_addr: Option<String>,
|
||||
|
||||
/// Outbound connect timeout in milliseconds. Defaults to 300000 (300s).
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub connect_timeout_ms: Option<u64>,
|
||||
|
||||
/// Default MTU for Nym connections. Defaults to 1400.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub mtu: Option<u16>,
|
||||
|
||||
/// Seconds to wait for nym-socks5-client to become ready at startup.
|
||||
/// Defaults to 120.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub startup_timeout_secs: Option<u64>,
|
||||
}
|
||||
|
||||
impl NymConfig {
|
||||
/// Get the SOCKS5 proxy address. Default: "127.0.0.1:1080".
|
||||
pub fn socks5_addr(&self) -> &str {
|
||||
self.socks5_addr
|
||||
.as_deref()
|
||||
.unwrap_or(DEFAULT_NYM_SOCKS5_ADDR)
|
||||
}
|
||||
|
||||
/// Get the connect timeout in milliseconds. Default: 300000.
|
||||
pub fn connect_timeout_ms(&self) -> u64 {
|
||||
self.connect_timeout_ms
|
||||
.unwrap_or(DEFAULT_NYM_CONNECT_TIMEOUT_MS)
|
||||
}
|
||||
|
||||
/// Get the default MTU. Default: 1400.
|
||||
pub fn mtu(&self) -> u16 {
|
||||
self.mtu.unwrap_or(DEFAULT_NYM_MTU)
|
||||
}
|
||||
|
||||
/// Get the startup timeout in seconds. Default: 120.
|
||||
pub fn startup_timeout_secs(&self) -> u64 {
|
||||
self.startup_timeout_secs
|
||||
.unwrap_or(DEFAULT_NYM_STARTUP_TIMEOUT_SECS)
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// TransportsConfig
|
||||
// ============================================================================
|
||||
@@ -827,6 +898,10 @@ pub struct TransportsConfig {
|
||||
#[serde(default, skip_serializing_if = "is_transport_empty")]
|
||||
pub tor: TransportInstances<TorConfig>,
|
||||
|
||||
/// Nym transport instances.
|
||||
#[serde(default, skip_serializing_if = "is_transport_empty")]
|
||||
pub nym: TransportInstances<NymConfig>,
|
||||
|
||||
/// BLE transport instances.
|
||||
#[serde(default, skip_serializing_if = "is_transport_empty")]
|
||||
pub ble: TransportInstances<BleConfig>,
|
||||
@@ -844,6 +919,7 @@ impl TransportsConfig {
|
||||
&& self.ethernet.is_empty()
|
||||
&& self.tcp.is_empty()
|
||||
&& self.tor.is_empty()
|
||||
&& self.nym.is_empty()
|
||||
&& self.ble.is_empty()
|
||||
}
|
||||
|
||||
@@ -863,6 +939,9 @@ impl TransportsConfig {
|
||||
if !other.tor.is_empty() {
|
||||
self.tor = other.tor;
|
||||
}
|
||||
if !other.nym.is_empty() {
|
||||
self.nym = other.nym;
|
||||
}
|
||||
if !other.ble.is_empty() {
|
||||
self.ble = other.ble;
|
||||
}
|
||||
|
||||
+1
-1
@@ -30,7 +30,7 @@ pub use identity::{
|
||||
};
|
||||
|
||||
// Re-export config types
|
||||
pub use config::{Config, ConfigError, IdentityConfig, TorConfig, UdpConfig};
|
||||
pub use config::{Config, ConfigError, IdentityConfig, NymConfig, TorConfig, UdpConfig};
|
||||
pub use upper::config::{DnsConfig, TunConfig};
|
||||
|
||||
// Re-export discovery types
|
||||
|
||||
@@ -50,6 +50,7 @@ use crate::node::session::SessionEntry;
|
||||
use crate::peer::{ActivePeer, PeerConnection};
|
||||
#[cfg(unix)]
|
||||
use crate::transport::ethernet::EthernetTransport;
|
||||
use crate::transport::nym::NymTransport;
|
||||
use crate::transport::tcp::TcpTransport;
|
||||
use crate::transport::tor::TorTransport;
|
||||
use crate::transport::udp::UdpTransport;
|
||||
@@ -982,6 +983,21 @@ impl Node {
|
||||
transports.push(TransportHandle::Tor(tor));
|
||||
}
|
||||
|
||||
// Create Nym transport instances
|
||||
let nym_instances: Vec<_> = self
|
||||
.config()
|
||||
.transports
|
||||
.nym
|
||||
.iter()
|
||||
.map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
|
||||
.collect();
|
||||
|
||||
for (name, nym_config) in nym_instances {
|
||||
let transport_id = self.allocate_transport_id();
|
||||
let nym = NymTransport::new(transport_id, name, nym_config, packet_tx.clone());
|
||||
transports.push(TransportHandle::Nym(nym));
|
||||
}
|
||||
|
||||
// Create BLE transport instances
|
||||
#[cfg(bluer_available)]
|
||||
{
|
||||
|
||||
+34
-2
@@ -6,6 +6,7 @@
|
||||
|
||||
#[cfg(test)]
|
||||
pub mod loopback;
|
||||
pub mod nym;
|
||||
pub mod tcp;
|
||||
pub mod tor;
|
||||
pub mod udp;
|
||||
@@ -22,6 +23,7 @@ use ble::DefaultBleTransport;
|
||||
use ethernet::EthernetTransport;
|
||||
#[cfg(test)]
|
||||
use loopback::LoopbackTransport;
|
||||
use nym::NymTransport;
|
||||
use secp256k1::XOnlyPublicKey;
|
||||
use std::fmt;
|
||||
use std::net::SocketAddr;
|
||||
@@ -259,6 +261,13 @@ impl TransportType {
|
||||
reliable: true, // in-process channel delivery is lossless
|
||||
};
|
||||
|
||||
/// Nym mixnet transport (via SOCKS5).
|
||||
pub const NYM: TransportType = TransportType {
|
||||
name: "nym",
|
||||
connection_oriented: true,
|
||||
reliable: true,
|
||||
};
|
||||
|
||||
/// Check if the transport is connectionless.
|
||||
pub fn is_connectionless(&self) -> bool {
|
||||
!self.connection_oriented
|
||||
@@ -891,6 +900,8 @@ pub enum TransportHandle {
|
||||
Tcp(TcpTransport),
|
||||
/// Tor transport (via SOCKS5).
|
||||
Tor(TorTransport),
|
||||
/// Nym mixnet transport (via SOCKS5).
|
||||
Nym(NymTransport),
|
||||
/// BLE L2CAP transport.
|
||||
#[cfg(target_os = "linux")]
|
||||
Ble(DefaultBleTransport),
|
||||
@@ -908,6 +919,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.start_async().await,
|
||||
TransportHandle::Tcp(t) => t.start_async().await,
|
||||
TransportHandle::Tor(t) => t.start_async().await,
|
||||
TransportHandle::Nym(t) => t.start_async().await,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.start_async().await,
|
||||
#[cfg(test)]
|
||||
@@ -923,6 +935,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.stop_async().await,
|
||||
TransportHandle::Tcp(t) => t.stop_async().await,
|
||||
TransportHandle::Tor(t) => t.stop_async().await,
|
||||
TransportHandle::Nym(t) => t.stop_async().await,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.stop_async().await,
|
||||
#[cfg(test)]
|
||||
@@ -938,6 +951,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.send_async(addr, data).await,
|
||||
TransportHandle::Tcp(t) => t.send_async(addr, data).await,
|
||||
TransportHandle::Tor(t) => t.send_async(addr, data).await,
|
||||
TransportHandle::Nym(t) => t.send_async(addr, data).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.send_async(addr, data).await,
|
||||
#[cfg(test)]
|
||||
@@ -953,6 +967,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.transport_id(),
|
||||
TransportHandle::Tcp(t) => t.transport_id(),
|
||||
TransportHandle::Tor(t) => t.transport_id(),
|
||||
TransportHandle::Nym(t) => t.transport_id(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.transport_id(),
|
||||
#[cfg(test)]
|
||||
@@ -968,6 +983,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.name(),
|
||||
TransportHandle::Tcp(t) => t.name(),
|
||||
TransportHandle::Tor(t) => t.name(),
|
||||
TransportHandle::Nym(t) => t.name(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.name(),
|
||||
#[cfg(test)]
|
||||
@@ -983,6 +999,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.transport_type(),
|
||||
TransportHandle::Tcp(t) => t.transport_type(),
|
||||
TransportHandle::Tor(t) => t.transport_type(),
|
||||
TransportHandle::Nym(t) => t.transport_type(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.transport_type(),
|
||||
#[cfg(test)]
|
||||
@@ -998,6 +1015,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.state(),
|
||||
TransportHandle::Tcp(t) => t.state(),
|
||||
TransportHandle::Tor(t) => t.state(),
|
||||
TransportHandle::Nym(t) => t.state(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.state(),
|
||||
#[cfg(test)]
|
||||
@@ -1013,6 +1031,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.mtu(),
|
||||
TransportHandle::Tcp(t) => t.mtu(),
|
||||
TransportHandle::Tor(t) => t.mtu(),
|
||||
TransportHandle::Nym(t) => t.mtu(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.mtu(),
|
||||
#[cfg(test)]
|
||||
@@ -1031,6 +1050,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.link_mtu(addr),
|
||||
TransportHandle::Tcp(t) => t.link_mtu(addr),
|
||||
TransportHandle::Tor(t) => t.link_mtu(addr),
|
||||
TransportHandle::Nym(t) => t.link_mtu(addr),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.link_mtu(addr),
|
||||
#[cfg(test)]
|
||||
@@ -1046,6 +1066,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(_) => None,
|
||||
TransportHandle::Tcp(t) => t.local_addr(),
|
||||
TransportHandle::Tor(_) => None,
|
||||
TransportHandle::Nym(_) => None,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(_) => None,
|
||||
#[cfg(test)]
|
||||
@@ -1061,6 +1082,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => Some(t.interface_name()),
|
||||
TransportHandle::Tcp(_) => None,
|
||||
TransportHandle::Tor(_) => None,
|
||||
TransportHandle::Nym(_) => None,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(_) => None,
|
||||
#[cfg(test)]
|
||||
@@ -1100,6 +1122,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.discover(),
|
||||
TransportHandle::Tcp(t) => t.discover(),
|
||||
TransportHandle::Tor(t) => t.discover(),
|
||||
TransportHandle::Nym(t) => t.discover(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.discover(),
|
||||
#[cfg(test)]
|
||||
@@ -1115,6 +1138,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.auto_connect(),
|
||||
TransportHandle::Tcp(t) => t.auto_connect(),
|
||||
TransportHandle::Tor(t) => t.auto_connect(),
|
||||
TransportHandle::Nym(t) => t.auto_connect(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.auto_connect(),
|
||||
#[cfg(test)]
|
||||
@@ -1130,6 +1154,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.accept_connections(),
|
||||
TransportHandle::Tcp(t) => t.accept_connections(),
|
||||
TransportHandle::Tor(t) => t.accept_connections(),
|
||||
TransportHandle::Nym(t) => t.accept_connections(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.accept_connections(),
|
||||
#[cfg(test)]
|
||||
@@ -1139,7 +1164,7 @@ impl TransportHandle {
|
||||
|
||||
/// Initiate a non-blocking connection to a remote address.
|
||||
///
|
||||
/// For connection-oriented transports (TCP, Tor), spawns a background
|
||||
/// For connection-oriented transports (TCP, Tor, Nym), spawns a background
|
||||
/// task to establish the connection. For connectionless transports
|
||||
/// (UDP, Ethernet), this is a no-op that returns Ok immediately.
|
||||
///
|
||||
@@ -1151,6 +1176,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(_) => Ok(()), // connectionless
|
||||
TransportHandle::Tcp(t) => t.connect_async(addr).await,
|
||||
TransportHandle::Tor(t) => t.connect_async(addr).await,
|
||||
TransportHandle::Nym(t) => t.connect_async(addr).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.connect_async(addr).await,
|
||||
#[cfg(test)]
|
||||
@@ -1170,6 +1196,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(_) => ConnectionState::Connected,
|
||||
TransportHandle::Tcp(t) => t.connection_state_sync(addr),
|
||||
TransportHandle::Tor(t) => t.connection_state_sync(addr),
|
||||
TransportHandle::Nym(t) => t.connection_state_sync(addr),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.connection_state_sync(addr),
|
||||
#[cfg(test)]
|
||||
@@ -1179,7 +1206,7 @@ impl TransportHandle {
|
||||
|
||||
/// Close a specific connection on this transport.
|
||||
///
|
||||
/// No-op for connectionless transports. For TCP/Tor, removes the
|
||||
/// No-op for connectionless transports. For TCP/Tor/Nym, removes the
|
||||
/// connection from the pool and drops the stream.
|
||||
pub async fn close_connection(&self, addr: &TransportAddr) {
|
||||
match self {
|
||||
@@ -1188,6 +1215,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(t) => t.close_connection(addr),
|
||||
TransportHandle::Tcp(t) => t.close_connection_async(addr).await,
|
||||
TransportHandle::Tor(t) => t.close_connection_async(addr).await,
|
||||
TransportHandle::Nym(t) => t.close_connection_async(addr).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => t.close_connection_async(addr).await,
|
||||
#[cfg(test)]
|
||||
@@ -1212,6 +1240,7 @@ impl TransportHandle {
|
||||
TransportHandle::Ethernet(_) => TransportCongestion::default(),
|
||||
TransportHandle::Tcp(_) => TransportCongestion::default(),
|
||||
TransportHandle::Tor(_) => TransportCongestion::default(),
|
||||
TransportHandle::Nym(_) => TransportCongestion::default(),
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(_) => TransportCongestion::default(),
|
||||
#[cfg(test)]
|
||||
@@ -1249,6 +1278,9 @@ impl TransportHandle {
|
||||
TransportHandle::Tor(t) => {
|
||||
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
|
||||
}
|
||||
TransportHandle::Nym(t) => {
|
||||
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
|
||||
}
|
||||
#[cfg(target_os = "linux")]
|
||||
TransportHandle::Ble(t) => {
|
||||
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
//! Mock SOCKS5 server for testing the Nym transport's connect path.
|
||||
//!
|
||||
//! A copy of the Tor transport's mock, implementing just enough of the
|
||||
//! SOCKS5 protocol (RFC 1928) to support the no-auth (and, defensively,
|
||||
//! username/password) CONNECT flow, then proxying bytes bidirectionally to a
|
||||
//! fixed target.
|
||||
//!
|
||||
//! Difference from the Tor mock: this one accepts connections in a loop and
|
||||
//! handles each on its own task. `NymTransport::start_async` first probes the
|
||||
//! proxy port for readiness (opening and immediately dropping a connection);
|
||||
//! looping lets the mock shrug that probe off — its handler returns on the
|
||||
//! short first read — and still serve the real data connection that follows.
|
||||
|
||||
use std::net::SocketAddr;
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
use tokio::net::TcpListener;
|
||||
use tokio::task::JoinHandle;
|
||||
|
||||
/// SOCKS5 protocol constants.
|
||||
const SOCKS_VERSION: u8 = 0x05;
|
||||
const AUTH_NONE: u8 = 0x00;
|
||||
const AUTH_PASSWORD: u8 = 0x02;
|
||||
const CMD_CONNECT: u8 = 0x01;
|
||||
const ATYP_IPV4: u8 = 0x01;
|
||||
const ATYP_DOMAIN: u8 = 0x03;
|
||||
const REP_SUCCESS: u8 = 0x00;
|
||||
|
||||
/// Username/password auth sub-negotiation version (RFC 1929).
|
||||
const AUTH_SUBNEG_VERSION: u8 = 0x01;
|
||||
const AUTH_SUBNEG_SUCCESS: u8 = 0x00;
|
||||
|
||||
/// A minimal mock SOCKS5 proxy server for testing.
|
||||
///
|
||||
/// Accepts connections in a loop, performs the SOCKS5 handshake (supporting
|
||||
/// both no-auth and username/password auth), then connects to a fixed target
|
||||
/// address and proxies bytes bidirectionally.
|
||||
pub struct MockSocks5Server {
|
||||
/// Address the mock proxy is listening on.
|
||||
addr: SocketAddr,
|
||||
/// The real target address to connect to (ignores SOCKS5 requested target).
|
||||
target_addr: SocketAddr,
|
||||
/// Listener handle.
|
||||
listener: Option<TcpListener>,
|
||||
}
|
||||
|
||||
impl MockSocks5Server {
|
||||
/// Create a new mock SOCKS5 server that forwards to the given target.
|
||||
///
|
||||
/// Binds to `127.0.0.1:0` (OS-assigned port).
|
||||
pub async fn new(target_addr: SocketAddr) -> std::io::Result<Self> {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await?;
|
||||
let addr = listener.local_addr()?;
|
||||
Ok(Self {
|
||||
addr,
|
||||
target_addr,
|
||||
listener: Some(listener),
|
||||
})
|
||||
}
|
||||
|
||||
/// Get the proxy's listen address (for `NymConfig.socks5_addr`).
|
||||
pub fn addr(&self) -> SocketAddr {
|
||||
self.addr
|
||||
}
|
||||
|
||||
/// Run the proxy, accepting connections in a loop and proxying each.
|
||||
///
|
||||
/// Returns a JoinHandle for the accept loop.
|
||||
pub fn spawn(mut self) -> JoinHandle<()> {
|
||||
let listener = self.listener.take().expect("listener already consumed");
|
||||
let target_addr = self.target_addr;
|
||||
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
let (client, _) = match listener.accept().await {
|
||||
Ok(c) => c,
|
||||
Err(_) => break,
|
||||
};
|
||||
// Handle each connection independently so the readiness probe
|
||||
// (which opens and drops a connection) cannot block the real
|
||||
// data connection behind it.
|
||||
tokio::spawn(handle_conn(client, target_addr));
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// Handle a single accepted connection: SOCKS5 handshake then byte proxy.
|
||||
async fn handle_conn(mut client: tokio::net::TcpStream, target_addr: SocketAddr) {
|
||||
// === Method negotiation ===
|
||||
// Client sends: [version, nmethods, methods...]
|
||||
let mut ver_nmethods = [0u8; 2];
|
||||
if client.read_exact(&mut ver_nmethods).await.is_err() {
|
||||
// Readiness probe (or any early close) — nothing to serve.
|
||||
return;
|
||||
}
|
||||
assert_eq!(ver_nmethods[0], SOCKS_VERSION, "expected SOCKS5");
|
||||
let nmethods = ver_nmethods[1] as usize;
|
||||
|
||||
let mut methods = vec![0u8; nmethods];
|
||||
client.read_exact(&mut methods).await.expect("read methods");
|
||||
|
||||
// Prefer username/password auth if offered, fall back to no-auth.
|
||||
let selected = if methods.contains(&AUTH_PASSWORD) {
|
||||
AUTH_PASSWORD
|
||||
} else if methods.contains(&AUTH_NONE) {
|
||||
AUTH_NONE
|
||||
} else {
|
||||
panic!("no supported auth method offered");
|
||||
};
|
||||
|
||||
// Reply: [version, selected_method]
|
||||
client
|
||||
.write_all(&[SOCKS_VERSION, selected])
|
||||
.await
|
||||
.expect("write method reply");
|
||||
|
||||
// === Username/password sub-negotiation (RFC 1929) ===
|
||||
if selected == AUTH_PASSWORD {
|
||||
// Client sends: [ver(1), ulen(1), uname(ulen), plen(1), passwd(plen)]
|
||||
let mut subneg_header = [0u8; 2];
|
||||
client
|
||||
.read_exact(&mut subneg_header)
|
||||
.await
|
||||
.expect("read subneg header");
|
||||
assert_eq!(
|
||||
subneg_header[0], AUTH_SUBNEG_VERSION,
|
||||
"expected auth subneg v1"
|
||||
);
|
||||
|
||||
let ulen = subneg_header[1] as usize;
|
||||
let mut uname = vec![0u8; ulen];
|
||||
client.read_exact(&mut uname).await.expect("read username");
|
||||
|
||||
let mut plen_buf = [0u8; 1];
|
||||
client.read_exact(&mut plen_buf).await.expect("read plen");
|
||||
let plen = plen_buf[0] as usize;
|
||||
let mut passwd = vec![0u8; plen];
|
||||
client.read_exact(&mut passwd).await.expect("read password");
|
||||
|
||||
client
|
||||
.write_all(&[AUTH_SUBNEG_VERSION, AUTH_SUBNEG_SUCCESS])
|
||||
.await
|
||||
.expect("write subneg reply");
|
||||
}
|
||||
|
||||
// === Connect request ===
|
||||
// Client sends: [version, cmd, rsv, atyp, addr..., port]
|
||||
let mut header = [0u8; 4];
|
||||
client
|
||||
.read_exact(&mut header)
|
||||
.await
|
||||
.expect("read connect header");
|
||||
assert_eq!(header[0], SOCKS_VERSION);
|
||||
assert_eq!(header[1], CMD_CONNECT);
|
||||
|
||||
// Read and skip the address (we connect to target_addr regardless).
|
||||
match header[3] {
|
||||
ATYP_IPV4 => {
|
||||
let mut addr_port = [0u8; 6]; // 4 IP + 2 port
|
||||
client
|
||||
.read_exact(&mut addr_port)
|
||||
.await
|
||||
.expect("read IPv4 addr");
|
||||
}
|
||||
ATYP_DOMAIN => {
|
||||
let mut len_buf = [0u8; 1];
|
||||
client
|
||||
.read_exact(&mut len_buf)
|
||||
.await
|
||||
.expect("read domain len");
|
||||
let domain_len = len_buf[0] as usize;
|
||||
let mut domain_port = vec![0u8; domain_len + 2]; // domain + 2 port
|
||||
client
|
||||
.read_exact(&mut domain_port)
|
||||
.await
|
||||
.expect("read domain addr");
|
||||
}
|
||||
other => panic!("unsupported ATYP: {}", other),
|
||||
}
|
||||
|
||||
// Connect to the real target.
|
||||
let mut target = tokio::net::TcpStream::connect(target_addr)
|
||||
.await
|
||||
.expect("connect to target");
|
||||
|
||||
// Reply: success, bind addr = 0.0.0.0:0
|
||||
let reply = [
|
||||
SOCKS_VERSION,
|
||||
REP_SUCCESS,
|
||||
0x00, // RSV
|
||||
ATYP_IPV4,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
0, // bind addr
|
||||
0,
|
||||
0, // bind port
|
||||
];
|
||||
client.write_all(&reply).await.expect("write connect reply");
|
||||
|
||||
// Proxy bytes bidirectionally.
|
||||
let _ = tokio::io::copy_bidirectional(&mut client, &mut target).await;
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,119 @@
|
||||
//! Nym transport statistics.
|
||||
|
||||
use portable_atomic::{AtomicU64, Ordering};
|
||||
|
||||
use serde::Serialize;
|
||||
|
||||
/// Statistics for a Nym transport instance.
|
||||
///
|
||||
/// Uses atomic counters for lock-free updates from per-connection
|
||||
/// receive loops and the send path concurrently.
|
||||
pub struct NymStats {
|
||||
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 connect_timeouts: AtomicU64,
|
||||
pub socks5_errors: AtomicU64,
|
||||
}
|
||||
|
||||
impl NymStats {
|
||||
/// 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),
|
||||
connect_timeouts: AtomicU64::new(0),
|
||||
socks5_errors: 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 connect timeout.
|
||||
pub fn record_connect_timeout(&self) {
|
||||
self.connect_timeouts.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Record a SOCKS5 protocol error.
|
||||
pub fn record_socks5_error(&self) {
|
||||
self.socks5_errors.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Take a snapshot of all counters.
|
||||
pub fn snapshot(&self) -> NymStatsSnapshot {
|
||||
NymStatsSnapshot {
|
||||
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),
|
||||
connect_timeouts: self.connect_timeouts.load(Ordering::Relaxed),
|
||||
socks5_errors: self.socks5_errors.load(Ordering::Relaxed),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for NymStats {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
/// Point-in-time snapshot of Nym stats (non-atomic, copyable).
|
||||
#[derive(Clone, Debug, Default, Serialize)]
|
||||
pub struct NymStatsSnapshot {
|
||||
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 connect_timeouts: u64,
|
||||
pub socks5_errors: u64,
|
||||
}
|
||||
Reference in New Issue
Block a user