diff --git a/service-providers/ip-packet-router/src/connected_client_handler.rs b/service-providers/ip-packet-router/src/connected_client_handler.rs new file mode 100644 index 0000000000..6b955a7b70 --- /dev/null +++ b/service-providers/ip-packet-router/src/connected_client_handler.rs @@ -0,0 +1,92 @@ +use std::{collections::HashMap, net::IpAddr}; + +use nym_ip_packet_requests::response::IpPacketResponse; +use nym_sdk::mixnet::{MixnetMessageSender, Recipient}; +use nym_task::TaskClient; +#[cfg(target_os = "linux")] +use tokio::io::AsyncReadExt; + +use crate::{ + error::{IpPacketRouterError, Result}, + mixnet_listener::{self}, + util::{create_message::create_input_message, parse_ip::parse_dst_addr}, +}; + +// Data flow +// Out: mixnet_listener -> decode -> handle_packet -> write_to_tun +// In: tun_listener -> [connected_client_handler -> encode] -> mixnet_sender + +// This handler is spawned as a task, and it listens to IP packets passed from the tun_listener, +// encodes it, and then sends to mixnet. +pub(crate) struct ConnectedClientHandler { + nym_address: Recipient, + mix_hops: Option, + tun_rx: tokio::sync::mpsc::UnboundedReceiver>, + mixnet_client_sender: nym_sdk::mixnet::MixnetClientSender, + close_rx: tokio::sync::oneshot::Receiver<()>, +} + +impl ConnectedClientHandler { + pub(crate) fn new( + nym_address: Recipient, + mix_hops: Option, + tun_rx: tokio::sync::mpsc::UnboundedReceiver>, + mixnet_client_sender: nym_sdk::mixnet::MixnetClientSender, + close_rx: tokio::sync::oneshot::Receiver<()>, + ) -> Self { + ConnectedClientHandler { + nym_address, + mix_hops, + tun_rx, + mixnet_client_sender, + close_rx, + } + } + + async fn handle_packet(&mut self, packet: Vec) -> Result<()> { + let response_packet = IpPacketResponse::new_ip_packet(packet.into()) + .to_bytes() + .map_err(|err| IpPacketRouterError::FailedToSerializeResponsePacket { source: err })?; + let input_message = create_input_message(self.nym_address, response_packet, self.mix_hops); + + self.mixnet_client_sender + .send(input_message) + .await + .map_err(|err| IpPacketRouterError::FailedToSendPacketToMixnet { source: err })?; + + Ok(()) + } + + async fn run(mut self) -> Result<()> { + loop { + tokio::select! { + _ = &mut self.close_rx => { + log::warn!("ConnectedClientHandler: received shutdown"); + break; + }, + packet = self.tun_rx.recv() => match packet { + Some(packet) => { + if let Err(err) = self.handle_packet(packet).await { + log::error!("connected client handler: failed to handle packet: {err}"); + } + }, + None => { + log::error!("connected client handler: tun channel closed"); + break; + } + } + } + } + + log::warn!("ConnectedClientHandler: exiting"); + Ok(()) + } + + pub(crate) fn start(self) { + tokio::spawn(async move { + if let Err(err) = self.run().await { + log::error!("connected client handler has failed: {err}") + } + }); + } +} diff --git a/service-providers/ip-packet-router/src/lib.rs b/service-providers/ip-packet-router/src/lib.rs index f8aa663be8..3d464b05c3 100644 --- a/service-providers/ip-packet-router/src/lib.rs +++ b/service-providers/ip-packet-router/src/lib.rs @@ -5,6 +5,7 @@ pub use crate::config::Config; pub use ip_packet_router::{IpPacketRouter, OnStartData}; pub mod config; +mod connected_client_handler; mod constants; pub mod error; mod ip_packet_router; diff --git a/service-providers/ip-packet-router/src/mixnet_listener.rs b/service-providers/ip-packet-router/src/mixnet_listener.rs index 0031683e53..c380b78660 100644 --- a/service-providers/ip-packet-router/src/mixnet_listener.rs +++ b/service-providers/ip-packet-router/src/mixnet_listener.rs @@ -19,9 +19,11 @@ use tokio::io::AsyncWriteExt; use crate::{ config::Config, + connected_client_handler, constants::{CLIENT_INACTIVITY_TIMEOUT, DISCONNECT_TIMER_INTERVAL}, error::{IpPacketRouterError, Result}, request_filter::{self}, + tun_listener, util::generate_new_ip, util::{ create_message::create_input_message, @@ -44,54 +46,15 @@ pub(crate) struct ConnectedClients { connected_client_tx: tokio::sync::mpsc::UnboundedSender, } -pub(crate) struct ConnectedClientsListener { - clients: HashMap, - pub(crate) connected_client_rx: tokio::sync::mpsc::UnboundedReceiver, -} - -impl ConnectedClientsListener { - pub(crate) fn get(&self, ip: &IpAddr) -> Option<&ConnectedClient> { - self.clients.get(ip) - } - - pub(crate) fn update(&mut self, event: ConnectedClientEvent) { - match event { - ConnectedClientEvent::Connect(connected_event) => { - let ConnectEvent { - ip, - nym_address, - mix_hops, - } = *connected_event; - log::trace!("Connect client: {ip}"); - self.clients.insert( - ip, - ConnectedClient { - nym_address, - mix_hops, - last_activity: std::time::Instant::now(), - }, - ); - } - ConnectedClientEvent::Disconnect(DisconnectEvent(ip)) => { - log::trace!("Disconnect client: {ip}"); - self.clients.remove(&ip); - } - } - } -} - impl ConnectedClients { - pub(crate) fn new() -> (Self, ConnectedClientsListener) { + pub(crate) fn new() -> (Self, tun_listener::ConnectedClientsListener) { let (connected_client_tx, connected_client_rx) = tokio::sync::mpsc::unbounded_channel(); ( Self { clients: Default::default(), connected_client_tx, }, - ConnectedClientsListener { - clients: Default::default(), - connected_client_rx, - }, + tun_listener::ConnectedClientsListener::new(connected_client_rx), ) } @@ -125,20 +88,33 @@ impl ConnectedClients { .find(|client| client.nym_address == *nym_address) } - fn connect(&mut self, ip: IpAddr, nym_address: Recipient, mix_hops: Option) { + fn connect( + &mut self, + ip: IpAddr, + nym_address: Recipient, + mix_hops: Option, + forward_from_tun_tx: tokio::sync::mpsc::UnboundedSender>, + close_tx: tokio::sync::oneshot::Sender<()>, + ) { + // The map of connected clients that the mixnet listener keeps track of. It monitors + // activity and disconnects clients that have been inactive for too long. self.clients.insert( ip, ConnectedClient { nym_address, mix_hops, last_activity: std::time::Instant::now(), + close_tx: Some(close_tx), }, ); + // Send the connected client info to the tun listener, which will use it to forward packets + // to the connected client handler. self.connected_client_tx .send(ConnectedClientEvent::Connect(Box::new(ConnectEvent { ip, nym_address, mix_hops, + forward_from_tun_tx, }))) .unwrap(); } @@ -183,6 +159,19 @@ pub(crate) struct ConnectedClient { pub(crate) nym_address: Recipient, pub(crate) mix_hops: Option, pub(crate) last_activity: std::time::Instant, + // Send to connected clients listener to stop + // This is inside an Option only because we want to send in Drop + pub(crate) close_tx: Option>, +} + +impl Drop for ConnectedClient { + fn drop(&mut self) { + log::info!("Dropping connected client: {}", self.nym_address); + if let Some(close_tx) = self.close_tx.take() { + log::info!("Sending close signal to connected client handler"); + close_tx.send(()).unwrap(); + } + } } impl ConnectedClient { @@ -230,8 +219,28 @@ impl MixnetListener { } (false, false) => { log::info!("Connecting a new client"); - self.connected_clients - .connect(requested_ip, reply_to, reply_to_hops); + + // Start the ConnectedClientHandler for the new client + let (close_tx, close_rx) = tokio::sync::oneshot::channel(); + let (forward_from_tun_tx, forward_from_tun_rx) = + tokio::sync::mpsc::unbounded_channel(); + let connected_client_handler = + connected_client_handler::ConnectedClientHandler::new( + reply_to, + reply_to_hops, + forward_from_tun_rx, + self.mixnet_client.split_sender(), + close_rx, + ); + connected_client_handler.start(); + + self.connected_clients.connect( + requested_ip, + reply_to, + reply_to_hops, + forward_from_tun_tx, + close_tx, + ); Ok(Some(IpPacketResponse::new_static_connect_success( request_id, reply_to, ))) @@ -298,8 +307,25 @@ impl MixnetListener { ))); }; - self.connected_clients - .connect(new_ip, reply_to, reply_to_hops); + // Start the ConnectedClientHandler for the new client + let (close_tx, close_rx) = tokio::sync::oneshot::channel(); + let (forward_from_tun_tx, forward_from_tun_rx) = tokio::sync::mpsc::unbounded_channel(); + let connected_client_handler = connected_client_handler::ConnectedClientHandler::new( + reply_to, + reply_to_hops, + forward_from_tun_rx, + self.mixnet_client.split_sender(), + close_rx, + ); + connected_client_handler.start(); + + self.connected_clients.connect( + new_ip, + reply_to, + reply_to_hops, + forward_from_tun_tx, + close_tx, + ); Ok(Some(IpPacketResponse::new_dynamic_connect_success( request_id, reply_to, new_ip, ))) @@ -504,4 +530,5 @@ pub(crate) struct ConnectEvent { pub(crate) ip: IpAddr, pub(crate) nym_address: Recipient, pub(crate) mix_hops: Option, + pub(crate) forward_from_tun_tx: tokio::sync::mpsc::UnboundedSender>, } diff --git a/service-providers/ip-packet-router/src/tun_listener.rs b/service-providers/ip-packet-router/src/tun_listener.rs index c50eb09247..d235b20df4 100644 --- a/service-providers/ip-packet-router/src/tun_listener.rs +++ b/service-providers/ip-packet-router/src/tun_listener.rs @@ -1,5 +1,7 @@ +use std::{collections::HashMap, net::IpAddr}; + use nym_ip_packet_requests::response::IpPacketResponse; -use nym_sdk::mixnet::MixnetMessageSender; +use nym_sdk::mixnet::{MixnetMessageSender, Recipient}; use nym_task::TaskClient; #[cfg(target_os = "linux")] use tokio::io::AsyncReadExt; @@ -10,13 +12,73 @@ use crate::{ util::{create_message::create_input_message, parse_ip::parse_dst_addr}, }; +pub(crate) struct ConnectedClientMirror { + pub(crate) nym_address: Recipient, + pub(crate) mix_hops: Option, + pub(crate) last_activity: std::time::Instant, + // Forward packets we read from the TUN device to the connected clients listener + pub(crate) forward_from_tun_tx: tokio::sync::mpsc::UnboundedSender>, +} + +pub(crate) struct ConnectedClientsListener { + clients: HashMap, + connected_client_rx: + tokio::sync::mpsc::UnboundedReceiver, +} + +impl ConnectedClientsListener { + pub(crate) fn new( + connected_client_rx: tokio::sync::mpsc::UnboundedReceiver< + mixnet_listener::ConnectedClientEvent, + >, + ) -> Self { + ConnectedClientsListener { + clients: HashMap::new(), + connected_client_rx, + } + } + + pub(crate) fn get(&self, ip: &IpAddr) -> Option<&ConnectedClientMirror> { + self.clients.get(ip) + } + + pub(crate) fn update(&mut self, event: mixnet_listener::ConnectedClientEvent) { + match event { + mixnet_listener::ConnectedClientEvent::Connect(connected_event) => { + let mixnet_listener::ConnectEvent { + ip, + nym_address, + mix_hops, + forward_from_tun_tx, + } = *connected_event; + log::trace!("Connect client: {ip}"); + self.clients.insert( + ip, + ConnectedClientMirror { + nym_address, + mix_hops, + last_activity: std::time::Instant::now(), + forward_from_tun_tx, + }, + ); + } + mixnet_listener::ConnectedClientEvent::Disconnect( + mixnet_listener::DisconnectEvent(ip), + ) => { + log::trace!("Disconnect client: {ip}"); + self.clients.remove(&ip); + } + } + } +} + // Reads packet from TUN and writes to mixnet client #[cfg(target_os = "linux")] pub(crate) struct TunListener { pub(crate) tun_reader: tokio::io::ReadHalf, pub(crate) mixnet_client_sender: nym_sdk::mixnet::MixnetClientSender, pub(crate) task_client: TaskClient, - pub(crate) connected_clients: mixnet_listener::ConnectedClientsListener, + pub(crate) connected_clients: ConnectedClientsListener, } #[cfg(target_os = "linux")] @@ -27,24 +89,27 @@ impl TunListener { return Ok(()); }; - if let Some(mixnet_listener::ConnectedClient { + if let Some(ConnectedClientMirror { nym_address, mix_hops, - .. + last_activity, + forward_from_tun_tx, }) = self.connected_clients.get(&dst_addr) { let packet = buf[..len].to_vec(); - let response_packet = IpPacketResponse::new_ip_packet(packet.into()) - .to_bytes() - .map_err(|err| IpPacketRouterError::FailedToSerializeResponsePacket { - source: err, - })?; - let input_message = create_input_message(*nym_address, response_packet, *mix_hops); + forward_from_tun_tx.send(packet).unwrap(); - self.mixnet_client_sender - .send(input_message) - .await - .map_err(|err| IpPacketRouterError::FailedToSendPacketToMixnet { source: err })?; + // let response_packet = IpPacketResponse::new_ip_packet(packet.into()) + // .to_bytes() + // .map_err(|err| IpPacketRouterError::FailedToSerializeResponsePacket { + // source: err, + // })?; + // let input_message = create_input_message(*nym_address, response_packet, *mix_hops); + // + // self.mixnet_client_sender + // .send(input_message) + // .await + // .map_err(|err| IpPacketRouterError::FailedToSendPacketToMixnet { source: err })?; } else { log::info!("No registered nym-address for packet - dropping"); }