diff --git a/common/socks5-client-core/src/lib.rs b/common/socks5-client-core/src/lib.rs index daea560837..aa85a8a00f 100644 --- a/common/socks5-client-core/src/lib.rs +++ b/common/socks5-client-core/src/lib.rs @@ -23,6 +23,7 @@ use nym_client_core::client::key_manager::KeyManager; use nym_client_core::config::persistence::key_pathfinder::ClientKeyPathfinder; use nym_credential_storage::storage::Storage; use nym_sphinx::addressing::clients::Recipient; +use nym_sphinx::params::PacketSize; use nym_task::{TaskClient, TaskManager}; use nym_validator_client::nyxd::QueryNyxdClient; use nym_validator_client::Client; @@ -126,6 +127,14 @@ impl NymClient { .. } = client_status; + // FIXME: use correct value from the config once https://github.com/nymtech/nym/pull/3217 + // is merged + let packet_size = config + .get_base() + .get_use_extended_packet_size() + .map(Into::into) + .unwrap_or(PacketSize::RegularPacket); + let authenticator = Authenticator::new(auth_methods, allowed_users); let mut sphinx_socks = SphinxSocksServer::new( socks5_config.get_listening_port(), @@ -134,6 +143,7 @@ impl NymClient { self_address, shared_lane_queue_lengths, socks::client::Config::new( + packet_size, socks5_config.get_provider_interface_version(), socks5_config.get_socks5_protocol_version(), socks5_config.get_send_anonymously(), diff --git a/common/socks5-client-core/src/socks/client.rs b/common/socks5-client-core/src/socks/client.rs index e6f19c62d4..d1d4b10697 100644 --- a/common/socks5-client-core/src/socks/client.rs +++ b/common/socks5-client-core/src/socks/client.rs @@ -17,6 +17,7 @@ use nym_socks5_requests::{ ConnectionId, RemoteAddress, Socks5ProtocolVersion, Socks5ProviderRequest, Socks5Request, }; use nym_sphinx::addressing::clients::Recipient; +use nym_sphinx::params::PacketSize; use nym_task::connections::{LaneQueueLengths, TransmissionLane}; use nym_task::TaskClient; use pin_project::pin_project; @@ -131,6 +132,7 @@ impl AsyncWrite for StreamState { #[derive(Debug, Copy, Clone)] pub(crate) struct Config { + biggest_packet_size: PacketSize, provider_interface_version: ProviderInterfaceVersion, socks5_protocol_version: Socks5ProtocolVersion, use_surbs_for_responses: bool, @@ -140,6 +142,7 @@ pub(crate) struct Config { impl Config { pub(crate) fn new( + biggest_packet_size: PacketSize, provider_interface_version: ProviderInterfaceVersion, socks5_protocol_version: Socks5ProtocolVersion, use_surbs_for_responses: bool, @@ -147,6 +150,7 @@ impl Config { per_request_surbs: u32, ) -> Self { Self { + biggest_packet_size, provider_interface_version, socks5_protocol_version, use_surbs_for_responses, @@ -410,6 +414,9 @@ impl SocksClient { remote_proxy_target, conn_receiver, input_sender, + // FIXME: this does NOT include overhead due to acks or chunking + // (so actual true plaintext is smaller) + self.config.biggest_packet_size.plaintext_size(), connection_id, Some(self.lane_queue_lengths.clone()), self.shutdown_listener.clone(), diff --git a/common/socks5/proxy-helpers/src/available_reader.rs b/common/socks5/proxy-helpers/src/available_reader.rs index 018486e68b..1f169678ac 100644 --- a/common/socks5/proxy-helpers/src/available_reader.rs +++ b/common/socks5/proxy-helpers/src/available_reader.rs @@ -14,7 +14,7 @@ use tokio::io::AsyncRead; use tokio::time::{sleep, Duration, Instant, Sleep}; use tokio_util::io::poll_read_buf; -const MAX_READ_AMOUNT: usize = 500 * 1000; // 0.5MB +const DEFAULT_MAX_READ_AMOUNT: usize = 128 * 1024; // 128kB const GRACE_DURATION: Duration = Duration::from_millis(1); const READ_TIMEOUT: Duration = Duration::from_millis(10); @@ -25,6 +25,7 @@ pub struct AvailableReader<'a, R: AsyncRead + Unpin> { inner: RefCell<&'a mut R>, grace_period: Option>>, read_deadline: Option>>, + max_read: usize, } impl<'a, R> AvailableReader<'a, R> @@ -33,11 +34,11 @@ where { const BUF_INCREMENT: usize = 4096; - pub fn new(reader: &'a mut R) -> Self { - // pub fn new(reader: &'a mut R, buffer_size: usize) -> Self { + pub fn new(reader: &'a mut R, max_read: Option) -> Self { AvailableReader { buf: RefCell::new(BytesMut::with_capacity(Self::BUF_INCREMENT)), inner: RefCell::new(reader), + max_read: max_read.unwrap_or(DEFAULT_MAX_READ_AMOUNT), grace_period: None, read_deadline: None, } @@ -141,7 +142,7 @@ impl<'a, R: AsyncRead + Unpin> Stream for AvailableReader<'a, R> { // if we reached our maximum amount or we've been trying to read the data for too long // return what we have let read_bytes_len = self.buf.borrow().len(); - if read_bytes_len >= MAX_READ_AMOUNT || deadline_poll_res.is_ready() { + if read_bytes_len >= self.max_read || deadline_poll_res.is_ready() { self.return_buf() } else { Poll::Pending @@ -166,7 +167,7 @@ mod tests { let data = vec![42u8; 100]; let mut reader = Cursor::new(data.clone()); - let mut available_reader = AvailableReader::new(&mut reader); + let mut available_reader = AvailableReader::new(&mut reader, None); let read_data = available_reader.next().await.unwrap().unwrap(); assert_eq!(read_data, data); @@ -178,7 +179,7 @@ mod tests { let data = vec![42u8; AvailableReader::>>::BUF_INCREMENT + 100]; let mut reader = Cursor::new(data.clone()); - let mut available_reader = AvailableReader::new(&mut reader); + let mut available_reader = AvailableReader::new(&mut reader, None); let read_data = available_reader.next().await.unwrap().unwrap(); assert_eq!(read_data, data); @@ -196,7 +197,7 @@ mod tests { .read(&second_data_chunk) .build(); - let mut available_reader = AvailableReader::new(&mut reader_mock); + let mut available_reader = AvailableReader::new(&mut reader_mock, None); let read_data = available_reader.next().await.unwrap().unwrap(); assert_eq!(read_data, first_data_chunk); @@ -215,7 +216,7 @@ mod tests { .read(&data) .build(); - let mut available_reader = AvailableReader::new(&mut reader_mock); + let mut available_reader = AvailableReader::new(&mut reader_mock, None); let read_data = available_reader.next().await.unwrap().unwrap(); assert_eq!(read_data, data); diff --git a/common/socks5/proxy-helpers/src/proxy_runner/inbound.rs b/common/socks5/proxy-helpers/src/proxy_runner/inbound.rs index 548d1a0971..584ace090f 100644 --- a/common/socks5/proxy-helpers/src/proxy_runner/inbound.rs +++ b/common/socks5/proxy-helpers/src/proxy_runner/inbound.rs @@ -167,6 +167,7 @@ pub(super) async fn run_inbound( remote_source_address: String, connection_id: ConnectionId, mix_sender: MixProxySender, + available_plaintext_per_mix_packet: usize, adapter_fn: F, shutdown_notify: Arc, lane_queue_lengths: Option, @@ -176,7 +177,9 @@ where F: Fn(ConnectionId, Vec, bool) -> S + Send + 'static, S: Debug, { - let mut available_reader = AvailableReader::new(&mut reader); + // TODO: this multiplication by 4 is completely arbitrary here + let mut available_reader = + AvailableReader::new(&mut reader, Some(available_plaintext_per_mix_packet * 4)); let mut message_sender = OrderedMessageSender::new(); let shutdown_future = shutdown_notify.notified().then(|_| sleep(SHUTDOWN_TIMEOUT)); diff --git a/common/socks5/proxy-helpers/src/proxy_runner/mod.rs b/common/socks5/proxy-helpers/src/proxy_runner/mod.rs index c9d925ad4f..c8978141bf 100644 --- a/common/socks5/proxy-helpers/src/proxy_runner/mod.rs +++ b/common/socks5/proxy-helpers/src/proxy_runner/mod.rs @@ -1,4 +1,4 @@ -// Copyright 2021 - Nym Technologies SA +// Copyright 2021-2023 - Nym Technologies SA // SPDX-License-Identifier: Apache-2.0 use crate::connection_controller::ConnectionReceiver; @@ -49,6 +49,8 @@ pub struct ProxyRunner { connection_id: ConnectionId, lane_queue_lengths: Option, + available_plaintext_per_mix_packet: usize, + // Listens to shutdown commands from higher up shutdown_listener: TaskClient, } @@ -64,6 +66,7 @@ where remote_source_address: String, mix_receiver: ConnectionReceiver, mix_sender: MixProxySender, + available_plaintext_per_mix_packet: usize, connection_id: ConnectionId, lane_queue_lengths: Option, shutdown_listener: TaskClient, @@ -76,6 +79,7 @@ where remote_source_address, connection_id, lane_queue_lengths, + available_plaintext_per_mix_packet, shutdown_listener, } } @@ -96,6 +100,7 @@ where self.remote_source_address.clone(), self.connection_id, self.mix_sender.clone(), + self.available_plaintext_per_mix_packet, adapter_fn, Arc::clone(&shutdown_notify), self.lane_queue_lengths.clone(), diff --git a/service-providers/network-requester/src/core.rs b/service-providers/network-requester/src/core.rs index d8665a1306..39060e9e55 100644 --- a/service-providers/network-requester/src/core.rs +++ b/service-providers/network-requester/src/core.rs @@ -29,6 +29,7 @@ use nym_socks5_requests::{ }; use nym_sphinx::addressing::clients::Recipient; use nym_sphinx::anonymous_replies::requests::AnonymousSenderTag; +use nym_sphinx::params::PacketSize; use nym_statistics_common::collector::StatisticsSender; use nym_task::connections::LaneQueueLengths; use nym_task::{TaskClient, TaskManager}; @@ -57,6 +58,8 @@ pub struct NRServiceProviderBuilder { } struct NRServiceProvider { + config: Config, + outbound_request_filter: OutboundRequestFilter, open_proxy: bool, mixnet_client: nym_sdk::mixnet::MixnetClient, @@ -242,6 +245,7 @@ impl NRServiceProviderBuilder { start_allowed_list_reloader(self.allowed_hosts, shutdown.subscribe()).await; let service_provider = NRServiceProvider { + config: self.config, outbound_request_filter: self.outbound_request_filter, open_proxy: self.open_proxy, mixnet_client, @@ -329,6 +333,7 @@ impl NRServiceProvider { connection_id: ConnectionId, remote_addr: String, return_address: reply::MixnetAddress, + biggest_packet_size: PacketSize, controller_sender: ControllerSender, mix_input_sender: MixProxySender, lane_queue_lengths: LaneQueueLengths, @@ -385,6 +390,7 @@ impl NRServiceProvider { // run the proxy on the connection conn.run_proxy( remote_version, + biggest_packet_size, mix_receiver, mix_input_sender, lane_queue_lengths, @@ -437,6 +443,15 @@ impl NRServiceProvider { return; } + // FIXME: use correct value from the config once https://github.com/nymtech/nym/pull/3217 + // is merged + let packet_size = self + .config + .get_base() + .get_use_extended_packet_size() + .map(Into::into) + .unwrap_or(PacketSize::RegularPacket); + let controller_sender_clone = self.controller_sender.clone(); let mix_input_sender_clone = self.mix_input_sender.clone(); let lane_queue_lengths_clone = self.mixnet_client.shared_lane_queue_lengths(); @@ -449,6 +464,7 @@ impl NRServiceProvider { conn_id, remote_addr, return_address, + packet_size, controller_sender_clone, mix_input_sender_clone, lane_queue_lengths_clone, diff --git a/service-providers/network-requester/src/socks5/tcp.rs b/service-providers/network-requester/src/socks5/tcp.rs index 18281b0fb7..2a50bbf3a3 100644 --- a/service-providers/network-requester/src/socks5/tcp.rs +++ b/service-providers/network-requester/src/socks5/tcp.rs @@ -7,6 +7,7 @@ use nym_service_providers_common::interface::RequestVersion; use nym_socks5_proxy_helpers::connection_controller::ConnectionReceiver; use nym_socks5_proxy_helpers::proxy_runner::{MixProxySender, ProxyRunner}; use nym_socks5_requests::{ConnectionId, RemoteAddress, Socks5Request}; +use nym_sphinx::params::PacketSize; use nym_task::connections::LaneQueueLengths; use nym_task::TaskClient; use std::io; @@ -42,6 +43,7 @@ impl Connection { pub(crate) async fn run_proxy( &mut self, remote_version: RequestVersion, + biggest_packet_size: PacketSize, mix_receiver: ConnectionReceiver, mix_sender: MixProxySender, lane_queue_lengths: LaneQueueLengths, @@ -57,6 +59,9 @@ impl Connection { remote_source_address, mix_receiver, mix_sender, + // FIXME: this does NOT include overhead due to acks or chunking + // (so actual true plaintext is smaller) + biggest_packet_size.plaintext_size(), connection_id, Some(lane_queue_lengths), shutdown,