diff --git a/Cargo.lock b/Cargo.lock index 32f4505cdd..da284c6ce4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3731,6 +3731,7 @@ dependencies = [ "number_prefix", "portable-atomic", "unicode-width 0.2.0", + "vt100", "web-time", ] @@ -4911,6 +4912,7 @@ dependencies = [ "schemars", "serde", "serde_json", + "tracing", "tracing-opentelemetry", "tracing-subscriber", "tracing-tree", @@ -6108,22 +6110,25 @@ version = "1.7.0" dependencies = [ "anyhow", "arc-swap", + "arrayref", "async-trait", "axum 0.7.9", - "axum-extra", "bip39", + "blake2 0.8.1", "bs58", "cargo_metadata 0.18.1", "celes", + "chacha", "clap", "colored", + "csv", "cupid", - "dashmap", "futures", - "headers", "human-repr", "humantime-serde", + "indicatif", "ipnetwork", + "lioness", "nym-authenticator", "nym-bin-common", "nym-client-core-config-types", @@ -6145,6 +6150,8 @@ dependencies = [ "nym-sphinx-addressing", "nym-sphinx-forwarding", "nym-sphinx-framing", + "nym-sphinx-params", + "nym-sphinx-routing", "nym-sphinx-types", "nym-task", "nym-topology", @@ -6154,10 +6161,8 @@ dependencies = [ "nym-wireguard", "nym-wireguard-types", "rand 0.8.5", - "semver 1.0.26", "serde", "serde_json", - "si-scale", "sysinfo", "thiserror 2.0.12", "time", @@ -6166,6 +6171,7 @@ dependencies = [ "toml 0.8.20", "tower-http", "tracing", + "tracing-indicatif", "tracing-subscriber", "url", "utoipa", @@ -10419,6 +10425,18 @@ dependencies = [ "tracing", ] +[[package]] +name = "tracing-indicatif" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8201ca430e0cd893ef978226fd3516c06d9c494181c8bf4e5b32e30ed4b40aa1" +dependencies = [ + "indicatif", + "tracing", + "tracing-core", + "tracing-subscriber", +] + [[package]] name = "tracing-log" version = "0.1.4" @@ -11013,6 +11031,39 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "vt100" +version = "0.15.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "84cd863bf0db7e392ba3bd04994be3473491b31e66340672af5d11943c6274de" +dependencies = [ + "itoa", + "log", + "unicode-width 0.1.14", + "vte", +] + +[[package]] +name = "vte" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f5022b5fbf9407086c180e9557be968742d839e68346af7792b8592489732197" +dependencies = [ + "arrayvec", + "utf8parse", + "vte_generate_state_changes", +] + +[[package]] +name = "vte_generate_state_changes" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e369bee1b05d510a7b4ed645f5faa90619e05437111783ea5848f28d97d3c2e" +dependencies = [ + "proc-macro2", + "quote", +] + [[package]] name = "waker-fn" version = "1.2.0" diff --git a/Cargo.toml b/Cargo.toml index a55daa56a0..15a2722b85 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -347,6 +347,7 @@ tracing-log = "0.2" tracing-opentelemetry = "0.19.0" tracing-subscriber = "0.3.19" tracing-tree = "0.2.2" +tracing-indicatif = "0.3.9" ts-rs = "10.1.0" tungstenite = { version = "0.20.1", default-features = false } uniffi = "0.29.1" diff --git a/common/bin-common/Cargo.toml b/common/bin-common/Cargo.toml index 78e8e6ca93..51cfa31b19 100644 --- a/common/bin-common/Cargo.toml +++ b/common/bin-common/Cargo.toml @@ -21,6 +21,7 @@ serde_json = { workspace = true, optional = true } ## tracing tracing-subscriber = { workspace = true, features = ["env-filter"], optional = true } tracing-tree = { workspace = true, optional = true } +tracing = { workspace = true, optional = true } opentelemetry-jaeger = { workspace = true, features = ["rt-tokio", "collector_client", "isahc_collector_client"], optional = true } tracing-opentelemetry = { workspace = true, optional = true } utoipa = { workspace = true, optional = true } @@ -35,7 +36,7 @@ default = [] openapi = ["utoipa"] output_format = ["serde_json", "dep:clap"] bin_info_schema = ["schemars"] -basic_tracing = ["tracing-subscriber"] +basic_tracing = ["dep:tracing", "tracing-subscriber"] tracing = [ "basic_tracing", "tracing-tree", diff --git a/common/bin-common/src/logging/mod.rs b/common/bin-common/src/logging/mod.rs index bd9a76374b..ffbe6a6cb0 100644 --- a/common/bin-common/src/logging/mod.rs +++ b/common/bin-common/src/logging/mod.rs @@ -44,10 +44,38 @@ pub fn setup_logging() { .init(); } +// don't call init so that we could attach additional layers #[cfg(feature = "basic_tracing")] -pub fn setup_tracing_logger() { - let log_builder = tracing_subscriber::fmt() - .with_writer(std::io::stderr) +pub fn build_tracing_logger() -> impl tracing_subscriber::layer::SubscriberExt { + use tracing_subscriber::prelude::*; + + tracing_subscriber::registry() + .with(default_tracing_fmt_layer(std::io::stderr)) + .with(default_tracing_env_filter()) +} + +#[cfg(feature = "basic_tracing")] +pub fn default_tracing_env_filter() -> tracing_subscriber::filter::EnvFilter { + if ::std::env::var("RUST_LOG").is_ok() { + tracing_subscriber::filter::EnvFilter::from_default_env() + } else { + // if the env value was not found, default to `INFO` level rather than `ERROR` + tracing_subscriber::filter::EnvFilter::builder() + .with_default_directive(tracing_subscriber::filter::LevelFilter::INFO.into()) + .parse_lossy("") + } +} + +#[cfg(feature = "basic_tracing")] +pub fn default_tracing_fmt_layer( + writer: W, +) -> impl tracing_subscriber::Layer + Sync + Send + 'static +where + S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>, + W: for<'writer> tracing_subscriber::fmt::MakeWriter<'writer> + Sync + Send + 'static, +{ + tracing_subscriber::fmt::layer() + .with_writer(writer) // Use a more compact, abbreviated log format .compact() // Display source code file paths @@ -55,18 +83,13 @@ pub fn setup_tracing_logger() { // Display source code line numbers .with_line_number(true) // Don't display the event's target (module path) - .with_target(false); + .with_target(false) +} - if ::std::env::var("RUST_LOG").is_ok() { - log_builder - .with_env_filter(tracing_subscriber::filter::EnvFilter::from_default_env()) - .init() - } else { - // default to 'Info - log_builder - .with_max_level(tracing_subscriber::filter::LevelFilter::INFO) - .init() - } +#[cfg(feature = "basic_tracing")] +pub fn setup_tracing_logger() { + use tracing_subscriber::util::SubscriberInitExt; + build_tracing_logger().init() } // TODO: This has to be a macro, running it as a function does not work for the file_appender for some reason diff --git a/common/client-libs/mixnet-client/src/client.rs b/common/client-libs/mixnet-client/src/client.rs index 94e32f8404..8788a0aee7 100644 --- a/common/client-libs/mixnet-client/src/client.rs +++ b/common/client-libs/mixnet-client/src/client.rs @@ -24,10 +24,10 @@ use tracing::*; #[derive(Clone, Copy)] pub struct Config { - initial_reconnection_backoff: Duration, - maximum_reconnection_backoff: Duration, - initial_connection_timeout: Duration, - maximum_connection_buffer_size: usize, + pub initial_reconnection_backoff: Duration, + pub maximum_reconnection_backoff: Duration, + pub initial_connection_timeout: Duration, + pub maximum_connection_buffer_size: usize, } impl Config { @@ -50,7 +50,7 @@ pub trait SendWithoutResponse { // Without response in this context means we will not listen for anything we might get back (not // that we should get anything), including any possible io errors fn send_without_response( - &mut self, + &self, address: NymNodeRoutingAddress, packet: NymPacket, packet_type: PacketType, @@ -196,7 +196,7 @@ impl Client { } } - fn make_connection(&mut self, address: NymNodeRoutingAddress, pending_packet: FramedNymPacket) { + fn make_connection(&self, address: NymNodeRoutingAddress, pending_packet: FramedNymPacket) { let (sender, receiver) = mpsc::channel(self.config.maximum_connection_buffer_size); // this CAN'T fail because we just created the channel which has a non-zero capacity @@ -247,7 +247,7 @@ impl Client { impl SendWithoutResponse for Client { fn send_without_response( - &mut self, + &self, address: NymNodeRoutingAddress, packet: NymPacket, packet_type: PacketType, diff --git a/common/nymsphinx/types/src/lib.rs b/common/nymsphinx/types/src/lib.rs index 517e45e783..7afb5df1fe 100644 --- a/common/nymsphinx/types/src/lib.rs +++ b/common/nymsphinx/types/src/lib.rs @@ -8,7 +8,7 @@ use thiserror::Error; use nym_outfox::packet::{OutfoxPacket, OutfoxProcessedPacket}; #[cfg(feature = "sphinx")] -use sphinx_packet::{SphinxPacket, SphinxPacketBuilder}; +pub use sphinx_packet::{SphinxPacket, SphinxPacketBuilder}; #[cfg(feature = "outfox")] pub use nym_outfox::{ @@ -166,4 +166,20 @@ impl NymPacket { } } } + + #[cfg(feature = "sphinx")] + pub fn sphinx_packet_ref(&self) -> Option<&SphinxPacket> { + match self { + NymPacket::Sphinx(packet) => Some(packet), + _ => None, + } + } + + #[cfg(feature = "sphinx")] + pub fn as_sphinx_packet(self) -> Option { + match self { + NymPacket::Sphinx(packet) => Some(packet), + _ => None, + } + } } diff --git a/nym-node/Cargo.toml b/nym-node/Cargo.toml index 45f65f3b8d..b80ec18064 100644 --- a/nym-node/Cargo.toml +++ b/nym-node/Cargo.toml @@ -21,18 +21,19 @@ bip39 = { workspace = true, features = ["zeroize"] } bs58.workspace = true celes = { workspace = true } # country codes colored = { workspace = true } +csv = { workspace = true } clap = { workspace = true, features = ["cargo", "env"] } -dashmap = { workspace = true } futures = { workspace = true } humantime-serde = { workspace = true } human-repr = { workspace = true } ipnetwork = { workspace = true } +indicatif = { workspace = true } rand = { workspace = true } serde = { workspace = true, features = ["derive"] } serde_json.workspace = true -si-scale = { workspace = true } thiserror.workspace = true tracing.workspace = true +tracing-indicatif = { workspace = true } tracing-subscriber.workspace = true tokio = { workspace = true, features = ["macros", "sync", "rt-multi-thread"] } tokio-util = { workspace = true, features = ["codec"] } @@ -40,9 +41,6 @@ toml = { workspace = true } url = { workspace = true, features = ["serde"] } zeroize = { workspace = true, features = ["zeroize_derive"] } -# temporary bonding information v1 (to grab and parse nym-mixnode and nym-gateway package versions) -semver = { workspace = true } - # system info: cupid = { workspace = true } sysinfo = { workspace = true } @@ -62,6 +60,8 @@ nym-sphinx-addressing = { path = "../common/nymsphinx/addressing" } nym-sphinx-framing = { path = "../common/nymsphinx/framing" } nym-sphinx-types = { path = "../common/nymsphinx/types" } nym-sphinx-forwarding = { path = "../common/nymsphinx/forwarding" } +nym-sphinx-routing = { path = "../common/nymsphinx/routing" } +nym-sphinx-params = { path = "../common/nymsphinx/params" } nym-task = { path = "../common/task" } nym-types = { path = "../common/types" } nym-validator-client = { path = "../common/client-libs/validator-client" } @@ -77,8 +77,6 @@ nym-http-api-client = { path = "../common/http-api-client" } # useful for `#[axum_macros::debug_handler]` #axum-macros = "0.3.8" axum.workspace = true -axum-extra = { workspace = true, features = ["typed-header"] } -headers.workspace = true time = { workspace = true, features = ["serde"] } tower-http = { workspace = true, features = ["fs"] } utoipa = { workspace = true, features = ["axum_extras", "time"] } @@ -95,6 +93,22 @@ nym-network-requester = { path = "../service-providers/network-requester" } nym-ip-packet-router = { path = "../service-providers/ip-packet-router" } +# throughput tester to recreate lioness +# we don't care about particular versions - just pull whatever is used by sphinx +[dependencies.lioness] +version = "*" + +[dependencies.chacha] +version = "*" + +[dependencies.arrayref] +version = "*" + +[dependencies.blake2] +version = "*" + + + [build-dependencies] # temporary bonding information v1 (to grab and parse nym-mixnode and nym-gateway package versions) cargo_metadata = { workspace = true } diff --git a/nym-node/src/cli/commands/bonding_information.rs b/nym-node/src/cli/commands/bonding_information.rs index a864e14cbb..6edc7fd7c4 100644 --- a/nym-node/src/cli/commands/bonding_information.rs +++ b/nym-node/src/cli/commands/bonding_information.rs @@ -3,7 +3,6 @@ use crate::cli::helpers::ConfigArgs; use crate::config::upgrade_helpers::try_load_current_config; -use crate::error::NymNodeError; use crate::node::bonding_information::BondingInformation; use nym_bin_common::output_format::OutputFormat; @@ -21,7 +20,7 @@ pub struct Args { pub(crate) output: OutputFormat, } -pub async fn execute(args: Args) -> Result<(), NymNodeError> { +pub async fn execute(args: Args) -> anyhow::Result<()> { let config = try_load_current_config(args.config.config_path()).await?; let info = BondingInformation::try_load(&config)?; args.output.to_stdout(&info); diff --git a/nym-node/src/cli/commands/build_info.rs b/nym-node/src/cli/commands/build_info.rs index 612f681e4e..7bddf236f3 100644 --- a/nym-node/src/cli/commands/build_info.rs +++ b/nym-node/src/cli/commands/build_info.rs @@ -1,7 +1,6 @@ // Copyright 2024 - Nym Technologies SA // SPDX-License-Identifier: GPL-3.0-only -use crate::error::NymNodeError; use nym_bin_common::bin_info_owned; use nym_bin_common::output_format::OutputFormat; @@ -11,7 +10,7 @@ pub(crate) struct Args { output: OutputFormat, } -pub(crate) fn execute(args: Args) -> Result<(), NymNodeError> { +pub(crate) fn execute(args: Args) -> anyhow::Result<()> { println!("{}", args.output.format(&bin_info_owned!())); Ok(()) } diff --git a/nym-node/src/cli/commands/migrate.rs b/nym-node/src/cli/commands/migrate.rs index b47ee6853c..4575d7d020 100644 --- a/nym-node/src/cli/commands/migrate.rs +++ b/nym-node/src/cli/commands/migrate.rs @@ -11,7 +11,7 @@ pub(crate) struct Args { _args: Vec, } -pub(crate) async fn execute(_args: Args) -> Result<(), NymNodeError> { +pub(crate) fn execute(_args: Args) -> Result<(), NymNodeError> { let orange = TrueColor { r: 251, g: 110, diff --git a/nym-node/src/cli/commands/mod.rs b/nym-node/src/cli/commands/mod.rs index 16d6e3611b..dea762daf9 100644 --- a/nym-node/src/cli/commands/mod.rs +++ b/nym-node/src/cli/commands/mod.rs @@ -7,3 +7,4 @@ pub(super) mod migrate; pub(crate) mod node_details; pub(super) mod run; pub(super) mod sign; +pub(crate) mod test_throughput; diff --git a/nym-node/src/cli/commands/node_details.rs b/nym-node/src/cli/commands/node_details.rs index 1b8ec69360..11ef51f120 100644 --- a/nym-node/src/cli/commands/node_details.rs +++ b/nym-node/src/cli/commands/node_details.rs @@ -3,7 +3,6 @@ use crate::cli::helpers::ConfigArgs; use crate::config::upgrade_helpers::try_load_current_config; -use crate::error::NymNodeError; use crate::node::NymNode; use nym_bin_common::output_format::OutputFormat; @@ -21,7 +20,7 @@ pub(crate) struct Args { pub(crate) output: OutputFormat, } -pub async fn execute(args: Args) -> Result<(), NymNodeError> { +pub async fn execute(args: Args) -> anyhow::Result<()> { let config = try_load_current_config(args.config.config_path()).await?; let details = NymNode::new(config).await?.display_details()?; args.output.to_stdout(&details); diff --git a/nym-node/src/cli/commands/sign.rs b/nym-node/src/cli/commands/sign.rs index 0b9fc56d5b..a9326df24f 100644 --- a/nym-node/src/cli/commands/sign.rs +++ b/nym-node/src/cli/commands/sign.rs @@ -3,7 +3,6 @@ use crate::cli::helpers::ConfigArgs; use crate::config::upgrade_helpers::try_load_current_config; -use crate::error::NymNodeError; use crate::node::helpers::load_ed25519_identity_keypair; use nym_bin_common::output_format::OutputFormat; use nym_crypto::asymmetric::identity; @@ -70,7 +69,7 @@ fn print_signed_contract_msg( // SAFETY: clippy ArgGroup ensures only a single branch is actually called #[allow(clippy::unreachable)] -pub async fn execute(args: Args) -> Result<(), NymNodeError> { +pub async fn execute(args: Args) -> anyhow::Result<()> { let config = try_load_current_config(args.config.config_path()).await?; let identity_keypair = load_ed25519_identity_keypair(config.storage_paths.keys.ed25519_identity_storage_paths())?; diff --git a/nym-node/src/cli/commands/test_throughput.rs b/nym-node/src/cli/commands/test_throughput.rs new file mode 100644 index 0000000000..af7cc8ad42 --- /dev/null +++ b/nym-node/src/cli/commands/test_throughput.rs @@ -0,0 +1,87 @@ +// Copyright 2025 - Nym Technologies SA +// SPDX-License-Identifier: GPL-3.0-only + +use crate::cli::helpers::ConfigArgs; +use crate::logging::granual_filtered_env; +use crate::throughput_tester::test_mixing_throughput; +use anyhow::bail; +use humantime_serde::re::humantime; +use indicatif::ProgressStyle; +use nym_bin_common::logging::default_tracing_fmt_layer; +use std::env::temp_dir; +use std::path::PathBuf; +use std::time::Duration; +use time::OffsetDateTime; +use tracing_indicatif::IndicatifLayer; +use tracing_subscriber::layer::SubscriberExt; +use tracing_subscriber::util::SubscriberInitExt; + +#[derive(Debug, clap::Args)] +pub struct Args { + #[clap(flatten)] + config: ConfigArgs, + + #[clap(long, default_value_t = 10)] + senders: usize, + + /// target packet latency, if current value is below threshold, clients will increase their sending rates + /// and similarly if it's above it, they will decrease it + #[clap(long, default_value = "15ms", value_parser = humantime::parse_duration)] + packet_latency_threshold: Duration, + + #[clap(long, default_value_t = 50)] + starting_sending_batch_size: usize, + + #[clap(long, default_value = "50ms", value_parser = humantime::parse_duration)] + starting_sending_delay: Duration, + + #[clap(long, short)] + output_directory: Option, +} + +fn init_test_logger() -> anyhow::Result<()> { + let indicatif_layer = IndicatifLayer::new() + .with_progress_style(ProgressStyle::with_template( + "{span_child_prefix}{spinner} {span_fields} -- {span_name} {wide_msg}", + )?) + .with_span_child_prefix_symbol("↳ ") + .with_span_child_prefix_indent(" "); + + tracing_subscriber::registry() + .with(default_tracing_fmt_layer( + indicatif_layer.get_stderr_writer(), + )) + .with(indicatif_layer) + .with(granual_filtered_env()?) + .init(); + Ok(()) +} + +pub fn execute(args: Args) -> anyhow::Result<()> { + init_test_logger()?; + + let output_dir = match args.output_directory { + Some(output_dir) => { + if !output_dir.is_dir() { + bail!("'{}' is not a directory", output_dir.display()); + } + + output_dir + } + None => { + let now = OffsetDateTime::now_utc().unix_timestamp(); + temp_dir() + .join("nym-node-throughput-testing") + .join(now.to_string()) + } + }; + + test_mixing_throughput( + args.config.config_path(), + args.senders, + args.packet_latency_threshold, + args.starting_sending_batch_size, + args.starting_sending_delay, + output_dir, + ) +} diff --git a/nym-node/src/cli/mod.rs b/nym-node/src/cli/mod.rs index d3374c7fdf..40b05b4a84 100644 --- a/nym-node/src/cli/mod.rs +++ b/nym-node/src/cli/mod.rs @@ -1,14 +1,17 @@ // Copyright 2024 - Nym Technologies SA // SPDX-License-Identifier: GPL-3.0-only -use crate::cli::commands::{bonding_information, build_info, migrate, node_details, run, sign}; +use crate::cli::commands::{ + bonding_information, build_info, migrate, node_details, run, sign, test_throughput, +}; use crate::env::vars::{NYMNODE_CONFIG_ENV_FILE_ARG, NYMNODE_NO_BANNER_ARG}; -use crate::error::NymNodeError; +use crate::logging::setup_tracing_logger; use clap::{Parser, Subcommand}; use nym_bin_common::bin_info; +use std::future::Future; use std::sync::OnceLock; -mod commands; +pub(crate) mod commands; mod helpers; pub const DEFAULT_NYMNODE_ID: &str = "default-nym-node"; @@ -42,15 +45,31 @@ pub(crate) struct Cli { } impl Cli { - pub(crate) async fn execute(self) -> Result<(), NymNodeError> { - match self.command { - Commands::BuildInfo(args) => build_info::execute(args), - Commands::BondingInformation(args) => bonding_information::execute(args).await, - Commands::NodeDetails(args) => node_details::execute(args).await, - Commands::Run(args) => run::execute(*args).await, - Commands::Migrate(args) => migrate::execute(*args).await, - Commands::Sign(args) => sign::execute(args).await, + fn execute_async(fut: F) -> anyhow::Result { + Ok(tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build()? + .block_on(fut)) + } + + pub(crate) fn execute(self) -> anyhow::Result<()> { + // NOTE: `test_throughput` sets up its own logger as it has to include additional layers + if !matches!(self.command, Commands::TestThroughput(..)) { + setup_tracing_logger()?; } + + match self.command { + Commands::BuildInfo(args) => build_info::execute(args)?, + Commands::BondingInformation(args) => { + { Self::execute_async(bonding_information::execute(args))? }? + } + Commands::NodeDetails(args) => { Self::execute_async(node_details::execute(args))? }?, + Commands::Run(args) => { Self::execute_async(run::execute(*args))? }?, + Commands::Migrate(args) => migrate::execute(*args)?, + Commands::Sign(args) => { Self::execute_async(sign::execute(args))? }?, + Commands::TestThroughput(args) => test_throughput::execute(args)?, + } + Ok(()) } } @@ -73,6 +92,11 @@ pub(crate) enum Commands { /// Use identity key of this node to sign provided message. Sign(sign::Args), + + /// Attempt to approximate the maximum mixnet throughput if nym-node + /// was running on this machine in mixnet mode + #[clap(hide = true)] + TestThroughput(test_throughput::Args), } #[cfg(test)] diff --git a/nym-node/src/logging.rs b/nym-node/src/logging.rs index 2fd15d94b6..332e274944 100644 --- a/nym-node/src/logging.rs +++ b/nym-node/src/logging.rs @@ -1,35 +1,34 @@ // Copyright 2024 - Nym Technologies SA // SPDX-License-Identifier: GPL-3.0-only -use tracing::level_filters::LevelFilter; -use tracing_subscriber::{filter::Directive, EnvFilter}; +use nym_bin_common::logging::{default_tracing_env_filter, default_tracing_fmt_layer}; +use tracing_subscriber::filter::Directive; +use tracing_subscriber::layer::SubscriberExt; +use tracing_subscriber::util::SubscriberInitExt; -pub(crate) fn setup_tracing_logger() -> anyhow::Result<()> { +pub(crate) fn granual_filtered_env() -> anyhow::Result { fn directive_checked(directive: impl Into) -> anyhow::Result { directive.into().parse().map_err(From::from) } - let log_builder = tracing_subscriber::fmt() - // Use a more compact, abbreviated log format - .compact() - // Display source code file paths - .with_file(true) - // Display source code line numbers - .with_line_number(true) - // Don't display the event's target (module path) - .with_target(false); + let mut filter = default_tracing_env_filter(); - let mut filter = EnvFilter::builder() - // if RUST_LOG isn't set, set default level - .with_default_directive(LevelFilter::INFO.into()) - .from_env_lossy(); // these crates are more granularly filtered let filter_crates = ["defguard_wireguard_rs"]; for crate_name in filter_crates { filter = filter.add_directive(directive_checked(format!("{}=warn", crate_name))?); } + Ok(filter) +} - log_builder.with_env_filter(filter).init(); +pub(crate) fn build_tracing_logger() -> anyhow::Result { + Ok(tracing_subscriber::registry() + .with(default_tracing_fmt_layer(std::io::stderr)) + .with(granual_filtered_env()?)) +} + +pub(crate) fn setup_tracing_logger() -> anyhow::Result<()> { + build_tracing_logger()?.init(); Ok(()) } diff --git a/nym-node/src/main.rs b/nym-node/src/main.rs index 9d66d778d2..04ff85bb91 100644 --- a/nym-node/src/main.rs +++ b/nym-node/src/main.rs @@ -4,7 +4,7 @@ #![warn(clippy::expect_used)] #![warn(clippy::unwrap_used)] -use crate::{cli::Cli, logging::setup_tracing_logger}; +use crate::cli::Cli; use clap::{crate_name, crate_version, Parser}; use nym_bin_common::logging::maybe_print_banner; use nym_config::defaults::setup_env; @@ -15,10 +15,10 @@ mod env; pub(crate) mod error; mod logging; pub(crate) mod node; +pub(crate) mod throughput_tester; pub(crate) mod wireguard; -#[tokio::main] -async fn main() -> anyhow::Result<()> { +fn main() -> anyhow::Result<()> { // std::env::set_var( // "RUST_LOG", // "trace,handlebars=warn,tendermint_rpc=warn,h2=warn,hyper=warn,rustls=warn,reqwest=warn,tungstenite=warn,async_tungstenite=warn,tokio_util=warn,tokio_tungstenite=warn,tokio-util=warn", @@ -26,13 +26,12 @@ async fn main() -> anyhow::Result<()> { let cli = Cli::parse(); setup_env(cli.config_env_file.as_ref()); - setup_tracing_logger()?; if !cli.no_banner { maybe_print_banner(crate_name!(), crate_version!()); } - cli.execute().await?; + cli.execute()?; Ok(()) } diff --git a/nym-node/src/node/mixnet/packet_forwarding/mod.rs b/nym-node/src/node/mixnet/packet_forwarding/mod.rs index a61f5517ec..db120e23a7 100644 --- a/nym-node/src/node/mixnet/packet_forwarding/mod.rs +++ b/nym-node/src/node/mixnet/packet_forwarding/mod.rs @@ -1,7 +1,6 @@ // Copyright 2024 - Nym Technologies SA // SPDX-License-Identifier: GPL-3.0-only -use crate::node::mixnet::packet_forwarding::global::is_global_ip; use crate::node::shared_network::RoutingFilter; use futures::StreamExt; use nym_mixnet_client::forwarder::{ @@ -13,38 +12,33 @@ use nym_nonexhaustive_delayqueue::{Expired, NonExhaustiveDelayQueue}; use nym_sphinx_forwarding::packet::MixPacket; use nym_task::ShutdownToken; use std::io; -use std::net::IpAddr; use tokio::time::Instant; use tracing::{debug, error, trace, warn}; pub(crate) mod global; -pub struct PacketForwarder { - testnet: bool, - +pub struct PacketForwarder { delay_queue: NonExhaustiveDelayQueue, mixnet_client: C, metrics: NymNodeMetrics, - routing_filter: RoutingFilter, + routing_filter: F, packet_sender: MixForwardingSender, packet_receiver: MixForwardingReceiver, shutdown: ShutdownToken, } -impl PacketForwarder { +impl PacketForwarder { pub fn new( client: C, - testnet: bool, - routing_filter: RoutingFilter, + routing_filter: F, metrics: NymNodeMetrics, shutdown: ShutdownToken, ) -> Self { let (packet_sender, packet_receiver) = mix_forwarding_channels(); PacketForwarder { - testnet, delay_queue: NonExhaustiveDelayQueue::new(), mixnet_client: client, metrics, @@ -59,22 +53,14 @@ impl PacketForwarder { self.packet_sender.clone() } - fn should_route(&mut self, ip_addr: IpAddr) -> bool { - // only allow non-global ips on testnets - if self.testnet && !is_global_ip(&ip_addr) { - return true; - } - - self.routing_filter.attempt_resolve(ip_addr).should_route() - } - fn forward_packet(&mut self, packet: MixPacket) where C: SendWithoutResponse, + F: RoutingFilter, { let next_hop = packet.next_hop(); - if !self.should_route(next_hop.as_ref().ip()) { + if !self.routing_filter.should_route(next_hop.as_ref().ip()) { debug!("dropping packet as the egress address does not belong to any known node"); self.metrics .mixnet @@ -113,6 +99,7 @@ impl PacketForwarder { fn handle_done_delaying(&mut self, packet: Expired) where C: SendWithoutResponse, + F: RoutingFilter, { let delayed_packet = packet.into_inner(); self.forward_packet(delayed_packet); @@ -121,6 +108,7 @@ impl PacketForwarder { fn handle_new_packet(&mut self, new_packet: PacketToForward) where C: SendWithoutResponse, + F: RoutingFilter, { // in case of a zero delay packet, don't bother putting it in the delay queue, // just forward it immediately @@ -152,6 +140,7 @@ impl PacketForwarder { pub async fn run(&mut self) where C: SendWithoutResponse, + F: RoutingFilter, { let mut processed = 0; trace!("starting PacketForwarder"); diff --git a/nym-node/src/node/mod.rs b/nym-node/src/node/mod.rs index f68a2d5553..abf53ee102 100644 --- a/nym-node/src/node/mod.rs +++ b/nym-node/src/node/mod.rs @@ -29,7 +29,7 @@ use crate::node::mixnet::packet_forwarding::PacketForwarder; use crate::node::mixnet::shared::ProcessingConfig; use crate::node::mixnet::SharedFinalHopData; use crate::node::shared_network::{ - CachedNetwork, CachedTopologyProvider, NetworkRefresher, RoutingFilter, + CachedNetwork, CachedTopologyProvider, NetworkRefresher, OpenFilter, RoutingFilter, }; use nym_bin_common::bin_info; use nym_crypto::asymmetric::{ed25519, x25519}; @@ -461,6 +461,10 @@ impl NymNode { &self.config } + pub(crate) fn shutdown_token>(&self, child_suffix: S) -> ShutdownToken { + self.shutdown_manager.clone_token(child_suffix) + } + pub(crate) fn with_accepted_operator_terms_and_conditions( mut self, accepted_operator_terms_and_conditions: bool, @@ -528,12 +532,17 @@ impl NymNode { self.x25519_sphinx_keys.public_key() } + pub(crate) fn x25519_sphinx_keys(&self) -> Arc { + self.x25519_sphinx_keys.clone() + } + pub(crate) fn x25519_noise_key(&self) -> &x25519::PublicKey { self.x25519_noise_keys.public_key() } async fn build_network_refresher(&self) -> Result { NetworkRefresher::initialise_new( + self.config.debug.testnet, self.user_agent(), self.config.mixnet.nym_api_urls.clone(), self.config.debug.topology_cache_ttl, @@ -970,12 +979,15 @@ impl NymNode { events_sender } - pub(crate) fn start_mixnet_listener( + pub(crate) fn start_mixnet_listener( &self, active_clients_store: &ActiveClientsStore, - routing_filter: RoutingFilter, + routing_filter: F, shutdown: ShutdownToken, - ) -> (MixForwardingSender, ActiveConnections) { + ) -> (MixForwardingSender, ActiveConnections) + where + F: RoutingFilter + Send + Sync + 'static, + { let processing_config = ProcessingConfig::new(&self.config); // we're ALWAYS listening for mixnet packets, either for forward or final hops (or both) @@ -1002,7 +1014,6 @@ impl NymNode { let mut packet_forwarder = PacketForwarder::new( mixnet_client, - self.config.debug.testnet, routing_filter, self.metrics.clone(), shutdown.clone_with_suffix("mix-packet-forwarder"), @@ -1028,6 +1039,19 @@ impl NymNode { (mix_packet_sender, active_connections) } + pub(crate) async fn run_minimal_mixnet_processing(self) -> Result<(), NymNodeError> { + self.start_mixnet_listener( + &ActiveClientsStore::new(), + OpenFilter, + self.shutdown_manager.clone_token("mixnet-traffic"), + ); + + self.shutdown_manager.close(); + self.shutdown_manager.wait_for_shutdown_signal().await; + + Ok(()) + } + pub(crate) async fn run(mut self) -> Result<(), NymNodeError> { info!("starting Nym Node {} with the following modes: mixnode: {}, entry: {}, exit: {}, wireguard: {}", self.ed25519_identity_key(), diff --git a/nym-node/src/node/shared_network.rs b/nym-node/src/node/shared_network.rs index aa2e378541..3abd2f843a 100644 --- a/nym-node/src/node/shared_network.rs +++ b/nym-node/src/node/shared_network.rs @@ -2,6 +2,7 @@ // SPDX-License-Identifier: GPL-3.0-only use crate::error::NymNodeError; +use crate::node::mixnet::packet_forwarding::global::is_global_ip; use arc_swap::ArcSwap; use async_trait::async_trait; use nym_gateway::node::UserAgent; @@ -22,8 +23,34 @@ use tracing::log::error; use tracing::{debug, trace, warn}; use url::Url; +pub(crate) trait RoutingFilter { + fn should_route(&self, ip: IpAddr) -> bool; +} + +#[derive(Debug, Copy, Clone, Default)] +pub(crate) struct OpenFilter; + +impl RoutingFilter for OpenFilter { + fn should_route(&self, _: IpAddr) -> bool { + true + } +} + +impl RoutingFilter for NetworkRoutingFilter { + fn should_route(&self, ip: IpAddr) -> bool { + // only allow non-global ips on testnets + if self.testnet_mode && !is_global_ip(&ip) { + return true; + } + + self.attempt_resolve(ip).should_route() + } +} + #[derive(Clone)] -pub(crate) struct RoutingFilter { +pub(crate) struct NetworkRoutingFilter { + testnet_mode: bool, + resolved: KnownNodes, // while this is technically behind a lock, it should not be called too often as once resolved it will @@ -31,9 +58,10 @@ pub(crate) struct RoutingFilter { pending: UnknownNodes, } -impl RoutingFilter { - fn new_empty() -> Self { - RoutingFilter { +impl NetworkRoutingFilter { + fn new_empty(testnet_mode: bool) -> Self { + NetworkRoutingFilter { + testnet_mode, resolved: Default::default(), pending: Default::default(), } @@ -254,11 +282,12 @@ pub struct NetworkRefresher { shutdown_token: ShutdownToken, network: CachedNetwork, - routing_filter: RoutingFilter, + routing_filter: NetworkRoutingFilter, } impl NetworkRefresher { pub(crate) async fn initialise_new( + testnet: bool, user_agent: UserAgent, nym_api_urls: Vec, full_refresh_interval: Duration, @@ -280,7 +309,7 @@ impl NetworkRefresher { pending_check_interval, shutdown_token, network: CachedNetwork::new_empty(), - routing_filter: RoutingFilter::new_empty(), + routing_filter: NetworkRoutingFilter::new_empty(testnet), }; this.obtain_initial_network().await?; @@ -403,7 +432,7 @@ impl NetworkRefresher { .map_err(|source| NymNodeError::InitialTopologyQueryFailure { source }) } - pub(crate) fn routing_filter(&self) -> RoutingFilter { + pub(crate) fn routing_filter(&self) -> NetworkRoutingFilter { self.routing_filter.clone() } diff --git a/nym-node/src/throughput_tester/client.rs b/nym-node/src/throughput_tester/client.rs new file mode 100644 index 0000000000..3031ce5252 --- /dev/null +++ b/nym-node/src/throughput_tester/client.rs @@ -0,0 +1,438 @@ +// Copyright 2025 - Nym Technologies SA +// SPDX-License-Identifier: GPL-3.0-only + +use crate::throughput_tester::stats::ClientStats; +use anyhow::bail; +use arrayref::array_ref; +use blake2::VarBlake2b; +use chacha::ChaCha; +use futures::{stream, SinkExt, Stream, StreamExt}; +use human_repr::{HumanCount, HumanDuration, HumanThroughput}; +use lioness::Lioness; +use nym_crypto::asymmetric::x25519; +use nym_sphinx_addressing::nodes::NymNodeRoutingAddress; +use nym_sphinx_framing::codec::{NymCodec, NymCodecError}; +use nym_sphinx_framing::packet::FramedNymPacket; +use nym_sphinx_params::PacketSize; +use nym_sphinx_routing::generate_hop_delays; +use nym_sphinx_types::header::keys::PayloadKey; +use nym_sphinx_types::{ + Destination, DestinationAddressBytes, Node, NymPacket, SphinxHeader, + DESTINATION_ADDRESS_LENGTH, IDENTIFIER_LENGTH, +}; +use nym_task::ShutdownToken; +use rand::rngs::OsRng; +use std::net::SocketAddr; +use std::pin::Pin; +use std::task::{Context, Poll, Waker}; +use std::time::Duration; +use time::OffsetDateTime; +use tokio::net::{TcpListener, TcpStream}; +use tokio::select; +use tokio::time::{interval, sleep, Instant}; +use tokio_util::codec::Framed; +use tracing::{debug, error, info, Span}; +use tracing_indicatif::span_ext::IndicatifSpanExt; + +struct PacketTag { + sending_timestamp: OffsetDateTime, + batch_id: u64, + index: u64, +} + +impl PacketTag { + const SIZE: usize = 32; + + fn elapsed(&self) -> time::Duration { + OffsetDateTime::now_utc() - self.sending_timestamp + } + + fn elapsed_nanos(&self) -> u64 { + // here we're making few assumptions: the latency is lower than u64::MAX + // and it's strictly positive (which are rather valid...) + self.elapsed().whole_nanoseconds() as u64 + } + + fn to_bytes(&self) -> Vec { + self.sending_timestamp + .unix_timestamp_nanos() + .to_be_bytes() + .into_iter() + .chain(self.batch_id.to_be_bytes()) + .chain(self.index.to_be_bytes()) + .collect() + } + + #[allow(clippy::unwrap_used)] + fn from_bytes(bytes: &[u8]) -> PacketTag { + let sending_timestamp = i128::from_be_bytes(bytes[0..16].try_into().unwrap()); + let sending_timestamp = + OffsetDateTime::from_unix_timestamp_nanos(sending_timestamp).unwrap(); + + let batch_id = u64::from_be_bytes(bytes[8..16].try_into().unwrap()); + let index = u64::from_be_bytes(bytes[16..24].try_into().unwrap()); + PacketTag { + sending_timestamp, + batch_id, + index, + } + } +} + +pub(crate) struct ThroughputTestingClient { + stats: ClientStats, + last_received_update: Instant, + last_received_at_update: usize, + current_batch: u64, + sending_delay: Duration, + latency_threshold: Duration, + current_batch_size: usize, + forward_header_bytes: Vec, + unwrapped_forward_payload_bytes: Vec, + shutdown_token: ShutdownToken, + local_address: SocketAddr, + listener: TcpListener, + forward_connection: Framed, + payload_key: PayloadKey, +} + +impl ThroughputTestingClient { + pub(crate) async fn try_create( + initial_sending_delay: Duration, + initial_batch_size: usize, + latency_threshold: Duration, + node_keys: &x25519::KeyPair, + node_listener: SocketAddr, + stats: ClientStats, + cancellation_token: ShutdownToken, + ) -> anyhow::Result { + // attempt to bind to some port to receive processed packets + let listener = TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0))).await?; + let local_address = listener.local_addr()?; + + info!("listening on {local_address}"); + + // create the sphinx packet we're going to be repeatedly sending + // (next hop has to be our mixnode, then this client, and then it doesn't matter since the packet won't + // get further processed) + let mut rng = OsRng; + // keys of this client + let ephemeral_keys = x25519::KeyPair::new(&mut rng); + + let route = [ + Node::new( + NymNodeRoutingAddress::from(node_listener).try_into()?, + (*node_keys.public_key()).into(), + ), + Node::new( + NymNodeRoutingAddress::from(local_address).try_into()?, + (*ephemeral_keys.public_key()).into(), + ), + Node::new( + NymNodeRoutingAddress::from(local_address).try_into()?, + (*ephemeral_keys.public_key()).into(), + ), + ]; + let destination = Destination::new( + DestinationAddressBytes::from_bytes([0u8; DESTINATION_ADDRESS_LENGTH]), + [0u8; IDENTIFIER_LENGTH], + ); + let delays = generate_hop_delays(Duration::default(), 3); + let payload = PacketSize::RegularPacket.payload_size(); + + let forward_packet = + NymPacket::sphinx_build(payload, b"foomp", &route, &destination, &delays)?; + + // SAFETY: we constructed a sphinx packet... + #[allow(clippy::unwrap_used)] + let sphinx_packet = forward_packet.as_sphinx_packet().unwrap(); + let header = &sphinx_packet.header; + + // derive the routing keys of our node so we could tag the payload to figure out latency + // by tagging the packet + let routing_keys = SphinxHeader::compute_routing_keys( + &header.shared_secret, + (&node_keys.private_key()).as_ref(), + ); + let payload_key = routing_keys.payload_key; + let unwrapped_payload = sphinx_packet.payload.unwrap(&payload_key)?; + let unwrapped_forward_payload_bytes = unwrapped_payload.into_bytes(); + + let start = Instant::now(); + let forward_connection = loop { + if let Ok(connection) = TcpStream::connect(node_listener).await { + break connection; + } + // fallback + sleep(Duration::from_secs(1)).await; + if start.elapsed() > Duration::from_secs(10) { + bail!("failed to connect to local nym-node") + } + }; + + Ok(ThroughputTestingClient { + stats, + last_received_update: Instant::now(), + last_received_at_update: 0, + current_batch: 0, + sending_delay: initial_sending_delay, + latency_threshold, + current_batch_size: initial_batch_size, + forward_header_bytes: sphinx_packet.header.to_bytes(), + unwrapped_forward_payload_bytes, + shutdown_token: cancellation_token, + local_address, + listener, + forward_connection: Framed::new(forward_connection, NymCodec), + payload_key, + }) + } + + fn update_progress_bar(&mut self) { + let received = self.stats.received(); + let sent = self.stats.sent(); + let latency = self.stats.average_latency_duration(); + let received_since_update = received - self.last_received_at_update; + + let time_delta_secs = self.last_received_update.elapsed().as_secs_f64(); + let receive_rate = received_since_update as f64 / time_delta_secs; + + self.last_received_at_update = received; + self.last_received_update = Instant::now(); + // I couldn't figure out how to directly pull it from span fields without duplication, + // so that's a second best + Span::current().pb_set_message(&format!( + "{}: CURRENT SENDING DELAY/BATCH: {} / {} | received: {} sent: {} (avg packet latency: {}, avg receive rate: {})", + self.local_address, + self.sending_delay.human_duration(), + self.current_batch_size, + received.human_count_bare(), + sent.human_count_bare(), + latency.human_duration(), + receive_rate.human_throughput("packets") + )); + } + + fn lioness_encrypt(&self, block: &mut [u8]) -> anyhow::Result<()> { + let lioness_cipher = Lioness::::new_raw(array_ref!( + self.payload_key, + 0, + lioness::RAW_KEY_SIZE + )); + lioness_cipher.encrypt(block)?; + Ok(()) + } + + fn tag_framed_packet(&self, tag: PacketTag) -> anyhow::Result { + let tag_bytes = tag.to_bytes(); + + let mut payload_bytes = self.unwrapped_forward_payload_bytes.clone(); + payload_bytes[..PacketTag::SIZE].copy_from_slice(&tag_bytes); + + self.lioness_encrypt(&mut payload_bytes)?; + + let mut packet_bytes = self.forward_header_bytes.clone(); + packet_bytes.append(&mut payload_bytes); + + let forward_packet = NymPacket::sphinx_from_bytes(&packet_bytes)?; + Ok(FramedNymPacket::new(forward_packet, Default::default())) + } + + async fn send_packets(&mut self) -> anyhow::Result<()> { + // mess with our payload in such a way that upon unwrapping by the first hop, + // we'll get our tag + let mut batch = Vec::with_capacity(self.current_batch_size); + let now = OffsetDateTime::now_utc(); + for i in 0..self.current_batch_size { + let tag = PacketTag { + sending_timestamp: now, + batch_id: self.current_batch, + index: i as u64, + }; + let framed_packet = self.tag_framed_packet(tag)?; + batch.push(Ok(framed_packet)); + } + + self.current_batch += 1; + self.forward_connection + .send_all(&mut stream::iter(batch)) + .await?; + self.stats.new_sent_batch(self.current_batch_size); + + Ok(()) + } + + // don't bother processing packets, just increment the count because that's the only thing that matters + fn handle_received(&mut self, maybe_packet: Result) { + let Ok(received) = maybe_packet else { + error!("FAILED TO RECEIVE PACKET"); + return; + }; + let inner = received.into_inner(); + // safety: we sent a sphinx packet... + #[allow(clippy::unwrap_used)] + let sphinx = inner.as_sphinx_packet().unwrap(); + let tag = PacketTag::from_bytes(sphinx.payload.as_bytes()); + + self.stats.new_received(tag.elapsed_nanos()); + } + + fn update_sending_rates(&mut self) { + let current = self.stats.average_latency_nanos() as f64; + let threshold = self.latency_threshold.as_nanos() as f64; + + let saturation = current / threshold; + + let sending_delay_nanos = self.sending_delay.as_nanos(); + let batch_size = self.current_batch_size; + + let diff = 1. - saturation; + + if saturation > 1. { + debug!("saturation {saturation:.2}, packet latency over threshold: need to decrease sending rate"); + } else { + debug!("saturation {saturation:.2}, packet latency under threshold: can increase sending rate"); + } + + // be conservative and only apply 50% of the diff + // (and split it equally between sending delay and batch size) + // but also make sure the current values don't increase by more than 5% + let mut new_batch_size = (batch_size as f64 * (1. + 0.25 * diff)).floor() as u64; + let mut new_sending_delay_nanos = + (sending_delay_nanos as f64 * (1. - 0.25 * diff)).floor() as u64; + + if (new_batch_size as f64) > (batch_size as f64 * 1.05) { + new_batch_size = ((batch_size as f64) * 1.05) as u64; + } + if (new_batch_size as f64) < (batch_size as f64 * 0.95) { + new_batch_size = ((batch_size as f64) * 0.95) as u64; + } + + if (new_sending_delay_nanos as f64) > (sending_delay_nanos as f64 * 1.05) { + new_sending_delay_nanos = ((sending_delay_nanos as f64) * 1.05) as u64; + } + if (new_sending_delay_nanos as f64) < (sending_delay_nanos as f64 * 0.95) { + new_sending_delay_nanos = ((sending_delay_nanos as f64) * 0.95) as u64; + } + + // normalize values + if new_batch_size < 20 { + new_batch_size = 20; + } + let mut new_sending_delay = Duration::from_nanos(new_sending_delay_nanos); + + if new_sending_delay.is_zero() { + new_sending_delay = Duration::from_micros(500); + } + if new_sending_delay.as_millis() > 100 { + new_sending_delay = Duration::from_millis(100); + } + + debug!( + "changing sending delay from {} to {}", + self.sending_delay.human_duration(), + new_sending_delay.human_duration() + ); + debug!("changing sending batch from {batch_size} to {new_batch_size}"); + + self.sending_delay = new_sending_delay; + self.current_batch_size = new_batch_size as usize; + } + + #[allow(clippy::panic)] + pub(crate) async fn run(mut self) -> anyhow::Result<()> { + let mut ingress_connection = StreamWrapper::default(); + + let mut sending_interval = interval(self.sending_delay); + sending_interval.reset(); + + // quite arbitrary + let mut update_interval = interval(Duration::from_millis(500)); + update_interval.reset(); + + let mut last_rate_update = Instant::now(); + + loop { + select! { + biased; + _ = self.shutdown_token.cancelled() => { + info!("cancelled"); + return Ok(()); + } + _ = update_interval.tick() => { + self.update_progress_bar(); + + // every 500ms attempt to adjust sending rates + if last_rate_update.elapsed() > Duration::from_millis(500) { + last_rate_update = Instant::now(); + self.update_sending_rates(); + sending_interval = interval(self.sending_delay); + sending_interval.reset(); + } + + } + accepted = self.listener.accept() => { + info!("accepted connection"); + if ingress_connection.inner.is_some() { + // this should never happen under local settings + // (and since it's not exposed to 'proper' traffic, it's fine to panic and shutdown) + panic!("attempted to overwrite existing connection") + } + let (stream, _) = accepted?; + let framed = Framed::new(stream, NymCodec); + ingress_connection.set(framed); + } + received = ingress_connection.next() => { + let Some(received) = received else { + // if the stream has terminated, we return + if ingress_connection.inner.is_some() { + return Ok(()) + } + continue; + }; + self.handle_received(received) + } + _ = sending_interval.tick() => { + self.send_packets().await?; + } + } + } + } +} + +// I must be blind, because I couldn't find something to do equivalent of `OptionStream`... +#[derive(Default)] +struct StreamWrapper { + inner: Option>, + maybe_initial_waker: Option, +} + +impl StreamWrapper { + fn set(&mut self, inner: Framed) { + self.inner = Some(inner); + if let Some(waker) = self.maybe_initial_waker.take() { + waker.wake(); + } + } +} + +impl Stream for StreamWrapper { + type Item = as Stream>::Item; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + match self.inner.as_mut() { + None => { + self.maybe_initial_waker = Some(cx.waker().clone()); + Poll::Pending + } + Some(inner) => Pin::new(inner).poll_next(cx), + } + } + + fn size_hint(&self) -> (usize, Option) { + match &self.inner { + None => (0, None), + Some(inner) => inner.size_hint(), + } + } +} diff --git a/nym-node/src/throughput_tester/global_stats.rs b/nym-node/src/throughput_tester/global_stats.rs new file mode 100644 index 0000000000..97239ddcee --- /dev/null +++ b/nym-node/src/throughput_tester/global_stats.rs @@ -0,0 +1,183 @@ +// Copyright 2025 - Nym Technologies SA +// SPDX-License-Identifier: GPL-3.0-only + +use crate::throughput_tester::stats::ClientStats; +use colored::Colorize; +use human_repr::{HumanCount, HumanDuration, HumanThroughput}; +use nym_task::ShutdownToken; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use std::fs::create_dir_all; +use std::path::PathBuf; +use std::time::Duration; +use sysinfo::System; +use time::OffsetDateTime; +use tokio::select; +use tokio::time::{interval, Instant}; +use tracing::{error, info, Span}; +use tracing_indicatif::span_ext::IndicatifSpanExt; + +#[derive(Default, Serialize, Deserialize)] +pub(crate) struct StatsRecord { + timestamp: i64, + receive_rate: f64, + latency: u64, + sent: usize, + received: usize, +} + +pub(crate) struct GlobalStatsUpdater { + system: System, + last_update: Instant, + last_total_received: usize, + header_span: Span, + client_stats: Vec, + + global_records: Vec, + records: HashMap>, + + output_directory: PathBuf, + shutdown: ShutdownToken, +} + +impl GlobalStatsUpdater { + pub(crate) fn new( + header_span: Span, + client_stats: Vec, + output_directory: PathBuf, + shutdown: ShutdownToken, + ) -> Self { + let mut system_info = System::new_all(); + system_info.refresh_cpu_usage(); + + // pre-allocate vecs + let mut records = HashMap::new(); + for (i, _) in client_stats.iter().enumerate() { + records.insert(i, vec![]); + } + + GlobalStatsUpdater { + system: system_info, + last_update: Instant::now(), + last_total_received: 0, + header_span, + client_stats, + global_records: vec![], + records, + output_directory, + shutdown, + } + } + + fn update_stats_span(&mut self) { + let now = OffsetDateTime::now_utc().unix_timestamp(); + let time_delta_secs = self.last_update.elapsed().as_secs_f64(); + + let mut all_received = 0; + let mut all_sent = 0; + let mut all_latencies = 0; + for (i, stat) in self.client_stats.iter().enumerate() { + // SAFETY: we create all entries during initialisation + #[allow(clippy::unwrap_used)] + let records = self.records.get_mut(&i).unwrap(); + + let mut client_record = StatsRecord::default(); + + let sent = stat.sent(); + let received = stat.received(); + let latency = stat.average_latency_nanos(); + + client_record.timestamp = now; + client_record.received = received; + client_record.sent = sent; + client_record.latency = latency; + if let Some(last) = records.last() { + let receive_rate = (received - last.received) as f64 / time_delta_secs; + client_record.receive_rate = receive_rate; + } + records.push(client_record); + + all_sent += sent; + all_received += received; + all_latencies += latency; + } + + let receive_rate = (all_received - self.last_total_received) as f64 / time_delta_secs; + let avg_rate = receive_rate.human_throughput("packets"); + let avg_latency = all_latencies as f64 / self.client_stats.len() as f64; + + self.global_records.push(StatsRecord { + timestamp: now, + receive_rate, + latency: avg_latency as u64, + sent: all_sent, + received: all_received, + }); + + self.system.refresh_cpu_usage(); + let cpu_usage = self.system.global_cpu_usage(); + let cpu_count = self.system.cpus().len(); + let usage_per_cpu = cpu_usage / cpu_count as f32; + + let formatted_usage = if usage_per_cpu < 0.3 { + format!("{:.2}%", usage_per_cpu * 100.).green().bold() + } else if usage_per_cpu < 0.7 { + format!("{:.2}%", usage_per_cpu * 100.).yellow().bold() + } else { + format!("{:.2}%", usage_per_cpu * 100.).red().bold() + }; + + self.header_span.pb_set_message(&format!( + "active_clients: {} | total received: {} total sent {} (avg packet latency: {}, total receive rate: {avg_rate}), avg core load: {formatted_usage}", + self.client_stats.len(), + all_received.human_count_bare(), + all_sent.human_count_bare(), + Duration::from_nanos(avg_latency as u64).human_duration() + )); + self.last_total_received = all_received; + self.last_update = Instant::now(); + } + + fn save_results_to_files(&self) -> anyhow::Result<()> { + create_dir_all(self.output_directory.as_path())?; + + let global = self.output_directory.join("global_stats.csv"); + let mut writer = csv::Writer::from_path(&global)?; + for record in &self.global_records { + writer.serialize(record)?; + } + + info!("wrote global stats to {}", global.display()); + + for (sender_id, records) in self.records.iter() { + let output = self + .output_directory + .join(format!("sender{}.csv", sender_id)); + let mut writer = csv::Writer::from_path(&output)?; + for record in records { + writer.serialize(record)?; + } + info!("wrote client records to {}", output.display()); + } + Ok(()) + } + + pub(crate) async fn run(&mut self) { + let mut update_interval = interval(Duration::from_millis(500)); + + loop { + select! { + biased; + _ = self.shutdown.cancelled() => { + if let Err(err) = self.save_results_to_files() { + error!("failed to save measurement results to files: {err}") + } + break; + } + _ = update_interval.tick() => { + self.update_stats_span(); + } + } + } + } +} diff --git a/nym-node/src/throughput_tester/mod.rs b/nym-node/src/throughput_tester/mod.rs new file mode 100644 index 0000000000..68a7719fbd --- /dev/null +++ b/nym-node/src/throughput_tester/mod.rs @@ -0,0 +1,185 @@ +// Copyright 2025 - Nym Technologies SA +// SPDX-License-Identifier: GPL-3.0-only + +use crate::config::upgrade_helpers::try_load_current_config; +use crate::node::NymNode; +use crate::throughput_tester::client::ThroughputTestingClient; +use crate::throughput_tester::global_stats::GlobalStatsUpdater; +use crate::throughput_tester::stats::ClientStats; +use futures::future::join_all; +use human_repr::HumanDuration; +use indicatif::{ProgressState, ProgressStyle}; +use nym_crypto::asymmetric::x25519; +use nym_task::ShutdownToken; +use rand::{thread_rng, Rng}; +use std::net::{IpAddr, Ipv4Addr, SocketAddr}; +use std::path::PathBuf; +use std::sync::Arc; +use std::time::Duration; +use tokio::runtime; +use tokio::runtime::Runtime; +use tokio::time::sleep; +use tracing::{info, info_span, instrument}; +use tracing_indicatif::span_ext::IndicatifSpanExt; + +pub(crate) mod client; +pub(crate) mod global_stats; +mod stats; + +pub struct ThroughputTest { + node_runtime: Runtime, + clients_runtime: Runtime, +} + +impl ThroughputTest { + fn new(senders: usize) -> anyhow::Result { + Ok(ThroughputTest { + node_runtime: runtime::Builder::new_multi_thread() + .enable_all() + .thread_name("nym-node-pool") + .build()?, + clients_runtime: runtime::Builder::new_multi_thread() + .enable_all() + .worker_threads(senders) + .thread_name("testing-clients-pool") + .build()?, + }) + } + + fn prepare_nymnode(&self, config_path: PathBuf) -> anyhow::Result { + self.node_runtime.block_on(async { + let mut config = try_load_current_config(config_path).await?; + + // make sure to change bind address to localhost! + config + .mixnet + .bind_address + .set_ip(IpAddr::V4(Ipv4Addr::LOCALHOST)); + + let nym_node = NymNode::new(config).await?; + Ok(nym_node) + }) + } +} + +#[instrument( + skip_all, + fields( + sender_id = %sender_id + ) + +)] +#[allow(clippy::too_many_arguments)] +async fn run_testing_client( + sender_id: usize, + node_keys: Arc, + node_listener: SocketAddr, + packet_latency_threshold: Duration, + starting_sending_batch_size: usize, + starting_sending_delay: Duration, + stats: ClientStats, + shutdown_token: ShutdownToken, +) -> anyhow::Result<()> { + let _ = sender_id; + let client = ThroughputTestingClient::try_create( + starting_sending_delay, + starting_sending_batch_size, + packet_latency_threshold, + &node_keys, + node_listener, + stats, + shutdown_token, + ) + .await?; + + // wait a random amount of time before actually starting to desync the clients a bit + // (so they wouldn't update their rates at the same time) + let delay = Duration::from_millis(thread_rng().gen_range(10..200)); + info!( + "waiting for {} before attempting to start the processing loop", + delay.human_duration() + ); + sleep(delay).await; + + client.run().await +} + +pub(crate) fn test_mixing_throughput( + config_path: PathBuf, + senders: usize, + packet_latency_threshold: Duration, + starting_sending_batch_size: usize, + starting_sending_delay: Duration, + output_directory: PathBuf, +) -> anyhow::Result<()> { + let tester = ThroughputTest::new(senders)?; + + let nym_node = tester.prepare_nymnode(config_path)?; + let listener = nym_node.config().mixnet.bind_address; + + let sphinx_keys = nym_node.x25519_sphinx_keys(); + + let mut stats = Vec::with_capacity(senders); + for _ in 0..senders { + stats.push(ClientStats::default()) + } + + let header_span = info_span!("header"); + header_span.pb_set_style( + &ProgressStyle::with_template( + "testing mixing throughput of this machine... {wide_msg} {elapsed}\n{wide_bar}", + )? + .with_key( + "elapsed", + |state: &ProgressState, writer: &mut dyn std::fmt::Write| { + let _ = writer.write_str(&format!("{}", state.elapsed().human_duration())); + }, + ) + .progress_chars("---"), + ); + header_span.pb_start(); + + // Bit of a hack to show a full "-----" line underneath the header. + header_span.pb_set_length(1); + header_span.pb_set_position(1); + + let mut tasks_handles = Vec::new(); + + for (sender_id, stats) in stats.iter().enumerate() { + let token = nym_node.shutdown_token(format!("dummy-load-client-{sender_id}")); + + let client_future = run_testing_client( + sender_id, + sphinx_keys.clone(), + listener, + packet_latency_threshold, + starting_sending_batch_size, + starting_sending_delay, + stats.clone(), + token, + ); + let handle = tester.clients_runtime.spawn(client_future); + tasks_handles.push(handle); + } + + let mut global_stats = GlobalStatsUpdater::new( + header_span, + stats, + output_directory, + nym_node.shutdown_token("global-stats"), + ); + + let stats_handle = tester.clients_runtime.spawn(async move { + global_stats.run().await; + Ok(()) + }); + tasks_handles.push(stats_handle); + + tester + .node_runtime + .block_on(async move { nym_node.run_minimal_mixnet_processing().await })?; + + tester.clients_runtime.block_on(join_all(tasks_handles)); + + Ok(()) +} diff --git a/nym-node/src/throughput_tester/stats.rs b/nym-node/src/throughput_tester/stats.rs new file mode 100644 index 0000000000..b28a8942e5 --- /dev/null +++ b/nym-node/src/throughput_tester/stats.rs @@ -0,0 +1,58 @@ +// Copyright 2025 - Nym Technologies SA +// SPDX-License-Identifier: GPL-3.0-only + +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +#[derive(Clone, Default)] +pub(crate) struct ClientStats { + inner: Arc, +} + +impl ClientStats { + pub(crate) fn new_received(&self, new_measurement: u64) { + const ALPHA: f64 = 0.0005; + const ONE_SUB_ALPHA: f64 = 1.0 - ALPHA; + + let old_average_latency = self.inner.average_latency_nanos.load(Ordering::SeqCst); + + let new_average = if old_average_latency == 0 { + new_measurement + } else { + ((ALPHA * new_measurement as f64) + ONE_SUB_ALPHA * old_average_latency as f64) as u64 + }; + + self.inner.received.fetch_add(1, Ordering::SeqCst); + self.inner + .average_latency_nanos + .store(new_average, Ordering::SeqCst); + } + + pub(crate) fn new_sent_batch(&self, batch_size: usize) { + self.inner.sent.fetch_add(batch_size, Ordering::SeqCst); + } + + pub(crate) fn received(&self) -> usize { + self.inner.received.load(Ordering::SeqCst) + } + + pub(crate) fn sent(&self) -> usize { + self.inner.sent.load(Ordering::SeqCst) + } + + pub(crate) fn average_latency_nanos(&self) -> u64 { + self.inner.average_latency_nanos.load(Ordering::SeqCst) + } + + pub(crate) fn average_latency_duration(&self) -> Duration { + Duration::from_nanos(self.average_latency_nanos()) + } +} + +#[derive(Default)] +struct ClientStatsInner { + sent: AtomicUsize, + received: AtomicUsize, + average_latency_nanos: AtomicU64, +}