worker.cc 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160
  1. extern "C" {
  2. #include <event2/bufferevent.h>
  3. #include <event2/buffer.h>
  4. #include <event2/event.h>
  5. }
  6. #include <vector>
  7. #include <stdio.h>
  8. #include "evutils.h"
  9. #include "subproblem.h"
  10. typedef enum {
  11. WRKCCSTATE_AWAITCMD,
  12. WRKCCSTATE_RDPROBLEM,
  13. WRKCCSTATE_RDDPNODES,
  14. WRKCCSTATE_END
  15. } WrkCCState;
  16. struct WrkControllerConnInfo {
  17. WrkCCState state;
  18. unsigned char subproblem[SUBPROBLEM_DESC_LEN];
  19. unsigned short num_dpnodes;
  20. WrkControllerConnInfo() : state(WRKCCSTATE_AWAITCMD) {}
  21. };
  22. static struct WrkControllerState {
  23. Subproblem *current_problem;
  24. unsigned short num_expected_dpnodes;
  25. vector<struct bufferevent *> dpnodes;
  26. unsigned short num_connected_dpnodes;
  27. WrkControllerState(): current_problem(NULL) {}
  28. } wrkctrlstate;
  29. static void dpconn_event_cb(struct bufferevent *bev, short events,
  30. void *ctx)
  31. {
  32. if (events & BEV_EVENT_CONNECTED) {
  33. // We have successfully connected to the dpnode
  34. ++wrkctrlstate.num_connected_dpnodes;
  35. if (wrkctrlstate.num_connected_dpnodes ==
  36. wrkctrlstate.num_expected_dpnodes) {
  37. cerr << "Starting work\n";
  38. // start_working();
  39. }
  40. } else if (events & (BEV_EVENT_EOF|BEV_EVENT_ERROR)) {
  41. fprintf(stderr, "Closing connection to controller and restarting\n");
  42. bufferevent_free(bev);
  43. cerr << "Stopping work\n";
  44. // stop_working();
  45. }
  46. }
  47. static void controllerconn_reader(struct bufferevent *bev, void *ctx)
  48. {
  49. struct evbuffer *input = bufferevent_get_input(bev);
  50. WrkControllerConnInfo *info = (WrkControllerConnInfo *)ctx;
  51. unsigned char cmd[1];
  52. while(1) {
  53. size_t len = evbuffer_get_length(input);
  54. switch(info->state) {
  55. case WRKCCSTATE_AWAITCMD:
  56. if (len < 1) return;
  57. bufferevent_read(bev, cmd, 1);
  58. switch(cmd[0]) {
  59. case 'P':
  60. info->state = WRKCCSTATE_RDPROBLEM;
  61. break;
  62. case 'S':
  63. // stop_working();
  64. break;
  65. default:
  66. /* Unknown command received */
  67. fprintf(stderr, "Unknown command in "
  68. "controllerconn_reader: %c\n", cmd[0]);
  69. info->state = WRKCCSTATE_END;
  70. break;
  71. }
  72. break;
  73. case WRKCCSTATE_RDPROBLEM:
  74. if (len < SUBPROBLEM_DESC_LEN+2) return;
  75. // stop_working();
  76. bufferevent_read(bev, info->subproblem, SUBPROBLEM_DESC_LEN);
  77. bufferevent_read(bev, &(info->num_dpnodes), 2);
  78. info->state = WRKCCSTATE_RDDPNODES;
  79. /* FALLTHROUGH */
  80. case WRKCCSTATE_RDDPNODES:
  81. if (len < 6*(info->num_dpnodes)) return;
  82. wrkctrlstate.num_expected_dpnodes = info->num_dpnodes;
  83. {
  84. unsigned short i;
  85. for(i=0;i<info->num_dpnodes;++i) {
  86. unsigned char ipport[6];
  87. bufferevent_read(bev, ipport, 6);
  88. // XXX: Start a connection to this DPnode
  89. cerr << "Connecting to DPnode " <<
  90. int(ipport[0]) << "." <<
  91. int(ipport[1]) << "." <<
  92. int(ipport[2]) << "." <<
  93. int(ipport[3]) << ":" <<
  94. ((ipport[4] << 8) + ipport[5]) <<
  95. "\n";
  96. struct bufferevent *dpbev = client_create(
  97. bufferevent_get_base(bev), ipport,
  98. dpconn_event_cb);
  99. if (dpbev) {
  100. wrkctrlstate.dpnodes.push_back(dpbev);
  101. } else {
  102. // stop_working();
  103. }
  104. }
  105. }
  106. info->state = WRKCCSTATE_AWAITCMD;
  107. break;
  108. case WRKCCSTATE_END:
  109. // Shut down
  110. delete info;
  111. event_base_loopbreak(bufferevent_get_base(bev));
  112. bufferevent_free(bev);
  113. return;
  114. }
  115. }
  116. }
  117. static void controllerconn_event_cb(struct bufferevent *bev, short events,
  118. void *ctx)
  119. {
  120. if (events & BEV_EVENT_CONNECTED) {
  121. // We have successfully connected to the controller
  122. char id[1] = { 'W' };
  123. bufferevent_enable(bev, EV_READ|EV_WRITE);
  124. bufferevent_write(bev, id, 1);
  125. bufferevent_setcb(bev, controllerconn_reader, NULL,
  126. controllerconn_event_cb, new WrkControllerConnInfo());
  127. } else if (events & (BEV_EVENT_EOF|BEV_EVENT_ERROR)) {
  128. fprintf(stderr, "Closing connection to controller and exiting\n");
  129. event_base_loopbreak(bufferevent_get_base(bev));
  130. bufferevent_free(bev);
  131. }
  132. }
  133. int main(int argc, char **argv)
  134. {
  135. if (argc != 3) {
  136. fprintf(stderr, "Usage: %s controller_host controller_port\n", argv[0]);
  137. return 1;
  138. }
  139. unsigned short controller_port = strtoul(argv[2], NULL, 10);
  140. return controller_client(argv[1], controller_port,
  141. controllerconn_event_cb);
  142. }