使用异步tun
This commit is contained in:
@@ -1,187 +1,11 @@
|
|||||||
use std::io;
|
|
||||||
use std::net::Ipv4Addr;
|
|
||||||
|
|
||||||
use parking_lot::Mutex;
|
|
||||||
|
|
||||||
use packet::ip::ipv4::packet::IpV4Packet;
|
|
||||||
use packet::ip::ipv4::protocol::Protocol;
|
|
||||||
|
|
||||||
use crate::channel::context::ChannelContext;
|
|
||||||
use crate::cipher::Cipher;
|
|
||||||
use crate::external_route::ExternalRoute;
|
|
||||||
use crate::handle::{check_dest, CurrentDeviceInfo, PeerDeviceInfo};
|
|
||||||
#[cfg(feature = "ip_proxy")]
|
|
||||||
use crate::ip_proxy::{IpProxyMap, ProxyHandler};
|
|
||||||
use crate::protocol;
|
|
||||||
use crate::protocol::body::ENCRYPTION_RESERVED;
|
|
||||||
use crate::protocol::ip_turn_packet::BroadcastPacket;
|
|
||||||
use crate::protocol::{ip_turn_packet, NetPacket, MAX_TTL};
|
|
||||||
|
|
||||||
mod channel_group;
|
mod channel_group;
|
||||||
pub mod tun_handler;
|
pub mod tun_handler;
|
||||||
|
|
||||||
fn broadcast(
|
#[cfg(unix)]
|
||||||
server_cipher: &Cipher,
|
mod unix;
|
||||||
sender: &ChannelContext,
|
#[cfg(unix)]
|
||||||
net_packet: &mut NetPacket<&mut [u8]>,
|
pub(crate) use unix::*;
|
||||||
current_device: &CurrentDeviceInfo,
|
#[cfg(target_os = "windows")]
|
||||||
device_list: &Mutex<(u16, Vec<PeerDeviceInfo>)>,
|
mod windows;
|
||||||
) -> io::Result<()> {
|
#[cfg(target_os = "windows")]
|
||||||
let list: Vec<Ipv4Addr> = device_list
|
pub(crate) use windows::*;
|
||||||
.lock()
|
|
||||||
.1
|
|
||||||
.iter()
|
|
||||||
.filter(|info| info.status.is_online())
|
|
||||||
.map(|info| info.virtual_ip)
|
|
||||||
.collect();
|
|
||||||
const MAX_COUNT: usize = 8;
|
|
||||||
let mut p2p_ips = Vec::with_capacity(8);
|
|
||||||
let mut relay_ips = Vec::with_capacity(8);
|
|
||||||
let mut overflow = false;
|
|
||||||
for (index, peer_ip) in list.into_iter().enumerate() {
|
|
||||||
if index > MAX_COUNT {
|
|
||||||
overflow = true;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
if let Some(route) = sender.route_table.route_one_p2p(&peer_ip) {
|
|
||||||
if sender
|
|
||||||
.send_by_key(net_packet.buffer(), route.route_key())
|
|
||||||
.is_ok()
|
|
||||||
{
|
|
||||||
p2p_ips.push(peer_ip);
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
relay_ips.push(peer_ip);
|
|
||||||
}
|
|
||||||
if !overflow && relay_ips.is_empty() {
|
|
||||||
//全部p2p,不需要服务器中转
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
|
|
||||||
if p2p_ips.is_empty() {
|
|
||||||
//都没有p2p则直接由服务器转发
|
|
||||||
if current_device.status.online() {
|
|
||||||
sender.send_default(net_packet.buffer(), current_device.connect_server)?;
|
|
||||||
}
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
if !overflow && relay_ips.len() == 2 {
|
|
||||||
// 如果转发的ip数不多就直接发
|
|
||||||
for peer_ip in relay_ips {
|
|
||||||
//非直连的广播要改变目的地址,不然服务端收到了会再次广播
|
|
||||||
net_packet.set_destination(peer_ip);
|
|
||||||
sender.send_ipv4_by_id(
|
|
||||||
net_packet.buffer(),
|
|
||||||
&peer_ip,
|
|
||||||
current_device.connect_server,
|
|
||||||
current_device.status.online(),
|
|
||||||
)?;
|
|
||||||
}
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
if current_device.status.offline() {
|
|
||||||
//离线的不再转发
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
let buf = vec![0u8; 12 + 1 + p2p_ips.len() * 4 + net_packet.data_len() + ENCRYPTION_RESERVED];
|
|
||||||
//剩余的发送到服务端,需要告知哪些已发送过
|
|
||||||
let mut server_packet = NetPacket::new_encrypt(buf)?;
|
|
||||||
server_packet.set_default_version();
|
|
||||||
server_packet.set_gateway_flag(true);
|
|
||||||
server_packet.first_set_ttl(MAX_TTL);
|
|
||||||
server_packet.set_source(net_packet.source());
|
|
||||||
//使用对应的目的地址
|
|
||||||
server_packet.set_destination(net_packet.destination());
|
|
||||||
server_packet.set_protocol(protocol::Protocol::IpTurn);
|
|
||||||
server_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4Broadcast.into());
|
|
||||||
|
|
||||||
let mut broadcast = BroadcastPacket::unchecked(server_packet.payload_mut());
|
|
||||||
broadcast.set_address(&p2p_ips)?;
|
|
||||||
broadcast.set_data(net_packet.buffer())?;
|
|
||||||
server_cipher.encrypt_ipv4(&mut server_packet)?;
|
|
||||||
sender.send_default(server_packet.buffer(), current_device.connect_server)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// 实现一个原地发送,必须保证是如下结构
|
|
||||||
/// |12字节开头|ip报文|至少1024字节结尾|
|
|
||||||
///
|
|
||||||
#[inline]
|
|
||||||
pub fn base_handle(
|
|
||||||
context: &ChannelContext,
|
|
||||||
buf: &mut [u8],
|
|
||||||
data_len: usize, //数据总长度=12+ip包长度
|
|
||||||
current_device: CurrentDeviceInfo,
|
|
||||||
ip_route: &ExternalRoute,
|
|
||||||
#[cfg(feature = "ip_proxy")] proxy_map: &Option<IpProxyMap>,
|
|
||||||
client_cipher: &Cipher,
|
|
||||||
server_cipher: &Cipher,
|
|
||||||
device_list: &Mutex<(u16, Vec<PeerDeviceInfo>)>,
|
|
||||||
) -> io::Result<()> {
|
|
||||||
let ipv4_packet = IpV4Packet::new(&buf[12..data_len])?;
|
|
||||||
let protocol = ipv4_packet.protocol();
|
|
||||||
let src_ip = ipv4_packet.source_ip();
|
|
||||||
let mut dest_ip = ipv4_packet.destination_ip();
|
|
||||||
let mut net_packet = NetPacket::new0(data_len, buf)?;
|
|
||||||
net_packet.set_default_version();
|
|
||||||
net_packet.set_protocol(protocol::Protocol::IpTurn);
|
|
||||||
net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into());
|
|
||||||
net_packet.first_set_ttl(6);
|
|
||||||
net_packet.set_source(src_ip);
|
|
||||||
net_packet.set_destination(dest_ip);
|
|
||||||
if dest_ip == current_device.virtual_gateway {
|
|
||||||
// 发到网关的加密方式不一样,要单独处理
|
|
||||||
if protocol == Protocol::Icmp {
|
|
||||||
net_packet.set_gateway_flag(true);
|
|
||||||
server_cipher.encrypt_ipv4(&mut net_packet)?;
|
|
||||||
context.send_default(net_packet.buffer(), current_device.connect_server)?;
|
|
||||||
}
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
if dest_ip.is_multicast() {
|
|
||||||
//当作广播处理
|
|
||||||
dest_ip = Ipv4Addr::BROADCAST;
|
|
||||||
net_packet.set_destination(Ipv4Addr::BROADCAST);
|
|
||||||
}
|
|
||||||
if dest_ip.is_broadcast() || current_device.broadcast_ip == dest_ip {
|
|
||||||
// 广播 发送到直连目标
|
|
||||||
client_cipher.encrypt_ipv4(&mut net_packet)?;
|
|
||||||
broadcast(
|
|
||||||
server_cipher,
|
|
||||||
context,
|
|
||||||
&mut net_packet,
|
|
||||||
¤t_device,
|
|
||||||
device_list,
|
|
||||||
)?;
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
if !check_dest(
|
|
||||||
dest_ip,
|
|
||||||
current_device.virtual_netmask,
|
|
||||||
current_device.virtual_network,
|
|
||||||
) {
|
|
||||||
if let Some(r_dest_ip) = ip_route.route(&dest_ip) {
|
|
||||||
//路由的目标不能是自己
|
|
||||||
if r_dest_ip == src_ip {
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
//需要修改目的地址
|
|
||||||
dest_ip = r_dest_ip;
|
|
||||||
net_packet.set_destination(r_dest_ip);
|
|
||||||
} else {
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
#[cfg(feature = "ip_proxy")]
|
|
||||||
if let Some(proxy_map) = proxy_map {
|
|
||||||
let mut ipv4_packet = IpV4Packet::new(net_packet.payload_mut())?;
|
|
||||||
proxy_map.send_handle(&mut ipv4_packet)?;
|
|
||||||
}
|
|
||||||
client_cipher.encrypt_ipv4(&mut net_packet)?;
|
|
||||||
context.send_ipv4_by_id(
|
|
||||||
net_packet.buffer(),
|
|
||||||
&dest_ip,
|
|
||||||
current_device.connect_server,
|
|
||||||
current_device.status.online(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
|
use std::net::Ipv4Addr;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::{io, thread};
|
use std::{io, thread};
|
||||||
|
|
||||||
@@ -6,22 +7,27 @@ use parking_lot::Mutex;
|
|||||||
|
|
||||||
use packet::icmp::icmp::IcmpPacket;
|
use packet::icmp::icmp::IcmpPacket;
|
||||||
use packet::icmp::Kind;
|
use packet::icmp::Kind;
|
||||||
use packet::ip::ipv4;
|
|
||||||
use packet::ip::ipv4::packet::IpV4Packet;
|
use packet::ip::ipv4::packet::IpV4Packet;
|
||||||
|
use packet::ip::ipv4::protocol::Protocol;
|
||||||
use tun::device::IFace;
|
use tun::device::IFace;
|
||||||
use tun::Device;
|
use tun::Device;
|
||||||
|
|
||||||
use crate::channel::context::ChannelContext;
|
use crate::channel::context::ChannelContext;
|
||||||
use crate::cipher::Cipher;
|
use crate::cipher::Cipher;
|
||||||
use crate::external_route::ExternalRoute;
|
use crate::external_route::ExternalRoute;
|
||||||
use crate::handle::tun_tap::channel_group::{channel_group, GroupSyncSender};
|
use crate::handle::tun_tap::channel_group::channel_group;
|
||||||
use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo};
|
use crate::handle::{check_dest, CurrentDeviceInfo, PeerDeviceInfo};
|
||||||
#[cfg(feature = "ip_proxy")]
|
#[cfg(feature = "ip_proxy")]
|
||||||
use crate::ip_proxy::IpProxyMap;
|
use crate::ip_proxy::IpProxyMap;
|
||||||
|
use crate::ip_proxy::ProxyHandler;
|
||||||
|
use crate::protocol;
|
||||||
|
use crate::protocol::body::ENCRYPTION_RESERVED;
|
||||||
|
use crate::protocol::ip_turn_packet::BroadcastPacket;
|
||||||
|
use crate::protocol::{ip_turn_packet, NetPacket, MAX_TTL};
|
||||||
use crate::util::{SingleU64Adder, StopManager};
|
use crate::util::{SingleU64Adder, StopManager};
|
||||||
|
|
||||||
fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> io::Result<()> {
|
fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> io::Result<()> {
|
||||||
if ipv4_packet.protocol() == ipv4::protocol::Protocol::Icmp {
|
if ipv4_packet.protocol() == Protocol::Icmp {
|
||||||
let mut icmp = IcmpPacket::new(ipv4_packet.payload_mut())?;
|
let mut icmp = IcmpPacket::new(ipv4_packet.payload_mut())?;
|
||||||
if icmp.kind() == Kind::EchoRequest {
|
if icmp.kind() == Kind::EchoRequest {
|
||||||
icmp.set_kind(Kind::EchoReply);
|
icmp.set_kind(Kind::EchoReply);
|
||||||
@@ -37,7 +43,7 @@ fn icmp(device_writer: &Device, mut ipv4_packet: IpV4Packet<&mut [u8]>) -> io::R
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// 接收tun数据,并且转发到udp上
|
/// 接收tun数据,并且转发到udp上
|
||||||
fn handle(
|
pub(crate) fn handle(
|
||||||
context: &ChannelContext,
|
context: &ChannelContext,
|
||||||
data: &mut [u8],
|
data: &mut [u8],
|
||||||
len: usize,
|
len: usize,
|
||||||
@@ -59,7 +65,7 @@ fn handle(
|
|||||||
if src_ip == dest_ip {
|
if src_ip == dest_ip {
|
||||||
return icmp(&device_writer, ipv4_packet);
|
return icmp(&device_writer, ipv4_packet);
|
||||||
}
|
}
|
||||||
return crate::handle::tun_tap::base_handle(
|
return base_handle(
|
||||||
context,
|
context,
|
||||||
data,
|
data,
|
||||||
len,
|
len,
|
||||||
@@ -86,24 +92,6 @@ pub fn start(
|
|||||||
mut up_counter: SingleU64Adder,
|
mut up_counter: SingleU64Adder,
|
||||||
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
||||||
) -> io::Result<()> {
|
) -> io::Result<()> {
|
||||||
#[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))]
|
|
||||||
let worker = {
|
|
||||||
#[cfg(target_os = "macos")]
|
|
||||||
let current_device = current_device.clone();
|
|
||||||
let device = device.clone();
|
|
||||||
stop_manager.add_listener("tun_device".into(), move || {
|
|
||||||
if let Err(e) = device.shutdown() {
|
|
||||||
log::warn!("{:?}", e);
|
|
||||||
}
|
|
||||||
#[cfg(target_os = "macos")]
|
|
||||||
{
|
|
||||||
let ip = current_device.load().virtual_ip;
|
|
||||||
if let Ok(udp) = std::net::UdpSocket::bind("0.0.0.0:0") {
|
|
||||||
let _ = udp.send_to(b"stop", format!("{:?}:1234", ip));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})?
|
|
||||||
};
|
|
||||||
if parallel > 1 {
|
if parallel > 1 {
|
||||||
let (sender, receivers) = channel_group::<(Vec<u8>, usize)>(parallel, 16);
|
let (sender, receivers) = channel_group::<(Vec<u8>, usize)>(parallel, 16);
|
||||||
for (index, receiver) in receivers.into_iter().enumerate() {
|
for (index, receiver) in receivers.into_iter().enumerate() {
|
||||||
@@ -148,17 +136,20 @@ pub fn start(
|
|||||||
thread::Builder::new()
|
thread::Builder::new()
|
||||||
.name("tunHandlerM".into())
|
.name("tunHandlerM".into())
|
||||||
.spawn(move || {
|
.spawn(move || {
|
||||||
if let Err(e) = start_multi(stop_manager, device, sender, &mut up_counter) {
|
if let Err(e) = crate::handle::tun_tap::start_multi(
|
||||||
|
stop_manager,
|
||||||
|
device,
|
||||||
|
sender,
|
||||||
|
&mut up_counter,
|
||||||
|
) {
|
||||||
log::warn!("stop:{}", e);
|
log::warn!("stop:{}", e);
|
||||||
}
|
}
|
||||||
#[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))]
|
|
||||||
worker.stop_all();
|
|
||||||
})?;
|
})?;
|
||||||
} else {
|
} else {
|
||||||
thread::Builder::new()
|
thread::Builder::new()
|
||||||
.name("tunHandlerS".into())
|
.name("tunHandlerS".into())
|
||||||
.spawn(move || {
|
.spawn(move || {
|
||||||
if let Err(e) = start_simple(
|
if let Err(e) = crate::handle::tun_tap::start_simple(
|
||||||
stop_manager,
|
stop_manager,
|
||||||
&context,
|
&context,
|
||||||
device,
|
device,
|
||||||
@@ -173,74 +164,173 @@ pub fn start(
|
|||||||
) {
|
) {
|
||||||
log::warn!("stop:{}", e);
|
log::warn!("stop:{}", e);
|
||||||
}
|
}
|
||||||
#[cfg(any(target_os = "windows", target_os = "linux", target_os = "macos"))]
|
|
||||||
worker.stop_all();
|
|
||||||
})?;
|
})?;
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn start_simple(
|
fn broadcast(
|
||||||
stop_manager: StopManager,
|
server_cipher: &Cipher,
|
||||||
context: &ChannelContext,
|
sender: &ChannelContext,
|
||||||
device: Arc<Device>,
|
net_packet: &mut NetPacket<&mut [u8]>,
|
||||||
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
|
current_device: &CurrentDeviceInfo,
|
||||||
ip_route: ExternalRoute,
|
device_list: &Mutex<(u16, Vec<PeerDeviceInfo>)>,
|
||||||
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
|
|
||||||
client_cipher: Cipher,
|
|
||||||
server_cipher: Cipher,
|
|
||||||
up_counter: &mut SingleU64Adder,
|
|
||||||
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
|
||||||
) -> io::Result<()> {
|
) -> io::Result<()> {
|
||||||
let mut buf = [0; 1024 * 16];
|
let list: Vec<Ipv4Addr> = device_list
|
||||||
loop {
|
.lock()
|
||||||
if stop_manager.is_stop() {
|
.1
|
||||||
return Ok(());
|
.iter()
|
||||||
|
.filter(|info| info.status.is_online())
|
||||||
|
.map(|info| info.virtual_ip)
|
||||||
|
.collect();
|
||||||
|
const MAX_COUNT: usize = 8;
|
||||||
|
let mut p2p_ips = Vec::with_capacity(8);
|
||||||
|
let mut relay_ips = Vec::with_capacity(8);
|
||||||
|
let mut overflow = false;
|
||||||
|
for (index, peer_ip) in list.into_iter().enumerate() {
|
||||||
|
if index > MAX_COUNT {
|
||||||
|
overflow = true;
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
let len = device.read(&mut buf[12..])? + 12;
|
if let Some(route) = sender.route_table.route_one_p2p(&peer_ip) {
|
||||||
//单线程的
|
if sender
|
||||||
up_counter.add(len as u64);
|
.send_by_key(net_packet.buffer(), route.route_key())
|
||||||
#[cfg(any(target_os = "macos"))]
|
.is_ok()
|
||||||
let mut buf = &mut buf[4..];
|
{
|
||||||
// buf是重复利用的,需要重置头部
|
p2p_ips.push(peer_ip);
|
||||||
buf[..12].fill(0);
|
continue;
|
||||||
match handle(
|
|
||||||
context,
|
|
||||||
&mut buf,
|
|
||||||
len,
|
|
||||||
&device,
|
|
||||||
current_device.load(),
|
|
||||||
&ip_route,
|
|
||||||
#[cfg(feature = "ip_proxy")]
|
|
||||||
&ip_proxy_map,
|
|
||||||
&client_cipher,
|
|
||||||
&server_cipher,
|
|
||||||
&device_list,
|
|
||||||
) {
|
|
||||||
Ok(_) => {}
|
|
||||||
Err(e) => {
|
|
||||||
log::warn!("{:?}", e)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
relay_ips.push(peer_ip);
|
||||||
}
|
}
|
||||||
|
if !overflow && relay_ips.is_empty() {
|
||||||
|
//全部p2p,不需要服务器中转
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
if p2p_ips.is_empty() {
|
||||||
|
//都没有p2p则直接由服务器转发
|
||||||
|
if current_device.status.online() {
|
||||||
|
sender.send_default(net_packet.buffer(), current_device.connect_server)?;
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
if !overflow && relay_ips.len() == 2 {
|
||||||
|
// 如果转发的ip数不多就直接发
|
||||||
|
for peer_ip in relay_ips {
|
||||||
|
//非直连的广播要改变目的地址,不然服务端收到了会再次广播
|
||||||
|
net_packet.set_destination(peer_ip);
|
||||||
|
sender.send_ipv4_by_id(
|
||||||
|
net_packet.buffer(),
|
||||||
|
&peer_ip,
|
||||||
|
current_device.connect_server,
|
||||||
|
current_device.status.online(),
|
||||||
|
)?;
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
if current_device.status.offline() {
|
||||||
|
//离线的不再转发
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let buf = vec![0u8; 12 + 1 + p2p_ips.len() * 4 + net_packet.data_len() + ENCRYPTION_RESERVED];
|
||||||
|
//剩余的发送到服务端,需要告知哪些已发送过
|
||||||
|
let mut server_packet = NetPacket::new_encrypt(buf)?;
|
||||||
|
server_packet.set_default_version();
|
||||||
|
server_packet.set_gateway_flag(true);
|
||||||
|
server_packet.first_set_ttl(MAX_TTL);
|
||||||
|
server_packet.set_source(net_packet.source());
|
||||||
|
//使用对应的目的地址
|
||||||
|
server_packet.set_destination(net_packet.destination());
|
||||||
|
server_packet.set_protocol(protocol::Protocol::IpTurn);
|
||||||
|
server_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4Broadcast.into());
|
||||||
|
|
||||||
|
let mut broadcast = BroadcastPacket::unchecked(server_packet.payload_mut());
|
||||||
|
broadcast.set_address(&p2p_ips)?;
|
||||||
|
broadcast.set_data(net_packet.buffer())?;
|
||||||
|
server_cipher.encrypt_ipv4(&mut server_packet)?;
|
||||||
|
sender.send_default(server_packet.buffer(), current_device.connect_server)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn start_multi(
|
/// 实现一个原地发送,必须保证是如下结构
|
||||||
stop_manager: StopManager,
|
/// |12字节开头|ip报文|至少1024字节结尾|
|
||||||
device: Arc<Device>,
|
///
|
||||||
mut group_sync_sender: GroupSyncSender<(Vec<u8>, usize)>,
|
#[inline]
|
||||||
up_counter: &mut SingleU64Adder,
|
fn base_handle(
|
||||||
|
context: &ChannelContext,
|
||||||
|
buf: &mut [u8],
|
||||||
|
data_len: usize, //数据总长度=12+ip包长度
|
||||||
|
current_device: CurrentDeviceInfo,
|
||||||
|
ip_route: &ExternalRoute,
|
||||||
|
#[cfg(feature = "ip_proxy")] proxy_map: &Option<IpProxyMap>,
|
||||||
|
client_cipher: &Cipher,
|
||||||
|
server_cipher: &Cipher,
|
||||||
|
device_list: &Mutex<(u16, Vec<PeerDeviceInfo>)>,
|
||||||
) -> io::Result<()> {
|
) -> io::Result<()> {
|
||||||
loop {
|
let ipv4_packet = IpV4Packet::new(&buf[12..data_len])?;
|
||||||
if stop_manager.is_stop() {
|
let protocol = ipv4_packet.protocol();
|
||||||
return Ok(());
|
let src_ip = ipv4_packet.source_ip();
|
||||||
|
let mut dest_ip = ipv4_packet.destination_ip();
|
||||||
|
let mut net_packet = NetPacket::new0(data_len, buf)?;
|
||||||
|
net_packet.set_default_version();
|
||||||
|
net_packet.set_protocol(protocol::Protocol::IpTurn);
|
||||||
|
net_packet.set_transport_protocol(ip_turn_packet::Protocol::Ipv4.into());
|
||||||
|
net_packet.first_set_ttl(6);
|
||||||
|
net_packet.set_source(src_ip);
|
||||||
|
net_packet.set_destination(dest_ip);
|
||||||
|
if dest_ip == current_device.virtual_gateway {
|
||||||
|
// 发到网关的加密方式不一样,要单独处理
|
||||||
|
if protocol == Protocol::Icmp {
|
||||||
|
net_packet.set_gateway_flag(true);
|
||||||
|
server_cipher.encrypt_ipv4(&mut net_packet)?;
|
||||||
|
context.send_default(net_packet.buffer(), current_device.connect_server)?;
|
||||||
}
|
}
|
||||||
let mut buf = vec![0; 1024 * 16];
|
return Ok(());
|
||||||
let len = device.read(&mut buf[12..])? + 12;
|
}
|
||||||
//单线程的
|
if dest_ip.is_multicast() {
|
||||||
up_counter.add(len as u64);
|
//当作广播处理
|
||||||
if group_sync_sender.send((buf, len)).is_err() {
|
dest_ip = Ipv4Addr::BROADCAST;
|
||||||
|
net_packet.set_destination(Ipv4Addr::BROADCAST);
|
||||||
|
}
|
||||||
|
if dest_ip.is_broadcast() || current_device.broadcast_ip == dest_ip {
|
||||||
|
// 广播 发送到直连目标
|
||||||
|
client_cipher.encrypt_ipv4(&mut net_packet)?;
|
||||||
|
broadcast(
|
||||||
|
server_cipher,
|
||||||
|
context,
|
||||||
|
&mut net_packet,
|
||||||
|
¤t_device,
|
||||||
|
device_list,
|
||||||
|
)?;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
if !check_dest(
|
||||||
|
dest_ip,
|
||||||
|
current_device.virtual_netmask,
|
||||||
|
current_device.virtual_network,
|
||||||
|
) {
|
||||||
|
if let Some(r_dest_ip) = ip_route.route(&dest_ip) {
|
||||||
|
//路由的目标不能是自己
|
||||||
|
if r_dest_ip == src_ip {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
//需要修改目的地址
|
||||||
|
dest_ip = r_dest_ip;
|
||||||
|
net_packet.set_destination(r_dest_ip);
|
||||||
|
} else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
#[cfg(feature = "ip_proxy")]
|
||||||
|
if let Some(proxy_map) = proxy_map {
|
||||||
|
let mut ipv4_packet = IpV4Packet::new(net_packet.payload_mut())?;
|
||||||
|
proxy_map.send_handle(&mut ipv4_packet)?;
|
||||||
|
}
|
||||||
|
client_cipher.encrypt_ipv4(&mut net_packet)?;
|
||||||
|
context.send_ipv4_by_id(
|
||||||
|
net_packet.buffer(),
|
||||||
|
&dest_ip,
|
||||||
|
current_device.connect_server,
|
||||||
|
current_device.status.online(),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,183 @@
|
|||||||
|
use crate::channel::context::ChannelContext;
|
||||||
|
use crate::cipher::Cipher;
|
||||||
|
use crate::external_route::ExternalRoute;
|
||||||
|
use crate::handle::tun_tap::channel_group::GroupSyncSender;
|
||||||
|
use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo};
|
||||||
|
use crate::ip_proxy::IpProxyMap;
|
||||||
|
use crate::util::{SingleU64Adder, StopManager};
|
||||||
|
use crossbeam_utils::atomic::AtomicCell;
|
||||||
|
use mio::event::Source;
|
||||||
|
use mio::unix::SourceFd;
|
||||||
|
use mio::{Events, Interest, Poll, Token, Waker};
|
||||||
|
use parking_lot::Mutex;
|
||||||
|
use std::io;
|
||||||
|
use std::os::fd::AsRawFd;
|
||||||
|
use std::sync::Arc;
|
||||||
|
use tun::Device;
|
||||||
|
|
||||||
|
const STOP: Token = Token(0);
|
||||||
|
const FD: Token = Token(1);
|
||||||
|
|
||||||
|
pub(crate) fn start_simple(
|
||||||
|
stop_manager: StopManager,
|
||||||
|
context: &ChannelContext,
|
||||||
|
device: Arc<Device>,
|
||||||
|
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
|
||||||
|
ip_route: ExternalRoute,
|
||||||
|
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
|
||||||
|
client_cipher: Cipher,
|
||||||
|
server_cipher: Cipher,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let poll = Poll::new()?;
|
||||||
|
let waker = Arc::new(Waker::new(poll.registry(), STOP)?);
|
||||||
|
let _waker = waker.clone();
|
||||||
|
let worker = stop_manager.add_listener("tun_device".into(), move || {
|
||||||
|
let _ = waker.wake();
|
||||||
|
})?;
|
||||||
|
if let Err(e) = start_simple0(
|
||||||
|
poll,
|
||||||
|
context,
|
||||||
|
device,
|
||||||
|
current_device,
|
||||||
|
ip_route,
|
||||||
|
#[cfg(feature = "ip_proxy")]
|
||||||
|
ip_proxy_map,
|
||||||
|
client_cipher,
|
||||||
|
server_cipher,
|
||||||
|
up_counter,
|
||||||
|
device_list,
|
||||||
|
) {
|
||||||
|
log::error!("{:?}", e);
|
||||||
|
};
|
||||||
|
worker.stop_all();
|
||||||
|
drop(_waker);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn start_simple0(
|
||||||
|
mut poll: Poll,
|
||||||
|
context: &ChannelContext,
|
||||||
|
device: Arc<Device>,
|
||||||
|
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
|
||||||
|
ip_route: ExternalRoute,
|
||||||
|
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
|
||||||
|
client_cipher: Cipher,
|
||||||
|
server_cipher: Cipher,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let mut buf = [0; 1024 * 16];
|
||||||
|
let fd = device.as_tun_fd();
|
||||||
|
fd.set_nonblock()?;
|
||||||
|
SourceFd(&fd.as_raw_fd()).register(poll.registry(), FD, Interest::READABLE)?;
|
||||||
|
let mut evnets = Events::with_capacity(4);
|
||||||
|
#[cfg(not(target_os = "macos"))]
|
||||||
|
let start = 12;
|
||||||
|
#[cfg(target_os = "macos")]
|
||||||
|
let start = 12 - 4;
|
||||||
|
loop {
|
||||||
|
poll.poll(&mut evnets, None)?;
|
||||||
|
for event in evnets.iter() {
|
||||||
|
if event.token() == STOP {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
loop {
|
||||||
|
let len = match fd.read(&mut buf[start..]) {
|
||||||
|
Ok(len) => len + start,
|
||||||
|
Err(e) => {
|
||||||
|
if e.kind() == io::ErrorKind::WouldBlock {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Err(e)?
|
||||||
|
}
|
||||||
|
};
|
||||||
|
//单线程的
|
||||||
|
up_counter.add(len as u64);
|
||||||
|
// buf是重复利用的,需要重置头部
|
||||||
|
buf[..12].fill(0);
|
||||||
|
match crate::handle::tun_tap::tun_handler::handle(
|
||||||
|
context,
|
||||||
|
&mut buf,
|
||||||
|
len,
|
||||||
|
&device,
|
||||||
|
current_device.load(),
|
||||||
|
&ip_route,
|
||||||
|
#[cfg(feature = "ip_proxy")]
|
||||||
|
&ip_proxy_map,
|
||||||
|
&client_cipher,
|
||||||
|
&server_cipher,
|
||||||
|
&device_list,
|
||||||
|
) {
|
||||||
|
Ok(_) => {}
|
||||||
|
Err(e) => {
|
||||||
|
log::warn!("{:?}", e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn start_multi(
|
||||||
|
stop_manager: StopManager,
|
||||||
|
device: Arc<Device>,
|
||||||
|
group_sync_sender: GroupSyncSender<(Vec<u8>, usize)>,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let poll = Poll::new()?;
|
||||||
|
let waker = Arc::new(Waker::new(poll.registry(), STOP)?);
|
||||||
|
let _waker = waker.clone();
|
||||||
|
let worker = stop_manager.add_listener("tun_device".into(), move || {
|
||||||
|
let _ = waker.wake();
|
||||||
|
})?;
|
||||||
|
if let Err(e) = start_multi0(poll, device, group_sync_sender, up_counter) {
|
||||||
|
log::error!("{:?}", e);
|
||||||
|
};
|
||||||
|
worker.stop_all();
|
||||||
|
drop(_waker);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn start_multi0(
|
||||||
|
mut poll: Poll,
|
||||||
|
device: Arc<Device>,
|
||||||
|
mut group_sync_sender: GroupSyncSender<(Vec<u8>, usize)>,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let fd = device.as_tun_fd();
|
||||||
|
fd.set_nonblock()?;
|
||||||
|
SourceFd(&fd.as_raw_fd()).register(poll.registry(), FD, Interest::READABLE)?;
|
||||||
|
let mut evnets = Events::with_capacity(4);
|
||||||
|
let mut buf = vec![0; 1024 * 16];
|
||||||
|
#[cfg(not(target_os = "macos"))]
|
||||||
|
let start = 12;
|
||||||
|
#[cfg(target_os = "macos")]
|
||||||
|
let start = 12 - 4;
|
||||||
|
loop {
|
||||||
|
poll.poll(&mut evnets, None)?;
|
||||||
|
for event in evnets.iter() {
|
||||||
|
if event.token() == STOP {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
loop {
|
||||||
|
let len = match fd.read(&mut buf[start..]) {
|
||||||
|
Ok(len) => len + start,
|
||||||
|
Err(e) => {
|
||||||
|
if e.kind() == io::ErrorKind::WouldBlock {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Err(e)?
|
||||||
|
}
|
||||||
|
};
|
||||||
|
//单线程的
|
||||||
|
up_counter.add(len as u64);
|
||||||
|
if group_sync_sender.send((buf, len)).is_err() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
buf = vec![0; 1024 * 16];
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,124 @@
|
|||||||
|
use crate::channel::context::ChannelContext;
|
||||||
|
use crate::cipher::Cipher;
|
||||||
|
use crate::external_route::ExternalRoute;
|
||||||
|
use crate::handle::tun_tap::channel_group::GroupSyncSender;
|
||||||
|
use crate::handle::{CurrentDeviceInfo, PeerDeviceInfo};
|
||||||
|
use crate::ip_proxy::IpProxyMap;
|
||||||
|
use crate::util::{SingleU64Adder, StopManager};
|
||||||
|
use crossbeam_utils::atomic::AtomicCell;
|
||||||
|
use parking_lot::Mutex;
|
||||||
|
use std::io;
|
||||||
|
use std::sync::Arc;
|
||||||
|
use tun::device::IFace;
|
||||||
|
use tun::Device;
|
||||||
|
|
||||||
|
pub(crate) fn start_simple(
|
||||||
|
stop_manager: StopManager,
|
||||||
|
context: &ChannelContext,
|
||||||
|
device: Arc<Device>,
|
||||||
|
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
|
||||||
|
ip_route: ExternalRoute,
|
||||||
|
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
|
||||||
|
client_cipher: Cipher,
|
||||||
|
server_cipher: Cipher,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let worker = {
|
||||||
|
let device = device.clone();
|
||||||
|
stop_manager.add_listener("tun_device".into(), move || {
|
||||||
|
if let Err(e) = device.shutdown() {
|
||||||
|
log::warn!("{:?}", e);
|
||||||
|
}
|
||||||
|
})?
|
||||||
|
};
|
||||||
|
if let Err(e) = start_simple0(
|
||||||
|
context,
|
||||||
|
device,
|
||||||
|
current_device,
|
||||||
|
ip_route,
|
||||||
|
#[cfg(feature = "ip_proxy")]
|
||||||
|
ip_proxy_map,
|
||||||
|
client_cipher,
|
||||||
|
server_cipher,
|
||||||
|
up_counter,
|
||||||
|
device_list,
|
||||||
|
) {
|
||||||
|
log::error!("{:?}", e);
|
||||||
|
}
|
||||||
|
worker.stop_all();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
fn start_simple0(
|
||||||
|
context: &ChannelContext,
|
||||||
|
device: Arc<Device>,
|
||||||
|
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
|
||||||
|
ip_route: ExternalRoute,
|
||||||
|
#[cfg(feature = "ip_proxy")] ip_proxy_map: Option<IpProxyMap>,
|
||||||
|
client_cipher: Cipher,
|
||||||
|
server_cipher: Cipher,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
device_list: Arc<Mutex<(u16, Vec<PeerDeviceInfo>)>>,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let mut buf = [0; 1024 * 16];
|
||||||
|
loop {
|
||||||
|
let len = device.read(&mut buf[12..])? + 12;
|
||||||
|
//单线程的
|
||||||
|
up_counter.add(len as u64);
|
||||||
|
// buf是重复利用的,需要重置头部
|
||||||
|
buf[..12].fill(0);
|
||||||
|
match crate::handle::tun_tap::tun_handler::handle(
|
||||||
|
context,
|
||||||
|
&mut buf,
|
||||||
|
len,
|
||||||
|
&device,
|
||||||
|
current_device.load(),
|
||||||
|
&ip_route,
|
||||||
|
#[cfg(feature = "ip_proxy")]
|
||||||
|
&ip_proxy_map,
|
||||||
|
&client_cipher,
|
||||||
|
&server_cipher,
|
||||||
|
&device_list,
|
||||||
|
) {
|
||||||
|
Ok(_) => {}
|
||||||
|
Err(e) => {
|
||||||
|
log::warn!("tun/tap {:?}", e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
pub(crate) fn start_multi(
|
||||||
|
stop_manager: StopManager,
|
||||||
|
device: Arc<Device>,
|
||||||
|
group_sync_sender: GroupSyncSender<(Vec<u8>, usize)>,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
let worker = {
|
||||||
|
let device = device.clone();
|
||||||
|
stop_manager.add_listener("tun_device_multi".into(), move || {
|
||||||
|
if let Err(e) = device.shutdown() {
|
||||||
|
log::warn!("{:?}", e);
|
||||||
|
}
|
||||||
|
})?
|
||||||
|
};
|
||||||
|
if let Err(e) = start_multi0(device, group_sync_sender, up_counter) {
|
||||||
|
log::error!("{:?}", e);
|
||||||
|
};
|
||||||
|
worker.stop_all();
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
fn start_multi0(
|
||||||
|
device: Arc<Device>,
|
||||||
|
mut group_sync_sender: GroupSyncSender<(Vec<u8>, usize)>,
|
||||||
|
up_counter: &mut SingleU64Adder,
|
||||||
|
) -> io::Result<()> {
|
||||||
|
loop {
|
||||||
|
let mut buf = vec![0; 1024 * 16];
|
||||||
|
let len = device.read(&mut buf[12..])? + 12;
|
||||||
|
//单线程的
|
||||||
|
up_counter.add(len as u64);
|
||||||
|
if group_sync_sender.send((buf, len)).is_err() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user