Fix disconnection logic in beacon::client
This commit is contained in:
parent
162f8c0c29
commit
5b649541b6
1 changed files with 38 additions and 29 deletions
|
|
@ -113,11 +113,14 @@ async fn run(
|
||||||
let version = semver::Version::parse(env!("CARGO_PKG_VERSION"))
|
let version = semver::Version::parse(env!("CARGO_PKG_VERSION"))
|
||||||
.expect("Parse package version");
|
.expect("Parse package version");
|
||||||
|
|
||||||
let mut buffer = Vec::new();
|
let command_sender = {
|
||||||
|
// Discard by default
|
||||||
|
let (sender, _receiver) = mpsc::channel(1);
|
||||||
|
|
||||||
|
Arc::new(Mutex::new(sender))
|
||||||
|
};
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let command_sender = Arc::new(Mutex::new(None));
|
|
||||||
|
|
||||||
match _connect().await {
|
match _connect().await {
|
||||||
Ok(stream) => {
|
Ok(stream) => {
|
||||||
is_connected.store(true, atomic::Ordering::Relaxed);
|
is_connected.store(true, atomic::Ordering::Relaxed);
|
||||||
|
|
@ -138,39 +141,45 @@ async fn run(
|
||||||
let command_sender = command_sender.clone();
|
let command_sender = command_sender.clone();
|
||||||
|
|
||||||
drop(task::spawn(async move {
|
drop(task::spawn(async move {
|
||||||
while let Some(action) = receiver.recv().await {
|
let mut buffer = Vec::new();
|
||||||
match action {
|
|
||||||
Action::Send(message) => {
|
loop {
|
||||||
match send(&mut writer, message).await {
|
match receive(&mut reader, &mut buffer).await {
|
||||||
Ok(()) => {}
|
Ok(command) => {
|
||||||
Err(error) => {
|
let sender = command_sender.lock().await;
|
||||||
if error.kind()
|
let _ = sender.send(command).await;
|
||||||
!= io::ErrorKind::BrokenPipe
|
|
||||||
{
|
|
||||||
log::warn!(
|
|
||||||
"Error sending message to server: {error}"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Action::Forward(sender) => {
|
|
||||||
*command_sender.lock().await = Some(sender);
|
|
||||||
}
|
}
|
||||||
|
Err(Error::DecodingFailed(_)) => {}
|
||||||
|
Err(Error::IOFailed(_)) => break,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}))
|
}))
|
||||||
};
|
};
|
||||||
|
|
||||||
loop {
|
while let Some(action) = receiver.recv().await {
|
||||||
let Ok(command) = receive(&mut reader, &mut buffer).await
|
match action {
|
||||||
else {
|
Action::Send(message) => {
|
||||||
continue;
|
match send(&mut writer, message).await {
|
||||||
};
|
Ok(()) => {}
|
||||||
|
Err(error) => {
|
||||||
|
if error.kind() != io::ErrorKind::BrokenPipe
|
||||||
|
{
|
||||||
|
log::warn!(
|
||||||
|
"Error sending message to server: {error}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
if let Some(sender) = command_sender.lock().await.as_ref() {
|
is_connected.store(
|
||||||
let _ = sender.send(command).await;
|
false,
|
||||||
|
atomic::Ordering::Relaxed,
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Action::Forward(sender) => {
|
||||||
|
*command_sender.lock().await = sender;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue