-
Notifications
You must be signed in to change notification settings - Fork 0
/
pairs.cpp
executable file
·171 lines (131 loc) · 3.83 KB
/
pairs.cpp
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
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
// Copyright 2017 by Bill Torpey. All Rights Reserved.
// This work is licensed under a Creative Commons Attribution-NonCommercial-NoDerivs 3.0 United States License.
// http://creativecommons.org/licenses/by-nc-nd/3.0/us/deed.en
#include <zmq.h>
#include "common.h"
// tls for control socket
pthread_key_t key;
pthread_mutex_t sendMutex = PTHREAD_MUTEX_INITIALIZER;
// send a command to the control socket
void sendCommand(void* context, zmqControlMsg* msg, int msgSize)
{
int rc;
void* controlPub = pthread_getspecific(key);
if (controlPub == NULL) {
controlPub = createSocket(theContext, ZMQ_PAIR);
checkVoid(controlPub);
rc = pthread_setspecific(key, controlPub);
checkInt(rc);
}
if (mutexFlag) {
checkInt(pthread_mutex_lock(&sendMutex));
}
rc = zmq_connect(controlPub, CONTROL_ENDPOINT);
checkInt(rc);
if (pollFlag) {
zmq_pollitem_t items[] = { { controlPub, 0, ZMQ_POLLIN, 0} };
rc = zmq_poll(items, 1, 1);
checkInt(rc);
}
int i = zmq_send(controlPub, msg, msgSize, dontWaitFlag ? ZMQ_DONTWAIT : 0);
checkInt(i);
assert(i == msgSize);
rc = zmq_disconnect(controlPub, CONTROL_ENDPOINT);
checkInt(rc);
if (mutexFlag) {
checkInt(pthread_mutex_unlock(&sendMutex));
}
}
void* mainLoop(void*)
{
int rc;
// control socket
void* controlSub = createSocket(theContext, ZMQ_PAIR);
rc = zmq_bind(controlSub, CONTROL_ENDPOINT);
checkInt(rc);
// data socket
void* dataSub = createSocket(theContext, ZMQ_SUB);
rc = zmq_bind(dataSub, DATA_ENDPOINT);
checkInt(rc);
zmq_msg_t zmsg;
zmq_msg_init(&zmsg);
while (keepGoing) {
zmq_pollitem_t items[] = {
{ controlSub, 0, ZMQ_POLLIN, 0},
{ dataSub, 0, ZMQ_POLLIN, 0}
};
int rc = zmq_poll(items, 2, -1);
if (rc < 0) {
break;
}
// got command msg?
if (items[0].revents & ZMQ_POLLIN) {
// controlSub
int size = zmq_msg_recv(&zmsg, controlSub, 0);
if (size != -1) {
processControlMsg(&zmsg);
}
continue;
}
// got normal msg?
if (items[1].revents & ZMQ_POLLIN) {
int size = zmq_msg_recv(&zmsg, dataSub, 0);
if (size != -1) {
processDataMsg(&zmsg);
}
continue;
}
}
closeSocket(controlSub);
closeSocket(dataSub);
return NULL;
}
int main(int argc, char** argv)
{
int rc;
printVersion();
fprintf(stderr, "Using pair sockets\n");
parseParams(argc, argv);
// setup signal handler
signal(SIGINT, &onSignal);
// initialize zmq
theContext = zmq_ctx_new();
checkVoid(theContext);
rc = zmq_ctx_set(theContext, ZMQ_BLOCKY, false);
checkInt(rc);
// create the key for socket tls
pthread_key_create(&key, cleanupSocket);
// start main loop
pthread_t mainThread;
pthread_create(&mainThread, NULL, mainLoop, NULL);
// start command loops
pthread_t* commandThread = new pthread_t[numThreads];
for (int i = 0; i < numThreads; ++i) {
maxControlMsgSent[&commandThread[i]] = 0;
maxControlMsgReceived[&commandThread[i]] = 0;
}
for (int i = 0; i < numThreads; ++i) {
pthread_create(&commandThread[i], NULL, commandLoop, &commandThread[i]);
}
// wait for command threads to finish
for (int i = 0; i < numThreads; ++i) {
pthread_join(commandThread[i], NULL);
}
// sleep a bit
fprintf(stderr, "Waiting...");
usleep(sleepDuration*1000000);
// send shutdown command
pthread_t shutdownThread;
pthread_create(&shutdownThread, NULL, shutdownFunc, NULL);
pthread_join(shutdownThread, NULL);
// wait for main thread to finish
pthread_join(mainThread, NULL);
rc = zmq_ctx_destroy(theContext);
checkInt(rc);
checkResults();
printResults();
if (failed)
return 1;
else
return 0;
}