-
Notifications
You must be signed in to change notification settings - Fork 237
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
22 changed files
with
758 additions
and
397 deletions.
There are no files selected for viewing
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,9 +1,8 @@ | ||
// Copyright 2022-2023 - Nym Technologies SA <[email protected]> | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
use super::packet_statistics_control::PacketStatisticsReporter; | ||
use super::received_buffer::ReceivedBufferMessage; | ||
use super::statistics::StatisticsControl; | ||
use super::statistics_control::StatisticsControl; | ||
use super::topology_control::geo_aware_provider::GeoAwareTopologyProvider; | ||
use crate::client::base_client::storage::helpers::store_client_keys; | ||
use crate::client::base_client::storage::MixnetClientStorage; | ||
|
@@ -13,7 +12,6 @@ use crate::client::key_manager::persistence::KeyStore; | |
use crate::client::key_manager::ClientKeys; | ||
use crate::client::mix_traffic::transceiver::{GatewayReceiver, GatewayTransceiver, RemoteGateway}; | ||
use crate::client::mix_traffic::{BatchMixMessageSender, MixTrafficController}; | ||
use crate::client::packet_statistics_control::PacketStatisticsControl; | ||
use crate::client::real_messages_control; | ||
use crate::client::real_messages_control::RealMessagesController; | ||
use crate::client::received_buffer::{ | ||
|
@@ -50,7 +48,7 @@ use nym_sphinx::addressing::clients::Recipient; | |
use nym_sphinx::addressing::nodes::NodeIdentity; | ||
use nym_sphinx::params::PacketType; | ||
use nym_sphinx::receiver::{ReconstructedMessage, SphinxMessageReceiver}; | ||
use nym_statistics_common::clients::ClientStatsReporter; | ||
use nym_statistics_common::clients::ClientStatsSender; | ||
use nym_task::connections::{ConnectionCommandReceiver, ConnectionCommandSender, LaneQueueLengths}; | ||
use nym_task::{TaskClient, TaskHandle}; | ||
use nym_topology::provider_trait::TopologyProvider; | ||
|
@@ -61,6 +59,7 @@ use std::fmt::Debug; | |
use std::os::raw::c_int as RawFd; | ||
use std::path::Path; | ||
use std::sync::Arc; | ||
use tokio::sync::mpsc::Sender; | ||
use url::Url; | ||
|
||
#[cfg(all( | ||
|
@@ -283,7 +282,7 @@ where | |
self_address: Recipient, | ||
topology_accessor: TopologyAccessor, | ||
mix_tx: BatchMixMessageSender, | ||
stats_tx: PacketStatisticsReporter, | ||
stats_tx: ClientStatsSender, | ||
shutdown: TaskClient, | ||
) { | ||
info!("Starting loop cover traffic stream..."); | ||
|
@@ -316,7 +315,7 @@ where | |
client_connection_rx: ConnectionCommandReceiver, | ||
shutdown: TaskClient, | ||
packet_type: PacketType, | ||
stats_tx: PacketStatisticsReporter, | ||
stats_tx: ClientStatsSender, | ||
) { | ||
info!("Starting real traffic stream..."); | ||
|
||
|
@@ -345,7 +344,7 @@ where | |
reply_key_storage: SentReplyKeys, | ||
reply_controller_sender: ReplyControllerSender, | ||
shutdown: TaskClient, | ||
packet_statistics_control: PacketStatisticsReporter, | ||
metrics_reporter: ClientStatsSender, | ||
) { | ||
info!("Starting received messages buffer controller..."); | ||
let controller: ReceivedMessagesBufferController<SphinxMessageReceiver> = | ||
|
@@ -355,7 +354,7 @@ where | |
mixnet_receiver, | ||
reply_key_storage, | ||
reply_controller_sender, | ||
packet_statistics_control, | ||
metrics_reporter, | ||
); | ||
controller.start_with_shutdown(shutdown) | ||
} | ||
|
@@ -596,11 +595,21 @@ where | |
Ok(()) | ||
} | ||
|
||
fn start_packet_statistics_control(shutdown: TaskClient) -> PacketStatisticsReporter { | ||
fn start_statistics_control( | ||
stats_reporting_addres: Option<Recipient>, | ||
input_sender: Sender<InputMessage>, | ||
shutdown: TaskClient, | ||
) -> ClientStatsSender { | ||
info!("Starting packet statistics control..."); | ||
let (packet_statistics_control, packet_stats_reporter) = PacketStatisticsControl::new(); | ||
packet_statistics_control.start_with_shutdown(shutdown); | ||
packet_stats_reporter | ||
match stats_reporting_addres { | ||
Some(reporting_address) => { | ||
let (stats_control, stats_reporter) = | ||
StatisticsControl::new(reporting_address, input_sender.clone()); | ||
stats_control.start_with_shutdown(shutdown.fork("statistics_control")); | ||
stats_reporter | ||
} | ||
None => ClientStatsSender::sink(), | ||
} | ||
} | ||
|
||
fn start_mix_traffic_controller( | ||
|
@@ -730,6 +739,12 @@ where | |
self.user_agent.clone(), | ||
); | ||
|
||
let stats_reporter = Self::start_statistics_control( | ||
self.stats_reporting_address, | ||
input_sender.clone(), | ||
shutdown.fork("statistics_control"), | ||
); | ||
|
||
// needs to be started as the first thing to block if required waiting for the gateway | ||
Self::start_topology_refresher( | ||
topology_provider, | ||
|
@@ -741,9 +756,6 @@ where | |
) | ||
.await?; | ||
|
||
let packet_stats_reporter = | ||
Self::start_packet_statistics_control(shutdown.fork("packet_statistics_control")); | ||
|
||
let gateway_packet_router = PacketRouter::new( | ||
ack_sender, | ||
mixnet_messages_sender, | ||
|
@@ -775,7 +787,7 @@ where | |
reply_storage.key_storage(), | ||
reply_controller_sender.clone(), | ||
shutdown.fork("received_messages_buffer"), | ||
packet_stats_reporter.clone(), | ||
stats_reporter.clone(), | ||
); | ||
|
||
// The message_sender is the transmitter for any component generating sphinx packets | ||
|
@@ -814,7 +826,7 @@ where | |
client_connection_rx, | ||
shutdown.fork("real_traffic_controller"), | ||
self.config.debug.traffic.packet_type, | ||
packet_stats_reporter.clone(), | ||
stats_reporter.clone(), | ||
); | ||
|
||
if !self | ||
|
@@ -829,21 +841,11 @@ where | |
self_address, | ||
shared_topology_accessor.clone(), | ||
message_sender, | ||
packet_stats_reporter, | ||
stats_reporter.clone(), | ||
shutdown.fork("cover_traffic_stream"), | ||
); | ||
} | ||
|
||
let stats_reporter = match self.stats_reporting_address { | ||
Some(reporting_address) => { | ||
let (stats_control, stats_reporter) = | ||
StatisticsControl::new(reporting_address, input_sender.clone()); | ||
stats_control.start_with_shutdown(shutdown.fork("statistics_control")); | ||
Some(stats_reporter) | ||
} | ||
None => None, | ||
}; | ||
|
||
debug!("Core client startup finished!"); | ||
debug!("The address of this client is: {self_address}"); | ||
|
||
|
@@ -879,7 +881,7 @@ pub struct BaseClient { | |
pub client_input: ClientInputStatus, | ||
pub client_output: ClientOutputStatus, | ||
pub client_state: ClientState, | ||
pub stats_reporter: Option<ClientStatsReporter>, | ||
pub stats_reporter: ClientStatsSender, | ||
|
||
pub task_handle: TaskHandle, | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,9 +1,9 @@ | ||
// Copyright 2021 - Nym Technologies SA <[email protected]> | ||
// SPDX-License-Identifier: Apache-2.0 | ||
|
||
use crate::client::packet_statistics_control::{PacketStatisticsEvent, PacketStatisticsReporter}; | ||
|
||
use super::action_controller::{AckActionSender, Action}; | ||
use nym_statistics_common::clients::{packet_statistics::PacketStatisticsEvent, ClientStatsSender}; | ||
|
||
use futures::StreamExt; | ||
use log::*; | ||
use nym_gateway_client::AcknowledgementReceiver; | ||
|
@@ -19,15 +19,15 @@ pub(super) struct AcknowledgementListener { | |
ack_key: Arc<AckKey>, | ||
ack_receiver: AcknowledgementReceiver, | ||
action_sender: AckActionSender, | ||
stats_tx: PacketStatisticsReporter, | ||
stats_tx: ClientStatsSender, | ||
} | ||
|
||
impl AcknowledgementListener { | ||
pub(super) fn new( | ||
ack_key: Arc<AckKey>, | ||
ack_receiver: AcknowledgementReceiver, | ||
action_sender: AckActionSender, | ||
stats_tx: PacketStatisticsReporter, | ||
stats_tx: ClientStatsSender, | ||
) -> Self { | ||
AcknowledgementListener { | ||
ack_key, | ||
|
@@ -40,7 +40,7 @@ impl AcknowledgementListener { | |
async fn on_ack(&mut self, ack_content: Vec<u8>) { | ||
trace!("Received an ack"); | ||
self.stats_tx | ||
.report(PacketStatisticsEvent::AckReceived(ack_content.len())); | ||
.report(PacketStatisticsEvent::AckReceived(ack_content.len()).into()); | ||
|
||
let frag_id = match recover_identifier(&self.ack_key, &ack_content) | ||
.map(FragmentIdentifier::try_from_bytes) | ||
|
@@ -57,13 +57,13 @@ impl AcknowledgementListener { | |
if frag_id == COVER_FRAG_ID { | ||
trace!("Received an ack for a cover message - no need to do anything"); | ||
self.stats_tx | ||
.report(PacketStatisticsEvent::CoverAckReceived(ack_content.len())); | ||
.report(PacketStatisticsEvent::CoverAckReceived(ack_content.len()).into()); | ||
return; | ||
} | ||
|
||
trace!("Received {} from the mix network", frag_id); | ||
self.stats_tx | ||
.report(PacketStatisticsEvent::RealAckReceived(ack_content.len())); | ||
.report(PacketStatisticsEvent::RealAckReceived(ack_content.len()).into()); | ||
self.action_sender | ||
.unbounded_send(Action::new_remove(frag_id)) | ||
.unwrap(); | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.