parallel为1时不另起任务

This commit is contained in:
lubeilin
2023-08-27 22:55:40 +08:00
parent f9217625e1
commit cdf5c3a508
2 changed files with 136 additions and 65 deletions
+66 -34
View File
@@ -22,6 +22,9 @@ use crate::handle::tun_tap::channel_group::{buf_channel_group, BufSenderGroup};
use crate::igmp_server::IgmpServer; use crate::igmp_server::IgmpServer;
use crate::ip_proxy::IpProxyMap; use crate::ip_proxy::IpProxyMap;
use crate::tun_tap_device::{DeviceReader, DeviceWriter}; use crate::tun_tap_device::{DeviceReader, DeviceWriter};
lazy_static! {
static ref POOL:BytePool<Vec<u8>> = BytePool::<Vec<u8>>::new();
}
pub fn start(worker: VntWorker, sender: ChannelSender, pub fn start(worker: VntWorker, sender: ChannelSender,
device_reader: DeviceReader, device_reader: DeviceReader,
@@ -30,47 +33,59 @@ pub fn start(worker: VntWorker, sender: ChannelSender,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>, current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
ip_route: Option<ExternalRoute>, ip_route: Option<ExternalRoute>,
ip_proxy_map: Option<IpProxyMap>, ip_proxy_map: Option<IpProxyMap>,
client_cipher: Cipher, server_cipher: Cipher,parallel:usize) { client_cipher: Cipher, server_cipher: Cipher, parallel: usize) {
let (buf_sender, buf_receiver) = buf_channel_group(parallel); if parallel == 1 {
for mut buf_receiver in buf_receiver.0 { thread::Builder::new().name("tap_handler".into()).spawn(move || {
let sender = sender.clone(); tokio::runtime::Builder::new_current_thread()
let device_writer = device_writer.clone(); .enable_all().build().unwrap()
let igmp_server = igmp_server.clone(); .block_on(async move {
let current_device = current_device.clone(); if let Err(e) = start_simple(sender, device_reader,
let ip_route = ip_route.clone(); device_writer, igmp_server,
let ip_proxy_map = ip_proxy_map.clone(); current_device, ip_route, ip_proxy_map, client_cipher, server_cipher).await {
let client_cipher = client_cipher.clone(); log::warn!("tap:{:?}",e);
let server_cipher = server_cipher.clone(); }
tokio::spawn(async move { worker.stop_all();
while let Some((mut buf, _, len)) = buf_receiver.recv().await { });
match handle(&mut buf, len, &igmp_server, &current_device, &device_writer, &sender, }).unwrap();
&ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { } else {
Ok(_) => {} let (buf_sender, buf_receiver) = buf_channel_group(parallel);
Err(e) => { for mut buf_receiver in buf_receiver.0 {
log::warn!("{:?}", e) 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, &current_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(); }
} thread::Builder::new().name("tap_handler".into()).spawn(move || {
lazy_static!{ tokio::runtime::Builder::new_current_thread()
static ref POOL:BytePool<Vec<u8>> = BytePool::<Vec<u8>>::new(); .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, async fn start_(sender: ChannelSender,
device_reader: DeviceReader, device_reader: DeviceReader,
mut buf_sender: BufSenderGroup) -> io::Result<()> { mut buf_sender: BufSenderGroup) -> io::Result<()> {
loop { loop {
let mut buf = POOL.alloc(4096); let mut buf = POOL.alloc(4096);
if sender.is_close() { 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<IgmpServer>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
ip_route: Option<ExternalRoute>,
ip_proxy_map: Option<IpProxyMap>,
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, &current_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<IgmpServer>, current_device: &AtomicCell<CurrentDeviceInfo>, async fn handle(buf: &mut [u8], len: usize, igmp_server: &Option<IgmpServer>, current_device: &AtomicCell<CurrentDeviceInfo>,
device_writer: &DeviceWriter, sender: &ChannelSender, ip_route: &Option<ExternalRoute>, device_writer: &DeviceWriter, sender: &ChannelSender, ip_route: &Option<ExternalRoute>,
proxy_map: &Option<IpProxyMap>, client_cipher: &Cipher, server_cipher: &Cipher) -> crate::Result<()> { proxy_map: &Option<IpProxyMap>, client_cipher: &Cipher, server_cipher: &Cipher) -> crate::Result<()> {
+70 -31
View File
@@ -69,40 +69,54 @@ pub async fn start(worker: VntWorker, sender: ChannelSender,
ip_route: Option<ExternalRoute>, ip_route: Option<ExternalRoute>,
ip_proxy_map: Option<IpProxyMap>, ip_proxy_map: Option<IpProxyMap>,
client_cipher: Cipher, server_cipher: Cipher, parallel: usize) { client_cipher: Cipher, server_cipher: Cipher, parallel: usize) {
let (buf_sender, buf_receiver) = buf_channel_group(parallel); if parallel == 1 {
for mut buf_receiver in buf_receiver.0 { thread::Builder::new().name("tun_handler".into()).spawn(move || {
let sender = sender.clone(); tokio::runtime::Builder::new_current_thread()
let device_writer = device_writer.clone(); .enable_all().build().unwrap()
let igmp_server = igmp_server.clone(); .block_on(async move {
let current_device = current_device.clone(); 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 {
let ip_route = ip_route.clone(); log::warn!("stop:{}",e);
let ip_proxy_map = ip_proxy_map.clone(); }
let client_cipher = client_cipher.clone(); let _ = device_writer.close();
let server_cipher = server_cipher.clone(); worker.stop_all();
tokio::spawn(async move { })
while let Some((mut buf, start, len)) = buf_receiver.recv().await { }).unwrap();
match handle(&sender, &mut buf[start..], len, &device_writer, &igmp_server, current_device.load(), } else {
&ip_route, &ip_proxy_map, &client_cipher, &server_cipher).await { let (buf_sender, buf_receiver) = buf_channel_group(parallel);
Ok(_) => {} for mut buf_receiver in buf_receiver.0 {
Err(e) => { let sender = sender.clone();
log::warn!("{:?}", e) 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 || { thread::Builder::new().name("tun_handler".into()).spawn(move || {
tokio::runtime::Builder::new_current_thread() tokio::runtime::Builder::new_current_thread()
.enable_all().build().unwrap() .enable_all().build().unwrap()
.block_on(async move { .block_on(async move {
if let Err(e) = start_(sender, device_reader, buf_sender).await { if let Err(e) = start_(sender, device_reader, buf_sender).await {
log::warn!("stop:{}",e); log::warn!("stop:{}",e);
} }
let _ = device_writer.close(); let _ = device_writer.close();
worker.stop_all(); worker.stop_all();
}) })
}).unwrap(); }).unwrap();
}
} }
async fn start_(sender: ChannelSender, device_reader: DeviceReader, mut buf_sender: BufSenderGroup) -> io::Result<()> { 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<IgmpServer>,
current_device: Arc<AtomicCell<CurrentDeviceInfo>>,
ip_route: Option<ExternalRoute>,
ip_proxy_map: Option<IpProxyMap>,
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)
}
}
}
}