use crate::ip_packet_client::{
discovery::{create_nym_api_client, get_best_ipr},
handle_ipr_response,
listener::check_ipr_message_version,
MixnetMessageOutcome,
};
use crate::mixnet::{MixnetClient, MixnetStream, Recipient};
use crate::Error;
use bytes::Bytes;
use current_ipr::response::IpPacketResponse;
use nym_ip_packet_requests::response_helpers;
use nym_ip_packet_requests::{v9 as current_ipr, IpPair};
use nym_network_defaults::NymNetworkDetails;
use std::time::Duration;
use tokio::io::AsyncWriteExt;
use tracing::{debug, info};
const IPR_CONNECT_TIMEOUT: Duration = Duration::from_secs(60);
pub struct IpMixStream {
stream: MixnetStream,
client: MixnetClient,
allocated_ips: IpPair,
connected: bool,
}
impl IpMixStream {
pub async fn new() -> Result<Self, Error> {
let network_defaults = NymNetworkDetails::new_mainnet();
let api_client =
create_nym_api_client(network_defaults.nym_api_urls.ok_or(Error::NoNymAPIUrl)?)?;
let ipr_address = get_best_ipr(api_client).await?;
Self::new_with_ipr(ipr_address).await
}
pub async fn new_with_ipr(ipr_address: Recipient) -> Result<Self, Error> {
nym_network_defaults::setup_env(None::<&str>);
let mut client = MixnetClient::connect_new().await?;
let mut stream = client.open_stream(ipr_address, Some(10)).await?;
info!("Connecting to IP packet router at {ipr_address}");
let allocated_ips = Self::connect_tunnel(&mut stream).await?;
info!(
"Connected — IPv4: {}, IPv6: {}",
allocated_ips.ipv4, allocated_ips.ipv6
);
Ok(Self {
stream,
client,
allocated_ips,
connected: true,
})
}
pub fn nym_address(&self) -> &Recipient {
self.client.nym_address()
}
pub fn allocated_ips(&self) -> &IpPair {
&self.allocated_ips
}
pub fn is_connected(&self) -> bool {
self.connected
}
pub fn check_connected(&self) -> Result<(), Error> {
if self.connected {
Ok(())
} else {
Err(Error::IprStreamClientNotConnected)
}
}
async fn connect_tunnel(stream: &mut MixnetStream) -> Result<IpPair, Error> {
let (request, request_id) = current_ipr::new_connect_request(None);
debug!("Sending connect request with ID: {}", request_id);
let request_bytes = request.to_bytes()?;
stream
.write_all(&request_bytes)
.await
.map_err(|_| Error::MessageSendingFailure)?;
let timeout = tokio::time::sleep(IPR_CONNECT_TIMEOUT);
tokio::pin!(timeout);
loop {
tokio::select! {
_ = &mut timeout => {
return Err(Error::IPRConnectResponseTimeout);
}
result = stream.recv() => {
let data = result.ok_or(Error::IPRClientStreamClosed)?;
check_ipr_message_version(&data)?;
if let Ok(response) = IpPacketResponse::from_bytes(&data) {
if response.id() == Some(request_id) {
return response_helpers::parse_connect_response(response)
.map_err(|e| match e {
response_helpers::IprResponseError::ConnectDenied(r) => Error::ConnectDenied(r),
response_helpers::IprResponseError::UnexpectedResponse(d) => Error::UnexpectedResponseType(d),
other => Error::IPRMessageVersionCheckFailed(other.to_string()),
});
}
}
}
}
}
}
pub async fn send_ip_packet(&mut self, packet: &[u8]) -> Result<(), Error> {
self.check_connected()?;
let request = current_ipr::new_data_request(packet.to_vec().into());
let request_bytes = request.to_bytes()?;
self.stream
.write_all(&request_bytes)
.await
.map_err(|_| Error::MessageSendingFailure)
}
pub async fn handle_incoming(&mut self) -> Result<Vec<Bytes>, Error> {
let data = match tokio::time::timeout(Duration::from_secs(10), self.stream.recv()).await {
Err(_) => return Ok(Vec::new()),
Ok(None) => {
self.connected = false;
return Err(Error::IPRClientStreamClosed);
}
Ok(Some(data)) => data,
};
match handle_ipr_response(&data) {
Ok(Some(MixnetMessageOutcome::IpPackets(packets))) => {
debug!("Extracted {} IP packets", packets.len());
Ok(packets)
}
Ok(Some(MixnetMessageOutcome::Disconnect)) => {
info!("Received disconnect");
self.connected = false;
Err(Error::IprTunnelDisconnected)
}
Ok(None) => Ok(Vec::new()),
Err(e) => Err(e),
}
}
pub async fn disconnect(self) {
debug!("Disconnecting");
self.client.disconnect().await;
debug!("Disconnected");
}
}