#[cfg(not(target_os = "linux"))]
fn main() -> Result<(), BoxError> {
Err(BoxError::from_static_str(
"the linux_tproxy_tcp example only supports Linux",
))
}
#[cfg(target_os = "linux")]
use ::{
rama::{
Layer, Service,
net::{
address::SocketAddress,
proxy::IoForwardService,
socket::{
SocketOptions,
linux::ConnectorTargetFromGetSocketnameLayer,
opts::{Domain, TcpKeepAlive},
},
stream::Socket,
},
rt::Executor,
service::service_fn,
tcp::{TcpStream, proxy::IoToProxyBridgeIoLayer, server::TcpListener},
telemetry::tracing::{
self,
level_filters::LevelFilter,
subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt},
},
},
std::time::Duration,
};
#[cfg(not(target_os = "linux"))]
use rama::error::{BoxError, BoxErrorExt};
#[cfg(target_os = "linux")]
const LISTEN_ADDR_V4: SocketAddress = SocketAddress::default_ipv4(62052);
#[cfg(target_os = "linux")]
const LISTEN_ADDR_V6: SocketAddress = SocketAddress::default_ipv6(62052);
#[cfg(target_os = "linux")]
fn tcp_keep_alive() -> TcpKeepAlive {
TcpKeepAlive {
time: Some(Duration::from_mins(2)),
interval: Some(Duration::from_secs(30)),
retries: Some(5),
}
}
#[cfg(target_os = "linux")]
#[tokio::main]
async fn main() {
tracing::subscriber::registry()
.with(fmt::layer())
.with(
EnvFilter::builder()
.with_default_directive(LevelFilter::INFO.into())
.from_env_lossy(),
)
.init();
if let Err(err) = run().await {
tracing::error!(error = %err, "linux tproxy tcp example failed");
std::process::exit(1);
}
}
#[cfg(target_os = "linux")]
async fn run() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let exec = Executor::default();
let base = SocketOptions {
reuse_address: Some(true),
reuse_port: Some(true),
tcp_no_delay: Some(true),
tcp_keep_alive: Some(tcp_keep_alive()),
..SocketOptions::default_tcp()
};
let socket_v4 = SocketOptions {
address: Some(LISTEN_ADDR_V4),
ip_transparent: Some(true),
freebind: Some(true),
..base.clone()
}
.try_build_socket(Domain::IPv4)?;
socket_v4.listen(32_768)?;
let listener_v4 = TcpListener::bind_socket(socket_v4, exec.clone()).await?;
let socket_v6 = SocketOptions {
address: Some(LISTEN_ADDR_V6),
ip_transparent_v6: Some(true),
freebind_ipv6: Some(true),
..base
}
.try_build_socket(Domain::IPv6)?;
socket_v6.listen(32_768)?;
let listener_v6 = TcpListener::bind_socket(socket_v6, exec.clone()).await?;
tracing::info!(
listen.v4 = %LISTEN_ADDR_V4,
listen.v6 = %LISTEN_ADDR_V6,
"dual-stack transparent tcp proxy listening"
);
tracing::info!(
"make sure Linux policy routing and TPROXY rules are installed first (for both IPv4 and IPv6)"
);
let service = ConnectorTargetFromGetSocketnameLayer::new().into_layer(service_fn({
let forward = IoToProxyBridgeIoLayer::extension_connector_target()
.into_layer(IoForwardService::new(exec));
move |stream: TcpStream| {
let forward = forward.clone();
async move {
let original_dst = stream.local_addr()?;
let peer_addr = stream.peer_addr()?;
tracing::info!(
network.peer.address = %peer_addr.ip_addr,
network.peer.port = peer_addr.port,
network.original.address = %original_dst.ip_addr,
network.original.port = original_dst.port,
"accepted intercepted tcp flow"
);
forward.serve(stream).await
}
}
}));
tokio::join!(
listener_v4.serve(service.clone()),
listener_v6.serve(service),
);
Ok(())
}