universal.cc 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493
  1. /* Gearman server and library
  2. * Copyright (C) 2008 Brian Aker, Eric Day
  3. * All rights reserved.
  4. *
  5. * Use and distribution licensed under the BSD license. See
  6. * the COPYING file in the parent directory for full text.
  7. */
  8. /**
  9. * @file
  10. * @brief Gearman State Definitions
  11. */
  12. #include <libgearman/common.h>
  13. #include <libgearman/connection.h>
  14. #include <libgearman/packet.hpp>
  15. #include <libgearman/universal.hpp>
  16. #include <cassert>
  17. #include <cerrno>
  18. #include <cstdarg>
  19. #include <cstdio>
  20. #include <cstdlib>
  21. #include <cstring>
  22. void gearman_universal_initialize(gearman_universal_st &self, const gearman_universal_options_t *options)
  23. {
  24. { // Set defaults on all options.
  25. self.options.dont_track_packets= false;
  26. self.options.non_blocking= false;
  27. self.options.stored_non_blocking= false;
  28. }
  29. if (options)
  30. {
  31. while (*options != GEARMAN_MAX)
  32. {
  33. /**
  34. @note Check for bad value, refactor gearman_add_options().
  35. */
  36. gearman_universal_add_options(self, *options);
  37. options++;
  38. }
  39. }
  40. self.verbose= GEARMAN_VERBOSE_NEVER;
  41. self.con_count= 0;
  42. self.packet_count= 0;
  43. self.pfds_size= 0;
  44. self.sending= 0;
  45. self.timeout= -1;
  46. self.con_list= NULL;
  47. self.packet_list= NULL;
  48. self.pfds= NULL;
  49. self.log_fn= NULL;
  50. self.log_context= NULL;
  51. self.workload_malloc_fn= NULL;
  52. self.workload_malloc_context= NULL;
  53. self.workload_free_fn= NULL;
  54. self.workload_free_context= NULL;
  55. self.error.rc= GEARMAN_SUCCESS;
  56. self.error.last_errno= 0;
  57. self.error.last_error[0]= 0;
  58. }
  59. void gearman_universal_clone(gearman_universal_st &destination, const gearman_universal_st &source)
  60. {
  61. gearman_universal_initialize(destination);
  62. (void)gearman_universal_set_option(destination, GEARMAN_NON_BLOCKING, source.options.non_blocking);
  63. (void)gearman_universal_set_option(destination, GEARMAN_DONT_TRACK_PACKETS, source.options.dont_track_packets);
  64. destination.timeout= source.timeout;
  65. for (gearman_connection_st *con= source.con_list; con; con= con->next)
  66. {
  67. if (gearman_connection_clone(&destination, NULL, con) == NULL)
  68. {
  69. gearman_universal_free(destination);
  70. return;
  71. }
  72. }
  73. }
  74. void gearman_universal_free(gearman_universal_st &universal)
  75. {
  76. gearman_free_all_cons(&universal);
  77. gearman_free_all_packets(universal);
  78. if (universal.pfds)
  79. {
  80. // created realloc()
  81. free(universal.pfds);
  82. }
  83. }
  84. gearman_return_t gearman_universal_set_option(gearman_universal_st &self, gearman_universal_options_t option, bool value)
  85. {
  86. switch (option)
  87. {
  88. case GEARMAN_NON_BLOCKING:
  89. self.options.non_blocking= value;
  90. break;
  91. case GEARMAN_DONT_TRACK_PACKETS:
  92. self.options.dont_track_packets= value;
  93. break;
  94. case GEARMAN_MAX:
  95. default:
  96. return GEARMAN_INVALID_COMMAND;
  97. }
  98. return GEARMAN_SUCCESS;
  99. }
  100. int gearman_universal_timeout(gearman_universal_st &self)
  101. {
  102. return self.timeout;
  103. }
  104. void gearman_universal_set_timeout(gearman_universal_st &self, int timeout)
  105. {
  106. self.timeout= timeout;
  107. }
  108. void gearman_set_log_fn(gearman_universal_st &self, gearman_log_fn *function,
  109. void *context, gearman_verbose_t verbose)
  110. {
  111. self.log_fn= function;
  112. self.log_context= context;
  113. self.verbose= verbose;
  114. }
  115. void gearman_set_workload_malloc_fn(gearman_universal_st *universal,
  116. gearman_malloc_fn *function,
  117. void *context)
  118. {
  119. universal->workload_malloc_fn= function;
  120. universal->workload_malloc_context= context;
  121. }
  122. void gearman_set_workload_free_fn(gearman_universal_st *universal,
  123. gearman_free_fn *function,
  124. void *context)
  125. {
  126. universal->workload_free_fn= function;
  127. universal->workload_free_context= context;
  128. }
  129. void gearman_free_all_cons(gearman_universal_st *universal)
  130. {
  131. while (universal->con_list != NULL)
  132. gearman_connection_free(universal->con_list);
  133. }
  134. gearman_return_t gearman_flush_all(gearman_universal_st *universal)
  135. {
  136. for (gearman_connection_st *con= universal->con_list; con != NULL; con= con->next)
  137. {
  138. if (con->events & POLLOUT)
  139. continue;
  140. gearman_return_t ret= gearman_connection_flush(con);
  141. if (gearman_failed(ret) and ret != GEARMAN_IO_WAIT)
  142. return ret;
  143. }
  144. return GEARMAN_SUCCESS;
  145. }
  146. gearman_return_t gearman_wait(gearman_universal_st *universal)
  147. {
  148. struct pollfd *pfds;
  149. nfds_t x;
  150. int ret;
  151. gearman_return_t gret;
  152. if (universal->pfds_size < universal->con_count)
  153. {
  154. pfds= static_cast<pollfd*>(realloc(universal->pfds, universal->con_count * sizeof(struct pollfd)));
  155. if (pfds == NULL)
  156. {
  157. gearman_perror(*universal, "pollfd realloc");
  158. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  159. }
  160. universal->pfds= pfds;
  161. universal->pfds_size= universal->con_count;
  162. }
  163. else
  164. {
  165. pfds= universal->pfds;
  166. }
  167. x= 0;
  168. for (gearman_connection_st *con= universal->con_list; con != NULL; con= con->next)
  169. {
  170. if (con->events == 0)
  171. continue;
  172. pfds[x].fd= con->fd;
  173. pfds[x].events= con->events;
  174. pfds[x].revents= 0;
  175. x++;
  176. }
  177. if (x == 0)
  178. {
  179. gearman_error(universal, GEARMAN_NO_ACTIVE_FDS, "no active file descriptors");
  180. return GEARMAN_NO_ACTIVE_FDS;
  181. }
  182. while (1)
  183. {
  184. ret= poll(pfds, x, universal->timeout);
  185. if (ret == -1)
  186. {
  187. if (errno == EINTR)
  188. continue;
  189. gearman_perror(*universal, "poll");
  190. return GEARMAN_ERRNO;
  191. }
  192. break;
  193. }
  194. if (ret == 0)
  195. {
  196. gearman_error(universal, GEARMAN_TIMEOUT, "timeout reached");
  197. return GEARMAN_TIMEOUT;
  198. }
  199. x= 0;
  200. for (gearman_connection_st *con= universal->con_list; con != NULL; con= con->next)
  201. {
  202. if (con->events == 0)
  203. continue;
  204. gret= gearman_connection_set_revents(con, pfds[x].revents);
  205. if (gret != GEARMAN_SUCCESS)
  206. return gret;
  207. x++;
  208. }
  209. return GEARMAN_SUCCESS;
  210. }
  211. gearman_connection_st *gearman_ready(gearman_universal_st *universal)
  212. {
  213. gearman_connection_st *con;
  214. /* We can't keep universal between calls since connections may be removed during
  215. processing. If this list ever gets big, we may want something faster. */
  216. for (con= universal->con_list; con != NULL; con= con->next)
  217. {
  218. if (con->options.ready)
  219. {
  220. con->options.ready= false;
  221. return con;
  222. }
  223. }
  224. return NULL;
  225. }
  226. /**
  227. @note _push_blocking is only used for echo (and should be fixed
  228. when tricky flip/flop in IO is fixed).
  229. */
  230. static inline void _push_blocking(gearman_universal_st *universal, bool *orig_block_universal)
  231. {
  232. *orig_block_universal= universal->options.non_blocking;
  233. universal->options.non_blocking= false;
  234. }
  235. static inline void _pop_non_blocking(gearman_universal_st *universal, bool orig_block_universal)
  236. {
  237. universal->options.non_blocking= orig_block_universal;
  238. }
  239. gearman_return_t gearman_echo(gearman_universal_st *universal,
  240. const void *workload,
  241. size_t workload_size)
  242. {
  243. gearman_connection_st *con;
  244. gearman_packet_st packet;
  245. gearman_return_t ret;
  246. bool orig_block_universal;
  247. ret= gearman_packet_create_args(*universal, &packet, GEARMAN_MAGIC_REQUEST,
  248. GEARMAN_COMMAND_ECHO_REQ,
  249. &workload, &workload_size, 1);
  250. if (gearman_failed(ret))
  251. {
  252. return ret;
  253. }
  254. _push_blocking(universal, &orig_block_universal);
  255. for (con= universal->con_list; con != NULL; con= con->next)
  256. {
  257. gearman_packet_st *packet_ptr;
  258. ret= gearman_connection_send(con, &packet, true);
  259. if (gearman_failed(ret))
  260. {
  261. goto exit;
  262. }
  263. packet_ptr= gearman_connection_recv(con, &(con->packet), &ret, true);
  264. if (gearman_failed(ret))
  265. {
  266. goto exit;
  267. }
  268. assert(packet_ptr);
  269. if (con->packet.data_size != workload_size ||
  270. memcmp(workload, con->packet.data, workload_size))
  271. {
  272. gearman_packet_free(&(con->packet));
  273. gearman_error(universal, GEARMAN_ECHO_DATA_CORRUPTION, "corruption during echo");
  274. ret= GEARMAN_ECHO_DATA_CORRUPTION;
  275. goto exit;
  276. }
  277. gearman_packet_free(&(con->packet));
  278. }
  279. ret= GEARMAN_SUCCESS;
  280. exit:
  281. gearman_packet_free(&packet);
  282. _pop_non_blocking(universal, orig_block_universal);
  283. return ret;
  284. }
  285. bool gearman_request_option(gearman_universal_st &universal,
  286. gearman_string_t &option)
  287. {
  288. gearman_connection_st *con;
  289. gearman_packet_st packet;
  290. gearman_return_t ret;
  291. bool orig_block_universal;
  292. const void *args[]= { gearman_c_str(option) };
  293. size_t args_size[]= { gearman_size(option) };
  294. ret= gearman_packet_create_args(universal, &packet, GEARMAN_MAGIC_REQUEST,
  295. GEARMAN_COMMAND_OPTION_REQ,
  296. args, args_size, 1);
  297. if (gearman_failed(ret))
  298. {
  299. gearman_error(&universal, GEARMAN_MEMORY_ALLOCATION_FAILURE, "gearman_packet_create_args()");
  300. return ret;
  301. }
  302. _push_blocking(&universal, &orig_block_universal);
  303. for (con= universal.con_list; con != NULL; con= con->next)
  304. {
  305. gearman_packet_st *packet_ptr;
  306. ret= gearman_connection_send(con, &packet, true);
  307. if (gearman_failed(ret))
  308. {
  309. gearman_packet_free(&(con->packet));
  310. goto exit;
  311. }
  312. packet_ptr= gearman_connection_recv(con, &(con->packet), &ret, true);
  313. if (gearman_failed(ret))
  314. {
  315. gearman_packet_free(&(con->packet));
  316. goto exit;
  317. }
  318. assert(packet_ptr);
  319. if (packet_ptr->command == GEARMAN_COMMAND_ERROR)
  320. {
  321. gearman_packet_free(&(con->packet));
  322. gearman_error(&universal, GEARMAN_INVALID_ARGUMENT, "invalid server option");
  323. ret= GEARMAN_INVALID_ARGUMENT;;
  324. goto exit;
  325. }
  326. gearman_packet_free(&(con->packet));
  327. }
  328. ret= GEARMAN_SUCCESS;
  329. exit:
  330. gearman_packet_free(&packet);
  331. _pop_non_blocking(&universal, orig_block_universal);
  332. return gearman_success(ret);
  333. }
  334. void gearman_free_all_packets(gearman_universal_st &universal)
  335. {
  336. while (universal.packet_list)
  337. gearman_packet_free(universal.packet_list);
  338. }
  339. /*
  340. * Local Definitions
  341. */
  342. void gearman_universal_set_error(gearman_universal_st *universal,
  343. gearman_return_t rc,
  344. const char *function,
  345. const char *format, ...)
  346. {
  347. size_t size;
  348. char log_buffer[GEARMAN_MAX_ERROR_SIZE];
  349. char *ptr= log_buffer;
  350. va_list args;
  351. universal->error.rc= rc;
  352. if (rc != GEARMAN_ERRNO)
  353. {
  354. universal->error.last_errno= 0;
  355. }
  356. size= strlen(gearman_strerror(rc));
  357. ptr= static_cast<char *>(memcpy(ptr, gearman_strerror(rc), size));
  358. ptr+= size;
  359. ptr[0]= '-';
  360. size++;
  361. ptr++;
  362. ptr[0]= '>';
  363. size++;
  364. ptr++;
  365. size= strlen(function);
  366. ptr= static_cast<char *>(memcpy(ptr, static_cast<void*>(const_cast<char *>(function)), size));
  367. ptr+= size;
  368. ptr[0]= ':';
  369. size++;
  370. ptr++;
  371. ptr[0]= ' ';
  372. size++;
  373. ptr++;
  374. va_start(args, format);
  375. size+= size_t(vsnprintf(ptr, GEARMAN_MAX_ERROR_SIZE - size, format, args));
  376. va_end(args);
  377. if (universal->log_fn == NULL)
  378. {
  379. if (size >= GEARMAN_MAX_ERROR_SIZE)
  380. size= GEARMAN_MAX_ERROR_SIZE - 1;
  381. memcpy(universal->error.last_error, log_buffer, size + 1);
  382. }
  383. else
  384. {
  385. universal->log_fn(log_buffer, GEARMAN_VERBOSE_FATAL,
  386. static_cast<void *>(universal->log_context));
  387. }
  388. }
  389. void gearman_universal_set_perror(const char *position, gearman_universal_st &self, const char *message)
  390. {
  391. self.error.rc= GEARMAN_ERRNO;
  392. self.error.last_errno= errno;
  393. const char *errmsg_ptr;
  394. char errmsg[GEARMAN_MAX_ERROR_SIZE];
  395. errmsg[0]= 0;
  396. #ifdef STRERROR_R_CHAR_P
  397. errmsg_ptr= strerror_r(self.error.last_errno, errmsg, sizeof(errmsg));
  398. #else
  399. strerror_r(self.error.last_errno, errmsg, sizeof(errmsg));
  400. errmsg_ptr= errmsg;
  401. #endif
  402. char final[GEARMAN_MAX_ERROR_SIZE];
  403. snprintf(final, sizeof(final), "%s(%s)", message, errmsg_ptr);
  404. gearman_universal_set_error(&self, GEARMAN_ERRNO, position, final);
  405. }