Support timeouts while waiting for messages from robots
This commit is contained in:
@@ -8,6 +8,7 @@ pub mod robot_connector;
|
||||
pub mod robot_manager;
|
||||
|
||||
pub const ROBOT_MESSAGE_QUEUE_SIZE: usize = 10;
|
||||
const ROBOT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct ConnectorHandle {
|
||||
@@ -33,4 +34,6 @@ pub enum Error {
|
||||
ConnectionClosed(u32),
|
||||
#[error("No robots to handle requests")]
|
||||
NoRobotsToHandleRequest,
|
||||
#[error("Timed out waiting for response")]
|
||||
Timeout(#[from] tokio::time::error::Elapsed),
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use crate::robot::Error;
|
||||
use crate::robot::{Error, ROBOT_TIMEOUT};
|
||||
use log::{error, info};
|
||||
use prost::Message;
|
||||
use raas_types::raas::cmd::{Command, Request};
|
||||
@@ -58,7 +58,9 @@ impl RobotConnector {
|
||||
}
|
||||
|
||||
async fn get_response(&mut self) -> Result<Response, Error> {
|
||||
let recv = recv_raas_msg(&mut self.tcp_stream).await?;
|
||||
let recv_timeout = tokio::time::timeout(ROBOT_TIMEOUT, recv_raas_msg(&mut self.tcp_stream));
|
||||
let recv = recv_timeout.await??;
|
||||
|
||||
let resp = Response::decode(recv.msg.as_slice()).unwrap();
|
||||
|
||||
info!("Worker (bot_id={}) got resp", self.bot_id);
|
||||
@@ -79,53 +81,60 @@ impl RobotConnector {
|
||||
self.is_running = true;
|
||||
|
||||
loop {
|
||||
info!("Worker (bot_id={}) is waiting for requests", self.bot_id);
|
||||
let request = match self.wait_for_request().await {
|
||||
Ok(r) => r,
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to get request: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
break;
|
||||
}
|
||||
};
|
||||
|
||||
match self.send_command(request).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to send command to bot: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
let resp = match self.get_response().await {
|
||||
Ok(r) => r,
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to get response from bot: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
break;
|
||||
}
|
||||
};
|
||||
|
||||
match self.respond_to_request(resp).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to send response: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
break;
|
||||
}
|
||||
if let Err(res) = self.handle_request().await {
|
||||
error!("Closing robot connection: {}", res);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
error!("Exiting Worker (bot_id={})", self.bot_id);
|
||||
self.is_running = false;
|
||||
}
|
||||
|
||||
async fn handle_request(&mut self) -> Result<(), Error> {
|
||||
info!("Worker (bot_id={}) is waiting for requests", self.bot_id);
|
||||
let request = match self.wait_for_request().await {
|
||||
Ok(r) => r,
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to get request: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
|
||||
match self.send_command(request).await {
|
||||
Ok(_) => {}
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to send command to bot: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
return Err(err);
|
||||
}
|
||||
}
|
||||
|
||||
let resp = match self.get_response().await {
|
||||
Ok(r) => r,
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to get response from bot: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
return Err(err);
|
||||
}
|
||||
};
|
||||
|
||||
match self.respond_to_request(resp).await {
|
||||
Ok(_) => Ok(()),
|
||||
Err(err) => {
|
||||
error!(
|
||||
"Worker (bot_id={}) failed to send response: {:?}",
|
||||
self.bot_id, err
|
||||
);
|
||||
Err(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use crate::robot::{ConnectorManager, Error};
|
||||
use crate::robot::{ConnectorManager, Error, ROBOT_TIMEOUT};
|
||||
use log::{error, info};
|
||||
use prost::Message;
|
||||
use raas_types::raas;
|
||||
@@ -92,7 +92,9 @@ impl RobotManager {
|
||||
return Err(Error::ConnectionClosed(*id));
|
||||
}
|
||||
|
||||
let ret = recv_raas_msg(&mut connector.stream).await;
|
||||
let timeout = tokio::time::timeout(ROBOT_TIMEOUT, recv_raas_msg(&mut connector.stream));
|
||||
|
||||
let ret = timeout.await?;
|
||||
|
||||
let resp = match ret {
|
||||
Ok(r) => r,
|
||||
@@ -157,6 +159,11 @@ impl RobotManager {
|
||||
raas::error::ErrorType::NoRobotsToHandleRequest,
|
||||
"No robots to handle request".to_string(),
|
||||
),
|
||||
Error::Timeout(_) => (
|
||||
0,
|
||||
raas::error::ErrorType::RobotError,
|
||||
"Timed out waiting for response from Robot".to_string(),
|
||||
),
|
||||
};
|
||||
|
||||
self.send_error(id, error, error_msg.as_str()).await?;
|
||||
|
||||
Reference in New Issue
Block a user