修改对称网络处理逻辑

This commit is contained in:
lubeilin
2023-12-31 13:42:32 +08:00
parent 858ca9bbe7
commit 76eaead0b9
+29 -82
View File
@@ -523,7 +523,6 @@ impl Channel {
let context = context.clone(); let context = context.clone();
let main_channel = main_channel.try_clone().unwrap(); let main_channel = main_channel.try_clone().unwrap();
let handler = handler.clone(); let handler = handler.clone();
let buf_sender = buf_sender.clone();
thread::Builder::new() thread::Builder::new()
.name("channel_udp".into()) .name("channel_udp".into())
.spawn(move || { .spawn(move || {
@@ -570,7 +569,7 @@ impl Channel {
Ok(udp) => { Ok(udp) => {
let udp = Arc::new(udp); let udp = Arc::new(udp);
let context = context.clone(); let context = context.clone();
tokio::spawn(Self::start_(worker.worker("symmetric_channel"),context, udp,handler.clone(),buf_sender.clone(), head_reserve, false)); tokio::spawn(Self::start_(worker.worker("symmetric_channel"),context, udp,handler.clone(), head_reserve));
} }
Err(e) => { Err(e) => {
log::error!("{}",e); log::error!("{}",e);
@@ -653,9 +652,7 @@ impl Channel {
context: Context, context: Context,
udp: Arc<UdpSocket>, udp: Arc<UdpSocket>,
handler: ChannelDataHandler, handler: ChannelDataHandler,
buf_sender: Option<BufSenderGroup>,
head_reserve: usize, head_reserve: usize,
is_core: bool,
) { ) {
let mut status_receiver = context.inner.status_receiver.clone(); let mut status_receiver = context.inner.status_receiver.clone();
#[cfg(target_os = "windows")] #[cfg(target_os = "windows")]
@@ -668,92 +665,42 @@ impl Channel {
let id = 3 + udp.as_raw_fd() as usize; let id = 3 + udp.as_raw_fd() as usize;
context.insert_udp(id, udp.clone()); context.insert_udp(id, udp.clone());
match buf_sender { let mut buf = [0; 4096];
None => { loop {
let mut buf = [0; 4096]; tokio::select! {
loop { rs=udp.recv_from(&mut buf[head_reserve..])=>{
tokio::select! { match rs {
rs=udp.recv_from(&mut buf[head_reserve..])=>{ Ok((len, addr)) => {
match rs { handler.handle(&mut buf, head_reserve, head_reserve + len, RouteKey::new(id, addr), &context);
Ok((len, addr)) => {
handler.handle(&mut buf, head_reserve, head_reserve + len, RouteKey::new(id, addr), &context);
}
Err(e) => {
log::error!("{:?}",e)
}
}
} }
changed=status_receiver.changed()=>{ Err(e) => {
match changed { log::error!("{:?}",e)
Ok(_) => { }
match *status_receiver.borrow() { }
Status::Cone => { }
if !is_core{ changed=status_receiver.changed()=>{
break; match changed {
} Ok(_) => {
} match *status_receiver.borrow() {
Status::Close=>{ Status::Cone => {
break;
}
Status::Symmetric => {}
}
}
Err(_) => {
break; break;
} }
Status::Close=>{
break;
}
Status::Symmetric => {}
} }
}
Err(_) => {
break;
}
} }
_=worker.stop_wait()=>{ }
break; _=worker.stop_wait()=>{
} break;
}
} }
} }
Some(mut buf_sender) => loop {
let mut buf = vec![0; 4096];
tokio::select! {
rs=udp.recv_from(&mut buf[head_reserve..])=>{
match rs {
Ok((len, addr)) => {
if !buf_sender.send((buf,head_reserve,head_reserve+len,RouteKey::new(id, addr))){
log::error!("udp buf_sender发送数据失败");
break;
}
}
Err(e) => {
log::error!("{:?}",e)
}
}
}
changed=status_receiver.changed()=>{
match changed {
Ok(_) => {
match *status_receiver.borrow() {
Status::Cone => {
if !is_core{
break;
}
}
Status::Close=>{
break;
}
Status::Symmetric => {}
}
}
Err(_) => {
break;
}
}
}
_=worker.stop_wait()=>{
break;
}
}
},
} }
context.remove_udp(id); context.remove_udp(id);
if is_core {
worker.stop_all();
}
} }
} }