Extract out mixnet_listener.rs
This commit is contained in:
@@ -1,29 +1,19 @@
|
||||
#![cfg_attr(not(target_os = "linux"), allow(dead_code))]
|
||||
#![cfg_attr(not(target_os = "linux"), allow(unused_imports))]
|
||||
|
||||
use std::{collections::HashMap, net::IpAddr, path::Path};
|
||||
use std::path::Path;
|
||||
|
||||
use futures::{channel::oneshot, StreamExt};
|
||||
use futures::channel::oneshot;
|
||||
use nym_client_core::{
|
||||
client::mix_traffic::transceiver::GatewayTransceiver, HardcodedTopologyProvider,
|
||||
TopologyProvider,
|
||||
};
|
||||
use nym_ip_packet_requests::{
|
||||
DynamicConnectFailureReason, IpPacketRequest, IpPacketRequestData, IpPacketResponse,
|
||||
StaticConnectFailureReason,
|
||||
};
|
||||
use nym_sdk::mixnet::{InputMessage, MixnetMessageSender, Recipient};
|
||||
use nym_sphinx::receiver::ReconstructedMessage;
|
||||
use nym_task::{connections::TransmissionLane, TaskClient, TaskHandle};
|
||||
#[cfg(target_os = "linux")]
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use nym_sdk::mixnet::Recipient;
|
||||
use nym_task::{TaskClient, TaskHandle};
|
||||
|
||||
use crate::{
|
||||
constants::{CLIENT_INACTIVITY_TIMEOUT, DISCONNECT_TIMER_INTERVAL},
|
||||
error::IpPacketRouterError,
|
||||
request_filter::{self, RequestFilter},
|
||||
util::generate_new_ip,
|
||||
util::parse_ip::{parse_packet, ParsedPacket},
|
||||
Config,
|
||||
};
|
||||
|
||||
@@ -118,6 +108,8 @@ impl IpPacketRouterBuilder {
|
||||
#[cfg(target_os = "linux")]
|
||||
pub async fn run_service_provider(self) -> Result<(), IpPacketRouterError> {
|
||||
// Used to notify tasks to shutdown. Not all tasks fully supports this (yet).
|
||||
|
||||
use crate::{mixnet_listener, tun_listener};
|
||||
let task_handle: TaskHandle = self.shutdown.map(Into::into).unwrap_or_default();
|
||||
|
||||
// Connect to the mixnet
|
||||
@@ -146,7 +138,7 @@ impl IpPacketRouterBuilder {
|
||||
// TunListener
|
||||
let (connected_client_tx, connected_client_rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
|
||||
let tun_listener = crate::tun_listener::TunListener {
|
||||
let tun_listener = tun_listener::TunListener {
|
||||
tun_reader,
|
||||
mixnet_client_sender: mixnet_client.split_sender(),
|
||||
task_client: task_handle.get_handle(),
|
||||
@@ -158,7 +150,7 @@ impl IpPacketRouterBuilder {
|
||||
let request_filter = request_filter::RequestFilter::new(&self.config).await?;
|
||||
request_filter.start_update_tasks().await;
|
||||
|
||||
let ip_packet_router_service = MixnetListener {
|
||||
let mixnet_listener = mixnet_listener::MixnetListener {
|
||||
_config: self.config,
|
||||
request_filter: request_filter.clone(),
|
||||
tun_writer,
|
||||
@@ -181,320 +173,6 @@ impl IpPacketRouterBuilder {
|
||||
}
|
||||
}
|
||||
|
||||
ip_packet_router_service.run().await
|
||||
mixnet_listener.run().await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
struct MixnetListener {
|
||||
_config: Config,
|
||||
request_filter: request_filter::RequestFilter,
|
||||
tun_writer: tokio::io::WriteHalf<tokio_tun::Tun>,
|
||||
mixnet_client: nym_sdk::mixnet::MixnetClient,
|
||||
task_handle: TaskHandle,
|
||||
|
||||
connected_clients: HashMap<IpAddr, ConnectedClient>,
|
||||
connected_client_tx: tokio::sync::mpsc::UnboundedSender<ConnectedClientEvent>,
|
||||
}
|
||||
|
||||
pub(crate) struct ConnectedClient {
|
||||
pub(crate) nym_address: Recipient,
|
||||
pub(crate) last_activity: std::time::Instant,
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
impl MixnetListener {
|
||||
async fn on_static_connect_request(
|
||||
&mut self,
|
||||
connect_request: nym_ip_packet_requests::StaticConnectRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::info!(
|
||||
"Received static connect request from {sender_address}",
|
||||
sender_address = connect_request.reply_to
|
||||
);
|
||||
|
||||
let request_id = connect_request.request_id;
|
||||
let requested_ip = connect_request.ip;
|
||||
let reply_to = connect_request.reply_to;
|
||||
// TODO: ignoring reply_to_hops and reply_to_avg_mix_delays for now
|
||||
|
||||
// Check that the IP is available in the set of connected clients
|
||||
let is_ip_taken = self.connected_clients.contains_key(&requested_ip);
|
||||
|
||||
// Check that the nym address isn't already registered
|
||||
let is_nym_address_taken = self
|
||||
.connected_clients
|
||||
.values()
|
||||
.any(|client| client.nym_address == reply_to);
|
||||
|
||||
match (is_ip_taken, is_nym_address_taken) {
|
||||
(true, true) => {
|
||||
log::info!("Connecting an already connected client");
|
||||
// Update the last activity time for the client
|
||||
if let Some(client) = self.connected_clients.get_mut(&requested_ip) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::error!("Failed to update last activity time for client");
|
||||
}
|
||||
Ok(Some(IpPacketResponse::new_static_connect_success(
|
||||
request_id, reply_to,
|
||||
)))
|
||||
}
|
||||
(false, false) => {
|
||||
log::info!("Connecting a new client");
|
||||
self.connected_clients.insert(
|
||||
requested_ip,
|
||||
ConnectedClient {
|
||||
nym_address: reply_to,
|
||||
last_activity: std::time::Instant::now(),
|
||||
},
|
||||
);
|
||||
self.connected_client_tx
|
||||
.send(ConnectedClientEvent::Connect(
|
||||
requested_ip,
|
||||
Box::new(reply_to),
|
||||
))
|
||||
.unwrap();
|
||||
Ok(Some(IpPacketResponse::new_static_connect_success(
|
||||
request_id, reply_to,
|
||||
)))
|
||||
}
|
||||
(true, false) => {
|
||||
log::info!("Requested IP is not available");
|
||||
Ok(Some(IpPacketResponse::new_static_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
StaticConnectFailureReason::RequestedIpAlreadyInUse,
|
||||
)))
|
||||
}
|
||||
(false, true) => {
|
||||
log::info!("Nym address is already registered");
|
||||
Ok(Some(IpPacketResponse::new_static_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
StaticConnectFailureReason::RequestedNymAddressAlreadyInUse,
|
||||
)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn on_dynamic_connect_request(
|
||||
&mut self,
|
||||
connect_request: nym_ip_packet_requests::DynamicConnectRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::info!(
|
||||
"Received dynamic connect request from {sender_address}",
|
||||
sender_address = connect_request.reply_to
|
||||
);
|
||||
|
||||
let request_id = connect_request.request_id;
|
||||
let reply_to = connect_request.reply_to;
|
||||
// TODO: ignoring reply_to_hops and reply_to_avg_mix_delays for now
|
||||
|
||||
// Check if it's the same client connecting again, then we just reuse the same IP
|
||||
// TODO: this is problematic. Until we sign connect requests this means you can spam people
|
||||
// with return traffic
|
||||
let existing_ip = self.connected_clients.iter().find_map(|(ip, client)| {
|
||||
if client.nym_address == reply_to {
|
||||
Some(*ip)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
});
|
||||
|
||||
if let Some(existing_ip) = existing_ip {
|
||||
log::info!("Found existing client for nym address");
|
||||
// Update the last activity time for the client
|
||||
if let Some(client) = self.connected_clients.get_mut(&existing_ip) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::error!("Failed to update last activity time for client");
|
||||
}
|
||||
return Ok(Some(IpPacketResponse::new_dynamic_connect_success(
|
||||
request_id,
|
||||
reply_to,
|
||||
existing_ip,
|
||||
)));
|
||||
}
|
||||
|
||||
let Some(new_ip) = generate_new_ip::find_new_ip(&self.connected_clients) else {
|
||||
log::info!("No available IP address");
|
||||
return Ok(Some(IpPacketResponse::new_dynamic_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
DynamicConnectFailureReason::NoAvailableIp,
|
||||
)));
|
||||
};
|
||||
|
||||
self.connected_clients.insert(
|
||||
new_ip,
|
||||
ConnectedClient {
|
||||
nym_address: reply_to,
|
||||
last_activity: std::time::Instant::now(),
|
||||
},
|
||||
);
|
||||
self.connected_client_tx
|
||||
.send(ConnectedClientEvent::Connect(new_ip, Box::new(reply_to)))
|
||||
.unwrap();
|
||||
Ok(Some(IpPacketResponse::new_dynamic_connect_success(
|
||||
request_id, reply_to, new_ip,
|
||||
)))
|
||||
}
|
||||
|
||||
async fn on_data_request(
|
||||
&mut self,
|
||||
data_request: nym_ip_packet_requests::DataRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::trace!("Received data request");
|
||||
|
||||
// We don't forward packets that we are not able to parse. BUT, there might be a good
|
||||
// reason to still forward them.
|
||||
//
|
||||
// For example, if we are running in a mode where we are only supposed to forward
|
||||
// packets to a specific destination, we might want to forward them anyway.
|
||||
//
|
||||
// TODO: look into this
|
||||
let ParsedPacket {
|
||||
packet_type,
|
||||
src_addr,
|
||||
dst_addr,
|
||||
dst,
|
||||
} = parse_packet(&data_request.ip_packet)?;
|
||||
|
||||
let dst_str = dst.map_or(dst_addr.to_string(), |dst| dst.to_string());
|
||||
log::info!("Received packet: {packet_type}: {src_addr} -> {dst_str}");
|
||||
|
||||
// Check if there is a connected client for this src_addr. If there is, update the last activity time
|
||||
// for the client. If there isn't, drop the packet.
|
||||
if let Some(client) = self.connected_clients.get_mut(&src_addr) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::info!("Dropping packet: no connected client for {src_addr}");
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
// Filter check
|
||||
if let Some(dst) = dst {
|
||||
if !self.request_filter.check_address(&dst).await {
|
||||
log::warn!("Failed filter check: {dst}");
|
||||
// TODO: we could consider sending back a response here
|
||||
return Err(IpPacketRouterError::AddressFailedFilterCheck { addr: dst });
|
||||
}
|
||||
} else {
|
||||
// TODO: we should also filter packets without port number
|
||||
log::warn!("Ignoring filter check for packet without port number! TODO!");
|
||||
}
|
||||
|
||||
// TODO: consider changing from Vec<u8> to bytes::Bytes?
|
||||
let packet = data_request.ip_packet;
|
||||
self.tun_writer
|
||||
.write_all(&packet)
|
||||
.await
|
||||
.map_err(|_| IpPacketRouterError::FailedToWritePacketToTun)?;
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
async fn on_reconstructed_message(
|
||||
&mut self,
|
||||
reconstructed: ReconstructedMessage,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::debug!(
|
||||
"Received message with sender_tag: {:?}",
|
||||
reconstructed.sender_tag
|
||||
);
|
||||
|
||||
// Check version of request
|
||||
if let Some(version) = reconstructed.message.first() {
|
||||
// The idea is that in the future we can add logic here to parse older versions to stay
|
||||
// backwards compatible.
|
||||
if *version != nym_ip_packet_requests::CURRENT_VERSION {
|
||||
log::warn!("Received packet with invalid version");
|
||||
return Err(IpPacketRouterError::InvalidPacketVersion(*version));
|
||||
}
|
||||
}
|
||||
|
||||
let request = IpPacketRequest::from_reconstructed_message(&reconstructed)
|
||||
.map_err(|err| IpPacketRouterError::FailedToDeserializeTaggedPacket { source: err })?;
|
||||
|
||||
match request.data {
|
||||
IpPacketRequestData::StaticConnect(connect_request) => {
|
||||
self.on_static_connect_request(connect_request).await
|
||||
}
|
||||
IpPacketRequestData::DynamicConnect(connect_request) => {
|
||||
self.on_dynamic_connect_request(connect_request).await
|
||||
}
|
||||
IpPacketRequestData::Data(data_request) => self.on_data_request(data_request).await,
|
||||
}
|
||||
}
|
||||
|
||||
async fn run(mut self) -> Result<(), IpPacketRouterError> {
|
||||
let mut task_client = self.task_handle.fork("main_loop");
|
||||
let mut disconnect_timer = tokio::time::interval(DISCONNECT_TIMER_INTERVAL);
|
||||
|
||||
while !task_client.is_shutdown() {
|
||||
tokio::select! {
|
||||
_ = task_client.recv() => {
|
||||
log::debug!("IpPacketRouter [main loop]: received shutdown");
|
||||
},
|
||||
_ = disconnect_timer.tick() => {
|
||||
let now = std::time::Instant::now();
|
||||
let inactive_clients: Vec<IpAddr> = self.connected_clients.iter()
|
||||
.filter_map(|(ip, client)| {
|
||||
if now.duration_since(client.last_activity) > CLIENT_INACTIVITY_TIMEOUT {
|
||||
Some(*ip)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
for ip in inactive_clients {
|
||||
log::info!("Disconnect inactive client: {ip}");
|
||||
self.connected_clients.remove(&ip);
|
||||
self.connected_client_tx.send(ConnectedClientEvent::Disconnect(ip)).unwrap();
|
||||
}
|
||||
},
|
||||
msg = self.mixnet_client.next() => {
|
||||
if let Some(msg) = msg {
|
||||
match self.on_reconstructed_message(msg).await {
|
||||
Ok(Some(response)) => {
|
||||
let Some(recipient) = response.recipient() else {
|
||||
log::error!("IpPacketRouter [main loop]: failed to get recipient from response");
|
||||
continue;
|
||||
};
|
||||
let response_packet = response.to_bytes();
|
||||
let Ok(response_packet) = response_packet else {
|
||||
log::error!("Failed to serialize response packet");
|
||||
continue;
|
||||
};
|
||||
let lane = TransmissionLane::General;
|
||||
let packet_type = None;
|
||||
let input_message = InputMessage::new_regular(*recipient, response_packet, lane, packet_type);
|
||||
if let Err(err) = self.mixnet_client.send(input_message).await {
|
||||
log::error!("IpPacketRouter [main loop]: failed to send packet to mixnet: {err}");
|
||||
};
|
||||
},
|
||||
Ok(None) => {
|
||||
continue;
|
||||
},
|
||||
Err(err) => {
|
||||
log::error!("Error handling mixnet message: {err}");
|
||||
}
|
||||
|
||||
};
|
||||
} else {
|
||||
log::trace!("IpPacketRouter [main loop]: stopping since channel closed");
|
||||
break;
|
||||
};
|
||||
},
|
||||
|
||||
}
|
||||
}
|
||||
log::debug!("IpPacketRouter: stopping");
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) enum ConnectedClientEvent {
|
||||
Disconnect(IpAddr),
|
||||
Connect(IpAddr, Box<Recipient>),
|
||||
}
|
||||
|
||||
@@ -4,11 +4,12 @@
|
||||
pub use crate::config::Config;
|
||||
pub use ip_packet_router::{IpPacketRouterBuilder, OnStartData};
|
||||
|
||||
pub mod config;
|
||||
mod constants;
|
||||
pub mod error;
|
||||
mod ip_packet_router;
|
||||
mod mixnet_client;
|
||||
mod mixnet_listener;
|
||||
mod request_filter;
|
||||
mod tun_listener;
|
||||
mod util;
|
||||
pub mod config;
|
||||
pub mod error;
|
||||
|
||||
@@ -0,0 +1,335 @@
|
||||
use std::{collections::HashMap, net::IpAddr};
|
||||
|
||||
use futures::StreamExt;
|
||||
use nym_ip_packet_requests::{
|
||||
DynamicConnectFailureReason, IpPacketRequest, IpPacketRequestData, IpPacketResponse,
|
||||
StaticConnectFailureReason,
|
||||
};
|
||||
use nym_sdk::mixnet::{InputMessage, MixnetMessageSender, Recipient};
|
||||
use nym_sphinx::receiver::ReconstructedMessage;
|
||||
use nym_task::{connections::TransmissionLane, TaskHandle};
|
||||
#[cfg(target_os = "linux")]
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
use crate::{
|
||||
constants::{CLIENT_INACTIVITY_TIMEOUT, DISCONNECT_TIMER_INTERVAL},
|
||||
error::IpPacketRouterError,
|
||||
request_filter::{self},
|
||||
util::generate_new_ip,
|
||||
util::parse_ip::{parse_packet, ParsedPacket},
|
||||
Config,
|
||||
};
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
pub(crate) struct MixnetListener {
|
||||
pub(crate) _config: Config,
|
||||
pub(crate) request_filter: request_filter::RequestFilter,
|
||||
pub(crate) tun_writer: tokio::io::WriteHalf<tokio_tun::Tun>,
|
||||
pub(crate) mixnet_client: nym_sdk::mixnet::MixnetClient,
|
||||
pub(crate) task_handle: TaskHandle,
|
||||
|
||||
pub(crate) connected_clients: HashMap<IpAddr, ConnectedClient>,
|
||||
pub(crate) connected_client_tx: tokio::sync::mpsc::UnboundedSender<ConnectedClientEvent>,
|
||||
}
|
||||
|
||||
pub(crate) struct ConnectedClient {
|
||||
pub(crate) nym_address: Recipient,
|
||||
pub(crate) last_activity: std::time::Instant,
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
impl MixnetListener {
|
||||
async fn on_static_connect_request(
|
||||
&mut self,
|
||||
connect_request: nym_ip_packet_requests::StaticConnectRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::info!(
|
||||
"Received static connect request from {sender_address}",
|
||||
sender_address = connect_request.reply_to
|
||||
);
|
||||
|
||||
let request_id = connect_request.request_id;
|
||||
let requested_ip = connect_request.ip;
|
||||
let reply_to = connect_request.reply_to;
|
||||
// TODO: ignoring reply_to_hops and reply_to_avg_mix_delays for now
|
||||
|
||||
// Check that the IP is available in the set of connected clients
|
||||
let is_ip_taken = self.connected_clients.contains_key(&requested_ip);
|
||||
|
||||
// Check that the nym address isn't already registered
|
||||
let is_nym_address_taken = self
|
||||
.connected_clients
|
||||
.values()
|
||||
.any(|client| client.nym_address == reply_to);
|
||||
|
||||
match (is_ip_taken, is_nym_address_taken) {
|
||||
(true, true) => {
|
||||
log::info!("Connecting an already connected client");
|
||||
// Update the last activity time for the client
|
||||
if let Some(client) = self.connected_clients.get_mut(&requested_ip) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::error!("Failed to update last activity time for client");
|
||||
}
|
||||
Ok(Some(IpPacketResponse::new_static_connect_success(
|
||||
request_id, reply_to,
|
||||
)))
|
||||
}
|
||||
(false, false) => {
|
||||
log::info!("Connecting a new client");
|
||||
self.connected_clients.insert(
|
||||
requested_ip,
|
||||
ConnectedClient {
|
||||
nym_address: reply_to,
|
||||
last_activity: std::time::Instant::now(),
|
||||
},
|
||||
);
|
||||
self.connected_client_tx
|
||||
.send(ConnectedClientEvent::Connect(
|
||||
requested_ip,
|
||||
Box::new(reply_to),
|
||||
))
|
||||
.unwrap();
|
||||
Ok(Some(IpPacketResponse::new_static_connect_success(
|
||||
request_id, reply_to,
|
||||
)))
|
||||
}
|
||||
(true, false) => {
|
||||
log::info!("Requested IP is not available");
|
||||
Ok(Some(IpPacketResponse::new_static_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
StaticConnectFailureReason::RequestedIpAlreadyInUse,
|
||||
)))
|
||||
}
|
||||
(false, true) => {
|
||||
log::info!("Nym address is already registered");
|
||||
Ok(Some(IpPacketResponse::new_static_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
StaticConnectFailureReason::RequestedNymAddressAlreadyInUse,
|
||||
)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn on_dynamic_connect_request(
|
||||
&mut self,
|
||||
connect_request: nym_ip_packet_requests::DynamicConnectRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::info!(
|
||||
"Received dynamic connect request from {sender_address}",
|
||||
sender_address = connect_request.reply_to
|
||||
);
|
||||
|
||||
let request_id = connect_request.request_id;
|
||||
let reply_to = connect_request.reply_to;
|
||||
// TODO: ignoring reply_to_hops and reply_to_avg_mix_delays for now
|
||||
|
||||
// Check if it's the same client connecting again, then we just reuse the same IP
|
||||
// TODO: this is problematic. Until we sign connect requests this means you can spam people
|
||||
// with return traffic
|
||||
let existing_ip = self.connected_clients.iter().find_map(|(ip, client)| {
|
||||
if client.nym_address == reply_to {
|
||||
Some(*ip)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
});
|
||||
|
||||
if let Some(existing_ip) = existing_ip {
|
||||
log::info!("Found existing client for nym address");
|
||||
// Update the last activity time for the client
|
||||
if let Some(client) = self.connected_clients.get_mut(&existing_ip) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::error!("Failed to update last activity time for client");
|
||||
}
|
||||
return Ok(Some(IpPacketResponse::new_dynamic_connect_success(
|
||||
request_id,
|
||||
reply_to,
|
||||
existing_ip,
|
||||
)));
|
||||
}
|
||||
|
||||
let Some(new_ip) = generate_new_ip::find_new_ip(&self.connected_clients) else {
|
||||
log::info!("No available IP address");
|
||||
return Ok(Some(IpPacketResponse::new_dynamic_connect_failure(
|
||||
request_id,
|
||||
reply_to,
|
||||
DynamicConnectFailureReason::NoAvailableIp,
|
||||
)));
|
||||
};
|
||||
|
||||
self.connected_clients.insert(
|
||||
new_ip,
|
||||
ConnectedClient {
|
||||
nym_address: reply_to,
|
||||
last_activity: std::time::Instant::now(),
|
||||
},
|
||||
);
|
||||
self.connected_client_tx
|
||||
.send(ConnectedClientEvent::Connect(new_ip, Box::new(reply_to)))
|
||||
.unwrap();
|
||||
Ok(Some(IpPacketResponse::new_dynamic_connect_success(
|
||||
request_id, reply_to, new_ip,
|
||||
)))
|
||||
}
|
||||
|
||||
async fn on_data_request(
|
||||
&mut self,
|
||||
data_request: nym_ip_packet_requests::DataRequest,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::trace!("Received data request");
|
||||
|
||||
// We don't forward packets that we are not able to parse. BUT, there might be a good
|
||||
// reason to still forward them.
|
||||
//
|
||||
// For example, if we are running in a mode where we are only supposed to forward
|
||||
// packets to a specific destination, we might want to forward them anyway.
|
||||
//
|
||||
// TODO: look into this
|
||||
let ParsedPacket {
|
||||
packet_type,
|
||||
src_addr,
|
||||
dst_addr,
|
||||
dst,
|
||||
} = parse_packet(&data_request.ip_packet)?;
|
||||
|
||||
let dst_str = dst.map_or(dst_addr.to_string(), |dst| dst.to_string());
|
||||
log::info!("Received packet: {packet_type}: {src_addr} -> {dst_str}");
|
||||
|
||||
// Check if there is a connected client for this src_addr. If there is, update the last activity time
|
||||
// for the client. If there isn't, drop the packet.
|
||||
if let Some(client) = self.connected_clients.get_mut(&src_addr) {
|
||||
client.last_activity = std::time::Instant::now();
|
||||
} else {
|
||||
log::info!("Dropping packet: no connected client for {src_addr}");
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
// Filter check
|
||||
if let Some(dst) = dst {
|
||||
if !self.request_filter.check_address(&dst).await {
|
||||
log::warn!("Failed filter check: {dst}");
|
||||
// TODO: we could consider sending back a response here
|
||||
return Err(IpPacketRouterError::AddressFailedFilterCheck { addr: dst });
|
||||
}
|
||||
} else {
|
||||
// TODO: we should also filter packets without port number
|
||||
log::warn!("Ignoring filter check for packet without port number! TODO!");
|
||||
}
|
||||
|
||||
// TODO: consider changing from Vec<u8> to bytes::Bytes?
|
||||
let packet = data_request.ip_packet;
|
||||
self.tun_writer
|
||||
.write_all(&packet)
|
||||
.await
|
||||
.map_err(|_| IpPacketRouterError::FailedToWritePacketToTun)?;
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
async fn on_reconstructed_message(
|
||||
&mut self,
|
||||
reconstructed: ReconstructedMessage,
|
||||
) -> Result<Option<IpPacketResponse>, IpPacketRouterError> {
|
||||
log::debug!(
|
||||
"Received message with sender_tag: {:?}",
|
||||
reconstructed.sender_tag
|
||||
);
|
||||
|
||||
// Check version of request
|
||||
if let Some(version) = reconstructed.message.first() {
|
||||
// The idea is that in the future we can add logic here to parse older versions to stay
|
||||
// backwards compatible.
|
||||
if *version != nym_ip_packet_requests::CURRENT_VERSION {
|
||||
log::warn!("Received packet with invalid version");
|
||||
return Err(IpPacketRouterError::InvalidPacketVersion(*version));
|
||||
}
|
||||
}
|
||||
|
||||
let request = IpPacketRequest::from_reconstructed_message(&reconstructed)
|
||||
.map_err(|err| IpPacketRouterError::FailedToDeserializeTaggedPacket { source: err })?;
|
||||
|
||||
match request.data {
|
||||
IpPacketRequestData::StaticConnect(connect_request) => {
|
||||
self.on_static_connect_request(connect_request).await
|
||||
}
|
||||
IpPacketRequestData::DynamicConnect(connect_request) => {
|
||||
self.on_dynamic_connect_request(connect_request).await
|
||||
}
|
||||
IpPacketRequestData::Data(data_request) => self.on_data_request(data_request).await,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn run(mut self) -> Result<(), IpPacketRouterError> {
|
||||
let mut task_client = self.task_handle.fork("main_loop");
|
||||
let mut disconnect_timer = tokio::time::interval(DISCONNECT_TIMER_INTERVAL);
|
||||
|
||||
while !task_client.is_shutdown() {
|
||||
tokio::select! {
|
||||
_ = task_client.recv() => {
|
||||
log::debug!("IpPacketRouter [main loop]: received shutdown");
|
||||
},
|
||||
_ = disconnect_timer.tick() => {
|
||||
let now = std::time::Instant::now();
|
||||
let inactive_clients: Vec<IpAddr> = self.connected_clients.iter()
|
||||
.filter_map(|(ip, client)| {
|
||||
if now.duration_since(client.last_activity) > CLIENT_INACTIVITY_TIMEOUT {
|
||||
Some(*ip)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
for ip in inactive_clients {
|
||||
log::info!("Disconnect inactive client: {ip}");
|
||||
self.connected_clients.remove(&ip);
|
||||
self.connected_client_tx.send(ConnectedClientEvent::Disconnect(ip)).unwrap();
|
||||
}
|
||||
},
|
||||
msg = self.mixnet_client.next() => {
|
||||
if let Some(msg) = msg {
|
||||
match self.on_reconstructed_message(msg).await {
|
||||
Ok(Some(response)) => {
|
||||
let Some(recipient) = response.recipient() else {
|
||||
log::error!("IpPacketRouter [main loop]: failed to get recipient from response");
|
||||
continue;
|
||||
};
|
||||
let response_packet = response.to_bytes();
|
||||
let Ok(response_packet) = response_packet else {
|
||||
log::error!("Failed to serialize response packet");
|
||||
continue;
|
||||
};
|
||||
let lane = TransmissionLane::General;
|
||||
let packet_type = None;
|
||||
let input_message = InputMessage::new_regular(*recipient, response_packet, lane, packet_type);
|
||||
if let Err(err) = self.mixnet_client.send(input_message).await {
|
||||
log::error!("IpPacketRouter [main loop]: failed to send packet to mixnet: {err}");
|
||||
};
|
||||
},
|
||||
Ok(None) => {
|
||||
continue;
|
||||
},
|
||||
Err(err) => {
|
||||
log::error!("Error handling mixnet message: {err}");
|
||||
}
|
||||
|
||||
};
|
||||
} else {
|
||||
log::trace!("IpPacketRouter [main loop]: stopping since channel closed");
|
||||
break;
|
||||
};
|
||||
},
|
||||
|
||||
}
|
||||
}
|
||||
log::debug!("IpPacketRouter: stopping");
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) enum ConnectedClientEvent {
|
||||
Disconnect(IpAddr),
|
||||
Connect(IpAddr, Box<Recipient>),
|
||||
}
|
||||
@@ -7,7 +7,7 @@ use nym_task::{connections::TransmissionLane, TaskClient};
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::mpsc::UnboundedReceiver;
|
||||
|
||||
use crate::{error::IpPacketRouterError, ip_packet_router, util::parse_ip::parse_dst_addr};
|
||||
use crate::{error::IpPacketRouterError, mixnet_listener, util::parse_ip::parse_dst_addr};
|
||||
|
||||
// Reads packet from TUN and writes to mixnet client
|
||||
#[cfg(target_os = "linux")]
|
||||
@@ -17,8 +17,8 @@ pub(crate) struct TunListener {
|
||||
pub(crate) task_client: TaskClient,
|
||||
|
||||
// A mirror of the one in IpPacketRouter
|
||||
pub(crate) connected_clients: HashMap<IpAddr, ip_packet_router::ConnectedClient>,
|
||||
pub(crate) connected_client_rx: UnboundedReceiver<ip_packet_router::ConnectedClientEvent>,
|
||||
pub(crate) connected_clients: HashMap<IpAddr, mixnet_listener::ConnectedClient>,
|
||||
pub(crate) connected_client_rx: UnboundedReceiver<mixnet_listener::ConnectedClientEvent>,
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
@@ -31,14 +31,14 @@ impl TunListener {
|
||||
log::trace!("TunListener: received shutdown");
|
||||
},
|
||||
event = self.connected_client_rx.recv() => match event {
|
||||
Some(ip_packet_router::ConnectedClientEvent::Connect(ip, nym_addr)) => {
|
||||
Some(mixnet_listener::ConnectedClientEvent::Connect(ip, nym_addr)) => {
|
||||
log::trace!("Connect client: {ip}");
|
||||
self.connected_clients.insert(ip, ip_packet_router::ConnectedClient {
|
||||
self.connected_clients.insert(ip, mixnet_listener::ConnectedClient {
|
||||
nym_address: *nym_addr,
|
||||
last_activity: std::time::Instant::now(),
|
||||
});
|
||||
},
|
||||
Some(ip_packet_router::ConnectedClientEvent::Disconnect(ip)) => {
|
||||
Some(mixnet_listener::ConnectedClientEvent::Disconnect(ip)) => {
|
||||
log::trace!("Disconnect client: {ip}");
|
||||
self.connected_clients.remove(&ip);
|
||||
},
|
||||
|
||||
@@ -3,7 +3,7 @@ use std::{
|
||||
net::{IpAddr, Ipv4Addr},
|
||||
};
|
||||
|
||||
use crate::{constants::TUN_DEVICE_ADDRESS, ip_packet_router::ConnectedClient};
|
||||
use crate::{constants::TUN_DEVICE_ADDRESS, mixnet_listener::ConnectedClient};
|
||||
|
||||
// Find an available IP address in self.connected_clients
|
||||
// TODO: make this nicer
|
||||
|
||||
Reference in New Issue
Block a user