|
@@ -13,7 +13,7 @@ use zeromq::*;
|
|
|
|
|
|
|
|
pub type Slabs = Vec<Vec<u8>>;
|
|
pub type Slabs = Vec<Vec<u8>>;
|
|
|
|
|
|
|
|
-pub struct GatewayService{
|
|
|
|
|
|
|
+pub struct GatewayService {
|
|
|
slabs: Slabs,
|
|
slabs: Slabs,
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -40,7 +40,7 @@ impl GatewayService {
|
|
|
NetEvent::RECEIVE(request) => {
|
|
NetEvent::RECEIVE(request) => {
|
|
|
ex2.spawn(Self::handle_request(send_queue_s.clone(), request))
|
|
ex2.spawn(Self::handle_request(send_queue_s.clone(), request))
|
|
|
.detach();
|
|
.detach();
|
|
|
- }
|
|
|
|
|
|
|
+ }
|
|
|
NetEvent::SEND(reply) => {
|
|
NetEvent::SEND(reply) => {
|
|
|
worker.send(reply).await?;
|
|
worker.send(reply).await?;
|
|
|
}
|
|
}
|
|
@@ -77,19 +77,20 @@ struct GatewayClient {
|
|
|
impl GatewayClient {
|
|
impl GatewayClient {
|
|
|
pub fn new() -> GatewayClient {
|
|
pub fn new() -> GatewayClient {
|
|
|
let sender = zeromq::ReqSocket::new();
|
|
let sender = zeromq::ReqSocket::new();
|
|
|
- GatewayClient { slabs: vec![], sender}
|
|
|
|
|
|
|
+ GatewayClient {
|
|
|
|
|
+ slabs: vec![],
|
|
|
|
|
+ sender,
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
pub async fn start(&mut self) -> Result<()> {
|
|
pub async fn start(&mut self) -> Result<()> {
|
|
|
-
|
|
|
|
|
self.sender.connect("tcp://127.0.0.1:3333").await?;
|
|
self.sender.connect("tcp://127.0.0.1:3333").await?;
|
|
|
Ok(())
|
|
Ok(())
|
|
|
-
|
|
|
|
|
}
|
|
}
|
|
|
- async fn request(&mut self, command: GatewayCommand, data: Vec<u8>) -> Result<Vec<u8>> {
|
|
|
|
|
|
|
+ async fn request(&mut self, command: GatewayCommand, data: Vec<u8>) -> Result<Vec<u8>> {
|
|
|
let request = Request::new(command as u8, data);
|
|
let request = Request::new(command as u8, data);
|
|
|
let req = serialize(&request);
|
|
let req = serialize(&request);
|
|
|
let req = bytes::Bytes::from(req);
|
|
let req = bytes::Bytes::from(req);
|
|
|
-
|
|
|
|
|
|
|
+
|
|
|
self.sender.send(req.into()).await?;
|
|
self.sender.send(req.into()).await?;
|
|
|
|
|
|
|
|
let rep: zeromq::ZmqMessage = self.sender.recv().await?;
|
|
let rep: zeromq::ZmqMessage = self.sender.recv().await?;
|
|
@@ -97,9 +98,9 @@ impl GatewayClient {
|
|
|
let rep: Vec<u8> = rep.to_vec();
|
|
let rep: Vec<u8> = rep.to_vec();
|
|
|
|
|
|
|
|
let reply: Reply = deserialize(&rep)?;
|
|
let reply: Reply = deserialize(&rep)?;
|
|
|
-
|
|
|
|
|
|
|
+
|
|
|
if reply.has_error() {
|
|
if reply.has_error() {
|
|
|
- return Err(ServicesError::ResonseError("response has an error").into());
|
|
|
|
|
|
|
+ return Err(ServicesError::ResonseError("response has an error").into());
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
assert!(reply.get_id() == request.get_id());
|
|
assert!(reply.get_id() == request.get_id());
|
|
@@ -108,14 +109,15 @@ impl GatewayClient {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
pub async fn get_slab(&mut self, index: u32) -> Result<Vec<u8>> {
|
|
pub async fn get_slab(&mut self, index: u32) -> Result<Vec<u8>> {
|
|
|
- self.request(GatewayCommand::GETSLAB, index.to_be_bytes().to_vec()).await
|
|
|
|
|
|
|
+ self.request(GatewayCommand::GETSLAB, index.to_be_bytes().to_vec())
|
|
|
|
|
+ .await
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- pub async fn put_slab(&mut self, data: Vec<u8>) -> Result<()>{
|
|
|
|
|
|
|
+ pub async fn put_slab(&mut self, data: Vec<u8>) -> Result<()> {
|
|
|
self.request(GatewayCommand::GETSLAB, data).await?;
|
|
self.request(GatewayCommand::GETSLAB, data).await?;
|
|
|
Ok(())
|
|
Ok(())
|
|
|
}
|
|
}
|
|
|
- pub async fn get_last_index(&mut self) -> Result<u32>{
|
|
|
|
|
|
|
+ pub async fn get_last_index(&mut self) -> Result<u32> {
|
|
|
let rep = self.request(GatewayCommand::GETLASTINDEX, vec![]).await?;
|
|
let rep = self.request(GatewayCommand::GETLASTINDEX, vec![]).await?;
|
|
|
let rep: [u8; 4] = rep.try_into().unwrap();
|
|
let rep: [u8; 4] = rep.try_into().unwrap();
|
|
|
Ok(u32::from_be_bytes(rep))
|
|
Ok(u32::from_be_bytes(rep))
|