worker.cc 34 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195
  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 <libgearman/common.h>
  39. #include <libgearman/connection.h>
  40. #include <libgearman/packet.hpp>
  41. #include <libgearman/allocator.hpp>
  42. #include <libgearman/universal.hpp>
  43. #include <libgearman/function/base.hpp>
  44. #include <libgearman/function/make.hpp>
  45. #include <cassert>
  46. #include <cstdio>
  47. #include <cstdlib>
  48. #include <cstring>
  49. #include <memory>
  50. #include <unistd.h>
  51. /**
  52. * @addtogroup gearman_worker_static Static Worker Declarations
  53. * @ingroup gearman_worker
  54. * @{
  55. */
  56. static inline struct _worker_function_st *_function_exist(gearman_worker_st *worker, const char *function_name, size_t function_length)
  57. {
  58. struct _worker_function_st *function;
  59. for (function= worker->function_list; function;
  60. function= function->next)
  61. {
  62. if (function_length == function->function_length)
  63. {
  64. if (not memcmp(function_name, function->function_name, function_length))
  65. break;
  66. }
  67. }
  68. return function;
  69. }
  70. /**
  71. * Allocate a worker structure.
  72. */
  73. static gearman_worker_st *_worker_allocate(gearman_worker_st *worker, bool is_clone);
  74. /**
  75. * Initialize common packets for later use.
  76. */
  77. static gearman_return_t _worker_packet_init(gearman_worker_st *worker);
  78. /**
  79. * Callback function used when parsing server lists.
  80. */
  81. static gearman_return_t _worker_add_server(const char *host, in_port_t port, void *context);
  82. /**
  83. * Allocate and add a function to the register list.
  84. */
  85. static gearman_return_t _worker_function_create(gearman_worker_st *worker,
  86. const char *function_name, size_t function_length,
  87. const gearman_function_t &function,
  88. uint32_t timeout,
  89. void *context);
  90. /**
  91. * Free a function.
  92. */
  93. static void _worker_function_free(gearman_worker_st *worker,
  94. struct _worker_function_st *function);
  95. /** @} */
  96. /*
  97. * Public Definitions
  98. */
  99. gearman_worker_st *gearman_worker_create(gearman_worker_st *worker)
  100. {
  101. worker= _worker_allocate(worker, false);
  102. if (not worker)
  103. return NULL;
  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 *from)
  113. {
  114. if (not from)
  115. {
  116. return _worker_allocate(worker, false);
  117. }
  118. worker= _worker_allocate(worker, true);
  119. if (not worker)
  120. return worker;
  121. worker->options.non_blocking= from->options.non_blocking;
  122. worker->options.change= from->options.change;
  123. worker->options.grab_uniq= from->options.grab_uniq;
  124. worker->options.grab_all= from->options.grab_all;
  125. worker->options.timeout_return= from->options.timeout_return;
  126. gearman_universal_clone(worker->universal, from->universal);
  127. if (gearman_failed(_worker_packet_init(worker)))
  128. {
  129. gearman_worker_free(worker);
  130. return NULL;
  131. }
  132. return worker;
  133. }
  134. void gearman_worker_free(gearman_worker_st *worker)
  135. {
  136. if (not worker)
  137. return;
  138. gearman_worker_unregister_all(worker);
  139. if (worker->options.packet_init)
  140. {
  141. gearman_packet_free(&worker->grab_job);
  142. gearman_packet_free(&worker->pre_sleep);
  143. }
  144. gearman_job_free(worker->job);
  145. worker->work_job= NULL;
  146. if (worker->work_result)
  147. {
  148. gearman_free(worker->universal, worker->work_result);
  149. }
  150. while (worker->function_list)
  151. {
  152. _worker_function_free(worker, worker->function_list);
  153. }
  154. gearman_job_free_all(worker);
  155. gearman_universal_free(worker->universal);
  156. if (worker->options.allocated)
  157. delete worker;
  158. }
  159. const char *gearman_worker_error(const gearman_worker_st *worker)
  160. {
  161. if (not worker)
  162. return NULL;
  163. return gearman_universal_error(worker->universal);
  164. }
  165. int gearman_worker_errno(gearman_worker_st *worker)
  166. {
  167. if (not worker)
  168. return 0;
  169. return gearman_universal_errno(worker->universal);
  170. }
  171. gearman_worker_options_t gearman_worker_options(const gearman_worker_st *worker)
  172. {
  173. if (not worker)
  174. return gearman_worker_options_t();
  175. int options;
  176. memset(&options, 0, sizeof(gearman_worker_options_t));
  177. if (worker->options.allocated)
  178. options|= int(GEARMAN_WORKER_ALLOCATED);
  179. if (worker->options.non_blocking)
  180. options|= int(GEARMAN_WORKER_NON_BLOCKING);
  181. if (worker->options.packet_init)
  182. options|= int(GEARMAN_WORKER_PACKET_INIT);
  183. if (worker->options.change)
  184. options|= int(GEARMAN_WORKER_CHANGE);
  185. if (worker->options.grab_uniq)
  186. options|= int(GEARMAN_WORKER_GRAB_UNIQ);
  187. if (worker->options.grab_all)
  188. options|= int(GEARMAN_WORKER_GRAB_ALL);
  189. if (worker->options.timeout_return)
  190. options|= int(GEARMAN_WORKER_TIMEOUT_RETURN);
  191. return gearman_worker_options_t(options);
  192. }
  193. void gearman_worker_set_options(gearman_worker_st *worker,
  194. gearman_worker_options_t options)
  195. {
  196. if (not worker)
  197. return;
  198. gearman_worker_options_t usable_options[]= {
  199. GEARMAN_WORKER_NON_BLOCKING,
  200. GEARMAN_WORKER_GRAB_UNIQ,
  201. GEARMAN_WORKER_GRAB_ALL,
  202. GEARMAN_WORKER_TIMEOUT_RETURN,
  203. GEARMAN_WORKER_MAX
  204. };
  205. gearman_worker_options_t *ptr;
  206. for (ptr= usable_options; *ptr != GEARMAN_WORKER_MAX ; ptr++)
  207. {
  208. if (options & *ptr)
  209. {
  210. gearman_worker_add_options(worker, *ptr);
  211. }
  212. else
  213. {
  214. gearman_worker_remove_options(worker, *ptr);
  215. }
  216. }
  217. }
  218. void gearman_worker_add_options(gearman_worker_st *worker,
  219. gearman_worker_options_t options)
  220. {
  221. if (not worker)
  222. return;
  223. if (options & GEARMAN_WORKER_NON_BLOCKING)
  224. {
  225. gearman_universal_add_options(worker->universal, GEARMAN_NON_BLOCKING);
  226. worker->options.non_blocking= true;
  227. }
  228. if (options & GEARMAN_WORKER_GRAB_UNIQ)
  229. {
  230. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB_UNIQ;
  231. gearman_return_t rc= gearman_packet_pack_header(&(worker->grab_job));
  232. assert(gearman_success(rc));
  233. worker->options.grab_uniq= true;
  234. }
  235. if (options & GEARMAN_WORKER_GRAB_ALL)
  236. {
  237. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB_ALL;
  238. gearman_return_t rc= gearman_packet_pack_header(&(worker->grab_job));
  239. assert(gearman_success(rc));
  240. worker->options.grab_all= true;
  241. }
  242. if (options & GEARMAN_WORKER_TIMEOUT_RETURN)
  243. {
  244. worker->options.timeout_return= true;
  245. }
  246. }
  247. void gearman_worker_remove_options(gearman_worker_st *worker,
  248. gearman_worker_options_t options)
  249. {
  250. if (not worker)
  251. return;
  252. if (options & GEARMAN_WORKER_NON_BLOCKING)
  253. {
  254. gearman_universal_remove_options(worker->universal, GEARMAN_NON_BLOCKING);
  255. worker->options.non_blocking= false;
  256. }
  257. if (options & GEARMAN_WORKER_TIMEOUT_RETURN)
  258. {
  259. worker->options.timeout_return= false;
  260. gearman_universal_set_timeout(worker->universal, GEARMAN_WORKER_WAIT_TIMEOUT);
  261. }
  262. if (options & GEARMAN_WORKER_GRAB_UNIQ)
  263. {
  264. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB;
  265. (void)gearman_packet_pack_header(&(worker->grab_job));
  266. worker->options.grab_uniq= false;
  267. }
  268. if (options & GEARMAN_WORKER_GRAB_ALL)
  269. {
  270. worker->grab_job.command= GEARMAN_COMMAND_GRAB_JOB;
  271. (void)gearman_packet_pack_header(&(worker->grab_job));
  272. worker->options.grab_all= false;
  273. }
  274. }
  275. int gearman_worker_timeout(gearman_worker_st *worker)
  276. {
  277. if (not worker)
  278. return 0;
  279. return gearman_universal_timeout(worker->universal);
  280. }
  281. void gearman_worker_set_timeout(gearman_worker_st *worker, int timeout)
  282. {
  283. if (not worker)
  284. return;
  285. gearman_worker_add_options(worker, GEARMAN_WORKER_TIMEOUT_RETURN);
  286. gearman_universal_set_timeout(worker->universal, timeout);
  287. }
  288. void *gearman_worker_context(const gearman_worker_st *worker)
  289. {
  290. if (not worker)
  291. return NULL;
  292. return worker->context;
  293. }
  294. void gearman_worker_set_context(gearman_worker_st *worker, void *context)
  295. {
  296. if (not worker)
  297. return;
  298. worker->context= context;
  299. }
  300. void gearman_worker_set_log_fn(gearman_worker_st *worker,
  301. gearman_log_fn *function, void *context,
  302. gearman_verbose_t verbose)
  303. {
  304. gearman_set_log_fn(worker->universal, function, context, verbose);
  305. }
  306. void gearman_worker_set_workload_malloc_fn(gearman_worker_st *worker,
  307. gearman_malloc_fn *function,
  308. void *context)
  309. {
  310. if (not worker)
  311. return;
  312. gearman_set_workload_malloc_fn(worker->universal, function, context);
  313. }
  314. void gearman_worker_set_workload_free_fn(gearman_worker_st *worker,
  315. gearman_free_fn *function,
  316. void *context)
  317. {
  318. if (not worker)
  319. return;
  320. gearman_set_workload_free_fn(worker->universal, function, context);
  321. }
  322. gearman_return_t gearman_worker_add_server(gearman_worker_st *worker,
  323. const char *host, in_port_t port)
  324. {
  325. if (not gearman_connection_create_args(worker->universal, host, port))
  326. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  327. return GEARMAN_SUCCESS;
  328. }
  329. gearman_return_t gearman_worker_add_servers(gearman_worker_st *worker, const char *servers)
  330. {
  331. return gearman_parse_servers(servers, _worker_add_server, worker);
  332. }
  333. void gearman_worker_remove_servers(gearman_worker_st *worker)
  334. {
  335. if (not worker)
  336. return;
  337. gearman_free_all_cons(worker->universal);
  338. }
  339. gearman_return_t gearman_worker_wait(gearman_worker_st *worker)
  340. {
  341. if (not worker)
  342. return GEARMAN_INVALID_ARGUMENT;
  343. return gearman_wait(worker->universal);
  344. }
  345. gearman_return_t gearman_worker_register(gearman_worker_st *worker,
  346. const char *function_name,
  347. uint32_t timeout)
  348. {
  349. gearman_function_t null_func= gearman_function_create_null();
  350. return _worker_function_create(worker, function_name, strlen(function_name), null_func, timeout, NULL);
  351. }
  352. bool gearman_worker_function_exist(gearman_worker_st *worker,
  353. const char *function_name,
  354. size_t function_length)
  355. {
  356. struct _worker_function_st *function;
  357. function= _function_exist(worker, function_name, function_length);
  358. return (function && function->options.remove == false) ? true : false;
  359. }
  360. static inline gearman_return_t _worker_unregister(gearman_worker_st *worker,
  361. const char *function_name, size_t function_length)
  362. {
  363. struct _worker_function_st *function;
  364. gearman_return_t ret;
  365. const void *args[1];
  366. size_t args_size[1];
  367. function= _function_exist(worker, function_name, function_length);
  368. if (function == NULL || function->options.remove)
  369. {
  370. return GEARMAN_NO_REGISTERED_FUNCTION;
  371. }
  372. gearman_packet_free(&(function->packet));
  373. args[0]= function->name();
  374. args_size[0]= function->length();
  375. ret= gearman_packet_create_args(worker->universal, function->packet,
  376. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_CANT_DO,
  377. args, args_size, 1);
  378. if (gearman_failed(ret))
  379. {
  380. function->options.packet_in_use= false;
  381. return ret;
  382. }
  383. function->options.change= true;
  384. function->options.remove= true;
  385. worker->options.change= true;
  386. return GEARMAN_SUCCESS;
  387. }
  388. gearman_return_t gearman_worker_unregister(gearman_worker_st *worker,
  389. const char *function_name)
  390. {
  391. return _worker_unregister(worker, function_name, strlen(function_name));
  392. }
  393. gearman_return_t gearman_worker_unregister_all(gearman_worker_st *worker)
  394. {
  395. struct _worker_function_st *function;
  396. uint32_t count= 0;
  397. if (not worker->function_list)
  398. return GEARMAN_NO_REGISTERED_FUNCTIONS;
  399. /* Lets find out if we have any functions left that are valid */
  400. for (function= worker->function_list; function;
  401. function= function->next)
  402. {
  403. if (function->options.remove == false)
  404. count++;
  405. }
  406. if (count == 0)
  407. return GEARMAN_NO_REGISTERED_FUNCTIONS;
  408. gearman_packet_free(&(worker->function_list->packet));
  409. gearman_return_t ret= gearman_packet_create_args(worker->universal,
  410. worker->function_list->packet,
  411. GEARMAN_MAGIC_REQUEST,
  412. GEARMAN_COMMAND_RESET_ABILITIES,
  413. NULL, NULL, 0);
  414. if (gearman_failed(ret))
  415. {
  416. worker->function_list->options.packet_in_use= false;
  417. return ret;
  418. }
  419. while (worker->function_list->next)
  420. _worker_function_free(worker, worker->function_list->next);
  421. worker->function_list->options.change= true;
  422. worker->function_list->options.remove= true;
  423. worker->options.change= true;
  424. return GEARMAN_SUCCESS;
  425. }
  426. gearman_job_st *gearman_worker_grab_job(gearman_worker_st *worker,
  427. gearman_job_st *job,
  428. gearman_return_t *ret_ptr)
  429. {
  430. struct _worker_function_st *function;
  431. uint32_t active;
  432. gearman_return_t unused;
  433. if (not ret_ptr)
  434. ret_ptr= &unused;
  435. while (1)
  436. {
  437. switch (worker->state)
  438. {
  439. case GEARMAN_WORKER_STATE_START:
  440. /* If there are any new functions changes, send them now. */
  441. if (worker->options.change)
  442. {
  443. worker->function= worker->function_list;
  444. while (worker->function)
  445. {
  446. if (not (worker->function->options.change))
  447. {
  448. worker->function= worker->function->next;
  449. continue;
  450. }
  451. for (worker->con= (&worker->universal)->con_list; worker->con;
  452. worker->con= worker->con->next)
  453. {
  454. if (worker->con->fd == -1)
  455. continue;
  456. case GEARMAN_WORKER_STATE_FUNCTION_SEND:
  457. *ret_ptr= worker->con->send(worker->function->packet, true);
  458. if (gearman_failed(*ret_ptr))
  459. {
  460. if (*ret_ptr == GEARMAN_IO_WAIT)
  461. {
  462. worker->state= GEARMAN_WORKER_STATE_FUNCTION_SEND;
  463. }
  464. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  465. {
  466. continue;
  467. }
  468. return NULL;
  469. }
  470. }
  471. if (worker->function->options.remove)
  472. {
  473. function= worker->function->prev;
  474. _worker_function_free(worker, worker->function);
  475. if (function == NULL)
  476. worker->function= worker->function_list;
  477. else
  478. worker->function= function;
  479. }
  480. else
  481. {
  482. worker->function->options.change= false;
  483. worker->function= worker->function->next;
  484. }
  485. }
  486. worker->options.change= false;
  487. }
  488. if (not worker->function_list)
  489. {
  490. gearman_error(worker->universal, GEARMAN_NO_REGISTERED_FUNCTIONS, "no functions have been registered");
  491. *ret_ptr= GEARMAN_NO_REGISTERED_FUNCTIONS;
  492. return NULL;
  493. }
  494. for (worker->con= (&worker->universal)->con_list; worker->con;
  495. worker->con= worker->con->next)
  496. {
  497. /* If the connection to the job server is not active, start it. */
  498. if (worker->con->fd == -1)
  499. {
  500. for (worker->function= worker->function_list;
  501. worker->function;
  502. worker->function= worker->function->next)
  503. {
  504. case GEARMAN_WORKER_STATE_CONNECT:
  505. *ret_ptr= worker->con->send(worker->function->packet, true);
  506. if (gearman_failed(*ret_ptr))
  507. {
  508. if (*ret_ptr == GEARMAN_IO_WAIT)
  509. {
  510. worker->state= GEARMAN_WORKER_STATE_CONNECT;
  511. }
  512. else if (*ret_ptr == GEARMAN_COULD_NOT_CONNECT or *ret_ptr == GEARMAN_LOST_CONNECTION)
  513. {
  514. break;
  515. }
  516. return NULL;
  517. }
  518. }
  519. if (*ret_ptr == GEARMAN_COULD_NOT_CONNECT)
  520. {
  521. continue;
  522. }
  523. }
  524. case GEARMAN_WORKER_STATE_GRAB_JOB_SEND:
  525. if (worker->con->fd == -1)
  526. continue;
  527. *ret_ptr= worker->con->send(worker->grab_job, true);
  528. if (gearman_failed(*ret_ptr))
  529. {
  530. if (*ret_ptr == GEARMAN_IO_WAIT)
  531. {
  532. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_SEND;
  533. }
  534. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  535. {
  536. continue;
  537. }
  538. return NULL;
  539. }
  540. if (not worker->job)
  541. {
  542. worker->job= gearman_job_create(worker, job);
  543. if (not worker->job)
  544. {
  545. *ret_ptr= GEARMAN_MEMORY_ALLOCATION_FAILURE;
  546. return NULL;
  547. }
  548. }
  549. while (1)
  550. {
  551. case GEARMAN_WORKER_STATE_GRAB_JOB_RECV:
  552. assert(worker->job);
  553. (void)worker->con->receiving(worker->job->assigned, *ret_ptr, true);
  554. if (gearman_failed(*ret_ptr))
  555. {
  556. if (*ret_ptr == GEARMAN_IO_WAIT)
  557. {
  558. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_RECV;
  559. }
  560. else
  561. {
  562. gearman_job_free(worker->job);
  563. worker->job= NULL;
  564. if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  565. break;
  566. }
  567. return NULL;
  568. }
  569. if (worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN or
  570. worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN_ALL or
  571. worker->job->assigned.command == GEARMAN_COMMAND_JOB_ASSIGN_UNIQ)
  572. {
  573. worker->job->options.assigned_in_use= true;
  574. worker->job->con= worker->con;
  575. worker->state= GEARMAN_WORKER_STATE_GRAB_JOB_SEND;
  576. job= worker->job;
  577. worker->job= NULL;
  578. return job;
  579. }
  580. if (worker->job->assigned.command == GEARMAN_COMMAND_NO_JOB)
  581. {
  582. gearman_packet_free(&(worker->job->assigned));
  583. break;
  584. }
  585. if (worker->job->assigned.command != GEARMAN_COMMAND_NOOP)
  586. {
  587. gearman_universal_set_error(worker->universal, GEARMAN_UNEXPECTED_PACKET, AT,
  588. "unexpected packet:%s",
  589. gearman_command_info(worker->job->assigned.command)->name);
  590. gearman_packet_free(&(worker->job->assigned));
  591. gearman_job_free(worker->job);
  592. worker->job= NULL;
  593. *ret_ptr= GEARMAN_UNEXPECTED_PACKET;
  594. return NULL;
  595. }
  596. gearman_packet_free(&(worker->job->assigned));
  597. }
  598. }
  599. case GEARMAN_WORKER_STATE_PRE_SLEEP:
  600. for (worker->con= (&worker->universal)->con_list; worker->con;
  601. worker->con= worker->con->next)
  602. {
  603. if (worker->con->fd == INVALID_SOCKET)
  604. {
  605. continue;
  606. }
  607. *ret_ptr= worker->con->send(worker->pre_sleep, true);
  608. if (gearman_failed(*ret_ptr))
  609. {
  610. if (*ret_ptr == GEARMAN_IO_WAIT)
  611. {
  612. worker->state= GEARMAN_WORKER_STATE_PRE_SLEEP;
  613. }
  614. else if (*ret_ptr == GEARMAN_LOST_CONNECTION)
  615. {
  616. continue;
  617. }
  618. return NULL;
  619. }
  620. }
  621. worker->state= GEARMAN_WORKER_STATE_START;
  622. /* Set a watch on all active connections that we sent a PRE_SLEEP to. */
  623. active= 0;
  624. for (worker->con= worker->universal.con_list; worker->con;
  625. worker->con= worker->con->next)
  626. {
  627. if (worker->con->fd == INVALID_SOCKET)
  628. continue;
  629. worker->con->set_events(POLLIN);
  630. active++;
  631. }
  632. if ((&worker->universal)->options.non_blocking)
  633. {
  634. *ret_ptr= GEARMAN_NO_JOBS;
  635. return NULL;
  636. }
  637. if (active == 0)
  638. {
  639. if (worker->universal.timeout < 0)
  640. {
  641. gearman_nap(GEARMAN_WORKER_WAIT_TIMEOUT);
  642. }
  643. else
  644. {
  645. if (worker->universal.timeout > 0)
  646. {
  647. gearman_nap(worker->universal);
  648. }
  649. if (worker->options.timeout_return)
  650. {
  651. gearman_error(worker->universal, GEARMAN_TIMEOUT, "Option timeout return reached");
  652. *ret_ptr= GEARMAN_TIMEOUT;
  653. return NULL;
  654. }
  655. }
  656. }
  657. else
  658. {
  659. *ret_ptr= gearman_wait(worker->universal);
  660. if (gearman_failed(*ret_ptr) and (*ret_ptr != GEARMAN_TIMEOUT or worker->options.timeout_return))
  661. {
  662. return NULL;
  663. }
  664. }
  665. break;
  666. }
  667. }
  668. }
  669. void gearman_job_free_all(gearman_worker_st *worker)
  670. {
  671. while (worker->job_list)
  672. gearman_job_free(worker->job_list);
  673. }
  674. gearman_return_t gearman_worker_add_function(gearman_worker_st *worker,
  675. const char *function_name,
  676. uint32_t timeout,
  677. gearman_worker_fn *worker_fn,
  678. void *context)
  679. {
  680. if (not worker)
  681. return GEARMAN_INVALID_ARGUMENT;
  682. if (not function_name)
  683. {
  684. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function name not given");
  685. }
  686. if (not worker_fn)
  687. {
  688. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function not given");
  689. }
  690. gearman_function_t local= gearman_function_create_v1(worker_fn);
  691. return _worker_function_create(worker,
  692. function_name, strlen(function_name),
  693. local,
  694. timeout,
  695. context);
  696. }
  697. gearman_return_t gearman_worker_define_function(gearman_worker_st *worker,
  698. const char *function_name, const size_t function_name_length,
  699. const gearman_function_t function,
  700. const uint32_t timeout,
  701. void *context)
  702. {
  703. if (not worker)
  704. {
  705. return GEARMAN_INVALID_ARGUMENT;
  706. }
  707. if (not function_name or function_name_length == 0)
  708. {
  709. return gearman_error(worker->universal, GEARMAN_INVALID_ARGUMENT, "function name not given");
  710. }
  711. return _worker_function_create(worker,
  712. function_name, function_name_length,
  713. function,
  714. timeout,
  715. context);
  716. return GEARMAN_INVALID_ARGUMENT;
  717. }
  718. gearman_return_t gearman_worker_work(gearman_worker_st *worker)
  719. {
  720. bool shutdown= false;
  721. if (not worker)
  722. {
  723. return GEARMAN_INVALID_ARGUMENT;
  724. }
  725. switch (worker->work_state)
  726. {
  727. case GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB:
  728. {
  729. gearman_return_t ret;
  730. worker->work_job= gearman_worker_grab_job(worker, NULL, &ret);
  731. if (gearman_failed(ret))
  732. {
  733. if (ret == GEARMAN_COULD_NOT_CONNECT)
  734. {
  735. gearman_reset(worker->universal);
  736. }
  737. return ret;
  738. }
  739. assert(worker->work_job);
  740. for (worker->work_function= worker->function_list;
  741. worker->work_function;
  742. worker->work_function= worker->work_function->next)
  743. {
  744. if (not strcmp(gearman_job_function_name(worker->work_job),
  745. worker->work_function->function_name))
  746. {
  747. break;
  748. }
  749. }
  750. if (not worker->work_function)
  751. {
  752. gearman_job_free(worker->work_job);
  753. worker->work_job= NULL;
  754. return gearman_error(worker->universal, GEARMAN_INVALID_FUNCTION_NAME, "Function not found");
  755. }
  756. if (not worker->work_function->has_callback())
  757. {
  758. gearman_job_free(worker->work_job);
  759. worker->work_job= NULL;
  760. return gearman_error(worker->universal, GEARMAN_INVALID_FUNCTION_NAME, "Neither a gearman_worker_fn, or gearman_function_fn callback was supplied");
  761. }
  762. worker->work_result_size= 0;
  763. }
  764. case GEARMAN_WORKER_WORK_UNIVERSAL_FUNCTION:
  765. {
  766. switch (worker->work_function->callback(worker->work_job,
  767. static_cast<void *>(worker->work_function->context)))
  768. {
  769. case GEARMAN_FUNCTION_INVALID_ARGUMENT:
  770. worker->work_job->error_code= GEARMAN_INVALID_ARGUMENT;
  771. case GEARMAN_FUNCTION_FATAL:
  772. 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
  773. {
  774. worker->work_job->error_code= GEARMAN_LOST_CONNECTION;
  775. break;
  776. }
  777. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_FAIL;
  778. return worker->work_job->error_code;
  779. case GEARMAN_FUNCTION_ERROR: // retry
  780. worker->work_job->error_code= GEARMAN_LOST_CONNECTION;
  781. break;
  782. case GEARMAN_FUNCTION_SHUTDOWN:
  783. shutdown= true;
  784. case GEARMAN_FUNCTION_SUCCESS:
  785. break;
  786. }
  787. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  788. break;
  789. }
  790. case GEARMAN_WORKER_WORK_UNIVERSAL_COMPLETE:
  791. {
  792. worker->work_job->error_code= gearman_job_send_complete_fin(worker->work_job,
  793. worker->work_result, worker->work_result_size);
  794. if (worker->work_job->error_code == GEARMAN_IO_WAIT)
  795. {
  796. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_COMPLETE;
  797. return gearman_error(worker->universal, worker->work_job->error_code,
  798. "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.");
  799. }
  800. if (worker->work_result)
  801. {
  802. gearman_free(worker->universal, worker->work_result);
  803. worker->work_result= NULL;
  804. }
  805. // If we lost the connection, we retry the work, otherwise we error
  806. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  807. {
  808. break;
  809. }
  810. else if (gearman_failed(worker->work_job->error_code))
  811. {
  812. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_FAIL;
  813. return worker->work_job->error_code;
  814. }
  815. }
  816. break;
  817. case GEARMAN_WORKER_WORK_UNIVERSAL_FAIL:
  818. {
  819. if (gearman_failed(worker->work_job->error_code= gearman_job_send_fail_fin(worker->work_job)))
  820. {
  821. if (worker->work_job->error_code == GEARMAN_LOST_CONNECTION)
  822. {
  823. break;
  824. }
  825. return worker->work_job->error_code;
  826. }
  827. }
  828. break;
  829. }
  830. gearman_job_free(worker->work_job);
  831. worker->work_job= NULL;
  832. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB;
  833. return shutdown ? GEARMAN_SHUTDOWN : GEARMAN_SUCCESS;
  834. }
  835. gearman_return_t gearman_worker_echo(gearman_worker_st *worker,
  836. const void *workload,
  837. size_t workload_size)
  838. {
  839. if (not worker)
  840. return GEARMAN_INVALID_ARGUMENT;
  841. return gearman_echo(worker->universal, workload, workload_size);
  842. }
  843. /*
  844. * Static Definitions
  845. */
  846. static gearman_worker_st *_worker_allocate(gearman_worker_st *worker, bool is_clone)
  847. {
  848. if (not worker)
  849. {
  850. worker= new (std::nothrow) gearman_worker_st;
  851. if (worker == NULL)
  852. {
  853. gearman_perror(worker->universal, "gearman_worker_st new");
  854. return NULL;
  855. }
  856. worker->options.allocated= true;
  857. }
  858. else
  859. {
  860. worker->options.allocated= false;
  861. }
  862. worker->options.non_blocking= false;
  863. worker->options.packet_init= false;
  864. worker->options.change= false;
  865. worker->options.grab_uniq= true;
  866. worker->options.grab_all= true;
  867. worker->options.timeout_return= false;
  868. worker->state= GEARMAN_WORKER_STATE_START;
  869. worker->work_state= GEARMAN_WORKER_WORK_UNIVERSAL_GRAB_JOB;
  870. worker->function_count= 0;
  871. worker->job_count= 0;
  872. worker->work_result_size= 0;
  873. worker->context= NULL;
  874. worker->con= NULL;
  875. worker->job= NULL;
  876. worker->job_list= NULL;
  877. worker->function= NULL;
  878. worker->function_list= NULL;
  879. worker->work_function= NULL;
  880. worker->work_result= NULL;
  881. if (not is_clone)
  882. {
  883. gearman_universal_initialize(worker->universal);
  884. #if 0
  885. gearman_universal_set_timeout(worker->universal, GEARMAN_WORKER_WAIT_TIMEOUT);
  886. #endif
  887. }
  888. return worker;
  889. }
  890. static gearman_return_t _worker_packet_init(gearman_worker_st *worker)
  891. {
  892. gearman_return_t ret= gearman_packet_create_args(worker->universal, worker->grab_job,
  893. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_GRAB_JOB_ALL,
  894. NULL, NULL, 0);
  895. if (gearman_failed(ret))
  896. return ret;
  897. ret= gearman_packet_create_args(worker->universal, worker->pre_sleep,
  898. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_PRE_SLEEP,
  899. NULL, NULL, 0);
  900. if (gearman_failed(ret))
  901. {
  902. gearman_packet_free(&(worker->grab_job));
  903. return ret;
  904. }
  905. worker->options.packet_init= true;
  906. return GEARMAN_SUCCESS;
  907. }
  908. static gearman_return_t _worker_add_server(const char *host, in_port_t port, void *context)
  909. {
  910. return gearman_worker_add_server(static_cast<gearman_worker_st *>(context), host, port);
  911. }
  912. static gearman_return_t _worker_function_create(gearman_worker_st *worker,
  913. const char *function_name, size_t function_length,
  914. const gearman_function_t &function_arg,
  915. uint32_t timeout,
  916. void *context)
  917. {
  918. const void *args[2];
  919. size_t args_size[2];
  920. _worker_function_st *function= make(worker->universal._namespace, function_name, function_length, function_arg, context);
  921. if (not function)
  922. {
  923. gearman_perror(worker->universal, "_worker_function_st::new()");
  924. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  925. }
  926. gearman_return_t ret;
  927. if (timeout > 0)
  928. {
  929. char timeout_buffer[11];
  930. snprintf(timeout_buffer, sizeof(timeout_buffer), "%u", timeout);
  931. args[0]= function->name();
  932. args_size[0]= function->length() + 1;
  933. args[1]= timeout_buffer;
  934. args_size[1]= strlen(timeout_buffer);
  935. ret= gearman_packet_create_args(worker->universal, function->packet,
  936. GEARMAN_MAGIC_REQUEST,
  937. GEARMAN_COMMAND_CAN_DO_TIMEOUT,
  938. args, args_size, 2);
  939. }
  940. else
  941. {
  942. args[0]= function->name();
  943. args_size[0]= function->length();
  944. ret= gearman_packet_create_args(worker->universal, function->packet,
  945. GEARMAN_MAGIC_REQUEST, GEARMAN_COMMAND_CAN_DO,
  946. args, args_size, 1);
  947. }
  948. if (gearman_failed(ret))
  949. {
  950. delete function;
  951. return ret;
  952. }
  953. if (worker->function_list)
  954. worker->function_list->prev= function;
  955. function->next= worker->function_list;
  956. function->prev= NULL;
  957. worker->function_list= function;
  958. worker->function_count++;
  959. worker->options.change= true;
  960. return GEARMAN_SUCCESS;
  961. }
  962. static void _worker_function_free(gearman_worker_st *worker,
  963. struct _worker_function_st *function)
  964. {
  965. if (worker->function_list == function)
  966. worker->function_list= function->next;
  967. if (function->prev)
  968. function->prev->next= function->next;
  969. if (function->next)
  970. function->next->prev= function->prev;
  971. worker->function_count--;
  972. delete function;
  973. }
  974. gearman_return_t gearman_worker_set_memory_allocators(gearman_worker_st *worker,
  975. gearman_malloc_fn *malloc_fn,
  976. gearman_free_fn *free_fn,
  977. gearman_realloc_fn *realloc_fn,
  978. gearman_calloc_fn *calloc_fn,
  979. void *context)
  980. {
  981. if (not worker)
  982. return GEARMAN_INVALID_ARGUMENT;
  983. return gearman_set_memory_allocator(worker->universal.allocator, malloc_fn, free_fn, realloc_fn, calloc_fn, context);
  984. }
  985. bool gearman_worker_set_server_option(gearman_worker_st *self, const char *option_arg, size_t option_arg_size)
  986. {
  987. gearman_string_t option= { option_arg, option_arg_size };
  988. return gearman_request_option(self->universal, option);
  989. }
  990. void gearman_worker_set_namespace(gearman_worker_st *self, const char *namespace_key, size_t namespace_key_size)
  991. {
  992. if (not self)
  993. return;
  994. gearman_universal_set_namespace(self->universal, namespace_key, namespace_key_size);
  995. }