1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
use std::sync::atomic::{AtomicBool, Ordering};
use zmq::{Error, SocketType};
use crate::socket::{MessageBuf, Sender, SocketBuilder, ZmqSocket};
pub fn request(endpoint: &str) -> Result<SocketBuilder<'_, Request>, zmq::Error> {
let socket = zmq::Context::new().socket(SocketType::REQ)?;
Ok(SocketBuilder::new(socket, endpoint))
}
pub struct Request {
inner: Sender,
received: AtomicBool,
}
impl From<zmq::Socket> for Request {
fn from(socket: zmq::Socket) -> Self {
Self {
inner: Sender {
socket: ZmqSocket::from(socket),
buffer: MessageBuf::default(),
},
received: AtomicBool::new(false),
}
}
}
impl Request {
pub async fn send<T: Into<MessageBuf>>(&self, msg: T) -> Result<(), Error> {
let mut msg = msg.into();
let res =
async_std::future::poll_fn(move |cx| self.inner.socket.send(cx, &mut msg)).await?;
self.received.store(false, Ordering::Relaxed);
Ok(res)
}
pub async fn recv(&self) -> Result<MessageBuf, Error> {
let msg = async_std::future::poll_fn(|cx| self.inner.socket.recv(cx)).await?;
self.received.store(true, Ordering::Relaxed);
Ok(msg)
}
pub async fn as_raw_socket(&self) -> &zmq::Socket {
&self.inner.socket.get_ref().0
}
}