#[cfg(any(target_os = "linux", target_os = "android"))]
use std::os::fd::RawFd;
use std::{
net::{Ipv4Addr, SocketAddr},
time::{Duration, Instant},
};
use tokio::io::{AsyncRead, AsyncWrite};
use tokio_util::sync::CancellationToken;
use tracing::*;
use crate::transport::{quic, ssh, tls};
use crate::{config::ClientConfig, error::TransportError};
use nym_bridges_types::TransportAssociation;
const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
pub(crate) fn make_socket(addr: Option<SocketAddr>) -> std::io::Result<std::net::UdpSocket> {
let addr = addr.unwrap_or((Ipv4Addr::UNSPECIFIED, 0).into());
let socket = std::net::UdpSocket::bind(addr)?;
socket.set_nonblocking(true)?;
Ok(socket)
}
#[cfg(any(target_os = "linux", target_os = "android"))]
#[allow(non_snake_case)]
pub fn SOCKET_OPEN_NOP(_: RawFd) {}
pub struct BridgeConn {
#[allow(unused)] pub(crate) params: ClientConfig,
pub(crate) endpoint: SocketAddr,
pub(crate) reader: Box<dyn AsyncRead + Send + Unpin>,
pub(crate) writer: Box<dyn AsyncWrite + Send + Unpin>,
pub(crate) closer: Box<dyn TransportCloser>,
}
pub trait TransportCloser: Send {
fn close(self: Box<Self>) -> futures::future::BoxFuture<'static, ()>;
}
impl TransportCloser for () {
fn close(self: Box<Self>) -> futures::future::BoxFuture<'static, ()> {
Box::pin(async {})
}
}
impl BridgeConn {
pub async fn try_connect(
params: ClientConfig,
token: CancellationToken,
#[cfg(any(target_os = "linux", target_os = "android"))] on_socket_open: impl Fn(RawFd),
conn_timeout: Option<Duration>,
) -> Result<Self, TransportError> {
let start = Instant::now();
let connect_timeout = conn_timeout.unwrap_or(DEFAULT_CONNECT_TIMEOUT);
match params {
ClientConfig::QuicPlain(ref opts) => {
let conn = token
.run_until_cancelled(quic::transport_conn(
opts,
#[cfg(any(target_os = "linux", target_os = "android"))]
on_socket_open,
connect_timeout,
))
.await
.ok_or(TransportError::Cancelled)??;
let endpoint = conn.remote_address();
let (writer, reader) = token
.run_until_cancelled(conn.open_bi())
.await
.ok_or(TransportError::Cancelled)??;
info!("quic transport connected in {:?}", start.elapsed());
Ok(Self {
reader: Box::new(reader),
writer: Box::new(writer),
params,
endpoint,
closer: Box::new(conn),
})
}
ClientConfig::TlsPlain(ref opts) => {
let conn = token
.run_until_cancelled(tls::transport_conn(
opts,
#[cfg(any(target_os = "linux", target_os = "android"))]
on_socket_open,
connect_timeout,
))
.await
.ok_or(TransportError::Cancelled)??;
let endpoint = conn.get_ref().0.peer_addr()?;
info!(
"{} transport connected in {:?}",
opts.transport_name(),
start.elapsed()
);
let (reader, writer) = tokio::io::split(conn);
Ok(Self {
reader: Box::new(reader),
writer: Box::new(writer),
params,
endpoint,
closer: Box::new(()),
})
}
ClientConfig::SshPlain(ref opts) => {
let endpoint = *opts.addresses.first().ok_or_else(|| {
TransportError::config_err("no ssh bridge address configured")
})?;
let stream = token
.run_until_cancelled(ssh::transport_conn(opts, connect_timeout))
.await
.ok_or(TransportError::Cancelled)??;
let (reader, writer) = tokio::io::split(stream);
info!("ssh transport connected in {:?}", start.elapsed());
Ok(Self {
reader: Box::new(reader),
writer: Box::new(writer),
params,
endpoint,
closer: Box::new(()),
})
}
}
}
pub fn endpoint(&self) -> SocketAddr {
self.endpoint
}
pub fn params(&self) -> &ClientConfig {
&self.params
}
pub fn into_parts(
self,
) -> (
Box<dyn AsyncRead + Send + Unpin>,
Box<dyn AsyncWrite + Send + Unpin>,
Box<dyn TransportCloser>,
) {
(self.reader, self.writer, self.closer)
}
}