PhoenixZMQ
8.2.0
Library which integrates zeromq use
Toggle main menu visibility
Loading...
Searching...
No Matches
phoenix_zmq.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 <sstream>
8
#include "
phoenix_zmq.h
"
9
10
12
18
zmq::socket_t*
pzmq_createClientSocket
(zmq::context_t & context,
int
type,
const
std::string & address,
size_t
port,
int
immediate){
19
std::stringstream socketAddressData;
20
socketAddressData <<
"tcp://"
<< address <<
":"
<< port;
21
zmq::socket_t *socket =
new
zmq::socket_t(context, type);
22
if
(socket == NULL){
return
NULL;}
23
24
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
25
socket->set(zmq::sockopt::immediate, immediate);
26
socket->set(zmq::sockopt::linger, -1);
//1 ms to stop
27
#else
28
int
immediateVal(immediate);
29
socket->setsockopt(ZMQ_IMMEDIATE, &immediateVal,
sizeof
(
int
));
30
socket->setsockopt(ZMQ_LINGER, -1);
//1 ms to stop
31
#endif
32
socket->connect(socketAddressData.str());
// ← connect APRÈS les options
33
34
if
(type == ZMQ_SUB){
35
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
36
socket->set(zmq::sockopt::subscribe,
""
);
37
socket->set(zmq::sockopt::conflate, 1);
38
#else
39
socket->setsockopt(ZMQ_SUBSCRIBE,
""
, 0);
40
int
conflate(1);
41
socket->setsockopt(ZMQ_CONFLATE, &conflate,
sizeof
(
int
));
42
#endif
43
}
44
return
socket;
45
}
46
47
52
zmq::socket_t*
pzmq_createServerSocket
(zmq::context_t & context,
int
type,
size_t
port){
53
std::stringstream socketAddressData;
54
socketAddressData <<
"tcp://127.0.0.1:"
<< port;
55
zmq::socket_t *socket =
new
zmq::socket_t(context, type);
56
if
(socket != NULL){
57
socket->bind(socketAddressData.str());
58
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
59
socket->set(zmq::sockopt::linger, -1);
//1 ms to stop
60
#else
61
socket->setsockopt(ZMQ_LINGER, -1);
//1 ms to stop
62
#endif
63
}
64
return
socket;
65
}
66
68
77
zmq::socket_t*
pzmq_createServerSocket
(zmq::context_t & context,
size_t
port,
int
type,
int
nbBufferMessage,
int
bufferSizeByte,
78
size_t
threadAffinity, ssize_t dataRate)
79
{
80
zmq::socket_t* socket =
pzmq_createServerSocket
(context, type, port);
81
bool
b(socket != NULL);
82
if
(b){
83
pzmq_setBufferSize
(socket, type, nbBufferMessage, dataRate, bufferSizeByte);
84
pzmq_setThreadAffinity
(socket, threadAffinity);
85
}
86
return
socket;
87
}
88
89
91
101
zmq::socket_t*
pzmq_createClientSocket
(zmq::context_t & context,
const
std::string & address,
size_t
port,
int
type,
int
nbBufferMessage,
102
int
bufferSizeByte,
size_t
threadAffinity, ssize_t dataRate,
int
immediate)
103
{
104
zmq::socket_t* socket =
pzmq_createClientSocket
(context, type, address, port, immediate);
105
bool
b(socket != NULL);
106
if
(b){
107
pzmq_setBufferSize
(socket, type, nbBufferMessage, dataRate, bufferSizeByte);
108
pzmq_setThreadAffinity
(socket, threadAffinity);
109
}
110
return
socket;
111
}
112
113
115
117
void
pzmq_closeServerSocket
(zmq::socket_t *& socket){
118
if
(socket != NULL){
119
socket->close();
120
delete
socket;
121
socket = NULL;
122
}
123
}
124
126
129
void
pzmq_setNbMessageBuffer
(zmq::socket_t* socket,
int
nbBufferMessage){
130
if
(socket != NULL){
131
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
132
socket->set(zmq::sockopt::rcvhwm, nbBufferMessage);
133
#else
134
//See doc at http://api.zeromq.org/3-1:zmq-setsockopt
135
socket->setsockopt(ZMQ_RCVHWM, &nbBufferMessage,
sizeof
(
int
));
136
socket->setsockopt(ZMQ_SNDHWM, &nbBufferMessage,
sizeof
(
int
));
137
#endif
138
}
139
}
140
142
146
void
pzmq_setDataRate
(zmq::socket_t* socket,
int
type,
int
dataRate){
147
if
(socket != NULL && (type == ZMQ_PUB || type == ZMQ_SUB)){
148
int
dataRateKbit(dataRate*8l);
149
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
150
socket->set(zmq::sockopt::rate, dataRateKbit);
151
#else
152
//See doc at http://api.zeromq.org/3-1:zmq-setsockopt
153
socket->setsockopt(ZMQ_RATE, dataRateKbit);
154
#endif
155
}
156
}
157
159
162
void
pzmq_setRecvBufferSize
(zmq::socket_t* socket,
int
bufferSizeByte){
163
if
(socket != NULL){
164
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
165
socket->set(zmq::sockopt::rcvbuf, bufferSizeByte);
166
#else
167
//See doc at http://api.zeromq.org/3-1:zmq-setsockopt
168
socket->setsockopt(ZMQ_RCVBUF, bufferSizeByte);
169
#endif
170
}
171
}
172
174
177
void
pzmq_setSendBufferSize
(zmq::socket_t* socket,
int
bufferSizeByte){
178
if
(socket != NULL){
179
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
180
socket->set(zmq::sockopt::sndbuf, bufferSizeByte);
181
#else
182
//See doc at http://api.zeromq.org/3-1:zmq-setsockopt
183
socket->setsockopt(ZMQ_SNDBUF, bufferSizeByte);
184
#endif
185
}
186
}
187
189
192
void
pzmq_setThreadAffinity
(zmq::socket_t* socket,
size_t
threadAffinity){
193
if
(socket != NULL){
194
#if (CPPZMQ_VERSION_MAJOR*100 + CPPZMQ_VERSION_MINOR*10 + CPPZMQ_VERSION_PATCH) >= 471
195
socket->set(zmq::sockopt::affinity, threadAffinity);
196
#else
197
//See doc at http://api.zeromq.org/3-1:zmq-setsockopt
198
socket->setsockopt(ZMQ_AFFINITY, threadAffinity);
199
#endif
200
}
201
}
202
204
210
void
pzmq_setBufferSize
(zmq::socket_t* socket,
int
type,
int
nbBufferMessage,
int
dataRate,
size_t
bufferSizeByte){
211
pzmq_setNbMessageBuffer
(socket, nbBufferMessage);
212
pzmq_setDataRate
(socket, type, dataRate);
213
if
(type == ZMQ_PULL || type == ZMQ_SUB){
214
pzmq_setRecvBufferSize
(socket, bufferSizeByte);
215
}
else
if
(type == ZMQ_PUSH || type == ZMQ_PUB){
216
pzmq_setSendBufferSize
(socket, bufferSizeByte);
217
}
218
}
219
pzmq_setSendBufferSize
void pzmq_setSendBufferSize(zmq::socket_t *socket, int bufferSizeByte)
Set the size of the buffer to send messages.
Definition
phoenix_zmq.cpp:177
pzmq_createClientSocket
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.
Definition
phoenix_zmq.cpp:18
pzmq_closeServerSocket
void pzmq_closeServerSocket(zmq::socket_t *&socket)
Close the given server socket.
Definition
phoenix_zmq.cpp:117
pzmq_createServerSocket
zmq::socket_t * pzmq_createServerSocket(zmq::context_t &context, int type, size_t port)
Create a server socket to be used by the SocketManagerZMQ.
Definition
phoenix_zmq.cpp:52
pzmq_setThreadAffinity
void pzmq_setThreadAffinity(zmq::socket_t *socket, size_t threadAffinity)
Set the thread affinity of zmq.
Definition
phoenix_zmq.cpp:192
pzmq_setRecvBufferSize
void pzmq_setRecvBufferSize(zmq::socket_t *socket, int bufferSizeByte)
Set the size of the buffer to received messages.
Definition
phoenix_zmq.cpp:162
pzmq_setBufferSize
void pzmq_setBufferSize(zmq::socket_t *socket, int type, int nbBufferMessage, int dataRate, size_t bufferSizeByte)
Set the size of the buffer to send messages.
Definition
phoenix_zmq.cpp:210
pzmq_setDataRate
void pzmq_setDataRate(zmq::socket_t *socket, int type, int dataRate)
Set the data rate of the socket.
Definition
phoenix_zmq.cpp:146
pzmq_setNbMessageBuffer
void pzmq_setNbMessageBuffer(zmq::socket_t *socket, int nbBufferMessage)
Set the number of messages in the messages buffer.
Definition
phoenix_zmq.cpp:129
phoenix_zmq.h
src
phoenix_zmq.cpp
Generated by
1.17.0