22 return PSendStatus::OK;
24 int err = zmq_errno();
26 return PSendStatus::SOCKET_NOT_AVAILABLE;
28 std::cerr <<
"Unknown ZMQ error in send: '" << err <<
"'"<< std::endl;
29 return PSendStatus::BROKEN_BACKEND;
45 return PRecvStatus::OK;
47 int err = zmq_errno();
49 return PRecvStatus::NO_MESSAGE_RECEIVED;
51 return PRecvStatus::BROKEN_SOCKET;
61 if(flag == PSendFlag::NON_BLOCK){
return zmq::send_flags::dontwait;}
62 else{
return zmq::send_flags::none;}
70 if(flag == PRecvFlag::NON_BLOCK){
return zmq::recv_flags::dontwait;}
71 else{
return zmq::recv_flags::none;}
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);
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));
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);
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));
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;
194 std::cerr <<
"Unknown ZMQ error in send: " << e.what() << std::endl;
195 return PSendStatus::BROKEN_BACKEND;
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;
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;
218 else if(e.num() == EINTR){
219 std::cerr <<
"ZMQ error in recv: " << e.what() << std::endl;
220 return PRecvStatus::SIGNAL_INTERRUPTION;
223 std::cerr <<
"Unknown ZMQ error in recv: " << e.what() << std::endl;
224 return PRecvStatus::BROKEN_BACKEND;
295 size_t dataSize(msg.size());
296 mockMsg.resize(dataSize);
297 memcpy(mockMsg.data(), (
const void*)msg.data(), dataSize);
305 size_t dataSize(mockMsg.size());
306 msg.rebuild(dataSize);
307 memcpy((
void*)msg.data(), mockMsg.data(), dataSize);
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.
PZmqParam Param
Define the type of extra parameters which can be used to create a Socket used by the PAbstractSocketM...
PZmqSocket Socket
Define the socket of the backend used by the PAbstractSocketManager.
bool createServerSocket(Socket &socket, const PSocketParam &socketParam, const PZmqParam ¶m)
Create a server socket.
bool createClientSocket(Socket &socket, const PSocketParam &socketParam, const PZmqParam ¶m)
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.
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.
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...
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.
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.
ssize_t dataRate
Data rate.
int nbBufferMessage
Number of messages in the buffer.
int bufferSizeByte
Size of the message buffer in bytes.
int linger
linger period for socket shutdown
int immediate
Immediate flag for socket connexion (EAGAIN if no peer connected)
size_t threadAffinity
Mask of threads which deal with reconnection.