Преглед изворни кода

change address type to std::net::SocketAddr

ghassmo пре 5 година
родитељ
комит
baec13555c
3 измењених фајлова са 29 додато и 18 уклоњено
  1. 2 2
      src/bin/demowallet.rs
  2. 5 4
      src/service/gateway.rs
  3. 22 12
      src/service/reqrep.rs

+ 2 - 2
src/bin/demowallet.rs

@@ -6,7 +6,7 @@ use sapvi::service::{fetch_slabs_loop, GatewayClient};
 use sapvi::Result;
 use sapvi::Result;
 
 
 async fn start(executor: Arc<Executor<'_>>) -> Result<()> {
 async fn start(executor: Arc<Executor<'_>>) -> Result<()> {
-    let mut client = GatewayClient::new(String::from("tcp://127.0.0.1:3333"));
+    let mut client = GatewayClient::new("127.0.0.1:3333".parse()?);
 
 
     client.start().await?;
     client.start().await?;
     println!("connected to a server");
     println!("connected to a server");
@@ -14,7 +14,7 @@ async fn start(executor: Arc<Executor<'_>>) -> Result<()> {
     let slabs = Arc::new(Mutex::new(vec![]));
     let slabs = Arc::new(Mutex::new(vec![]));
 
 
     let subscriber = client
     let subscriber = client
-        .subscribe(String::from("tcp://127.0.0.1:4444"))
+        .subscribe("127.0.0.1:4444".parse()?)
         .await?;
         .await?;
 
 
     println!("subscription ready");
     println!("subscription ready");

+ 5 - 4
src/service/gateway.rs

@@ -1,5 +1,6 @@
 use async_std::sync::{Arc, Mutex};
 use async_std::sync::{Arc, Mutex};
 use std::convert::TryInto;
 use std::convert::TryInto;
+use std::net::SocketAddr; 
 
 
 use super::reqrep::{Publisher, RepProtocol, Reply, ReqProtocol, Request, Subscriber};
 use super::reqrep::{Publisher, RepProtocol, Reply, ReqProtocol, Request, Subscriber};
 use crate::{Error, Result};
 use crate::{Error, Result};
@@ -17,12 +18,12 @@ enum GatewayCommand {
 
 
 pub struct GatewayService {
 pub struct GatewayService {
     slabs: Mutex<Slabs>,
     slabs: Mutex<Slabs>,
-    addr: String,
+    addr: SocketAddr,
     publisher: Mutex<Publisher>,
     publisher: Mutex<Publisher>,
 }
 }
 
 
 impl GatewayService {
 impl GatewayService {
-    pub fn new(addr: String, pub_addr: String) -> Arc<GatewayService> {
+    pub fn new(addr: SocketAddr, pub_addr: SocketAddr) -> Arc<GatewayService> {
         let slabs = Mutex::new(vec![]);
         let slabs = Mutex::new(vec![]);
         let publisher = Mutex::new(Publisher::new(pub_addr));
         let publisher = Mutex::new(Publisher::new(pub_addr));
         Arc::new(GatewayService {
         Arc::new(GatewayService {
@@ -97,7 +98,7 @@ pub struct GatewayClient {
 }
 }
 
 
 impl GatewayClient {
 impl GatewayClient {
-    pub fn new(addr: String) -> GatewayClient {
+    pub fn new(addr: SocketAddr) -> GatewayClient {
         let protocol = ReqProtocol::new(addr);
         let protocol = ReqProtocol::new(addr);
         GatewayClient { protocol }
         GatewayClient { protocol }
     }
     }
@@ -106,7 +107,7 @@ impl GatewayClient {
         Ok(())
         Ok(())
     }
     }
 
 
-    pub async fn subscribe(&self, sub_addr: String) -> Result<Arc<Mutex<Subscriber>>> {
+    pub async fn subscribe(&self, sub_addr: SocketAddr) -> Result<Arc<Mutex<Subscriber>>> {
         let mut subscriber = Subscriber::new(sub_addr);
         let mut subscriber = Subscriber::new(sub_addr);
         subscriber.start().await?;
         subscriber.start().await?;
         Ok(Arc::new(Mutex::new(subscriber)))
         Ok(Arc::new(Mutex::new(subscriber)))

+ 22 - 12
src/service/reqrep.rs

@@ -1,4 +1,5 @@
 use std::io;
 use std::io;
+use std::net::SocketAddr;
 
 
 use crate::serial::{deserialize, serialize};
 use crate::serial::{deserialize, serialize};
 use crate::{Decodable, Encodable, Result};
 use crate::{Decodable, Encodable, Result};
@@ -13,8 +14,13 @@ enum NetEvent {
     Send(Reply),
     Send(Reply),
 }
 }
 
 
+
+pub fn addr_to_string (addr: SocketAddr) -> String {
+    format!("tcp://{}", addr.to_string())
+}
+
 pub struct RepProtocol {
 pub struct RepProtocol {
-    addr: String,
+    addr: SocketAddr,
     socket: zeromq::RepSocket,
     socket: zeromq::RepSocket,
     recv_queue: async_channel::Receiver<Reply>,
     recv_queue: async_channel::Receiver<Reply>,
     send_queue: async_channel::Sender<Request>,
     send_queue: async_channel::Sender<Request>,
@@ -25,7 +31,7 @@ pub struct RepProtocol {
 }
 }
 
 
 impl RepProtocol {
 impl RepProtocol {
-    pub fn new(addr: String) -> RepProtocol {
+    pub fn new(addr: SocketAddr) -> RepProtocol {
         let socket = zeromq::RepSocket::new();
         let socket = zeromq::RepSocket::new();
         let (send_queue, recv_channel) = async_channel::unbounded::<Request>();
         let (send_queue, recv_channel) = async_channel::unbounded::<Request>();
         let (send_channel, recv_queue) = async_channel::unbounded::<Reply>();
         let (send_channel, recv_queue) = async_channel::unbounded::<Reply>();
@@ -47,7 +53,8 @@ impl RepProtocol {
         async_channel::Sender<Reply>,
         async_channel::Sender<Reply>,
         async_channel::Receiver<Request>,
         async_channel::Receiver<Request>,
     )> {
     )> {
-        self.socket.bind(self.addr.as_str()).await?;
+        let addr = addr_to_string(self.addr);
+        self.socket.bind(addr.as_str()).await?;
         Ok(self.channels.clone())
         Ok(self.channels.clone())
     }
     }
 
 
@@ -76,18 +83,19 @@ impl RepProtocol {
 }
 }
 
 
 pub struct ReqProtocol {
 pub struct ReqProtocol {
-    addr: String,
+    addr: SocketAddr,
     socket: zeromq::ReqSocket,
     socket: zeromq::ReqSocket,
 }
 }
 
 
 impl ReqProtocol {
 impl ReqProtocol {
-    pub fn new(addr: String) -> ReqProtocol {
+    pub fn new(addr: SocketAddr) -> ReqProtocol {
         let socket = zeromq::ReqSocket::new();
         let socket = zeromq::ReqSocket::new();
         ReqProtocol { addr, socket }
         ReqProtocol { addr, socket }
     }
     }
 
 
     pub async fn start(&mut self) -> Result<()> {
     pub async fn start(&mut self) -> Result<()> {
-        self.socket.connect(self.addr.as_str()).await?;
+        let addr = addr_to_string(self.addr);
+        self.socket.connect(addr.as_str()).await?;
         Ok(())
         Ok(())
     }
     }
 
 
@@ -115,17 +123,18 @@ impl ReqProtocol {
 }
 }
 
 
 pub struct Publisher {
 pub struct Publisher {
-    addr: String,
+    addr: SocketAddr,
     socket: zeromq::PubSocket,
     socket: zeromq::PubSocket,
 }
 }
 
 
 impl Publisher {
 impl Publisher {
-    pub fn new(addr: String) -> Publisher {
+    pub fn new(addr: SocketAddr) -> Publisher {
         let socket = zeromq::PubSocket::new();
         let socket = zeromq::PubSocket::new();
         Publisher { addr, socket }
         Publisher { addr, socket }
     }
     }
     pub async fn start(&mut self) -> Result<()> {
     pub async fn start(&mut self) -> Result<()> {
-        self.socket.bind(self.addr.as_str()).await?;
+        let addr = addr_to_string(self.addr);
+        self.socket.bind(addr.as_str()).await?;
         Ok(())
         Ok(())
     }
     }
 
 
@@ -137,18 +146,19 @@ impl Publisher {
 }
 }
 
 
 pub struct Subscriber {
 pub struct Subscriber {
-    addr: String,
+    addr: SocketAddr,
     socket: zeromq::SubSocket,
     socket: zeromq::SubSocket,
 }
 }
 
 
 impl Subscriber {
 impl Subscriber {
-    pub fn new(addr: String) -> Subscriber {
+    pub fn new(addr: SocketAddr) -> Subscriber {
         let socket = zeromq::SubSocket::new();
         let socket = zeromq::SubSocket::new();
         Subscriber { addr, socket }
         Subscriber { addr, socket }
     }
     }
 
 
     pub async fn start(&mut self) -> Result<()> {
     pub async fn start(&mut self) -> Result<()> {
-        self.socket.connect(self.addr.as_str()).await?;
+        let addr = addr_to_string(self.addr);
+        self.socket.connect(addr.as_str()).await?;
 
 
         self.socket.subscribe("").await?;
         self.socket.subscribe("").await?;