worker.cc 36 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360
  1. /* vim:expandtab:shiftwidth=2:tabstop=2:smarttab:
  2. *
  3. * Gearmand client and server library.
  4. *
  5. * Copyright (C) 2011 Data Differential, http://datadifferential.com/
  6. * Copyright (C) 2008 Brian Aker, Eric Day
  7. * All rights reserved.
  8. *
  9. * Redistribution and use in source and binary forms, with or without
  10. * modification, are permitted provided that the following conditions are
  11. * met:
  12. *
  13. * * Redistributions of source code must retain the above copyright
  14. * notice, this list of conditions and the following disclaimer.
  15. *
  16. * * Redistributions in binary form must reproduce the above
  17. * copyright notice, this list of conditions and the following disclaimer
  18. * in the documentation and/or other materials provided with the
  19. * distribution.
  20. *
  21. * * The names of its contributors may not be used to endorse or
  22. * promote products derived from this software without specific prior
  23. * written permission.
  24. *
  25. * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
  26. * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
  27. * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
  28. * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
  29. * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
  30. * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
  31. * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
  32. * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
  33. * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
  34. * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
  35. * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
  36. *
  37. */
  38. #include <config.h>
  39. #include <libgearman/common.h>
  40. #include <libgearman/function/base.hpp>
  41. #include <libgearman/function/make.hpp>
  42. #include <cassert>
  43. #include <cstdio>
  44. #include <cstdlib>
  45. #include <cstring>
  46. #include <memory>
  47. #include <unistd.h>
  48. #include <fcntl.h>
  49. /**
  50. * @addtogroup gearman_worker_static Static Worker Declarations
  51. * @ingroup gearman_worker
  52. * @{
  53. */
  54. static inline struct _worker_function_st *_function_exist(gearman_worker_st *worker, const char *function_name, size_t function_length)
  55. {
  56. struct _worker_function_st *function;
  57. for (function= worker->function_list; function;
  58. function= function->next)
  59. {
  60. if (function_length == function->function_length)
  61. {
  62. if (not memcmp(function_name, function->function_name, function_length))
  63. break;
  64. }
  65. }
  66. return function;
  67. }
  68. /**
  69. * Allocate a worker structure.
  70. */
  71. static gearman_worker_st *_worker_allocate(gearman_worker_st *worker, bool is_clone);
  72. /**
  73. * Initialize common packets for later use.
  74. */
  75. static gearman_return_t _worker_packet_init(gearman_worker_st *worker);
  76. /**
  77. * Callback function used when parsing server lists.
  78. */
  79. static gearman_return_t _worker_add_server(const char *host, in_port_t port, void *context);
  80. /**
  81. * Allocate and add a function to the register list.
  82. */
  83. static gearman_return_t _worker_function_create(gearman_worker_st *worker,
  84. const char *function_name, size_t function_length,
  85. const gearman_function_t &function,
  86. uint32_t timeout,
  87. void *context);
  88. /**
  89. * Free a function.
  90. */
  91. static void _worker_function_free(gearman_worker_st *worker,
  92. struct _worker_function_st *function);
  93. /** @} */
  94. /*
  95. * Public Definitions
  96. */
  97. gearman_worker_st *gearman_worker_create(gearman_worker_st *worker)
  98. {
  99. worker= _worker_allocate(worker, false);
  100. if (worker == NULL)
  101. {
  102. return NULL;
  103. }
  104. if (gearman_failed(_worker_packet_init(worker)))
  105. {
  106. gearman_worker_free(worker);
  107. return NULL;
  108. }
  109. return worker;
  110. }
  111. gearman_worker_st *gearman_worker_clone(gearman_worker_st *worker,
  112. const gearman_worker_st *source)
  113. {
  114. if (source == NULL)
  115. {
  116. return _worker_allocate(worker, false);
  117. }
  118. worker= _worker_allocate(worker, true);
  119. if (worker == NULL)
  120. {
  121. return worker;
  122. }
  123. worker->options.non_blocking= source->options.non_blocking;
  124. worker->options.change= source->options.change;
  125. worker->options.grab_uniq= source->options.grab_uniq;
  126. worker->options.grab_all= source->options.grab_all;
  127. worker->options.timeout_return= source->options.timeout_return;
  128. gearman_universal_clone(worker->universal, source->universal, true);
  129. if (gearman_failed(_worker_packet_init(worker)))
  130. {
  131. gearman_worker_free(worker);
  132. return NULL;
  133. }
  134. return worker;
  135. }
  136. void gearman_worker_free(gearman_worker_st *worker)
  137. {
  138. if (worker == NULL)
  139. {
  140. return;
  141. }
  142. if (worker->universal.wakeup_fd[0] != INVALID_SOCKET)
  143. {
  144. close(worker->universal.wakeup_fd[0]);
  145. }
  146. if (worker->universal.wakeup_fd[1] != INVALID_SOCKET)
  147. {
  148. close(worker->universal.wakeup_fd[1]);
  149. }
  150. gearman_worker_unregister_all(worker);
  151. if (worker->options.packet_init)
  152. {
  153. gearman_packet_free(&worker->grab_job);
  154. gearman_packet_free(&worker->pre_sleep);
  155. }
  156. gearman_job_free(worker->job);
  157. worker->work_job= NULL;
  158. if (worker->work_result)
  159. {
  160. gearman_free(worker->universal, worker->work_result);
  161. }
  162. while (worker->function_list)
  163. {
  164. _worker_function_free(worker, worker->function_list);
  165. }
  166. gearman_job_free_all(worker);
  167. gearman_universal_free(worker->universal);
  168. if (worker->options.allocated)
  169. {
  170. delete worker;
  171. }
  172. }
  173. const char *gearman_worker_error(const gearman_worker_st *worker)
  174. {
  175. if (worker == NULL)
  176. {
  177. return NULL;
  178. }
  179. return gearman_universal_error(worker->universal);
  180. }
  181. int gearman_worker_errno(gearman_worker_st *worker)
  182. {
  183. if (worker == NULL)
  184. {
  185. return 0;
  186. }
  187. return gearman_universal_errno(worker->universal);
  188. }
  189. gearman_worker_options_t gearman_worker_options(const gearman_worker_st *worker)
  190. {
  191. if (worker == NULL)
  192. {
  193. return gearman_worker_options_t();
  194. }
  195. int options;
  196. memset(&options, 0, sizeof(gearman_worker_options_t));
  197. if (worker->options.allocated)
  198. options|= int(GEARMAN_WORKER_ALLOCATED);
  199. if (worker->options.non_blocking)
  200. options|= int(GEARMAN_WORKER_NON_BLOCKING);
  201. if (worker->options.packet_init)
  202. options|= int(GEARMAN_WORKER_PACKET_INIT);
  203. if (worker->options.change)
  204. options|= int(GEARMAN_WORKER_CHANGE);
  205. if (worker->options.grab_uniq)
  206. options|= int(GEARMAN_WORKER_GRAB_UNIQ);
  207. if (worker->options.grab_all)
  208. options|= int(GEARMAN_WORKER_GRAB_ALL);
  209. if (worker->options.timeout_return)
  210. options|= int(GEARMAN_WORKER_TIMEOUT_RETURN);
  211. return gearman_worker_options_t(options);
  212. }
  213. void gearman_worker_set_options(gearman_worker_st *worker,
  214. gearman_worker_options_t options)
  215. {
  216. if (worker == NULL)
  217. {
  218. return;
  219. }
  220. gearman_worker_options_t usable_options[]= {
  221. GEARMAN_WORKER_NON_BLOCKING,
  222. GEARMAN_WORKER_GRAB_UNIQ,
  223. GEARMAN_WORKER_GRAB_ALL,
  224. GEARMAN_WORKER_TIMEOUT_RETURN,
  225. GEARMAN_WORKER_MAX
  226. };
  227. gearman_worker_options_t *ptr;
  228. for (ptr= usable_options; *ptr != GEARMAN_WORKER_MAX ; ptr++)
  229. {
  230. if (options & *ptr)
  231. {
  232. gearman_worker_add_options(worker, *ptr);
  233. }
  234. else
  235. {
  236. gearman_worker_remove_options(worker, *ptr);
  237. }
  238. }
  239. }
  240. void gearman_worker_add_options(gearman_worker_st *worker,
  241. gearman_worker_options_t options)
  242. {
  243. if (worker == NULL)
  244. {
  245. return;
  246. }
  247. if (options & GEARMAN_WORKER_NON_BLOCKING)
  248. {
  249. gearman_universal_add_options(worker->universal, GEARMAN_NON_BLOCKING);
  250. worker->options.non_blocking= true;
  251. }
  252. if (options & GEARMAN_WORKER_GRAB_UNIQ)
  253. {
  254. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB_UNIQ;
  255. gearman_return_t rc= gearman_packet_pack_header(&(worker->grab_job));
  256. assert(gearman_success(rc));
  257. worker->options.grab_uniq= true;
  258. }
  259. if (options & GEARMAN_WORKER_GRAB_ALL)
  260. {
  261. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB_ALL;
  262. gearman_return_t rc= gearman_packet_pack_header(&(worker->grab_job));
  263. assert(gearman_success(rc));
  264. worker->options.grab_all= true;
  265. }
  266. if (options & GEARMAN_WORKER_TIMEOUT_RETURN)
  267. {
  268. worker->options.timeout_return= true;
  269. }
  270. }
  271. void gearman_worker_remove_options(gearman_worker_st *worker,
  272. gearman_worker_options_t options)
  273. {
  274. if (worker == NULL)
  275. {
  276. return;
  277. }
  278. if (options & GEARMAN_WORKER_NON_BLOCKING)
  279. {
  280. gearman_universal_remove_options(worker->universal, GEARMAN_NON_BLOCKING);
  281. worker->options.non_blocking= false;
  282. }
  283. if (options & GEARMAN_WORKER_TIMEOUT_RETURN)
  284. {
  285. worker->options.timeout_return= false;
  286. gearman_universal_set_timeout(worker->universal, GEARMAN_WORKER_WAIT_TIMEOUT);
  287. }
  288. if (options & GEARMAN_WORKER_GRAB_UNIQ)
  289. {
  290. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB;
  291. (void)gearman_packet_pack_header(&(worker->grab_job));
  292. worker->options.grab_uniq= false;
  293. }
  294. if (options & GEARMAN_WORKER_GRAB_ALL)
  295. {
  296. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB;
  297. (void)gearman_packet_pack_header(&(worker->grab_job));
  298. worker->options.grab_all= false;
  299. }
  300. }
  301. int gearman_worker_timeout(gearman_worker_st *worker)
  302. {
  303. if (worker == NULL)
  304. {
  305. return 0;
  306. }
  307. return gearman_universal_timeout(worker->universal);
  308. }
  309. void gearman_worker_set_timeout(gearman_worker_st *worker, int timeout)
  310. {
  311. if (worker == NULL)
  312. {
  313. return;
  314. }
  315. gearman_worker_add_options(worker, GEARMAN_WORKER_TIMEOUT_RETURN);
  316. gearman_universal_set_timeout(worker->universal, timeout);
  317. }
  318. void *gearman_worker_context(const gearman_worker_st *worker)
  319. {
  320. if (worker == NULL)
  321. {
  322. return NULL;
  323. }
  324. return worker->context;
  325. }
  326. void gearman_worker_set_context(gearman_worker_st *worker, void *context)
  327. {
  328. if (worker == NULL)
  329. {
  330. return;
  331. }
  332. worker->context= context;
  333. }
  334. void gearman_worker_set_log_fn(gearman_worker_st *worker,
  335. gearman_log_fn *function, void *context,
  336. gearman_verbose_t verbose)
  337. {
  338. gearman_set_log_fn(worker->universal, function, context, verbose);
  339. }
  340. void gearman_worker_set_workload_malloc_fn(gearman_worker_st *worker,
  341. gearman_malloc_fn *function,
  342. void *context)
  343. {
  344. if (worker == NULL)
  345. {
  346. return;
  347. }
  348. gearman_set_workload_malloc_fn(worker->universal, function, context);
  349. }
  350. void gearman_worker_set_workload_free_fn(gearman_worker_st *worker,
  351. gearman_free_fn *function,
  352. void *context)
  353. {
  354. if (worker == NULL)
  355. {
  356. return;
  357. }
  358. gearman_set_workload_free_fn(worker->universal, function, context);
  359. }
  360. gearman_return_t gearman_worker_add_server(gearman_worker_st *worker,
  361. const char *host, in_port_t port)
  362. {
  363. if (worker == NULL)
  364. {
  365. return GEARMAN_INVALID_ARGUMENT;
  366. }
  367. if (gearman_connection_create_args(worker->universal, host, port) == NULL)
  368. {
  369. return gearman_universal_error_code(worker->universal);
  370. }
  371. return GEARMAN_SUCCESS;
  372. }
  373. gearman_return_t gearman_worker_add_servers(gearman_worker_st *worker, const char *servers)
  374. {
  375. return gearman_parse_servers(servers, _worker_add_server, worker);
  376. }
  377. void gearman_worker_remove_servers(gearman_worker_st *worker)
  378. {
  379. if (worker == NULL)
  380. {
  381. return;
  382. }
  383. gearman_free_all_cons(worker->universal);
  384. }
  385. gearman_return_t gearman_worker_wait(gearman_worker_st *worker)
  386. {
  387. if (worker == NULL)
  388. {
  389. return GEARMAN_INVALID_ARGUMENT;
  390. }
  391. return gearman_wait(worker->universal);
  392. }
  393. gearman_return_t gearman_worker_register(gearman_worker_st *worker,
  394. const char *function_name,
  395. uint32_t timeout)
  396. {
  397. gearman_function_t null_func= gearman_function_create_null();
  398. return _worker_function_create(worker, function_name, strlen(function_name), null_func, timeout, NULL);
  399. }
  400. bool gearman_worker_function_exist(gearman_worker_st *worker,
  401. const char *function_name,
  402. size_t function_length)
  403. {
  404. struct _worker_function_st *function;
  405. function= _function_exist(worker, function_name, function_length);
  406. return (function && function->options.remove == false) ? true : false;
  407. }
  408. static inline gearman_return_t _worker_unregister(gearman_worker_st *worker,
  409. const char *function_name, size_t function_length)
  410. {
  411. struct _worker_function_st *function;
  412. gearman_return_t ret;
  413. const void *args[1];
  414. size_t args_size[1];
  415. function= _function_exist(worker, function_name, function_length);
  416. if (function == NULL || function->options.remove)
  417. {
  418. return GEARMAN_NO_REGISTERED_FUNCTION;
  419. }
  420. gearman_packet_free(&(function->packet));
  421. args[0]= function->name();
  422. args_size[0]= function->length();
  423. ret= gearman_packet_create_args(worker->universal, function->packet,
  424. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_CANT_DO,
  425. args, args_size, 1);
  426. if (gearman_failed(ret))
  427. {
  428. function->options.packet_in_use= false;
  429. return ret;
  430. }
  431. function->options.change= true;
  432. function->options.remove= true;
  433. worker->options.change= true;
  434. return GEARMAN_SUCCESS;
  435. }
  436. gearman_return_t gearman_worker_unregister(gearman_worker_st *worker,
  437. const char *function_name)
  438. {
  439. return _worker_unregister(worker, function_name, strlen(function_name));
  440. }
  441. gearman_return_t gearman_worker_unregister_all(gearman_worker_st *worker)
  442. {
  443. struct _worker_function_st *function;
  444. uint32_t count= 0;
  445. if (worker->function_list == NULL)
  446. {
  447. return GEARMAN_NO_REGISTERED_FUNCTIONS;
  448. }
  449. /* Lets find out if we have any functions left that are valid */
  450. for (function= worker->function_list; function;
  451. function= function->next)
  452. {
  453. if (function->options.remove == false)
  454. {
  455. count++;
  456. }
  457. }
  458. if (count == 0)
  459. {
  460. return GEARMAN_NO_REGISTERED_FUNCTIONS;
  461. }
  462. gearman_packet_free(&(worker->function_list->packet));
  463. gearman_return_t ret= gearman_packet_create_args(worker->universal,
  464. worker->function_list->packet,
  465. GEARMAN_MAGIC_REQUEST,
  466. GEARMAN_COMMAND_RESET_ABILITIES,
  467. NULL, NULL, 0);
  468. if (gearman_failed(ret))
  469. {
  470. worker->function_list->options.packet_in_use= false;
  471. return ret;
  472. }
  473. while (worker->function_list->next)
  474. {
  475. _worker_function_free(worker, worker->function_list->next);
  476. }
  477. worker->function_list->options.change= true;
  478. worker->function_list->options.remove= true;
  479. worker->options.change= true;
  480. return GEARMAN_SUCCESS;
  481. }
  482. gearman_job_st *gearman_worker_grab_job(gearman_worker_st *worker,
  483. gearman_job_st *job,
  484. gearman_return_t *ret_ptr)
  485. {
  486. struct _worker_function_st *function;
  487. uint32_t active;
  488. gearman_return_t unused;
  489. if (not ret_ptr)
  490. {
  491. ret_ptr= &unused;
  492. }
  493. while (1)
  494. {
  495. switch (worker->state)
  496. {
  497. case GEARMAN_WORKER_STATE_START:
  498. /* If there are any new functions changes, send them now. */
  499. if (worker->options.change)
  500. {
  501. worker->function= worker->function_list;
  502. while (worker->function)
  503. {
  504. if (not (worker->function->options.change))
  505. {
  506. worker->function= worker->function->next;
  507. continue;
  508. }
  509. for (worker->con= (&worker->universal)->con_list; worker->con;
  510. worker->con= worker->con->next)
  511. {
  512. if (worker->con->fd == -1)
  513. continue;
  514. case GEARMAN_WORKER_STATE_FUNCTION_SEND:
  515. *ret_ptr= worker->con->send_packet(worker->function->packet, true);
  516. if (gearman_failed(*ret_ptr))
  517. {
  518. if (*ret_ptr == GEARMAN_IO_WAIT)
  519. {
  520. worker->state= GEARMAN_WORKER_STATE_FUNCTION_SEND;
  521. }
  522. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  523. {
  524. continue;
  525. }
  526. return NULL;
  527. }
  528. }
  529. if (worker->function->options.remove)
  530. {
  531. function= worker->function->prev;
  532. _worker_function_free(worker, worker->function);
  533. if (function == NULL)
  534. worker->function= worker->function_list;
  535. else
  536. worker->function= function;
  537. }
  538. else
  539. {
  540. worker->function->options.change= false;
  541. worker->function= worker->function->next;
  542. }
  543. }
  544. worker->options.change= false;
  545. }
  546. if (not worker->function_list)
  547. {
  548. gearman_error(worker->universal, GEARMAN_NO_REGISTERED_FUNCTIONS, "no functions have been registered");
  549. *ret_ptr= GEARMAN_NO_REGISTERED_FUNCTIONS;
  550. return NULL;
  551. }
  552. for (worker->con= (&worker->universal)->con_list; worker->con;
  553. worker->con= worker->con->next)
  554. {
  555. /* If the connection to the job server is not active, start it. */
  556. if (worker->con->fd == -1)
  557. {
  558. for (worker->function= worker->function_list;
  559. worker->function;
  560. worker->function= worker->function->next)
  561. {
  562. case GEARMAN_WORKER_STATE_CONNECT:
  563. *ret_ptr= worker->con->send_packet(worker->function->packet, true);
  564. if (gearman_failed(*ret_ptr))
  565. {
  566. if (*ret_ptr == GEARMAN_IO_WAIT)
  567. {
  568. worker->state= GEARMAN_WORKER_STATE_CONNECT;
  569. }
  570. else if (*ret_ptr == GEARMAN_COULD_NOT_CONNECT or *ret_ptr == GEARMAN_LOST_CONNECTION)
  571. {
  572. break;
  573. }
  574. return NULL;
  575. }
  576. }
  577. if (*ret_ptr == GEARMAN_COULD_NOT_CONNECT)
  578. {
  579. continue;
  580. }
  581. }
  582. case GEARMAN_WORKER_STATE_GRAB_JOB_SEND:
  583. if (worker->con->fd == -1)
  584. continue;
  585. *ret_ptr= worker->con->send_packet(worker->grab_job, true);
  586. if (gearman_failed(*ret_ptr))
  587. {
  588. if (*ret_ptr == GEARMAN_IO_WAIT)
  589. {
  590. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_SEND;
  591. }
  592. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  593. {
  594. continue;
  595. }
  596. return NULL;
  597. }
  598. if (not worker->job)
  599. {
  600. worker->job= gearman_job_create(worker, job);
  601. if (not worker->job)
  602. {
  603. *ret_ptr= GEARMAN_MEMORY_ALLOCATION_FAILURE;
  604. return NULL;
  605. }
  606. }
  607. while (1)
  608. {
  609. case GEARMAN_WORKER_STATE_GRAB_JOB_RECV:
  610. assert(worker->job);
  611. (void)worker->con->receiving(worker->job->assigned, *ret_ptr, true);
  612. if (gearman_failed(*ret_ptr))
  613. {
  614. if (*ret_ptr == GEARMAN_IO_WAIT)
  615. {
  616. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_RECV;
  617. }
  618. else
  619. {
  620. gearman_job_free(worker->job);
  621. worker->job= NULL;
  622. if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  623. {
  624. break;
  625. }
  626. }
  627. return NULL;
  628. }
  629. if (worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN or
  630. worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN_ALL or
  631. worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN_UNIQ)
  632. {
  633. worker->job->options.assigned_in_use= true;
  634. worker->job->con= worker->con;
  635. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_SEND;
  636. job= worker->job;
  637. worker->job= NULL;
  638. return job;
  639. }
  640. if (worker->job->assigned.command == GEARMAN_COMMAND_NO_JOB or
  641. worker->job->assigned.command == GEARMAN_COMMAND_OPTION_RES)
  642. {
  643. gearman_packet_free(&(worker->job->assigned));
  644. break;
  645. }
  646. if (worker->job->assigned.command != GEARMAN_COMMAND_NOOP)
  647. {
  648. gearman_universal_set_error(worker->universal, GEARMAN_UNEXPECTED_PACKET, GEARMAN_AT,
  649. "unexpected packet:%s",
  650. gearman_command_info(worker->job->assigned.command)->name);
  651. gearman_packet_free(&(worker->job->assigned));
  652. gearman_job_free(worker->job);
  653. worker->job= NULL;
  654. *ret_ptr= GEARMAN_UNEXPECTED_PACKET;
  655. return NULL;
  656. }
  657. gearman_packet_free(&(worker->job->assigned));
  658. }
  659. }
  660. case GEARMAN_WORKER_STATE_PRE_SLEEP:
  661. for (worker->con= (&worker->universal)->con_list; worker->con;
  662. worker->con= worker->con->next)
  663. {
  664. if (worker->con->fd == INVALID_SOCKET)
  665. {
  666. continue;
  667. }
  668. *ret_ptr= worker->con->send_packet(worker->pre_sleep, true);
  669. if (gearman_failed(*ret_ptr))
  670. {
  671. if (*ret_ptr == GEARMAN_IO_WAIT)
  672. {
  673. worker->state= GEARMAN_WORKER_STATE_PRE_SLEEP;
  674. }
  675. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  676. {
  677. continue;
  678. }
  679. return NULL;
  680. }
  681. }
  682. worker->state= GEARMAN_WORKER_STATE_START;
  683. /* Set a watch on all active connections that we sent a PRE_SLEEP to. */
  684. active= 0;
  685. for (worker->con= worker->universal.con_list; worker->con; worker->con= worker->con->next)
  686. {
  687. if (worker->con->fd == INVALID_SOCKET)
  688. {
  689. continue;
  690. }
  691. worker->con->set_events(POLLIN);
  692. active++;
  693. }
  694. if ((&worker->universal)->options.non_blocking)
  695. {
  696. *ret_ptr= GEARMAN_NO_JOBS;
  697. return NULL;
  698. }
  699. if (active == 0)
  700. {
  701. if (worker->universal.timeout < 0)
  702. {
  703. gearman_nap(GEARMAN_WORKER_WAIT_TIMEOUT);
  704. }
  705. else
  706. {
  707. if (worker->universal.timeout > 0)
  708. {
  709. gearman_nap(worker->universal);
  710. }
  711. if (worker->options.timeout_return)
  712. {
  713. gearman_error(worker->universal, GEARMAN_TIMEOUT, "Option timeout return reached");
  714. *ret_ptr= GEARMAN_TIMEOUT;
  715. return NULL;
  716. }
  717. }
  718. }
  719. else
  720. {
  721. *ret_ptr= gearman_wait(worker->universal);
  722. if (gearman_failed(*ret_ptr) and (*ret_ptr != GEARMAN_TIMEOUT or worker->options.timeout_return))
  723. {
  724. return NULL;
  725. }
  726. }
  727. break;
  728. }
  729. }
  730. }
  731. void gearman_job_free_all(gearman_worker_st *worker)
  732. {
  733. while (worker->job_list)
  734. {
  735. gearman_job_free(worker->job_list);
  736. }
  737. }
  738. gearman_return_t gearman_worker_add_function(gearman_worker_st *worker,
  739. const char *function_name,
  740. uint32_t timeout,
  741. gearman_worker_fn *worker_fn,
  742. void *context)
  743. {
  744. if (worker == NULL)
  745. {
  746. return GEARMAN_INVALID_ARGUMENT;
  747. }
  748. if (function_name == NULL)
  749. {
  750. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function name not given");
  751. }
  752. if (worker_fn == NULL)
  753. {
  754. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function not given");
  755. }
  756. gearman_function_t local= gearman_function_create_v1(worker_fn);
  757. return _worker_function_create(worker,
  758. function_name, strlen(function_name),
  759. local,
  760. timeout,
  761. context);
  762. }
  763. gearman_return_t gearman_worker_define_function(gearman_worker_st *worker,
  764. const char *function_name, const size_t function_name_length,
  765. const gearman_function_t function,
  766. const uint32_t timeout,
  767. void *context)
  768. {
  769. if (worker == NULL)
  770. {
  771. return GEARMAN_INVALID_ARGUMENT;
  772. }
  773. if (function_name == NULL or function_name_length == 0)
  774. {
  775. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function name not given");
  776. }
  777. return _worker_function_create(worker,
  778. function_name, function_name_length,
  779. function,
  780. timeout,
  781. context);
  782. return GEARMAN_INVALID_ARGUMENT;
  783. }
  784. gearman_return_t gearman_worker_work(gearman_worker_st *worker)
  785. {
  786. bool shutdown= false;
  787. if (worker == NULL)
  788. {
  789. return GEARMAN_INVALID_ARGUMENT;
  790. }
  791. switch (worker->work_state)
  792. {
  793. case GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB:
  794. {
  795. gearman_return_t ret;
  796. worker->work_job= gearman_worker_grab_job(worker, NULL, &ret);
  797. if (gearman_failed(ret))
  798. {
  799. if (ret == GEARMAN_COULD_NOT_CONNECT)
  800. {
  801. gearman_reset(worker->universal);
  802. }
  803. return ret;
  804. }
  805. assert(worker->work_job);
  806. for (worker->work_function= worker->function_list;
  807. worker->work_function;
  808. worker->work_function= worker->work_function->next)
  809. {
  810. if (not strcmp(gearman_job_function_name(worker->work_job),
  811. worker->work_function->function_name))
  812. {
  813. break;
  814. }
  815. }
  816. if (not worker->work_function)
  817. {
  818. gearman_job_free(worker->work_job);
  819. worker->work_job= NULL;
  820. return gearman_error(worker->universal, GEARMAN_INVALID_FUNCTION_NAME, "Function not found");
  821. }
  822. if (not worker->work_function->has_callback())
  823. {
  824. gearman_job_free(worker->work_job);
  825. worker->work_job= NULL;
  826. return gearman_error(worker->universal, GEARMAN_INVALID_FUNCTION_NAME, "Neither a gearman_worker_fn, or gearman_function_fn callback was supplied");
  827. }
  828. worker->work_result_size= 0;
  829. }
  830. case GEARMAN_WORKER_WORK_UNIVERSAL_FUNCTION:
  831. {
  832. switch (worker->work_function->callback(worker->work_job,
  833. static_cast<void *>(worker->work_function->context)))
  834. {
  835. case GEARMAN_FUNCTION_INVALID_ARGUMENT:
  836. worker->work_job->error_code= GEARMAN_INVALID_ARGUMENT;
  837. case GEARMAN_FUNCTION_FATAL:
  838. if (gearman_job_send_fail_fin(worker->work_job) == GEARMAN_LOST_CONNECTION) // If we fail this, we have no connection, @note this causes us to lose the current error
  839. {
  840. worker->work_job->error_code= GEARMAN_LOST_CONNECTION;
  841. break;
  842. }
  843. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_FAIL;
  844. return worker->work_job->error_code;
  845. case GEARMAN_FUNCTION_ERROR: // retry
  846. worker->work_job->error_code= GEARMAN_LOST_CONNECTION;
  847. break;
  848. case GEARMAN_FUNCTION_SHUTDOWN:
  849. shutdown= true;
  850. case GEARMAN_FUNCTION_SUCCESS:
  851. break;
  852. }
  853. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  854. {
  855. break;
  856. }
  857. }
  858. case GEARMAN_WORKER_WORK_UNIVERSAL_COMPLETE:
  859. {
  860. worker->work_job->error_code= gearman_job_send_complete_fin(worker->work_job,
  861. worker->work_result, worker->work_result_size);
  862. if (worker->work_job->error_code == GEARMAN_IO_WAIT)
  863. {
  864. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_COMPLETE;
  865. return gearman_error(worker->universal, worker->work_job->error_code,
  866. "A failure occurred after worker had successful complete, unless gearman_job_send_complete() was called directly by worker, client has not been informed of success.");
  867. }
  868. if (worker->work_result)
  869. {
  870. gearman_free(worker->universal, worker->work_result);
  871. worker->work_result= NULL;
  872. }
  873. // If we lost the connection, we retry the work, otherwise we error
  874. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  875. {
  876. break;
  877. }
  878. else if (gearman_failed(worker->work_job->error_code))
  879. {
  880. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_FAIL;
  881. return worker->work_job->error_code;
  882. }
  883. }
  884. break;
  885. case GEARMAN_WORKER_WORK_UNIVERSAL_FAIL:
  886. {
  887. if (gearman_failed(worker->work_job->error_code= gearman_job_send_fail_fin(worker->work_job)))
  888. {
  889. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  890. {
  891. break;
  892. }
  893. return worker->work_job->error_code;
  894. }
  895. }
  896. break;
  897. }
  898. gearman_job_free(worker->work_job);
  899. worker->work_job= NULL;
  900. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB;
  901. return shutdown ? GEARMAN_SHUTDOWN : GEARMAN_SUCCESS;
  902. }
  903. gearman_return_t gearman_worker_echo(gearman_worker_st *worker,
  904. const void *workload,
  905. size_t workload_size)
  906. {
  907. if (worker == NULL)
  908. {
  909. return GEARMAN_INVALID_ARGUMENT;
  910. }
  911. return gearman_echo(worker->universal, workload, workload_size);
  912. }
  913. /*
  914. * Static Definitions
  915. */
  916. static gearman_worker_st *_worker_allocate(gearman_worker_st *worker, bool is_clone)
  917. {
  918. if (worker)
  919. {
  920. worker->options.allocated= false;
  921. }
  922. else
  923. {
  924. worker= new (std::nothrow) gearman_worker_st;
  925. if (worker == NULL)
  926. {
  927. return NULL;
  928. }
  929. worker->options.allocated= true;
  930. }
  931. worker->options.non_blocking= false;
  932. worker->options.packet_init= false;
  933. worker->options.change= false;
  934. worker->options.grab_uniq= true;
  935. worker->options.grab_all= true;
  936. worker->options.timeout_return= false;
  937. worker->state= GEARMAN_WORKER_STATE_START;
  938. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB;
  939. worker->function_count= 0;
  940. worker->job_count= 0;
  941. worker->work_result_size= 0;
  942. worker->context= NULL;
  943. worker->con= NULL;
  944. worker->job= NULL;
  945. worker->job_list= NULL;
  946. worker->function= NULL;
  947. worker->function_list= NULL;
  948. worker->work_function= NULL;
  949. worker->work_result= NULL;
  950. if (is_clone == false)
  951. {
  952. gearman_universal_initialize(worker->universal);
  953. #if 0
  954. gearman_universal_set_timeout(worker->universal, GEARMAN_WORKER_WAIT_TIMEOUT);
  955. #endif
  956. }
  957. if (pipe(worker->universal.wakeup_fd) != 0)
  958. {
  959. delete worker;
  960. return NULL;
  961. }
  962. int returned_flags;
  963. if ((returned_flags= fcntl(worker->universal.wakeup_fd[0], F_GETFL, 0)) < 0)
  964. {
  965. close(worker->universal.wakeup_fd[0]);
  966. close(worker->universal.wakeup_fd[1]);
  967. delete worker;
  968. return NULL;
  969. }
  970. if (fcntl(worker->universal.wakeup_fd[0], F_SETFL, returned_flags | O_NONBLOCK) < 0)
  971. {
  972. close(worker->universal.wakeup_fd[0]);
  973. close(worker->universal.wakeup_fd[1]);
  974. delete worker;
  975. return NULL;
  976. }
  977. #ifdef F_SETNOSIGPIPE
  978. if (fcntl(worker->universal.wakeup_fd[1], F_SETNOSIGPIPE, 0) < 0)
  979. {
  980. close(worker->universal.wakeup_fd[0]);
  981. close(worker->universal.wakeup_fd[1]);
  982. delete worker;
  983. return NULL;
  984. }
  985. #endif
  986. return worker;
  987. }
  988. static gearman_return_t _worker_packet_init(gearman_worker_st *worker)
  989. {
  990. gearman_return_t ret= gearman_packet_create_args(worker->universal, worker->grab_job,
  991. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_GRAB_JOB_ALL,
  992. NULL, NULL, 0);
  993. if (gearman_failed(ret))
  994. {
  995. return ret;
  996. }
  997. ret= gearman_packet_create_args(worker->universal, worker->pre_sleep,
  998. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_PRE_SLEEP,
  999. NULL, NULL, 0);
  1000. if (gearman_failed(ret))
  1001. {
  1002. gearman_packet_free(&(worker->grab_job));
  1003. return ret;
  1004. }
  1005. worker->options.packet_init= true;
  1006. return GEARMAN_SUCCESS;
  1007. }
  1008. static gearman_return_t _worker_add_server(const char *host, in_port_t port, void *context)
  1009. {
  1010. return gearman_worker_add_server(static_cast<gearman_worker_st *>(context), host, port);
  1011. }
  1012. static gearman_return_t _worker_function_create(gearman_worker_st *worker,
  1013. const char *function_name, size_t function_length,
  1014. const gearman_function_t &function_arg,
  1015. uint32_t timeout,
  1016. void *context)
  1017. {
  1018. const void *args[2];
  1019. size_t args_size[2];
  1020. if (worker == NULL)
  1021. {
  1022. return GEARMAN_INVALID_ARGUMENT;
  1023. }
  1024. if (function_length == 0 or function_name == NULL or function_length > GEARMAN_FUNCTION_MAX_SIZE)
  1025. {
  1026. if (function_length > GEARMAN_FUNCTION_MAX_SIZE)
  1027. {
  1028. gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function name longer then GEARMAN_MAX_FUNCTION_SIZE");
  1029. }
  1030. else
  1031. {
  1032. gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "invalid function");
  1033. }
  1034. return GEARMAN_INVALID_ARGUMENT;
  1035. }
  1036. _worker_function_st *function= make(worker->universal._namespace, function_name, function_length, function_arg, context);
  1037. if (function == NULL)
  1038. {
  1039. gearman_perror(worker->universal, "_worker_function_st::new()");
  1040. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  1041. }
  1042. gearman_return_t ret;
  1043. if (timeout > 0)
  1044. {
  1045. char timeout_buffer[11];
  1046. snprintf(timeout_buffer, sizeof(timeout_buffer), "%u", timeout);
  1047. args[0]= function->name();
  1048. args_size[0]= function->length() + 1;
  1049. args[1]= timeout_buffer;
  1050. args_size[1]= strlen(timeout_buffer);
  1051. ret= gearman_packet_create_args(worker->universal, function->packet,
  1052. GEARMAN_MAGIC_REQUEST,
  1053. GEARMAN_COMMAND_CAN_DO_TIMEOUT,
  1054. args, args_size, 2);
  1055. }
  1056. else
  1057. {
  1058. args[0]= function->name();
  1059. args_size[0]= function->length();
  1060. ret= gearman_packet_create_args(worker->universal, function->packet,
  1061. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_CAN_DO,
  1062. args, args_size, 1);
  1063. }
  1064. if (gearman_failed(ret))
  1065. {
  1066. delete function;
  1067. return ret;
  1068. }
  1069. if (worker->function_list)
  1070. {
  1071. worker->function_list->prev= function;
  1072. }
  1073. function->next= worker->function_list;
  1074. function->prev= NULL;
  1075. worker->function_list= function;
  1076. worker->function_count++;
  1077. worker->options.change= true;
  1078. return GEARMAN_SUCCESS;
  1079. }
  1080. static void _worker_function_free(gearman_worker_st *worker,
  1081. struct _worker_function_st *function)
  1082. {
  1083. if (worker->function_list == function)
  1084. {
  1085. worker->function_list= function->next;
  1086. }
  1087. if (function->prev)
  1088. {
  1089. function->prev->next= function->next;
  1090. }
  1091. if (function->next)
  1092. {
  1093. function->next->prev= function->prev;
  1094. }
  1095. worker->function_count--;
  1096. delete function;
  1097. }
  1098. gearman_return_t gearman_worker_set_memory_allocators(gearman_worker_st *worker,
  1099. gearman_malloc_fn *malloc_fn,
  1100. gearman_free_fn *free_fn,
  1101. gearman_realloc_fn *realloc_fn,
  1102. gearman_calloc_fn *calloc_fn,
  1103. void *context)
  1104. {
  1105. if (worker == NULL)
  1106. {
  1107. return GEARMAN_INVALID_ARGUMENT;
  1108. }
  1109. return gearman_set_memory_allocator(worker->universal.allocator, malloc_fn, free_fn, realloc_fn, calloc_fn, context);
  1110. }
  1111. bool gearman_worker_set_server_option(gearman_worker_st *self, const char *option_arg, size_t option_arg_size)
  1112. {
  1113. gearman_string_t option= { option_arg, option_arg_size };
  1114. return gearman_request_option(self->universal, option);
  1115. }
  1116. void gearman_worker_set_namespace(gearman_worker_st *self, const char *namespace_key, size_t namespace_key_size)
  1117. {
  1118. if (self == NULL)
  1119. {
  1120. return;
  1121. }
  1122. gearman_universal_set_namespace(self->universal, namespace_key, namespace_key_size);
  1123. }
  1124. gearman_id_t gearman_worker_id(gearman_worker_st *self)
  1125. {
  1126. if (self == NULL)
  1127. {
  1128. gearman_id_t handle= { INVALID_SOCKET };
  1129. return handle;
  1130. }
  1131. return gearman_universal_id(self->universal);
  1132. }
  1133. gearman_worker_st *gearman_job_clone_worker(gearman_job_st *job)
  1134. {
  1135. return gearman_worker_clone(NULL, job->worker);
  1136. }
  1137. gearman_return_t gearman_worker_set_identifier(gearman_worker_st *worker,
  1138. const char *id, size_t id_size)
  1139. {
  1140. return gearman_set_identifier(worker->universal, id, id_size);
  1141. }