Просмотр исходного кода

Enable libevent thread safety for the worker

Ian Goldberg 14 лет назад
Родитель
Сommit
4a0de9516a
6 измененных файлов с 26 добавлено и 17 удалено
  1. 1 1
      Makefile
  2. 1 1
      controller.cc
  3. 2 2
      dpnode.cc
  4. 10 7
      evutils.cc
  5. 5 3
      evutils.h
  6. 7 3
      worker.cc

+ 1 - 1
Makefile

@@ -58,7 +58,7 @@ dpnode: dpnode.o evutils.o
 	g++ -g -Wall $^ -o $@ -L$(LIBEVENT)/lib -levent -lntl -L$(GMP) -lgmp
 
 worker: worker.o evutils.o cudadl.o
-	g++ -g -Wall $^ -o $@ -L$(LIBEVENT)/lib -levent -lntl -L$(GMP) -lgmp -lpthread -L$(CUDA)/lib64 -lcudart
+	g++ -g -Wall $^ -o $@ -L$(LIBEVENT)/lib -levent -levent_pthreads -lntl -L$(GMP) -lgmp -lpthread -L$(CUDA)/lib64 -lcudart
 
 clean:
 	-rm -f $(OFILES)

+ 1 - 1
controller.cc

@@ -501,7 +501,7 @@ void *controller_create(struct event_base *evbase, unsigned short bindport,
     unsigned int *ip, unsigned short *boundport)
 {
     return listener_create(evbase, bindport, controller_accept_cb, NULL,
-	ip, boundport);
+	ip, boundport, false);
 }
 
 static unsigned short curproblemid = 0;

+ 2 - 2
dpnode.cc

@@ -64,7 +64,7 @@ struct evconnlistener *dpnode_create(struct event_base *evbase,
     unsigned int *ip, unsigned short *boundport)
 {
     return listener_create(evbase, 0, dpnode_accept_cb, NULL,
-	ip, boundport);
+	ip, boundport, false);
 }
 
 typedef enum {
@@ -201,5 +201,5 @@ int main(int argc, char **argv)
     unsigned short controller_port = strtoul(argv[2], NULL, 10);
 
     return controller_client(argv[1], controller_port,
-				controllerconn_event_cb);
+				controllerconn_event_cb, false);
 }

+ 10 - 7
evutils.cc

