dpnode.cc 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216
  1. extern "C" {
  2. #include <event2/listener.h>
  3. #include <event2/bufferevent.h>
  4. #include <event2/buffer.h>
  5. }
  6. #include <sys/socket.h>
  7. #include <netinet/in.h>
  8. #include <arpa/inet.h>
  9. #include <set>
  10. #include <stdlib.h>
  11. #include <string.h>
  12. #include "evutils.h"
  13. #include "subproblem.h"
  14. typedef enum {
  15. DPSTATE_START,
  16. DPSTATE_END
  17. } DPState;
  18. struct DPNodeConnInfo {
  19. DPState state;
  20. DPNodeConnInfo() : state(DPSTATE_START) {}
  21. };
  22. static void dpnode_event_cb(struct bufferevent *bev, short events,
  23. void *ctx)
  24. {
  25. if (events & (BEV_EVENT_EOF|BEV_EVENT_ERROR)) {
  26. fprintf(stderr, "Closing connection\n");
  27. DPNodeConnInfo *info = (DPNodeConnInfo*)ctx;
  28. delete info;
  29. bufferevent_free(bev);
  30. }
  31. }
  32. static void dpnode_reader(struct bufferevent *bev, void *ctx)
  33. {
  34. struct evbuffer *input = bufferevent_get_input(bev);
  35. unsigned int dp[WORDS+6];
  36. int num_read = 0;
  37. while(1) {
  38. size_t len = evbuffer_get_length(input);
  39. if (len < ((WORDS+6)*sizeof(unsigned int))) break;
  40. bufferevent_read(bev, dp, (WORDS+6)*sizeof(unsigned int));
  41. ++num_read;
  42. }
  43. // cerr << num_read << " DPs read\n";
  44. }
  45. static void dpnode_accept_cb(struct evconnlistener *listener,
  46. evutil_socket_t fd, struct sockaddr *address, int socklen,
  47. void *ctx)
  48. {
  49. DPNodeConnInfo *info = new DPNodeConnInfo();
  50. // Create a bufferevent for the new connection
  51. struct event_base *base = evconnlistener_get_base(listener);
  52. struct bufferevent *bev = bufferevent_socket_new(
  53. base, fd, BEV_OPT_CLOSE_ON_FREE);
  54. bufferevent_setcb(bev, dpnode_reader, NULL,
  55. dpnode_event_cb, info);
  56. bufferevent_enable(bev, EV_READ|EV_WRITE);
  57. }
  58. // Create a new DPnode socket. ip and boundport are set to the IP and
  59. // port of the socket, in network byte order.
  60. struct evconnlistener *dpnode_create(struct event_base *evbase,
  61. unsigned int *ip, unsigned short *boundport)
  62. {
  63. return listener_create(evbase, 0, dpnode_accept_cb, NULL,
  64. ip, boundport, false);
  65. }
  66. typedef enum {
  67. DPCCSTATE_AWAITCMD,
  68. DPCCSTATE_RDPROBLEM,
  69. DPCCSTATE_END
  70. } DPCCState;
  71. struct DPControllerConnInfo {
  72. DPCCState state;
  73. DPControllerConnInfo() : state(DPCCSTATE_AWAITCMD) {}
  74. };
  75. static struct DPControllerState {
  76. Subproblem *current_problem;
  77. struct evconnlistener *listener;
  78. std::set<struct bufferevent *> workers;
  79. DPControllerState() : current_problem(NULL), listener(NULL) {}
  80. } dpctrlstate;
  81. static void stop_problem(void)
  82. {
  83. if (dpctrlstate.current_problem) {
  84. dpctrlstate.current_problem = NULL;
  85. }
  86. if (dpctrlstate.listener) {
  87. evconnlistener_free(dpctrlstate.listener);
  88. dpctrlstate.listener = NULL;
  89. }
  90. std::set<struct bufferevent *>::iterator wit;
  91. for (wit = dpctrlstate.workers.begin(); wit != dpctrlstate.workers.end();
  92. ++wit) {
  93. bufferevent_free(*wit);
  94. }
  95. dpctrlstate.workers.clear();
  96. }
  97. static void start_problem(struct bufferevent *bev,
  98. const unsigned char *subproblem)
  99. {
  100. unsigned int myip;
  101. unsigned short myport;
  102. stop_problem();
  103. dpctrlstate.current_problem = new Subproblem(subproblem);
  104. dpctrlstate.current_problem->dump(cerr);
  105. // Create the DPNode server socket
  106. dpctrlstate.listener = dpnode_create(bufferevent_get_base(bev),
  107. &myip, &myport);
  108. struct in_addr myaddr = { myip };
  109. fprintf(stderr, "Bound to %s:%d\n", inet_ntoa(myaddr), ntohs(myport));
  110. unsigned char idstring[7];
  111. idstring[0] = 'L';
  112. memmove(idstring+1, &myip, 4);
  113. memmove(idstring+5, &myport, 2);
  114. bufferevent_write(bev, idstring, 7);
  115. }
  116. static void controllerconn_reader(struct bufferevent *bev, void *ctx)
  117. {
  118. struct evbuffer *input = bufferevent_get_input(bev);
  119. DPControllerConnInfo *info = (DPControllerConnInfo *)ctx;
  120. unsigned char cmd[1];
  121. unsigned char subproblem[SUBPROBLEM_DESC_LEN];
  122. while(1) {
  123. size_t len = evbuffer_get_length(input);
  124. switch(info->state) {
  125. case DPCCSTATE_AWAITCMD:
  126. if (len < 1) return;
  127. bufferevent_read(bev, cmd, 1);
  128. switch(cmd[0]) {
  129. case 'P':
  130. info->state = DPCCSTATE_RDPROBLEM;
  131. break;
  132. case 'S':
  133. stop_problem();
  134. break;
  135. default:
  136. /* Unknown command received */
  137. fprintf(stderr, "Unknown command in "
  138. "controllerconn_reader: %c\n", cmd[0]);
  139. info->state = DPCCSTATE_END;
  140. break;
  141. }
  142. break;
  143. case DPCCSTATE_RDPROBLEM:
  144. if (len < SUBPROBLEM_DESC_LEN) return;
  145. bufferevent_read(bev, subproblem, SUBPROBLEM_DESC_LEN);
  146. start_problem(bev, subproblem);
  147. info->state = DPCCSTATE_AWAITCMD;
  148. break;
  149. case DPCCSTATE_END:
  150. // Shut down
  151. delete info;
  152. event_base_loopbreak(bufferevent_get_base(bev));
  153. bufferevent_free(bev);
  154. return;
  155. }
  156. }
  157. }
  158. static void controllerconn_event_cb(struct bufferevent *bev, short events,
  159. void *ctx)
  160. {
  161. if (events & BEV_EVENT_CONNECTED) {
  162. // We have successfully connected to the controller
  163. char id[1] = { 'D' };
  164. bufferevent_enable(bev, EV_READ|EV_WRITE);
  165. bufferevent_write(bev, id, 1);
  166. bufferevent_setcb(bev, controllerconn_reader, NULL,
  167. controllerconn_event_cb, new DPControllerConnInfo());
  168. } else if (events & (BEV_EVENT_EOF|BEV_EVENT_ERROR)) {
  169. fprintf(stderr, "Closing connection to controller and exiting\n");
  170. event_base_loopbreak(bufferevent_get_base(bev));
  171. bufferevent_free(bev);
  172. }
  173. }
  174. int main(int argc, char **argv)
  175. {
  176. if (argc != 3) {
  177. fprintf(stderr, "Usage: %s controller_host controller_port\n", argv[0]);
  178. return 1;
  179. }
  180. unsigned short controller_port = strtoul(argv[2], NULL, 10);
  181. return controller_client(argv[1], controller_port,
  182. controllerconn_event_cb, false);
  183. }