mirror of
https://github.com/jmcorgan/fips.git
synced 2026-08-09 16:24:45 +00:00
Merge branch 'master' into next
Carries the Nym mixnet transport + single-container demo. Forward-merge adaptation: the Nym end-to-end SOCKS5 send/recv test was authored against master's IK FMP framing (msg1 wire size 114, version nibble 0). next uses the XX handshake and FMP v1, so the TCP read path rejected the test frame (UnknownVersion / HandshakeSizeMismatch), timing the test out. Adapted build_msg1_frame to next's framing (wire size 41, ver=1/phase=1 header byte). next-only; master keeps its IK form.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -64,6 +64,7 @@ use crate::peer::{ActivePeer, PeerConnection};
|
||||
use crate::protocol::NodeProfile;
|
||||
#[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;
|
||||
@@ -1007,6 +1008,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