PhoenixZMQ  8.1.3
Library which integrates zeromq use
Loading...
Searching...
No Matches
PZmqBackend.cpp
Go to the documentation of this file.
1/***************************************
2 Auteur : Pierre Aubert
3 Mail : pierre.aubert@lapp.in2p3.fr
4 Licence : CeCILL-C
5****************************************/
6
7#include "PZmqBackend.h"
8
10
20PSendStatus::PSendStatus checkSendStatus(zmq::send_result_t res){
21 if(res.has_value()){
22 return PSendStatus::OK;
23 }else{
24 int err = zmq_errno();
25 if(err == EAGAIN){
26 return PSendStatus::SOCKET_NOT_AVAILABLE;
27 }else{
28 std::cerr << "Unknown ZMQ error in send: '" << err << "'"<< std::endl;
29 return PSendStatus::BROKEN_BACKEND;
30 }
31 }
32}
33
35
43PRecvStatus::PRecvStatus checkRecvStatus(zmq::recv_result_t res){
44 if(res.has_value()){
45 return PRecvStatus::OK;
46 }else{
47 int err = zmq_errno();
48 if(err == EAGAIN){
49 return PRecvStatus::NO_MESSAGE_RECEIVED;
50 }else{
51 return PRecvStatus::BROKEN_SOCKET;
52 }
53 }
54}
55
57
60zmq::send_flags convertToSendFlag(PSendFlag::PSendFlag flag){
61 if(flag == PSendFlag::NON_BLOCK){return zmq::send_flags::dontwait;}
62 else{return zmq::send_flags::none;}
63}
64
66
69zmq::recv_flags convertToRecvFlag(PRecvFlag::PRecvFlag flag){
70 if(flag == PRecvFlag::NON_BLOCK){return zmq::recv_flags::dontwait;}
71 else{return zmq::recv_flags::none;}
72}
73
75
84PZmqParam pzmq_createParamClient(int type, int nbBufferMessage, int bufferSizeByte, size_t threadAffinity, ssize_t dataRate, int linger, int immediate){
85 PZmqParam param;
86 param.type = type;
87 param.nbBufferMessage = nbBufferMessage;
88 param.bufferSizeByte = bufferSizeByte;
89 param.threadAffinity = threadAffinity;
90 param.dataRate = dataRate;
91 param.linger = linger;
92 param.immediate = immediate; //EAGAIN if no peer connected
93 return param;
94}
95
97
105PZmqParam pzmq_createParamServer(int type, int nbBufferMessage, int bufferSizeByte, size_t threadAffinity, ssize_t dataRate, int linger, int immediate){
106 PZmqParam param;
107 param.type = type;
108 param.nbBufferMessage = nbBufferMessage;
109 param.bufferSizeByte = bufferSizeByte;
110 param.threadAffinity = threadAffinity;
111 param.dataRate = dataRate;
112 param.linger = linger;
113 param.immediate = immediate; //EAGAIN if no peer connected
114 return param;
115}
116
119 :p_socket(NULL)
120{
121
122}
123
128
130
135bool PZmqSocket::createClientSocket(zmq::context_t & context, const PSocketParam & socketParam, const Param & extraParam){
136 p_socket = pzmq_createClientSocket(context, socketParam.hostname, socketParam.port, extraParam.type, extraParam.nbBufferMessage,
137 extraParam.bufferSizeByte, extraParam.threadAffinity, extraParam.dataRate, extraParam.immediate);
138#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
139 p_socket->set(zmq::sockopt::rcvtimeo, socketParam.recvTimeOut);
140 p_socket->set(zmq::sockopt::sndtimeo, socketParam.sendTimeOut);
141 p_socket->set(zmq::sockopt::linger, extraParam.linger); //number of ms to stop
142 p_socket->set(zmq::sockopt::immediate, extraParam.immediate); //EAGAIN if no peer connected
143#else
144 p_socket->setsockopt(ZMQ_RCVTIMEO, &socketParam.recvTimeOut, sizeof(int));
145 p_socket->setsockopt(ZMQ_SNDTIMEO, &socketParam.sendTimeOut, sizeof(int));
146 p_socket->setsockopt(ZMQ_LINGER, &extraParam.linger, sizeof(int)); //number of ms to stop
147 p_socket->setsockopt(ZMQ_IMMEDIATE, &extraParam.immediate, sizeof(int)); //EAGAIN if no peer connected
148#endif
149 return p_socket != NULL;
150}
151
153
158bool PZmqSocket::createServerSocket(zmq::context_t & context, const PSocketParam & socketParam, const Param & extraParam){
159 p_socket = pzmq_createServerSocket(context, socketParam.port, extraParam.type, extraParam.nbBufferMessage, extraParam.bufferSizeByte,
160 extraParam.threadAffinity, extraParam.dataRate);
161#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
162 p_socket->set(zmq::sockopt::rcvtimeo, socketParam.recvTimeOut);
163 p_socket->set(zmq::sockopt::sndtimeo, socketParam.sendTimeOut);
164 p_socket->set(zmq::sockopt::linger, extraParam.linger); //number of ms to stop
165 p_socket->set(zmq::sockopt::immediate, extraParam.immediate); //EAGAIN if no peer connected
166#else
167 p_socket->setsockopt(ZMQ_RCVTIMEO, &socketParam.recvTimeOut, sizeof(int));
168 p_socket->setsockopt(ZMQ_SNDTIMEO, &socketParam.sendTimeOut, sizeof(int));
169 p_socket->setsockopt(ZMQ_LINGER, &extraParam.linger, sizeof(int)); //number of ms to stop
170 p_socket->setsockopt(ZMQ_IMMEDIATE, &extraParam.immediate, sizeof(int)); //EAGAIN if no peer connected
171#endif
172 return p_socket != NULL;
173}
174
176
180PSendStatus::PSendStatus PZmqSocket::sendMsg(Message & msg, PSendFlag::PSendFlag flag){
181 try {
182 zmq::send_result_t result = p_socket->send(msg, convertToSendFlag(flag));
183 return checkSendStatus(result);
184 }catch(const zmq::error_t& e) {
185 if(e.num() == ENOTSOCK){
186 return PSendStatus::BROKEN_SOCKET;
187 }else if(e.num() == ENOTSUP || e.num() == EFSM || e.num() == ETERM || e.num() == EINVAL){
188 return PSendStatus::BROKEN_BACKEND;
189 }else if(e.num() == EINTR){
190 return PSendStatus::SIGNAL_INTERRUPTION;
191 }else if(e.num() == EHOSTUNREACH){
192 return PSendStatus::NO_ROUTE_TO_RECEIVER;
193 }else{
194 std::cerr << "Unknown ZMQ error in send: " << e.what() << std::endl;
195 return PSendStatus::BROKEN_BACKEND;
196 }
197 }
198}
199
201
205PRecvStatus::PRecvStatus PZmqSocket::recvMsg(Message & msg, PRecvFlag::PRecvFlag flag){
206 try {
207 zmq::recv_result_t result = p_socket->recv(msg, convertToRecvFlag(flag));
208 return checkRecvStatus(result);
209 }catch(const zmq::error_t& e){
210 if(e.num() == ENOTSOCK){
211 std::cerr << "ZMQ error in recv: " << e.what() << std::endl;
212 return PRecvStatus::BROKEN_SOCKET;
213 }
214 else if(e.num() == ENOTSUP || e.num() == EFSM || e.num() == ETERM || e.num() == EINVAL){
215 std::cerr << "ZMQ error in recv: " << e.what() << std::endl;
216 return PRecvStatus::BROKEN_BACKEND;
217 }
218 else if(e.num() == EINTR){
219 std::cerr << "ZMQ error in recv: " << e.what() << std::endl;
220 return PRecvStatus::SIGNAL_INTERRUPTION;
221 }
222 else{
223 std::cerr << "Unknown ZMQ error in recv: " << e.what() << std::endl;
224 return PRecvStatus::BROKEN_BACKEND;
225 }
226 }
227}
228
230
233 if(p_socket != NULL){
234 return p_socket->handle() != NULL;
235 }
236 return false;
237}
238
241 if(p_socket != NULL){
242 p_socket->close();
243 delete p_socket;
244 p_socket = NULL;
245 }
246}
247
254
256
261
263
268
270
276bool PZmqSocketGenerator::createClientSocket(PZmqSocketGenerator::Socket & socket, const PSocketParam & socketParam, const PZmqParam & param){
277 return socket.createClientSocket(p_context, socketParam, param);
278}
279
281
286bool PZmqSocketGenerator::createServerSocket(PZmqSocketGenerator::Socket & socket, const PSocketParam & socketParam, const PZmqParam & param){
287 return socket.createServerSocket(p_context, socketParam, param);
288}
289
291
294void PZmqSocketGenerator::msgToMock(DataStreamMsg & mockMsg, const PZmqSocketGenerator::Message & msg){
295 size_t dataSize(msg.size());
296 mockMsg.resize(dataSize);
297 memcpy(mockMsg.data(), (const void*)msg.data(), dataSize);
298}
299
301
305 size_t dataSize(mockMsg.size());
306 msg.rebuild(dataSize);
307 memcpy((void*)msg.data(), mockMsg.data(), dataSize);
308}
zmq::send_flags convertToSendFlag(PSendFlag::PSendFlag flag)
Convert a send flag into zmq flag.
PZmqParam pzmq_createParamServer(int type, int nbBufferMessage, int bufferSizeByte, size_t threadAffinity, ssize_t dataRate, int linger, int immediate)
Create param for a client socket.
zmq::recv_flags convertToRecvFlag(PRecvFlag::PRecvFlag flag)
Convert a recv flag into zmq flag.
PRecvStatus::PRecvStatus checkRecvStatus(zmq::recv_result_t res)
Check the recv result and convert it into PRecvStatus.
PZmqParam pzmq_createParamClient(int type, int nbBufferMessage, int bufferSizeByte, size_t threadAffinity, ssize_t dataRate, int linger, int immediate)
Create param for a client socket.
PSendStatus::PSendStatus checkSendStatus(zmq::send_result_t res)
Check the send result and convert it into PSendStatus.
PZmqParam pzmq_createParamServer(int type, int nbBufferMessage=10000, int bufferSizeByte=1000000, size_t threadAffinity=0lu, ssize_t dataRate=200000l, int linger=-1, int immediate=0)
Create param for a client socket.
PRecvStatus::PRecvStatus checkRecvStatus(zmq::recv_result_t res)
Check the recv result and convert it into PRecvStatus.
PZmqParam pzmq_createParamClient(int type, int nbBufferMessage=10000, int bufferSizeByte=1000000, size_t threadAffinity=0lu, ssize_t dataRate=200000l, int linger=-1, int immediate=0)
Create param for a client socket.
PSendStatus::PSendStatus checkSendStatus(zmq::send_result_t res)
Check the send result and convert it into PSendStatus.
zmq::context_t p_context
Context ZMQ.
Definition PZmqBackend.h:93
PZmqParam Param
Define the type of extra parameters which can be used to create a Socket used by the PAbstractSocketM...
Definition PZmqBackend.h:78
PZmqSocket Socket
Define the socket of the backend used by the PAbstractSocketManager.
Definition PZmqBackend.h:74
bool createServerSocket(Socket &socket, const PSocketParam &socketParam, const PZmqParam &param)
Create a server socket.
bool createClientSocket(Socket &socket, const PSocketParam &socketParam, const PZmqParam &param)
Create a client socket.
static void mockToMsg(Message &msg, DataStreamMsg &mockMsg)
Copy mock message data into current backend message.
zmq::message_t Message
Define the type of message used by the PAbstractSocketManager.
Definition PZmqBackend.h:76
PZmqSocketGenerator()
Default constructor of PZmqSocketGenerator setting the number of threads for zmq I/O to 1.
static void msgToMock(DataStreamMsg &mockMsg, const Message &msg)
Copy current backend message data into mock message.
static Param server()
Create a server parameter.
static Param client()
Create a client parameter.
zmq::socket_t * p_socket
ZMQ Socket.
Definition PZmqBackend.h:67
PZmqSocket()
Default constructor of PZmqSocket.
bool isConnected() const
Check if the socket is connected.
bool createClientSocket(zmq::context_t &context, const PSocketParam &socketParam, const Param &extraParam)
Create a client socket.
PSendStatus::PSendStatus sendMsg(Message &msg, PSendFlag::PSendFlag flag=PSendFlag::BLOCK)
Send data with the socket.
void close()
Close the socket.
bool createServerSocket(zmq::context_t &context, const PSocketParam &socketParam, const Param &extraParam)
Create a server socket.
PZmqParam Param
Define the type of extra parameters which can be used to create a Socket used by the PAbstractSocketM...
Definition PZmqBackend.h:47
virtual ~PZmqSocket()
Destructor of PZmqSocket.
PRecvStatus::PRecvStatus recvMsg(Message &msg, PRecvFlag::PRecvFlag flag=PRecvFlag::BLOCK)
Receive data with the socket.
zmq::message_t Message
Define the type of message used by the PAbstractSocketManager.
Definition PZmqBackend.h:45
zmq::socket_t * pzmq_createClientSocket(zmq::context_t &context, int type, const std::string &address, size_t port, int immediate)
Create a client socket to be used by the SocketManagerZMQ.
zmq::socket_t * pzmq_createServerSocket(zmq::context_t &context, int type, size_t port)
Create a server socket to be used by the SocketManagerZMQ.
Set of parameters to be passed to create a socket with zmq backend.
Definition PZmqBackend.h:17
ssize_t dataRate
Data rate.
Definition PZmqBackend.h:27
int nbBufferMessage
Number of messages in the buffer.
Definition PZmqBackend.h:21
int bufferSizeByte
Size of the message buffer in bytes.
Definition PZmqBackend.h:23
int type
Socket type.
Definition PZmqBackend.h:19
int linger
linger period for socket shutdown
Definition PZmqBackend.h:29
int immediate
Immediate flag for socket connexion (EAGAIN if no peer connected)
Definition PZmqBackend.h:31
size_t threadAffinity
Mask of threads which deal with reconnection.
Definition PZmqBackend.h:25