@@ -36,7 +36,7 @@ static unsigned int getlocalIP(void)
 // to the IP and port of the socket, in network byte order.
 struct evconnlistener *listener_create(struct event_base *evbase,
     unsigned short bindport, evconnlistener_cb cb, void *ctx,
-    unsigned int *ip, unsigned short *boundport)
+    unsigned int *ip, unsigned short *boundport, bool multithread)
 {
     struct sockaddr_in sin;
     sin.sin_family = AF_INET;
@@ -52,8 +52,9 @@ struct evconnlistener *listener_create(struct event_base *evbase,
     }
 
     struct evconnlistener *listener = evconnlistener_new_bind(evbase,
-	cb, ctx, LEV_OPT_CLOSE_ON_FREE|LEV_OPT_REUSEABLE, -1,
-	(struct sockaddr*)&sin, sizeof(sin));
+	cb, ctx, LEV_OPT_CLOSE_ON_FREE|LEV_OPT_REUSEABLE|
+	         (multithread ? LEV_OPT_THREADSAFE : 0),
+        -1, (struct sockaddr*)&sin, sizeof(sin));
     if (!listener) {
 	perror("Unable to create listener");
 	return NULL;
@@ -74,7 +75,8 @@ struct evconnlistener *listener_create(struct event_base *evbase,
 // Create a client connection to the given 6-byte ipport (4 byte IP, 2
 // byte port in network byte order).
 struct bufferevent *client_create(struct event_base *evbase,
-    unsigned char ipport[6], bufferevent_event_cb event_handler)
+    unsigned char ipport[6], bufferevent_event_cb event_handler,
+    bool multithread)
 {
     // Create a Controller client socket
     struct sockaddr_in client_sin;
@@ -83,7 +85,7 @@ struct bufferevent *client_create(struct event_base *evbase,
     memmove(&client_sin.sin_port, ipport+4, 2);
 
     struct bufferevent *client_bev = bufferevent_socket_new(evbase,
-	-1, BEV_OPT_CLOSE_ON_FREE);
+	-1, BEV_OPT_CLOSE_ON_FREE|(multithread ? BEV_OPT_THREADSAFE : 0));
 
     if (bufferevent_socket_connect(client_bev,
 	    (struct sockaddr *)&client_sin, sizeof(client_sin)) < 0) {
@@ -102,7 +104,8 @@ struct bufferevent *client_create(struct event_base *evbase,
 // fails.  This function calls the libevent main loop, so it will only return
 // when the program is finished.
 int controller_client(const char *controller_host,
-    unsigned short controller_port, bufferevent_event_cb event_handler)
+    unsigned short controller_port, bufferevent_event_cb event_handler,
+    bool multithread)
 {
     struct event_base *evbase = event_base_new();
 
@@ -117,7 +120,7 @@ int controller_client(const char *controller_host,
     controller_sin.sin_port = htons(controller_port);
 
     struct bufferevent *controller_bev = bufferevent_socket_new(evbase,
-	-1, BEV_OPT_CLOSE_ON_FREE);
+	-1, BEV_OPT_CLOSE_ON_FREE|(multithread ? BEV_OPT_THREADSAFE : 0));
 
     if (bufferevent_socket_connect(controller_bev,
 	    (struct sockaddr *)&controller_sin, sizeof(controller_sin)) < 0) {

+ 5 - 3
evutils.h

@@ -14,18 +14,20 @@ unsigned int hostlookup(const char *hostname);
 // to the IP and port of the socket, in network byte order.
 struct evconnlistener *listener_create(struct event_base *evbase,
     unsigned short bindport, evconnlistener_cb cb, void *ctx,
-    unsigned int *ip, unsigned short *boundport);
+    unsigned int *ip, unsigned short *boundport, bool multithread);
 
 // Create a client connection to the given 6-byte ipport (4 byte IP, 2
 // byte port in network byte order).
 struct bufferevent *client_create(struct event_base *evbase,
-    unsigned char ipport[6], bufferevent_event_cb event_handler);
+    unsigned char ipport[6], bufferevent_event_cb event_handler,
+    bool multithread);
 
 // Create a client connection to a controller, with event_handler set as
 // the event callback.  It will be called when the connection succeeds or
 // fails.  This function calls the libevent main loop, so it will only return
 // when the program is finished.
 int controller_client(const char *controller_host,
-    unsigned short controller_port, bufferevent_event_cb event_handler);
+    unsigned short controller_port, bufferevent_event_cb event_handler,
+    bool multithread);
 
 #endif

+ 7 - 3
worker.cc

@@ -1,4 +1,7 @@
+#include <pthread.h>
+
 extern "C" {
+#include <event2/thread.h>
 #include <event2/bufferevent.h>
 #include <event2/buffer.h>
 #include <event2/event.h>
@@ -7,7 +10,6 @@ extern "C" {
 #include <NTL/ZZ_p.h>
 
 #include <vector>
-#include <pthread.h>
 #include <stdio.h>
 
 #include "cudadl.h"
@@ -200,7 +202,7 @@ static void controllerconn_reader(struct bufferevent *bev, void *ctx)
 				    "\n";
 			struct bufferevent *dpbev = client_create(
 				bufferevent_get_base(bev), ipport,
-				dpconn_event_cb);
+				dpconn_event_cb, true);
 			if (dpbev) {
 			    wrkctrlstate.dpnodes.push_back(dpbev);
 			} else {
@@ -248,6 +250,8 @@ int main(int argc, char **argv)
 
     unsigned short controller_port = strtoul(argv[2], NULL, 10);
 
+    evthread_use_pthreads();
+
     return controller_client(argv[1], controller_port,
-				controllerconn_event_cb);
+				controllerconn_event_cb, true);
 }