|
|
@@ -5,20 +5,18 @@ use std::net::SocketAddr;
|
|
|
use crate::serial::{deserialize, serialize};
|
|
|
use crate::{Decodable, Encodable, Result};
|
|
|
|
|
|
+use async_executor::Executor;
|
|
|
use bytes::Bytes;
|
|
|
use futures::FutureExt;
|
|
|
+use log::*;
|
|
|
use rand::Rng;
|
|
|
+use signal_hook::{consts::SIGINT, iterator::Signals};
|
|
|
use zeromq::*;
|
|
|
-use signal_hook::{iterator::Signals, consts::SIGINT};
|
|
|
-use async_executor::Executor;
|
|
|
-use log::*;
|
|
|
-
|
|
|
-
|
|
|
|
|
|
enum NetEvent {
|
|
|
Receive(zeromq::ZmqMessage),
|
|
|
Send(Reply),
|
|
|
- Stop
|
|
|
+ Stop,
|
|
|
}
|
|
|
|
|
|
pub fn addr_to_string(addr: SocketAddr) -> String {
|
|
|
@@ -58,8 +56,8 @@ impl RepProtocol {
|
|
|
pub async fn start(
|
|
|
&mut self,
|
|
|
) -> Result<(
|
|
|
- async_channel::Sender<Reply>,
|
|
|
- async_channel::Receiver<Request>,
|
|
|
+ async_channel::Sender<Reply>,
|
|
|
+ async_channel::Receiver<Request>,
|
|
|
)> {
|
|
|
let addr = addr_to_string(self.addr);
|
|
|
self.socket.bind(addr.as_str()).await?;
|
|
|
@@ -68,7 +66,6 @@ impl RepProtocol {
|
|
|
}
|
|
|
|
|
|
pub async fn run(&mut self, executor: Arc<Executor<'_>>) -> Result<()> {
|
|
|
-
|
|
|
info!("{} SERVICE: running", self.service_name);
|
|
|
|
|
|
let (stop_s, stop_r) = async_channel::unbounded::<()>();
|
|
|
@@ -93,19 +90,18 @@ impl RepProtocol {
|
|
|
match event {
|
|
|
NetEvent::Receive(request) => {
|
|
|
if let Some(req) = request.get(0) {
|
|
|
- let request: Vec<u8> = req.to_vec();
|
|
|
- let req: Request = deserialize(&request)?;
|
|
|
+ let req: Vec<u8> = req.to_vec();
|
|
|
+ let req: Request = deserialize(&req)?;
|
|
|
self.send_queue.send(req).await?;
|
|
|
}
|
|
|
}
|
|
|
NetEvent::Send(reply) => {
|
|
|
let reply: Vec<u8> = serialize(&reply);
|
|
|
let reply = Bytes::from(reply);
|
|
|
- self.socket.send(reply.into()).await?;
|
|
|
- }
|
|
|
- NetEvent::Stop => {
|
|
|
- break
|
|
|
+ let reply: zeromq::ZmqMessage = reply.into();
|
|
|
+ self.socket.send(reply).await?;
|
|
|
}
|
|
|
+ NetEvent::Stop => break,
|
|
|
}
|
|
|
}
|
|
|
let _ = stop_task.cancel().await;
|
|
|
@@ -135,14 +131,15 @@ impl ReqProtocol {
|
|
|
let request = Request::new(command, data);
|
|
|
let req = serialize(&request);
|
|
|
let req = bytes::Bytes::from(req);
|
|
|
+ let req: zeromq::ZmqMessage = req.into();
|
|
|
|
|
|
- self.socket.send(req.into()).await?;
|
|
|
+ self.socket.send(req).await?;
|
|
|
|
|
|
let rep: zeromq::ZmqMessage = self.socket.recv().await?;
|
|
|
- if let Some(rep) = rep.get(0) {
|
|
|
- let rep: Vec<u8> = rep.to_vec();
|
|
|
+ if let Some(reply) = rep.get(0) {
|
|
|
+ let reply: Vec<u8> = reply.to_vec();
|
|
|
|
|
|
- let reply: Reply = deserialize(&rep)?;
|
|
|
+ let reply: Reply = deserialize(&reply)?;
|
|
|
|
|
|
if reply.has_error() {
|
|
|
return Err(crate::Error::ServicesError("response has an error"));
|
|
|
@@ -152,7 +149,9 @@ impl ReqProtocol {
|
|
|
|
|
|
Ok(reply.get_payload())
|
|
|
} else {
|
|
|
- Err(crate::Error::ZMQError("Couldn't parse ZmqMessage".to_string()))
|
|
|
+ Err(crate::Error::ZMQError(
|
|
|
+ "Couldn't parse ZmqMessage".to_string(),
|
|
|
+ ))
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
@@ -207,11 +206,10 @@ impl Subscriber {
|
|
|
let data = d.to_vec();
|
|
|
Ok(data)
|
|
|
}
|
|
|
- None => {
|
|
|
- Err(crate::Error::ZMQError("Couldn't parse ZmqMessage".to_string()))
|
|
|
- }
|
|
|
+ None => Err(crate::Error::ZMQError(
|
|
|
+ "Couldn't parse ZmqMessage".to_string(),
|
|
|
+ )),
|
|
|
}
|
|
|
-
|
|
|
}
|
|
|
}
|
|
|
|