diff --git a/vnt/src/handle/tun_tap/tap_handler.rs b/vnt/src/handle/tun_tap/tap_handler.rs index 62bca57..670c856 100644 --- a/vnt/src/handle/tun_tap/tap_handler.rs +++ b/vnt/src/handle/tun_tap/tap_handler.rs @@ -22,6 +22,9 @@ use crate::handle::tun_tap::channel_group::{buf_channel_group, BufSenderGroup}; use crate::igmp_server::IgmpServer; use crate::ip_proxy::IpProxyMap; use crate::tun_tap_device::{DeviceReader, DeviceWriter}; +lazy_static! { + static ref POOL:BytePool> = BytePool::>::new(); +} pub fn start(worker: VntWorker, sender: ChannelSender, device_reader: DeviceReader, @@ -30,47 +33,59 @@ pub fn start(worker: VntWorker, sender: ChannelSender, current_device: Arc>, ip_route: Option, ip_proxy_map: Option, - client_cipher: Cipher, server_cipher: Cipher,parallel:usize) { - let (buf_sender, buf_receiver) = buf_channel_group(parallel); - for mut buf_receiver in buf_receiver.0 { - let sender = sender.clone(); - let device_writer = device_writer.clone(); - let igmp_server = igmp_server.clone(); - let current_device = current_device.clone(); - let ip_route = ip_route.clone(); - let ip_proxy_map = ip_proxy_map.clone(); - let client_cipher = client_cipher.clone(); - let server_cipher = server_cipher.clone(); - tokio::spawn(async move { - while let Some((mut buf, _, len)) = buf_receiver.recv().await { - match handle(&mut buf, len, &igmp_server, ¤t_device, &device_writer, &sender, - &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { - Ok(_) => {} - Err(e) => { - log::warn!("{:?}", e) + client_cipher: Cipher, server_cipher: Cipher, parallel: usize) { + if parallel == 1 { + thread::Builder::new().name("tap_handler".into()).spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all().build().unwrap() + .block_on(async move { + if let Err(e) = start_simple(sender, device_reader, + device_writer, igmp_server, + current_device, ip_route, ip_proxy_map, client_cipher, server_cipher).await { + log::warn!("tap:{:?}",e); + } + worker.stop_all(); + }); + }).unwrap(); + } else { + let (buf_sender, buf_receiver) = buf_channel_group(parallel); + for mut buf_receiver in buf_receiver.0 { + let sender = sender.clone(); + let device_writer = device_writer.clone(); + let igmp_server = igmp_server.clone(); + let current_device = current_device.clone(); + let ip_route = ip_route.clone(); + let ip_proxy_map = ip_proxy_map.clone(); + let client_cipher = client_cipher.clone(); + let server_cipher = server_cipher.clone(); + tokio::spawn(async move { + while let Some((mut buf, _, len)) = buf_receiver.recv().await { + match handle(&mut buf, len, &igmp_server, ¤t_device, &device_writer, &sender, + &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { + Ok(_) => {} + Err(e) => { + log::warn!("{:?}", e) + } } } - } - }); - } - thread::Builder::new().name("tap_handler".into()).spawn(move || { - tokio::runtime::Builder::new_current_thread() - .enable_all().build().unwrap() - .block_on(async move { - if let Err(e) = start_(sender, device_reader, buf_sender).await { - log::warn!("tap:{:?}",e); - } - worker.stop_all(); }); - }).unwrap(); -} -lazy_static!{ - static ref POOL:BytePool> = BytePool::>::new(); + } + thread::Builder::new().name("tap_handler".into()).spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all().build().unwrap() + .block_on(async move { + if let Err(e) = start_(sender, device_reader, buf_sender).await { + log::warn!("tap:{:?}",e); + } + worker.stop_all(); + }); + }).unwrap(); + } } + async fn start_(sender: ChannelSender, device_reader: DeviceReader, mut buf_sender: BufSenderGroup) -> io::Result<()> { - loop { let mut buf = POOL.alloc(4096); if sender.is_close() { @@ -84,6 +99,23 @@ async fn start_(sender: ChannelSender, } } +async fn start_simple(sender: ChannelSender, + device_reader: DeviceReader, + device_writer: DeviceWriter, + igmp_server: Option, + current_device: Arc>, + ip_route: Option, + ip_proxy_map: Option, + client_cipher: Cipher, server_cipher: Cipher) -> io::Result<()> { + let mut buf = [0; 4096]; + loop { + let len = device_reader.read(&mut buf)?; + if let Err(e) = handle(&mut buf, len, &igmp_server, ¤t_device, &device_writer, &sender, &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { + log::warn!("tap handle{:?}",e); + } + } +} + async fn handle(buf: &mut [u8], len: usize, igmp_server: &Option, current_device: &AtomicCell, device_writer: &DeviceWriter, sender: &ChannelSender, ip_route: &Option, proxy_map: &Option, client_cipher: &Cipher, server_cipher: &Cipher) -> crate::Result<()> { diff --git a/vnt/src/handle/tun_tap/tun_handler.rs b/vnt/src/handle/tun_tap/tun_handler.rs index e097b30..b69d752 100644 --- a/vnt/src/handle/tun_tap/tun_handler.rs +++ b/vnt/src/handle/tun_tap/tun_handler.rs @@ -69,40 +69,54 @@ pub async fn start(worker: VntWorker, sender: ChannelSender, ip_route: Option, ip_proxy_map: Option, client_cipher: Cipher, server_cipher: Cipher, parallel: usize) { - let (buf_sender, buf_receiver) = buf_channel_group(parallel); - for mut buf_receiver in buf_receiver.0 { - let sender = sender.clone(); - let device_writer = device_writer.clone(); - let igmp_server = igmp_server.clone(); - let current_device = current_device.clone(); - let ip_route = ip_route.clone(); - let ip_proxy_map = ip_proxy_map.clone(); - let client_cipher = client_cipher.clone(); - let server_cipher = server_cipher.clone(); - tokio::spawn(async move { - while let Some((mut buf, start, len)) = buf_receiver.recv().await { - match handle(&sender, &mut buf[start..], len, &device_writer, &igmp_server, current_device.load(), - &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { - Ok(_) => {} - Err(e) => { - log::warn!("{:?}", e) + if parallel == 1 { + thread::Builder::new().name("tun_handler".into()).spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all().build().unwrap() + .block_on(async move { + if let Err(e) = start_simple(sender, device_reader, &device_writer, igmp_server, current_device, ip_route, ip_proxy_map, client_cipher, server_cipher).await { + log::warn!("stop:{}",e); + } + let _ = device_writer.close(); + worker.stop_all(); + }) + }).unwrap(); + } else { + let (buf_sender, buf_receiver) = buf_channel_group(parallel); + for mut buf_receiver in buf_receiver.0 { + let sender = sender.clone(); + let device_writer = device_writer.clone(); + let igmp_server = igmp_server.clone(); + let current_device = current_device.clone(); + let ip_route = ip_route.clone(); + let ip_proxy_map = ip_proxy_map.clone(); + let client_cipher = client_cipher.clone(); + let server_cipher = server_cipher.clone(); + tokio::spawn(async move { + while let Some((mut buf, start, len)) = buf_receiver.recv().await { + match handle(&sender, &mut buf[start..], len, &device_writer, &igmp_server, current_device.load(), + &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { + Ok(_) => {} + Err(e) => { + log::warn!("{:?}", e) + } } } - } - }); - } + }); + } - thread::Builder::new().name("tun_handler".into()).spawn(move || { - tokio::runtime::Builder::new_current_thread() - .enable_all().build().unwrap() - .block_on(async move { - if let Err(e) = start_(sender, device_reader, buf_sender).await { - log::warn!("stop:{}",e); - } - let _ = device_writer.close(); - worker.stop_all(); - }) - }).unwrap(); + thread::Builder::new().name("tun_handler".into()).spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all().build().unwrap() + .block_on(async move { + if let Err(e) = start_(sender, device_reader, buf_sender).await { + log::warn!("stop:{}",e); + } + let _ = device_writer.close(); + worker.stop_all(); + }) + }).unwrap(); + } } async fn start_(sender: ChannelSender, device_reader: DeviceReader, mut buf_sender: BufSenderGroup) -> io::Result<()> { @@ -120,3 +134,28 @@ async fn start_(sender: ChannelSender, device_reader: DeviceReader, mut buf_send } } } + +async fn start_simple(sender: ChannelSender, + device_reader: DeviceReader, + device_writer: &DeviceWriter, + igmp_server: Option, + current_device: Arc>, + ip_route: Option, + ip_proxy_map: Option, + client_cipher: Cipher, server_cipher: Cipher) -> io::Result<()> { + let mut buf = [0; 4096]; + loop { + if sender.is_close() { + return Ok(()); + } + let len = device_reader.read(&mut buf[12..])? + 12; + #[cfg(any(target_os = "macos"))] + let mut buf = &mut buf[4..]; + match handle(&sender, &mut buf, len, device_writer, &igmp_server, current_device.load(), &ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { + Ok(_) => {} + Err(e) => { + log::warn!("{:?}", e) + } + } + } +}