Replace select! with into_split in beacon::client
This commit is contained in:
parent
d5d4479a53
commit
41d7487ab0
1 changed files with 28 additions and 20 deletions
|
|
@ -3,12 +3,12 @@ use crate::core::time::{Duration, SystemTime};
|
||||||
use crate::span;
|
use crate::span;
|
||||||
use crate::theme;
|
use crate::theme;
|
||||||
|
|
||||||
use futures::{FutureExt, select};
|
|
||||||
use semver::Version;
|
use semver::Version;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
|
use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
|
||||||
use tokio::net;
|
use tokio::net;
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::{RwLock, mpsc};
|
||||||
|
use tokio::task;
|
||||||
use tokio::time;
|
use tokio::time;
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
@ -116,14 +116,16 @@ async fn run(
|
||||||
let mut buffer = Vec::new();
|
let mut buffer = Vec::new();
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let mut command_sender = None;
|
let command_sender = Arc::new(RwLock::new(None));
|
||||||
|
|
||||||
match _connect().await {
|
match _connect().await {
|
||||||
Ok(mut stream) => {
|
Ok(stream) => {
|
||||||
is_connected.store(true, atomic::Ordering::Relaxed);
|
is_connected.store(true, atomic::Ordering::Relaxed);
|
||||||
|
|
||||||
|
let (mut reader, mut writer) = stream.into_split();
|
||||||
|
|
||||||
let _ = send(
|
let _ = send(
|
||||||
&mut stream,
|
&mut writer,
|
||||||
Message::Connected {
|
Message::Connected {
|
||||||
at: SystemTime::now(),
|
at: SystemTime::now(),
|
||||||
name: name.clone(),
|
name: name.clone(),
|
||||||
|
|
@ -132,17 +134,18 @@ async fn run(
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
loop {
|
{
|
||||||
select! {
|
let command_sender = command_sender.clone();
|
||||||
action = receiver.recv().fuse() => {
|
|
||||||
let Some(action) = action else { break; };
|
|
||||||
|
|
||||||
|
drop(task::spawn(async move {
|
||||||
|
while let Some(action) = receiver.recv().await {
|
||||||
match action {
|
match action {
|
||||||
Action::Send(message) => {
|
Action::Send(message) => {
|
||||||
match send(&mut stream, message).await {
|
match send(&mut writer, message).await {
|
||||||
Ok(()) => {}
|
Ok(()) => {}
|
||||||
Err(error) => {
|
Err(error) => {
|
||||||
if error.kind() != io::ErrorKind::BrokenPipe
|
if error.kind()
|
||||||
|
!= io::ErrorKind::BrokenPipe
|
||||||
{
|
{
|
||||||
log::warn!(
|
log::warn!(
|
||||||
"Error sending message to server: {error}"
|
"Error sending message to server: {error}"
|
||||||
|
|
@ -153,17 +156,22 @@ async fn run(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Action::Forward(sender) => {
|
Action::Forward(sender) => {
|
||||||
command_sender = Some(sender);
|
*command_sender.write().await =
|
||||||
|
Some(sender);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
command = receive(&mut stream, &mut buffer).fuse() => {
|
}))
|
||||||
let Ok(command) = command else { continue; };
|
};
|
||||||
|
|
||||||
if let Some(sender) = command_sender.as_mut() {
|
loop {
|
||||||
let _ = sender.send(command).await;
|
let Ok(command) = receive(&mut reader, &mut buffer).await
|
||||||
}
|
else {
|
||||||
}
|
continue;
|
||||||
|
};
|
||||||
|
|
||||||
|
if let Some(sender) = command_sender.read().await.as_ref() {
|
||||||
|
let _ = sender.send(command).await;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -186,7 +194,7 @@ async fn _connect() -> Result<net::TcpStream, io::Error> {
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn send(
|
async fn send(
|
||||||
stream: &mut net::TcpStream,
|
stream: &mut net::tcp::OwnedWriteHalf,
|
||||||
message: Message,
|
message: Message,
|
||||||
) -> Result<(), io::Error> {
|
) -> Result<(), io::Error> {
|
||||||
let bytes = bincode::serialize(&message).expect("Encode input message");
|
let bytes = bincode::serialize(&message).expect("Encode input message");
|
||||||
|
|
@ -200,7 +208,7 @@ async fn send(
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn receive(
|
async fn receive(
|
||||||
stream: &mut net::TcpStream,
|
stream: &mut net::tcp::OwnedReadHalf,
|
||||||
buffer: &mut Vec<u8>,
|
buffer: &mut Vec<u8>,
|
||||||
) -> Result<Command, Error> {
|
) -> Result<Command, Error> {
|
||||||
let size = stream.read_u64().await? as usize;
|
let size = stream.read_u64().await? as usize;
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue